use std::path::Path;
use std::sync::atomic::Ordering;
use std::sync::Arc;
use serde_json::json;
use super::budget::SessionDeadline;
use super::contract::{evaluate_contract_with_baselines, BaselineCaptures, OutcomeContract};
use super::native_loop::{LoopFailure, LoopOutcome, TurnGenerator};
use super::session::{CancelFlag, CoderEventKind, EventSink, IntegratedSubtask};
use super::shell_tool::{tail, WorktreeExecutor};
#[derive(Debug)]
pub enum ForemanFallback {
SingleSessionPreferred,
PlanInvalid(String),
NothingAccepted(String),
IntegrationRejected(String),
}
impl ForemanFallback {
pub fn reason(&self) -> String {
match self {
Self::SingleSessionPreferred => {
"plan prefers a single session (no parallel speedup)".into()
}
Self::PlanInvalid(e) => format!("decomposition invalid: {e}"),
Self::NothingAccepted(e) => format!("no subtask passed the merge gate: {e}"),
Self::IntegrationRejected(e) => format!("union integration rejected: {e}"),
}
}
}
fn union_goal_command(contract: &OutcomeContract) -> Option<Vec<String>> {
let chain: Vec<&str> = contract
.checks
.iter()
.filter(|c| c.expect_exit_zero && c.output_contains.is_none())
.map(|c| c.command.as_str())
.collect();
if chain.is_empty() {
return None;
}
Some(vec!["sh".into(), "-lc".into(), chain.join(" && ")])
}
fn regression_command(worktree: &Path) -> Option<Vec<String>> {
let candidates: [(&str, &str); 3] = [
("Cargo.toml", "cargo check"),
("go.mod", "go build ./..."),
("Package.swift", "swift build"),
];
candidates
.iter()
.find(|(marker, _)| worktree.join(marker).exists())
.map(|(_, cmd)| vec!["sh".into(), "-lc".into(), (*cmd).into()])
}
pub struct ForemanRun {
pub outcome: LoopOutcome,
pub integrated: Vec<IntegratedSubtask>,
}
impl ForemanRun {
fn nothing_integrated(outcome: LoopOutcome) -> Self {
Self {
outcome,
integrated: Vec::new(),
}
}
}
fn apply_patch(worktree: &Path, subtask_id: &str, patch: &str) -> Result<(), String> {
use std::io::Write;
let mut file = tempfile::NamedTempFile::new()
.map_err(|e| format!("temp patch file for {subtask_id}: {e}"))?;
file.write_all(patch.as_bytes())
.map_err(|e| format!("write patch {subtask_id}: {e}"))?;
let out = std::process::Command::new("git")
.arg("-C")
.arg(worktree)
.args(["apply", "--whitespace=nowarn"])
.arg(file.path())
.output()
.map_err(|e| format!("git apply {subtask_id}: {e}"))?;
if out.status.success() {
Ok(())
} else {
Err(format!(
"git apply {subtask_id} failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
))
}
}
async fn drain_gate_audit(
infra: &car_multi::SharedInfra,
sink: &Arc<EventSink>,
from: usize,
) -> usize {
let log = infra.log.lock().await;
let events = log.events();
for event in events.iter().skip(from) {
if infra.gate_audit_scope.as_deref().is_none()
|| event
.data
.get("gate_audit_scope")
.and_then(serde_json::Value::as_str)
!= infra.gate_audit_scope.as_deref()
{
continue;
}
let decision = match event.kind {
car_eventlog::EventKind::GateAccepted => "accepted",
car_eventlog::EventKind::GateRejected => "rejected",
_ => continue,
};
sink.record_gate_verdict(event.kind.clone(), event.data.clone());
let mut raw = serde_json::Map::new();
raw.insert("foreman".to_string(), json!("gate"));
raw.insert("decision".to_string(), json!(decision));
for (k, v) in &event.data {
raw.insert(k.clone(), v.clone());
}
sink.emit(CoderEventKind::ExternalEvent {
raw: serde_json::Value::Object(raw),
});
}
events.len()
}
pub async fn run_foreman_loop(
adapter_id: &str,
intent: &str,
contract: &OutcomeContract,
executor: &WorktreeExecutor,
sink: &Arc<EventSink>,
cancel: &CancelFlag,
generator: &Arc<dyn TurnGenerator>,
mcp_endpoint: Option<&str>,
mcp_config_dir: Option<&std::path::Path>,
infra: &car_multi::SharedInfra,
deadline: &Arc<SessionDeadline>,
workers: Option<&dyn car_multi::WorktreeAgent>,
baseline_captures: &BaselineCaptures,
) -> Result<ForemanRun, ForemanFallback> {
let worktree = executor.worktree().to_path_buf();
let cancelled = || {
LoopOutcome::lost(
LoopFailure::Cancelled,
Some("cancelled".into()),
0,
Vec::new(),
)
};
if let Some(reason) = deadline.admit() {
sink.emit(CoderEventKind::BudgetExhausted {
reason: reason.clone(),
elapsed_secs: deadline.elapsed_secs(),
iterations: 0,
});
return Ok(ForemanRun::nothing_integrated(LoopOutcome::lost(
LoopFailure::BudgetExhausted,
Some(reason),
0,
Vec::new(),
)));
}
if cancel.load(Ordering::SeqCst) {
return Ok(ForemanRun::nothing_integrated(cancelled()));
}
sink.emit(CoderEventKind::ExternalEvent {
raw: json!({ "foreman": "planning", "adapter": adapter_id }),
});
let plan_generator = generator.clone();
let plan = car_multi::decompose(&worktree, intent, 3, move |prompt| {
let generator = plan_generator.clone();
async move {
generator
.generate(car_inference::GenerateRequest {
prompt,
intent: car_inference::IntentHint::high_stakes_if(true),
..Default::default()
})
.await
.map(|r| r.text)
}
})
.await;
if !plan.is_valid() {
return Err(ForemanFallback::PlanInvalid(plan.issues.join("; ")));
}
sink.emit(CoderEventKind::ExternalEvent {
raw: json!({
"foreman": "planned",
"subtasks": plan.subtasks.len(),
"levels": plan.levels.len(),
"prefer_single_session": plan.prefer_single_session,
}),
});
if plan.prefer_single_session {
return Err(ForemanFallback::SingleSessionPreferred);
}
if cancel.load(Ordering::SeqCst) {
return Ok(ForemanRun::nothing_integrated(cancelled()));
}
let local = car_external_agents::ForemanExternalAgent::new(adapter_id.to_string());
let agent: &dyn car_multi::WorktreeAgent = workers.unwrap_or(&local);
let scoped_infra = infra.scoped_gate_audit(uuid::Uuid::new_v4().to_string());
let infra = &scoped_infra;
let gate_audit_from = infra.log.lock().await.events().len();
let config = car_multi::FarmOutConfig {
verify_command: regression_command(&worktree),
union_verify_command: union_goal_command(contract),
mcp_endpoint: mcp_endpoint.map(String::from),
mcp_config_dir: mcp_config_dir.map(std::path::Path::to_path_buf),
..Default::default()
};
let progress_sink: car_multi::ForemanProgressSink = {
let sink = Arc::clone(sink);
Arc::new(move |ev: car_multi::ForemanProgress| {
let raw = match ev {
car_multi::ForemanProgress::SubtaskStarted {
subtask_id,
index,
level,
total,
} => json!({
"foreman": "subtask_started",
"subtask_id": subtask_id,
"index": index,
"level": level,
"total": total,
}),
car_multi::ForemanProgress::SubtaskVerifying { subtask_id } => json!({
"foreman": "subtask_verifying",
"subtask_id": subtask_id,
}),
car_multi::ForemanProgress::SubtaskGated {
subtask_id,
accepted,
status,
} => json!({
"foreman": "subtask_gated",
"subtask_id": subtask_id,
"accepted": accepted,
"status": status,
}),
};
sink.emit(CoderEventKind::ExternalEvent { raw });
})
};
let farmed = car_multi::run_farm_out_with_progress(
&worktree,
&plan.subtasks,
agent,
&config,
infra,
progress_sink,
)
.await;
let audited = drain_gate_audit(infra, sink, gate_audit_from).await;
let accepted: Vec<(String, String)> = farmed
.outcomes
.iter()
.filter(|o| o.is_accepted())
.filter_map(|o| o.patch.clone().map(|p| (o.subtask_id.clone(), p)))
.collect();
sink.emit(CoderEventKind::ExternalEvent {
raw: json!({
"foreman": "farmed",
"accepted": accepted.len(),
"total": farmed.outcomes.len(),
}),
});
if accepted.is_empty() {
let detail = farmed
.outcomes
.iter()
.filter_map(|o| o.error.as_deref())
.take(3)
.collect::<Vec<_>>()
.join("; ");
return Err(ForemanFallback::NothingAccepted(if detail.is_empty() {
format!(
"{} subtask(s) all rejected or inconclusive",
farmed.outcomes.len()
)
} else {
detail
}));
}
if cancel.load(Ordering::SeqCst) {
return Ok(ForemanRun::nothing_integrated(cancelled()));
}
let label = format!("coder-{}", sink_label(&worktree));
let integration =
car_multi::integrate_and_verify(&worktree, &label, &accepted, &config, infra).await;
let _ = drain_gate_audit(infra, sink, audited).await;
let integration =
integration.map_err(|e| ForemanFallback::IntegrationRejected(e.to_string()))?;
if !integration.integrated_cleanly() {
if let Some(blame) = &integration.blame {
let reason = if !blame.apply_conflicts.is_empty() {
"patch conflict"
} else if !blame.duplicate_conflicts.is_empty() {
"duplicate declaration"
} else if blame.build_test.is_some() {
"build/test failed"
} else {
"rejected"
};
let detail = blame
.apply_conflicts
.first()
.map(|c| format!("{} did not apply", c.subtask_id))
.or_else(|| {
blame
.duplicate_conflicts
.first()
.map(|d| format!("duplicate `{}` in {}", d.symbol, d.file))
})
.or_else(|| blame.build_test.as_ref().map(|b| tail(&b.output_tail, 200)));
let implicated: Vec<String> = blame.implicated_subtasks().into_iter().collect();
sink.emit(CoderEventKind::ExternalEvent {
raw: json!({
"foreman": "union_rejected",
"reason": reason,
"implicated": implicated,
"detail": detail,
}),
});
}
return Err(ForemanFallback::IntegrationRejected(format!(
"applied {}, conflicts: [{}], union verdict accepting: {}",
integration.applied,
integration.apply_conflicts.join(", "),
integration
.verdict
.as_ref()
.is_some_and(|v| v.is_accepted()),
)));
}
sink.emit(CoderEventKind::ExternalEvent {
raw: json!({ "foreman": "union_verified", "applied": integration.applied }),
});
let mut integrated = Vec::with_capacity(accepted.len());
for (subtask_id, patch) in &accepted {
apply_patch(&worktree, subtask_id, patch).map_err(ForemanFallback::IntegrationRejected)?;
integrated.push(IntegratedSubtask {
subtask_id: subtask_id.clone(),
files: car_multi::files_in_patch(patch),
});
sink.emit(CoderEventKind::ToolResult {
tool: "foreman.apply".into(),
ok: true,
preview: format!("applied {subtask_id}"),
});
}
let last_results =
evaluate_contract_with_baselines(contract, executor, sink, baseline_captures).await;
let passed = last_results.iter().all(|r| r.passed);
let outcome = if passed {
LoopOutcome::green(1, last_results)
} else {
LoopOutcome::lost(LoopFailure::Verification, None, 1, last_results)
};
Ok(ForemanRun {
outcome,
integrated,
})
}
fn sink_label(worktree: &Path) -> String {
worktree
.file_name()
.map(|n| n.to_string_lossy().into_owned())
.unwrap_or_else(|| "session".into())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::coder::contract::ContractCheck;
#[tokio::test]
async fn gate_verdicts_reach_the_session_journal() {
use std::collections::HashMap;
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("s1.events.jsonl");
let sink = Arc::new(EventSink::new("s1", None, Some(journal.clone())));
let infra = car_multi::SharedInfra::new().scoped_gate_audit("s1-run".into());
{
let mut log = infra.log.lock().await;
let mut accepted = HashMap::new();
accepted.insert("subtask".to_string(), json!("a"));
accepted.insert("gate_audit_scope".to_string(), json!("s1-run"));
accepted.insert("build_test".to_string(), json!("passed"));
log.append(car_eventlog::EventKind::GateAccepted, None, None, accepted);
let mut rejected = HashMap::new();
rejected.insert("subtask".to_string(), json!("b"));
rejected.insert("gate_audit_scope".to_string(), json!("s1-run"));
rejected.insert("reasons".to_string(), json!(["containment"]));
log.append(car_eventlog::EventKind::GateRejected, None, None, rejected);
log.append(
car_eventlog::EventKind::RunStarted,
None,
None,
HashMap::new(),
);
}
let cursor = drain_gate_audit(&infra, &sink, 0).await;
assert_eq!(
cursor, 3,
"the cursor counts the whole log, not the matches"
);
let cursor2 = drain_gate_audit(&infra, &sink, cursor).await;
assert_eq!(cursor2, cursor);
drop(sink);
let body = std::fs::read_to_string(&journal).expect("the session journal exists");
let lines: Vec<&str> = body.lines().filter(|l| !l.trim().is_empty()).collect();
assert_eq!(
lines.len(),
2,
"both verdicts, once each, and nothing else: {body}"
);
assert!(
body.contains("gate_accepted") || body.contains("GateAccepted"),
"{body}"
);
assert!(
body.contains("gate_rejected") || body.contains("GateRejected"),
"{body}"
);
assert!(body.contains("containment"), "{body}");
assert!(!body.to_lowercase().contains("run_started"), "{body}");
}
#[tokio::test]
async fn concurrent_gate_runs_project_only_their_own_verdicts() {
let dir = tempfile::tempdir().unwrap();
let shared = car_multi::SharedInfra::new();
let a = shared.scoped_gate_audit("run-a".into());
let b = shared.scoped_gate_audit("run-b".into());
assert!(Arc::ptr_eq(&a.log, &b.log));
assert!(Arc::ptr_eq(&a.state, &b.state));
assert!(Arc::ptr_eq(&a.policies, &b.policies));
assert!(Arc::ptr_eq(&a.budget, &b.budget));
let a_path = dir.path().join("a.events.jsonl");
let b_path = dir.path().join("b.events.jsonl");
let a_sink = Arc::new(EventSink::new("a", None, Some(a_path.clone())));
let b_sink = Arc::new(EventSink::new("b", None, Some(b_path.clone())));
let a_from = shared.log.lock().await.events().len();
let b_from = a_from;
let command = if cfg!(windows) {
vec!["cmd".into(), "/C".into(), "exit 0".into()]
} else {
vec!["sh".into(), "-c".into(), "exit 0".into()]
};
let accepted =
car_multi::GateConfig::new("same-subtask", dir.path()).with_verify_command(command);
let rejected = car_multi::GateConfig::new("same-subtask", dir.path());
let footprint = car_multi::DeclaredFootprint::unconstrained();
let (a_verdict, b_verdict) = tokio::join!(
car_multi::verify_changes(&accepted, &[], &footprint, &a),
car_multi::verify_changes(&rejected, &[], &footprint, &b),
);
assert!(a_verdict.is_accepted());
assert!(!b_verdict.is_accepted());
assert_eq!(shared.log.lock().await.events().len(), 2);
let a_cursor = drain_gate_audit(&a, &a_sink, a_from).await;
let b_cursor = drain_gate_audit(&b, &b_sink, b_from).await;
assert_eq!(a_cursor, 2);
assert_eq!(b_cursor, 2);
drain_gate_audit(&a, &a_sink, a_cursor).await;
drain_gate_audit(&b, &b_sink, b_cursor).await;
drop(a_sink);
drop(b_sink);
let a_body = std::fs::read_to_string(a_path).unwrap();
let b_body = std::fs::read_to_string(b_path).unwrap();
assert_eq!(a_body.lines().count(), 1, "{a_body}");
assert_eq!(b_body.lines().count(), 1, "{b_body}");
assert!(
a_body.contains("run-a") && !a_body.contains("run-b"),
"{a_body}"
);
assert!(
b_body.contains("run-b") && !b_body.contains("run-a"),
"{b_body}"
);
}
#[tokio::test]
async fn a_gate_tagged_stream_event_cannot_forge_a_verdict() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("s2.events.jsonl");
let sink = Arc::new(EventSink::new("s2", None, Some(journal.clone())));
sink.emit(CoderEventKind::StateChanged {
from: "created".to_string(),
to: "running".to_string(),
});
sink.emit(CoderEventKind::ExternalEvent {
raw: json!({
"type": "system",
"subtype": "init",
"session_id": "s2",
"foreman": "gate",
"decision": "accepted",
"subtask": "peer-authored-patch",
"build_test": "passed",
}),
});
drop(sink);
let body = std::fs::read_to_string(&journal)
.expect("the positive control wrote a line, so the journal exists");
let lowered = body.to_lowercase();
assert!(
lowered.contains("state_changed"),
"the control did not journal, so this test proves nothing: {body}"
);
assert!(
!lowered.contains("gate_accepted") && !lowered.contains("gateaccepted"),
"a stream event forged a gate verdict into the audit record: {body}"
);
assert!(
!lowered.contains("gate_rejected") && !lowered.contains("gaterejected"),
"a stream event forged a gate verdict into the audit record: {body}"
);
}
#[derive(Default)]
struct RecordingAgent {
called: std::sync::atomic::AtomicUsize,
}
#[async_trait::async_trait]
impl car_multi::WorktreeAgent for RecordingAgent {
async fn run_in(
&self,
_req: &car_multi::WorktreeAgentRequest<'_>,
) -> Result<car_multi::AgentRunSummary, car_multi::ForemanError> {
self.called
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Ok(car_multi::AgentRunSummary {
answer: "recorded".into(),
})
}
}
fn git_repo() -> tempfile::TempDir {
let dir = tempfile::tempdir().unwrap();
for args in [
vec!["init", "-q", "-b", "main"],
vec!["config", "user.email", "t@t.t"],
vec!["config", "user.name", "t"],
] {
let out = std::process::Command::new("git")
.args(&args)
.current_dir(dir.path())
.output()
.expect("git");
assert!(out.status.success(), "git {args:?}");
}
std::fs::write(dir.path().join("seed.txt"), "seed\n").unwrap();
for args in [vec!["add", "-A"], vec!["commit", "-qm", "seed"]] {
std::process::Command::new("git")
.args(&args)
.current_dir(dir.path())
.output()
.expect("git");
}
dir
}
#[tokio::test]
async fn session_deny_tool_blocks_delivery_and_journals_rejection() {
struct WriteDeclaredFile;
#[async_trait::async_trait]
impl car_multi::WorktreeAgent for WriteDeclaredFile {
async fn run_in(
&self,
req: &car_multi::WorktreeAgentRequest<'_>,
) -> Result<car_multi::AgentRunSummary, car_multi::ForemanError> {
let path = req.cwd.join(format!("src/{}.rs", req.subtask.id));
std::fs::write(&path, format!("pub fn {}() {{}}\n", req.subtask.id))
.map_err(|e| car_multi::ForemanError::Agent(e.to_string()))?;
Ok(car_multi::AgentRunSummary::default())
}
}
struct TwoFilePlan;
#[async_trait::async_trait]
impl TurnGenerator for TwoFilePlan {
async fn generate(
&self,
_req: car_inference::GenerateRequest,
) -> Result<car_inference::InferenceResult, String> {
Ok(serde_json::from_value(serde_json::json!({
"text": r#"{"subtasks":[
{"id":"x","prompt":"x","writes":[{"file":"src/x.rs","symbol":"x"}]},
{"id":"y","prompt":"y","writes":[{"file":"src/y.rs","symbol":"y"}]}
]}"#,
"tool_calls": [],
"trace_id": "shared-infra-policy-test",
"model_used": "scripted",
"latency_ms": 0,
}))
.expect("scripted InferenceResult shape"))
}
}
let repo = git_repo();
std::fs::create_dir_all(repo.path().join("src")).unwrap();
std::fs::write(
repo.path().join("Cargo.toml"),
"[package]\nname = \"shared-infra-test\"\nversion = \"0.1.0\"\nedition = \"2021\"\n",
)
.unwrap();
std::fs::write(repo.path().join("src/lib.rs"), "pub fn seed() {}\n").unwrap();
for args in [vec!["add", "-A"], vec!["commit", "-qm", "cargo seed"]] {
let out = std::process::Command::new("git")
.args(&args)
.current_dir(repo.path())
.output()
.expect("git");
assert!(out.status.success(), "git {args:?}");
}
let audit_dir = tempfile::tempdir().unwrap();
let journal = audit_dir.path().join("client-session.jsonl");
let shared_state = Arc::new(car_state::StateStore::new());
let shared_log = Arc::new(tokio::sync::Mutex::new(
car_eventlog::EventLog::with_journal(journal.clone()),
));
let shared_policies = Arc::new(tokio::sync::RwLock::new(car_policy::PolicyEngine::new()));
{
let mut policy_engine = shared_policies.write().await;
car_policy::PolicyRules::from_toml(r#"deny_tool = ["foreman.integrate"]"#)
.unwrap()
.apply(&mut policy_engine);
}
let infra = car_multi::SharedInfra::with_shared(
Arc::clone(&shared_state),
Arc::clone(&shared_log),
Arc::clone(&shared_policies),
);
let sink = Arc::new(EventSink::test_sink());
let result = run_foreman_loop(
"scripted",
"write x and y",
&OutcomeContract {
allow_credentials: false,
description: "both files exist".into(),
checks: vec![check("goal", "true", true, None)],
},
&WorktreeExecutor::new(repo.path().to_path_buf()),
&sink,
&CancelFlag::default(),
&(Arc::new(TwoFilePlan) as Arc<dyn TurnGenerator>),
None, None, &infra,
&Arc::new(SessionDeadline::new(Some(300))),
Some(&WriteDeclaredFile),
&Default::default(), )
.await;
assert!(
matches!(result, Err(ForemanFallback::NothingAccepted(_))),
"the session deny rule must reject every patch"
);
assert!(
!repo.path().join("src/x.rs").exists() && !repo.path().join("src/y.rs").exists(),
"a policy-rejected patch must not reach the delivery worktree"
);
{
let log = shared_log.lock().await;
assert_eq!(log.events().len(), 2, "one decision per patch");
assert!(
log.events()
.iter()
.all(|event| event.kind == car_eventlog::EventKind::GateRejected),
"every runtime audit decision must be a rejection"
);
assert!(
log.events()
.iter()
.all(
|event| event.data.get("reasons").is_some_and(|reasons| reasons
.to_string()
.contains("deny_tool:foreman.integrate"))
),
"the journal must retain the loaded rule as the rejection reason"
);
}
drop(infra);
drop(shared_log);
let body = std::fs::read_to_string(&journal).expect("session audit journal exists");
let rows: Vec<serde_json::Value> = body
.lines()
.map(|line| serde_json::from_str(line).expect("valid journal JSONL"))
.collect();
assert_eq!(rows.len(), 2, "one durable decision per patch: {body}");
assert!(
rows.iter().all(|row| {
row.get("kind")
.is_some_and(|kind| kind.to_string().to_lowercase().contains("gate_rejected"))
&& row.get("data").is_some_and(|data| {
data.to_string().contains("deny_tool:foreman.integrate")
})
}),
"both durable rows must be policy gate rejections: {body}"
);
}
#[tokio::test]
async fn the_supplied_worker_is_the_one_that_runs_the_subtasks() {
let repo = git_repo();
let recorder = RecordingAgent::default();
struct Plan;
#[async_trait::async_trait]
impl TurnGenerator for Plan {
async fn generate(
&self,
_req: car_inference::GenerateRequest,
) -> Result<car_inference::InferenceResult, String> {
Ok(serde_json::from_value(serde_json::json!({
"text": r#"{"subtasks":[
{"id":"x","prompt":"x","writes":[{"file":"x.rs","symbol":"x"}]},
{"id":"y","prompt":"y","writes":[{"file":"y.rs","symbol":"y"}]}
]}"#,
"tool_calls": [],
"trace_id": "foreman-pool-test",
"model_used": "scripted",
"latency_ms": 0,
}))
.expect("scripted InferenceResult shape"))
}
}
let sink = Arc::new(EventSink::test_sink());
let contract = OutcomeContract {
allow_credentials: false,
description: "two things".into(),
checks: vec![check("c", "true", true, None)],
};
let executor = WorktreeExecutor::new(repo.path().to_path_buf());
let _ = run_foreman_loop(
"claude-code",
"two things",
&contract,
&executor,
&sink,
&CancelFlag::default(),
&(Arc::new(Plan) as Arc<dyn TurnGenerator>),
None, None, &car_multi::SharedInfra::new(),
&Arc::new(SessionDeadline::new(Some(300))),
Some(&recorder),
&BaselineCaptures::new(),
)
.await;
assert!(
recorder.called.load(std::sync::atomic::Ordering::SeqCst) > 0,
"the supplied worker must be the one that runs the subtasks"
);
}
fn check(name: &str, command: &str, exit_zero: bool, contains: Option<&str>) -> ContractCheck {
ContractCheck {
name: name.into(),
command: command.into(),
expect_exit_zero: exit_zero,
output_contains: contains.map(String::from),
timeout_secs: 60,
baseline: false,
differential: None,
}
}
#[test]
fn union_goal_chains_plain_exit_zero_checks_only() {
let contract = OutcomeContract {
allow_credentials: false,
description: "d".into(),
checks: vec![
check("build", "cargo build", true, None),
check("tests", "cargo test", true, None),
check("output", "cat x.txt", true, Some("needle")), check("inverted", "grep -q bad src/", false, Some("x")), ],
};
let cmd = union_goal_command(&contract).unwrap();
assert_eq!(cmd[0], "sh");
assert_eq!(cmd[2], "cargo build && cargo test");
}
#[test]
fn no_expressible_checks_means_no_union_goal_command() {
let contract = OutcomeContract {
allow_credentials: false,
description: "d".into(),
checks: vec![check("output", "cat x.txt", true, Some("needle"))],
};
assert!(union_goal_command(&contract).is_none());
}
#[test]
fn regression_command_maps_known_build_systems_only() {
let dir = tempfile::tempdir().unwrap();
assert!(
regression_command(dir.path()).is_none(),
"unknown repo → None (fail-closed)"
);
std::fs::write(dir.path().join("Cargo.toml"), "[package]").unwrap();
let cmd = regression_command(dir.path()).unwrap();
assert_eq!(cmd[2], "cargo check");
}
#[test]
fn apply_patch_lands_changes_in_worktree() {
let dir = tempfile::tempdir().unwrap();
for args in [
vec!["init", "-q", "-b", "main"],
vec![
"-c",
"user.name=t",
"-c",
"user.email=t@t",
"commit",
"-q",
"--allow-empty",
"-m",
"init",
],
] {
assert!(std::process::Command::new("git")
.arg("-C")
.arg(dir.path())
.args(&args)
.output()
.unwrap()
.status
.success());
}
let patch = "diff --git a/new.txt b/new.txt\nnew file mode 100644\n--- /dev/null\n+++ b/new.txt\n@@ -0,0 +1 @@\n+from foreman\n";
apply_patch(dir.path(), "s1", patch).unwrap();
assert_eq!(
std::fs::read_to_string(dir.path().join("new.txt"))
.unwrap()
.replace("\r\n", "\n"),
"from foreman\n"
);
}
#[test]
fn apply_patch_conflict_is_reported_not_panicked() {
let dir = tempfile::tempdir().unwrap();
assert!(std::process::Command::new("git")
.arg("-C")
.arg(dir.path())
.args(["init", "-q"])
.output()
.unwrap()
.status
.success());
let err = apply_patch(dir.path(), "s1", "not a patch").unwrap_err();
assert!(err.contains("git apply s1 failed"), "{err}");
}
#[test]
fn fallback_reasons_are_descriptive() {
assert!(ForemanFallback::SingleSessionPreferred
.reason()
.contains("single session"));
assert!(ForemanFallback::PlanInvalid("x".into())
.reason()
.contains("decomposition"));
assert!(ForemanFallback::NothingAccepted("y".into())
.reason()
.contains("merge gate"));
assert!(ForemanFallback::IntegrationRejected("z".into())
.reason()
.contains("integration"));
}
}