use crate::CliError;
use crate::config_parse::gate_timeout_secs;
use crate::pipeline_launch::launch_stage;
use crate::pipeline_outcomes::{run_checkout_hooks, truncate_reason};
use devflow_core::gates::{self, GateAction, Gates};
use devflow_core::hooks;
use devflow_core::mode;
use devflow_core::prompt::{self, FixType};
use devflow_core::stage::Stage;
use devflow_core::state::State;
use devflow_core::{events, workflow};
use std::path::Path;
use tracing::info;
pub(crate) fn transition(
project_root: &Path,
state: &mut State,
to: Stage,
) -> Result<(), CliError> {
let from = state.stage;
let _ = run_checkout_hooks(
project_root,
state,
&hooks::hooks_for_transition(from, to),
to,
);
state.stage = to;
if mode::transition_resets_consecutive_failures(from, to) {
state.consecutive_failures = 0;
}
state.infra_failures = 0;
state.gate_pending = false;
workflow::save_state(state)?;
events::emit(
project_root,
state.phase,
"transition",
serde_json::json!({
"from": from.to_string(),
"to": to.to_string(),
}),
);
launch_stage(state, None, Some(from))
}
pub(crate) fn loop_back_to_code(
project_root: &Path,
state: &mut State,
fix: FixType,
) -> Result<(), CliError> {
let from = state.stage;
let prompt = prepare_loop_back_to_code(project_root, state, fix)?;
launch_stage(state, Some(prompt), Some(from))
}
pub(crate) fn prepare_loop_back_to_code(
project_root: &Path,
state: &mut State,
fix: FixType,
) -> Result<String, CliError> {
let gate_stage = state.stage;
let _ = Gates::cleanup(project_root, state.phase, gate_stage);
state.stage = Stage::Code;
state.gate_pending = false;
workflow::save_state(state)?;
events::emit(
project_root,
state.phase,
"loop_back",
serde_json::json!({
"from": gate_stage.to_string(),
"consecutive_failures": state.consecutive_failures,
}),
);
println!(
"looping back to Code (validate failures: {})",
state.consecutive_failures
);
Ok(prompt::fix_prompt(fix, state.phase))
}
pub(crate) fn finish_workflow(project_root: &Path, state: &mut State) -> Result<(), CliError> {
loop {
if run_checkout_hooks(project_root, state, &hooks::hooks_after_ship(), Stage::Ship) {
break;
}
let _ = Gates::cleanup(project_root, state.phase, Stage::Ship);
let context = format!(
"[finalization failed] phase {} terminal hooks did not complete. Resolve the git/version error, then approve to retry; reject to loop back or abort.",
state.phase
);
match run_gate(project_root, state, Stage::Ship, &context)? {
GateAction::Advance => {
let _ = Gates::cleanup(project_root, state.phase, Stage::Ship);
}
GateAction::LoopBack(_) => {
return loop_back_to_code(project_root, state, FixType::AuditFix);
}
GateAction::Abort(reason) => return abort(project_root, state, &reason),
}
}
let _ = Gates::cleanup(project_root, state.phase, Stage::Validate);
let _ = Gates::cleanup(project_root, state.phase, Stage::Ship);
workflow::clear_state(project_root, state.phase)?;
events::emit(
project_root,
state.phase,
"workflow_finished",
serde_json::Value::Null,
);
println!("phase {} shipped — workflow complete", state.phase);
Ok(())
}
pub(crate) fn run_gate(
project_root: &Path,
state: &mut State,
stage: Stage,
context: &str,
) -> Result<GateAction, CliError> {
state.gate_pending = true;
workflow::save_state(state)?;
Gates::write_gate(project_root, state.phase, stage, context)?;
println!(
"gate written: .devflow/gates/{:02}-{stage}.json — awaiting response",
state.phase
);
let unexpected = !state.mode.should_gate(stage, state.consecutive_failures);
if unexpected {
info!(
"never-silent gate: {stage} failed in {:?} mode — surfacing an unattended gate this mode would not normally fire",
state.mode
);
}
events::emit(
project_root,
state.phase,
"gate_fired",
serde_json::json!({
"stage": stage.to_string(),
"unexpected": unexpected,
"context": context,
}),
);
gates::fire_gate_notify(state.phase, stage, context, unexpected);
events::emit(
project_root,
state.phase,
"notify_fired",
serde_json::json!({ "stage": stage.to_string(), "unexpected": unexpected }),
);
match Gates::poll_response(project_root, state.phase, stage, gate_timeout_secs()) {
Some(response) => {
state.gate_pending = false;
workflow::save_state(state)?;
Gates::ack(project_root, state.phase, stage)?;
let action = GateAction::from_response(&response);
events::emit(
project_root,
state.phase,
"gate_resolved",
serde_json::json!({
"stage": stage.to_string(),
"approved": response.approved,
"action": match &action {
GateAction::Advance => "advance",
GateAction::LoopBack(_) => "loop_back",
GateAction::Abort(_) => "abort",
},
"responded_by": response.responded_by,
}),
);
Ok(action)
}
None => {
events::emit(
project_root,
state.phase,
"gate_timeout",
serde_json::json!({ "stage": stage.to_string() }),
);
Err(CliError::Message(format!(
"gate for stage {stage} timed out awaiting a response"
)))
}
}
}
pub(crate) fn abort(project_root: &Path, state: &State, reason: &str) -> Result<(), CliError> {
println!("workflow aborted for phase {}: {reason}", state.phase);
let _ = Gates::cleanup(project_root, state.phase, state.stage);
let _ = workflow::clear_state(project_root, state.phase);
events::emit(
project_root,
state.phase,
"workflow_aborted",
serde_json::json!({ "reason": truncate_reason(reason) }),
);
Ok(())
}
pub(crate) fn print_dry_run(state: &State) {
println!(
"dry run — phase {} | agent {} | mode {}",
state.phase, state.agent, state.mode
);
println!("\nstage pipeline:");
let mut stage = Some(Stage::Define);
while let Some(s) = stage {
let command = s.gsd_command().replace("{N}", &state.phase.to_string());
let gate = if state.mode.should_gate(s, 0) {
" [GATE]".to_string()
} else if state.mode.should_gate(s, mode::MAX_CONSECUTIVE_FAILURES) {
format!(" [GATE after {} failures]", mode::MAX_CONSECUTIVE_FAILURES)
} else {
String::new()
};
println!(" {s:<9} {command}{gate}");
if let Some(next) = s.next() {
let transition_hooks = hooks::hooks_for_transition(s, next);
if !transition_hooks.is_empty() {
println!(" ↳ hooks: {transition_hooks:?}");
}
}
stage = s.next();
}
println!("\nafter ship: {:?}", hooks::hooks_after_ship());
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pipeline_launch::advance;
use crate::pipeline_outcomes::{
ValidateOutcome, handle_infra_outcome, handle_validate_outcome,
};
use crate::test_support::*;
use devflow_core::agent_result;
use devflow_core::gates::GateResponse;
use devflow_core::mode::Mode;
use devflow_core::state::AgentKind;
#[test]
fn advance_ship_success_runs_finish_workflow() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 21;
let branch = format!("feature/phase-{phase:02}");
let branch_created = std::process::Command::new("git")
.args(["branch", &branch, "develop"])
.current_dir(root)
.status()
.unwrap()
.success();
assert!(branch_created);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
std::fs::write(
agent_result::stdout_path(root, phase),
"DEVFLOW_RESULT: {\"status\":\"success\"}\n",
)
.unwrap();
let response_path = Gates::response_path(root, phase, Stage::Ship);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(
&response_path,
r#"{"approved":true,"note":null,"responded_by":"test"}"#,
)
.unwrap();
advance(root, Some(phase)).unwrap();
let err = workflow::load_state(root, phase).unwrap_err();
assert!(matches!(err, workflow::WorkflowError::MissingState(_)));
assert!(!Gates::gate_path(root, phase, Stage::Ship).exists());
assert!(!Gates::response_path(root, phase, Stage::Ship).exists());
assert!(!Gates::ack_path(root, phase, Stage::Ship).exists());
assert!(!Gates::gate_path(root, phase, Stage::Validate).exists());
}
#[test]
fn terminal_merge_failure_reopens_actionable_gate_and_never_reports_finished() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let git = |args: &[&str]| {
let output = std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap();
assert!(output.status.success(), "git {args:?} failed");
};
git(&["checkout", "-q", "-b", "feature/phase-22"]);
std::fs::write(root.join("conflict.txt"), "feature\n").unwrap();
git(&["add", "conflict.txt"]);
git(&["commit", "-q", "-m", "feature change"]);
git(&["checkout", "-q", "develop"]);
std::fs::write(root.join("conflict.txt"), "develop\n").unwrap();
git(&["add", "conflict.txt"]);
git(&["commit", "-q", "-m", "develop change"]);
let mut state = State::new(22, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
let root_owned = root.to_path_buf();
let handle = std::thread::spawn(move || {
let mut state = workflow::load_state(&root_owned, 22).unwrap();
finish_workflow(&root_owned, &mut state)
});
let gate_path = Gates::gate_path(root, 22, Stage::Ship);
for _ in 0..100 {
if gate_path.exists() {
break;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
assert!(
gate_path.exists(),
"finalization failure must reopen Ship gate"
);
assert!(workflow::load_state(root, 22).unwrap().gate_pending);
Gates::respond(
root,
22,
Stage::Ship,
&GateResponse {
approved: false,
note: Some("abort after merge conflict".into()),
responded_by: Some("test".into()),
},
)
.unwrap();
handle.join().unwrap().unwrap();
assert_ne!(
events::last_event_for_phase(root, 22)
.and_then(|event| event["event"].as_str().map(str::to_owned))
.as_deref(),
Some("workflow_finished")
);
let tags = std::process::Command::new("git")
.arg("tag")
.current_dir(root)
.output()
.unwrap();
assert!(tags.stdout.is_empty());
}
#[test]
fn concurrent_ship_advances_finish_both_phases_independently() {
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 phases = [31u32, 32u32];
for &phase in &phases {
let branch = format!("feature/phase-{phase:02}");
let branch_created = std::process::Command::new("git")
.args(["branch", &branch, "develop"])
.current_dir(root)
.status()
.unwrap()
.success();
assert!(branch_created);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
std::fs::write(
agent_result::stdout_path(root, phase),
"DEVFLOW_RESULT: {\"status\":\"success\"}\n",
)
.unwrap();
let response_path = Gates::response_path(root, phase, Stage::Ship);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(
&response_path,
r#"{"approved":true,"note":null,"responded_by":"test"}"#,
)
.unwrap();
}
let results: Vec<(u32, Result<(), CliError>)> = std::thread::scope(|scope| {
let handles: Vec<_> = phases
.iter()
.map(|&phase| (phase, scope.spawn(move || advance(root, Some(phase)))))
.collect();
handles
.into_iter()
.map(|(phase, handle)| (phase, handle.join().expect("advance thread")))
.collect()
});
unsafe {
match &original_gate_timeout {
Some(value) => std::env::set_var("DEVFLOW_GATE_TIMEOUT_SECS", value),
None => std::env::remove_var("DEVFLOW_GATE_TIMEOUT_SECS"),
}
}
let succeeded = results.iter().filter(|(_, r)| r.is_ok()).count();
assert!(
succeeded == 1 || succeeded == 2,
"at least one phase must finish independently of the other; got {succeeded}/2 successes"
);
for (phase, result) in &results {
match result {
Ok(()) => {
assert!(
matches!(
workflow::load_state(root, *phase),
Err(workflow::WorkflowError::MissingState(_))
),
"phase {phase} must be finished (state cleared)"
);
assert!(!Gates::gate_path(root, *phase, Stage::Ship).exists());
let last = devflow_core::events::last_event_for_phase(root, *phase)
.expect("events recorded for phase");
assert_eq!(
last["event"], "workflow_finished",
"phase {phase}'s own event stream must end in workflow_finished"
);
}
Err(err) => {
assert!(
err.to_string().contains("timed out"),
"phase {phase}'s only non-success outcome must be a bounded gate \
timeout, not some other failure: {err}"
);
let state = workflow::load_state(root, *phase)
.expect("a timed-out gate leaves state intact, not cleared");
assert!(
state.gate_pending,
"phase {phase} must leave an actionable, still-open gate for a human"
);
assert!(
Gates::gate_path(root, *phase, Stage::Ship).exists(),
"phase {phase}'s reopened Ship gate file must remain on disk"
);
}
}
}
}
#[test]
fn abort_cleans_up_gate_files_so_a_later_gate_does_not_reuse_stale_response() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 23;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
state.consecutive_failures = mode::MAX_CONSECUTIVE_FAILURES - 1;
workflow::save_state(&state).unwrap();
let response_path = Gates::response_path(root, phase, Stage::Validate);
std::fs::create_dir_all(response_path.parent().unwrap()).unwrap();
std::fs::write(
&response_path,
r#"{"approved":false,"note":"abort: requirements changed","responded_by":"test"}"#,
)
.unwrap();
handle_validate_outcome(root, &mut state, ValidateOutcome::Failed).unwrap();
assert!(!Gates::gate_path(root, phase, Stage::Validate).exists());
assert!(
!Gates::response_path(root, phase, Stage::Validate).exists(),
"stale response file must not survive an aborted gate"
);
assert!(!Gates::ack_path(root, phase, Stage::Validate).exists());
Gates::write_gate(root, phase, Stage::Validate, "re-fired gate").unwrap();
let started = std::time::Instant::now();
let got = Gates::poll_response(root, phase, Stage::Validate, 1);
assert!(
got.is_none(),
"poll_response must not instantly resolve from a stale response after cleanup"
);
assert!(started.elapsed() >= std::time::Duration::from_secs(1));
}
#[test]
fn transition_resets_infra_failures() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 80;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
state.infra_failures = mode::MAX_INFRA_FAILURES - 1;
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 _ = transition(root, &mut state, Stage::Validate);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
assert_eq!(
state.infra_failures, 0,
"transition() must reset infra_failures in-memory, not just consecutive_failures"
);
let reloaded = workflow::load_state(root, phase).unwrap();
assert_eq!(
reloaded.infra_failures, 0,
"transition() must persist the infra_failures reset to state.json"
);
let response_path = Gates::response_path(root, phase, Stage::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();
handle_infra_outcome(root, &mut state, Stage::Validate, Some("killed".into())).unwrap();
assert_eq!(state.infra_failures, 1);
}
#[test]
fn repeated_code_to_validate_transition_is_idempotent_on_the_counter() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 83;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
state.consecutive_failures = 2;
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 _ = transition(root, &mut state, Stage::Validate);
state.stage = Stage::Code;
let _ = transition(root, &mut state, Stage::Validate);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
assert_eq!(state.consecutive_failures, 2);
}
#[test]
fn consecutive_failures_are_independent_across_phases() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let mut state_a = State::new(84, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state_a.stage = Stage::Code;
state_a.consecutive_failures = 1;
workflow::save_state(&state_a).unwrap();
let mut state_b = State::new(85, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state_b.stage = Stage::Code;
state_b.consecutive_failures = 2;
workflow::save_state(&state_b).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 _ = transition(root, &mut state_a, Stage::Validate);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
let reloaded_a = workflow::load_state(root, 84).unwrap();
let reloaded_b = workflow::load_state(root, 85).unwrap();
assert_eq!(
reloaded_a.consecutive_failures, 1,
"the Code->Validate hop must not reset consecutive_failures"
);
assert_eq!(
reloaded_b.consecutive_failures, 2,
"an untouched sibling phase's counter must be unaffected"
);
}
}