use crate::CliError;
use crate::commands::phase_artifact_on_develop;
use crate::pipeline_gate::{abort, run_gate};
use crate::pipeline_launch::launch_stage;
use crate::pipeline_launch::launch_stage_inner;
use crate::pipeline_outcomes::truncate_reason;
use devflow_core::gates::{GateAction, Gates};
use devflow_core::mode::{self, Mode};
use devflow_core::stage::Stage;
use devflow_core::state::{AgentKind, State};
use devflow_core::{agents, events, workflow};
use std::path::{Path, PathBuf};
pub(crate) fn worktree_writable_roots(project_root: &Path, worktree: &Path) -> Vec<PathBuf> {
let git_dir = project_root.join(".git");
let admin = std::fs::read_to_string(worktree.join(".git"))
.ok()
.and_then(|s| {
s.trim()
.strip_prefix("gitdir:")
.map(|p| PathBuf::from(p.trim()))
})
.unwrap_or_else(|| {
git_dir
.join("worktrees")
.join(worktree.file_name().unwrap_or_default())
});
vec![git_dir, admin]
}
fn agent_binary_available(program: &str) -> bool {
use std::os::unix::fs::PermissionsExt;
let executable = |path: &Path| {
path.is_file()
&& std::fs::metadata(path)
.map(|m| m.permissions().mode() & 0o111 != 0)
.unwrap_or(false)
};
if program.contains('/') {
return executable(Path::new(program));
}
std::env::var_os("PATH")
.map(|paths| std::env::split_paths(&paths).any(|dir| executable(&dir.join(program))))
.unwrap_or(false)
}
pub(crate) fn agent_program(agent: AgentKind) -> &'static str {
agents::adapter_for(agent).exec_command(0, "", &[]).0
}
pub(crate) fn ensure_agent_binary(program: &str) -> Result<(), CliError> {
if agent_binary_available(program) {
return Ok(());
}
Err(CliError::Message(format!(
"agent binary `{program}` not found — is it installed? (run `devflow doctor`)"
)))
}
fn preflight_interactivity_check(project_root: &Path, state: &State) -> Result<(), String> {
if state.agent == AgentKind::Codex
&& state.mode == Mode::Auto
&& state.stage == Stage::Define
&& !phase_artifact_on_develop(project_root, state.phase, "-CONTEXT.md")
{
return Err(format!(
"phase {} has no CONTEXT.md on develop — codex cannot run Define's \
discuss-phase interview headlessly in auto mode",
state.phase
));
}
Ok(())
}
fn gh_auth_check_applies(stage: Stage) -> bool {
stage == Stage::Ship
}
fn preflight_gh_auth_check(state: &State) -> Result<(), String> {
if !gh_auth_check_applies(state.stage) {
return Ok(());
}
match std::process::Command::new("gh")
.args(["auth", "status"])
.output()
{
Ok(output) if output.status.success() => Ok(()),
Ok(_) => Err("gh auth status reports not authenticated".to_string()),
Err(_) => {
println!(
"warning: `gh` binary not found — cannot verify GitHub credential validity \
before Ship (fail-soft, not a preflight failure)"
);
Ok(())
}
}
}
fn generic_preflight_checks(project_root: &Path, state: &State) -> Result<(), String> {
preflight_interactivity_check(project_root, state)?;
preflight_gh_auth_check(state)
}
pub(crate) fn run_preflight(
project_root: &Path,
state: &mut State,
adapter: &dyn agents::AgentAdapter,
) -> Result<bool, CliError> {
let stage = state.stage;
if let Err(reason) =
generic_preflight_checks(project_root, state).and_then(|()| adapter.preflight(state))
{
if state.preflight_retries >= mode::MAX_PREFLIGHT_RETRIES {
let ceiling_reason = format!(
"preflight retry ceiling ({}) reached for stage {stage}: {}",
mode::MAX_PREFLIGHT_RETRIES,
truncate_reason(&reason)
);
events::emit(
project_root,
state.phase,
"preflight_retry_ceiling_reached",
serde_json::json!({
"stage": stage.to_string(),
"reason": truncate_reason(&reason),
"ceiling": mode::MAX_PREFLIGHT_RETRIES,
}),
);
abort(project_root, state, &ceiling_reason)?;
return Ok(false);
}
state.preflight_retries = state.preflight_retries.saturating_add(1);
workflow::save_state(state)?;
let context = format!(
"[never-silent] preflight failed for stage {stage}: {} — human review needed \
(retry, loop-to-code, or abort)",
truncate_reason(&reason)
);
match run_gate(project_root, state, stage, &context)? {
GateAction::Advance => {
let _ = Gates::cleanup(project_root, state.phase, stage);
state.gate_pending = false;
state.preflight_retries = 0;
workflow::save_state(state)?;
launch_stage_inner(state, None, None)?;
}
GateAction::LoopBack(_) => {
let _ = Gates::cleanup(project_root, state.phase, stage);
launch_stage(state, None, None)?;
}
GateAction::Abort(reason) => abort(project_root, state, &reason)?,
}
return Ok(false);
}
if state.preflight_retries != 0 {
state.preflight_retries = 0;
workflow::save_state(state)?;
}
Ok(true)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_support::*;
#[test]
fn ensure_agent_binary_diagnoses_missing_program() {
assert!(ensure_agent_binary("sh").is_ok());
assert!(ensure_agent_binary("/bin/sh").is_ok());
let err = ensure_agent_binary("definitely-not-a-real-agent-xyz").unwrap_err();
let msg = err.to_string();
assert!(msg.contains("not found — is it installed?"), "{msg}");
assert!(msg.contains("devflow doctor"), "{msg}");
assert!(ensure_agent_binary("/nonexistent/path/agent").is_err());
}
#[test]
fn preflight_interactivity_check_flags_auto_define_without_context_md() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let mut state = State::new(60, AgentKind::Codex, Mode::Auto, root.to_path_buf());
state.stage = Stage::Define;
assert!(preflight_interactivity_check(root, &state).is_err());
state.mode = Mode::Supervise;
assert!(preflight_interactivity_check(root, &state).is_ok());
state.mode = Mode::Auto;
state.stage = Stage::Plan;
assert!(preflight_interactivity_check(root, &state).is_ok());
state.stage = Stage::Define;
state.agent = AgentKind::Claude;
assert!(
preflight_interactivity_check(root, &state).is_ok(),
"Claude/OpenCode can complete Define headlessly — only Codex is flagged"
);
state.agent = AgentKind::Codex;
let git = |args: &[&str]| {
assert!(
devflow_core::test_support::git_command(root)
.args(args)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
std::fs::create_dir_all(root.join(".planning/phases/60-widget")).unwrap();
std::fs::write(root.join(".planning/phases/60-widget/60-CONTEXT.md"), "ctx").unwrap();
git(&["add", "-A"]);
git(&["commit", "-q", "-m", "context"]);
state.stage = Stage::Define;
assert!(preflight_interactivity_check(root, &state).is_ok());
}
#[test]
fn gh_auth_check_applies_only_to_ship_stage() {
assert!(gh_auth_check_applies(Stage::Ship));
for stage in [Stage::Define, Stage::Plan, Stage::Code, Stage::Validate] {
assert!(!gh_auth_check_applies(stage));
}
}
#[test]
fn run_preflight_failing_check_gates_and_never_reaches_spawn_monitor() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 61;
let mut state = State::new(phase, AgentKind::Codex, Mode::Auto, root.to_path_buf());
state.stage = Stage::Define;
workflow::save_state(&state).unwrap();
let response_path = Gates::response_path(root, phase, Stage::Define);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(
&response_path,
r#"{"approved":false,"note":"abort: test cleanup","responded_by":"test"}"#,
)
.unwrap();
let adapter = agents::adapter_for(AgentKind::Codex);
let should_continue = run_preflight(root, &mut state, adapter.as_ref()).unwrap();
assert!(
!should_continue,
"an aborted preflight must tell its caller not to continue launch_stage"
);
assert!(
workflow::load_state(root, phase).is_err(),
"abort() must clear state — spawn_monitor was never reached"
);
let last = devflow_core::events::last_event_for_phase(root, phase)
.expect("gate_fired/gate_resolved must have been recorded");
assert_ne!(last["event"], "stage_launched");
}
#[test]
fn run_preflight_adapter_hook_override_fires() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 62;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Plan;
workflow::save_state(&state).unwrap();
let response_path = Gates::response_path(root, phase, Stage::Plan);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(
&response_path,
r#"{"approved":false,"note":"abort: test cleanup","responded_by":"test"}"#,
)
.unwrap();
let should_continue = run_preflight(root, &mut state, &AlwaysFailAdapter).unwrap();
assert!(
!should_continue,
"an aborted preflight must tell its caller not to continue launch_stage"
);
assert!(workflow::load_state(root, phase).is_err());
let last = devflow_core::events::last_event_for_phase(root, phase).unwrap();
assert_eq!(last["event"], "workflow_aborted");
}
#[test]
fn run_preflight_advance_gate_launches_agent_exactly_once() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 63;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Plan;
workflow::save_state(&state).unwrap();
let response_path = Gates::response_path(root, phase, Stage::Plan);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(&response_path, r#"{"approved":true,"responded_by":"test"}"#).unwrap();
let stub_dir = stub_agent_binary("claude");
let original_path = std::env::var_os("PATH");
let stubbed_path = prepend_path(&stub_dir, &original_path);
unsafe {
std::env::set_var("PATH", &stubbed_path);
}
let adapter = FailOnceAdapter::new();
let should_continue = run_preflight(root, &mut state, &adapter).unwrap();
if should_continue {
launch_stage(&mut state, None, None).unwrap();
}
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
assert!(
!should_continue,
"an Advance-resolved preflight failure must tell its caller not \
to continue launch_stage — the recursive retry already did"
);
let launches = stage_launched_count(root, phase);
assert_eq!(
launches, 1,
"a preflight failure resolved by Advance must launch the agent \
exactly once, not {launches}"
);
}
#[test]
fn run_preflight_loopback_gate_launches_agent_exactly_once() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 64;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Plan;
workflow::save_state(&state).unwrap();
let response_path = Gates::response_path(root, phase, Stage::Plan);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(
&response_path,
r#"{"approved":false,"note":"retry","responded_by":"test"}"#,
)
.unwrap();
let stub_dir = stub_agent_binary("claude");
let original_path = std::env::var_os("PATH");
let stubbed_path = prepend_path(&stub_dir, &original_path);
unsafe {
std::env::set_var("PATH", &stubbed_path);
}
let adapter = FailOnceAdapter::new();
let should_continue = run_preflight(root, &mut state, &adapter).unwrap();
if should_continue {
launch_stage(&mut state, None, None).unwrap();
}
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
assert!(
!should_continue,
"a LoopBack-resolved preflight failure must tell its caller not \
to continue launch_stage — the recursive retry already did"
);
let launches = stage_launched_count(root, phase);
assert_eq!(
launches, 1,
"a preflight failure resolved by LoopBack must launch the agent \
exactly once, not {launches}"
);
}
#[test]
fn run_preflight_advance_skips_recheck_on_idempotently_failing_check() {
let _guard = ENV_MUTEX.lock().unwrap();
let original_gate_timeout = std::env::var_os("DEVFLOW_GATE_TIMEOUT_SECS");
unsafe {
std::env::set_var("DEVFLOW_GATE_TIMEOUT_SECS", "2");
}
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 620;
let mut state = State::new(phase, AgentKind::Codex, Mode::Auto, root.to_path_buf());
state.stage = Stage::Define;
workflow::save_state(&state).unwrap();
let response_path = Gates::response_path(root, phase, Stage::Define);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(&response_path, r#"{"approved":true,"responded_by":"test"}"#).unwrap();
let agent_dir = agent_free_dir_with_agent_stub("codex");
let original_path = std::env::var_os("PATH");
unsafe {
std::env::set_var("PATH", agent_dir.path());
}
let result = run_preflight(root, &mut state, &AlwaysFailAdapter);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
match &original_gate_timeout {
Some(value) => std::env::set_var("DEVFLOW_GATE_TIMEOUT_SECS", value),
None => std::env::remove_var("DEVFLOW_GATE_TIMEOUT_SECS"),
}
}
assert!(
matches!(result, Ok(false)),
"Advance on a preflight gate must skip the just-adjudicated \
check and return Ok(false), not {result:?}"
);
assert!(
!Gates::gate_path(root, phase, Stage::Define).exists(),
"no second gate should ever be written once Advance skips the recheck"
);
assert_eq!(
state.preflight_retries, 0,
"a human Advance must reset the retry counter"
);
}
#[test]
fn run_preflight_loopback_bounds_recursion() {
let _guard = ENV_MUTEX.lock().unwrap();
let original_gate_timeout = std::env::var_os("DEVFLOW_GATE_TIMEOUT_SECS");
unsafe {
std::env::set_var("DEVFLOW_GATE_TIMEOUT_SECS", "2");
}
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 621;
let mut state = State::new(phase, AgentKind::Codex, Mode::Auto, root.to_path_buf());
state.stage = Stage::Define;
state.preflight_retries = mode::MAX_PREFLIGHT_RETRIES - 1;
workflow::save_state(&state).unwrap();
let response_path = Gates::response_path(root, phase, Stage::Define);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(
&response_path,
r#"{"approved":false,"note":"retry","responded_by":"test"}"#,
)
.unwrap();
let agent_dir = agent_free_dir_with_agent_stub("codex");
let original_path = std::env::var_os("PATH");
unsafe {
std::env::set_var("PATH", agent_dir.path());
}
let result = run_preflight(root, &mut state, &AlwaysFailAdapter);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
match &original_gate_timeout {
Some(value) => std::env::set_var("DEVFLOW_GATE_TIMEOUT_SECS", value),
None => std::env::remove_var("DEVFLOW_GATE_TIMEOUT_SECS"),
}
}
assert!(
matches!(result, Ok(false)),
"the ceiling must abort cleanly, not error out, got {result:?}"
);
assert!(
workflow::load_state(root, phase).is_err(),
"the ceiling must abort() and clear state, not leave it gate_pending forever"
);
let last = devflow_core::events::last_event_for_phase(root, phase)
.expect("a ceiling or abort event must have been recorded");
assert!(
last["event"] == "preflight_retry_ceiling_reached"
|| last["event"] == "workflow_aborted",
"expected a ceiling or abort event, got {last:?}"
);
}
#[test]
fn preflight_retries_reset_on_pass() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 622;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Plan;
state.preflight_retries = 2;
workflow::save_state(&state).unwrap();
let adapter = agents::adapter_for(AgentKind::Claude);
let result = run_preflight(root, &mut state, adapter.as_ref());
assert!(
matches!(result, Ok(true)),
"a passing preflight must return Ok(true), got {result:?}"
);
assert_eq!(
state.preflight_retries, 0,
"the in-memory counter must reset immediately on a pass"
);
let reloaded = workflow::load_state(root, phase).unwrap();
assert_eq!(
reloaded.preflight_retries, 0,
"the reset must be persisted to disk, not just held in memory"
);
}
}