use crate::CliError;
use crate::pipeline_gate::transition;
use crate::pipeline_outcomes::{
ValidateOutcome, classify_validate_outcome, handle_infra_outcome, handle_rate_limited_outcome,
handle_ship_failure, handle_ship_outcome, handle_stage_failure, handle_validate_outcome,
truncate_reason,
};
use crate::preflight::{ensure_agent_binary, run_preflight, worktree_writable_roots};
use crate::staleness::enforce_build_staleness;
use devflow_core::config::{GitFlowConfig, capture_retention};
use devflow_core::outcome_policy::{self, Action};
use devflow_core::prompt;
use devflow_core::stage::Stage;
use devflow_core::state::State;
use devflow_core::{agent_result, agents, events, lock, monitor, workflow};
use std::path::Path;
pub(crate) fn launch_stage_inner(
state: &mut State,
prompt_override: Option<String>,
archived_stage: Option<Stage>,
) -> Result<(), CliError> {
state.monitor_pid = None;
workflow::save_state(state)?;
let prompt = prompt_override.unwrap_or_else(|| {
prompt::stage_prompt_for_project(state.stage, state.phase, &state.project_root)
});
let adapter = agents::adapter_for(state.agent);
let roots = state
.worktree_path
.as_deref()
.map(|wt| worktree_writable_roots(&state.project_root, wt))
.unwrap_or_default();
let (program, args) = adapter.exec_command(state.phase, &prompt, &roots);
ensure_agent_binary(program)?;
let project_root = state.project_root.clone();
enforce_build_staleness(
&project_root,
state,
env!("DEVFLOW_BUILD_COMMIT"),
env!("DEVFLOW_BUILD_DIRTY") == "true",
)?;
if let Some(stamp) = agent_result::archive_phase_files(
&state.project_root,
state
.worktree_path
.as_deref()
.unwrap_or(&state.project_root),
state.phase,
capture_retention(&state.project_root),
)
.map_err(|err| {
CliError::Message(format!(
"could not archive phase {} capture before rollover: {err}",
state.phase
))
})? {
events::emit(
&state.project_root,
state.phase,
"capture_archived",
serde_json::json!({
"stage": archived_stage.unwrap_or(state.stage).to_string(),
"to_stage": state.stage.to_string(),
"stamp": stamp,
}),
);
}
let pid = monitor::spawn_monitor(state, program, &args, &adapter.extra_env())
.map_err(|err| CliError::Message(format!("could not spawn monitor: {err}")))?;
state.monitor_pid = Some(pid);
workflow::save_state(state)?;
let _ = devflow_core::registry::register(&state.project_root, state.phase);
events::emit(
&state.project_root,
state.phase,
"stage_launched",
serde_json::json!({
"stage": state.stage.to_string(),
"agent": state.agent.to_string(),
"monitor_pid": pid,
}),
);
println!(
"stage {} → launched {} (monitor pid {pid})",
state.stage,
adapter.name()
);
Ok(())
}
pub(crate) fn launch_stage(
state: &mut State,
prompt_override: Option<String>,
archived_stage: Option<Stage>,
) -> Result<(), CliError> {
let adapter = agents::adapter_for(state.agent);
let prompt = prompt_override.clone().unwrap_or_else(|| {
prompt::stage_prompt_for_project(state.stage, state.phase, &state.project_root)
});
let roots = state
.worktree_path
.as_deref()
.map(|wt| worktree_writable_roots(&state.project_root, wt))
.unwrap_or_default();
let (program, _args) = adapter.exec_command(state.phase, &prompt, &roots);
ensure_agent_binary(program)?;
let project_root = state.project_root.clone();
if !run_preflight(&project_root, state, adapter.as_ref())? {
return Ok(());
}
launch_stage_inner(state, prompt_override, archived_stage)
}
pub(crate) fn resume(project_root: &Path, phase: u32) -> Result<(), CliError> {
let _lock = match lock::acquire(project_root, phase) {
Ok(guard) => guard,
Err(lock::LockError::Contended { pid, path: _ }) => {
return Err(CliError::Message(format!(
"another devflow process (pid {pid}) is already running"
)));
}
Err(err) => return Err(CliError::Message(format!("lock error: {err}"))),
};
let mut state = workflow::load_state(project_root, phase)?;
state.stopped = false;
state.stop_reason = None;
state.stop_until = None;
workflow::save_state(&state)?;
launch_stage(&mut state, None, None)
}
pub(crate) fn single_active_phase(project_root: &Path) -> Result<Option<u32>, CliError> {
let states = workflow::list_states(project_root);
match states.as_slice() {
[] => Ok(None),
[one] => Ok(Some(one.phase)),
many => Err(CliError::Message(format!(
"multiple active phases ({}) — pass --phase to pick one",
many.iter()
.map(|s| s.phase.to_string())
.collect::<Vec<_>>()
.join(", ")
))),
}
}
pub(crate) fn resolve_sole_active_phase(project_root: &Path) -> Result<u32, CliError> {
single_active_phase(project_root)?
.ok_or_else(|| CliError::Message("no active DevFlow state — nothing to advance".into()))
}
pub(crate) fn advance(project_root: &Path, phase: Option<u32>) -> Result<(), CliError> {
let phase = match phase {
Some(phase) => phase,
None => match resolve_sole_active_phase(project_root) {
Ok(phase) => phase,
Err(err) => {
events::emit(
project_root,
0,
"advance_failed",
serde_json::json!({ "reason": err.to_string() }),
);
return Err(err);
}
},
};
let _lock = match lock::acquire(project_root, phase) {
Ok(guard) => guard,
Err(lock::LockError::Contended { pid, path: _ }) => {
return Err(CliError::Message(format!(
"another devflow process (pid {pid}) is already running"
)));
}
Err(err) => return Err(CliError::Message(format!("lock error: {err}"))),
};
let mut state = workflow::load_state(project_root, phase)?;
let git_flow = GitFlowConfig::default();
let result = agent_result::evaluate_agent_result(project_root, &state, &git_flow)
.map_err(|err| CliError::Message(format!("could not evaluate agent result: {err}")))?;
let stage = state.stage;
println!("stage {stage} finished with status {:?}", result.status);
if let Some(reason) = &result.reason {
println!(" detail: {reason}");
}
events::emit(
project_root,
phase,
"advance_evaluated",
serde_json::json!({
"stage": stage.to_string(),
"status": result.status.as_wire_str(),
"verdict": result.verdict.map(|v| format!("{v:?}").to_ascii_lowercase()),
"decided_by_layer": result.decided_by_layer,
"reason": result.reason.as_deref().map(truncate_reason),
}),
);
match outcome_policy::decide_action(stage, result.status) {
Action::Advance => match stage {
Stage::Define => transition(project_root, &mut state, Stage::Plan),
Stage::Plan => transition(project_root, &mut state, Stage::Code),
Stage::Code => transition(project_root, &mut state, Stage::Validate),
Stage::Validate => {
handle_validate_outcome(
project_root,
&mut state,
classify_validate_outcome(&result),
)
}
Stage::Ship => handle_ship_outcome(project_root, &mut state),
},
Action::GateReview => match stage {
Stage::Validate => {
handle_validate_outcome(project_root, &mut state, ValidateOutcome::Failed)
}
Stage::Ship => handle_ship_failure(project_root, &mut state, result.reason),
_ => handle_stage_failure(project_root, &mut state, stage, result.reason),
},
Action::GateInfra => handle_infra_outcome(project_root, &mut state, stage, result.reason),
Action::AutoResume => {
handle_rate_limited_outcome(project_root, &mut state, phase, stage, result.reason)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_support::*;
use devflow_core::gates::Gates;
use devflow_core::mode::Mode;
use devflow_core::state::AgentKind;
#[test]
fn launch_stage_persists_monitor_pid_for_reload() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 65;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
workflow::save_state(&state).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 result = launch_stage(&mut state, None, None);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
result.unwrap();
assert!(
state.monitor_pid.is_some(),
"launch_stage must record the monitor pid on the in-memory state"
);
let reloaded = workflow::load_state(root, phase).unwrap();
assert_eq!(
reloaded.monitor_pid, state.monitor_pid,
"the monitor pid recorded by launch_stage must be persisted to disk, \
since transition() saves state before launch_stage runs"
);
}
#[test]
fn resume_clears_stop_marker_and_advances_past_stop_point() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 66;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Plan;
state.stop_until = Some(Stage::Plan);
state.stopped = true;
state.stop_reason = Some("stopped after plan completed (--until plan)".to_string());
workflow::save_state(&state).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 result = resume(root, phase);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
result.unwrap();
let reloaded = workflow::load_state(root, phase).unwrap();
assert!(
!reloaded.stopped,
"resume must clear stopped so the phase is no longer marked halted"
);
assert_eq!(
reloaded.stop_reason, None,
"resume must clear stop_reason alongside stopped"
);
assert_eq!(
reloaded.stop_until, None,
"resume must clear stop_until so the phase does not immediately re-stop \
the next time it advances past Plan"
);
}
#[test]
fn code_unknown_does_not_transition_to_validate() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 72;
let branch = format!("feature/phase-{phase:02}");
let git = |args: &[&str]| {
assert!(
devflow_core::test_support::git_command(root)
.args(args)
.status()
.unwrap()
.success(),
"git {args:?} failed"
);
};
git(&["checkout", "-q", "-b", &branch, "develop"]);
std::fs::write(root.join("work.txt"), "wip\n").unwrap();
git(&["add", "work.txt"]);
git(&["commit", "-q", "-m", "wip commit"]);
git(&["checkout", "-q", "develop"]);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
workflow::save_state(&state).unwrap();
let code_gate = Gates::gate_path(root, phase, Stage::Code);
let validate_gate = Gates::gate_path(root, phase, Stage::Validate);
let response_path = Gates::response_path(root, phase, Stage::Code);
std::thread::scope(|scope| {
scope.spawn(|| {
advance(root, Some(phase)).unwrap();
});
let mut seen = false;
for _ in 0..150 {
if code_gate.exists() {
seen = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
assert!(
seen,
"an Unknown Code outcome must fire a never-silent gate, not advance silently"
);
assert!(
!validate_gate.exists(),
"an Unknown Code outcome must never transition to Validate"
);
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();
});
}
#[test]
fn launch_stage_inner_clears_monitor_pid_on_early_failure() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 93;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
state.monitor_pid = Some(999_999);
workflow::save_state(&state).unwrap();
let neutral_path_dir = agent_free_git_only_path_dir();
let original_path = std::env::var_os("PATH");
unsafe {
std::env::set_var("PATH", neutral_path_dir.path());
}
let result = launch_stage_inner(&mut state, None, None);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
assert!(
result.is_err(),
"ensure_agent_binary must fail against the neutralized, agent-free PATH"
);
assert_eq!(
state.monitor_pid, None,
"an early launch failure must clear the stale monitor_pid in-memory, not carry it \
forward from the previous stage"
);
let reloaded = workflow::load_state(root, phase).unwrap();
assert_eq!(
reloaded.monitor_pid, None,
"the monitor_pid clear must be persisted to state.json, not just in-memory"
);
}
#[test]
fn advance_evaluated_emits_wire_status_and_decided_by_layer_for_resource_killed() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 78;
std::fs::create_dir_all(root.join(".devflow")).unwrap();
std::fs::write(agent_result::exit_code_path(root, phase), "137").unwrap();
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
workflow::save_state(&state).unwrap();
let response_path = Gates::response_path(root, phase, Stage::Code);
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();
advance(root, Some(phase)).unwrap();
let contents = std::fs::read_to_string(events::events_path(root)).unwrap();
let event = contents
.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.find(|e| e["event"] == "advance_evaluated")
.expect("advance_evaluated event recorded");
assert_eq!(event["status"], "resource_killed");
assert_ne!(event["status"], "resourcekilled");
assert_eq!(event["decided_by_layer"], 2);
}
}