use std::collections::{BTreeMap, BTreeSet};
use std::fs;
use std::path::{Path, PathBuf};
use std::process::Command;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::render::to_json;
use crate::{EXIT_FAILURE, EXIT_OK, EXIT_USAGE, Rendered, evidence, flows, run};
const CURSOR_SUFFIX: &str = ".cursor.json";
const EVENTS_SUBDIR: &str = "events";
const EVENTS_EXT: &str = "ndjson";
const SIM_CRASH_EXIT_CODE: i32 = 137;
const DEFAULT_MAX_RESTARTS: u32 = 8;
const LEASE_GRACE_MS: u64 = 200;
#[derive(Debug, Clone, Deserialize)]
struct SimPlan {
target: String,
#[serde(default)]
args: Vec<String>,
#[serde(default = "default_max_restarts")]
max_restarts: u32,
#[serde(default)]
flow_lease_ms: Option<u64>,
#[serde(default)]
assert: SimAssertions,
}
const fn default_max_restarts() -> u32 {
DEFAULT_MAX_RESTARTS
}
#[derive(Debug, Clone, Default, Deserialize)]
struct SimAssertions {
#[serde(default)]
max_attempts: BTreeMap<String, u32>,
#[serde(default)]
breaker_open: Vec<String>,
#[serde(default)]
no_breaker_open: Vec<String>,
#[serde(default)]
flow_status: Option<String>,
}
#[derive(Debug, Serialize)]
struct SimFinding {
detail: String,
level: &'static str,
topic: &'static str,
}
#[derive(Debug, Serialize)]
struct SimReport {
exit_code: i32,
findings: Vec<SimFinding>,
ok: bool,
plan: String,
restarts: u32,
}
pub fn run(project: &Path, plan_path: &str) -> Rendered {
run_with(project, plan_path, &|_cmd| {})
}
pub fn run_with(project: &Path, plan_path: &str, configure: &dyn Fn(&mut Command)) -> Rendered {
let path = Path::new(plan_path);
let text = match fs::read_to_string(path) {
Ok(t) => t,
Err(err) => return soft_error(&format!("could not read {plan_path}: {err}.")),
};
let plan: SimPlan = match serde_json::from_str(&text) {
Ok(p) => p,
Err(err) => {
return soft_error(&format!(
"{plan_path} is not a valid sim plan (docs/sim-format.md): {err}."
));
}
};
let abs_plan = fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf());
let _ = fs::remove_file(cursor_path_for(&abs_plan));
let run_plan = match run::plan(&plan.target, &plan.args, false) {
Ok(p) => p,
Err(e) => return e.render(),
};
let events_dir = project.join(".keel").join(EVENTS_SUBDIR);
let before = existing_event_files(&events_dir);
let result = drive(
&run_plan,
&abs_plan,
plan.max_restarts,
plan.flow_lease_ms,
configure,
);
let observed = scan_events(&new_event_files(&events_dir, &before));
let mut findings = Vec::new();
check_assertions(&plan.assert, &observed, project, &mut findings);
if result.crashed {
findings.push(SimFinding {
level: "error",
topic: "crash-restart",
detail: format!(
"the target kept crashing after {} restart(s) (max_restarts={}); it never \
reached a terminal state.",
result.restarts, plan.max_restarts
),
});
}
let report = SimReport {
exit_code: result.exit_code,
ok: findings.is_empty(),
plan: plan_path.to_owned(),
restarts: result.restarts,
findings,
};
let exit = if report.ok { EXIT_OK } else { EXIT_USAGE };
let human = human(&report);
Rendered::ok(human, to_json(&report)).with_exit(exit)
}
fn cursor_path_for(plan_path: &Path) -> PathBuf {
let mut s = plan_path.as_os_str().to_owned();
s.push(CURSOR_SUFFIX);
PathBuf::from(s)
}
struct DriveResult {
exit_code: i32,
restarts: u32,
crashed: bool,
}
fn drive(
plan: &run::RunPlan,
abs_plan_path: &Path,
max_restarts: u32,
flow_lease_ms: Option<u64>,
configure: &dyn Fn(&mut Command),
) -> DriveResult {
let mut restarts = 0u32;
loop {
let mut cmd = Command::new(&plan.program);
cmd.args(&plan.argv);
if plan.disable {
cmd.env("KEEL_DISABLE", "1");
}
cmd.env("KEEL_SIM_PLAN", abs_plan_path);
cmd.env("KEEL_EVENTS", "1");
if let Some(lease_ms) = flow_lease_ms {
cmd.env("KEEL_FLOW_LEASE_MS", lease_ms.to_string());
}
configure(&mut cmd);
let Ok(status) = cmd.status() else {
return DriveResult {
exit_code: EXIT_FAILURE,
restarts,
crashed: false,
};
};
if !crashed(status) {
return DriveResult {
exit_code: status.code().unwrap_or(EXIT_FAILURE),
restarts,
crashed: false,
};
}
if restarts >= max_restarts {
return DriveResult {
exit_code: SIM_CRASH_EXIT_CODE,
restarts,
crashed: true,
};
}
restarts += 1;
if let Some(lease_ms) = flow_lease_ms {
std::thread::sleep(std::time::Duration::from_millis(lease_ms + LEASE_GRACE_MS));
}
}
}
#[cfg(unix)]
fn crashed(status: std::process::ExitStatus) -> bool {
use std::os::unix::process::ExitStatusExt;
status.signal().is_some()
}
#[cfg(not(unix))]
fn crashed(_status: std::process::ExitStatus) -> bool {
false
}
fn existing_event_files(dir: &Path) -> BTreeSet<String> {
fs::read_dir(dir)
.into_iter()
.flatten()
.flatten()
.filter_map(|e| e.file_name().into_string().ok())
.collect()
}
fn new_event_files(dir: &Path, before: &BTreeSet<String>) -> Vec<PathBuf> {
let mut names: Vec<String> = fs::read_dir(dir)
.into_iter()
.flatten()
.flatten()
.filter_map(|e| e.file_name().into_string().ok())
.filter(|n| n.ends_with(&format!(".{EVENTS_EXT}")) && !before.contains(n))
.collect();
names.sort();
names.into_iter().map(|n| dir.join(n)).collect()
}
#[derive(Debug, Default)]
struct Observed {
max_attempts: BTreeMap<String, u32>,
breaker_open: BTreeSet<String>,
saw_any_event: bool,
}
fn scan_events(files: &[PathBuf]) -> Observed {
let mut out = Observed::default();
for path in files {
let Ok(text) = fs::read_to_string(path) else {
continue;
};
for line in text.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
let Ok(v) = serde_json::from_str::<Value>(line) else {
continue;
};
match v.get("event").and_then(Value::as_str) {
Some("call_end") => {
out.saw_any_event = true;
let Some(target) = v.get("target").and_then(Value::as_str) else {
continue;
};
let attempts =
u32::try_from(v.get("attempts").and_then(Value::as_u64).unwrap_or(0))
.unwrap_or(u32::MAX);
let entry = out.max_attempts.entry(target.to_owned()).or_insert(0);
if attempts > *entry {
*entry = attempts;
}
}
Some("breaker_open") => {
out.saw_any_event = true;
if let Some(target) = v.get("target").and_then(Value::as_str) {
out.breaker_open.insert(target.to_owned());
}
}
_ => {}
}
}
}
out
}
fn check_assertions(
assert: &SimAssertions,
observed: &Observed,
project: &Path,
findings: &mut Vec<SimFinding>,
) {
let wants_events = !assert.max_attempts.is_empty() || !assert.breaker_open.is_empty();
if wants_events && !observed.saw_any_event {
findings.push(SimFinding {
level: "error",
topic: "no-events",
detail: "max_attempts/breaker_open assertions were requested, but no Tier 1 events \
were observed at all. The event sink is a native-core-only feature \
(crates/keel-core/src/events.rs) — set KEEL_BACKEND=native (a built \
keel_core), or drop these assertions for a pure stub/dev-backend sim."
.to_owned(),
});
}
for (target, cap) in &assert.max_attempts {
if let Some(&seen) = observed.max_attempts.get(target)
&& seen > *cap
{
findings.push(SimFinding {
level: "error",
topic: "max-attempts",
detail: format!(
"{target} made {seen} attempt(s) on one call, exceeding the configured cap \
of {cap}."
),
});
}
}
for target in &assert.breaker_open {
if !observed.breaker_open.contains(target) {
findings.push(SimFinding {
level: "error",
topic: "breaker",
detail: format!(
"expected the breaker for {target} to open under the fault plan, but it \
never did."
),
});
}
}
for target in &assert.no_breaker_open {
if observed.breaker_open.contains(target) {
findings.push(SimFinding {
level: "error",
topic: "breaker",
detail: format!(
"the breaker for {target} opened, but the plan asserted it must not."
),
});
}
}
if let Some(want) = &assert.flow_status {
match newest_flow_status(project) {
Some(got) if &got == want => {}
Some(got) => findings.push(SimFinding {
level: "error",
topic: "flow-status",
detail: format!("expected the flow to end {want}, but it is {got}."),
}),
None => findings.push(SimFinding {
level: "error",
topic: "flow-status",
detail: format!(
"expected the flow to end {want}, but no flow was found in \
.keel/journal.db."
),
}),
}
}
}
fn newest_flow_status(project: &Path) -> Option<String> {
let path = evidence::resolved_journal(project).path;
if !path.exists() {
return None;
}
let conn = flows::open_ro(&path).ok()?;
conn.query_row(
"SELECT status FROM flows ORDER BY updated_at DESC, flow_id DESC LIMIT 1",
[],
|row| row.get::<_, String>(0),
)
.ok()
}
fn human(report: &SimReport) -> String {
let restarts_note = if report.restarts > 0 {
format!(", {} restart(s)", report.restarts)
} else {
String::new()
};
if report.ok {
return format!(
"keel \u{25b8} sim {} passed \u{2014} exit {}{restarts_note}.",
report.plan, report.exit_code
);
}
let mut lines = vec![format!(
"keel \u{25b8} sim {} found {} problem(s) (exit {}{restarts_note}):\n",
report.plan,
report.findings.len(),
report.exit_code
)];
for f in &report.findings {
lines.push(format!(" [{}] {}: {}\n", f.level, f.topic, f.detail));
}
lines.concat()
}
fn soft_error(message: &str) -> Rendered {
#[derive(Serialize)]
struct ErrReport<'a> {
error: &'a str,
}
Rendered {
human: format!("keel \u{25b8} {message}"),
json: to_json(&ErrReport { error: message }),
exit: EXIT_FAILURE,
to_stderr: true,
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write as _;
use std::os::unix::fs::PermissionsExt;
use tempfile::TempDir;
fn project() -> TempDir {
TempDir::new().unwrap()
}
fn write_plan(dir: &Path, name: &str, json: &str) -> PathBuf {
let path = dir.join(name);
fs::write(&path, json).unwrap();
path
}
fn write_script(dir: &Path, name: &str, body: &str) -> PathBuf {
let path = dir.join(name);
fs::write(&path, format!("#!/bin/sh\n{body}\n")).unwrap();
let mut perms = fs::metadata(&path).unwrap().permissions();
perms.set_mode(0o755);
fs::set_permissions(&path, perms).unwrap();
path
}
#[test]
fn missing_plan_file_is_a_precise_error() {
let dir = project();
let r = run(dir.path(), "does-not-exist.json");
assert_eq!(r.exit, EXIT_FAILURE);
assert!(r.human.contains("could not read"));
}
#[test]
fn malformed_plan_json_is_a_precise_error() {
let dir = project();
let plan = write_plan(dir.path(), "plan.json", "not json");
let r = run(dir.path(), &plan.to_string_lossy());
assert_eq!(r.exit, EXIT_FAILURE);
assert!(r.human.contains("not a valid sim plan"));
}
#[test]
fn unresolvable_target_reports_the_same_error_keel_run_would() {
let dir = project();
let plan = write_plan(
dir.path(),
"plan.json",
r#"{"v":1,"target":"does-not-exist.py"}"#,
);
let r = run(dir.path(), &plan.to_string_lossy());
assert_eq!(r.exit, EXIT_USAGE);
assert!(r.human.contains("no such file or directory"));
}
#[test]
fn clean_run_with_no_assertions_passes() {
let dir = project();
write_script(dir.path(), "ok.sh", "exit 0");
let plan = write_plan(
dir.path(),
"plan.json",
&format!(
r#"{{"v":1,"target":{:?}}}"#,
dir.path().join("ok.sh").to_string_lossy()
),
);
let run_plan = run::RunPlan {
program: dir.path().join("ok.sh").to_string_lossy().into_owned(),
argv: vec![],
disable: false,
};
let result = drive(&run_plan, &plan, 4, None, &|_cmd| {});
assert_eq!(result.exit_code, 0);
assert_eq!(result.restarts, 0);
assert!(!result.crashed);
}
#[test]
fn a_self_killed_child_is_retried_up_to_max_restarts() {
let dir = project();
let counter = dir.path().join("count");
fs::write(&counter, "0").unwrap();
let script = write_script(
dir.path(),
"flaky.sh",
&format!(
"n=$(cat {0})\nn=$((n+1))\necho $n > {0}\nif [ \"$n\" -le 2 ]; then kill -9 $$; fi\nexit 0",
counter.display()
),
);
let plan = write_plan(dir.path(), "plan.json", r#"{"v":1,"target":"x"}"#);
let run_plan = run::RunPlan {
program: script.to_string_lossy().into_owned(),
argv: vec![],
disable: false,
};
let result = drive(&run_plan, &plan, 4, None, &|_cmd| {});
assert_eq!(result.exit_code, 0);
assert_eq!(result.restarts, 2);
assert!(!result.crashed);
}
#[test]
fn exhausting_max_restarts_while_still_crashing_is_flagged() {
let dir = project();
let script = write_script(dir.path(), "always_dies.sh", "kill -9 $$");
let plan = write_plan(dir.path(), "plan.json", r#"{"v":1,"target":"x"}"#);
let run_plan = run::RunPlan {
program: script.to_string_lossy().into_owned(),
argv: vec![],
disable: false,
};
let result = drive(&run_plan, &plan, 2, None, &|_cmd| {});
assert_eq!(result.restarts, 2);
assert!(result.crashed);
assert_eq!(result.exit_code, SIM_CRASH_EXIT_CODE);
}
#[test]
fn cursor_sidecar_is_reset_before_the_first_spawn() {
let dir = project();
write_script(dir.path(), "ok.sh", "exit 0");
let target = dir.path().join("ok.sh").to_string_lossy().into_owned();
let plan_path = write_plan(
dir.path(),
"plan.json",
&format!(r#"{{"v":1,"target":{target:?}}}"#),
);
let abs = fs::canonicalize(&plan_path).unwrap();
let cursor = cursor_path_for(&abs);
fs::write(&cursor, r#"{"stale":"leftover"}"#).unwrap();
assert!(cursor.exists());
let r = run(dir.path(), &plan_path.to_string_lossy());
assert_eq!(r.exit, EXIT_USAGE, "{r:?}"); assert!(!cursor.exists(), "stale cursor must be cleared regardless");
}
#[test]
fn scan_events_finds_max_attempts_and_breaker_open() {
let dir = project();
let events = dir.path().join("run.ndjson");
let mut f = fs::File::create(&events).unwrap();
writeln!(
f,
r#"{{"v":1,"seq":0,"ms":0,"event":"run_start","run":"r"}}"#
)
.unwrap();
writeln!(
f,
r#"{{"v":1,"seq":1,"ms":1,"event":"call_end","call":"t-1","target":"api.a","result":"ok","attempts":3}}"#
)
.unwrap();
writeln!(
f,
r#"{{"v":1,"seq":2,"ms":2,"event":"breaker_open","call":"t-2","target":"api.b","cooldown_ms":500}}"#
)
.unwrap();
let observed = scan_events(&[events]);
assert_eq!(observed.max_attempts.get("api.a"), Some(&3));
assert!(observed.breaker_open.contains("api.b"));
}
#[test]
fn max_attempts_violation_is_a_finding() {
let observed = Observed {
max_attempts: BTreeMap::from([("api.a".to_owned(), 5)]),
breaker_open: BTreeSet::new(),
saw_any_event: true,
};
let assert = SimAssertions {
max_attempts: BTreeMap::from([("api.a".to_owned(), 3)]),
..Default::default()
};
let mut findings = Vec::new();
check_assertions(&assert, &observed, Path::new("."), &mut findings);
assert_eq!(findings.len(), 1, "{findings:?}");
assert_eq!(findings[0].topic, "max-attempts");
assert!(findings[0].detail.contains("5 attempt"));
}
#[test]
fn breaker_open_expected_but_absent_is_a_finding() {
let observed = Observed {
saw_any_event: true,
..Default::default()
};
let assert = SimAssertions {
breaker_open: vec!["api.flaky".to_owned()],
..Default::default()
};
let mut findings = Vec::new();
check_assertions(&assert, &observed, Path::new("."), &mut findings);
assert_eq!(findings.len(), 1, "{findings:?}");
assert_eq!(findings[0].topic, "breaker");
}
#[test]
fn no_breaker_open_violated_is_a_finding() {
let observed = Observed {
max_attempts: BTreeMap::new(),
breaker_open: BTreeSet::from(["api.pay".to_owned()]),
saw_any_event: true,
};
let assert = SimAssertions {
no_breaker_open: vec!["api.pay".to_owned()],
..Default::default()
};
let mut findings = Vec::new();
check_assertions(&assert, &observed, Path::new("."), &mut findings);
assert_eq!(findings.len(), 1, "{findings:?}");
}
#[test]
fn no_events_at_all_is_a_finding_when_event_based_assertions_were_requested() {
let observed = Observed::default();
let assert = SimAssertions {
max_attempts: BTreeMap::from([("api.a".to_owned(), 3)]),
..Default::default()
};
let mut findings = Vec::new();
check_assertions(&assert, &observed, Path::new("."), &mut findings);
assert_eq!(findings.len(), 1, "{findings:?}");
assert_eq!(findings[0].topic, "no-events");
}
#[test]
fn no_breaker_open_assertion_alone_never_needs_events() {
let observed = Observed::default();
let assert = SimAssertions {
no_breaker_open: vec!["api.pay".to_owned()],
..Default::default()
};
let mut findings = Vec::new();
check_assertions(&assert, &observed, Path::new("."), &mut findings);
assert!(findings.is_empty(), "{findings:?}");
}
#[test]
fn flow_status_mismatch_and_missing_are_findings() {
let dir = project();
let observed = Observed::default();
let assert = SimAssertions {
flow_status: Some("completed".to_owned()),
..Default::default()
};
let mut findings = Vec::new();
check_assertions(&assert, &observed, dir.path(), &mut findings);
assert_eq!(findings.len(), 1);
assert!(findings[0].detail.contains("no flow was found"));
}
fn python3_present() -> bool {
Command::new("python3")
.arg("--version")
.output()
.is_ok_and(|o| o.status.success())
}
fn python_path() -> String {
let manifest = Path::new(env!("CARGO_MANIFEST_DIR"));
format!(
"{}:{}",
manifest.join("../../python/keel/src").display(),
manifest.join("../../python/keel-core-stub").display(),
)
}
#[test]
fn end_to_end_report_is_ok_for_a_clean_run() {
if !python3_present() {
eprintln!("skip: python3 not available");
return;
}
let dir = project();
fs::write(dir.path().join("app.py"), "print(\"hi\")\n").unwrap();
let target = dir.path().join("app.py").to_string_lossy().into_owned();
let plan_path = write_plan(
dir.path(),
"plan.json",
&format!(r#"{{"v":1,"target":{target:?}}}"#),
);
let pythonpath = python_path();
let r = run_with(dir.path(), &plan_path.to_string_lossy(), &move |cmd| {
cmd.env("PYTHONPATH", &pythonpath);
cmd.env("KEEL_BACKEND", "stub");
cmd.env("KEEL_QUIET", "1");
});
assert_eq!(r.exit, EXIT_OK, "{r:?}");
assert_eq!(r.json["ok"], true);
assert_eq!(r.json["exit_code"], 0);
assert_eq!(r.json["restarts"], 0);
assert!(r.human.contains("passed"));
}
fn write_wrappable_target(dir: &Path) -> String {
fs::write(dir.join("keel.toml"), "[target.\"py:lib.call_api\"]\n").unwrap();
fs::write(
dir.join("lib.py"),
"import os\n\nCALLS_FILE = os.path.join(os.path.dirname(__file__), \"calls.txt\")\n\n\ndef call_api():\n with open(CALLS_FILE, \"a\", encoding=\"utf-8\") as f:\n f.write(\"call\\n\")\n return {\"ok\": True}\n",
)
.unwrap();
fs::write(
dir.join("app.py"),
"import lib\n\n\ndef main():\n print(lib.call_api())\n\n\nif __name__ == \"__main__\":\n main()\n",
)
.unwrap();
dir.join("app.py").to_string_lossy().into_owned()
}
fn calls_count(dir: &Path) -> usize {
fs::read_to_string(dir.join("calls.txt"))
.map_or(0, |t| t.lines().filter(|l| !l.trim().is_empty()).count())
}
#[test]
fn front_end_fault_injection_exhausts_real_retries_and_fails() {
let bin_dir = venv_bin_dir();
if !python3_present() || !native_core_present(bin_dir.as_deref()) {
eprintln!("skip: native core (keel_core) not available");
return;
}
let dir = project();
let target = write_wrappable_target(dir.path());
let plan_path = write_plan(
dir.path(),
"plan.json",
&format!(
r#"{{"v":1,"target":{target:?},"faults":{{"py:lib.call_api":[{{"kind":"timeout"}},{{"kind":"timeout"}},{{"kind":"timeout"}},{{"kind":"timeout"}}]}},"assert":{{"max_attempts":{{"py:lib.call_api":3}}}}}}"#
),
);
let pythonpath = python_path();
let path_env = bin_dir.map_or_else(
|| std::env::var("PATH").unwrap_or_default(),
|d| {
format!(
"{}:{}",
d.display(),
std::env::var("PATH").unwrap_or_default()
)
},
);
let project_dir = dir.path().to_path_buf();
let r = run_with(dir.path(), &plan_path.to_string_lossy(), &move |cmd| {
cmd.current_dir(&project_dir);
cmd.env("PATH", &path_env);
cmd.env("PYTHONPATH", &pythonpath);
cmd.env("KEEL_BACKEND", "native");
cmd.env("KEEL_QUIET", "1");
});
assert_ne!(r.json["exit_code"], 0, "{r:?}");
assert_eq!(calls_count(dir.path()), 0, "the real effect must never run");
assert_eq!(r.json["ok"], true, "{r:?}"); }
#[test]
fn front_end_fault_injection_absorbs_retries_then_succeeds() {
let bin_dir = venv_bin_dir();
if !python3_present() || !native_core_present(bin_dir.as_deref()) {
eprintln!("skip: native core (keel_core) not available");
return;
}
let dir = project();
let target = write_wrappable_target(dir.path());
let plan_path = write_plan(
dir.path(),
"plan.json",
&format!(
r#"{{"v":1,"target":{target:?},"faults":{{"py:lib.call_api":[{{"kind":"timeout"}},{{"kind":"5xx","status":503}},{{"kind":"ok"}}]}},"assert":{{"max_attempts":{{"py:lib.call_api":3}}}}}}"#
),
);
let pythonpath = python_path();
let path_env = bin_dir.map_or_else(
|| std::env::var("PATH").unwrap_or_default(),
|d| {
format!(
"{}:{}",
d.display(),
std::env::var("PATH").unwrap_or_default()
)
},
);
let project_dir = dir.path().to_path_buf();
let r = run_with(dir.path(), &plan_path.to_string_lossy(), &move |cmd| {
cmd.current_dir(&project_dir);
cmd.env("PATH", &path_env);
cmd.env("PYTHONPATH", &pythonpath);
cmd.env("KEEL_BACKEND", "native");
cmd.env("KEEL_QUIET", "1");
});
assert_eq!(r.exit, EXIT_OK, "{r:?}");
assert_eq!(r.json["ok"], true, "{r:?}");
assert_eq!(r.json["exit_code"], 0);
assert_eq!(
calls_count(dir.path()),
1,
"the real effect ran exactly once"
);
}
#[test]
fn front_end_fault_injection_exhausts_real_retries_on_the_stub() {
if !python3_present() {
eprintln!("skip: python3 not available");
return;
}
let dir = project();
let target = write_wrappable_target(dir.path());
let plan_path = write_plan(
dir.path(),
"plan.json",
&format!(
r#"{{"v":1,"target":{target:?},"faults":{{"py:lib.call_api":[{{"kind":"timeout"}},{{"kind":"timeout"}},{{"kind":"timeout"}},{{"kind":"timeout"}}]}}}}"#
),
);
let pythonpath = python_path();
let project_dir = dir.path().to_path_buf();
let r = run_with(dir.path(), &plan_path.to_string_lossy(), &move |cmd| {
cmd.current_dir(&project_dir);
cmd.env("PYTHONPATH", &pythonpath);
cmd.env("KEEL_BACKEND", "stub");
cmd.env("KEEL_QUIET", "1");
});
assert_ne!(r.json["exit_code"], 0, "{r:?}");
assert_eq!(calls_count(dir.path()), 0, "the real effect must never run");
}
fn venv_bin_dir() -> Option<PathBuf> {
let out = Command::new("git")
.args(["rev-parse", "--git-common-dir"])
.current_dir(env!("CARGO_MANIFEST_DIR"))
.output()
.ok()?;
if !out.status.success() {
return None;
}
let git_dir = PathBuf::from(String::from_utf8_lossy(&out.stdout).trim());
let repo_root = fs::canonicalize(git_dir).ok()?.parent()?.to_path_buf();
let bin = repo_root.join(".venv/bin");
bin.join("python3").is_file().then_some(bin)
}
fn native_core_present(bin_dir: Option<&Path>) -> bool {
let mut cmd = Command::new("python3");
cmd.arg("-c").arg("import keel_core");
if let Some(dir) = bin_dir {
let path = std::env::var("PATH").unwrap_or_default();
cmd.env("PATH", format!("{}:{path}", dir.display()));
}
cmd.status().is_ok_and(|s| s.success())
}
#[test]
fn crash_restart_resumes_a_real_tier_2_flow() {
let bin_dir = venv_bin_dir();
if !python3_present() || !native_core_present(bin_dir.as_deref()) {
eprintln!("skip: native core (keel_core) not available");
return;
}
let dir = project();
fs::write(
dir.path().join("keel.toml"),
"[flows]\nentrypoints = [\"py:pipeline:main\"]\n\n[target.\"py:pipeline.do_step\"]\n",
)
.unwrap();
fs::write(
dir.path().join("pipeline.py"),
"import os\n\n_LOG = os.path.join(os.path.dirname(__file__), \"steps.log\")\n\n\ndef do_step(n):\n with open(_LOG, \"a\", encoding=\"utf-8\") as f:\n f.write(f\"step-{n}\\n\")\n return {\"step\": n}\n\n\ndef main():\n for n in range(1, 6):\n do_step(n)\n print(\"PIPELINE_COMPLETE\")\n",
)
.unwrap();
let target = dir
.path()
.join("pipeline.py")
.to_string_lossy()
.into_owned();
let plan_path = write_plan(
dir.path(),
"plan.json",
&format!(
r#"{{"v":1,"target":{target:?},"max_restarts":2,"flow_lease_ms":300,"faults":{{"py:pipeline.do_step":[{{"kind":"ok"}},{{"kind":"ok"}},{{"kind":"ok"}},{{"kind":"crash"}}]}},"assert":{{"flow_status":"completed"}}}}"#
),
);
let pythonpath = python_path();
let project_dir = dir.path().to_path_buf();
let path_env = bin_dir.map_or_else(
|| std::env::var("PATH").unwrap_or_default(),
|d| {
format!(
"{}:{}",
d.display(),
std::env::var("PATH").unwrap_or_default()
)
},
);
let r = run_with(dir.path(), &plan_path.to_string_lossy(), &move |cmd| {
cmd.current_dir(&project_dir);
cmd.env("PATH", &path_env);
cmd.env("PYTHONPATH", &pythonpath);
cmd.env("KEEL_BACKEND", "native");
cmd.env("KEEL_QUIET", "1");
});
assert_eq!(r.json["restarts"], 1, "{r:?}");
assert_eq!(r.json["ok"], true, "{r:?}");
assert_eq!(r.json["exit_code"], 0, "{r:?}");
let log = fs::read_to_string(dir.path().join("steps.log")).unwrap();
let mut lines: Vec<&str> = log.lines().collect();
lines.sort_unstable();
assert_eq!(
lines,
vec!["step-1", "step-2", "step-3", "step-4", "step-5"]
);
}
}