use kranz_engine::auth_verify::AuthVerdict;
use kranz_engine::backend::{AgentBackend, PromptMode};
use kranz_engine::backend_claude::parse_stream_line;
use kranz_engine::backend_mock::{
mock_init, mock_result_error, mock_result_json, mock_result_text, mock_text, MockBackend,
MockScript,
};
use kranz_engine::control;
use kranz_engine::cost;
use kranz_engine::event_log::{EventLog, LockForce};
use kranz_engine::events::{Event, EventKind};
use kranz_engine::git_ops::GitRepo;
use kranz_engine::orchestrator::{synthesize_conflict_resolution, MissionEngine, PlanRequest};
use kranz_engine::paths::MissionPaths;
use kranz_engine::reducer;
use kranz_engine::types::*;
use serde_json::json;
use std::io::Write as _;
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::{Arc, Once};
use std::time::Duration;
use tempfile::TempDir;
use tokio::time::timeout;
fn file_exists_cmd(path: &str) -> String {
if cfg!(windows) {
format!("if exist {path} (exit 0) else (exit 1)")
} else {
format!("test -f {path}")
}
}
fn write_line_cmd(text: &str, path: &str) -> String {
if cfg!(windows) {
format!("echo {text}>{path}")
} else {
format!("echo {text} > {path}")
}
}
fn if_file_exists_cmd(path: &str, then: &str) -> String {
if cfg!(windows) {
format!("if exist {path} ({then}) else (exit 1)")
} else {
format!("test -f {path} && {then}")
}
}
const TEST_TIMEOUT: Duration = Duration::from_secs(60);
static ENV_ISOLATION: Once = Once::new();
fn isolate_git_env() {
ENV_ISOLATION.call_once(|| {
let missing = std::env::temp_dir().join(format!(
"kranz-mission-test-no-config-{}",
std::process::id()
));
std::env::set_var("GIT_CONFIG_GLOBAL", &missing);
std::env::set_var("GIT_CONFIG_SYSTEM", &missing);
if let Ok(ceiling) = std::fs::canonicalize(std::env::temp_dir()) {
std::env::set_var("GIT_CEILING_DIRECTORIES", ceiling);
}
});
}
fn git_available() -> bool {
Command::new("git")
.arg("--version")
.output()
.map(|o| o.status.success())
.unwrap_or(false)
}
fn setup() -> bool {
isolate_git_env();
if git_available() {
true
} else {
kranz_engine::test_capability::skip(
kranz_engine::test_capability::capability::GIT,
"git is not on PATH",
);
false
}
}
fn raw_git(dir: &Path, args: &[&str]) -> String {
let out = Command::new("git")
.args(args)
.current_dir(dir)
.output()
.expect("spawn git");
assert!(
out.status.success(),
"git {args:?} failed: {}",
String::from_utf8_lossy(&out.stderr)
);
String::from_utf8_lossy(&out.stdout).into_owned()
}
fn interpret_head_trailers(dir: &Path) -> String {
let message = raw_git(dir, &["log", "-1", "--format=%B"]);
let mut child = Command::new("git")
.args(["interpret-trailers", "--parse"])
.current_dir(dir)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("spawn git interpret-trailers");
child
.stdin
.as_mut()
.expect("stdin piped")
.write_all(message.as_bytes())
.expect("write commit message");
let out = child.wait_with_output().expect("wait for trailers");
assert!(
out.status.success(),
"git interpret-trailers failed: {}",
String::from_utf8_lossy(&out.stderr)
);
String::from_utf8_lossy(&out.stdout).into_owned()
}
fn init_repo() -> (TempDir, PathBuf) {
let dir = tempfile::tempdir().expect("create tempdir");
let init = Command::new("git")
.args(["init", "-b", "main"])
.current_dir(dir.path())
.output()
.expect("spawn git init");
if !init.status.success() {
raw_git(dir.path(), &["init"]);
raw_git(dir.path(), &["symbolic-ref", "HEAD", "refs/heads/main"]);
}
raw_git(dir.path(), &["config", "user.name", "test"]);
raw_git(dir.path(), &["config", "user.email", "test@example.com"]);
std::fs::write(dir.path().join("README.md"), "seed\n").unwrap();
raw_git(dir.path(), &["add", "-A"]);
raw_git(dir.path(), &["commit", "-m", "seed"]);
let root = std::fs::canonicalize(dir.path()).expect("canonicalize repo root");
(dir, root)
}
const GOAL: &str = "ship the demo feature";
fn test_cfg() -> MissionConfig {
MissionConfig {
skip_scrutiny: true,
skip_functional: true,
worker_isolation: WorkerIsolation::Checkout,
validator_allow_uncontained_degrade: true,
..MissionConfig::default()
}
}
fn make_engine(backend: &Arc<MockBackend>, root: &Path, cfg: MissionConfig) -> MissionEngine {
let backend: Arc<dyn AgentBackend> = Arc::clone(backend) as Arc<dyn AgentBackend>;
let mut engine =
MissionEngine::create(backend, root, GOAL, cfg).expect("create mission engine");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
engine
}
fn simple_plan(features: usize, contract: Vec<Assertion>) -> Plan {
Plan {
goal: GOAL.to_string(),
validation_contract: contract,
milestones: vec![PlanMilestone {
title: "M1".to_string(),
features: (1..=features)
.map(|i| PlanFeature {
title: format!("feature {i}"),
spec: format!("build part {i}"),
validation_criteria: vec![format!("part {i} works")],
})
.collect(),
}],
considered_alternatives: None,
command_grants: vec![],
touch_set: vec![],
standards_manifest: None,
reviewer_independence: None,
}
}
fn considered_alternatives() -> ConsideredAlternatives {
ConsideredAlternatives {
chosen: "single integrated slice with tests at the acceptance boundary".to_string(),
rejected: vec![
RejectedAlternative {
approach: "big-bang rewrite".to_string(),
trade_off: "too much review surface for one approval".to_string(),
},
RejectedAlternative {
approach: "docs-only spike".to_string(),
trade_off: "would not deliver the requested behavior".to_string(),
},
],
}
}
fn worker_pass() -> MockScript {
static COUNTER: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let path = format!("delivered-{n}.txt");
MockScript::single_shot_json(&json!({
"result": "pass",
"summary": "implemented and tested",
"filesTouched": [path],
"testsAdded": [],
"testEvidence": "all green",
"commits": []
}))
.writes_file(&path, "delivered by the mock worker\n")
}
fn worker_pass_no_write() -> MockScript {
MockScript::single_shot_json(&json!({
"result": "pass",
"summary": "implemented and tested",
"filesTouched": [],
"testsAdded": [],
"testEvidence": "all green",
"commits": []
}))
}
fn worker_fail() -> MockScript {
MockScript::single_shot_json(&json!({
"result": "fail",
"summary": "could not make the tests pass",
}))
}
fn worker_auth_death() -> MockScript {
MockScript {
events: vec![mock_init("mock-session")],
exit: kranz_engine::backend::SessionExit::Failed(
"claude exited with exit status: 1 without emitting a result message; \
stderr tail: Error: Not logged in"
.to_string(),
),
..Default::default()
}
}
fn validator_with(findings: serde_json::Value) -> MockScript {
MockScript::single_shot_json(&json!({ "findings": findings, "summary": "validated" }))
}
fn validator_denied(command: &str) -> MockScript {
let tool_use_line = json!({
"type": "assistant",
"message": { "id": "vm1", "content": [
{ "type": "tool_use", "name": "Bash", "input": { "command": command } }
] }
})
.to_string();
let denied_line = json!({
"type": "user",
"message": { "role": "user", "content": [
{ "type": "tool_result", "tool_use_id": "vt1",
"content": format!("Permission denied: Bash({command}) requires approval"),
"is_error": true }
] }
})
.to_string();
let mut events = vec![mock_init("mock-session")];
events.extend(parse_stream_line(&tool_use_line));
events.extend(parse_stream_line(&denied_line));
events.push(mock_result_error(&format!(
"stopped: `{command}` was denied and the checks could not run"
)));
MockScript {
events,
..Default::default()
}
}
fn validator_denied_but_passing(command: &str) -> MockScript {
let tool_use_line = json!({
"type": "assistant",
"message": { "id": "vm2", "content": [
{ "type": "tool_use", "name": "Bash", "input": { "command": command } }
] }
})
.to_string();
let denied_line = json!({
"type": "user",
"message": { "role": "user", "content": [
{ "type": "tool_result", "tool_use_id": "vt2",
"content": format!("Permission denied: Bash({command}) requires approval"),
"is_error": true }
] }
})
.to_string();
let report = json!({ "findings": [], "summary": "validated despite one denied probe" });
let mut events = vec![mock_init("mock-session")];
events.extend(parse_stream_line(&tool_use_line));
events.extend(parse_stream_line(&denied_line));
events.push(mock_text(&report.to_string()));
events.push(mock_result_json(&report));
MockScript {
events,
..Default::default()
}
}
fn validator_untrusted_no_denial() -> MockScript {
MockScript {
events: vec![
mock_init("mock-session"),
mock_result_error("stopped: backend produced no trusted report"),
],
..Default::default()
}
}
#[cfg(target_os = "macos")]
fn validator_egress_denied(host: &str, port: u16) -> MockScript {
MockScript {
events: vec![
mock_init("mock-session"),
mock_result_error(&format!(
"stopped: egress to `{host}:{port}` was denied and the checks could not run"
)),
],
..Default::default()
}
.connects_via_proxy(host, port)
}
async fn wait_for_pending_grant(paths: &MissionPaths) {
for _ in 0..400 {
if let Ok(snap) = reducer::read_snapshot(&paths.state_file()) {
if snap.pending_grant_request.is_some() {
return;
}
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
panic!("timed out waiting for a parked grant request");
}
fn orch_script(replies: Vec<String>) -> MockScript {
MockScript::streaming(vec![mock_init("orch-session"), mock_result_text("ready")]).responding(
replies
.iter()
.map(|reply| vec![mock_text(reply), mock_result_text(reply)])
.collect(),
)
}
fn dirty_tree_commit_as_is() -> String {
json!({ "action": "commit-as-is", "note": "worker delivered files" }).to_string()
}
fn judgement(decision: &str, guidance: &str) -> String {
json!({ "decision": decision, "guidance": guidance, "summary": format!("worker judged: {decision}") })
.to_string()
}
fn fix_features(n: usize) -> String {
let features: Vec<serde_json::Value> = (1..=n)
.map(|i| {
json!({
"title": format!("fix issue {i}"),
"spec": format!("resolve validation finding {i}"),
"validationCriteria": [format!("finding {i} resolved")]
})
})
.collect();
json!({ "fixFeatures": features, "summary": format!("{n} fix feature(s)") }).to_string()
}
fn no_lesson() -> String {
"NONE".to_string()
}
fn waive_reply(subject: &str, reason: &str) -> String {
json!({
"fixFeatures": [],
"waived": [{ "subject": subject, "reason": reason }],
"summary": "not worth a fix round"
})
.to_string()
}
fn command_broken_reply(subject: &str, diagnosis: &str) -> String {
json!({
"fixFeatures": [],
"waived": [],
"commandBroken": [{ "subject": subject, "diagnosis": diagnosis }],
"summary": "escalating a possibly author-broken assertion"
})
.to_string()
}
fn parallel_plan(ids: &[&str]) -> String {
json!({
"independent": ids,
"mergeOrder": ids,
"summary": format!("{} features are independent", ids.len())
})
.to_string()
}
fn verdicts_pass(ids: &[&str]) -> String {
let verdicts: Vec<serde_json::Value> = ids
.iter()
.map(|id| json!({ "id": id, "pass": true, "evidence": "verified" }))
.collect();
json!({ "verdicts": verdicts, "summary": "all assertions hold" }).to_string()
}
fn assertion(id: &str, statement: &str, command: Option<&str>) -> Assertion {
Assertion {
id: id.to_string(),
statement: statement.to_string(),
check: if command.is_some() {
AssertionCheck::Command
} else {
AssertionCheck::AgentJudgement
},
command: command.map(str::to_string),
negative_control: None,
pty_script: None,
}
}
fn read_log(paths: &MissionPaths) -> Vec<Event> {
EventLog::read_events(&paths.events_file()).expect("read events.jsonl")
}
fn event_types(events: &[Event]) -> Vec<&'static str> {
events.iter().map(|e| e.kind.type_name()).collect()
}
fn seq_of(events: &[Event], type_name: &str) -> u64 {
events
.iter()
.find(|e| e.kind.type_name() == type_name)
.unwrap_or_else(|| panic!("no {type_name} event in {:?}", event_types(events)))
.seq
}
#[tokio::test(flavor = "multi_thread")]
async fn happy_path_completes_mission_with_tag_and_contract_gate() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![
assertion("a-1", "the build command succeeds", Some("cd .")),
assertion("a-2", "error messages are actionable", None),
];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
dirty_tree_commit_as_is(),
judgement("complete", ""),
verdicts_pass(&["a-2"]),
"Always add a regression test alongside the fix it covers.".to_string(),
]),
worker_pass(),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(2, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
for expected in [
"mission.created",
"plan.approved",
"milestone.started",
"feature.started",
"worker.spawned",
"worker.completed",
"orchestrator.decision",
"feature.completed",
"milestone.validating",
"milestone.completed",
"mission.validating",
"mission.completed",
] {
assert!(types.contains(&expected), "missing {expected}: {types:?}");
}
assert!(seq_of(&events, "milestone.completed") < seq_of(&events, "mission.validating"));
assert!(seq_of(&events, "mission.validating") < seq_of(&events, "mission.completed"));
let tag_name = format!("kranz/{mission_id}/ms-1");
let tags = raw_git(&root, &["tag", "-l"]);
assert!(
tags.contains(&tag_name),
"tag {tag_name} missing from: {tags}"
);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::MilestoneCompleted { tag: Some(t), .. } if *t == tag_name
)));
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::FeatureCompleted { feature_id, commits } if feature_id == "f-1-2" && !commits.is_empty()
)));
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
assert!(state.mission.milestones[0]
.features
.iter()
.all(|f| f.status == FeatureStatus::Complete));
let report = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("report.md"),
)
.expect("report.md written at completion");
assert!(
report.starts_with(&format!("# Mission report — {mission_id}")),
"{report}"
);
assert!(report.contains("## What shipped"), "{report}");
assert!(report.contains("## Workspace"), "{report}");
assert!(report.contains("**Isolation:** `checkout`"), "{report}");
assert!(report.contains("**Worker/validator cwd:**"), "{report}");
assert!(report.contains("**Sandbox:**"), "{report}");
assert!(report.contains("**Preflight:**"), "{report}");
assert!(report.contains("feature 1"), "{report}");
assert!(report.contains("## Validation history"), "{report}");
assert!(report.contains("actual vs"), "{report}");
assert!(report.contains("## Contract outcomes"), "{report}");
assert!(
report.contains("**[a-2]**"),
"both assertions listed: {report}"
);
let approved: kranz_engine::cost::CostEstimate = serde_json::from_str(
&std::fs::read_to_string(paths.estimate_file())
.expect("estimate.json persisted at approval"),
)
.unwrap();
assert!(
report.contains(&format!("(expected ${:.2})", approved.expected_usd)),
"report must cite the approved estimate (expected ${:.2}): {report}",
approved.expected_usd
);
let subject = raw_git(&root, &["log", "-1", "--format=%s"]);
assert_eq!(
subject.trim(),
format!("[kranz] mission report for {mission_id}")
);
let trailers = interpret_head_trailers(&root);
assert!(
trailers.contains(&format!("Kranz-Mission: {mission_id}")),
"{trailers}"
);
assert!(
trailers.contains(&format!("Kranz-Cost-USD: {:.4}", state.total_cost_usd)),
"{trailers}"
);
assert!(
trailers.contains(&format!("Kranz-Tokens-Input: {}", state.totals.input))
&& trailers.contains(&format!("Kranz-Tokens-Output: {}", state.totals.output))
&& trailers.contains(&format!(
"Kranz-Tokens-Cache-Read: {}",
state.totals.cache_read
))
&& trailers.contains(&format!(
"Kranz-Tokens-Cache-Write: {}",
state.totals.cache_write
)),
"{trailers}"
);
let files = raw_git(&root, &["show", "--name-only", "--format=", "HEAD"]);
let mut files: Vec<&str> = files.lines().filter(|l| !l.trim().is_empty()).collect();
files.sort_unstable();
assert_eq!(
files,
vec![
".kranz/lessons/index.md".to_string(),
format!(".kranz/lessons/{mission_id}.md"),
".kranz/missions/index.md".to_string(),
format!(".kranz/missions/{mission_id}/report.md"),
],
"the FINDINGS-EMPTY completion's prose lesson lands in the SAME report commit \
as report.md (not a separate commit)"
);
let lesson = std::fs::read_to_string(
root.join(".kranz")
.join("lessons")
.join(format!("{mission_id}.md")),
)
.expect("lesson file written");
assert!(lesson.contains("regression test"), "{lesson}");
let lessons_index =
std::fs::read_to_string(root.join(".kranz").join("lessons").join("index.md")).unwrap();
assert!(
lessons_index.contains(&format!("{mission_id}.md")),
"{lessons_index}"
);
let index =
std::fs::read_to_string(root.join(".kranz").join("missions").join("index.md")).unwrap();
let line = index
.lines()
.find(|l| l.contains(&format!("[{mission_id}](")))
.expect("mission line in index.md");
assert!(
line.contains(&format!("({mission_id}/plan.md)")),
"plan link kept: {line}"
);
assert!(
line.contains(&format!("[report]({mission_id}/report.md)")),
"report link: {line}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn completed_mission_worker_spawns_record_frontier_provenance() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let spawns: Vec<_> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::WorkerSpawned {
run_id,
quant,
weight_hash,
..
} => Some((run_id.as_str(), quant.as_str(), weight_hash.clone())),
_ => None,
})
.collect();
assert!(
!spawns.is_empty(),
"at least one worker.spawned must be on the log"
);
for (run_id, quant, weight_hash) in &spawns {
assert_eq!(
*quant, "n/a",
"run {run_id} must record the frontier quant sentinel"
);
assert_eq!(
*weight_hash, None,
"run {run_id} must record no weight hash (frontier regime)"
);
}
let state = reducer::fold(&events).unwrap();
for (run_id, _, _) in &spawns {
let run = state.runs.get(*run_id).expect("run recorded");
assert_eq!(run.quant, "n/a");
assert_eq!(run.weight_hash, None);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn invalid_workspace_contract_refused_at_approve() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{"schemaVersion": 1, "secrets": ["sk-live-value-not-a-name"]}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let err = engine
.approve_plan(simple_plan(1, vec![]))
.expect_err("invalid workspace contract must refuse approval");
let msg = err.to_string();
assert!(msg.contains("workspace contract"), "{msg}");
assert!(msg.contains("repo-setup"), "{msg}");
assert!(msg.contains("not a secret NAME"), "{msg}");
let branches = raw_git(&root, &["branch", "--list"]);
assert!(
!branches.contains("kranz/mission-"),
"refused approve must not create the mission branch: {branches}"
);
assert_eq!(
raw_git(&root, &["branch", "--show-current"]).trim(),
"main",
"refused approve must not move the primary checkout"
);
let events = read_log(&engine.paths().clone());
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::PlanApproved { .. })),
"refused approve must not emit plan.approved"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn no_workspace_contract_approves_unchanged() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
assert!(
!root.join(".kranz/workspace.json").exists(),
"fixture must start without a contract"
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine
.approve_plan(simple_plan(1, vec![]))
.expect("approve without a contract behaves as today");
let events = read_log(&engine.paths().clone());
assert!(
events
.iter()
.any(|e| matches!(e.kind, EventKind::PlanApproved { .. })),
"plan.approved must land on the log"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn valid_workspace_contract_approves() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"bootstrap": ["cargo fetch"],
"services": [{"name": "db", "start": "docker compose up db", "port": {"policy": {"fixed": 5432}}}],
"readiness": ["pg_isready"],
"previews": [{"name": "app", "urlTemplate": "http://localhost:{port}/"}],
"secrets": ["DATABASE_URL"],
"mounts": ["/var/cache/cargo"]
}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine
.approve_plan(simple_plan(1, vec![]))
.expect("valid workspace contract must approve");
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_provider_pin_recorded_at_approval() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
assert!(!root.join(".kranz/workspace.json").exists());
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine
.approve_plan(simple_plan(1, vec![]))
.expect("approve pins the provider");
let pin = engine
.state()
.workspace_pin
.as_ref()
.expect("the pin folded into mission state");
assert_eq!(pin.provider, "local-worktree");
assert_eq!(pin.template, "checkout");
assert_eq!(pin.version, "none");
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let pin_pos = events
.iter()
.position(|e| matches!(e.kind, EventKind::WorkspaceProviderPinned { .. }))
.expect("workspace.provider.pinned on the log");
let approved_pos = events
.iter()
.position(|e| matches!(e.kind, EventKind::PlanApproved { .. }))
.expect("plan.approved on the log");
assert_eq!(
approved_pos,
pin_pos + 1,
"the log reads: provider pinned → plan approved"
);
match &events[pin_pos].kind {
EventKind::WorkspaceProviderPinned {
provider,
template,
version,
} => {
assert_eq!(provider, "local-worktree");
assert_eq!(template, "checkout");
assert_eq!(version, "none");
}
other => panic!("wrong variant: {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_provider_pin_with_contract_records_schema_version() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{"schemaVersion": 1, "readiness": ["pg_isready"]}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine
.approve_plan(simple_plan(1, vec![]))
.expect("approve with a contract pins the provider");
let pin = engine
.state()
.workspace_pin
.as_ref()
.expect("the pin folded into mission state");
assert_eq!(pin.provider, "local-worktree");
assert_eq!(pin.template, "checkout");
assert_eq!(
pin.version, "1",
"the contract's schemaVersion is the pin version"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn unknown_workspace_provider_refused_at_approve() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let cfg = MissionConfig {
workspace: WorkspaceConfig {
provider: Some("codr".to_string()), ..Default::default()
},
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
let err = engine
.approve_plan(simple_plan(1, vec![]))
.expect_err("an unknown workspace provider must refuse approval");
let msg = err.to_string();
assert!(msg.contains("workspace.provider"), "{msg}");
assert!(msg.contains("\"codr\""), "{msg}");
assert!(msg.contains("owner: operator"), "{msg}");
let branches = raw_git(&root, &["branch", "--list"]);
assert!(
!branches.contains("kranz/mission-"),
"refused approve must not create the mission branch: {branches}"
);
assert_eq!(
raw_git(&root, &["branch", "--show-current"]).trim(),
"main",
"refused approve must not move the primary checkout"
);
let events = read_log(&engine.paths().clone());
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::PlanApproved { .. })),
"refused approve must not emit plan.approved"
);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkspaceProviderPinned { .. })),
"refused approve must not emit workspace.provider.pinned"
);
}
fn remote_workspace_cfg(base_url: &str, token_env: &str) -> MissionConfig {
MissionConfig {
workspace: WorkspaceConfig {
provider: Some("remote".to_string()),
remote: Some(RemoteWorkspaceConfig {
base_url: Some(base_url.to_string()),
template: Some("tmpl-baked-ami".to_string()),
token_env: Some(token_env.to_string()),
idle_after_hours: None,
}),
teardown_mode: None,
},
..test_cfg()
}
}
struct MockSubstrate {
base_url: String,
requests: Arc<std::sync::Mutex<Vec<String>>>,
}
impl MockSubstrate {
fn requests(&self) -> Vec<String> {
self.requests.lock().unwrap().clone()
}
}
async fn read_request(socket: &mut tokio::net::TcpStream) -> String {
use tokio::io::AsyncReadExt;
let mut buf = Vec::new();
let mut chunk = [0u8; 4096];
loop {
let n = socket.read(&mut chunk).await.expect("read request");
assert!(n > 0, "connection closed before the full request arrived");
buf.extend_from_slice(&chunk[..n]);
if let Some(pos) = buf.windows(4).position(|window| window == b"\r\n\r\n") {
let headers = String::from_utf8_lossy(&buf[..pos]).to_string();
let content_length = headers
.lines()
.find_map(|line| {
line.to_ascii_lowercase()
.strip_prefix("content-length:")
.and_then(|value| value.trim().parse::<usize>().ok())
})
.unwrap_or(0);
if buf.len() >= pos + 4 + content_length {
break;
}
}
}
String::from_utf8_lossy(&buf).to_string()
}
fn spawn_mock_substrate(status: &str) -> MockSubstrate {
spawn_mock_substrate_impl(status, false)
}
fn spawn_mock_substrate_impl(status: &str, fail_transitions: bool) -> MockSubstrate {
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind loopback");
listener
.set_nonblocking(true)
.expect("nonblocking listener");
let base_url = format!("http://{}", listener.local_addr().expect("addr"));
let requests = Arc::new(std::sync::Mutex::new(Vec::new()));
let task_requests = Arc::clone(&requests);
let status = status.to_string();
tokio::spawn(async move {
use tokio::io::AsyncWriteExt;
let listener = tokio::net::TcpListener::from_std(listener).expect("tokio listener");
loop {
let Ok((mut socket, _)) = listener.accept().await else {
return;
};
let request = read_request(&mut socket).await;
task_requests.lock().unwrap().push(request.clone());
let (status_line, body) = if request.starts_with("POST /api/v2/users/me/workspaces ") {
("200 OK", r#"{"id":"ws-m1","urls":[{"name":"app","url":"https://app--m-1.coder.example.com","auth":true}],"takeover":"https://coder.example.com/@me/ws-m1"}"#.to_string())
} else if request.starts_with("GET /api/v2/workspaces/") {
(
"200 OK",
format!(r#"{{"latest_build":{{"status":"{status}"}}}}"#),
)
} else if request.starts_with("POST /api/v2/workspaces/")
&& request.contains("/builds ")
{
if fail_transitions {
(
"500 Internal Server Error",
r#"{"error":"substrate transition failed"}"#.to_string(),
)
} else {
("200 OK", "{}".to_string())
}
} else {
task_requests
.lock()
.unwrap()
.push(format!("UNEXPECTED: {request}"));
("200 OK", "{}".to_string())
};
let response = format!(
"HTTP/1.1 {status_line}\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{body}",
body.len()
);
if socket.write_all(response.as_bytes()).await.is_err() {
return;
}
}
});
MockSubstrate { base_url, requests }
}
#[tokio::test(flavor = "multi_thread")]
async fn remote_workspace_provider_pin_recorded_at_approval() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(&root, r#"{"schemaVersion": 1, "readiness": ["exit 0"]}"#);
let backend = Arc::new(MockBackend::new());
let cfg = remote_workspace_cfg(
"https://coder.internal.example.com",
"KRANZ_TEST_REMOTE_TOKEN_PIN_NEVER_READ",
);
let mut engine = make_engine(&backend, &root, cfg);
engine
.approve_plan(simple_plan(1, vec![]))
.expect("approve pins the remote provider without contacting a substrate");
let pin = engine
.state()
.workspace_pin
.as_ref()
.expect("the pin folded into mission state");
assert_eq!(pin.provider, "remote");
assert_eq!(pin.template, "tmpl-baked-ami", "the configured template id");
assert_eq!(
pin.version, "coder-v1",
"the adapter version, not a contract schema"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let pin_pos = events
.iter()
.position(|e| matches!(e.kind, EventKind::WorkspaceProviderPinned { .. }))
.expect("workspace.provider.pinned on the log");
let approved_pos = events
.iter()
.position(|e| matches!(e.kind, EventKind::PlanApproved { .. }))
.expect("plan.approved on the log");
assert_eq!(
approved_pos,
pin_pos + 1,
"the log reads: provider pinned → plan approved"
);
match &events[pin_pos].kind {
EventKind::WorkspaceProviderPinned {
provider,
template,
version,
} => {
assert_eq!(provider, "remote");
assert_eq!(template, "tmpl-baked-ami");
assert_eq!(version, "coder-v1");
}
other => panic!("wrong variant: {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn remote_workspace_incomplete_config_refused_at_approve() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let cfg = MissionConfig {
workspace: WorkspaceConfig {
provider: Some("remote".to_string()),
remote: Some(RemoteWorkspaceConfig {
base_url: None, template: Some("tmpl-baked-ami".to_string()),
token_env: Some("CODER_SESSION_TOKEN".to_string()),
idle_after_hours: None,
}),
teardown_mode: None,
},
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
let err = engine
.approve_plan(simple_plan(1, vec![]))
.expect_err("incomplete remote config must refuse approval");
let msg = err.to_string();
assert!(msg.contains("workspace.remote.baseUrl"), "{msg}");
assert!(msg.contains("owner: operator"), "{msg}");
assert!(
msg.contains("refusing rather than silently falling back"),
"{msg}"
);
let branches = raw_git(&root, &["branch", "--list"]);
assert!(
!branches.contains("kranz/mission-"),
"refused approve must not create the mission branch: {branches}"
);
let events = read_log(&engine.paths().clone());
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::PlanApproved { .. })),
"refused approve must not emit plan.approved"
);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkspaceProviderPinned { .. })),
"refused approve must not emit workspace.provider.pinned"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn remote_workspace_ready_mission_completes_and_records_substrate_urls() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"bootstrap": ["echo must-not-run > .remote-bootstrap-marker"],
"readiness": ["exit 99"],
"previews": [{"name": "app", "urlTemplate": "http://localhost:{port}/"}],
"secrets": ["DATABASE_URL"]
}"#,
);
std::env::set_var("KRANZ_TEST_REMOTE_TOKEN_READY", "test-token-ready");
let substrate = spawn_mock_substrate("running");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let cfg = remote_workspace_cfg(&substrate.base_url, "KRANZ_TEST_REMOTE_TOKEN_READY");
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
assert!(
!root.join(".remote-bootstrap-marker").exists(),
"contract bootstrap/readiness never executes on the remote substrate in v1"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let (takeover, previews, detail) = events
.iter()
.find_map(|e| match &e.kind {
EventKind::WorkspaceProvisioned {
provider,
takeover,
previews,
detail,
..
} if provider == "remote" => Some((takeover.clone(), previews.clone(), detail.clone())),
_ => None,
})
.expect("workspace.provisioned (remote) on the log");
assert_eq!(
takeover.as_deref(),
Some("https://coder.example.com/@me/ws-m1"),
"the substrate's takeover URL rides the provisioned event"
);
assert_eq!(
previews,
Some(vec![ProvisionedPreview {
name: "app".to_string(),
url: "https://app--m-1.coder.example.com".to_string(),
auth: Some(true),
}]),
"the substrate-reported preview URL, name-matched, with its auth report"
);
assert!(
detail.as_deref().unwrap().contains("DATABASE_URL"),
"the injected secret NAMES are recorded (never values): {detail:?}"
);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::WorkspaceReadinessReport { outcome, .. } if outcome == "ready"
)),
"substrate-reported readiness recorded"
);
let remote_lines = gate_decisions(&events, "workspace remote:");
assert_eq!(remote_lines.len(), 1);
assert!(
remote_lines[0].contains("substrate-reported readiness only"),
"honest wording — no implied contract gate: {}",
remote_lines[0]
);
assert!(
gate_decisions(&events, "workspace bootstrap:").is_empty()
&& gate_decisions(&events, "workspace readiness:").is_empty(),
"no gate phase lines on the remote path (the commands never ran)"
);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::WorkspaceTeardown { mode, .. } if mode == "keep"
)),
"Keep teardown recorded (the workspace stays live for takeover)"
);
let requests = substrate.requests();
let create = requests
.iter()
.find(|r| r.starts_with("POST /api/v2/users/me/workspaces "))
.expect("the create call");
assert!(
create.contains(r#""template_id":"tmpl-baked-ami""#),
"{create}"
);
assert!(
create.contains(r#""env_names":["DATABASE_URL"]"#),
"{create}"
);
assert!(
requests
.iter()
.any(|r| r.starts_with("GET /api/v2/workspaces/ws-m1 ")),
"the readiness poll: {requests:?}"
);
assert!(
!requests.iter().any(|r| r.contains("/builds ")),
"Keep ⇒ no stop/delete transition: {requests:?}"
);
assert!(
!requests.iter().any(|r| r.starts_with("UNEXPECTED:")),
"no unexpected substrate calls: {requests:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn remote_workspace_failed_status_blocks_with_provider_owner() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"previews": [{"name": "app", "urlTemplate": "http://localhost:{port}/"}],
"secrets": ["DATABASE_URL"]
}"#,
);
std::env::set_var("KRANZ_TEST_REMOTE_TOKEN_FAILED", "test-token-failed");
let substrate = spawn_mock_substrate("failed");
let backend = Arc::new(MockBackend::new());
let cfg = remote_workspace_cfg(&substrate.base_url, "KRANZ_TEST_REMOTE_TOKEN_FAILED");
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
assert_eq!(
engine.state().mission.milestones[0].status,
MilestoneStatus::Blocked
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let reason = gate_block_reason(&events, "ms-1");
assert!(
reason.starts_with("workspace provider:"),
"the provider-owned block shape: {reason}"
);
assert!(reason.contains("owner: provider"), "{reason}");
assert!(
reason.contains("kranz-remote-"),
"names the workspace: {reason}"
);
assert!(
!reason.contains("repo-setup") && !reason.contains("owner: operator"),
"distinct from the contract and config owners: {reason}"
);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::WorkspaceReadinessReport { outcome, .. } if outcome == "failed"
)),
"the readiness report records the provider failure"
);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"no spend on a failed workspace"
);
assert!(
substrate
.requests()
.iter()
.any(|r| r.starts_with("POST /api/v2/users/me/workspaces ")),
"the substrate was asked to create the workspace"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn remote_workspace_missing_token_fails_closed_at_run_start() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(&root, r#"{"schemaVersion": 1, "readiness": ["exit 0"]}"#);
let var = "KRANZ_TEST_REMOTE_TOKEN_NEVER_SET";
std::env::remove_var(var); let backend = Arc::new(MockBackend::new());
let cfg = remote_workspace_cfg("http://127.0.0.1:1", var);
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let err = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect_err("missing creds fail closed at provision");
let msg = err.to_string();
assert!(msg.contains(var), "names the env var NAME: {msg}");
assert!(msg.contains("workspace.remote.tokenEnv"), "{msg}");
assert!(msg.contains("owner: operator"), "{msg}");
let events = read_log(&engine.paths().clone());
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkspaceProvisioned { .. })),
"no workspace was provisioned"
);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"no spend"
);
}
fn remote_workspace_cfg_teardown(
base_url: &str,
token_env: &str,
teardown_mode: Option<&str>,
idle_after_hours: Option<f64>,
) -> MissionConfig {
MissionConfig {
workspace: WorkspaceConfig {
provider: Some("remote".to_string()),
remote: Some(RemoteWorkspaceConfig {
base_url: Some(base_url.to_string()),
template: Some("tmpl-baked-ami".to_string()),
token_env: Some(token_env.to_string()),
idle_after_hours,
}),
teardown_mode: teardown_mode.map(str::to_string),
},
..test_cfg()
}
}
fn teardown_outcome(events: &[Event]) -> (String, Option<String>, chrono::DateTime<chrono::Utc>) {
let teardown = workspace_lifecycle_events(events, "workspace.teardown");
assert_eq!(teardown.len(), 1, "exactly one teardown per run()");
match &teardown[0].kind {
EventKind::WorkspaceTeardown { mode, state } => {
(mode.clone(), state.clone(), teardown[0].ts)
}
other => panic!("wrong variant: {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn remote_workspace_terminal_hibernate_stops_the_workspace_and_records_lifecycle() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"readiness": ["exit 0"],
"secrets": ["DATABASE_URL"]
}"#,
);
std::env::set_var("KRANZ_TEST_REMOTE_TOKEN_HIBERNATE", "test-token-hibernate");
let substrate = spawn_mock_substrate("running");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let cfg = remote_workspace_cfg_teardown(
&substrate.base_url,
"KRANZ_TEST_REMOTE_TOKEN_HIBERNATE",
Some("hibernate"),
Some(24.0),
);
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let lifecycle = engine
.state()
.workspace_lifecycle
.clone()
.expect("the teardown outcome folded into state");
assert_eq!(lifecycle.state, "stopped");
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let (mode, state, ts) = teardown_outcome(&events);
assert_eq!(mode, "hibernate");
assert_eq!(state.as_deref(), Some("stopped"));
assert_eq!(
lifecycle.ts, ts,
"the folded lifecycle ts IS the transition event's own ts (workspace-hours anchor)"
);
let detail = events
.iter()
.find_map(|e| match &e.kind {
EventKind::WorkspaceProvisioned {
provider, detail, ..
} if provider == "remote" => detail.clone(),
_ => None,
})
.expect("workspace.provisioned (remote) on the log");
assert!(
detail.contains("idle policy: hibernate after 24h (substrate-owned)"),
"the provisioned event records the substrate-owned idle policy: {detail:?}"
);
let requests = substrate.requests();
let create = requests
.iter()
.find(|r| r.starts_with("POST /api/v2/users/me/workspaces "))
.expect("the create call");
assert!(
create.contains(r#""idle_after_hours":24.0"#),
"the idle policy VALUE rides the create call: {create}"
);
let transitions: Vec<&String> = requests.iter().filter(|r| r.contains("/builds ")).collect();
assert_eq!(transitions.len(), 1, "exactly one transition: {requests:?}");
assert!(
transitions[0].contains(r#""transition":"stop""#),
"hibernate is the stop transition: {transitions:?}"
);
assert!(
!requests.iter().any(|r| r.starts_with("UNEXPECTED:")),
"no unexpected substrate calls: {requests:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn remote_workspace_terminal_destroy_deletes_the_workspace_and_records_lifecycle() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(&root, r#"{"schemaVersion": 1, "readiness": ["exit 0"]}"#);
std::env::set_var("KRANZ_TEST_REMOTE_TOKEN_DESTROY", "test-token-destroy");
let substrate = spawn_mock_substrate("running");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let cfg = remote_workspace_cfg_teardown(
&substrate.base_url,
"KRANZ_TEST_REMOTE_TOKEN_DESTROY",
Some("destroy"),
None,
);
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let lifecycle = engine
.state()
.workspace_lifecycle
.clone()
.expect("the teardown outcome folded into state");
assert_eq!(lifecycle.state, "destroyed");
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let (mode, state, _) = teardown_outcome(&events);
assert_eq!(mode, "destroy");
assert_eq!(state.as_deref(), Some("destroyed"));
let requests = substrate.requests();
let transitions: Vec<&String> = requests.iter().filter(|r| r.contains("/builds ")).collect();
assert_eq!(transitions.len(), 1, "exactly one transition: {requests:?}");
assert!(
transitions[0].contains(r#""transition":"delete""#),
"destroy is the delete transition: {transitions:?}"
);
assert!(
!requests.iter().any(|r| r.starts_with("UNEXPECTED:")),
"no unexpected substrate calls: {requests:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn remote_workspace_teardown_failure_keeps_the_terminal_outcome() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(&root, r#"{"schemaVersion": 1, "readiness": ["exit 0"]}"#);
std::env::set_var(
"KRANZ_TEST_REMOTE_TOKEN_TEARDOWN_FAIL",
"test-token-teardown-fail",
);
let substrate = spawn_mock_substrate_impl("running", true);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let cfg = remote_workspace_cfg_teardown(
&substrate.base_url,
"KRANZ_TEST_REMOTE_TOKEN_TEARDOWN_FAIL",
Some("destroy"),
None,
);
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(
status,
MissionStatus::Complete,
"the teardown failure must not change the mission's terminal outcome"
);
let lifecycle = engine
.state()
.workspace_lifecycle
.clone()
.expect("the failed outcome still folds");
assert_eq!(lifecycle.state, "failed");
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let (mode, state, _) = teardown_outcome(&events);
assert_eq!(mode, "destroy");
assert_eq!(state.as_deref(), Some("failed"));
assert!(
events
.iter()
.any(|e| matches!(e.kind, EventKind::MissionCompleted {})),
"mission.completed still on the log"
);
let decision = events
.iter()
.find_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, detail }
if summary.starts_with("workspace teardown (destroy) failed") =>
{
Some((summary.clone(), detail.clone()))
}
_ => None,
})
.expect("the teardown failure decision: {events:?}");
assert!(
decision.0.contains("the mission outcome stands"),
"{}",
decision.0
);
assert!(decision.0.contains("owner: operator"), "{}", decision.0);
let detail = decision.1.expect("the reason rides the decision detail");
assert!(
detail.contains("delete_workspace") && detail.contains("HTTP 500"),
"the scrubbed provider reason: {detail}"
);
let requests = substrate.requests();
assert!(
requests
.iter()
.any(|r| r.contains("/builds ") && r.contains(r#""transition":"delete""#)),
"the delete transition was attempted: {requests:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_teardown_mode_local_worktree_is_always_keep() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let cfg = MissionConfig {
workspace: WorkspaceConfig {
teardown_mode: Some("destroy".to_string()),
..Default::default()
},
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let lifecycle = engine
.state()
.workspace_lifecycle
.clone()
.expect("the keep outcome folds");
assert_eq!(lifecycle.state, "kept");
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let (mode, state, _) = teardown_outcome(&events);
assert_eq!(
mode, "keep",
"local-worktree NEVER destroys: a configured destroy records an honest keep"
);
assert_eq!(state.as_deref(), Some("kept"));
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_teardown_mode_unknown_fails_closed_at_run_start() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
{
let mut log = EventLog::acquire(&paths, &mission_id, Duration::ZERO, LockForce::No)
.expect("acquire log");
log.append(EventKind::ConfigChanged {
patch: json!({"workspace": {"teardownMode": "purge"}}),
})
.expect("append config.changed");
}
let backend: Arc<dyn AgentBackend> = backend;
let mut engine =
MissionEngine::resume(backend, &root, &mission_id, LockForce::No).expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let err = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect_err("an unknown teardown mode must fail closed");
let msg = err.to_string();
assert!(msg.contains("workspace.teardownMode"), "{msg}");
assert!(msg.contains("\"purge\""), "{msg}");
assert!(msg.contains("owner: operator"), "{msg}");
drop(engine);
let events = read_log(&paths);
for wire_name in [
"workspace.provisioned",
"workspace.readiness",
"workspace.teardown",
] {
assert!(
workspace_lifecycle_events(&events, wire_name).is_empty(),
"no {wire_name} event on a run that failed closed"
);
}
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"no worker may spawn when the teardown mode fails closed"
);
}
fn portable_shell_json(contract_json: &str) -> String {
if !cfg!(windows) {
return contract_json.to_string();
}
contract_json
.split('"')
.enumerate()
.map(|(index, segment)| {
if index % 2 == 0 {
return segment.to_string();
}
match segment.strip_prefix("test -f ") {
None => segment.to_string(),
Some(rest) => match rest.split_once(" && ") {
Some((path, then)) => format!("if exist {path} ({then}) else (exit 1)"),
None => format!("if exist {rest} (exit 0) else (exit 1)"),
},
}
})
.collect::<Vec<_>>()
.join("\"")
}
fn commit_workspace_contract(root: &Path, contract_json: &str) {
let contract_json = portable_shell_json(contract_json);
std::fs::create_dir_all(root.join(".kranz")).unwrap();
std::fs::write(root.join(".kranz").join("workspace.json"), &contract_json).unwrap();
std::fs::write(root.join(".gitignore"), ".boot-marker\n.boot-count\n").unwrap();
raw_git(root, &["add", ".kranz/workspace.json", ".gitignore"]);
raw_git(root, &["commit", "-m", "workspace contract"]);
}
fn gate_decisions<'a>(events: &'a [Event], prefix: &str) -> Vec<&'a str> {
events
.iter()
.filter_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, .. } if summary.starts_with(prefix) => {
Some(summary.as_str())
}
_ => None,
})
.collect()
}
fn gate_block_reason(events: &[Event], milestone_id: &str) -> String {
events
.iter()
.rev()
.find_map(|e| match &e.kind {
EventKind::MilestoneBlocked {
milestone_id: id,
reason,
..
} if id == milestone_id => Some(reason.clone()),
_ => None,
})
.unwrap_or_else(|| panic!("no milestone.blocked for {milestone_id} on the log"))
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_bootstrap_readiness_gate_runs_before_workers() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"bootstrap": ["echo boot > .boot-marker"],
"readiness": ["test -f .boot-marker"]
}"#,
);
let contract = vec![assertion(
"a-1",
"the workspace marker exists",
Some(&file_exists_cmd(".boot-marker")),
)];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
let marker = std::fs::read_to_string(root.join(".boot-marker"))
.expect("bootstrap wrote the marker into the workspace cwd");
assert!(marker.contains("boot"), "{marker}");
let events = read_log(&paths);
assert_eq!(
gate_decisions(&events, "workspace bootstrap:"),
vec![
"workspace bootstrap: running 1 commands",
"workspace bootstrap: 1/1 commands ok"
]
);
assert_eq!(
gate_decisions(&events, "workspace readiness:"),
vec![
"workspace readiness: running 1 checks",
"workspace readiness: 1/1 checks ok"
]
);
let readiness_ok_seq = events
.iter()
.find(|e| matches!(&e.kind, EventKind::OrchestratorDecision { summary, .. } if summary == "workspace readiness: 1/1 checks ok"))
.map(|e| e.seq)
.unwrap();
assert!(
readiness_ok_seq < seq_of(&events, "worker.spawned"),
"readiness must pass before the first worker spawns"
);
let report = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("report.md"),
)
.expect("report.md written at completion");
assert!(
report.contains("- **Bootstrap:** 1/1 commands ok"),
"{report}"
);
assert!(
report.contains("- **Readiness:** 1/1 checks ok"),
"{report}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_bootstrap_failure_blocks_before_any_worker() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"bootstrap": [
"echo boot > .boot-marker",
"echo leaking sk-ant-api03-a1b2c3d4e5f6 1>&2 && exit 42"
],
"readiness": ["test -f .boot-marker"]
}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
assert_eq!(engine.state().mission.status, MissionStatus::Blocked);
assert_eq!(
engine.state().mission.milestones[0].status,
MilestoneStatus::Blocked
);
let paths = engine.paths().clone();
drop(engine);
assert!(root.join(".boot-marker").exists());
let events = read_log(&paths);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"no worker may spawn on a failed bootstrap: {:?}",
event_types(&events)
);
assert!(seq_of(&events, "milestone.started") < seq_of(&events, "milestone.blocked"));
assert_eq!(
gate_decisions(&events, "workspace bootstrap:"),
vec![
"workspace bootstrap: running 2 commands",
"workspace bootstrap: FAILED at command 2/2 — blocking mission (owner: repo-setup)"
]
);
assert!(gate_decisions(&events, "workspace readiness:").is_empty());
let reason = gate_block_reason(&events, "ms-1");
assert!(
reason.contains("workspace gate: bootstrap command 2/2 failed"),
"{reason}"
);
assert!(reason.contains("owner: repo-setup"), "{reason}");
assert!(reason.contains("exit code 42"), "{reason}");
assert!(reason.contains("echo leaking"), "{reason}");
assert!(
!reason.contains("sk-ant-api03-a1b2c3d4e5f6"),
"the output tail must be scrubbed: {reason}"
);
assert!(reason.contains("[REDACTED]"), "{reason}");
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_readiness_failure_blocks_after_bootstrap() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"bootstrap": ["echo boot > .boot-marker"],
"readiness": ["test -f .no-such-readiness-file"]
}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
let paths = engine.paths().clone();
drop(engine);
assert!(root.join(".boot-marker").exists());
let events = read_log(&paths);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"no worker may spawn on a failed readiness check"
);
assert_eq!(
gate_decisions(&events, "workspace bootstrap:"),
vec![
"workspace bootstrap: running 1 commands",
"workspace bootstrap: 1/1 commands ok"
]
);
assert_eq!(
gate_decisions(&events, "workspace readiness:"),
vec![
"workspace readiness: running 1 checks",
"workspace readiness: FAILED at check 1/1 — blocking mission (owner: repo-setup)"
]
);
let reason = gate_block_reason(&events, "ms-1");
assert!(
reason.contains("workspace gate: readiness check 1/1 failed"),
"{reason}"
);
assert!(reason.contains("owner: repo-setup"), "{reason}");
assert!(
reason.contains(&file_exists_cmd(".no-such-readiness-file")),
"the reason names the failing check: {reason}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn no_contract_run_has_no_workspace_gate_events() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
assert!(
!root.join(".kranz/workspace.json").exists(),
"fixture must start without a contract"
);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(gate_decisions(&events, "workspace bootstrap:").is_empty());
assert!(gate_decisions(&events, "workspace readiness:").is_empty());
let report = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("report.md"),
)
.expect("report.md written at completion");
assert!(
report.contains("- **Workspace contract:** no workspace contract (source isolation only)"),
"{report}"
);
assert!(!report.contains("- **Bootstrap:**"), "{report}");
assert!(!report.contains("- **Readiness:**"), "{report}");
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_bootstrap_reruns_on_resume_after_crash() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"bootstrap": ["echo run >> .boot-count"],
"readiness": ["test -f .boot-count"]
}"#,
);
let backend1 = Arc::new(MockBackend::with_scripts(vec![
worker_pass_no_write(),
orch_script(vec![]),
]));
let mut engine = make_engine(&backend1, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
engine.set_orch_stall_timeout(Duration::from_millis(400));
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect_err("phase 1 must error out (simulated crash)");
drop(engine);
let count = std::fs::read_to_string(root.join(".boot-count"))
.expect("phase 1 bootstrap wrote the count file");
assert_eq!(
count.lines().count(),
1,
"phase 1 ran bootstrap once: {count:?}"
);
let backend2 = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let backend2_dyn: Arc<dyn AgentBackend> = Arc::clone(&backend2) as Arc<dyn AgentBackend>;
let mut engine = MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No)
.expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
drop(engine);
let count = std::fs::read_to_string(root.join(".boot-count"))
.expect("count file persists across the resume");
assert_eq!(
count.lines().count(),
2,
"resume re-ran bootstrap (once per run() invocation): {count:?}"
);
let events = read_log(&paths);
assert_eq!(
gate_decisions(&events, "workspace bootstrap: running").len(),
2,
"one bootstrap start decision per run() invocation"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_gate_block_lifts_once_environment_is_fixed() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let ready_dir = root.join("env-ready");
let contract = serde_json::json!({
"schemaVersion": 1,
"readiness": ["cd env-ready"],
})
.to_string();
commit_workspace_contract(&root, &contract);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
drop(engine);
std::fs::create_dir_all(&ready_dir).unwrap();
let backend2 = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let backend2_dyn: Arc<dyn AgentBackend> = Arc::clone(&backend2) as Arc<dyn AgentBackend>;
let mut engine = MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No)
.expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
let paths2 = engine.paths().clone();
drop(engine);
let last_block_reason = read_log(&paths2)
.iter()
.rev()
.find_map(|e| match &e.kind {
EventKind::MilestoneBlocked { reason, .. } => Some(reason.clone()),
_ => None,
})
.unwrap_or_else(|| "(none)".to_string());
assert_eq!(
status,
MissionStatus::Complete,
"a fixed environment must resume without an unblock consultation (last block reason: {last_block_reason})"
);
let events = read_log(&paths);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::MilestoneUnblocked { milestone_id, reason, .. }
if milestone_id == "ms-1" && reason.contains("workspace gate now passing")
)));
assert!(seq_of(&events, "milestone.unblocked") < seq_of(&events, "worker.spawned"));
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_gate_reads_contract_from_base_branch_not_mission_branch() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"readiness": ["test -f .base-required-marker"]
}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine
.approve_plan(simple_plan(1, vec![]))
.expect("base contract is valid");
std::fs::write(
root.join(".kranz").join("workspace.json"),
br#"{"schemaVersion": 1}"#,
)
.unwrap();
raw_git(&root, &["add", ".kranz/workspace.json"]);
raw_git(&root, &["commit", "-m", "weaken the workspace contract"]);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(
status,
MissionStatus::Blocked,
"the gate must apply the BASE branch's contract, not the weakened mission-branch copy"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let reason = gate_block_reason(&events, "ms-1");
assert!(
reason.contains(&file_exists_cmd(".base-required-marker")),
"the base branch's readiness check gated the run: {reason}"
);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"the weakened contract must never let a worker spawn"
);
}
fn workspace_lifecycle_events<'a>(events: &'a [Event], wire_name: &str) -> Vec<&'a Event> {
events
.iter()
.filter(|e| e.kind.type_name() == wire_name)
.collect()
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_provider_events_land_and_fold_into_state() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"bootstrap": ["echo boot > .boot-marker"],
"readiness": ["test -f .boot-marker"]
}"#,
);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
assert_eq!(
engine.state().workspace_provider.as_deref(),
Some("local-worktree")
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let provisioned = workspace_lifecycle_events(&events, "workspace.provisioned");
assert_eq!(provisioned.len(), 1, "one provision per run()");
match &provisioned[0].kind {
EventKind::WorkspaceProvisioned {
provider,
cwd,
detail,
..
} => {
assert_eq!(provider, "local-worktree");
assert_eq!(
cwd,
&root.display().to_string(),
"checkout mode provisions the repo root as the workspace cwd"
);
assert_eq!(detail, &None, "local-worktree carries no detail");
}
other => panic!("wrong variant: {other:?}"),
}
let readiness = workspace_lifecycle_events(&events, "workspace.readiness");
assert_eq!(readiness.len(), 1);
match &readiness[0].kind {
EventKind::WorkspaceReadinessReport { outcome, detail } => {
assert_eq!(outcome, "ready");
assert_eq!(detail, &None);
}
other => panic!("wrong variant: {other:?}"),
}
let teardown = workspace_lifecycle_events(&events, "workspace.teardown");
assert_eq!(teardown.len(), 1);
match &teardown[0].kind {
EventKind::WorkspaceTeardown { mode, .. } => assert_eq!(mode, "keep"),
other => panic!("wrong variant: {other:?}"),
}
let first_gate_seq = events
.iter()
.find(|e| matches!(&e.kind, EventKind::OrchestratorDecision { summary, .. } if summary.starts_with("workspace bootstrap:")))
.map(|e| e.seq)
.expect("a gate decision on the log");
assert!(provisioned[0].seq < first_gate_seq);
assert!(readiness[0].seq > first_gate_seq);
assert!(readiness[0].seq < seq_of(&events, "worker.spawned"));
assert!(teardown[0].seq > readiness[0].seq);
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_provider_no_contract_provisions_without_readiness_report() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
assert!(!root.join(".kranz/workspace.json").exists());
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
assert_eq!(
engine.state().workspace_provider.as_deref(),
Some("local-worktree")
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert_eq!(
workspace_lifecycle_events(&events, "workspace.provisioned").len(),
1
);
assert!(
workspace_lifecycle_events(&events, "workspace.readiness").is_empty(),
"no readiness artifact without a contract (D-H): {:?}",
event_types(&events)
);
assert_eq!(
workspace_lifecycle_events(&events, "workspace.teardown").len(),
1
);
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_provider_readiness_failure_blocks_with_failed_report() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_workspace_contract(
&root,
r#"{
"schemaVersion": 1,
"bootstrap": ["echo boot > .boot-marker"],
"readiness": ["test -f .no-such-readiness-file"]
}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let reason = gate_block_reason(&events, "ms-1");
assert!(
reason.contains("workspace gate: readiness check 1/1 failed"),
"{reason}"
);
let readiness = workspace_lifecycle_events(&events, "workspace.readiness");
assert_eq!(readiness.len(), 1);
match &readiness[0].kind {
EventKind::WorkspaceReadinessReport { outcome, detail } => {
assert_eq!(outcome, "failed");
assert_eq!(
detail.as_deref(),
Some(reason.as_str()),
"the report carries the same scrubbed reason the block records"
);
}
other => panic!("wrong variant: {other:?}"),
}
assert!(readiness[0].seq < seq_of(&events, "milestone.blocked"));
assert_eq!(
workspace_lifecycle_events(&events, "workspace.teardown").len(),
1
);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"no worker may spawn on a failed readiness check"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_provider_unknown_name_fails_closed_at_run_start() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
{
let mut log = EventLog::acquire(&paths, &mission_id, Duration::ZERO, LockForce::No)
.expect("acquire log");
log.append(EventKind::ConfigChanged {
patch: json!({"workspace": {"provider": "coder"}}),
})
.expect("append config.changed");
}
let backend: Arc<dyn AgentBackend> = backend;
let mut engine =
MissionEngine::resume(backend, &root, &mission_id, LockForce::No).expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let err = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect_err("an unknown workspace provider must fail closed");
let msg = err.to_string();
assert!(msg.contains("workspace.provider"), "{msg}");
assert!(msg.contains("\"coder\""), "{msg}");
assert!(msg.contains("local-worktree"), "{msg}");
assert_eq!(engine.state().workspace_provider, None);
drop(engine);
let events = read_log(&paths);
for wire_name in [
"workspace.provisioned",
"workspace.readiness",
"workspace.teardown",
] {
assert!(
workspace_lifecycle_events(&events, wire_name).is_empty(),
"no {wire_name} event on a run that failed closed"
);
}
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"no worker may spawn when the provider fails closed"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn workspace_provider_resume_reprovisions_idempotently() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend1 = Arc::new(MockBackend::with_scripts(vec![
worker_pass_no_write(),
orch_script(vec![]),
]));
let mut engine = make_engine(&backend1, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
engine.set_orch_stall_timeout(Duration::from_millis(400));
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect_err("phase 1 must error out (simulated crash)");
drop(engine);
let events = read_log(&paths);
assert_eq!(
workspace_lifecycle_events(&events, "workspace.provisioned").len(),
1,
"the crashed run provisioned exactly once"
);
assert!(
workspace_lifecycle_events(&events, "workspace.teardown").is_empty(),
"a crashed run records no teardown (the resume sweep owns leftovers)"
);
let backend2 = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let backend2_dyn: Arc<dyn AgentBackend> = Arc::clone(&backend2) as Arc<dyn AgentBackend>;
let mut engine = MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No)
.expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
assert_eq!(
engine.state().workspace_provider.as_deref(),
Some("local-worktree")
);
drop(engine);
let events = read_log(&paths);
let provisioned = workspace_lifecycle_events(&events, "workspace.provisioned");
assert_eq!(
provisioned.len(),
2,
"resume re-provisions: one workspace.provisioned per run() invocation"
);
for event in provisioned {
match &event.kind {
EventKind::WorkspaceProvisioned { provider, .. } => {
assert_eq!(provider, "local-worktree")
}
other => panic!("wrong variant: {other:?}"),
}
}
assert_eq!(
workspace_lifecycle_events(&events, "workspace.teardown").len(),
1,
"only the completing run records teardown"
);
}
fn commit_data_contract(root: &Path, contract_json: &str) {
let contract_json = portable_shell_json(contract_json);
std::fs::create_dir_all(root.join(".kranz")).unwrap();
std::fs::write(root.join(".kranz").join("workspace.json"), &contract_json).unwrap();
std::fs::write(
root.join(".gitignore"),
".boot-marker\n.data-clone-marker\n.data-migrate-marker\n.data-skew-marker\n.data-reset-count\n",
)
.unwrap();
raw_git(root, &["add", ".kranz/workspace.json", ".gitignore"]);
raw_git(root, &["commit", "-m", "workspace contract"]);
}
fn decision_seq(events: &[Event], prefix: &str) -> u64 {
events
.iter()
.find(|e| matches!(&e.kind, EventKind::OrchestratorDecision { summary, .. } if summary.starts_with(prefix)))
.map(|e| e.seq)
.unwrap_or_else(|| panic!("no decision starting with {prefix:?}"))
}
#[tokio::test(flavor = "multi_thread")]
async fn golden_data_hooks_run_in_provision_order_before_workers() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_data_contract(
&root,
r#"{
"schemaVersion": 1,
"data": {
"clone": "echo cloned > .data-clone-marker",
"migrate": "test -f .data-clone-marker && echo mig > .data-migrate-marker",
"skewCheck": "test -f .boot-marker && echo checked > .data-skew-marker"
},
"bootstrap": ["test -f .data-migrate-marker && echo boot > .boot-marker"],
"readiness": ["test -f .boot-marker"]
}"#,
);
let contract = vec![assertion(
"a-1",
"the skew check ran",
Some(&file_exists_cmd(".data-skew-marker")),
)];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert_eq!(
gate_decisions(&events, "workspace data:"),
vec![
"workspace data: clone `echo cloned > .data-clone-marker` → ok (exit code 0)"
.to_string(),
format!(
"workspace data: migrate `{}` → ok (exit code 0)",
if_file_exists_cmd(".data-clone-marker", "echo mig > .data-migrate-marker")
),
format!(
"workspace data: skewCheck `{}` → ok (exit code 0)",
if_file_exists_cmd(".boot-marker", "echo checked > .data-skew-marker")
),
]
);
assert!(
decision_seq(&events, "workspace data: migrate")
< decision_seq(&events, "workspace bootstrap:")
);
assert!(
decision_seq(&events, "workspace readiness: 1/1")
< decision_seq(&events, "workspace data: skewCheck")
);
assert!(decision_seq(&events, "workspace data: skewCheck") < seq_of(&events, "worker.spawned"));
let readiness = workspace_lifecycle_events(&events, "workspace.readiness");
assert_eq!(readiness.len(), 1);
match &readiness[0].kind {
EventKind::WorkspaceReadinessReport { outcome, .. } => assert_eq!(outcome, "ready"),
other => panic!("wrong variant: {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn golden_data_skew_blocks_owned_actionable_and_never_a_finding() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_data_contract(
&root,
r#"{
"schemaVersion": 1,
"data": {
"migrate": "echo mig > .data-migrate-marker",
"skewCheck": "echo skewed sk-ant-api03-a1b2c3d4e5f6 1>&2 && exit 1"
},
"bootstrap": ["echo boot > .boot-marker"],
"readiness": ["test -f .boot-marker"]
}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
assert_eq!(
engine.state().mission.milestones[0].status,
MilestoneStatus::Blocked
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert_eq!(
gate_decisions(&events, "workspace bootstrap:"),
vec![
"workspace bootstrap: running 1 commands",
"workspace bootstrap: 1/1 commands ok"
]
);
assert_eq!(
gate_decisions(&events, "workspace readiness:"),
vec![
"workspace readiness: running 1 checks",
"workspace readiness: 1/1 checks ok"
]
);
let reason = gate_block_reason(&events, "ms-1");
assert!(
reason.contains("workspace gate: data skewCheck failed"),
"{reason}"
);
assert!(reason.contains("owner: repo-setup"), "{reason}");
assert!(reason.contains("exit code 1"), "{reason}");
assert!(
reason
.contains("run the data migrate hook (`echo mig > .data-migrate-marker`), then resume"),
"the action names the declared migrate hook: {reason}"
);
assert!(
!reason.contains("readiness check"),
"the skew reason never presents as a readiness flake: {reason}"
);
assert!(
!reason.contains("sk-ant-api03-a1b2c3d4e5f6"),
"the output tail must be scrubbed: {reason}"
);
assert!(reason.contains("[REDACTED]"), "{reason}");
let readiness = workspace_lifecycle_events(&events, "workspace.readiness");
assert_eq!(readiness.len(), 1);
match &readiness[0].kind {
EventKind::WorkspaceReadinessReport { outcome, detail } => {
assert_eq!(outcome, "skew");
assert_eq!(detail.as_deref(), Some(reason.as_str()));
}
other => panic!("wrong variant: {other:?}"),
}
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::ValidationFinding { .. })),
"skew must not present as a validator finding: {:?}",
event_types(&events)
);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"no worker may spawn on a skewed dataset"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn golden_data_reset_runs_before_each_validation_round() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_data_contract(
&root,
r#"{
"schemaVersion": 1,
"data": {
"reset": "echo reset >> .data-reset-count",
"resetBetweenRounds": true
}
}"#,
);
let finding = json!([{
"subject": "part 1 works",
"severity": "major",
"evidence": "the endpoint returns 500 on empty input",
"suggestedFix": "guard empty input"
}]);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1),
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_with(finding),
worker_pass(),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let count = std::fs::read_to_string(root.join(".data-reset-count"))
.expect("the reset hook appended into the workspace cwd");
assert_eq!(
count.lines().count(),
2,
"one reset per validation round: {count:?}"
);
let events = read_log(&paths);
let resets = gate_decisions(&events, "workspace data: reset");
assert_eq!(
resets,
vec!["workspace data: reset `echo reset >> .data-reset-count` → ok (exit code 0)"; 2],
"one reset decision per round"
);
let validating_seqs: Vec<u64> = events
.iter()
.filter(|e| matches!(e.kind, EventKind::MilestoneValidating { .. }))
.map(|e| e.seq)
.collect();
assert_eq!(validating_seqs.len(), 2, "two validation rounds");
let reset_seqs: Vec<u64> = events
.iter()
.filter(|e| matches!(&e.kind, EventKind::OrchestratorDecision { summary, .. } if summary.starts_with("workspace data: reset")))
.map(|e| e.seq)
.collect();
assert_eq!(reset_seqs.len(), 2, "one reset decision per round");
assert!(
validating_seqs[0] < reset_seqs[0] && reset_seqs[0] < validating_seqs[1],
"round 1's reset precedes round 2: {validating_seqs:?} vs {reset_seqs:?}"
);
assert!(
validating_seqs[1] < reset_seqs[1],
"round 2's reset follows its validating event: {validating_seqs:?} vs {reset_seqs:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn golden_data_reset_failure_blocks_owned_before_any_validator() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_data_contract(
&root,
r#"{
"schemaVersion": 1,
"data": {
"reset": "exit 7",
"resetBetweenRounds": true
}
}"#,
);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![dirty_tree_commit_as_is(), judgement("complete", "")]),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
assert_eq!(
engine.state().mission.milestones[0].status,
MilestoneStatus::Blocked
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
seq_of(&events, "milestone.validating") < seq_of(&events, "milestone.blocked"),
"the block lands inside the validation round"
);
assert_eq!(
gate_decisions(&events, "workspace data: reset"),
vec![
"workspace data: reset `exit 7` → FAILED (exit code 7) — blocking mission (owner: repo-setup)"
]
);
let reason = gate_block_reason(&events, "ms-1");
assert!(
reason.contains("workspace gate: data reset hook failed"),
"{reason}"
);
assert!(reason.contains("owner: repo-setup"), "{reason}");
assert!(reason.contains("exit code 7"), "{reason}");
assert!(reason.contains("`exit 7`"), "{reason}");
assert!(reason.contains("then resume"), "{reason}");
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::ValidationFinding { .. })),
"a reset failure must not present as a validator finding: {:?}",
event_types(&events)
);
}
#[tokio::test(flavor = "multi_thread")]
async fn validation_round_creates_fix_feature_then_completes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let finding = json!([{
"subject": "part 1 works",
"severity": "major",
"evidence": "the endpoint returns 500 on empty input",
"suggestedFix": "guard empty input"
}]);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1),
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_with(finding),
worker_pass(),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let state = engine.state();
let ms = &state.mission.milestones[0];
assert_eq!(ms.fix_cycles, 1);
let fix = ms
.features
.iter()
.find(|f| f.origin == FeatureOrigin::Fix)
.expect("fix feature exists");
assert_eq!(fix.id, "ms-1-fix-1-1");
assert_eq!(fix.status, FeatureStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
types.contains(&"validation.finding"),
"finding event: {types:?}"
);
assert!(
types.contains(&"fixfeature.created"),
"fixfeature event: {types:?}"
);
assert_eq!(
types
.iter()
.filter(|t| **t == "milestone.validating")
.count(),
2,
"two validation rounds: {types:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn loop_guard_blocks_milestone_after_max_fix_cycles() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let finding = json!([{
"subject": "part 1 works",
"severity": "major",
"evidence": "still failing",
"suggestedFix": ""
}]);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1),
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1), ]),
validator_with(finding.clone()),
worker_pass(),
validator_with(finding),
]));
let cfg = MissionConfig {
skip_functional: false,
max_fix_cycles_per_milestone: 1,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
assert_eq!(engine.state().mission.status, MissionStatus::Blocked);
assert_eq!(
engine.state().mission.milestones[0].status,
MilestoneStatus::Blocked
);
assert_eq!(engine.state().mission.milestones[0].fix_cycles, 1);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::MilestoneBlocked { reason, .. } if reason.contains("fix-cycle cap")
)));
assert_eq!(
event_types(&events)
.iter()
.filter(|t| **t == "fixfeature.created")
.count(),
1,
"only round 1 created a fix feature"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn unblock_guidance_survives_restart_and_reaches_validator() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let finding = json!([{
"subject": "part 1 works",
"severity": "major",
"evidence": "still failing",
"suggestedFix": ""
}]);
let backend1 = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1),
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1),
]),
validator_with(finding.clone()),
worker_pass(),
validator_with(finding),
]));
let cfg = MissionConfig {
skip_functional: false,
max_fix_cycles_per_milestone: 1,
..test_cfg()
};
let mut engine = make_engine(&backend1, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
control::enqueue(
&paths,
&ControlCommand::Msg {
text: "unblock and guide the validator".into(),
interrupt: false,
},
)
.unwrap();
let guidance = "FMT FIRST: run cargo fmt before the contract gate";
let unblock = json!({
"action": "unblock-raise-cap",
"note": "cap raised with guidance",
"validatorGuidance": guidance,
})
.to_string();
let backend2 = Arc::new(MockBackend::with_scripts(vec![
orch_script(vec![unblock, no_lesson()]),
validator_with(json!([])),
]));
let backend2_dyn: Arc<dyn AgentBackend> = Arc::clone(&backend2) as Arc<dyn AgentBackend>;
let mut engine = MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No)
.expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
drop(engine);
let events = read_log(&paths);
let unblock_idx = events
.iter()
.position(|e| matches!(&e.kind, EventKind::MilestoneUnblocked { .. }))
.expect("unblock event present");
match &events[unblock_idx].kind {
EventKind::MilestoneUnblocked {
validator_guidance, ..
} => assert_eq!(validator_guidance.as_deref(), Some(guidance)),
_ => unreachable!(),
}
let prefix_state = reducer::fold(&events[..=unblock_idx]).unwrap();
assert_eq!(
prefix_state.mission.milestones[0]
.validator_guidance
.as_deref(),
Some(guidance),
"folded state must carry the guidance across a restart"
);
let specs = backend2.started_specs();
let validator_task = specs
.iter()
.find_map(|s| match &s.prompt {
PromptMode::SingleShot(task) if task.contains("Validate milestone") => Some(task),
_ => None,
})
.expect("a validator session ran in phase 2");
assert!(
validator_task.contains(guidance),
"validator task must carry the operator guidance verbatim: {validator_task}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn unblock_add_fix_schedules_repair_before_revalidation() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let finding = json!([{
"subject": "part 1 works",
"severity": "major",
"evidence": "fmt check fails",
"suggestedFix": ""
}]);
let backend1 = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1),
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1),
]),
validator_with(finding.clone()),
worker_pass(),
validator_with(finding),
]));
let cfg = MissionConfig {
skip_functional: false,
max_fix_cycles_per_milestone: 1,
..test_cfg()
};
let mut engine = make_engine(&backend1, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
control::enqueue(
&paths,
&ControlCommand::Msg {
text: "it is just rustfmt — repair then re-validate".into(),
interrupt: false,
},
)
.unwrap();
let decision = json!({
"action": "unblock-add-fix",
"note": "schedule a fmt repair",
"fix": {
"title": "run cargo fmt --all",
"spec": "run cargo fmt --all and commit the result",
"validationCriteria": ["cargo fmt --all --check exits clean"]
}
})
.to_string();
let backend2 = Arc::new(MockBackend::with_scripts(vec![
orch_script(vec![
decision,
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
worker_pass(),
validator_with(json!([])),
]));
let backend2_dyn: Arc<dyn AgentBackend> = Arc::clone(&backend2) as Arc<dyn AgentBackend>;
let mut engine = MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No)
.expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let state = engine.state();
assert_eq!(
state.mission.milestones[0].fix_cycles, 1,
"a blocked-state repair must not spend a fix cycle"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let unblock_seq = events
.iter()
.find(|e| matches!(&e.kind, EventKind::MilestoneUnblocked { .. }))
.expect("unblock event")
.seq;
let post_unblock_fixes: Vec<_> = events
.iter()
.filter(|e| e.seq > unblock_seq && matches!(&e.kind, EventKind::FixFeatureCreated { .. }))
.collect();
assert_eq!(
post_unblock_fixes.len(),
1,
"exactly one repair feature follows the unblock"
);
match &post_unblock_fixes[0].kind {
EventKind::FixFeatureCreated { feature, .. } => {
assert_eq!(feature.title, "run cargo fmt --all");
assert_eq!(feature.origin, FeatureOrigin::Fix);
}
_ => unreachable!(),
}
}
#[tokio::test(flavor = "multi_thread")]
async fn contract_commands_run_engine_side_and_reach_validator_task() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![assertion("a-1", "the build succeeds", Some("git status"))];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
drop(engine);
let specs = backend.started_specs();
let validator_task = specs
.iter()
.find_map(|s| match &s.prompt {
PromptMode::SingleShot(task) if task.contains("Validate milestone") => Some(task),
_ => None,
})
.expect("a validator session ran");
assert!(
validator_task.contains("[a-1] `git status` → PASS"),
"engine-captured PASS must ride the validator task: {validator_task}"
);
assert!(
validator_task.contains("do NOT re-run"),
"the results block must forbid re-running: {validator_task}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn waive_completes_milestone() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let finding = json!([{
"subject": "part 1 works",
"severity": "minor",
"evidence": "the new helper's docstring omits the error case",
"suggestedFix": "extend the docstring"
}]);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
waive_reply("part 1 works", "docstring nitpick"),
no_lesson(),
]),
validator_with(finding),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let ms = &engine.state().mission.milestones[0];
assert_eq!(ms.status, MilestoneStatus::Complete);
assert_eq!(ms.fix_cycles, 0, "a waived round consumes no fix cycle");
assert!(ms.features.iter().all(|f| f.origin == FeatureOrigin::Plan));
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
types.contains(&"validation.finding"),
"finding still surfaced: {types:?}"
);
assert!(
!types.contains(&"fixfeature.created"),
"no fix feature: {types:?}"
);
assert_eq!(
types
.iter()
.filter(|t| **t == "milestone.validating")
.count(),
1,
"exactly one validation round: {types:?}"
);
assert!(
types.contains(&"mission.completed"),
"mission completed: {types:?}"
);
let (summary, detail) = events
.iter()
.find_map(|e| match &e.kind {
EventKind::OrchestratorDecision {
summary,
detail: Some(detail),
} if summary.starts_with("waived") => Some((summary.clone(), detail.clone())),
_ => None,
})
.expect("a waive orchestrator.decision exists");
assert!(
summary.contains("waived 1 finding(s)"),
"summary: {summary}"
);
assert!(
summary.contains("part 1 works"),
"summary names the subject: {summary}"
);
assert!(
detail.contains("docstring nitpick"),
"detail carries the reason: {detail}"
);
let report = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("report.md"),
)
.expect("report.md written at completion");
assert!(report.contains("## Validation history"), "{report}");
assert!(
report.contains("[minor] part 1 works"),
"finding listed: {report}"
);
assert!(
report.contains("docstring omits the error case"),
"evidence listed: {report}"
);
assert!(report.contains("Disposition: waived."), "{report}");
assert!(
report.contains("part 1 works: docstring nitpick"),
"waiver reason: {report}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn waive_at_cap_completes_instead_of_blocking() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let major = json!([{
"subject": "part 1 works",
"severity": "major",
"evidence": "the endpoint returns 500 on empty input",
"suggestedFix": "guard empty input"
}]);
let minor = json!([{
"subject": "helper docs",
"severity": "minor",
"evidence": "docstring omits the error case",
"suggestedFix": "extend the docstring"
}]);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
fix_features(1),
dirty_tree_commit_as_is(),
judgement("complete", ""),
waive_reply("helper docs", "cosmetic; outside the contract"),
no_lesson(),
]),
validator_with(major),
worker_pass(),
validator_with(minor),
]));
let cfg = MissionConfig {
skip_functional: false,
max_fix_cycles_per_milestone: 1,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let ms = &engine.state().mission.milestones[0];
assert_eq!(ms.status, MilestoneStatus::Complete);
assert_eq!(ms.fix_cycles, 1, "the waived round consumed no extra cycle");
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
!types.contains(&"milestone.blocked"),
"must not block: {types:?}"
);
assert!(
types.contains(&"mission.completed"),
"mission completed: {types:?}"
);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.starts_with("waived 1 finding(s)") && summary.contains("helper docs")
)));
}
#[tokio::test(flavor = "multi_thread")]
async fn command_assertion_at_final_gate_is_non_waivable() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![assertion(
"a-1",
"the build succeeds",
Some("cd kranz-no-such-dir"),
)];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
waive_reply("a-1", "command not runnable in this environment"),
dirty_tree_commit_as_is(),
judgement("complete", ""),
waive_reply("a-1", "still not runnable"),
]),
worker_pass(), ]));
let cfg = MissionConfig {
max_fix_cycles_per_milestone: 1,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
let decision_summaries: Vec<&str> = events
.iter()
.filter_map(|event| match &event.kind {
EventKind::OrchestratorDecision { summary, .. } => Some(summary.as_str()),
_ => None,
})
.collect();
assert!(
types.contains(&"validation.finding"),
"gate finding surfaced: {types:?}"
);
assert!(
types.contains(&"fixfeature.created"),
"refused waive must synthesize a fix feature: {types:?}"
);
assert!(
types.contains(&"milestone.blocked"),
"second refused waive at the fix-cycle cap must block: {types:?}"
);
assert!(
!types.contains(&"mission.completed"),
"command assertion must not COMPLETE via waive: {types:?}"
);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.contains("refused model waive") && summary.contains("a-1")
)),
"must surface the refuse-waive decision: {types:?}; decisions: {decision_summaries:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn command_broken_assertion_escalates_to_operator() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![assertion(
"a-1",
"the build succeeds",
Some("cd kranz-no-such-dir"),
)];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
command_broken_reply("a-1", "grep can only match pre-change; false negative"),
]),
]));
let cfg = MissionConfig {
max_fix_cycles_per_milestone: 1,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
let ms = &engine.state().mission.milestones[0];
assert_eq!(ms.fix_cycles, 0, "escalation must not spend a fix cycle");
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
!types.contains(&"fixfeature.created"),
"escalation must not synthesize a fix feature: {types:?}"
);
assert!(
!types.contains(&"mission.completed"),
"escalation must not complete the mission: {types:?}"
);
let blocked = events
.iter()
.find_map(|e| match &e.kind {
EventKind::MilestoneBlocked { reason, .. } => Some(reason.clone()),
_ => None,
})
.expect("milestone.blocked event present");
assert!(
blocked.contains("a-1"),
"blocked reason names the assertion id: {blocked}"
);
let lower = blocked.to_lowercase();
assert!(
lower.contains("evidence"),
"blocked reason mentions evidence: {blocked}"
);
assert!(
lower.contains("buggy") || lower.contains("false negative"),
"blocked reason indicates the assertion appears broken: {blocked}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn noncommand_finding_marked_command_broken_does_not_escalate() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let finding = json!([{
"subject": "part 1 works",
"severity": "major",
"evidence": "the endpoint returns 500 on empty input",
"suggestedFix": "guard empty input"
}]);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
command_broken_reply("part 1 works", "wrongly claimed author-broken"),
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_with(finding),
worker_pass(),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let state = engine.state();
let ms = &state.mission.milestones[0];
assert_eq!(
ms.fix_cycles, 1,
"mislabelled escalation must be treated as a normal fix"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
types.contains(&"fixfeature.created"),
"mislabelled escalation must still route to a fix feature: {types:?}"
);
assert!(
!types.contains(&"milestone.blocked"),
"mislabelled escalation must not block: {types:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn declared_pty_script_that_never_executes_cannot_green() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![Assertion {
id: "a-pty".to_string(),
statement: "the REPL echoes input back".to_string(),
check: AssertionCheck::PtyScript,
command: None,
negative_control: None,
pty_script: None,
}];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
command_broken_reply("a-pty", "declared pty-script never executed on this host"),
]),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(
status,
MissionStatus::Blocked,
"a declared pty-script that never executed must not green"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
!types.contains(&"mission.completed"),
"no vacuous green off a green round: {types:?}"
);
let finding = events
.iter()
.find_map(|e| match &e.kind {
EventKind::ValidationFinding { finding, .. } => Some(finding.clone()),
_ => None,
})
.expect("final-gate finding surfaced for the unexecuted pty-script");
assert_eq!(finding.subject, "a-pty");
assert_eq!(
finding.class, "command-assertion",
"non-waivable class: {finding:?}"
);
assert!(
finding.evidence.contains("validation.pty.transcript")
&& finding.evidence.contains("never executed"),
"the finding names the missing verdict: {}",
finding.evidence
);
let blocked = events
.iter()
.find_map(|e| match &e.kind {
EventKind::MilestoneBlocked { reason, .. } => Some(reason.clone()),
_ => None,
})
.expect("milestone.blocked event present");
assert!(
blocked.contains("a-pty"),
"blocked reason names the assertion id: {blocked}"
);
}
#[cfg(unix)]
#[tokio::test]
async fn cancelling_pty_validation_stops_writes_before_mission_unlock() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![Assertion {
id: "a-cancel-pty".into(),
statement: "the target completes".into(),
check: AssertionCheck::PtyScript,
command: None,
negative_control: None,
pty_script: Some(PtyScript {
command: "printf '%s' $$ > pty.pid; exec sleep 30".into(),
steps: vec![PtyStep::Expect {
pattern: "never printed".into(),
regex: false,
timeout_ms: Some(30_000),
}],
timeout_secs: Some(30),
}),
}];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![dirty_tree_commit_as_is(), judgement("complete", "")]),
]));
let mut engine = make_engine(
&backend,
&root,
MissionConfig {
skip_functional: false,
..test_cfg()
},
);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let paths = engine.paths().clone();
let mut run = Box::pin(engine.run());
let pid: i32 = tokio::select! {
status = &mut run => panic!("mission ended before cancellation: {status:?}"),
pid = async {
for _ in 0..1000 {
if let Ok(text) = std::fs::read_to_string(root.join("pty.pid")) {
if let Ok(pid) = text.parse() { return pid; }
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
panic!("PTY validation did not start");
} => pid,
};
let started = std::time::Instant::now();
drop(run);
assert!(started.elapsed() < Duration::from_secs(5));
assert!(paths.lock_file().exists(), "the engine still owns its lock");
assert!(
unsafe { libc::kill(pid, 0) } != 0,
"the PTY target outlived cancellation and could race a resumed mission"
);
drop(engine);
assert!(!paths.lock_file().exists());
let events = read_log(&paths);
assert!(!events
.iter()
.any(|event| matches!(event.kind, EventKind::ValidationPtyTranscript { .. })));
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread")]
async fn declared_pty_script_executes_and_passes_greens() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![Assertion {
id: "a-pty".to_string(),
statement: "the REPL echoes input back".to_string(),
check: AssertionCheck::PtyScript,
command: None,
negative_control: None,
pty_script: Some(PtyScript {
command: "printf '> '; while IFS= read -r line; do case \"$line\" in quit) \
printf 'bye\\n'; exit 0;; *) printf 'echo:%s\\n> ' \"$line\";; esac; done"
.to_string(),
steps: vec![
PtyStep::Expect {
pattern: "> ".to_string(),
regex: false,
timeout_ms: Some(10_000),
},
PtyStep::Send {
text: "hello\n".to_string(),
},
PtyStep::Expect {
pattern: "echo:hello".to_string(),
regex: false,
timeout_ms: Some(10_000),
},
PtyStep::Send {
text: "quit\n".to_string(),
},
PtyStep::Expect {
pattern: "bye".to_string(),
regex: false,
timeout_ms: Some(10_000),
},
],
timeout_secs: Some(30),
}),
}];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(
status,
MissionStatus::Complete,
"a declared pty-script that executed and passed still greens"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let verdict = events.iter().find_map(|e| match &e.kind {
EventKind::ValidationPtyTranscript {
assertion_id,
verdict,
..
} if assertion_id == "a-pty" => Some(*verdict),
_ => None,
});
assert_eq!(
verdict,
Some(kranz_engine::gate::GateVerdict::Pass),
"the executed session's verdict is on the log"
);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary == "pty-script assertions not re-run at the final gate"
)));
assert!(
!events.iter().any(|e| matches!(
&e.kind,
EventKind::ValidationFinding { finding, .. } if finding.subject == "a-pty"
)),
"no finding for an executed pty-script: {:?}",
event_types(&events)
);
let transcripts = paths.runs_dir().join("pty-transcripts");
assert!(
transcripts.is_dir() && std::fs::read_dir(&transcripts).unwrap().next().is_some(),
"transcript artifact written under {}",
transcripts.display()
);
}
#[tokio::test(flavor = "multi_thread")]
async fn capture_turn_error_still_completes_mission() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
MockScript::streaming(vec![mock_init("orch-session"), mock_result_text("ready")])
.responding(vec![
vec![
mock_text(&dirty_tree_commit_as_is()),
mock_result_text(&dirty_tree_commit_as_is()),
],
vec![
mock_text(&judgement("complete", "")),
mock_result_text(&judgement("complete", "")),
],
vec![mock_text("boom"), mock_result_error("boom")],
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(
status,
MissionStatus::Complete,
"completion must proceed despite the capture-turn error"
);
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
types.contains(&"mission.completed"),
"mission completed: {types:?}"
);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. } if summary == "no cross-mission lesson captured"
)));
assert!(
!root
.join(".kranz")
.join("lessons")
.join(format!("{mission_id}.md"))
.exists(),
"no lesson file written when the capture turn errors"
);
let subject = raw_git(&root, &["log", "-1", "--format=%s"]);
assert_eq!(
subject.trim(),
format!("[kranz] mission report for {mission_id}")
);
}
#[tokio::test(flavor = "multi_thread")]
async fn retry_retains_prior_checkpoint_commit_receipts() {
if !setup() {
return;
}
for isolation in [WorkerIsolation::Checkout, WorkerIsolation::Worktree] {
for succeeds in [true, false] {
let (_dir, root) = init_repo();
let partial = MockScript::single_shot_json(&json!({
"result": "partial",
"summary": "implementation left for the next attempt to verify",
"commits": []
}))
.writes_file("attempt.txt", "work from the first attempt\n");
let mut replies = vec![dirty_tree_commit_as_is()];
if succeeds {
replies.push(judgement("complete", "verified the retained work"));
}
replies.push(no_lesson());
let backend = Arc::new(MockBackend::with_scripts(vec![
partial,
orch_script(replies),
if succeeds {
worker_pass_no_write()
} else {
worker_fail()
},
]));
let cfg = MissionConfig {
max_respawns: 1,
worker_isolation: isolation,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
timeout(TEST_TIMEOUT, engine.run()).await.unwrap().unwrap();
let feature = &engine.state().mission.milestones[0].features[0];
assert_eq!(
feature.status,
if succeeds {
FeatureStatus::Complete
} else {
FeatureStatus::Failed
}
);
let committed = raw_git(
&root,
&[
"log",
&engine.state().mission.mission_branch,
"--format=%H",
"--",
"attempt.txt",
],
);
assert_eq!(committed.lines().count(), 1);
assert_eq!(feature.commits.len(), 1, "{isolation:?}, succeeds={succeeds}: prior attempt's checkpoint must remain attributed");
assert_eq!(
feature.commits[0].split_whitespace().next(),
Some(committed.trim())
);
}
}
}
#[tokio::test(flavor = "multi_thread")]
async fn respawn_bounded_fails_feature_then_mission_continues() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_fail(), worker_fail(), worker_pass(), orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""), no_lesson(),
]),
]));
let cfg = MissionConfig {
max_respawns: 1,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(2, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let ms = &engine.state().mission.milestones[0];
assert_eq!(ms.features[0].status, FeatureStatus::Failed);
assert_eq!(ms.features[0].respawns, 1, "exactly one respawn honoured");
assert_eq!(ms.features[0].worker_runs.len(), 2, "two worker runs total");
assert_eq!(ms.features[1].status, FeatureStatus::Complete);
let specs = backend.started_specs();
match &specs[1].prompt {
PromptMode::SingleShot(task) => {
assert!(
task.contains("worker run not trusted"),
"runner failure guidance passed: {task}"
)
}
other => panic!("respawned worker must be single-shot, got {other:?}"),
}
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::FeatureFailed { feature_id, reason, .. }
if feature_id == "f-1-1" && reason.contains("respawn budget exhausted")
)));
}
#[tokio::test(flavor = "multi_thread")]
async fn worker_auth_death_blocks_milestone_without_burning_respawn_budget() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_auth_death(), ]));
let cfg = MissionConfig {
max_respawns: 3,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
let ms = &engine.state().mission.milestones[0];
assert_eq!(
ms.features[0].status,
FeatureStatus::Active,
"the feature stays active (re-runs on re-auth), not failed"
);
assert_eq!(
ms.features[0].respawns, 0,
"an auth death must not consume the respawn budget"
);
assert_eq!(
ms.features[0].worker_runs.len(),
1,
"exactly one worker run — no respawn was spawned"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::MilestoneBlocked { reason, .. }
if reason.contains("unauthenticated") && reason.contains("claude")
)),
"a milestone.blocked naming the backend and the re-auth action"
);
assert!(
!events.iter().any(|e| matches!(
&e.kind,
EventKind::FeatureFailed { feature_id, .. } if feature_id == "f-1-1"
)),
"an auth death must not fail the feature"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn pause_resume_and_user_message_flow() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
orch_script(vec![
"Acknowledged — I'll fold the request into the remaining feature.".to_string(),
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
worker_pass(),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
control::enqueue(&paths, &ControlCommand::Pause).unwrap();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
timeout(TEST_TIMEOUT, async {
loop {
if let Ok(snapshot) = reducer::read_snapshot(&paths.state_file()) {
if snapshot.mission.status == MissionStatus::Paused {
break;
}
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.expect("engine must apply Pause before proceeding");
control::enqueue(
&paths,
&ControlCommand::Msg {
text: "swap feature".to_string(),
interrupt: false,
},
)
.unwrap();
timeout(TEST_TIMEOUT, async {
loop {
if let Ok(snapshot) = reducer::read_snapshot(&paths.state_file()) {
if snapshot.mission.status == MissionStatus::Paused
&& snapshot
.pending_user_messages
.iter()
.any(|text| text == "swap feature")
{
break;
}
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.expect("user message must be applied while paused, before Resume");
control::enqueue(&paths, &ControlCommand::Resume).unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Complete);
let mission_id = engine.mission_id().to_string();
drop(engine);
let report = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("report.md"),
)
.expect("report.md written at completion");
assert!(report.contains("paused)"), "paused time surfaced: {report}");
let events = read_log(&paths);
let paused = seq_of(&events, "mission.paused");
let resumed = seq_of(&events, "mission.resumed");
let user_msg = seq_of(&events, "user.message");
let decision = events
.iter()
.filter(|e| e.kind.type_name() == "orchestrator.decision")
.map(|e| e.seq)
.find(|s| *s > user_msg)
.expect("an orchestrator.decision follows the user message");
assert!(
paused < user_msg,
"paused {paused} before user.message {user_msg}"
);
assert!(
user_msg < resumed,
"user.message {user_msg} queued while paused, before resumed {resumed}"
);
assert!(
resumed < decision,
"resumed {resumed} before the consult decision {decision}"
);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::UserMessage { text, interrupt: false } if text == "swap feature"
)));
let leftover: Vec<String> = std::fs::read_dir(paths.control_dir())
.map(|rd| {
rd.flatten()
.map(|e| e.file_name().to_string_lossy().into_owned())
.collect()
})
.unwrap_or_default();
assert!(
leftover.is_empty(),
"control inbox must be empty after apply: {leftover:?}"
);
}
#[tokio::test]
async fn validator_denial_grant_approved_extends_grants_and_completes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_cmd = "gc audit --deep";
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_denied(denied_cmd),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
let snap = reducer::read_snapshot(&paths.state_file()).unwrap();
let pending = snap.pending_grant_request.expect("parked grant request");
assert_eq!(pending.command, denied_cmd);
assert_eq!(pending.milestone_id, "ms-1");
control::enqueue(
&paths,
&ControlCommand::ApproveGrant {
command: denied_cmd.to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
for expected in [
"grant.requested",
"grant.approved",
"milestone.completed",
"mission.completed",
] {
assert!(types.contains(&expected), "missing {expected}: {types:?}");
}
assert!(events.iter().any(|e| matches!(&e.kind,
EventKind::GrantRequested { command, milestone_id, .. }
if command == denied_cmd && milestone_id == "ms-1")));
assert!(seq_of(&events, "grant.requested") < seq_of(&events, "grant.approved"));
assert!(seq_of(&events, "grant.approved") < seq_of(&events, "milestone.completed"));
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
assert!(
state
.mission
.command_grants
.contains(&denied_cmd.to_string()),
"approved command must join command_grants: {:?}",
state.mission.command_grants
);
assert!(state.pending_grant_request.is_none());
}
#[tokio::test]
async fn validator_denial_grant_denied_blocks_the_milestone() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_cmd = "gc audit --deep";
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![dirty_tree_commit_as_is(), judgement("complete", "")]),
validator_denied(denied_cmd),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
control::enqueue(
&paths,
&ControlCommand::DenyGrant {
command: denied_cmd.to_string(),
reason: "not authorized this run".to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Blocked);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
for expected in ["grant.requested", "grant.denied", "milestone.blocked"] {
assert!(types.contains(&expected), "missing {expected}: {types:?}");
}
assert!(seq_of(&events, "grant.denied") < seq_of(&events, "milestone.blocked"));
let state = reducer::fold(&events).unwrap();
assert!(
state.mission.command_grants.is_empty(),
"deny must never widen command_grants: {:?}",
state.mission.command_grants
);
assert!(state.pending_grant_request.is_none());
}
#[tokio::test]
async fn unanswered_grant_request_times_out_to_denied() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_cmd = "gc audit --deep";
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![dirty_tree_commit_as_is(), judgement("complete", "")]),
validator_denied(denied_cmd),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
engine.set_grant_request_timeout(Duration::from_millis(0));
let paths = engine.paths().clone();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
for expected in ["grant.requested", "grant.denied", "milestone.blocked"] {
assert!(types.contains(&expected), "missing {expected}: {types:?}");
}
assert!(events.iter().any(|e| matches!(&e.kind,
EventKind::GrantDenied { reason, .. } if reason.contains("timed out"))));
let state = reducer::fold(&events).unwrap();
assert!(state.mission.command_grants.is_empty());
}
#[tokio::test]
async fn incidental_denial_on_a_passing_validator_does_not_park() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_denied_but_passing("gc audit --deep"),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
!types.contains(&"grant.requested"),
"a passing validator must not park for a grant: {types:?}"
);
assert!(types.contains(&"milestone.completed"));
assert!(types.contains(&"mission.completed"));
}
#[tokio::test]
async fn grant_request_cap_blocks_instead_of_looping() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_cmd = "gc audit --deep";
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![dirty_tree_commit_as_is(), judgement("complete", "")]),
validator_denied(denied_cmd),
validator_denied(denied_cmd),
validator_denied(denied_cmd),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
engine.set_grant_request_cap(1);
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
control::enqueue(
&paths,
&ControlCommand::ApproveGrant {
command: denied_cmd.to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Blocked);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let requested = events
.iter()
.filter(|e| e.kind.type_name() == "grant.requested")
.count();
assert_eq!(requested, 1, "cap = 1 must offer exactly one grant");
assert!(event_types(&events).contains(&"milestone.blocked"));
}
#[tokio::test]
async fn grant_offered_from_the_claude_retry_when_primary_had_no_command() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_cmd = "gc audit --deep";
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_untrusted_no_denial(),
validator_denied(denied_cmd),
validator_with(json!([])),
]));
let cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
let snap = reducer::read_snapshot(&paths.state_file()).unwrap();
assert_eq!(
snap.pending_grant_request.expect("parked").command,
denied_cmd,
"the grant must name the command the retry surfaced"
);
control::enqueue(
&paths,
&ControlCommand::ApproveGrant {
command: denied_cmd.to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(event_types(&events).contains(&"grant.requested"));
let state = reducer::fold(&events).unwrap();
assert!(state
.mission
.command_grants
.contains(&denied_cmd.to_string()));
}
#[tokio::test]
async fn out_of_contract_write_parks_a_touch_grant_and_approve_completes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let plan = Plan {
goal: GOAL.to_string(),
validation_contract: vec![],
milestones: vec![PlanMilestone {
title: "M1".to_string(),
features: vec![PlanFeature {
title: "feature 1".to_string(),
spec: "build part 1".to_string(),
validation_criteria: vec!["part 1 works".to_string()],
}],
}],
considered_alternatives: None,
command_grants: vec![],
touch_set: vec!["src/**".to_string()],
standards_manifest: None,
reviewer_independence: None,
};
let worker = MockScript::single_shot_json(&json!({
"result": "pass",
"summary": "implemented",
"filesTouched": ["out-of-bounds.txt"],
"testsAdded": [],
"testEvidence": "ok",
"commits": []
}))
.writes_file("out-of-bounds.txt", "written outside the touch-set\n");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker,
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(plan).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
let snap = reducer::read_snapshot(&paths.state_file()).unwrap();
let pending = snap.pending_grant_request.expect("parked touch grant");
assert_eq!(pending.kind, GrantKind::TouchPath);
assert_eq!(pending.command, "out-of-bounds.txt");
control::enqueue(
&paths,
&ControlCommand::ApproveGrant {
command: "out-of-bounds.txt".to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
for expected in [
"grant.requested",
"grant.approved",
"milestone.completed",
"mission.completed",
] {
assert!(types.contains(&expected), "missing {expected}: {types:?}");
}
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
assert!(
state
.mission
.touch_set
.contains(&"out-of-bounds.txt".to_string()),
"approved path must join touch_set: {:?}",
state.mission.touch_set
);
assert!(state.mission.command_grants.is_empty());
}
#[tokio::test]
async fn out_of_contract_write_touch_grant_denied_flows_to_waive() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let plan = Plan {
goal: GOAL.to_string(),
validation_contract: vec![],
milestones: vec![PlanMilestone {
title: "M1".to_string(),
features: vec![PlanFeature {
title: "feature 1".to_string(),
spec: "build part 1".to_string(),
validation_criteria: vec!["part 1 works".to_string()],
}],
}],
considered_alternatives: None,
command_grants: vec![],
touch_set: vec!["src/**".to_string()],
standards_manifest: None,
reviewer_independence: None,
};
let worker = MockScript::single_shot_json(&json!({
"result": "pass",
"summary": "implemented",
"filesTouched": ["out-of-bounds.txt"],
"testsAdded": [],
"testEvidence": "ok",
"commits": []
}))
.writes_file("out-of-bounds.txt", "written outside the touch-set\n");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker,
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
waive_reply("out-of-bounds.txt", "acceptable scratch file"),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(plan).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
control::enqueue(
&paths,
&ControlCommand::DenyGrant {
command: "out-of-bounds.txt".to_string(),
reason: "keep it out of contract".to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(types.contains(&"grant.denied"), "{types:?}");
assert!(!types.contains(&"milestone.blocked"), "{types:?}");
assert!(types.contains(&"validation.finding"), "{types:?}");
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
assert!(
!state
.mission
.touch_set
.contains(&"out-of-bounds.txt".to_string()),
"deny must not widen touch_set: {:?}",
state.mission.touch_set
);
}
#[tokio::test]
async fn worker_deny_grant_lifts_the_rule_and_respawn_completes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_cmd = "git push origin main";
let tool_use = json!({
"type": "assistant",
"message": { "id": "w1", "content": [
{ "type": "tool_use", "name": "Bash", "input": { "command": denied_cmd } }
] }
})
.to_string();
let denied = json!({
"type": "user",
"message": { "role": "user", "content": [
{ "type": "tool_result", "tool_use_id": "t1",
"content": format!("Permission denied: Bash({denied_cmd})"), "is_error": true }
] }
})
.to_string();
let report1 = json!({
"result": "fail", "summary": "blocked from pushing",
"filesTouched": [], "testsAdded": [], "testEvidence": "", "commits": []
});
let mut w1 = vec![mock_init("w1")];
w1.extend(parse_stream_line(&tool_use));
w1.extend(parse_stream_line(&denied));
w1.push(mock_text(&report1.to_string()));
w1.push(mock_result_json(&report1));
let worker_denied = MockScript {
events: w1,
..Default::default()
};
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_denied,
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
let snap = reducer::read_snapshot(&paths.state_file()).unwrap();
let pending = snap
.pending_grant_request
.expect("parked worker-deny grant");
assert_eq!(pending.kind, GrantKind::WorkerDeny);
assert_eq!(pending.command, "Bash(git push*)");
control::enqueue(
&paths,
&ControlCommand::ApproveGrant {
command: "Bash(git push*)".to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
for expected in ["grant.requested", "grant.approved", "mission.completed"] {
assert!(types.contains(&expected), "missing {expected}: {types:?}");
}
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
assert!(
state
.mission
.deny_exceptions
.contains(&"Bash(git push*)".to_string()),
"the lifted rule must join deny_exceptions: {:?}",
state.mission.deny_exceptions
);
}
#[tokio::test]
async fn worker_denial_on_a_passing_run_does_not_park() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_cmd = "git push origin main";
let tool_use = json!({
"type": "assistant",
"message": { "id": "w1", "content": [
{ "type": "tool_use", "name": "Bash", "input": { "command": denied_cmd } }
] }
})
.to_string();
let denied = json!({
"type": "user",
"message": { "role": "user", "content": [
{ "type": "tool_result", "tool_use_id": "t1",
"content": format!("Permission denied: Bash({denied_cmd})"), "is_error": true }
] }
})
.to_string();
let report = json!({
"result": "pass", "summary": "shipped without the push",
"filesTouched": ["delivered.txt"], "testsAdded": [], "testEvidence": "ok", "commits": []
});
let mut w = vec![mock_init("w1")];
w.extend(parse_stream_line(&tool_use));
w.extend(parse_stream_line(&denied));
w.push(mock_text(&report.to_string()));
w.push(mock_result_json(&report));
let worker = MockScript {
events: w,
..Default::default()
}
.writes_file("delivered.txt", "shipped\n");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker,
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
!types.contains(&"grant.requested"),
"a passing worker must not park a deny-lift grant: {types:?}"
);
assert!(types.contains(&"mission.completed"));
}
#[tokio::test]
async fn worker_deny_grant_respawn_does_not_eat_the_respawn_budget() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_cmd = "git push origin main";
let tool_use = json!({
"type": "assistant",
"message": { "id": "w1", "content": [
{ "type": "tool_use", "name": "Bash", "input": { "command": denied_cmd } }
] }
})
.to_string();
let denied = json!({
"type": "user",
"message": { "role": "user", "content": [
{ "type": "tool_result", "tool_use_id": "t1",
"content": format!("Permission denied: Bash({denied_cmd})"), "is_error": true }
] }
})
.to_string();
let report1 = json!({ "result": "fail", "summary": "blocked from pushing" });
let mut w1 = vec![mock_init("w1")];
w1.extend(parse_stream_line(&tool_use));
w1.extend(parse_stream_line(&denied));
w1.push(mock_text(&report1.to_string()));
w1.push(mock_result_json(&report1));
let worker_denied = MockScript {
events: w1,
..Default::default()
};
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_denied,
worker_fail(),
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let cfg = MissionConfig {
max_respawns: 1,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
control::enqueue(
&paths,
&ControlCommand::ApproveGrant {
command: "Bash(git push*)".to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Complete);
let feature = &engine.state().mission.milestones[0].features[0];
assert_eq!(feature.status, FeatureStatus::Complete);
assert_eq!(feature.worker_runs.len(), 3, "initial + 2 respawns");
assert_eq!(feature.respawns, 2, "grant respawn + judgement respawn");
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(types.contains(&"mission.completed"), "{types:?}");
assert!(
!events.iter().any(|e| matches!(&e.kind,
EventKind::FeatureFailed { reason, .. } if reason.contains("respawn budget"))),
"the grant respawn must not exhaust the budget"
);
}
#[cfg(target_os = "macos")]
#[tokio::test]
async fn egress_denial_grant_approved_extends_egress_grants_and_completes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_host = "registry.npmjs.org";
let target = format!("{denied_host}:443");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
validator_egress_denied(denied_host, 443),
validator_with(json!([])),
]));
let mut cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
cfg.validator_functional.sandbox.enforce = SandboxEnforce::FsNet;
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
let snap = reducer::read_snapshot(&paths.state_file()).unwrap();
let pending = snap.pending_grant_request.expect("parked grant request");
assert_eq!(pending.kind, GrantKind::Egress);
assert_eq!(pending.command, target);
assert_eq!(pending.milestone_id, "ms-1");
control::enqueue(
&paths,
&ControlCommand::ApproveGrant {
command: target.clone(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
for expected in [
"grant.requested",
"grant.approved",
"milestone.completed",
"mission.completed",
] {
assert!(types.contains(&expected), "missing {expected}: {types:?}");
}
assert!(events.iter().any(|e| matches!(&e.kind,
EventKind::GrantRequested { kind: GrantKind::Egress, command, milestone_id }
if command == &target && milestone_id == "ms-1")));
assert!(seq_of(&events, "grant.requested") < seq_of(&events, "grant.approved"));
assert!(seq_of(&events, "grant.approved") < seq_of(&events, "milestone.completed"));
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
assert!(
state.mission.egress_grants.contains(&target),
"approved destination must join egress_grants: {:?}",
state.mission.egress_grants
);
assert!(state.mission.command_grants.is_empty());
assert!(state.mission.touch_set.is_empty());
assert!(state.pending_grant_request.is_none());
let revalidated = backend
.started_specs()
.last()
.expect("a re-run validator session started")
.clone();
let sandbox = revalidated.sandbox.expect("fs+net sandbox on the re-run");
assert!(
sandbox.inputs.egress.contains(&target),
"the re-run's proxy allowlist must contain the granted host: {:?}",
sandbox.inputs.egress
);
let content = std::fs::read_to_string(paths.egress_denials_file()).unwrap();
assert_eq!(
content.lines().count(),
1,
"only the pre-grant denial is recorded: {content}"
);
}
#[cfg(target_os = "macos")]
#[tokio::test]
async fn egress_denial_grant_denied_blocks_the_milestone() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let denied_host = "registry.npmjs.org";
let target = format!("{denied_host}:443");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![dirty_tree_commit_as_is(), judgement("complete", "")]),
validator_egress_denied(denied_host, 443),
]));
let mut cfg = MissionConfig {
skip_functional: false,
..test_cfg()
};
cfg.validator_functional.sandbox.enforce = SandboxEnforce::FsNet;
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let paths = engine.paths().clone();
let handle = tokio::spawn(async move {
let result = engine.run().await;
(engine, result)
});
wait_for_pending_grant(&paths).await;
control::enqueue(
&paths,
&ControlCommand::DenyGrant {
command: target.clone(),
reason: "not authorized this run".to_string(),
},
)
.unwrap();
let (engine, result) = timeout(TEST_TIMEOUT, handle)
.await
.expect("run must not hang")
.unwrap();
assert_eq!(result.unwrap(), MissionStatus::Blocked);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
for expected in ["grant.requested", "grant.denied", "milestone.blocked"] {
assert!(types.contains(&expected), "missing {expected}: {types:?}");
}
assert!(seq_of(&events, "grant.denied") < seq_of(&events, "milestone.blocked"));
assert!(
events.iter().any(|e| matches!(&e.kind,
EventKind::MilestoneBlocked { reason, .. }
if reason.starts_with("egress denied:")
&& reason.contains(&target)
&& reason.contains("not authorized this run"))),
"the block must carry the egress denial reason"
);
let state = reducer::fold(&events).unwrap();
assert!(
state.mission.egress_grants.is_empty(),
"deny must never widen egress_grants: {:?}",
state.mission.egress_grants
);
assert!(state.pending_grant_request.is_none());
}
#[test]
fn mission_created_secret_redacts_and_audits() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let secret = "sk-ant-api03-AbCdEf_123-xyz";
let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::with_scripts(vec![]));
let engine = MissionEngine::create(backend, &root, &format!("ship with {secret}"), test_cfg())
.expect("create mission");
let paths = engine.paths().clone();
drop(engine);
let raw_log = std::fs::read_to_string(paths.events_file()).unwrap();
assert!(
!raw_log.contains(secret),
"events.jsonl leaked secret: {raw_log}"
);
assert!(raw_log.contains("[REDACTED]"));
let events = read_log(&paths);
assert!(events.iter().any(|event| {
matches!(
&event.kind,
EventKind::SecretRedacted {
rule_id,
location,
..
} if rule_id == "anthropic-api-key" && location.contains("/payload/goal")
)
}));
}
#[test]
fn create_routes_execution_class_ticket_goal_to_local_executor() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::with_scripts(vec![]));
let ticket = kranz_engine::ticket::Ticket::parse(
"bump-dep",
"\
---
title: Bump a dependency
task-class: execution-class
---
## Goal
Bump the dependency to the latest patch release.
",
)
.expect("parse ticket");
let goal = ticket.mission_goal();
let mut cfg = test_cfg();
cfg.worker.base_url = Some("http://127.0.0.1:8080".to_string());
cfg.worker.context_budget = Some(16_384);
let engine = MissionEngine::create(backend, &root, &goal, cfg).expect("create routed mission");
assert_eq!(engine.state().executor_tier(), ExecutorTier::Local);
assert_eq!(
engine.state().config.worker.backend.as_deref(),
Some("local")
);
assert_ne!(
engine.state().config.validator_scrutiny.backend.as_deref(),
Some("local")
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.contains("executor routed local")
)),
"expected a recorded routing decision; events: {:?}",
event_types(&events)
);
}
#[test]
fn create_leaves_executor_frontier_when_goal_carries_no_task_class() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::with_scripts(vec![]));
let mut cfg = test_cfg();
cfg.worker.base_url = Some("http://127.0.0.1:8080".to_string());
cfg.worker.context_budget = Some(16_384);
let engine = MissionEngine::create(backend, &root, GOAL, cfg).expect("create unrouted mission");
assert_eq!(engine.state().executor_tier(), ExecutorTier::Frontier);
assert_eq!(engine.state().config.worker.backend, None);
}
fn commit_routing_rules(root: &Path, rules_json: &str) {
std::fs::create_dir_all(root.join(".kranz")).unwrap();
std::fs::write(root.join(".kranz").join("routing-rules.json"), rules_json).unwrap();
raw_git(root, &["add", ".kranz/routing-rules.json"]);
raw_git(root, &["commit", "-m", "routing rules"]);
}
fn execution_class_goal() -> String {
kranz_engine::ticket::Ticket::parse(
"bump-dep",
"\
---
title: Bump a dependency
task-class: execution-class
---
## Goal
Bump the dependency to the latest patch release.
",
)
.expect("parse ticket")
.mission_goal()
}
#[test]
fn routing_rules_config_create_loads_base_rules_and_routes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_routing_rules(
&root,
r#"{
"taskClassRules": [
{"taskClass": "docs-class", "tier": "frontier"}
],
"patternRules": [
{"pattern": "execution-*", "tier": "local"}
]
}"#,
);
let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::with_scripts(vec![]));
let mut cfg = test_cfg();
cfg.worker.base_url = Some("http://127.0.0.1:8080".to_string());
cfg.worker.context_budget = Some(16_384);
let engine =
MissionEngine::create(backend, &root, &execution_class_goal(), cfg).expect("create");
assert_eq!(engine.state().executor_tier(), ExecutorTier::Local);
assert_eq!(engine.state().config.routing.task_class_rules.len(), 1);
assert_eq!(engine.state().config.routing.pattern_rules.len(), 1);
assert_eq!(
engine.state().config.routing.pattern_rules[0].pattern,
"execution-*"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.contains("routing rules loaded from .kranz/routing-rules.json (base branch \"main\"): 1 task-class rule(s), 1 pattern rule(s)")
)),
"expected the rules-load record; events: {:?}",
event_types(&events)
);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.contains("executor routed local (routing-table rule)")
)),
"expected the table-rule routing decision; events: {:?}",
event_types(&events)
);
}
#[test]
fn routing_rules_config_invalid_rules_fail_draft_closed_naming_the_rule() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_routing_rules(
&root,
r#"{"taskClassRules": [{"taskClass": "ok-class", "tier": "local"}, {"taskClass": " ", "tier": "frontier"}]}"#,
);
let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::with_scripts(vec![]));
let err = match MissionEngine::create(backend, &root, &execution_class_goal(), test_cfg()) {
Ok(_) => panic!("an invalid rules file must refuse mission creation"),
Err(err) => err,
};
let text = format!("{err}");
assert!(text.contains(".kranz/routing-rules.json"), "{text}");
assert!(text.contains("taskClassRules[1].taskClass"), "{text}");
assert!(text.contains("owner: repo-setup"), "{text}");
assert!(
!root.join(".kranz").join("missions").exists(),
"a refused draft leaves no mission side effects"
);
}
#[test]
fn routing_rules_config_invalid_rules_fail_approve_closed() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_routing_rules(
&root,
r#"{"patternRules": [{"pattern": "*", "tier": "frontier"}]}"#,
);
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
std::fs::write(
root.join(".kranz").join("routing-rules.json"),
r#"{"patternRules": [{"pattern": "", "tier": "frontier"}]}"#,
)
.unwrap();
raw_git(&root, &["add", ".kranz/routing-rules.json"]);
raw_git(&root, &["commit", "-m", "break the routing rules"]);
let err = engine
.approve_plan(simple_plan(1, vec![]))
.expect_err("approve must fail closed on invalid base rules");
let text = format!("{err}");
assert!(text.contains(".kranz/routing-rules.json"), "{text}");
assert!(text.contains("patternRules[0].pattern"), "{text}");
assert!(text.contains("owner: repo-setup"), "{text}");
}
#[test]
fn routing_rules_config_no_file_keeps_legacy_floor_byte_identical() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::with_scripts(vec![]));
let mut cfg = test_cfg();
cfg.worker.base_url = Some("http://127.0.0.1:8080".to_string());
cfg.worker.context_budget = Some(16_384);
let engine =
MissionEngine::create(backend, &root, &execution_class_goal(), cfg).expect("create");
assert!(engine.state().config.routing.is_empty());
assert_eq!(engine.state().executor_tier(), ExecutorTier::Local);
assert_eq!(
engine.state().config.worker.backend.as_deref(),
Some("local")
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
!events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. } if summary.contains("routing rules loaded")
)),
"no file ⇒ no rules-load record"
);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.contains("executor routed local (execution-class)")
)),
"the legacy literal-floor decision is unchanged; events: {:?}",
event_types(&events)
);
}
#[tokio::test(flavor = "multi_thread")]
async fn routing_rules_config_mission_branch_edit_ignored_and_surfaced() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
commit_routing_rules(
&root,
r#"{"taskClassRules": [{"taskClass": "execution-class", "tier": "local"}]}"#,
);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let backend_dyn: Arc<dyn AgentBackend> = backend.clone();
let mut engine = MissionEngine::create(backend_dyn, &root, &execution_class_goal(), test_cfg())
.expect("create routed mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
commit_routing_rules(
&root,
r#"{"patternRules": [{"pattern": "*", "tier": "local"}]}"#,
);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.contains("edits .kranz/routing-rules.json — ignored: routing rules are base-branch-owned")
)),
"expected the branch-edit surface note; decisions: {:?}",
events.iter().filter_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, .. } => Some(summary),
_ => None,
}).collect::<Vec<_>>()
);
let worker_route = events
.iter()
.find_map(|e| match &e.kind {
EventKind::WorkerSpawned {
role: Role::Worker,
executor_route,
..
} => Some(executor_route.clone()),
_ => None,
})
.expect("a worker.spawned with a route record");
let route = worker_route.expect("the worker session carries route provenance");
assert_eq!(route.tier, ExecutorTier::Frontier);
assert_eq!(route.rule.as_deref(), Some("taskClassRules[0]"));
assert!(events.iter().all(|e| match &e.kind {
EventKind::WorkerSpawned {
role: Role::Orchestrator,
executor_route,
..
} => executor_route.is_none(),
_ => true,
}));
}
#[tokio::test(flavor = "multi_thread")]
async fn orchestrator_decision_detail_is_scrubbed() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let token = "ghp_AbCdEfGhIjKlMnOpQrStUvWxYz0123";
let leaky_judgement = json!({
"decision": "complete",
"guidance": "",
"summary": format!("looks good; noticed {token} in the env"),
})
.to_string();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
leaky_judgement,
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let raw_log = std::fs::read_to_string(paths.events_file()).unwrap();
assert!(!raw_log.contains(token), "events.jsonl leaked the token");
let events = read_log(&paths);
let (summary, detail) = events
.iter()
.find_map(|e| match &e.kind {
EventKind::OrchestratorDecision {
summary,
detail: Some(detail),
} if summary.starts_with("judgement for") => Some((summary.clone(), detail.clone())),
_ => None,
})
.expect("a judgement orchestrator.decision with detail exists");
assert!(summary.contains("[REDACTED]"), "summary: {summary}");
assert!(detail.contains("[REDACTED]"), "detail: {detail}");
assert!(!detail.contains(token));
}
#[tokio::test(flavor = "multi_thread")]
async fn resume_retains_prior_checkpoint_commit_receipts() {
if !setup() {
return;
}
for isolation in [WorkerIsolation::Checkout, WorkerIsolation::Worktree] {
let (_dir, root) = init_repo();
let backend1 = Arc::new(MockBackend::with_scripts(vec![
worker_pass_no_write().writes_file("before-crash.txt", "retained work\n"),
orch_script(vec![dirty_tree_commit_as_is()]),
]));
let mut engine = make_engine(
&backend1,
&root,
MissionConfig {
worker_isolation: isolation,
..test_cfg()
},
);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
engine.set_orch_stall_timeout(Duration::from_millis(400));
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
timeout(TEST_TIMEOUT, engine.run())
.await
.expect("first run must not hang")
.expect_err("judgement stalls after the checkpoint");
drop(engine);
let before = reducer::fold(&read_log(&paths)).unwrap();
let feature = &before.mission.milestones[0].features[0];
assert_eq!(feature.status, FeatureStatus::Active);
assert_eq!(feature.commits.len(), 1);
assert!(before.feature_base_shas.contains_key(&feature.id));
let backend2 = Arc::new(MockBackend::with_scripts(vec![
worker_pass_no_write(),
orch_script(vec![judgement("complete", "verified"), no_lesson()]),
]));
let backend2_dyn: Arc<dyn AgentBackend> = backend2;
let mut engine =
MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No).unwrap();
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
assert_eq!(
timeout(TEST_TIMEOUT, engine.run()).await.unwrap().unwrap(),
MissionStatus::Complete
);
let after = &engine.state().mission.milestones[0].features[0];
assert_eq!(after.commits, feature.commits, "{isolation:?}");
assert_eq!(engine.state().feature_base_shas, before.feature_base_shas);
assert_eq!(
raw_git(
&root,
&[
"show",
&format!("{}:before-crash.txt", before.mission.mission_branch)
]
),
"retained work\n"
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn kill_and_resume_completes_on_single_log() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend1 = Arc::new(MockBackend::with_scripts(vec![
worker_pass_no_write(),
orch_script(vec![]), ]));
let mut engine = make_engine(&backend1, &root, test_cfg());
engine.approve_plan(simple_plan(2, vec![])).unwrap();
engine.set_orch_stall_timeout(Duration::from_millis(400));
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
let err = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect_err("phase 1 must error out (simulated crash)");
eprintln!("phase 1 crashed as scripted: {err}");
let phase1_events = {
drop(engine); read_log(&paths)
};
let phase1_orch_sdk_id = phase1_events
.iter()
.find_map(|e| match &e.kind {
EventKind::WorkerSpawned {
role: Role::Orchestrator,
sdk_session_id,
..
} => Some(sdk_session_id.clone()),
_ => None,
})
.expect("phase 1 spawned an orchestrator run");
let state = reducer::fold(&phase1_events).unwrap();
assert_eq!(
state.mission.milestones[0].features[0].status,
FeatureStatus::Active,
"crash left f-1-1 mid-feature"
);
let backend2 = Arc::new(MockBackend::with_scripts(vec![
worker_pass(), orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
worker_pass(), ]));
let backend2_dyn: Arc<dyn AgentBackend> = Arc::clone(&backend2) as Arc<dyn AgentBackend>;
let mut engine = MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No)
.expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
drop(engine);
let specs = backend2.started_specs();
let orch_spec = specs
.iter()
.find(|s| matches!(s.prompt, PromptMode::Streaming(_)))
.expect("phase 2 started a streaming orchestrator session");
assert_eq!(
orch_spec.resume.as_deref(),
Some(phase1_orch_sdk_id.as_str())
);
let events = read_log(&paths);
assert_eq!(events.first().unwrap().seq, 1);
assert_eq!(events.last().unwrap().seq, events.len() as u64);
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
assert!(event_types(&events).contains(&"mission.completed"));
}
#[tokio::test(flavor = "multi_thread")]
async fn force_reseed_reseeds_with_digest_and_plan() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let plan_json = json!({
"goal": GOAL,
"validationContract": [],
"milestones": [{
"title": "M1",
"features": [
{ "title": "feature 1", "spec": "build part 1", "validationCriteria": ["part 1 works"] },
{ "title": "feature 2", "spec": "build part 2", "validationCriteria": ["part 2 works"] }
]
}]
})
.to_string();
let backend = Arc::new(MockBackend::with_scripts(vec![
orch_script(vec![
"Understood. Two features under one milestone; no open questions.".to_string(),
plan_json,
]),
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
worker_pass(),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
let reply = timeout(
TEST_TIMEOUT,
engine.planning_turn("plan two features please"),
)
.await
.expect("planning turn must not hang")
.unwrap();
assert!(reply.contains("no open questions"));
let plan = match timeout(TEST_TIMEOUT, engine.request_plan())
.await
.expect("request_plan must not hang")
.unwrap()
{
PlanRequest::Ready(plan) => plan,
PlanRequest::NotReady(text) => panic!("scripted plan JSON must parse, got: {text}"),
PlanRequest::WrongPlan { reason } => {
panic!("scripted plan JSON must parse, got a wrong-plan escalation: {reason}")
}
};
assert_eq!(plan.milestones.len(), 1);
engine.approve_plan(plan).unwrap();
engine.force_reseed();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let specs = backend.started_specs();
let streaming: Vec<_> = specs
.iter()
.filter(|s| matches!(s.prompt, PromptMode::Streaming(_)))
.collect();
assert_eq!(streaming.len(), 2, "exactly two orchestrator sessions");
assert!(
streaming[1].resume.is_none(),
"re-seed is a fresh session, not a resume"
);
match &streaming[1].prompt {
PromptMode::Streaming(seed) => {
assert!(
seed.starts_with("MISSION m-"),
"digest header first: {seed}"
);
assert!(
seed.contains("APPROVED PLAN (plan.json):"),
"plan.json embedded: {seed}"
);
assert!(
seed.contains("build part 1"),
"plan content present: {seed}"
);
}
other => panic!("orchestrator must be streaming, got {other:?}"),
}
let injected = backend.injected_messages();
let second_orch_injected = &injected[2]; assert!(!second_orch_injected.is_empty());
assert!(
second_orch_injected[0].starts_with("MISSION m-"),
"turn digest prefix: {}",
second_orch_injected[0]
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. } if summary.contains("re-seeded")
)));
}
#[tokio::test(flavor = "multi_thread")]
async fn plan_not_ready_returns_prose() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let plan_json = json!({
"goal": GOAL,
"validationContract": [],
"milestones": [{
"title": "M1",
"features": [
{ "title": "feature 1", "spec": "build part 1", "validationCriteria": ["part 1 works"] }
]
}]
})
.to_string();
let prose_first =
"Before I emit the plan I still need an answer: which database should the demo target?"
.to_string();
let prose_retry =
"I cannot emit the plan yet — please answer the database question first.".to_string();
let backend = Arc::new(MockBackend::with_scripts(vec![orch_script(vec![
prose_first,
prose_retry,
plan_json,
])]));
let mut engine = make_engine(&backend, &root, test_cfg());
let request = timeout(TEST_TIMEOUT, engine.request_plan())
.await
.expect("request_plan must not hang")
.expect("prose replies are a conversational state, not a backend error");
match request {
PlanRequest::NotReady(text) => assert!(
text.contains("database question"),
"NotReady carries the retry turn's prose: {text}"
),
PlanRequest::Ready(plan) => panic!("prose must not parse as a plan: {plan:?}"),
PlanRequest::WrongPlan { reason } => {
panic!("prose must not parse as a wrong-plan escalation: {reason}")
}
}
assert_eq!(
engine.state().mission.status,
MissionStatus::Planning,
"a not-ready plan request leaves the mission in Planning"
);
let request = timeout(TEST_TIMEOUT, engine.request_plan())
.await
.expect("second request_plan must not hang")
.unwrap();
match request {
PlanRequest::Ready(plan) => assert_eq!(plan.milestones.len(), 1),
PlanRequest::NotReady(text) => panic!("scripted plan JSON must parse, got: {text}"),
PlanRequest::WrongPlan { reason } => {
panic!("scripted plan JSON must parse, got a wrong-plan escalation: {reason}")
}
}
}
#[tokio::test(flavor = "multi_thread")]
async fn seed_reply_is_captured_once_for_fresh_and_resumed_sessions() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let seed_text = "Two questions before I plan: which auth flows are in scope, and is the \
dashboard part of this mission?";
let backend = Arc::new(MockBackend::with_scripts(vec![MockScript::streaming(
vec![
mock_init("orch-session"),
mock_text(seed_text),
mock_result_text(seed_text),
],
)
.responding(vec![vec![mock_text("noted"), mock_result_text("noted")]])]));
let mut engine = make_engine(&backend, &root, test_cfg());
let reply = timeout(TEST_TIMEOUT, engine.planning_turn("hello"))
.await
.expect("planning turn must not hang")
.unwrap();
assert_eq!(reply, "noted");
assert_eq!(
engine.take_seed_reply().as_deref(),
Some(seed_text),
"the seed turn's reply is captured, not discarded"
);
assert_eq!(
engine.take_seed_reply(),
None,
"the seed reply is taken exactly once"
);
let mission_id = engine.mission_id().to_string();
drop(engine);
let ack_text = "Acknowledged — resuming the planning conversation.";
let backend2 = Arc::new(MockBackend::with_scripts(vec![MockScript::streaming(
vec![
mock_init("orch-session"),
mock_text(ack_text),
mock_result_text(ack_text),
],
)
.responding(vec![vec![
mock_text("continuing"),
mock_result_text("continuing"),
]])]));
let backend2_dyn: Arc<dyn AgentBackend> = Arc::clone(&backend2) as Arc<dyn AgentBackend>;
let mut engine = MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No)
.expect("resume mission");
let reply = timeout(TEST_TIMEOUT, engine.planning_turn("go on"))
.await
.expect("resumed planning turn must not hang")
.unwrap();
assert_eq!(reply, "continuing");
assert_eq!(
engine.take_seed_reply().as_deref(),
Some(ack_text),
"the resume-ack seed reply is captured"
);
assert_eq!(engine.take_seed_reply(), None);
let specs = backend2.started_specs();
let orch_spec = specs
.iter()
.find(|s| matches!(s.prompt, PromptMode::Streaming(_)))
.expect("the resumed engine started a streaming orchestrator session");
assert!(
orch_spec.resume.is_some(),
"resume-ack path expected (--resume set)"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn plan_approval_writes_plan_branch_and_commit() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let mission_id = engine.mission_id().to_string();
let mut plan = simple_plan(
1,
vec![
assertion("", "tests pass", Some("cargo test")),
assertion("", "docs updated", None),
],
);
plan.milestones[0].title = "Milestone One".to_string();
engine.approve_plan(plan).unwrap();
let paths = engine.paths().clone();
let plan_text = std::fs::read_to_string(paths.plan_file()).expect("plan.json written");
let written: Plan = serde_json::from_str(&plan_text).expect("plan.json parses");
let ids: Vec<&str> = written
.validation_contract
.iter()
.map(|a| a.id.as_str())
.collect();
assert_eq!(ids, vec!["a-1", "a-2"]);
let branch = raw_git(&root, &["rev-parse", "--abbrev-ref", "HEAD"]);
assert_eq!(branch.trim(), format!("kranz/mission-{mission_id}"));
let subject = raw_git(&root, &["log", "-1", "--format=%s"]);
assert_eq!(
subject.trim(),
format!("[kranz] approved plan for {mission_id}")
);
let files = raw_git(&root, &["show", "--name-only", "--format=", "HEAD"]);
let mut files: Vec<&str> = files.lines().filter(|l| !l.trim().is_empty()).collect();
files.sort_unstable();
assert_eq!(
files,
vec![
".kranz/missions/index.md".to_string(),
format!(".kranz/missions/{mission_id}/plan.json"),
format!(".kranz/missions/{mission_id}/plan.md"),
],
"the approval commit contains plan.json + plan.md + the missions index"
);
let index =
std::fs::read_to_string(root.join(".kranz").join("missions").join("index.md")).unwrap();
assert!(index.starts_with("# Kranz missions"), "{index}");
assert!(
index.contains(&format!("[{mission_id}]({mission_id}/plan.md)")),
"{index}"
);
let md = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("plan.md"),
)
.expect("plan.md written");
assert!(
md.starts_with(&format!("# Mission plan — {mission_id}")),
"{md}"
);
let calibration = cost::calibrate(&root);
assert_eq!(calibration.missions_used, 0);
let expected_estimate = cost::estimate(&written, &test_cfg(), &calibration.params);
assert!(md.contains("## Cost estimate"), "{md}");
assert!(
md.contains(&format!("${:.2}", expected_estimate.low_usd)),
"{md}"
);
assert!(
md.contains(&format!("${:.2}", expected_estimate.expected_usd)),
"{md}"
);
assert!(
md.contains(&format!("${:.2}", expected_estimate.high_usd)),
"{md}"
);
assert!(
md.contains("built-in defaults — no completed missions yet"),
"{md}"
);
assert!(md.contains("## Validation contract"), "{md}");
assert!(md.contains("**[a-1]**"), "{md}");
assert!(md.contains("## Milestone 1 —"), "{md}");
assert!(md.contains("Done when:"), "{md}");
assert!(
md.lines().all(|line| !line.ends_with([' ', '\t'])),
"plan.md must not contain trailing whitespace:\n{md}"
);
assert!(
md.ends_with('\n') && !md.ends_with("\n\n"),
"plan.md must end with exactly one newline:\n{md}"
);
let main_subject = raw_git(&root, &["log", "-1", "--format=%s", "main"]);
assert_eq!(main_subject.trim(), "seed");
let state = engine.state();
assert_eq!(state.mission.status, MissionStatus::Approved);
assert_eq!(state.mission.milestones.len(), 1);
assert_eq!(state.mission.milestones[0].id, "ms-1");
assert_eq!(state.mission.milestones[0].features[0].id, "f-1-1");
assert_eq!(state.mission.validation_contract.len(), 2);
drop(engine);
let events = read_log(&paths);
assert!(event_types(&events).contains(&"plan.approved"));
let backend_dyn: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
let mut engine = MissionEngine::resume(backend_dyn, &root, &mission_id, LockForce::No).unwrap();
let err = engine.approve_plan(simple_plan(1, vec![])).unwrap_err();
assert!(err.to_string().contains("Planning"), "got: {err}");
}
#[tokio::test(flavor = "multi_thread")]
async fn approval_lint_surfaces_suspects_in_plan_md_and_decision() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let plan = simple_plan(
1,
vec![
assertion("", "vacuous assertion", Some("exit 0")),
assertion("", "not-yet-landed assertion", Some("exit 1")),
],
);
engine.approve_plan(plan).unwrap();
let paths = engine.paths().clone();
let md = std::fs::read_to_string(paths.plan_md_file()).expect("plan.md written");
assert!(md.contains("## Contract lint"), "{md}");
assert!(
md.contains("author-bug suspects (already pass / no verdict on the untouched base)"),
"{md}"
);
assert!(md.contains("[a-1] exit 0"), "{md}");
assert!(
md.contains("base-expected-to-fail (benign): [a-2] exit 1"),
"{md}"
);
drop(engine);
let events = read_log(&paths);
let decision = events.iter().find_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, detail }
if summary.contains("contract lint") =>
{
Some((summary.clone(), detail.clone()))
}
_ => None,
});
let (summary, detail) = decision.expect("contract lint orchestrator.decision emitted");
assert!(summary.contains("1 author-bug suspect"), "{summary}");
let detail = detail.expect("decision carries the full lint summary");
assert!(detail.contains("[a-1] exit 0"), "{detail}");
assert!(detail.contains("[a-2] exit 1"), "{detail}");
}
#[tokio::test(flavor = "multi_thread")]
async fn approval_lint_uses_disposable_base_and_preserves_primary_checkout() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let plan = simple_plan(
1,
vec![assertion(
"",
"hostile approval assertion",
Some(&write_line_cmd("tampered", "README.md")),
)],
);
engine.approve_plan(plan).unwrap();
let paths = engine.paths().clone();
assert_eq!(
std::fs::read_to_string(root.join("README.md")).unwrap(),
"seed\n",
"the approval command must not alter the primary checkout"
);
assert!(
!paths
.runs_dir()
.join("approval-contract-lint-worktree")
.exists(),
"the disposable approval worktree is cleaned"
);
assert!(
!raw_git(&root, &["worktree", "list", "--porcelain"])
.contains("approval-contract-lint-worktree"),
"the disposable worktree registration is pruned"
);
let md = std::fs::read_to_string(paths.plan_md_file()).expect("plan.md written");
assert!(
md.contains("author-bug suspects (already pass / no verdict on the untouched base)"),
"{md}"
);
assert!(
!md.contains("working tree with uncommitted changes"),
"the pinned disposable tree is clean: {md}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn approve_refuses_preexisting_mission_branch_commits() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let branch = engine.state().mission.mission_branch.clone();
raw_git(&root, &["checkout", "-b", &branch]);
std::fs::write(root.join("planted.txt"), "not approved\n").unwrap();
raw_git(&root, &["add", "planted.txt"]);
raw_git(&root, &["commit", "-m", "planted mission commit"]);
raw_git(&root, &["checkout", "main"]);
let error = engine
.approve_plan(simple_plan(1, vec![]))
.expect_err("pre-existing mission bytes must fail closed");
assert!(
error
.to_string()
.contains("refusing to approve pre-existing commits"),
"{error}"
);
assert!(
!read_log(engine.paths())
.iter()
.any(|event| matches!(event.kind, EventKind::PlanApproved { .. })),
"the refusal emits no approval authority"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn approval_lint_never_blocks() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let plan = simple_plan(
1,
vec![
assertion("", "vacuous assertion one", Some("exit 0")),
assertion("", "vacuous assertion two", Some("cd .")),
],
);
engine.approve_plan(plan).unwrap();
assert_eq!(engine.state().mission.status, MissionStatus::Approved);
}
#[tokio::test(flavor = "multi_thread")]
async fn approval_lint_no_nested_runtime_panic() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let plan = simple_plan(
1,
vec![
assertion("", "passes on base", Some("exit 0")),
assertion("", "fails on base", Some("exit 1")),
],
);
engine.approve_plan(plan).unwrap();
assert_eq!(engine.state().mission.status, MissionStatus::Approved);
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread")]
async fn contract_gate_named_verdicts_reach_plan_md_and_decision() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let plan = simple_plan(
1,
vec![
assertion(
"",
"marker absent",
Some("! grep -q landed-marker kranz-no-such-file.txt"),
),
assertion("", "not-yet-landed assertion", Some("exit 1")),
],
);
engine.approve_plan(plan).unwrap();
let paths = engine.paths().clone();
let md = std::fs::read_to_string(paths.plan_md_file()).expect("plan.md written");
assert!(md.contains("named contract gates"), "{md}");
assert!(md.contains("wrong-polarity: FAIL"), "{md}");
assert!(md.contains("passes-on-base: FAIL"), "{md}");
assert!(md.contains("vacuous-filter: PASS"), "{md}");
assert!(md.contains("env-sensitive: PASS"), "{md}");
drop(engine);
let events = read_log(&paths);
let decision = events.iter().find_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, detail }
if summary.contains("contract lint") =>
{
Some((summary.clone(), detail.clone()))
}
_ => None,
});
let (summary, detail) = decision.expect("contract lint orchestrator.decision emitted");
assert!(summary.contains("1 author-bug suspect"), "{summary}");
assert!(
summary.contains("named contract gate(s) failed: wrong-polarity, passes-on-base"),
"{summary}"
);
let detail = detail.expect("decision carries the lint summary and gate verdicts");
assert!(detail.contains("wrong-polarity: FAIL"), "{detail}");
assert!(detail.contains("passes-on-base: FAIL"), "{detail}");
assert!(detail.contains("[a-1]"), "{detail}");
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread")]
async fn contract_gate_final_gate_decision_names_vacuous_green() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![assertion(
"",
"marker absent",
Some("! grep -q landed-marker kranz-no-such-file.txt"),
)];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let decision = events.iter().find_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, detail }
if summary.contains("contract gates (final gate)") =>
{
Some((summary.clone(), detail.clone()))
}
_ => None,
});
let (summary, detail) = decision.expect("final-gate contract-gate decision emitted");
assert!(summary.contains("wrong-polarity"), "{summary}");
let detail = detail.expect("decision carries the gate verdicts");
assert!(detail.contains("wrong-polarity: FAIL"), "{detail}");
assert!(detail.contains("[a-1]"), "{detail}");
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread")]
async fn gate_result_events_record_the_approval_ladder() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
let plan = simple_plan(
1,
vec![
assertion(
"",
"marker absent",
Some("! grep -q landed-marker kranz-no-such-file.txt"),
),
assertion("", "not-yet-landed assertion", Some("exit 1")),
],
);
engine.approve_plan(plan).unwrap();
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let ladder: Vec<(String, u32, String, String)> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::GateResult {
gate,
surface,
kind,
index,
verdict,
artefact_ref,
..
} => {
assert_eq!(
*surface,
kranz_engine::gate::GateSurface::Approval,
"approval gates carry the approval surface"
);
assert_eq!(*kind, kranz_engine::gate::GateKind::Deterministic);
Some((
gate.clone(),
*index,
serde_json::to_value(verdict)
.unwrap()
.as_str()
.unwrap()
.to_string(),
artefact_ref.clone(),
))
}
_ => None,
})
.collect();
assert_eq!(
ladder,
vec![
(
"vacuous-filter".to_string(),
0,
"pass".to_string(),
"contract gate vacuous-filter".to_string()
),
(
"wrong-polarity".to_string(),
1,
"fail".to_string(),
"contract gate wrong-polarity".to_string()
),
(
"passes-on-base".to_string(),
2,
"fail".to_string(),
"contract gate passes-on-base".to_string()
),
(
"env-sensitive".to_string(),
3,
"pass".to_string(),
"contract gate env-sensitive".to_string()
),
],
"one gate.result per approval gate, in pipeline order"
);
let wrong_polarity = events.iter().find_map(|e| match &e.kind {
EventKind::GateResult {
gate,
artefact_detail,
..
} if gate == "wrong-polarity" => artefact_detail.clone(),
_ => None,
});
assert!(
wrong_polarity
.as_deref()
.unwrap_or_default()
.contains("negated grep targets missing path"),
"{wrong_polarity:?}"
);
let seq_of = |pred: &dyn Fn(&Event) -> bool| {
events
.iter()
.find(|e| pred(e))
.map(|e| e.seq)
.expect("event present")
};
let approved_seq = seq_of(&|e| matches!(e.kind, EventKind::PlanApproved { .. }));
let first_gate_seq = seq_of(&|e| matches!(e.kind, EventKind::GateResult { .. }));
let decision_seq = seq_of(
&|e| matches!(&e.kind, EventKind::OrchestratorDecision { summary, .. } if summary.contains("contract lint")),
);
assert!(
approved_seq < first_gate_seq,
"{approved_seq} < {first_gate_seq}"
);
assert!(
first_gate_seq < decision_seq,
"{first_gate_seq} < {decision_seq}"
);
}
#[cfg(unix)]
#[tokio::test(flavor = "multi_thread")]
async fn gate_result_events_record_the_final_gate_ladder() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![assertion(
"",
"marker absent",
Some("! grep -q landed-marker kranz-no-such-file.txt"),
)];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let ladder: Vec<(String, u32, String)> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::GateResult {
gate,
surface,
index,
verdict,
..
} if *surface == kranz_engine::gate::GateSurface::FinalGate => Some((
gate.clone(),
*index,
serde_json::to_value(verdict)
.unwrap()
.as_str()
.unwrap()
.to_string(),
)),
_ => None,
})
.collect();
assert_eq!(
ladder,
vec![
("vacuous-filter".to_string(), 0, "pass".to_string()),
("wrong-polarity".to_string(), 1, "fail".to_string()),
("env-sensitive".to_string(), 2, "pass".to_string()),
],
"the final-gate floor, in pipeline order, passes and failures alike"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn approve_plan_requires_considered_alternatives_for_large_scope() {
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let cfg = MissionConfig {
considered_alternatives_feature_threshold: 2,
considered_alternatives_touch_set_threshold: 0,
considered_alternatives_high_usd_threshold: 0.0,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
let err = engine.approve_plan(simple_plan(2, vec![])).unwrap_err();
assert!(
err.to_string().contains("considered alternatives required"),
"large-scope refusal names the missing review material: {err}"
);
assert_eq!(engine.state().mission.status, MissionStatus::Planning);
}
#[tokio::test(flavor = "multi_thread")]
async fn approve_plan_persists_considered_alternatives_when_required() {
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let cfg = MissionConfig {
considered_alternatives_feature_threshold: 2,
considered_alternatives_touch_set_threshold: 0,
considered_alternatives_high_usd_threshold: 0.0,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
let mission_id = engine.mission_id().to_string();
let mut plan = simple_plan(2, vec![]);
plan.considered_alternatives = Some(considered_alternatives());
engine.approve_plan(plan).unwrap();
let plan_text = std::fs::read_to_string(engine.paths().plan_file()).unwrap();
let written: Plan = serde_json::from_str(&plan_text).unwrap();
assert!(
written.considered_alternatives.is_some(),
"plan.json carries the review section"
);
let md = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("plan.md"),
)
.unwrap();
assert!(md.contains("## Considered alternatives"), "{md}");
assert!(md.contains("big-bang rewrite"), "{md}");
}
#[tokio::test(flavor = "multi_thread")]
async fn plan_md_cost_estimate_uses_calibration_once_a_mission_completes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut first = make_engine(&backend, &root, test_cfg());
first.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, first.run())
.await
.expect("first mission must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
drop(first);
raw_git(&root, &["checkout", "main"]);
let backend2 = Arc::new(MockBackend::new());
let mut second = make_engine(&backend2, &root, test_cfg());
let plan = simple_plan(1, vec![]);
let calibration = cost::calibrate(&root);
assert_eq!(
calibration.missions_used, 1,
"the completed first mission must calibrate the second's estimate"
);
let expected_estimate = cost::estimate(&plan, &test_cfg(), &calibration.params);
second.approve_plan(plan).unwrap();
let mission_id = second.mission_id().to_string();
let md = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("plan.md"),
)
.expect("plan.md written");
assert!(md.contains("based on 1 completed mission(s)"), "{md}");
assert!(
!md.contains("built-in defaults — no completed missions yet"),
"{md}"
);
assert!(
md.contains(&format!("${:.2}", expected_estimate.expected_usd)),
"{md}"
);
}
#[test]
fn approval_records_base_branch_sha() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let base_tip_before = raw_git(&root, &["rev-parse", "main"]).trim().to_string();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let state = engine.state();
assert_eq!(
state.mission.base_sha.as_deref(),
Some(base_tip_before.as_str()),
"folded state.mission.base_sha must equal the base branch tip recorded before approval"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let approved = events
.iter()
.find(|e| e.kind.type_name() == "plan.approved")
.expect("plan.approved event must be on the log");
match &approved.kind {
EventKind::PlanApproved { base_sha, .. } => {
assert_eq!(base_sha.as_deref(), Some(base_tip_before.as_str()));
}
other => panic!("expected PlanApproved, got: {other:?}"),
}
let base_tip_after = raw_git(&root, &["rev-parse", "main"]).trim().to_string();
assert_eq!(base_tip_after, base_tip_before, "base branch must not move");
}
#[tokio::test(flavor = "multi_thread")]
async fn base_sha_reaches_worker_and_validator_env() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let base_tip_before = raw_git(&root, &["rev-parse", "main"]).trim().to_string();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
verdicts_pass(&["a-2"]),
no_lesson(),
]),
validator_with(json!([])),
]));
let contract = vec![assertion("a-2", "error messages are actionable", None)];
let cfg = MissionConfig {
skip_scrutiny: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
assert_eq!(
engine.state().mission.base_sha.as_deref(),
Some(base_tip_before.as_str())
);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let specs = backend.started_specs();
let worker_spec = specs
.iter()
.find(|s| matches!(s.prompt, PromptMode::SingleShot(ref t) if t.contains("Implement feature")))
.expect("a worker spec was started");
assert_eq!(
worker_spec.env.get("KRANZ_BASE_SHA").map(String::as_str),
Some(base_tip_before.as_str()),
"worker env carries the recorded base sha"
);
let validator_spec = specs
.iter()
.find(|s| matches!(s.prompt, PromptMode::SingleShot(ref t) if t.contains("Validate milestone")))
.expect("a validator spec was started");
assert_eq!(
validator_spec.env.get("KRANZ_BASE_SHA").map(String::as_str),
Some(base_tip_before.as_str()),
"validator env carries the recorded base sha"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn worker_file_write_lands_a_non_meta_commit_and_mission_completes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let base_tip_before = raw_git(&root, &["rev-parse", "main"]).trim().to_string();
let worker_writes = MockScript::single_shot_json(&json!({
"result": "pass",
"summary": "implemented and tested",
"filesTouched": ["feature.txt"],
"testsAdded": [],
"testEvidence": "all green",
"commits": []
}))
.writes_file("feature.txt", "delivered by the mock worker\n");
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_writes,
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
verdicts_pass(&["a-2"]),
no_lesson(),
]),
validator_with(json!([])),
]));
let contract = vec![assertion("a-2", "error messages are actionable", None)];
let cfg = MissionConfig {
skip_scrutiny: false,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let mission_branch = engine.state().mission.mission_branch.clone();
let log = raw_git(
&root,
&[
"log",
&format!("{base_tip_before}..{mission_branch}"),
"--format=%s",
],
);
let subjects: Vec<&str> = log.lines().collect();
assert!(
subjects
.iter()
.any(|s| !kranz_engine::contract_sweep::is_meta_commit(s)),
"expected at least one non-meta commit on the mission branch, got: {subjects:?}"
);
assert!(
subjects
.iter()
.any(|s| s.contains("feature.txt") || s.contains("checkpoint") || s.contains("f-1-1")),
"the worker's dirty tree should have produced the feature checkpoint commit: {subjects:?}"
);
let feature_file = std::fs::read_to_string(root.join("feature.txt"))
.expect("the mock worker's write should survive on the mission branch");
assert_eq!(feature_file, "delivered by the mock worker\n");
}
#[test]
fn mission_index_upserts_by_id() {
use kranz_engine::planning::upsert_mission_index;
let d1 = chrono::NaiveDate::from_ymd_opt(2026, 7, 2).unwrap();
let d2 = chrono::NaiveDate::from_ymd_opt(2026, 7, 3).unwrap();
let one = upsert_mission_index("", "m-aaa", "first goal", d1);
assert!(one.starts_with("# Kranz missions"), "{one}");
assert!(
one.contains("- 2026-07-02 · [m-aaa](m-aaa/plan.md) — first goal"),
"{one}"
);
let two = upsert_mission_index(&one, "m-bbb", "second goal\nwith newline", d2);
assert!(two.contains("[m-aaa]("), "{two}");
assert!(
two.contains("- 2026-07-03 · [m-bbb](m-bbb/plan.md) — second goal with newline"),
"{two}"
);
assert!(
two.find("[m-aaa](").unwrap() < two.find("[m-bbb](").unwrap(),
"newest last"
);
let re = upsert_mission_index(&two, "m-aaa", "first goal, re-planned", d2);
assert_eq!(
re.matches("[m-aaa](").count(),
1,
"no duplicate on re-approval: {re}"
);
assert!(re.contains("first goal, re-planned"), "{re}");
}
#[test]
fn mission_index_report_link_appends_once() {
use kranz_engine::mission_catalog::mark_mission_index_report;
use kranz_engine::planning::upsert_mission_index;
let d = chrono::NaiveDate::from_ymd_opt(2026, 7, 3).unwrap();
let index = upsert_mission_index("", "m-aaa", "goal", d);
let index = upsert_mission_index(&index, "m-bbb", "other goal", d);
let marked = mark_mission_index_report(&index, "m-aaa");
assert!(
marked.contains("- 2026-07-03 · [m-aaa](m-aaa/plan.md) — goal · [report](m-aaa/report.md)"),
"{marked}"
);
assert!(
!marked.contains("[report](m-bbb/report.md)"),
"only the named mission: {marked}"
);
let again = mark_mission_index_report(&marked, "m-aaa");
assert_eq!(again, marked, "idempotent");
let unknown = mark_mission_index_report(&marked, "m-zzz");
assert_eq!(unknown, marked, "unknown id leaves the index unchanged");
}
#[test]
fn mission_index_prune_removes_only_named_line() {
use kranz_engine::mission_catalog::prune_mission_index;
use kranz_engine::planning::upsert_mission_index;
let d = chrono::NaiveDate::from_ymd_opt(2026, 7, 3).unwrap();
let index = upsert_mission_index("", "m-aaa", "first goal", d);
let index = upsert_mission_index(&index, "m-bbb", "second goal", d);
let pruned = prune_mission_index(&index, "m-aaa");
assert!(!pruned.contains("[m-aaa]("), "{pruned}");
assert!(pruned.contains("[m-bbb]("), "{pruned}");
assert!(pruned.starts_with("# Kranz missions"), "{pruned}");
assert!(
pruned.contains("- 2026-07-03 · [m-bbb](m-bbb/plan.md) — second goal"),
"{pruned}"
);
let noop = prune_mission_index(&pruned, "m-zzz");
assert_eq!(noop, pruned, "unknown id leaves the index unchanged");
let empty = prune_mission_index("", "m-aaa");
assert_eq!(empty, "", "empty input returns unchanged");
}
#[test]
fn mission_index_ids_lists_ids_in_order() {
use kranz_engine::mission_catalog::{mark_mission_index_report, mission_index_ids};
use kranz_engine::planning::upsert_mission_index;
let d = chrono::NaiveDate::from_ymd_opt(2026, 7, 3).unwrap();
let index = upsert_mission_index("", "m-aaa", "first goal", d);
let index = upsert_mission_index(&index, "m-bbb", "second goal", d);
let index = mark_mission_index_report(&index, "m-aaa");
assert_eq!(
mission_index_ids(&index),
vec!["m-aaa".to_string(), "m-bbb".to_string()]
);
}
#[test]
fn delete_prunes_missions_index() {
use kranz_engine::mission_catalog::prune_mission_index_file;
use kranz_engine::planning::upsert_mission_index;
let dir = tempfile::tempdir().expect("create tempdir");
let root = std::fs::canonicalize(dir.path()).expect("canonicalize repo root");
let d = chrono::NaiveDate::from_ymd_opt(2026, 7, 3).unwrap();
let index_dir = root.join(".kranz").join("missions");
std::fs::create_dir_all(&index_dir).unwrap();
let index_path = index_dir.join("index.md");
let index = upsert_mission_index("", "m-aaa", "first goal", d);
let index = upsert_mission_index(&index, "m-bbb", "second goal", d);
std::fs::write(&index_path, &index).unwrap();
prune_mission_index_file(&root, "m-aaa");
let after = std::fs::read_to_string(&index_path).unwrap();
assert!(!after.contains("[m-aaa]("), "{after}");
assert!(after.contains("[m-bbb]("), "{after}");
}
#[test]
fn delete_prunes_missions_index_missing_file_is_noop() {
use kranz_engine::mission_catalog::prune_mission_index_file;
let dir = tempfile::tempdir().expect("create tempdir");
let root = std::fs::canonicalize(dir.path()).expect("canonicalize repo root");
prune_mission_index_file(&root, "m-aaa");
assert!(!root
.join(".kranz")
.join("missions")
.join("index.md")
.exists());
}
#[tokio::test(flavor = "multi_thread")]
async fn abandon_planning_mission_sets_abandoned_status() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let engine = make_engine(&backend, &root, test_cfg());
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
let before = read_log(&paths);
assert_eq!(
reducer::fold(&before).unwrap().mission.status,
MissionStatus::Planning
);
kranz_engine::mission_catalog::abandon_mission(
&root,
&mission_id,
"no longer needed",
LockForce::No,
)
.expect("abandon a planning mission");
let after = read_log(&paths);
assert_eq!(after.len(), before.len() + 1, "one event appended");
assert_eq!(after.first().unwrap().seq, 1);
assert_eq!(
after.last().unwrap().seq,
after.len() as u64,
"contiguous seq"
);
assert!(matches!(
&after.last().unwrap().kind,
EventKind::MissionAbandoned { reason } if reason == "no longer needed"
));
assert_eq!(
reducer::fold(&after).unwrap().mission.status,
MissionStatus::Abandoned
);
let snapshot = reducer::read_snapshot(&paths.state_file()).expect("state.json");
assert_eq!(snapshot.mission.status, MissionStatus::Abandoned);
assert_eq!(snapshot.last_seq, after.last().unwrap().seq);
}
#[tokio::test(flavor = "multi_thread")]
async fn running_an_abandoned_mission_is_rejected() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let engine = make_engine(&backend, &root, test_cfg());
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
kranz_engine::mission_catalog::abandon_mission(&root, &mission_id, "stop", LockForce::No)
.expect("abandon");
let backend2: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
let mut resumed =
MissionEngine::resume(backend2, &root, &mission_id, LockForce::No).expect("resume");
let err = resumed
.run()
.await
.expect_err("running a terminal mission must be rejected");
assert!(
matches!(err, kranz_engine::error::EngineError::InvalidState(_)),
"expected InvalidState, got {err:?}"
);
let events = read_log(&paths);
assert!(
!events
.iter()
.any(|e| matches!(e.kind, EventKind::WorkerSpawned { .. })),
"a rejected run must not spawn workers"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn abandon_already_terminal_mission_errors() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let engine = make_engine(&backend, &root, test_cfg());
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
drop(engine);
kranz_engine::mission_catalog::abandon_mission(&root, &mission_id, "first", LockForce::No)
.unwrap();
let after_first = read_log(&paths);
let err =
kranz_engine::mission_catalog::abandon_mission(&root, &mission_id, "again", LockForce::No)
.expect_err("abandoning a terminal mission must error");
assert!(
err.to_string().contains("already terminal"),
"error should name the terminal state: {err}"
);
let after_second = read_log(&paths);
assert_eq!(
after_second.len(),
after_first.len(),
"no event appended on the rejected abandon"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn abandon_fails_while_engine_holds_lock() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let engine = make_engine(&backend, &root, test_cfg());
let mission_id = engine.mission_id().to_string();
let err =
kranz_engine::mission_catalog::abandon_mission(&root, &mission_id, "x", LockForce::No)
.expect_err("abandon must fail while the lock is held");
assert!(
matches!(err, kranz_engine::error::EngineError::LockHeld(_)),
"expected LockHeld, got: {err}"
);
drop(engine);
}
#[tokio::test(flavor = "multi_thread")]
async fn preflight_flags_missing_program_and_ignores_present_ones() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let missing = vec![assertion(
"a-1",
"the check passes",
Some("definitely-not-a-real-program-xyz --check"),
)];
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, missing)).unwrap();
let issues = engine.preflight();
let warn = issues
.iter()
.find(|i| i.message.contains("definitely-not-a-real-program-xyz"))
.expect("the missing program is flagged");
assert_eq!(
warn.severity, "warn",
"a missing program is a warning, not an error"
);
assert!(
!issues.iter().any(|i| i.severity == "error"),
"no false hard errors: {issues:?}"
);
drop(engine);
let (_dir2, root2) = init_repo();
let backend2 = Arc::new(MockBackend::new());
let present = vec![
assertion("a-1", "trivially true", Some("exit 0")),
assertion("a-2", "a builtin", Some("cd .")),
];
let mut engine2 = make_engine(&backend2, &root2, test_cfg());
engine2.approve_plan(simple_plan(1, present)).unwrap();
assert!(
engine2.preflight().is_empty(),
"true / cd . resolve, so preflight is clean: {:?}",
engine2.preflight()
);
}
#[tokio::test(flavor = "multi_thread")]
async fn run_emits_preflight_decision_when_issues_exist() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let contract = vec![assertion(
"a-1",
"the check passes",
Some("definitely-not-a-real-program-xyz --check"),
)];
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
waive_reply("a-1", "command program unavailable in this environment"),
dirty_tree_commit_as_is(),
judgement("complete", ""),
waive_reply("a-1", "still unavailable"),
]),
worker_pass(),
]));
let cfg = MissionConfig {
max_fix_cycles_per_milestone: 1,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Blocked);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let preflight_seqs: Vec<u64> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, .. }
if summary.starts_with("preflight:") =>
{
Some(e.seq)
}
_ => None,
})
.collect();
assert_eq!(
preflight_seqs.len(),
1,
"exactly one preflight decision: {preflight_seqs:?}"
);
let preflight_seq = preflight_seqs[0];
let first_spawn = seq_of(&events, "worker.spawned");
assert!(
preflight_seq < first_spawn,
"preflight decision {preflight_seq} precedes the first worker spawn {first_spawn}"
);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.starts_with("preflight:")
&& summary.contains("definitely-not-a-real-program-xyz")
)));
}
#[tokio::test(flavor = "multi_thread")]
async fn sandbox_preflight_inert_when_enforce_off() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let contract = vec![assertion(
"a-1",
"writes outside the allowlist",
Some("sh -c 'echo x > $HOME/kranz_pf_should_not_run'"),
)];
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, contract)).unwrap();
let issues = engine.preflight();
assert!(
!issues
.iter()
.any(|i| i.message.contains("fs sandbox profile")),
"enforce:off must add zero sandbox preflight issues: {issues:?}"
);
}
#[cfg(target_os = "linux")]
#[tokio::test(flavor = "multi_thread")]
async fn sandbox_preflight_emits_no_macos_profile_issue_on_linux() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut cfg = test_cfg();
cfg.worker.sandbox.enforce = kranz_engine::types::SandboxEnforce::Fs;
let contract = vec![assertion(
"a-1",
"writes outside the allowlist",
Some("sh -c 'echo x > $HOME/kranz_pf_should_not_run'"),
)];
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let issues = engine.preflight();
assert!(
!issues
.iter()
.any(|i| i.message.contains("fs sandbox profile")),
"Linux must add no macOS-profile preflight issue: {issues:?}"
);
}
#[cfg(windows)]
#[test]
fn sandbox_preflight_windows_process_backend_is_appcontainer() {
assert_eq!(
kranz_engine::sandbox::platform_support(kranz_engine::types::SandboxEnforce::Fs, "windows"),
kranz_engine::sandbox::SandboxDecision::Enforce(
kranz_engine::sandbox::SandboxBackend::AppContainer
)
);
}
#[cfg(target_os = "macos")]
#[tokio::test(flavor = "multi_thread")]
async fn sandbox_preflight_flags_command_that_writes_outside_allowlist() {
if !setup() {
return;
}
if std::process::Command::new("which")
.arg("sandbox-exec")
.output()
.map(|o| !o.status.success())
.unwrap_or(true)
{
kranz_engine::test_capability::skip(
kranz_engine::test_capability::capability::SANDBOX_EXEC,
"sandbox-exec not found on this host",
);
return;
}
match std::process::Command::new("sandbox-exec")
.arg("-p")
.arg("(version 1)\n(allow default)\n")
.arg("/usr/bin/true")
.output()
{
Ok(output) if output.status.success() => {}
Ok(output) => {
eprintln!(
"SKIP-UNDER-WRAP (gate-sandbox-supervision-dogfood): \
sandbox-exec cannot apply a smoke profile here (nested apply is denied \
inside the gate sandbox wrap); skipping: {}",
String::from_utf8_lossy(&output.stderr)
);
return;
}
Err(e) => {
eprintln!("sandbox-exec smoke probe failed; skipping: {e}");
return;
}
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut cfg = test_cfg();
cfg.worker.sandbox.enforce = kranz_engine::types::SandboxEnforce::Fs;
let marker = format!("kranz_pf_{}", uuid::Uuid::new_v4());
let real_home = std::env::var("HOME").expect("HOME must be set for this test");
let contract = vec![
assertion(
"a-outside",
"writes outside the allowlist",
Some(&format!("echo x > '{real_home}/{marker}'")),
),
assertion("a-benign", "trivially true", Some("exit 0")),
];
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, contract)).unwrap();
let issues = engine.preflight();
assert!(
!issues.iter().any(|i| i.severity == "error"),
"sandbox preflight must never escalate to error: {issues:?}"
);
let outside_issue = issues
.iter()
.find(|i| i.message.contains("[a-outside]") && i.message.contains("fs sandbox profile"));
assert!(
outside_issue.is_some(),
"expected a sandbox warn for the out-of-allowlist command: {issues:?}"
);
assert_eq!(outside_issue.unwrap().severity, "warn");
assert!(
!issues
.iter()
.any(|i| i.message.contains("[a-benign]") && i.message.contains("fs sandbox profile")),
"benign command must not produce a sandbox issue: {issues:?}"
);
if let Ok(home) = std::env::var("HOME") {
let _ = std::fs::remove_file(std::path::Path::new(&home).join(&marker));
}
}
#[tokio::test(flavor = "multi_thread")]
async fn request_and_approve_revised_plan_drops_and_adds_features() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let revised_json = json!({
"goal": GOAL,
"validationContract": [],
"milestones": [{
"title": "M1",
"features": [
{ "title": "feature 1", "spec": "build part 1", "validationCriteria": ["part 1 works"] },
{ "title": "extra feature", "spec": "build the newly-needed part", "validationCriteria": ["extra works"] }
]
}]
})
.to_string();
let backend = Arc::new(MockBackend::with_scripts(vec![orch_script(vec![
revised_json,
])]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(3, vec![])).unwrap();
assert_eq!(engine.state().mission.status, MissionStatus::Approved);
let request = timeout(TEST_TIMEOUT, engine.request_revised_plan())
.await
.expect("request_revised_plan must not hang")
.expect("scripted plan JSON is not a backend error");
let plan = match request {
PlanRequest::Ready(plan) => plan,
PlanRequest::NotReady(text) => panic!("scripted revised plan must parse: {text}"),
PlanRequest::WrongPlan { reason } => {
panic!("scripted revised plan must parse, got a wrong-plan escalation: {reason}")
}
};
assert_eq!(plan.milestones[0].features.len(), 2);
engine
.approve_revised_plan(plan)
.expect("apply the revised plan");
let ms = &engine.state().mission.milestones[0];
let by_id = |id: &str| ms.features.iter().find(|f| f.id == id).cloned();
assert_eq!(
by_id("f-1-1").unwrap().status,
FeatureStatus::Pending,
"kept feature untouched"
);
assert_eq!(
by_id("f-1-2").unwrap().status,
FeatureStatus::Skipped,
"dropped feature skipped"
);
assert_eq!(
by_id("f-1-3").unwrap().status,
FeatureStatus::Skipped,
"dropped feature skipped"
);
let added = by_id("ms-1-replan-1-1").expect("added feature exists with re-plan id");
assert_eq!(
added.origin,
FeatureOrigin::Fix,
"added feature is fix-origin"
);
assert_eq!(added.status, FeatureStatus::Pending);
assert_eq!(added.title, "extra feature");
let mission_id = engine.mission_id().to_string();
let md = std::fs::read_to_string(
root.join(".kranz")
.join("missions")
.join(&mission_id)
.join("revised-plan.md"),
)
.expect("revised-plan.md written");
assert!(
md.starts_with(&format!("# Revised mission plan — {mission_id}")),
"{md}"
);
assert!(md.contains("## Re-plan changes applied"), "{md}");
assert!(md.contains("extra feature"), "{md}");
let subject = raw_git(&root, &["log", "-1", "--format=%s"]);
assert_eq!(
subject.trim(),
format!("[kranz] revised plan for {mission_id}")
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. } if summary.starts_with("re-plan for ms-1")
)));
assert_eq!(
events
.iter()
.filter(
|e| matches!(&e.kind, EventKind::FeatureSkipped { reason, .. }
if reason.contains("re-plan"))
)
.count(),
2,
"both dropped features skipped by the re-plan"
);
assert_eq!(
event_types(&events)
.iter()
.filter(|t| **t == "fixfeature.created")
.count(),
1,
"exactly one feature added by the re-plan"
);
let state = reducer::fold(&events).unwrap();
assert_eq!(
state.mission.milestones[0].features.len(),
4,
"3 planned + 1 added"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn approve_revised_plan_rejects_dropping_a_completed_milestone() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let finding = json!([{
"subject": "part 2 works",
"severity": "major",
"evidence": "still failing",
"suggestedFix": ""
}]);
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""), dirty_tree_commit_as_is(),
judgement("complete", ""), fix_features(1), dirty_tree_commit_as_is(),
judgement("complete", ""), fix_features(1), ]),
validator_with(json!([])), worker_pass(), validator_with(finding.clone()), worker_pass(), validator_with(finding), ]));
let cfg = MissionConfig {
skip_functional: false,
max_fix_cycles_per_milestone: 1, ..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
let plan = Plan {
goal: GOAL.to_string(),
validation_contract: vec![],
milestones: vec![
PlanMilestone {
title: "M1".to_string(),
features: vec![PlanFeature {
title: "feature 1".to_string(),
spec: "build part 1".to_string(),
validation_criteria: vec!["part 1 works".to_string()],
}],
},
PlanMilestone {
title: "M2".to_string(),
features: vec![PlanFeature {
title: "feature 2".to_string(),
spec: "build part 2".to_string(),
validation_criteria: vec!["part 2 works".to_string()],
}],
},
],
considered_alternatives: None,
command_grants: vec![],
touch_set: vec![],
standards_manifest: None,
reviewer_independence: None,
};
engine.approve_plan(plan).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(
status,
MissionStatus::Blocked,
"M2 blocked at the fix-cycle cap"
);
assert_eq!(
engine.state().mission.milestones[0].status,
MilestoneStatus::Complete,
"M1 is complete"
);
let drops_completed = Plan {
goal: GOAL.to_string(),
validation_contract: vec![],
milestones: vec![PlanMilestone {
title: "M2".to_string(),
features: vec![PlanFeature {
title: "feature 2".to_string(),
spec: "build part 2".to_string(),
validation_criteria: vec!["part 2 works".to_string()],
}],
}],
considered_alternatives: None,
command_grants: vec![],
touch_set: vec![],
standards_manifest: None,
reviewer_independence: None,
};
let err = engine
.approve_revised_plan(drops_completed)
.expect_err("dropping a completed milestone must be rejected");
assert!(
matches!(err, kranz_engine::error::EngineError::InvalidState(_)),
"expected InvalidState, got: {err}"
);
assert!(
err.to_string().contains("M1"),
"error names the dropped completed milestone: {err}"
);
let alters_completed = Plan {
goal: GOAL.to_string(),
validation_contract: vec![],
milestones: vec![
PlanMilestone {
title: "M1".to_string(),
features: vec![PlanFeature {
title: "feature 1 RENAMED".to_string(),
spec: "build part 1".to_string(),
validation_criteria: vec!["part 1 works".to_string()],
}],
},
PlanMilestone {
title: "M2".to_string(),
features: vec![PlanFeature {
title: "feature 2".to_string(),
spec: "build part 2".to_string(),
validation_criteria: vec!["part 2 works".to_string()],
}],
},
],
considered_alternatives: None,
command_grants: vec![],
touch_set: vec![],
standards_manifest: None,
reviewer_independence: None,
};
let err = engine
.approve_revised_plan(alters_completed)
.expect_err("altering a completed milestone's features must be rejected");
assert!(
err.to_string().contains("alters"),
"error explains the alteration: {err}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn parallel_batch_runs_both_features_and_leaks_no_worktrees() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
orch_script(vec![
parallel_plan(&["f-1-1", "f-1-2"]),
judgement("complete", ""),
judgement("complete", ""),
no_lesson(),
]),
worker_pass(),
worker_pass(),
]));
let cfg = MissionConfig {
max_parallel_workers: 2,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(2, vec![])).unwrap();
let mission_id = engine.mission_id().to_string();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let ms = &engine.state().mission.milestones[0];
assert_eq!(
ms.features[0].status,
FeatureStatus::Complete,
"f-1-1 complete"
);
assert_eq!(
ms.features[1].status,
FeatureStatus::Complete,
"f-1-2 complete"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let worker_spawns = events
.iter()
.filter(|e| {
matches!(
&e.kind,
EventKind::WorkerSpawned {
role: Role::Worker,
..
}
)
})
.count();
assert_eq!(worker_spawns, 2, "both features ran a worker");
let types = event_types(&events);
assert_eq!(
types.iter().filter(|t| **t == "feature.completed").count(),
2,
"both features completed: {types:?}"
);
assert!(
types.contains(&"milestone.completed"),
"milestone completed: {types:?}"
);
assert!(
types.contains(&"mission.completed"),
"mission completed: {types:?}"
);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.starts_with("parallel:") && summary.contains("2 workers")
)),
"a parallel summary decision naming 2 workers exists"
);
let plan_seq = events
.iter()
.find_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, .. }
if summary.starts_with("parallel plan for") =>
{
Some(e.seq)
}
_ => None,
})
.expect("parallel plan decision exists");
let first_worker_spawn = events
.iter()
.find(|e| {
matches!(
&e.kind,
EventKind::WorkerSpawned {
role: Role::Worker,
..
}
)
})
.expect("a worker spawned")
.seq;
assert!(
plan_seq < first_worker_spawn,
"parallel plan precedes the first worker"
);
let repo = GitRepo::open(&root).unwrap();
let worktrees = repo.list_worktrees().unwrap();
assert_eq!(
worktrees.len(),
1,
"only the primary worktree remains: {worktrees:?}"
);
assert!(
!repo
.branch_exists(&format!("kranz/wt/{mission_id}/f-1-1"))
.unwrap_or(false),
"per-feature worktree branch must be deleted"
);
assert!(
!repo
.branch_exists(&format!("kranz/wt/{mission_id}/f-1-2"))
.unwrap_or(false),
"per-feature worktree branch must be deleted"
);
let leak_prefix = format!("kranz-wt-{mission_id}-");
for entry in std::fs::read_dir(std::env::temp_dir()).unwrap().flatten() {
let name = entry.file_name();
let name = name.to_string_lossy();
assert!(
!name.starts_with(&leak_prefix),
"a parallel worktree dir leaked into temp: {name}"
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn parallel_batch_sessions_overlap_in_wall_clock() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
orch_script(vec![
parallel_plan(&["f-1-1", "f-1-2"]),
judgement("complete", ""),
judgement("complete", ""),
no_lesson(),
]),
worker_pass().rendezvous(3),
worker_pass().rendezvous(3),
]));
let cfg = MissionConfig {
max_parallel_workers: 2,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(2, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.starts_with("parallel:") && summary.contains("peak 2 concurrent")
)),
"the batch summary must record peak 2 concurrent sessions: {:?}",
events
.iter()
.filter_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, .. } => Some(summary.clone()),
_ => None,
})
.collect::<Vec<_>>()
);
assert_eq!(events.first().unwrap().seq, 1);
assert_eq!(
events.last().unwrap().seq,
events.len() as u64,
"contiguous seq"
);
let worker_spawns = events
.iter()
.filter(|e| {
matches!(
&e.kind,
EventKind::WorkerSpawned {
role: Role::Worker,
..
}
)
})
.count();
let worker_completes = events
.iter()
.filter(|e| matches!(&e.kind, EventKind::WorkerCompleted { .. }))
.count();
assert_eq!(worker_spawns, 2, "both worker sessions spawned");
assert!(
worker_completes >= 2,
"both worker sessions completed: {worker_completes}"
);
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
}
#[tokio::test(flavor = "multi_thread")]
async fn crash_mid_parallel_batch_resumes_cleanly() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend1 = Arc::new(MockBackend::with_scripts(vec![
orch_script(vec![parallel_plan(&["f-1-1", "f-1-2"])]),
worker_pass(), ]));
let cfg = MissionConfig {
max_parallel_workers: 2,
..test_cfg()
};
let mut engine = make_engine(&backend1, &root, cfg);
engine.approve_plan(simple_plan(2, vec![])).unwrap();
let mission_id = engine.mission_id().to_string();
let paths = engine.paths().clone();
let err = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect_err("a starved concurrent worker aborts the batch (simulated crash)");
eprintln!("phase 1 crashed as scripted: {err}");
drop(engine);
let phase1 = read_log(&paths);
assert_eq!(phase1.first().unwrap().seq, 1);
assert_eq!(
phase1.last().unwrap().seq,
phase1.len() as u64,
"contiguous seq after crash"
);
let state1 = reducer::fold(&phase1).expect("crashed log still folds cleanly");
for f in &state1.mission.milestones[0].features {
assert_eq!(
f.status,
FeatureStatus::Active,
"features left Active by the crash"
);
assert!(
f.worker_runs.is_empty(),
"no worker run recorded before the crash"
);
}
assert!(
!phase1.iter().any(|e| matches!(
&e.kind,
EventKind::WorkerSpawned {
role: Role::Worker,
..
}
)),
"no buffered worker session survived the crash"
);
let backend2 = Arc::new(MockBackend::with_scripts(vec![
worker_pass(), orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
worker_pass(), ]));
let backend2_dyn: Arc<dyn AgentBackend> = Arc::clone(&backend2) as Arc<dyn AgentBackend>;
let mut engine = MissionEngine::resume(backend2_dyn, &root, &mission_id, LockForce::No)
.expect("resume mission");
engine.seed_worker_auth_verdict_for_test(AuthVerdict::Inconclusive);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
drop(engine);
let events = read_log(&paths);
assert_eq!(events.first().unwrap().seq, 1);
assert_eq!(
events.last().unwrap().seq,
events.len() as u64,
"one contiguous log"
);
let state = reducer::fold(&events).unwrap();
assert_eq!(state.mission.status, MissionStatus::Complete);
assert!(state.mission.milestones[0]
.features
.iter()
.all(|f| f.status == FeatureStatus::Complete));
let repo = GitRepo::open(&root).unwrap();
assert_eq!(
repo.list_worktrees().unwrap().len(),
1,
"no leaked worktrees: recovered clean"
);
assert!(!repo
.branch_exists(&format!("kranz/wt/{mission_id}/f-1-1"))
.unwrap_or(false));
assert!(!repo
.branch_exists(&format!("kranz/wt/{mission_id}/f-1-2"))
.unwrap_or(false));
}
#[tokio::test(flavor = "multi_thread")]
async fn resume_does_not_sweep_worktrees_while_lock_is_live() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::new());
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let mission_id = engine.mission_id().to_string();
let repo = GitRepo::open(&root).unwrap();
let branch = format!("kranz/wt/{mission_id}/f-1-1");
let wt_path = std::env::temp_dir().join(format!("kranz-wt-{mission_id}-f-1-1"));
let head = repo.head_sha().unwrap();
repo.add_worktree(&wt_path, &branch, &head).unwrap();
let backend2: Arc<dyn AgentBackend> = Arc::new(MockBackend::new());
let err = match MissionEngine::resume(backend2, &root, &mission_id, LockForce::No) {
Ok(_) => panic!("resume must fail while a live engine holds the lock"),
Err(e) => e,
};
assert!(
matches!(err, kranz_engine::error::EngineError::LockHeld(_)),
"expected LockHeld, got: {err}"
);
assert!(
wt_path.exists(),
"live worktree must not be swept by a lock-refused resume"
);
assert!(
repo.branch_exists(&branch).unwrap_or(false),
"live branch must not be -D'd by a lock-refused resume"
);
let _ = repo.remove_worktree(&wt_path);
let _ = repo.prune_worktrees();
let _ = repo.delete_branch_force(&branch);
drop(engine);
}
#[tokio::test(flavor = "multi_thread")]
async fn max_parallel_one_is_the_unchanged_sequential_path() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
worker_pass(),
]));
let cfg = MissionConfig {
max_parallel_workers: 1,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(2, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
!events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.starts_with("parallel plan for") || summary.starts_with("parallel:")
)),
"no parallel decisions with max_parallel_workers=1"
);
let types = event_types(&events);
assert_eq!(
types.iter().filter(|t| **t == "feature.completed").count(),
2,
"both features completed sequentially: {types:?}"
);
assert!(types.contains(&"mission.completed"));
let repo = GitRepo::open(&root).unwrap();
assert_eq!(
repo.list_worktrees().unwrap().len(),
1,
"no extra worktrees in sequential mode"
);
}
fn plan_feature(id: &str, title: &str, spec: &str) -> Feature {
Feature {
id: id.to_string(),
title: title.to_string(),
spec: spec.to_string(),
validation_criteria: vec![format!("{title} works")],
origin: FeatureOrigin::Plan,
status: FeatureStatus::Pending,
worker_runs: Vec::new(),
commits: Vec::new(),
respawns: 0,
}
}
#[test]
fn conflict_synthesizes_resolution_fix_feature() {
let f_1_1 = plan_feature("f-1-1", "feature one", "build part one");
let f_1_2 = plan_feature("f-1-2", "feature two", "build the widget in src/widget.rs");
let existing = vec![f_1_1, f_1_2.clone()];
let conflict_files = vec!["src/widget.rs".to_string(), "src/lib.rs".to_string()];
let resolution = synthesize_conflict_resolution("ms-1", &f_1_2, &conflict_files, &existing)
.expect("a Plan-origin conflict must synthesize a resolution feature");
assert_eq!(
resolution.id, "ms-1-conflict-1",
"namespaced conflict-resolution id"
);
assert_eq!(
resolution.origin,
FeatureOrigin::Fix,
"resolution is a fix feature"
);
assert_eq!(
resolution.status,
FeatureStatus::Pending,
"resolution starts Pending"
);
assert!(
resolution.worker_runs.is_empty(),
"fresh feature, no runs yet"
);
assert_eq!(resolution.respawns, 0);
let spec = &resolution.spec;
assert!(
spec.contains("build the widget in src/widget.rs"),
"original spec text: {spec}"
);
assert!(
spec.contains("src/widget.rs"),
"conflicting file listed: {spec}"
);
assert!(
spec.contains("src/lib.rs"),
"second conflicting file listed: {spec}"
);
assert!(
spec.contains("already contains") || spec.contains("merged first"),
"notes earlier features already merged: {spec}"
);
assert!(
resolution.title.contains("feature two"),
"title references original: {}",
resolution.title
);
assert_eq!(resolution.validation_criteria, f_1_2.validation_criteria);
}
#[test]
fn second_conflict_gets_a_fresh_namespaced_id() {
let f_1_2 = plan_feature("f-1-2", "feature two", "spec two");
let f_1_3 = plan_feature("f-1-3", "feature three", "spec three");
let mut resolution_1 = plan_feature(
"ms-1-conflict-1",
"Resolve merge conflict: feature two",
"…",
);
resolution_1.origin = FeatureOrigin::Fix;
let existing = vec![f_1_2, f_1_3.clone(), resolution_1];
let resolution =
synthesize_conflict_resolution("ms-1", &f_1_3, &["src/x.rs".to_string()], &existing)
.expect("second conflict synthesizes a resolution");
assert_eq!(
resolution.id, "ms-1-conflict-2",
"second conflict is -conflict-2, no collision"
);
}
#[test]
fn resolution_feature_does_not_spawn_another_resolution() {
let mut resolution = plan_feature(
"ms-1-conflict-1",
"Resolve merge conflict: feature two",
"spec",
);
resolution.origin = FeatureOrigin::Fix;
let existing = vec![resolution.clone()];
let again =
synthesize_conflict_resolution("ms-1", &resolution, &["src/x.rs".to_string()], &existing);
assert!(
again.is_none(),
"a -conflict- feature must not spawn another resolution"
);
}
#[test]
fn conflict_with_no_named_files_still_synthesizes() {
let f = plan_feature("f-1-1", "feature one", "the original work");
let resolution = synthesize_conflict_resolution("ms-2", &f, &[], std::slice::from_ref(&f))
.expect("empty file list still synthesizes");
assert_eq!(resolution.id, "ms-2-conflict-1");
assert!(
resolution.spec.contains("the original work"),
"original spec preserved"
);
assert!(
resolution.spec.contains("no specific files"),
"empty conflict list gets a placeholder: {}",
resolution.spec
);
}
#[tokio::test]
async fn run_reasserts_mission_branch_and_restores_base() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
"NONE".to_string(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine
.approve_plan(simple_plan(
1,
vec![assertion("a-1", "the build command succeeds", Some("cd ."))],
))
.unwrap();
raw_git(&root, &["checkout", "main"]);
let main_before = raw_git(&root, &["rev-parse", "main"]);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let branch = engine.state().mission.mission_branch.clone();
drop(engine);
assert_eq!(
main_before,
raw_git(&root, &["rev-parse", "main"]),
"main must not receive mission commits"
);
let ahead: u32 = raw_git(&root, &["rev-list", "--count", &format!("main..{branch}")])
.trim()
.parse()
.unwrap();
assert!(ahead > 0, "the mission branch must carry the work");
assert_eq!(
raw_git(&root, &["rev-parse", "--abbrev-ref", "HEAD"]).trim(),
branch,
"completed run keeps its artifacts visible on the mission branch"
);
}
#[tokio::test]
async fn create_refuses_another_missions_branch_as_base() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
raw_git(&root, &["checkout", "-b", "kranz/mission-m-fake01"]);
let backend: Arc<dyn AgentBackend> = Arc::new(MockBackend::with_scripts(vec![]));
let Err(err) = MissionEngine::create(backend, &root, GOAL, test_cfg()) else {
panic!("create must refuse a kranz/mission-* base branch");
};
let msg = err.to_string();
assert!(
msg.contains("another mission's branch"),
"refusal must explain the stacking hazard, got: {msg}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn empty_deliverable_mission_terminates_failed() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass_no_write(),
orch_script(vec![judgement("complete", "")]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Failed);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
types.contains(&"mission.failed"),
"expected mission.failed: {types:?}"
);
assert!(
!types.contains(&"mission.completed"),
"must never complete on an empty deliverable diff: {types:?}"
);
let reason = events
.iter()
.find_map(|e| match &e.kind {
EventKind::MissionFailed { reason } => Some(reason.clone()),
_ => None,
})
.expect("mission.failed event carries a reason");
assert!(!reason.trim().is_empty(), "reason must be non-empty");
}
#[tokio::test(flavor = "multi_thread")]
async fn delivering_mission_still_completes() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let types = event_types(&events);
assert!(
types.contains(&"mission.completed"),
"expected mission.completed: {types:?}"
);
assert!(
!types.contains(&"mission.failed"),
"delivering mission must not trip the empty-deliverable safety net: {types:?}"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn spoofed_meta_subject_commit_counts_as_deliverable() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_pass_no_write(),
orch_script(vec![judgement("complete", ""), no_lesson()]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
std::fs::write(root.join("smuggled.txt"), "real deliverable\n").unwrap();
raw_git(&root, &["add", "smuggled.txt"]);
raw_git(&root, &["commit", "-m", "[kranz] mission report cleanup"]);
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(
status,
MissionStatus::Complete,
"a spoofed-subject commit carrying a real file is a deliverable; \
the empty-deliverable net must count it, not hide it"
);
let paths = engine.paths().clone();
drop(engine);
let types = event_types(&read_log(&paths));
assert!(
!types.contains(&"mission.failed"),
"must not fail as empty-deliverable: {types:?}"
);
}
const LEAKED_SECRET: &str = "sk-ant-api03-aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
fn worker_leaks_secret() -> MockScript {
MockScript::single_shot_json(&json!({
"result": "pass",
"summary": "implemented and tested",
"filesTouched": ["leak.txt"],
"testsAdded": [],
"testEvidence": "all green",
"commits": []
}))
.writes_file("leak.txt", format!("ANTHROPIC_API_KEY={LEAKED_SECRET}\n"))
}
#[tokio::test(flavor = "multi_thread")]
async fn secret_scan_refusal_fails_feature_instead_of_erroring_run() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_leaks_secret(),
orch_script(vec![dirty_tree_commit_as_is()]),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect("a scan refusal must not error the run loop");
assert_eq!(status, MissionStatus::Blocked);
assert_eq!(engine.state().mission.status, MissionStatus::Blocked);
assert_eq!(
engine.state().mission.milestones[0].status,
MilestoneStatus::Blocked,
"the milestone blocks against the poisoned tree"
);
assert_eq!(
engine.state().mission.milestones[0].features[0].status,
FeatureStatus::Failed,
"the leaking feature is failed, not left Active"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let reason = events
.iter()
.find_map(|e| match &e.kind {
EventKind::FeatureFailed { reason, .. } => Some(reason.clone()),
_ => None,
})
.expect("feature.failed with the refusal reason is on the log");
assert!(
reason.contains("secret scan"),
"reason names the scan: {reason}"
);
assert!(
reason.contains("anthropic-api-key"),
"reason names the rule: {reason}"
);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.contains("refused by secret scan")
)),
"an orchestrator.decision records the refusal"
);
let blocked_reason = events
.iter()
.find_map(|e| match &e.kind {
EventKind::MilestoneBlocked { reason, .. } => Some(reason.clone()),
_ => None,
})
.expect("milestone.blocked with the cleanup cue is on the log");
assert!(
blocked_reason.contains("secret scan") && blocked_reason.contains("leak.txt"),
"blocked reason names the scan and the dirty path: {blocked_reason}"
);
let raw_log = std::fs::read_to_string(paths.events_file()).unwrap();
assert!(
!raw_log.contains(LEAKED_SECRET),
"the raw secret must never land in events.jsonl"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn secret_scan_refusal_blocks_milestone_so_later_features_never_run() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
worker_leaks_secret(),
orch_script(vec![dirty_tree_commit_as_is(), dirty_tree_commit_as_is()]),
worker_pass_no_write(),
]));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(2, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect("a scan refusal must not error the run loop");
assert_eq!(status, MissionStatus::Blocked);
let ms = &engine.state().mission.milestones[0];
assert_eq!(ms.status, MilestoneStatus::Blocked);
assert_eq!(ms.features[0].status, FeatureStatus::Failed);
assert_eq!(
ms.features[1].status,
FeatureStatus::Pending,
"feature B must stay Pending, not be failed against A's poisoned tree"
);
assert_eq!(
backend.started_specs().len(),
2,
"no session may spawn for feature B after the milestone blocked"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::WorkerSpawned { feature_id: Some(id), .. } if id == "f-1-1"
)),
"feature A's worker spawned"
);
assert!(
!events.iter().any(|e| matches!(
&e.kind,
EventKind::WorkerSpawned { feature_id: Some(id), .. } if id == "f-1-2"
)),
"feature B's worker must never spawn"
);
let failed: Vec<&String> = events
.iter()
.filter_map(|e| match &e.kind {
EventKind::FeatureFailed { feature_id, .. } => Some(feature_id),
_ => None,
})
.collect();
assert_eq!(
failed,
vec!["f-1-1"],
"exactly one feature.failed, for the feature that actually leaked"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn parallel_secret_scan_refusal_is_recorded_not_swallowed() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(vec![
orch_script(vec![parallel_plan(&["f-1-1", "f-1-2"])]),
worker_leaks_secret(),
worker_leaks_secret(),
]));
let cfg = MissionConfig {
max_parallel_workers: 2,
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(2, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.expect("a parallel scan refusal must not error the run loop");
assert_eq!(status, MissionStatus::Failed);
let ms = &engine.state().mission.milestones[0];
assert_eq!(ms.features[0].status, FeatureStatus::Failed);
assert_eq!(ms.features[1].status, FeatureStatus::Failed);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
for feature_id in ["f-1-1", "f-1-2"] {
assert!(
events.iter().any(|e| matches!(
&e.kind,
EventKind::OrchestratorDecision { summary, .. }
if summary.contains("refused by secret scan")
&& summary.contains(feature_id)
)),
"a refusal decision names {feature_id}"
);
}
let raw_log = std::fs::read_to_string(paths.events_file()).unwrap();
assert!(
!raw_log.contains(LEAKED_SECRET),
"the raw secret must never land in events.jsonl"
);
let repo = GitRepo::open(&root).unwrap();
let worktrees = repo.list_worktrees().unwrap();
assert_eq!(
worktrees.len(),
1,
"only the primary worktree remains: {worktrees:?}"
);
}
fn pack_contract_fixture_dir() -> PathBuf {
Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/pack-contract-synthetic")
}
fn pack_contract_scripts() -> Vec<MockScript> {
vec![
worker_pass(),
orch_script(vec![
dirty_tree_commit_as_is(),
judgement("complete", ""),
no_lesson(),
]),
]
}
fn worker_system_prompt(backend: &MockBackend) -> String {
let specs: Vec<_> = backend
.started_specs()
.into_iter()
.filter(|s| {
matches!(&s.prompt, PromptMode::SingleShot(task) if task.contains("Implement feature `f-1-1`"))
})
.collect();
assert_eq!(specs.len(), 1, "exactly one worker session for f-1-1");
specs[0]
.append_system_prompt
.clone()
.expect("worker sessions carry an append_system_prompt")
}
fn worker_prompt_hash(events: &[Event]) -> String {
events
.iter()
.find_map(|e| match &e.kind {
EventKind::WorkerSpawned {
role: Role::Worker,
prompt_hash,
..
} => Some(prompt_hash.clone()),
_ => None,
})
.expect("worker.spawned for the worker role")
}
fn pack_contract_decisions(events: &[Event]) -> Vec<(String, Option<String>)> {
events
.iter()
.filter_map(|e| match &e.kind {
EventKind::OrchestratorDecision { summary, detail } => {
Some((summary.clone(), detail.clone()))
}
_ => None,
})
.collect()
}
#[tokio::test(flavor = "multi_thread")]
async fn pack_contract_gate_runs_and_prompt_reaches_worker() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(pack_contract_scripts()));
let cfg = MissionConfig {
pack_dir: Some(pack_contract_fixture_dir().display().to_string()),
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let decisions = pack_contract_decisions(&events);
assert!(
decisions.iter().any(|(summary, detail)| summary
.starts_with("pack contract: pack `zz-synthetic-fixture-pack` (schema 3) registered:")
&& detail
.as_deref()
.is_some_and(|d| d.contains("zz-pack-gate-synthetic")
&& d.contains("zz-pack-prompt-synthetic"))),
"run-start pack decision missing: {decisions:?}"
);
let (summary, detail) = decisions
.iter()
.find(|(summary, _)| {
summary.starts_with("pack `zz-synthetic-fixture-pack` gates (final gate):")
})
.expect("final-gate pack decision missing");
assert!(
summary.contains("1 deterministic gate(s) passed"),
"{summary}"
);
let detail = detail.as_deref().expect("verdict detail");
assert!(detail.contains("zz-pack-gate-synthetic: PASS"), "{detail}");
let worker_prompt = worker_system_prompt(&backend);
assert!(
worker_prompt.contains("ZZ-SYNTHETIC-PACK-MARKER"),
"pack guidance missing from the worker prompt"
);
assert!(
worker_prompt
.contains("pack `zz-synthetic-fixture-pack`, prompt `zz-pack-prompt-synthetic`"),
"the injection is marked with its pack/prompt provenance"
);
for spec in backend.started_specs() {
let is_worker = matches!(&spec.prompt, PromptMode::SingleShot(task) if task.contains("Implement feature `f-1-1`"));
if !is_worker {
let prompt = spec.append_system_prompt.as_deref().unwrap_or("");
assert!(
!prompt.contains("ZZ-SYNTHETIC-PACK-MARKER"),
"pack guidance leaked into a non-target role's prompt"
);
}
}
assert_ne!(
worker_prompt_hash(&events),
kranz_engine::prompts::hash(Role::Worker),
"a pack-extended prompt must not record the bare template hash"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn pack_contract_failing_gate_is_advisory_named_and_never_blocks() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let pack_dir = root.join("zz-failing-pack");
std::fs::create_dir_all(&pack_dir).unwrap();
std::fs::write(
pack_dir.join("pack.toml"),
"[pack]\nname = \"zz-failing-pack\"\nschema = 3\n\n\
[[gate]]\nname = \"zz-pack-gate-failing\"\ncommand = \"exit 1\"\n",
)
.unwrap();
let backend = Arc::new(MockBackend::with_scripts(pack_contract_scripts()));
let cfg = MissionConfig {
pack_dir: Some("zz-failing-pack".to_string()),
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(
status,
MissionStatus::Complete,
"advisory: a failing pack gate never blocks completion"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let (summary, detail) = pack_contract_decisions(&events)
.into_iter()
.find(|(summary, _)| summary.starts_with("pack `zz-failing-pack` gates (final gate):"))
.expect("final-gate pack decision missing");
assert!(
summary.contains("named gate(s) failed: zz-pack-gate-failing"),
"{summary}"
);
let detail = detail.expect("verdict detail");
assert!(detail.contains("zz-pack-gate-failing: FAIL"), "{detail}");
}
#[tokio::test(flavor = "multi_thread")]
async fn pack_contract_no_pack_behavior_is_byte_identical() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let backend = Arc::new(MockBackend::with_scripts(pack_contract_scripts()));
let mut engine = make_engine(&backend, &root, test_cfg());
engine.approve_plan(simple_plan(1, vec![])).unwrap();
let status = timeout(TEST_TIMEOUT, engine.run())
.await
.expect("run must not hang")
.unwrap();
assert_eq!(status, MissionStatus::Complete);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
let decisions = pack_contract_decisions(&events);
assert!(
!decisions
.iter()
.any(|(summary, _)| summary.starts_with("pack contract:") || summary.contains("pack `")),
"no pack decisions without a pack: {decisions:?}"
);
let worker_prompt = worker_system_prompt(&backend);
assert!(
!worker_prompt.contains("Pack guidance"),
"no pack section without a pack"
);
assert_eq!(
worker_prompt_hash(&events),
kranz_engine::prompts::hash(Role::Worker),
"the bare template hash is recorded without a pack"
);
}
#[tokio::test(flavor = "multi_thread")]
async fn pack_contract_invalid_pack_fails_closed_before_approval() {
if !setup() {
return;
}
let (_dir, root) = init_repo();
let pack_dir = root.join("zz-broken-pack");
std::fs::create_dir_all(&pack_dir).unwrap();
std::fs::write(
pack_dir.join("pack.toml"),
"[pack]\nname = \"zz-broken-pack\"\nschema = 3\n\n\
[[gate]]\nname = \"zz-dup\"\ncommand = \"cd .\"\n\n\
[[gate]]\nname = \"zz-dup\"\ncommand = \"cd .\"\n",
)
.unwrap();
let backend = Arc::new(MockBackend::with_scripts(pack_contract_scripts()));
let cfg = MissionConfig {
pack_dir: Some("zz-broken-pack".to_string()),
..test_cfg()
};
let mut engine = make_engine(&backend, &root, cfg);
let err = engine
.approve_plan(simple_plan(1, vec![]))
.expect_err("an invalid pack must fail approval closed");
let msg = err.to_string();
assert!(
msg.contains("duplicate [[gate]] name `zz-dup`"),
"the error names the offending field: {msg}"
);
let paths = engine.paths().clone();
drop(engine);
let events = read_log(&paths);
assert!(
!events.iter().any(|e| matches!(
e.kind,
EventKind::PlanApproved { .. } | EventKind::WorkerSpawned { .. }
)),
"neither consent nor spend may occur when the pack cannot load"
);
}