use crate::CliError;
use crate::config_parse::{foreground_gate_timeout_secs, gate_timeout_secs};
use crate::pipeline_launch::launch_stage;
use crate::pipeline_outcomes::{run_checkout_hooks, truncate_reason};
use devflow_core::gates::{self, GateAction, GateError, GateResponse, Gates};
use devflow_core::hooks;
use devflow_core::mode;
use devflow_core::phase_id::PhaseId;
use devflow_core::prompt::{self, FixType};
use devflow_core::stage::Stage;
use devflow_core::state::State;
use devflow_core::{events, lock, registry, 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;
if state.stop_until == Some(from) {
state.stopped = true;
state.stop_reason = Some(format!("stopped after {from} completed (--until {from})"));
state.monitor_pid = None;
state.gate_pending = false;
workflow::save_state(state)?;
events::emit(
project_root,
state.phase,
"workflow_finished",
serde_json::json!({
"reason": "stopped_at",
"stage": from.to_string(),
}),
);
return Ok(());
}
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))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum LoopBackReason {
ValidateFailure,
ValidateFailureNoBaseline,
GateResponse,
}
impl LoopBackReason {
pub(crate) fn as_str(self) -> &'static str {
match self {
LoopBackReason::ValidateFailure => "validate_failure",
LoopBackReason::ValidateFailureNoBaseline => "validate_failure_no_commit_baseline",
LoopBackReason::GateResponse => "gate_response",
}
}
}
pub(crate) fn loop_back_to_code(
project_root: &Path,
state: &mut State,
fix: FixType,
reason: LoopBackReason,
) -> Result<(), CliError> {
let from = state.stage;
let prompt = prepare_loop_back_to_code(project_root, state, fix, reason)?;
launch_stage(state, Some(prompt), Some(from))
}
pub(crate) fn prepare_loop_back_to_code(
project_root: &Path,
state: &mut State,
fix: FixType,
reason: LoopBackReason,
) -> 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,
"phase_validate_failures": state.phase_validate_failures,
"reason": reason.as_str(),
"fix": format!("{fix:?}"),
}),
);
println!(
"looping back to Code ({} validate failure(s) this phase, {} in the current streak)",
state.phase_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> {
finish_workflow_with_gate_timeout(project_root, state, gate_timeout_secs())
}
pub(crate) fn finish_workflow_with_gate_timeout(
project_root: &Path,
state: &mut State,
gate_timeout_secs: u64,
) -> 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_with_timeout(
project_root,
state,
Stage::Ship,
&context,
gate_timeout_secs,
None,
)? {
GateAction::Advance => {
let _ = Gates::cleanup(project_root, state.phase, Stage::Ship);
}
GateAction::LoopBack(_) => {
return loop_back_to_code(
project_root,
state,
FixType::AuditFix,
LoopBackReason::GateResponse,
);
}
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)?;
registry::deregister(project_root, state.phase);
events::emit(
project_root,
state.phase,
"workflow_shipped",
serde_json::json!({
"stage": Stage::Ship.to_string(),
}),
);
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> {
run_gate_with_timeout(
project_root,
state,
stage,
context,
gate_timeout_secs(),
None,
)
}
pub(crate) fn run_gate_with_timeout(
project_root: &Path,
state: &mut State,
stage: Stage,
context: &str,
timeout_secs: u64,
auto_response: Option<&GateResponse>,
) -> 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/{}-{stage}.json — awaiting response",
state.phase.padded()
);
let unexpected = !state.mode.should_gate(
stage,
state.consecutive_failures,
state.phase_validate_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 }),
);
if let Some(response) = auto_response {
match Gates::respond(project_root, state.phase, stage, response) {
Ok(_) => {}
Err(GateError::AlreadyResponded { .. }) => {}
Err(err) => return Err(err.into()),
}
}
match Gates::poll_response(project_root, state.phase, stage, 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);
registry::deregister(project_root, state.phase);
events::emit(
project_root,
state.phase,
"workflow_aborted",
serde_json::json!({ "reason": truncate_reason(reason) }),
);
Ok(())
}
pub(crate) fn ship_override(
project_root: &Path,
phase: PhaseId,
force: bool,
) -> Result<(), CliError> {
let _lock = match lock::acquire(project_root, phase) {
Ok(guard) => guard,
Err(lock::LockError::Contended { pid, .. }) => {
return Err(CliError::Message(format!(
"phase {phase}: another devflow process (pid {pid}) holds the per-phase lock — \
refusing to race its poll of the Ship gate response"
)));
}
Err(err) => return Err(CliError::Message(format!("lock error: {err}"))),
};
let mut state = workflow::load_state(project_root, phase)?;
if state.stage != Stage::Ship {
return Err(CliError::Message(format!(
"phase {phase} is at stage {} — `devflow ship` requires state.stage == Stage::Ship; \
resolve stage {} first (--force does not skip stages)",
state.stage, state.stage
)));
}
if !Gates::gate_path(project_root, phase, Stage::Ship).exists()
|| !Gates::response_path(project_root, phase, Stage::Ship).exists()
{
return Err(CliError::Message(format!(
"phase {phase}: no Ship gate response written yet — nothing to ship (a dead monitor \
never wrote or received one; wait for `devflow gate approve` or resolve the pipeline first)"
)));
}
if Gates::ack_path(project_root, phase, Stage::Ship).exists() {
return Err(CliError::Message(format!(
"phase {phase}: the Ship gate response was already consumed (an ack file is present) \
— the phase may be mid-finalization from a monitor that died partway through; run \
`devflow doctor` to inspect it rather than re-running terminal hooks"
)));
}
let response_path = Gates::response_path(project_root, phase, Stage::Ship);
let contents = std::fs::read_to_string(&response_path).map_err(|err| {
CliError::Message(format!("could not read the Ship gate response: {err}"))
})?;
let response: GateResponse = serde_json::from_str(&contents).map_err(|err| {
CliError::Message(format!("could not parse the Ship gate response: {err}"))
})?;
println!(
"phase {phase}: manual ship override (--force={force}) — driving the already-written \
Ship response through the same terminal path the live monitor would have used"
);
match GateAction::from_response(&response) {
GateAction::Advance => {
let timeout = foreground_gate_timeout_secs();
println!(
"phase {phase}: if terminal-hook finalization fails, this foreground command \
will wait up to {timeout}s for the reopened Ship gate before failing (vs. the \
multi-day background default) — set DEVFLOW_FOREGROUND_GATE_TIMEOUT_SECS to \
change this"
);
finish_workflow_with_gate_timeout(project_root, &mut state, timeout)
}
GateAction::LoopBack(_) => {
println!(
"phase {phase}: Ship response loops back to Code — launching a new, detached \
monitor agent to drive the retry"
);
loop_back_to_code(
project_root,
&mut state,
FixType::AuditFix,
LoopBackReason::GateResponse,
)
}
GateAction::Abort(reason) => abort(project_root, &state, &reason),
}
}
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 unconditional = state.mode.should_gate(s, 0, 0);
let gate = if unconditional {
" [GATE]".to_string()
} else if state.mode.should_gate(s, mode::MAX_CONSECUTIVE_FAILURES, 0) {
format!(
" [GATE after {} consecutive failures]",
mode::MAX_CONSECUTIVE_FAILURES
)
} else {
String::new()
};
let phase_gate = if !unconditional
&& state
.mode
.should_gate(s, 0, mode::MAX_PHASE_VALIDATE_FAILURES)
{
format!(
" [GATE at {} validate failures for this phase]",
mode::MAX_PHASE_VALIDATE_FAILURES
)
} else {
String::new()
};
let stop_marker = if state.stop_until == Some(s) {
" [STOPS HERE — --until]"
} else {
""
};
println!(" {s:<9} {command}{gate}{phase_gate}{stop_marker}");
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();
}
if let Some(until) = state.stop_until {
println!("\nnote: --until {until} — this run will halt after {until} completes");
}
println!(
"\nship gate: {}",
if state.yes_ship {
"pre-authorized"
} else {
"not pre-authorized"
}
);
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 = PhaseId::new(21);
let branch = format!("feature/phase-{padded}", padded = phase.padded());
let branch_created = devflow_core::test_support::git_command(root)
.args(["branch", &branch, "develop"])
.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 advance_ship_success_emits_workflow_shipped_and_ship_evidence_reports_shipped() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = PhaseId::new(23);
let branch = format!("feature/phase-{padded}", padded = phase.padded());
let branch_created = devflow_core::test_support::git_command(root)
.args(["branch", &branch, "develop"])
.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();
assert!(
devflow_core::events::has_event_for_phase(root, phase, "workflow_shipped"),
"a real Ship finalization must emit the terminal-only workflow_shipped event"
);
let evidence = devflow_core::ship_evidence::collect(root, phase);
assert!(
evidence.shipped,
"ship_evidence must read the just-emitted workflow_shipped event as shipped"
);
let last = devflow_core::events::last_event_for_phase(root, phase)
.expect("events recorded for phase");
assert_eq!(last["event"], "workflow_finished");
}
#[test]
fn until_stop_never_emits_workflow_shipped_and_ship_evidence_reports_not_shipped() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = PhaseId::new(24);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Plan;
state.stop_until = Some(Stage::Plan);
workflow::save_state(&state).unwrap();
transition(root, &mut state, Stage::Code).unwrap();
assert!(
!devflow_core::events::has_event_for_phase(root, phase, "workflow_shipped"),
"the --until clean-stop branch must never emit workflow_shipped"
);
let evidence = devflow_core::ship_evidence::collect(root, phase);
assert!(
!evidence.shipped,
"a phase that only stopped after one stage must not read as shipped"
);
assert!(evidence.workflow_finished_seen);
assert_eq!(evidence.finished_reason.as_deref(), Some("stopped_at"));
assert!(devflow_core::ship_evidence::is_stopped_at(&evidence));
}
#[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 = devflow_core::test_support::git_command(root)
.args(args)
.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(
PhaseId::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, PhaseId::new(22)).unwrap();
finish_workflow(&root_owned, &mut state)
});
let gate_path = Gates::gate_path(root, PhaseId::new(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, PhaseId::new(22))
.unwrap()
.gate_pending
);
Gates::respond(
root,
PhaseId::new(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, PhaseId::new(22))
.and_then(|event| event["event"].as_str().map(str::to_owned))
.as_deref(),
Some("workflow_finished")
);
let tags = devflow_core::test_support::git_command(root)
.arg("tag")
.output()
.unwrap();
assert!(tags.stdout.is_empty());
}
#[test]
fn finalization_retry_gate_never_auto_approves_even_with_yes_ship_set() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let git_cmd = |args: &[&str]| {
assert!(
devflow_core::test_support::git_command(root)
.args(args)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git_cmd(&["checkout", "--orphan", "unreachable-release"]);
git_cmd(&["commit", "--allow-empty", "-q", "-m", "chore: orphan"]);
git_cmd(&["tag", "v9.9.9"]);
git_cmd(&["checkout", "develop"]);
assert!(
devflow_core::version::compute_version(root).is_err(),
"the orphan tag must make compute_version refuse (D-10) before \
VersionBump ever gets to its own git.tag(&tag) call"
);
let phase = PhaseId::new(60);
let branch = format!("feature/phase-{padded}", padded = phase.padded());
let branch_created = devflow_core::test_support::git_command(root)
.args(["branch", &branch, "develop"])
.status()
.unwrap()
.success();
assert!(branch_created);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
state.yes_ship = true;
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, phase).unwrap();
finish_workflow(&root_owned, &mut state)
});
let gate_path = Gates::gate_path(root, phase, Stage::Ship);
for _ in 0..150 {
if gate_path.exists() {
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
assert!(
gate_path.exists(),
"a terminal-hook failure must reopen the Ship gate even with yes_ship set"
);
std::thread::sleep(std::time::Duration::from_millis(50));
let response_path = Gates::response_path(root, phase, Stage::Ship);
assert!(
!response_path.exists(),
"yes_ship must NEVER auto-approve the reopened finalization-retry gate \
(T-23-91) — a failing finalization must wait for a human"
);
assert!(
workflow::load_state(root, phase).unwrap().gate_pending,
"the reopened gate must be recorded as pending, awaiting a human"
);
assert!(
devflow_core::events::has_event_for_phase(root, phase, "merge_result"),
"Merge must have succeeded before VersionBump failed — the reopened gate must be \
exercised in the state it actually occurs in"
);
Gates::respond(
root,
phase,
Stage::Ship,
&GateResponse {
approved: false,
note: Some("abort: test cleanup".into()),
responded_by: Some("test".into()),
},
)
.unwrap();
handle.join().unwrap().unwrap();
}
#[test]
fn concurrent_ship_advances_finish_both_phases_independently() {
let _guard = env_lock();
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 = [PhaseId::new(31), PhaseId::new(32)];
for &phase in &phases {
let branch = format!("feature/phase-{padded}", padded = phase.padded());
let branch_created = devflow_core::test_support::git_command(root)
.args(["branch", &branch, "develop"])
.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<(PhaseId, 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 = PhaseId::new(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;
state.last_validate_failure_commit_count = Some(0);
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();
{
let _guard = env_lock();
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());
}
handle_validate_outcome(root, &mut state, ValidateOutcome::Failed).unwrap();
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
}
assert!(
state.consecutive_failures >= mode::MAX_CONSECUTIVE_FAILURES,
"gate threshold must have been reached (got {}) — a reset means the gate never fired",
state.consecutive_failures
);
assert_eq!(
state.stage,
Stage::Validate,
"gate must have fired and aborted — Stage::Code means it silently looped back and tried to launch an agent"
);
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_lock();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(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_lock();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(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 ship_override_advances_via_written_response() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = PhaseId::new(90);
let branch = format!("feature/phase-{padded}", padded = phase.padded());
let branch_created = devflow_core::test_support::git_command(root)
.args(["branch", &branch, "develop"])
.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();
Gates::write_gate(root, phase, Stage::Ship, "Ship complete — approve merge?").unwrap();
Gates::respond(
root,
phase,
Stage::Ship,
&GateResponse {
approved: true,
note: None,
responded_by: Some("test".into()),
},
)
.unwrap();
ship_override(root, phase, false).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());
let last = devflow_core::events::last_event_for_phase(root, phase)
.expect("events recorded for phase");
assert_eq!(last["event"], "workflow_finished");
}
#[test]
fn ship_override_bounds_foreground_wait_on_terminal_hook_failure() {
let _guard = env_lock();
let original_foreground_timeout = std::env::var_os("DEVFLOW_FOREGROUND_GATE_TIMEOUT_SECS");
unsafe {
std::env::set_var("DEVFLOW_FOREGROUND_GATE_TIMEOUT_SECS", "2");
}
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let git = |args: &[&str]| {
let output = devflow_core::test_support::git_command(root)
.args(args)
.output()
.unwrap();
assert!(output.status.success(), "git {args:?} failed");
};
let phase = PhaseId::new(96);
let branch = format!("feature/phase-{padded}", padded = phase.padded());
git(&["checkout", "-q", "-b", &branch]);
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(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
Gates::write_gate(root, phase, Stage::Ship, "Ship complete — approve merge?").unwrap();
Gates::respond(
root,
phase,
Stage::Ship,
&GateResponse {
approved: true,
note: None,
responded_by: Some("test".into()),
},
)
.unwrap();
let started = std::time::Instant::now();
let result = ship_override(root, phase, false);
let elapsed = started.elapsed();
unsafe {
match &original_foreground_timeout {
Some(value) => {
std::env::set_var("DEVFLOW_FOREGROUND_GATE_TIMEOUT_SECS", value);
}
None => std::env::remove_var("DEVFLOW_FOREGROUND_GATE_TIMEOUT_SECS"),
}
}
assert!(
result.is_err(),
"an unresolved reopened Ship gate must fail closed, not silently advance"
);
assert!(
elapsed < std::time::Duration::from_secs(30),
"the foreground wait must be bounded by DEVFLOW_FOREGROUND_GATE_TIMEOUT_SECS, \
not gate_timeout_secs' multi-day default — took {elapsed:?}"
);
assert!(
Gates::gate_path(root, phase, Stage::Ship).exists(),
"the merge failure must reopen an actionable Ship gate for a human, not silently \
drop the phase"
);
assert!(
workflow::load_state(root, phase).is_ok(),
"state must NOT be cleared — finish_workflow_with_gate_timeout must fail before \
reaching workflow::clear_state"
);
}
#[test]
fn ship_override_abort_routes_through_abort() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(95);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
Gates::write_gate(root, phase, Stage::Ship, "Ship complete — approve merge?").unwrap();
Gates::respond(
root,
phase,
Stage::Ship,
&GateResponse {
approved: false,
note: Some("abort: found a blocking regression during manual review".into()),
responded_by: Some("test".into()),
},
)
.unwrap();
ship_override(root, phase, false).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());
let last = devflow_core::events::last_event_for_phase(root, phase)
.expect("events recorded for phase");
assert_eq!(last["event"], "workflow_aborted");
}
#[test]
fn ship_override_refuses_when_not_at_ship_stage() {
for stage in [Stage::Define, Stage::Plan, Stage::Code, Stage::Validate] {
for force in [true, false] {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(91);
let mut state =
State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = stage;
workflow::save_state(&state).unwrap();
let err = ship_override(root, phase, force).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains(&stage.to_string()),
"error for stage {stage} (force={force}) must name the stage: {msg}"
);
let reloaded = workflow::load_state(root, phase)
.expect("state must survive a stage-mismatch refusal, not be cleared");
assert_eq!(reloaded.stage, stage);
}
}
}
#[test]
fn ship_override_refuses_when_no_response_written() {
for force in [true, false] {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(92);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
let err = ship_override(root, phase, force).unwrap_err();
assert!(
err.to_string().contains("no Ship gate response"),
"force={force}: {err}"
);
Gates::write_gate(root, phase, Stage::Ship, "ctx").unwrap();
let err = ship_override(root, phase, force).unwrap_err();
assert!(
err.to_string().contains("no Ship gate response"),
"force={force}: {err}"
);
let _ = Gates::cleanup(root, phase, Stage::Ship);
}
}
#[test]
fn ship_override_refuses_when_lock_contended() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(93);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
let _held = lock::acquire(root, phase).expect("hold the per-phase lock");
let this_pid = std::process::id().to_string();
for force in [true, false] {
let err = ship_override(root, phase, force).unwrap_err();
let msg = err.to_string();
assert!(msg.contains(&this_pid), "force={force}: {msg}");
}
}
#[test]
fn ship_override_refuses_when_response_already_acked() {
for force in [true, false] {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(94);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
Gates::write_gate(root, phase, Stage::Ship, "ctx").unwrap();
Gates::respond(
root,
phase,
Stage::Ship,
&GateResponse {
approved: true,
note: None,
responded_by: Some("test".into()),
},
)
.unwrap();
Gates::ack(root, phase, Stage::Ship).unwrap();
let err = ship_override(root, phase, force).unwrap_err();
assert!(
err.to_string().contains("devflow doctor"),
"force={force}: {err}"
);
let reloaded = workflow::load_state(root, phase)
.expect("an already-acked refusal must never clear state");
assert_eq!(reloaded.stage, Stage::Ship);
}
}
#[test]
fn consecutive_failures_are_independent_across_phases() {
let _guard = env_lock();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let mut state_a = State::new(
PhaseId::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(
PhaseId::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, PhaseId::new(84)).unwrap();
let reloaded_b = workflow::load_state(root, PhaseId::new(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"
);
}
}