use clap::{Parser, Subcommand};
use devflow_core::agent;
use devflow_core::config::{DEVELOP, FEATURE_PREFIX, GitFlowConfig, capture_retention};
use devflow_core::gates::{self, GateAction, GateResponse, Gates, OpenGate};
use devflow_core::git::GitFlow;
use devflow_core::hooks::{self, HookContext};
use devflow_core::mode::{self, Mode};
use devflow_core::prompt::{self, FixType};
use devflow_core::stage::Stage;
use devflow_core::state::{AgentKind, State};
use devflow_core::{
agent_result, agents, events, history, lock, monitor, outcome_policy, recover, worktree,
};
use devflow_core::{
agent_result::{AgentStatus, Verdict},
outcome_policy::Action,
workflow,
};
use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
use tracing::info;
const GATE_ESCALATION_THRESHOLD_SECS: u64 = 30 * 60;
fn parse_gate_timeout(raw: Option<String>) -> u64 {
const SEVEN_DAYS: u64 = 7 * 24 * 60 * 60;
raw.and_then(|s| s.parse().ok()).unwrap_or(SEVEN_DAYS)
}
fn gate_timeout_secs() -> u64 {
parse_gate_timeout(std::env::var("DEVFLOW_GATE_TIMEOUT_SECS").ok())
}
fn parse_checkout_lock_timeout(raw: Option<String>) -> std::time::Duration {
const DEFAULT_SECS: u64 = 120;
std::time::Duration::from_secs(raw.and_then(|s| s.parse().ok()).unwrap_or(DEFAULT_SECS))
}
fn checkout_lock_timeout() -> std::time::Duration {
parse_checkout_lock_timeout(std::env::var("DEVFLOW_CHECKOUT_LOCK_TIMEOUT_SECS").ok())
}
#[derive(Debug, Parser)]
#[command(
name = "devflow",
version,
about = "Agent-agnostic, GSD-native development workflow automation"
)]
struct Cli {
#[command(subcommand)]
command: Command,
}
#[derive(Debug, Subcommand)]
enum Command {
Start {
#[arg(long)]
phase: u32,
#[arg(long, default_value = "claude")]
agent: AgentKind,
#[arg(long)]
mode: Mode,
#[arg(long)]
force: bool,
#[arg(long, hide = true)]
worktree: bool,
#[arg(long)]
no_worktree: bool,
#[arg(long)]
dry_run: bool,
#[arg(default_value = ".")]
project: PathBuf,
},
#[command(hide = true)]
Advance {
#[arg(default_value = ".")]
project: PathBuf,
#[arg(long)]
phase: Option<u32>,
},
Resume {
#[arg(long)]
phase: u32,
#[arg(default_value = ".")]
project: PathBuf,
},
Gate {
#[command(subcommand)]
action: GateCmd,
},
Logs {
#[arg(long)]
phase: Option<u32>,
#[arg(long, short = 'f')]
follow: bool,
#[arg(long)]
stderr: bool,
#[arg(default_value = ".")]
project: PathBuf,
},
History {
phase: Option<u32>,
#[arg(default_value = ".")]
project: PathBuf,
},
Parallel {
#[arg(long)]
phases: String,
#[arg(long)]
agents: Option<String>,
#[arg(long, default_value = "auto")]
mode: Mode,
#[arg(long)]
force: bool,
#[arg(default_value = ".")]
project: PathBuf,
},
Sequentagent {
#[arg(long)]
phase: u32,
#[arg(long)]
agents: String,
#[arg(long)]
force: bool,
#[arg(default_value = ".")]
project: PathBuf,
},
Reference {
#[arg(long)]
branch: Option<String>,
#[arg(long)]
refresh: bool,
#[arg(default_value = ".")]
project: PathBuf,
},
Cleanup {
#[arg(default_value = ".")]
project: PathBuf,
#[arg(long)]
force: bool,
},
Status {
#[arg(default_value = ".")]
project: PathBuf,
},
List {
#[arg(default_value = ".")]
project: PathBuf,
},
Recover {
#[arg(default_value = ".")]
project: PathBuf,
#[arg(long)]
clean: bool,
#[arg(long)]
phase: Option<u32>,
},
Test {
#[arg(default_value = ".")]
project: PathBuf,
},
Doctor {
#[arg(long)]
json: bool,
#[arg(default_value = ".")]
project: PathBuf,
},
}
#[derive(Debug, Subcommand)]
enum GateCmd {
List {
#[arg(default_value = ".")]
project: PathBuf,
},
Approve {
phase: u32,
#[arg(value_name = "STAGE_OR_PROJECT")]
stage: Option<String>,
#[arg(value_name = "PROJECT")]
legacy_project: Option<PathBuf>,
#[arg(long = "stage")]
stage_option: Option<Stage>,
#[arg(long)]
note: Option<String>,
#[arg(long, default_value = ".")]
project: PathBuf,
},
Reject {
phase: u32,
#[arg(value_name = "STAGE_OR_PROJECT")]
stage: Option<String>,
#[arg(value_name = "PROJECT")]
legacy_project: Option<PathBuf>,
#[arg(long = "stage")]
stage_option: Option<Stage>,
#[arg(long)]
note: String,
#[arg(long, default_value = ".")]
project: PathBuf,
},
}
#[derive(Debug, thiserror::Error)]
enum CliError {
#[error(transparent)]
Workflow(#[from] devflow_core::workflow::WorkflowError),
#[error(transparent)]
Recover(#[from] devflow_core::recover::RecoverError),
#[error(transparent)]
Git(#[from] devflow_core::git::GitError),
#[error(transparent)]
Worktree(#[from] devflow_core::worktree::WorktreeError),
#[error(transparent)]
Gate(#[from] devflow_core::gates::GateError),
#[error(transparent)]
Ship(#[from] devflow_core::ship::ShipError),
#[error("{0}")]
Message(String),
}
fn main() {
match std::env::var("DEVFLOW_LOG_FORMAT").as_deref() {
Ok("json") => {
let filter = tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info"));
tracing_subscriber::fmt()
.json()
.with_writer(std::io::stderr)
.with_env_filter(filter)
.init();
}
_ => {
let filter = tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info"));
tracing_subscriber::fmt()
.with_writer(std::io::stderr)
.with_env_filter(filter)
.init();
}
}
if let Err(err) = run() {
eprintln!("error: {err}");
std::process::exit(1);
}
}
fn run() -> Result<(), CliError> {
let cli = Cli::parse();
match cli.command {
Command::Start {
phase,
agent,
mode,
force,
worktree: _worktree,
no_worktree,
dry_run,
project,
} => {
let worktree = !no_worktree;
start(
&project_root(project)?,
phase,
agent,
mode,
force,
worktree,
dry_run,
)
}
Command::Advance { project, phase } => advance(&project_root(project)?, phase),
Command::Resume { phase, project } => resume(&project_root(project)?, phase),
Command::Gate { action } => match action {
GateCmd::List { project } => gate_list(&project_root(project)?),
GateCmd::Approve {
phase,
stage,
legacy_project,
stage_option,
note,
project,
} => {
let (stage, project) =
resolve_gate_target(stage, legacy_project, stage_option, project)?;
gate_respond(&project_root(project)?, phase, stage, true, note)
}
GateCmd::Reject {
phase,
stage,
legacy_project,
stage_option,
note,
project,
} => {
let (stage, project) =
resolve_gate_target(stage, legacy_project, stage_option, project)?;
gate_respond(&project_root(project)?, phase, stage, false, Some(note))
}
},
Command::Logs {
phase,
follow,
stderr,
project,
} => logs(&project_root(project)?, phase, follow, stderr),
Command::History { phase, project } => history_cmd(&project_root(project)?, phase),
Command::Parallel {
phases,
agents,
mode,
force,
project,
} => parallel(
&project_root(project)?,
&phases,
agents.as_deref(),
mode,
force,
),
Command::Sequentagent {
phase,
agents,
force,
project,
} => sequentagent(&project_root(project)?, phase, &agents, force),
Command::Reference {
branch,
refresh,
project,
} => reference(&project_root(project)?, branch, refresh),
Command::Cleanup { project, force } => cleanup(&project_root(project)?, force),
Command::Status { project } => status(&project_root(project)?),
Command::List { project } => list(&project_root(project)?),
Command::Recover {
project,
clean,
phase,
} => recover_cmd(&project_root(project)?, clean, phase),
Command::Test { project } => test_cmd(&project_root(project)?),
Command::Doctor { json, project } => doctor(&project_root(project)?, json),
}
}
fn resolve_gate_target(
positional: Option<String>,
legacy_project: Option<PathBuf>,
stage_option: Option<Stage>,
project: PathBuf,
) -> Result<(Option<Stage>, PathBuf), CliError> {
let Some(positional) = positional else {
return Ok((stage_option, project));
};
if let Ok(positional_stage) = positional.parse::<Stage>() {
if let Some(flagged_stage) = stage_option
&& flagged_stage != positional_stage
{
return Err(CliError::Message(format!(
"conflicting stages: positional {positional_stage} and --stage {flagged_stage}"
)));
}
let target = legacy_project.unwrap_or(project);
return Ok((Some(stage_option.unwrap_or(positional_stage)), target));
}
if legacy_project.is_some() {
return Err(CliError::Message(format!(
"unsupported stage `{positional}`; expected define, plan, code, validate, or ship"
)));
}
if project.as_path() != Path::new(".") {
return Err(CliError::Message(
"project was supplied both positionally and with --project".into(),
));
}
Ok((stage_option, PathBuf::from(positional)))
}
#[allow(clippy::too_many_arguments)]
fn phase_artifact_on_develop(project_root: &Path, phase: u32, suffix: &str) -> bool {
let prefix = format!(".planning/phases/{phase:02}-");
let output = std::process::Command::new("git")
.args([
"ls-tree",
"-r",
"--name-only",
"develop",
"--",
".planning/phases/",
])
.current_dir(project_root)
.output();
let Ok(out) = output else { return true };
if !out.status.success() {
return true;
}
String::from_utf8_lossy(&out.stdout).lines().any(|path| {
path.strip_prefix(&prefix)
.is_some_and(|rest| rest.contains('/') && rest.ends_with(suffix))
})
}
fn start(
project_root: &Path,
phase: u32,
agent: AgentKind,
mode: Mode,
force: bool,
worktree: bool,
dry_run: bool,
) -> Result<(), CliError> {
let mut state = State::new(phase, agent, mode, project_root.to_path_buf());
if dry_run {
print_dry_run(&state);
return Ok(());
}
ensure_agent_binary(agent_program(agent))?;
if agent == AgentKind::Codex {
if !phase_artifact_on_develop(project_root, phase, "-CONTEXT.md") {
return Err(CliError::Message(format!(
"phase {phase} has no CONTEXT.md on develop, and codex cannot run an \
interactive discussion headless. Run /gsd-discuss-phase {phase} \
interactively first (any agent), or use --agent claude."
)));
}
if !phase_artifact_on_develop(project_root, phase, "-PLAN.md") {
println!(
"warning: phase {phase} has no PLAN.md on develop — headless codex \
planning is untested and may need input; pre-writing plans is safer"
);
}
}
if !worktree && let Ok((_ahead, behind)) = GitFlow::new(project_root).divergence_from_develop()
{
if behind > 50 {
return Err(CliError::Message(format!(
"develop is {behind} commits ahead — your branch is too far behind. \
Rebase onto develop first, or use --force to override."
)));
}
if behind > 10 {
println!("warning: develop is {behind} commits ahead — consider rebasing first");
}
}
if worktree {
let wt = ensure_phase_worktree(project_root, phase, force)?;
println!(
"created worktree: {} (branch {FEATURE_PREFIX}phase-{phase:02})",
wt.display(),
);
state.worktree_path = Some(wt);
} else {
let git = GitFlow::new(project_root);
let result = if force {
git.feature_start_force(phase)
} else {
git.feature_start(phase)
};
match result {
Ok(branch) => println!("created feature branch: {branch}"),
Err(err) => {
if !force {
return Err(CliError::Message(format!(
"{err}\nUse --force to overwrite the existing branch."
)));
}
return Err(err.into());
}
}
}
workflow::save_state(&state)?;
events::emit(
project_root,
phase,
"workflow_started",
workflow_started_payload(&state),
);
if let Err(err) = launch_stage(&mut state, None, None) {
if let Err(clear_err) = workflow::clear_state(project_root, phase) {
eprintln!("warning: could not clear state after failed launch: {clear_err}");
}
return Err(err);
}
println!(
"started phase {} in {mode} mode at {} — monitor will auto-advance",
state.phase, state.started_at
);
println!(" watch live: devflow logs -f --phase {phase}");
Ok(())
}
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)
}
fn agent_program(agent: AgentKind) -> &'static str {
agents::adapter_for(agent).exec_command(0, "", &[]).0
}
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)
}
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)
}
fn workflow_started_payload(state: &State) -> serde_json::Value {
serde_json::json!({
"agent": state.agent.to_string(),
"mode": state.mode.to_string(),
"worktree": state.worktree_path.as_ref().map(|p| p.display().to_string()),
"version": env!("CARGO_PKG_VERSION"),
"commit": env!("DEVFLOW_BUILD_COMMIT"),
"dirty": env!("DEVFLOW_BUILD_DIRTY"),
"exe_path": std::env::current_exe()
.ok()
.map(|p| p.display().to_string()),
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Staleness {
Fresh,
Stale,
Ahead,
Indeterminate,
}
fn embedded_commit_is_stale(execution_root: &Path, embedded_commit: &str) -> Staleness {
if embedded_commit.is_empty() {
return Staleness::Indeterminate;
}
let output = std::process::Command::new("git")
.args(["merge-base", "--is-ancestor", embedded_commit, "HEAD"])
.current_dir(execution_root)
.output();
match output.map(|o| o.status.code()) {
Ok(Some(0)) => match run_git_stdout(execution_root, &["rev-parse", "HEAD"]) {
Some(head) if head.trim() == embedded_commit.trim() => Staleness::Fresh,
Some(_) => Staleness::Stale,
None => Staleness::Indeterminate,
},
Ok(Some(1)) => {
let reverse = std::process::Command::new("git")
.args(["merge-base", "--is-ancestor", "HEAD", embedded_commit])
.current_dir(execution_root)
.output();
match reverse.map(|o| o.status.code()) {
Ok(Some(0)) => Staleness::Ahead,
Ok(Some(1)) => Staleness::Stale,
_ => Staleness::Indeterminate,
}
}
_ => Staleness::Indeterminate,
}
}
fn run_git_stdout(project_root: &Path, args: &[&str]) -> Option<String> {
let output = std::process::Command::new("git")
.args(args)
.current_dir(project_root)
.output()
.ok()?;
output
.status
.success()
.then(|| String::from_utf8_lossy(&output.stdout).to_string())
}
fn tree_has_modified_build_inputs(execution_root: &Path) -> Option<bool> {
let status = run_git_stdout(execution_root, &["status", "--porcelain"])?;
if status.trim().is_empty() {
return Some(false);
}
Some(
status
.lines()
.any(|line| porcelain_tracked_path(line).is_some_and(affects_compiled_binary)),
)
}
fn porcelain_tracked_path(line: &str) -> Option<&str> {
if line.len() < 4 || line.starts_with("??") {
return None;
}
let path = &line[3..];
let path = path.rsplit(" -> ").next().unwrap_or(path);
Some(path.trim_matches('"'))
}
fn affects_compiled_binary(rel_path: &str) -> bool {
const BUILD_AFFECTING_FILES: [&str; 4] = [
"Cargo.toml",
"Cargo.lock",
"build.rs",
"rust-toolchain.toml",
];
rel_path.ends_with(".rs")
|| BUILD_AFFECTING_FILES
.iter()
.any(|name| rel_path == *name || rel_path.ends_with(&format!("/{name}")))
}
fn combined_staleness(
execution_root: &Path,
embedded_commit: &str,
build_dirty: bool,
) -> Staleness {
let ancestry = embedded_commit_is_stale(execution_root, embedded_commit);
if ancestry == Staleness::Stale {
return Staleness::Stale;
}
match tree_has_modified_build_inputs(execution_root) {
Some(true) if build_dirty => Staleness::Indeterminate,
Some(true) => Staleness::Stale,
_ => ancestry,
}
}
fn is_self_dogfood_workspace(project_root: &Path) -> bool {
let Ok(contents) = std::fs::read_to_string(project_root.join("Cargo.toml")) else {
return false;
};
let Some(members_start) = contents.match_indices("members").find_map(|(idx, _)| {
let preceded_by_ident = contents[..idx]
.chars()
.next_back()
.is_some_and(|ch| ch.is_alphanumeric() || ch == '_' || ch == '-');
(!preceded_by_ident).then_some(idx)
}) else {
return false;
};
let rest = &contents[members_start..];
let Some(open_rel) = rest.find('[') else {
return false;
};
let after_open = &rest[open_rel + 1..];
let Some(close_rel) = after_open.find(']') else {
return false;
};
let members = &after_open[..close_rel];
let has_member = |wanted: &str| {
members
.split(',')
.any(|entry| entry.trim().trim_matches(['"', '\'']).trim() == wanted)
};
has_member("crates/devflow-core") && has_member("crates/devflow-cli")
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum StalenessOutcome {
Block,
Warn,
Ok,
}
fn staleness_outcome(is_self_dogfood: bool, staleness: Staleness) -> StalenessOutcome {
match (is_self_dogfood, staleness) {
(true, Staleness::Stale) => StalenessOutcome::Block,
(false, Staleness::Stale) => StalenessOutcome::Warn,
(_, Staleness::Ahead) => StalenessOutcome::Warn,
(_, Staleness::Indeterminate) => StalenessOutcome::Warn,
(_, Staleness::Fresh) => StalenessOutcome::Ok,
}
}
fn enforce_build_staleness(
project_root: &Path,
state: &State,
embedded_commit: &str,
build_dirty: bool,
) -> Result<(), CliError> {
let execution_root = state.worktree_path.as_deref().unwrap_or(project_root);
let staleness = combined_staleness(execution_root, embedded_commit, build_dirty);
let self_dogfood = is_self_dogfood_workspace(project_root);
match staleness_outcome(self_dogfood, staleness) {
StalenessOutcome::Block => {
let message = format!(
"self-dogfood stale build blocked for stage {}: this devflow binary's \
embedded commit is not an ancestor of {}'s current HEAD (or its tracked \
source is newer than the build) — rebuild devflow before driving its own \
workspace (D-18; the Phase 16 false-evidence incident){}",
state.stage,
execution_root.display(),
if state.worktree_path.is_some() {
" — evaluated against this phase's WORKTREE HEAD, not the main checkout; \
rebuild and reinstall the binary before resuming"
} else {
""
}
);
gates::fire_gate_notify(state.phase, state.stage, &message, true);
events::emit(
project_root,
state.phase,
"self_dogfood_stale_blocked",
serde_json::json!({
"stage": state.stage.to_string(),
"reason": "stale_build_blocked",
"worktree": state.worktree_path.is_some(),
}),
);
Err(CliError::Message(message))
}
StalenessOutcome::Warn => {
println!(
"warning: build provenance staleness check did not confirm a fresh build for \
stage {} — proceeding (only DevFlow's own workspace is ever hard-blocked, D-18)",
state.stage
);
Ok(())
}
StalenessOutcome::Ok => Ok(()),
}
}
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)?;
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(())
}
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)
}
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)?;
launch_stage(&mut state, None, None)
}
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(", ")
))),
}
}
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()))
}
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)
}
}
}
fn handle_infra_outcome(
project_root: &Path,
state: &mut State,
stage: Stage,
reason: Option<String>,
) -> Result<(), CliError> {
state.infra_failures = state.infra_failures.saturating_add(1);
workflow::save_state(state)?;
gate_or_abort_infra(project_root, state, stage, reason)
}
fn gate_or_abort_infra(
project_root: &Path,
state: &mut State,
stage: Stage,
reason: Option<String>,
) -> Result<(), CliError> {
if state.infra_failures >= mode::MAX_INFRA_FAILURES {
return abort(
project_root,
state,
&format!(
"infrastructure failures reached the ceiling ({} of {}) — aborting rather than gating again",
state.infra_failures,
mode::MAX_INFRA_FAILURES
),
);
}
handle_stage_failure(project_root, state, stage, reason)
}
fn handle_rate_limited_outcome(
project_root: &Path,
state: &mut State,
phase: u32,
stage: Stage,
reason: Option<String>,
) -> Result<(), CliError> {
let retry_after = retry_after_from_reason(reason.as_deref());
let projected_infra_failures = state.infra_failures.saturating_add(1);
if projected_infra_failures >= mode::MAX_INFRA_FAILURES {
return handle_infra_outcome(project_root, state, stage, reason);
}
state.infra_failures = projected_infra_failures;
workflow::save_state(state)?;
let instructions =
devflow_core::ship::build_single_agent_cron_instructions(project_root, phase, &retry_after);
devflow_core::ship::write_cron_instructions(project_root, &instructions)?;
if instructions.hermes_cron.schedule.is_empty() {
return gate_or_abort_infra(
project_root,
state,
stage,
Some(format!(
"rate limited with no parseable retry time ({retry_after}) — auto-resume cron not scheduled; resume manually"
)),
);
}
println!(
"rate limited — wrote {}",
devflow_core::ship::cron_instructions_path(project_root, phase)
.strip_prefix(project_root)
.map(|p| p.display().to_string())
.unwrap_or_else(|_| {
devflow_core::ship::cron_instructions_path(project_root, phase)
.display()
.to_string()
})
);
events::emit(
project_root,
phase,
"rate_limit_resume_scheduled",
serde_json::json!({
"stage": stage.to_string(),
"retry_after": retry_after,
"infra_failures": state.infra_failures,
}),
);
Ok(())
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum ValidateOutcome {
Passed,
Failed,
Ambiguous(String),
}
fn classify_validate_outcome(result: &agent_result::AgentResult) -> ValidateOutcome {
let external = result.decided_by_layer == Some(0) && result.status == AgentStatus::Success;
match (external, result.verdict) {
(_, Some(Verdict::Pass)) => ValidateOutcome::Passed,
(true, Some(Verdict::Gaps)) => ValidateOutcome::Ambiguous(
"external verification passed but the agent reported gaps".to_string(),
),
(true, None) => ValidateOutcome::Ambiguous(
"external verification passed but no agent verdict arrived".to_string(),
),
_ => ValidateOutcome::Failed,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ValidateResult {
Passed,
Failed,
}
fn handle_validate_outcome(
project_root: &Path,
state: &mut State,
outcome: ValidateOutcome,
) -> Result<(), CliError> {
let result = match outcome {
ValidateOutcome::Ambiguous(detail) => {
let context = format!(
"[never-silent] validate ambiguous: {}",
truncate_reason(&detail)
);
return match run_gate(project_root, state, Stage::Validate, &context)? {
GateAction::Advance => transition(project_root, state, Stage::Ship),
GateAction::LoopBack(_) => {
loop_back_to_code(project_root, state, FixType::GapsOnly)
}
GateAction::Abort(reason) => abort(project_root, state, &reason),
};
}
ValidateOutcome::Passed => ValidateResult::Passed,
ValidateOutcome::Failed => ValidateResult::Failed,
};
if result == ValidateResult::Failed {
state.consecutive_failures = state.consecutive_failures.saturating_add(1);
workflow::save_state(state)?;
}
if state
.mode
.should_gate(Stage::Validate, state.consecutive_failures)
{
let context = match result {
ValidateResult::Passed => "Validation passed — approve to ship?".to_string(),
ValidateResult::Failed => format!(
"Validation failed {} time(s) — human review needed.",
state.consecutive_failures
),
};
return match run_gate(project_root, state, Stage::Validate, &context)? {
GateAction::Advance => transition(project_root, state, Stage::Ship),
GateAction::LoopBack(_) => loop_back_to_code(project_root, state, FixType::GapsOnly),
GateAction::Abort(reason) => abort(project_root, state, &reason),
};
}
match result {
ValidateResult::Passed => transition(project_root, state, Stage::Ship),
ValidateResult::Failed => loop_back_to_code(project_root, state, FixType::GapsOnly),
}
}
fn handle_ship_outcome(project_root: &Path, state: &mut State) -> Result<(), CliError> {
match run_gate(
project_root,
state,
Stage::Ship,
"Ship complete — approve merge?",
)? {
GateAction::Advance => finish_workflow(project_root, state),
GateAction::LoopBack(_) => loop_back_to_code(project_root, state, FixType::GapsOnly),
GateAction::Abort(reason) => abort(project_root, state, &reason),
}
}
fn truncate_reason(reason: &str) -> String {
render_gate_context(reason, 300)
}
fn render_gate_context(context: &str, max_chars: usize) -> String {
const TRUNCATED: &str = "… [truncated; full output in .devflow/]";
let sanitized: String = context
.chars()
.map(|character| {
if character.is_control() {
' '
} else {
character
}
})
.collect();
if sanitized.chars().count() <= max_chars {
return sanitized;
}
let suffix_len = TRUNCATED.chars().count().min(max_chars);
let head_len = max_chars.saturating_sub(suffix_len);
let head: String = sanitized.chars().take(head_len).collect();
let suffix: String = TRUNCATED.chars().take(suffix_len).collect();
format!("{head}{suffix}")
}
fn handle_stage_failure(
project_root: &Path,
state: &mut State,
stage: Stage,
reason: Option<String>,
) -> Result<(), CliError> {
let context = format!(
"[never-silent] stage {stage} failed: {} — human review needed (retry, loop-to-code, or abort)",
truncate_reason(&reason.unwrap_or_else(|| "no details available".into()))
);
match run_gate(project_root, state, stage, &context)? {
GateAction::Advance => {
let _ = Gates::cleanup(project_root, state.phase, stage);
state.gate_pending = false;
launch_stage(state, None, Some(stage))
}
GateAction::LoopBack(_) => {
let _ = Gates::cleanup(project_root, state.phase, stage);
launch_stage(state, None, Some(stage))
}
GateAction::Abort(reason) => abort(project_root, state, &reason),
}
}
fn handle_ship_failure(
project_root: &Path,
state: &mut State,
reason: Option<String>,
) -> Result<(), CliError> {
if is_ship_review_failure(&reason) {
return loop_back_to_code(project_root, state, FixType::AuditFix);
}
handle_stage_failure(project_root, state, Stage::Ship, reason)
}
fn is_ship_review_failure(reason: &Option<String>) -> bool {
reason
.as_deref()
.map(|r| r.trim().to_ascii_lowercase().starts_with("review:"))
.unwrap_or(false)
}
fn hook_context_root(project_root: &Path, state: &State, terminal_batch: bool) -> PathBuf {
if terminal_batch {
return project_root.to_path_buf();
}
state
.worktree_path
.as_ref()
.filter(|path| path.exists())
.map(|path| path.to_path_buf())
.unwrap_or_else(|| project_root.to_path_buf())
}
fn run_checkout_hooks(
project_root: &Path,
state: &State,
batch: &[hooks::Hook],
stage: Stage,
) -> bool {
if batch.is_empty() {
return true;
}
let _checkout_lock = match lock::acquire_project_blocking(project_root, checkout_lock_timeout())
{
Ok(guard) => guard,
Err(err) => {
println!(
"warning: could not acquire the checkout lock ({err}) — \
SKIPPING hooks {batch:?} rather than mutating the checkout \
unserialized. Re-run them once the holder finishes."
);
events::emit(
project_root,
state.phase,
"checkout_lock_timeout",
serde_json::json!({ "stage": stage.to_string(), "error": err.to_string() }),
);
for hook in batch {
events::emit(
project_root,
state.phase,
"hook_run",
serde_json::json!({
"hook": format!("{hook:?}"),
"ok": false,
"skipped": "checkout lock timeout",
}),
);
}
return false;
}
};
let git_flow = GitFlowConfig::default();
let mut all_succeeded = true;
let terminal_batch = batch == hooks::hooks_after_ship().as_slice();
let hook_root = hook_context_root(project_root, state, terminal_batch);
let mut ctx = HookContext {
phase: state.phase,
project_root: hook_root.clone(),
stage,
git_flow: git_flow.clone(),
shipped_version: None,
};
for hook in batch {
let outcome = hook.run(&mut ctx);
if let Err(ref err) = outcome {
println!("warning: hook {hook:?} failed: {err}");
all_succeeded = false;
}
events::emit(
project_root,
state.phase,
"hook_run",
serde_json::json!({
"hook": format!("{hook:?}"),
"ok": outcome.is_ok(),
}),
);
if terminal_batch && outcome.is_err() {
break;
}
}
all_succeeded
}
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))
}
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))
}
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))
}
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(())
}
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"
)))
}
}
}
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(())
}
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());
}
fn ensure_phase_worktree(
project_root: &Path,
phase: u32,
force: bool,
) -> Result<PathBuf, CliError> {
let wt = worktree::phase_path(project_root, phase);
let branch = format!("{FEATURE_PREFIX}phase-{phase:02}");
if force {
if wt.exists() {
worktree::remove(project_root, &wt, true)?;
}
let _ = GitFlow::new(project_root).delete_branch(&branch, true);
}
match worktree::add(project_root, &wt, &branch, DEVELOP, true) {
Ok(()) => Ok(wt),
Err(devflow_core::worktree::WorktreeError::Exists(path)) => {
Err(CliError::Message(format!(
"worktree already exists at {} — use --force to recreate it",
path.display()
)))
}
Err(err) => Err(err.into()),
}
}
fn parse_phase_agent_pairs(
phases: &str,
agents: Option<&str>,
) -> Result<Vec<(u32, AgentKind)>, CliError> {
let phases: Vec<u32> = phases
.split(',')
.map(|p| p.trim())
.filter(|p| !p.is_empty())
.map(|p| {
p.parse::<u32>()
.map_err(|_| CliError::Message(format!("invalid phase number `{p}`")))
})
.collect::<Result<_, _>>()?;
if phases.is_empty() {
return Err(CliError::Message("no phases given".into()));
}
let agents: Vec<AgentKind> = match agents {
Some(list) => list
.split(',')
.map(|a| a.trim())
.filter(|a| !a.is_empty())
.map(|a| {
a.parse::<AgentKind>()
.map_err(|err| CliError::Message(err.to_string()))
})
.collect::<Result<_, _>>()?,
None => Vec::new(),
};
if agents.len() > phases.len() {
return Err(CliError::Message(format!(
"got {} agents for {} phases — provide at most one agent per phase",
agents.len(),
phases.len()
)));
}
Ok(phases
.into_iter()
.enumerate()
.map(|(i, phase)| (phase, agents.get(i).copied().unwrap_or(AgentKind::Claude)))
.collect())
}
fn parallel(
project_root: &Path,
phases: &str,
agents: Option<&str>,
mode: Mode,
force: bool,
) -> Result<(), CliError> {
let pairs = parse_phase_agent_pairs(phases, agents)?;
println!("launching {} phase(s) in parallel worktrees", pairs.len());
for (phase, agent) in pairs {
println!("\n=== phase {phase} ({agent}) ===");
start(project_root, phase, agent, mode, force, true, false)?;
}
Ok(())
}
fn split_two_agents(agents: &str) -> Result<(AgentKind, AgentKind), CliError> {
let parsed: Vec<AgentKind> = agents
.split(',')
.map(|a| a.trim())
.filter(|a| !a.is_empty())
.map(|a| {
a.parse::<AgentKind>()
.map_err(|err| CliError::Message(err.to_string()))
})
.collect::<Result<_, _>>()?;
if parsed.len() != 2 {
return Err(CliError::Message(format!(
"sequentagent requires exactly two agents (e.g. claude,codex), got {}",
parsed.len()
)));
}
Ok((parsed[0], parsed[1]))
}
fn run_agent_blocking(
project_root: &Path,
phase: u32,
agent: AgentKind,
workdir: &Path,
) -> Result<Option<agent_result::AgentResult>, CliError> {
if let Some(stamp) = agent_result::archive_phase_files(
project_root,
workdir,
phase,
capture_retention(project_root),
)
.map_err(|err| {
CliError::Message(format!(
"could not archive phase {phase} capture before rollover: {err}"
))
})? {
events::emit(
project_root,
phase,
"capture_archived",
serde_json::json!({"stage": "code", "stamp": stamp}),
);
}
let adapter = agents::adapter_for(agent);
let prompt = prompt::stage_prompt_for_project(Stage::Code, phase, project_root);
let roots = if workdir == project_root {
Vec::new()
} else {
worktree_writable_roots(project_root, workdir)
};
let (program, args) = adapter.exec_command(phase, &prompt, &roots);
ensure_agent_binary(program)?;
let mut state = State::new(phase, agent, Mode::Auto, project_root.to_path_buf());
state.stage = Stage::Code;
if workdir != project_root {
state.worktree_path = Some(workdir.to_path_buf());
}
let monitor_pid =
monitor::spawn_monitor_no_advance(&state, program, &args, &adapter.extra_env())
.map_err(|err| CliError::Message(format!("could not spawn monitor: {err}")))?;
println!(
"launched {} (monitor pid {monitor_pid}) in {}",
adapter.name(),
workdir.display()
);
println!(" watch live: devflow logs -f --phase {phase} [--stderr]");
let exit_code = monitor::wait_for_agent_exit(project_root, phase, monitor_pid)
.map_err(|err| CliError::Message(format!("agent run did not complete: {err}")))?;
println!("agent {agent} exited with code {exit_code}");
let result = agent_result::evaluate_layer1(project_root, phase);
if result.is_none() && exit_code != 0 {
return Ok(Some(agent_result::AgentResult {
status: AgentStatus::Failed,
exit_code: Some(exit_code),
reason: Some(format!(
"agent exited with code {exit_code} without reporting a result"
)),
commits: None,
summary: None,
verdict: None,
decided_by_layer: Some(2),
}));
}
Ok(result)
}
fn integrate_agent_branch(
project_root: &Path,
git: &GitFlow,
base: &str,
agent_branch: &str,
) -> Result<(), CliError> {
let _checkout_lock = lock::acquire_project_blocking(project_root, checkout_lock_timeout())
.map_err(|err| {
CliError::Message(format!(
"could not lock checkout to integrate {agent_branch} into {base}: {err}. \
Earlier integrations into {base} are already in place; once the lock \
holder finishes, integrate manually with \
`git fetch . {agent_branch}:{base}` — do NOT re-run sequentagent --force."
))
})?;
git.fast_forward_branch(base, agent_branch)?;
println!("integrated {agent_branch} into {base}");
if git.has_remote() {
match git.push(base) {
Ok(()) => println!("pushed {base} to origin"),
Err(err) => println!("warning: could not push {base}: {err}"),
}
}
Ok(())
}
fn sequentagent(
project_root: &Path,
phase: u32,
agents: &str,
force: bool,
) -> Result<(), CliError> {
let _phase_lock = match lock::acquire(project_root, phase) {
Ok(guard) => guard,
Err(lock::LockError::Contended { pid, path: _ }) => {
return Err(CliError::Message(format!(
"phase {phase} is already being driven by another devflow process (pid {pid})"
)));
}
Err(err) => return Err(CliError::Message(format!("lock error: {err}"))),
};
if let Err(err) = devflow_core::ship::delete_cron_instructions(project_root, phase) {
println!("warning: could not remove stale cron-instructions file: {err}");
}
let (agent_a, agent_b) = split_two_agents(agents)?;
ensure_agent_binary(agent_program(agent_a))?;
ensure_agent_binary(agent_program(agent_b))?;
let git = GitFlow::new(project_root);
let base = format!("{FEATURE_PREFIX}phase-{phase:02}");
{
let _checkout_lock = lock::acquire_project_blocking(project_root, checkout_lock_timeout())
.map_err(|err| CliError::Message(format!("could not lock checkout: {err}")))?;
git.ensure_branch(&base, DEVELOP)?;
}
let branch_a = format!("{base}-{agent_a}");
let branch_b = format!("{base}-{agent_b}");
let wt_a = worktree::phase_agent_path(project_root, phase, &agent_a.to_string());
let wt_b = worktree::phase_agent_path(project_root, phase, &agent_b.to_string());
if force {
for (wt, branch) in [(&wt_a, &branch_a), (&wt_b, &branch_b)] {
if wt.exists() {
worktree::remove(project_root, wt, true)?;
}
let _ = git.delete_branch(branch, true);
}
}
add_or_explain(project_root, &wt_a, &branch_a, &base)?;
add_or_explain(project_root, &wt_b, &branch_b, &base)?;
println!("worktree A: {} ({branch_a})", wt_a.display());
println!("worktree B: {} ({branch_b})", wt_b.display());
println!("\n=== agent A: {agent_a} ===");
if let Some(result) = run_agent_blocking(project_root, phase, agent_a, &wt_a)? {
match result.status {
AgentStatus::Failed => {
return Err(CliError::Message(format!(
"agent A ({agent_a}) failed: {}",
result.reason.unwrap_or_else(|| "no details".into())
)));
}
AgentStatus::RateLimited => {
let retry_after = retry_after_from_reason(result.reason.as_deref());
write_rate_limit_cron(project_root, phase, &retry_after, agents)?;
let commits = count_commits_between(project_root, &base, &branch_a)?;
if commits == 0 {
println!(
"Agent A rate-limited with zero commits; paused — resume record at {}",
devflow_core::ship::cron_instructions_path(project_root, phase).display()
);
return Ok(());
}
println!("Agent A rate-limited; handing off to agent B");
}
_ => {}
}
}
integrate_agent_branch(project_root, &git, &base, &branch_a)?;
git.rebase_in(&wt_b, &base).map_err(|err| {
CliError::Message(format!(
"rebase of {branch_b} onto {base} hit conflicts — resolve them in {} \
then re-run sequentagent: {err}",
wt_b.display()
))
})?;
println!("rebased {branch_b} onto {base}");
println!("\n=== agent B: {agent_b} ===");
if let Some(result) = run_agent_blocking(project_root, phase, agent_b, &wt_b)?
&& matches!(
result.status,
AgentStatus::Failed | AgentStatus::RateLimited
)
{
let label = if result.status == AgentStatus::RateLimited {
"rate-limited"
} else {
"failed"
};
return Err(CliError::Message(format!(
"agent B ({agent_b}) {label}: {}",
result.reason.unwrap_or_else(|| "no details".into())
)));
}
integrate_agent_branch(project_root, &git, &base, &branch_b)?;
if let Err(err) = devflow_core::ship::delete_cron_instructions(project_root, phase) {
println!("warning: could not remove cron-instructions file: {err}");
}
println!("\nsequentagent complete — both agents integrated into {base}");
Ok(())
}
fn retry_after_from_reason(reason: Option<&str>) -> String {
reason
.and_then(|s| s.strip_prefix("rate limited until "))
.unwrap_or("unknown")
.to_string()
}
fn write_rate_limit_cron(
project_root: &Path,
phase: u32,
retry_after: &str,
agents: &str,
) -> Result<(), CliError> {
let instructions =
devflow_core::ship::build_cron_instructions(project_root, phase, retry_after, agents);
devflow_core::ship::write_cron_instructions(project_root, &instructions)?;
if instructions.hermes_cron.schedule.is_empty() {
println!("no parseable retry time — auto-resume cron not scheduled; resume manually");
} else {
println!(
"wrote {}",
devflow_core::ship::cron_instructions_path(project_root, phase)
.strip_prefix(project_root)
.map(|p| p.display().to_string())
.unwrap_or_else(|_| {
devflow_core::ship::cron_instructions_path(project_root, phase)
.display()
.to_string()
})
);
}
Ok(())
}
fn count_commits_between(project_root: &Path, base: &str, branch: &str) -> Result<u32, CliError> {
let range = format!("{base}..{branch}");
let output = std::process::Command::new("git")
.args(["rev-list", "--count", &range])
.current_dir(project_root)
.output()
.map_err(|err| CliError::Message(format!("could not count commits on {branch}: {err}")))?;
if !output.status.success() {
return Err(CliError::Message(format!(
"could not count commits on {branch}: {}",
String::from_utf8_lossy(&output.stderr).trim()
)));
}
String::from_utf8_lossy(&output.stdout)
.trim()
.parse::<u32>()
.map_err(|err| CliError::Message(format!("invalid commit count for {branch}: {err}")))
}
fn add_or_explain(
project_root: &Path,
path: &Path,
branch: &str,
base: &str,
) -> Result<(), CliError> {
match worktree::add(project_root, path, branch, base, true) {
Ok(()) => Ok(()),
Err(devflow_core::worktree::WorktreeError::Exists(p)) => Err(CliError::Message(format!(
"worktree already exists at {} — use --force to recreate it",
p.display()
))),
Err(err) => Err(err.into()),
}
}
fn reference(project_root: &Path, branch: Option<String>, refresh: bool) -> Result<(), CliError> {
let branch = branch.unwrap_or_else(|| DEVELOP.to_string());
let path = worktree::reference_path(project_root);
if path.exists() {
if !refresh {
println!(
"reference exists at {} (use --refresh to update it)",
path.display()
);
return Ok(());
}
worktree::remove(project_root, &path, true)?;
worktree::add_detached(project_root, &path, &branch)?;
println!(
"refreshed reference worktree at {} (snapshot of {branch})",
path.display()
);
} else {
worktree::add_detached(project_root, &path, &branch)?;
println!(
"created reference worktree at {} (snapshot of {branch})",
path.display()
);
}
Ok(())
}
fn cleanup(project_root: &Path, force: bool) -> Result<(), CliError> {
let git = GitFlow::new(project_root);
let worktrees_dir = worktree::worktrees_dir(project_root);
let reference = worktree::reference_path(project_root);
let worktrees = worktree::list(project_root)?;
let mut removed = 0usize;
for wt in &worktrees {
if !wt.path.starts_with(&worktrees_dir) {
continue;
}
if wt.path == reference && !force {
println!("keeping reference worktree (use --force to remove it)");
continue;
}
worktree::remove(project_root, &wt.path, force)?;
print!("removed worktree {}", wt.path.display());
match &wt.branch {
Some(branch) if branch.starts_with(FEATURE_PREFIX) => {
match git.delete_branch(branch, force) {
Ok(()) => println!(" + deleted branch {branch}"),
Err(err) => println!(" (branch {branch} kept: {err})"),
}
}
_ => println!(),
}
removed += 1;
}
worktree::prune(project_root)?;
if removed == 0 {
println!("no worktrees to clean up");
}
match git.cleanup_merged() {
Ok(merged) => {
for branch in merged {
println!("deleted merged branch {branch}");
}
}
Err(err) => println!("warning: could not prune merged branches: {err}"),
}
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Liveness {
Healthy,
BetweenStages,
Stuck,
Unknown,
}
impl Liveness {
fn describe(self) -> &'static str {
match self {
Liveness::Healthy => "healthy",
Liveness::BetweenStages => "between stages",
Liveness::Stuck => "stuck — needs devflow resume",
Liveness::Unknown => "unknown (no monitor recorded)",
}
}
}
fn liveness(monitor_pid: Option<u32>, monitor_alive: bool, agent_alive: bool) -> Liveness {
match monitor_pid {
None => Liveness::Unknown,
Some(_) => match (monitor_alive, agent_alive) {
(true, true) => Liveness::Healthy,
(true, false) => Liveness::BetweenStages,
(false, _) => Liveness::Stuck,
},
}
}
fn status(project_root: &Path) -> Result<(), CliError> {
let states = workflow::list_states(project_root);
let mut current_worktree: Option<PathBuf> = None;
if states.is_empty() {
println!("stage: idle");
println!("project_root: {}", project_root.display());
} else {
let mut last_events = events::last_events_by_phase(project_root);
println!("project_root: {}", project_root.display());
println!(
"active phases: {}",
states
.iter()
.map(|s| s.phase.to_string())
.collect::<Vec<_>>()
.join(", ")
);
for state in &states {
let gate = if state.gate_pending {
"pending"
} else {
"none"
};
println!("\nphase {}:", state.phase);
println!(
" stage: {} | mode: {} | gate: {}",
state.stage, state.mode, gate
);
println!(" agent: {}", agents::adapter_for(state.agent).name());
if state.consecutive_failures > 0 {
println!(" validate failures: {}", state.consecutive_failures);
}
println!(
" started: {} ({})",
state.started_at,
recover::format_age(&state.started_at)
);
if let Some(ref wt) = state.worktree_path {
println!(" worktree: {}", wt.display());
}
current_worktree = current_worktree.or_else(|| state.worktree_path.clone());
let agent_pid = agent_pid_from_file(project_root, state.phase);
match agent_pid {
Some(pid) => {
println!(
" agent_pid: {pid} (running: {})",
agent::agent_running(pid)
);
}
None => println!(" agent_pid: none"),
}
match state.monitor_pid {
Some(pid) => {
println!(
" monitor_pid: {pid} (running: {})",
agent::agent_running(pid)
);
}
None => println!(" monitor_pid: none"),
}
let agent_alive = agent_pid.is_some_and(agent::agent_running);
let monitor_alive = state.monitor_pid.is_some_and(agent::agent_running);
let phase_liveness = liveness(state.monitor_pid, monitor_alive, agent_alive);
println!(" liveness: {}", phase_liveness.describe());
if phase_liveness == Liveness::Stuck {
println!(" → devflow resume --phase {}", state.phase);
}
if let Some(event) = last_events.remove(&state.phase) {
let ago = event
.get("ts")
.and_then(|t| t.as_u64())
.map(|t| format!(" ({})", recover::format_age(&t.to_string())))
.unwrap_or_default();
println!(" last action: {}{ago}", events::describe(&event));
}
}
}
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
if let Some(banner) = render_pending_gate_banner(&Gates::list_open(project_root), now) {
println!("\n{banner}");
}
print_open_branches(project_root);
print_worktrees(project_root, current_worktree.as_deref());
for hint in cron_instruction_hints(project_root) {
println!("\n{hint}");
}
Ok(())
}
fn render_pending_gate_banner(open: &[OpenGate], now: u64) -> Option<String> {
if open.is_empty() {
return None;
}
let mut banner = String::from("==================== PENDING GATE ====================\n");
for gate in open {
let timestamp = gate.timestamp.parse::<u64>().ok();
let escalated = timestamp
.and_then(|timestamp| now.checked_sub(timestamp))
.is_some_and(|age| age >= GATE_ESCALATION_THRESHOLD_SECS);
let marker = if escalated { "!!! ESCALATED" } else { "!!!" };
let context = render_gate_context(&gate.context, 300);
let stage = gate.stage.to_string();
banner.push_str(&format!(
"{marker}: phase {} {stage} ({})\n {context}\n approve: devflow gate approve {} --stage {stage}\n reject: devflow gate reject {} --stage {stage} --note <reason>\n",
gate.phase,
recover::format_age(&gate.timestamp),
gate.phase,
gate.phase,
));
}
banner.push_str("======================================================");
Some(banner)
}
fn gate_list(project_root: &Path) -> Result<(), CliError> {
let open = Gates::list_open(project_root);
if open.is_empty() {
println!("no open gates");
return Ok(());
}
println!("{:<6} {:<9} {:<9} CONTEXT", "PHASE", "STAGE", "AGE");
for gate in &open {
let context = render_gate_context(&gate.context, 100);
println!(
"{:<6} {:<9} {:<9} {context}",
gate.phase,
gate.stage.to_string(),
recover::format_age(&gate.timestamp),
);
}
println!(
"\nanswer with: devflow gate approve <phase> [--note ...] | \
devflow gate reject <phase> --note ... (note with \"abort\" ends the phase)"
);
Ok(())
}
fn gate_respond(
project_root: &Path,
phase: u32,
stage: Option<Stage>,
approved: bool,
note: Option<String>,
) -> Result<(), CliError> {
let stage = match stage {
Some(stage) => stage,
None => {
let open: Vec<_> = Gates::list_open(project_root)
.into_iter()
.filter(|g| g.phase == phase)
.collect();
match open.as_slice() {
[] => {
return Err(CliError::Message(format!(
"no open gate for phase {phase} — see `devflow gate list`"
)));
}
[one] => one.stage,
many => {
return Err(CliError::Message(format!(
"phase {phase} has several open gates ({}) — pass --stage",
many.iter()
.map(|g| g.stage.to_string())
.collect::<Vec<_>>()
.join(", ")
)));
}
}
}
};
let responded_by = std::env::var("USER")
.ok()
.filter(|user| !user.is_empty())
.unwrap_or_else(|| "devflow-cli".into());
let response = GateResponse {
approved,
note,
responded_by: Some(responded_by),
};
let path = Gates::respond(project_root, phase, stage, &response)?;
events::emit(
project_root,
phase,
"gate_response_written",
serde_json::json!({
"stage": stage.to_string(),
"approved": approved,
"via": "cli",
}),
);
let outcome = match GateAction::from_response(&response) {
GateAction::Advance => "workflow will advance",
GateAction::LoopBack(_) => "workflow will loop back to Code",
GateAction::Abort(_) => "phase will abort",
};
println!(
"{} gate for phase {phase} {stage} — {outcome} once the waiting monitor polls it \
(response at {})",
if approved { "approved" } else { "rejected" },
path.display()
);
Ok(())
}
fn logs(
project_root: &Path,
phase: Option<u32>,
follow: bool,
stderr: bool,
) -> Result<(), CliError> {
let phase = match phase {
Some(p) => p,
None => default_logs_phase(project_root)?,
};
let path = if stderr {
agent_result::stderr_path(project_root, phase)
} else {
agent_result::stdout_path(project_root, phase)
};
if !path.exists() && !follow {
return Err(CliError::Message(format!(
"no capture file for phase {phase} at {}",
path.display()
)));
}
eprintln!("== phase {phase}: {} ==", path.display());
let mut offset = print_capture_from(&path, 0)?;
if !follow {
return Ok(());
}
let exit_path = agent_result::exit_code_path(project_root, phase);
loop {
std::thread::sleep(std::time::Duration::from_millis(500));
let base = rollover_offset(&path, offset);
if base != offset {
eprintln!("== capture restarted (next stage) — following from the top ==");
}
let new_offset = print_capture_from(&path, base)?;
if exit_path.exists() && base == offset && new_offset == offset {
if let Ok(code) = std::fs::read_to_string(&exit_path) {
eprintln!("== agent exited with code {} ==", code.trim());
}
return Ok(());
}
offset = new_offset;
}
}
fn history_cmd(project_root: &Path, phase: Option<u32>) -> Result<(), CliError> {
let phase = match phase {
Some(phase) => phase,
None => single_active_phase(project_root)?.ok_or_else(|| {
CliError::Message("no active phase — pass a phase number to `devflow history`".into())
})?,
};
println!(
"{}",
history::render_timeline(&history::attempt_timeline(project_root, phase))
);
Ok(())
}
fn rollover_offset(path: &Path, offset: u64) -> u64 {
match std::fs::metadata(path) {
Ok(meta) if meta.len() < offset => 0,
_ => offset,
}
}
fn print_capture_from(path: &Path, offset: u64) -> Result<u64, CliError> {
use std::io::{Read, Seek, SeekFrom, Write};
let Ok(mut file) = std::fs::File::open(path) else {
return Ok(offset);
};
file.seek(SeekFrom::Start(offset))
.map_err(|err| CliError::Message(format!("could not seek capture file: {err}")))?;
let mut buf = Vec::new();
file.read_to_end(&mut buf)
.map_err(|err| CliError::Message(format!("could not read capture file: {err}")))?;
if !buf.is_empty() {
let mut stdout = std::io::stdout().lock();
let _ = stdout.write_all(&buf);
let _ = stdout.flush();
}
Ok(offset + buf.len() as u64)
}
fn default_logs_phase(project_root: &Path) -> Result<u32, CliError> {
if let Some(phase) = single_active_phase(project_root)? {
return Ok(phase);
}
let devflow = workflow::devflow_dir(project_root);
let mut newest: Option<(std::time::SystemTime, u32)> = None;
if let Ok(entries) = std::fs::read_dir(&devflow) {
for entry in entries.flatten() {
let name = entry.file_name();
let Some(name) = name.to_str() else { continue };
let Some(phase) = name
.strip_prefix("phase-")
.and_then(|rest| rest.strip_suffix("-stdout"))
.and_then(|num| num.parse::<u32>().ok())
else {
continue;
};
let Ok(modified) = entry.metadata().and_then(|m| m.modified()) else {
continue;
};
if newest.is_none_or(|(when, _)| modified > when) {
newest = Some((modified, phase));
}
}
}
newest.map(|(_, phase)| phase).ok_or_else(|| {
CliError::Message("no active phase and no capture files — nothing to show".into())
})
}
fn agent_pid_from_file(project_root: &Path, phase: u32) -> Option<u32> {
let path = agent_result::agent_pid_path(project_root, phase);
std::fs::read_to_string(path).ok()?.trim().parse().ok()
}
fn cron_instruction_hints(project_root: &Path) -> Vec<String> {
devflow_core::ship::list_cron_instructions(project_root)
.iter()
.map(|instructions| {
format!(
"Cron instruction pending (phase {}): hermes cron create --from-devflow {}",
instructions.phase,
project_root.display()
)
})
.collect()
}
fn print_worktrees(project_root: &Path, current: Option<&Path>) {
let worktrees_dir = worktree::worktrees_dir(project_root);
let worktrees = match worktree::list(project_root) {
Ok(w) => w,
Err(_) => return,
};
let active: Vec<_> = worktrees
.iter()
.filter(|w| w.path.starts_with(&worktrees_dir))
.collect();
if active.is_empty() {
return;
}
println!("\nactive worktrees:");
for wt in active {
let label = wt
.path
.file_name()
.map(|n| describe_worktree_dir(&n.to_string_lossy()))
.unwrap_or_default();
let branch = wt.branch.as_deref().unwrap_or("(detached)");
let marker = if current == Some(wt.path.as_path()) {
" *"
} else {
""
};
println!(" {} [{branch}]{label}{marker}", wt.path.display());
}
}
fn describe_worktree_dir(name: &str) -> String {
let Some(rest) = name.strip_prefix("phase-") else {
return String::new();
};
match rest.split_once('-') {
Some((phase, agent)) => {
format!(" — phase {}, agent {agent}", phase.trim_start_matches('0'))
}
None => format!(" — phase {}", rest.trim_start_matches('0')),
}
}
fn list(project_root: &Path) -> Result<(), CliError> {
let git = GitFlow::new(project_root);
let branches = git.list_feature_branches()?;
if branches.is_empty() {
println!("no open feature branches");
return Ok(());
}
println!(
"{:<25} {:>6} {:>7} LAST COMMIT",
"BRANCH", "AHEAD", "BEHIND"
);
for b in &branches {
println!(
"{:<25} {:>6} {:>7} {}",
b.name, b.ahead, b.behind, b.last_commit
);
}
Ok(())
}
fn print_open_branches(project_root: &Path) {
let git = GitFlow::new(project_root);
let branches = match git.list_feature_branches() {
Ok(b) => b,
Err(_) => return,
};
if branches.is_empty() {
return;
}
println!("\nopen branches:");
for b in &branches {
let staleness = if b.behind > 0 {
format!(" ({} behind develop)", b.behind)
} else {
String::new()
};
println!(" {} — {} ahead{staleness}", b.name, b.ahead);
}
}
fn project_root(project: PathBuf) -> Result<PathBuf, CliError> {
if !project.exists() {
return Err(CliError::Message(format!(
"project path does not exist: {}",
project.display()
)));
}
let start = project
.canonicalize()
.map_err(|err| CliError::Message(format!("failed to resolve project path: {err}")))?;
let mut probe = start.as_path();
loop {
if probe.join(".devflow").is_dir() {
return Ok(probe.to_path_buf());
}
match probe.parent() {
Some(parent) => probe = parent,
None => return Ok(start),
}
}
}
fn recover_cmd(project_root: &Path, do_clean: bool, phase: Option<u32>) -> Result<(), CliError> {
if do_clean {
let warnings = match phase {
Some(phase) => recover::clean_phase(project_root, phase)?,
None => recover::clean(project_root)?,
};
for warning in &warnings {
println!("warning: {warning}");
}
match phase {
Some(phase) => println!("cleaned up workflow state for phase {phase}"),
None => println!("cleaned up stale workflow state"),
}
return Ok(());
}
let statuses = match recover::inspect_all(project_root) {
Ok(s) => s,
Err(recover::RecoverError::NothingToRecover) => {
println!("no state to recover — project is idle");
return Ok(());
}
Err(err) => {
return Err(CliError::Message(format!(
"recover inspection failed: {err}"
)));
}
};
let mut any_stale = false;
for status in &statuses {
if let Some(only) = phase
&& status.state.phase != only
{
continue;
}
println!("phase: {}", status.state.phase);
println!(" stage: {}", status.state.stage);
println!(" mode: {}", status.state.mode);
println!(
" agent: {}",
agents::adapter_for(status.state.agent).name()
);
println!(" started: {} ({})", status.state.started_at, status.age);
match agent_pid_from_file(project_root, status.state.phase) {
Some(pid) => {
let running = agent::agent_running(pid);
println!(" agent_pid: {pid} (running: {running})");
if !running {
println!(" agent is not running — the monitor may have already advanced");
}
}
None => println!(" agent_pid: none"),
}
if status.is_stale {
any_stale = true;
println!(" state is stale");
}
}
if any_stale {
println!(
"\nstale state found — `devflow recover --clean` clears stale phases only; \
use `--clean --phase N` for a specific phase"
);
}
Ok(())
}
fn test_cmd(project_root: &Path) -> Result<(), CliError> {
let checks = [
("cargo test", "cargo test"),
(
"cargo clippy",
"cargo clippy --workspace --all-targets -- -D warnings",
),
("cargo fmt --check", "cargo fmt --check"),
];
let mut failures = Vec::new();
for (label, cmd) in checks {
println!("=== {label} ===");
let status = std::process::Command::new("sh")
.arg("-c")
.arg(cmd)
.current_dir(project_root)
.status()
.map_err(|err| CliError::Message(format!("could not run `{cmd}`: {err}")))?;
if status.success() {
println!(" ✓ {label}");
} else {
println!(" ✗ {label}");
failures.push(label);
}
}
if failures.is_empty() {
println!("\nall checks passed");
Ok(())
} else {
Err(CliError::Message(format!(
"quality checks failed: {}",
failures.join(", ")
)))
}
}
struct Check {
name: String,
status: String,
version: Option<String>,
install_hint: Option<String>,
}
fn doctor(project_root: &Path, json: bool) -> Result<(), CliError> {
use std::process::Command;
fn cmd_check(name: &str, cmd: &str, version_arg: &str, install_hint: &str) -> Check {
match Command::new(cmd).arg(version_arg).output() {
Ok(out) if out.status.success() => {
let version = String::from_utf8_lossy(&out.stdout)
.lines()
.next()
.unwrap_or("unknown")
.trim()
.to_string();
Check {
name: name.into(),
status: "ok".into(),
version: Some(version),
install_hint: None,
}
}
Ok(out) => {
let detail = String::from_utf8_lossy(&out.stderr)
.lines()
.next()
.unwrap_or("unknown")
.trim()
.to_string();
Check {
name: name.into(),
status: "warn".into(),
version: Some(detail),
install_hint: Some(format!(
"`{cmd} {version_arg}` exited non-zero — reinstall or check PATH"
)),
}
}
Err(_) => Check {
name: name.into(),
status: "missing".into(),
version: None,
install_hint: Some(install_hint.into()),
},
}
}
fn bool_check(name: &str, ok: bool, version: &str, install_hint: &str) -> Check {
Check {
name: name.into(),
status: if ok { "ok".into() } else { "missing".into() },
version: Some(version.into()),
install_hint: if ok { None } else { Some(install_hint.into()) },
}
}
let devflow_version = env!("CARGO_PKG_VERSION");
let (rust_log_status, rust_log_version, rust_log_hint) = match std::env::var("RUST_LOG") {
Ok(ref val) if val.is_empty() => (
"warn",
Some("empty (logging disabled)".into()),
Some("Set RUST_LOG=info for better diagnostics".into()),
),
Ok(val) => {
let all_valid = val.split(',').all(|directive| {
let directive = directive.trim();
if let Some((_target, level)) = directive.split_once('=') {
matches!(level.trim(), "error" | "warn" | "info" | "debug" | "trace")
} else {
matches!(directive, "error" | "warn" | "info" | "debug" | "trace")
}
});
if all_valid {
("ok", Some(val), None)
} else {
(
"warn",
Some(val),
Some("RUST_LOG value may be invalid — expected error, warn, info, debug, or trace".into()),
)
}
}
Err(_) => (
"missing",
Some("not set — defaulting to info".into()),
Some("Set RUST_LOG=info for better diagnostics".into()),
),
};
let checks: Vec<Check> = vec![
cmd_check(
"git",
"git",
"--version",
"Install from https://git-scm.com/downloads",
),
bool_check("sh (POSIX shell)", cfg!(unix), "built-in", "Unsupported OS"),
cmd_check(
"cargo/rust",
"cargo",
"--version",
"curl https://sh.rustup.rs -sSf | sh",
),
cmd_check(
"gh CLI",
"gh",
"--version",
"brew install gh / apt install gh",
),
cmd_check(
"claude",
"claude",
"--version",
"npm i -g @anthropic-ai/claude-code",
),
cmd_check("codex", "codex", "--version", "npm i -g @openai/codex"),
cmd_check(
"opencode",
"opencode",
"--version",
"cargo install opencode",
),
Check {
name: format!("devflow v{devflow_version}"),
status: "ok".into(),
version: Some(devflow_version.into()),
install_hint: None,
},
Check {
name: "RUST_LOG".into(),
status: rust_log_status.into(),
version: rust_log_version,
install_hint: rust_log_hint,
},
];
let facts = collect_phase_facts(project_root);
if json {
let body = doctor_json_body(&checks, &facts);
println!(
"{}",
serde_json::to_string_pretty(&body).expect("doctor --json body must serialize")
);
} else {
for c in &checks {
let icon = match c.status.as_str() {
"ok" => "✓",
"missing" => "✗",
"warn" => "⚠",
_ => "?",
};
let version_str = c.version.as_deref().unwrap_or("-");
print!(" {:<20} {:<20} {}", c.name, version_str, icon);
#[allow(clippy::collapsible_if)]
if c.status == "missing" || c.status == "warn" {
if let Some(hint) = &c.install_hint {
print!(" — {}", hint);
}
}
println!();
}
print!("{}", render_reconciliation_text(&facts));
}
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Severity {
Ok,
Warn,
Problem,
}
impl Severity {
fn label(self) -> &'static str {
match self {
Severity::Ok => "ok",
Severity::Warn => "warn",
Severity::Problem => "problem",
}
}
}
struct PhaseFacts {
phase: u32,
stage: Stage,
gate_pending: bool,
agent_pid: Option<u32>,
agent_alive: bool,
monitor_pid: Option<u32>,
monitor_alive: bool,
last_event: Option<String>,
last_launched_stage: Option<Stage>,
open_gate_stages: Vec<Stage>,
feature_branch_exists: bool,
}
struct PhaseFinding {
phase: u32,
severity: Severity,
detail: String,
repair: Option<String>,
}
fn check_gate_pending_without_gate(facts: &PhaseFacts) -> Option<PhaseFinding> {
if !facts.gate_pending || !facts.open_gate_stages.is_empty() {
return None;
}
Some(PhaseFinding {
phase: facts.phase,
severity: Severity::Problem,
detail: format!(
"phase {}: gate_pending is true at stage {} but no gate file is open",
facts.phase, facts.stage
),
repair: Some(format!("devflow resume --phase {}", facts.phase)),
})
}
fn check_orphan_gate(facts: &PhaseFacts) -> Option<PhaseFinding> {
if facts.gate_pending || facts.open_gate_stages.is_empty() {
return None;
}
let gate_stage = facts.open_gate_stages[0];
Some(PhaseFinding {
phase: facts.phase,
severity: Severity::Problem,
detail: format!(
"phase {}: gate open for stage {} but state.gate_pending is false",
facts.phase, gate_stage
),
repair: Some(format!(
"devflow gate approve {} --stage {}",
facts.phase, gate_stage
)),
})
}
fn check_dead_agent(facts: &PhaseFacts) -> Option<PhaseFinding> {
let pid = facts.agent_pid?;
if facts.agent_alive || !facts.stage.is_agent_stage() {
return None;
}
Some(PhaseFinding {
phase: facts.phase,
severity: Severity::Problem,
detail: format!(
"phase {}: agent pid {pid} recorded but not running at stage {}",
facts.phase, facts.stage
),
repair: Some(format!("devflow resume --phase {}", facts.phase)),
})
}
fn check_dead_monitor(facts: &PhaseFacts) -> Option<PhaseFinding> {
if liveness(facts.monitor_pid, facts.monitor_alive, facts.agent_alive) != Liveness::Stuck {
return None;
}
let pid = facts.monitor_pid?;
Some(PhaseFinding {
phase: facts.phase,
severity: Severity::Problem,
detail: format!(
"phase {}: monitor pid {pid} recorded but not running at stage {}",
facts.phase, facts.stage
),
repair: Some(format!("devflow resume --phase {}", facts.phase)),
})
}
fn check_stage_event_drift(facts: &PhaseFacts) -> Option<PhaseFinding> {
let launched = facts.last_launched_stage?;
if launched == facts.stage {
return None;
}
Some(PhaseFinding {
phase: facts.phase,
severity: Severity::Warn,
detail: format!(
"phase {}: last stage_launched event named {launched} but state.stage is {}",
facts.phase, facts.stage
),
repair: None,
})
}
fn check_missing_branch(facts: &PhaseFacts) -> Option<PhaseFinding> {
if facts.feature_branch_exists || facts.stage == Stage::Define {
return None;
}
Some(PhaseFinding {
phase: facts.phase,
severity: Severity::Warn,
detail: format!(
"phase {}: feature/phase-{:02} does not exist but stage is {}",
facts.phase, facts.phase, facts.stage
),
repair: None,
})
}
fn reconcile_phase(facts: &PhaseFacts) -> Vec<PhaseFinding> {
[
check_gate_pending_without_gate(facts),
check_orphan_gate(facts),
check_dead_agent(facts),
check_dead_monitor(facts),
check_stage_event_drift(facts),
check_missing_branch(facts),
]
.into_iter()
.flatten()
.collect()
}
fn collect_phase_facts(project_root: &Path) -> Vec<PhaseFacts> {
let states = workflow::list_states(project_root);
let mut last_events = events::last_events_by_phase(project_root);
let open_gates = Gates::list_open(project_root);
let mut facts: Vec<PhaseFacts> = states
.into_iter()
.map(|state| build_phase_facts(project_root, state, &mut last_events, &open_gates))
.collect();
facts.sort_by_key(|f| f.phase);
facts
}
fn build_phase_facts(
project_root: &Path,
state: State,
last_events: &mut std::collections::HashMap<u32, serde_json::Value>,
open_gates: &[OpenGate],
) -> PhaseFacts {
let phase = state.phase;
let agent_pid = agent_pid_from_file(project_root, phase);
let agent_alive = agent_pid.is_some_and(agent::agent_running);
let monitor_pid = state.monitor_pid;
let monitor_alive = monitor_pid.is_some_and(agent::agent_running);
let last_event = last_events.remove(&phase);
let last_launched_stage = last_event.as_ref().and_then(last_launched_stage_from_event);
let last_event_name = last_event
.as_ref()
.and_then(|e| e.get("event"))
.and_then(|e| e.as_str())
.map(str::to_string);
let open_gate_stages = open_gates
.iter()
.filter(|g| g.phase == phase)
.map(|g| g.stage)
.collect();
let branch_ref = format!("refs/heads/feature/phase-{phase:02}");
let feature_branch_exists =
run_git_stdout(project_root, &["rev-parse", "--verify", &branch_ref]).is_some();
PhaseFacts {
phase,
stage: state.stage,
gate_pending: state.gate_pending,
agent_pid,
agent_alive,
monitor_pid,
monitor_alive,
last_event: last_event_name,
last_launched_stage,
open_gate_stages,
feature_branch_exists,
}
}
fn last_launched_stage_from_event(event: &serde_json::Value) -> Option<Stage> {
if event.get("event").and_then(|e| e.as_str()) != Some("stage_launched") {
return None;
}
event
.get("stage")
.and_then(|s| s.as_str())
.and_then(|s| s.parse::<Stage>().ok())
}
fn findings_for_display(facts: &PhaseFacts) -> Vec<PhaseFinding> {
let findings = reconcile_phase(facts);
if !findings.is_empty() {
return findings;
}
vec![PhaseFinding {
phase: facts.phase,
severity: Severity::Ok,
detail: format!("phase {}: ok", facts.phase),
repair: None,
}]
}
fn render_reconciliation_text(facts: &[PhaseFacts]) -> String {
let mut out = String::from("\nreconciliation:\n");
if facts.is_empty() {
out.push_str(" no active phases — nothing to reconcile\n");
return out;
}
for phase_facts in facts {
for finding in findings_for_display(phase_facts) {
out.push_str(&format!(" {}\n", finding.detail));
if let Some(repair) = &finding.repair {
out.push_str(&format!(" repair: {repair}\n"));
}
}
}
out
}
fn render_reconciliation_json(facts: &[PhaseFacts]) -> serde_json::Value {
let findings: Vec<(&PhaseFacts, PhaseFinding)> = facts
.iter()
.flat_map(|pf| findings_for_display(pf).into_iter().map(move |f| (pf, f)))
.collect();
serde_json::Value::Array(
findings
.iter()
.map(|(phase_facts, finding)| {
serde_json::json!({
"phase": finding.phase,
"severity": finding.severity.label(),
"detail": finding.detail,
"repair": finding.repair,
"last_event": phase_facts.last_event,
})
})
.collect(),
)
}
fn checks_json_value(checks: &[Check]) -> serde_json::Value {
serde_json::Value::Array(
checks
.iter()
.map(|c| {
serde_json::json!({
"name": c.name,
"status": c.status,
"version": c.version,
"install_hint": c.install_hint,
})
})
.collect(),
)
}
fn doctor_json_body(checks: &[Check], facts: &[PhaseFacts]) -> serde_json::Value {
serde_json::json!({
"environment": checks_json_value(checks),
"reconciliation": render_reconciliation_json(facts),
})
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
static ENV_MUTEX: Mutex<()> = Mutex::new(());
#[test]
fn project_root_walks_up_to_nearest_devflow_ancestor() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("project");
let nested = root.join(".worktrees/phase-16/deep");
std::fs::create_dir_all(root.join(".devflow")).unwrap();
std::fs::create_dir_all(&nested).unwrap();
assert_eq!(project_root(nested).unwrap(), root.canonicalize().unwrap());
let idle = dir.path().join("idle/nested");
std::fs::create_dir_all(&idle).unwrap();
assert_eq!(
project_root(idle.clone()).unwrap(),
idle.canonicalize().unwrap()
);
let missing = dir.path().join("missing");
let error = project_root(missing).unwrap_err().to_string();
assert!(error.contains("project path does not exist"));
}
#[test]
fn gate_approve_arg_parsing_accepts_positional_stage() {
let cli = Cli::try_parse_from(["devflow", "gate", "approve", "15", "ship"]).unwrap();
let Command::Gate {
action: GateCmd::Approve { stage, project, .. },
} = cli.command
else {
panic!("expected gate approve command");
};
assert_eq!(stage.as_deref(), Some("ship"));
assert_eq!(project, PathBuf::from("."));
let flagged =
Cli::try_parse_from(["devflow", "gate", "approve", "15", "--stage", "ship"]).unwrap();
let Command::Gate {
action:
GateCmd::Approve {
stage,
stage_option,
..
},
} = flagged.command
else {
panic!("expected flagged gate approve command");
};
assert_eq!(stage, None);
assert_eq!(stage_option, Some(Stage::Ship));
let bare = Cli::try_parse_from(["devflow", "gate", "approve", "15"]).unwrap();
let Command::Gate {
action:
GateCmd::Approve {
stage,
stage_option,
..
},
} = bare.command
else {
panic!("expected bare gate approve command");
};
assert_eq!(stage, None);
assert_eq!(stage_option, None);
let legacy =
Cli::try_parse_from(["devflow", "gate", "approve", "15", "/tmp/example-project"])
.unwrap();
let Command::Gate {
action:
GateCmd::Approve {
stage,
legacy_project,
stage_option,
project,
..
},
} = legacy.command
else {
panic!("expected legacy gate approve command");
};
let (stage, project) =
resolve_gate_target(stage, legacy_project, stage_option, project).unwrap();
assert_eq!(stage, None);
assert_eq!(project, PathBuf::from("/tmp/example-project"));
}
#[test]
fn pairs_default_missing_agents_to_claude() {
let pairs = parse_phase_agent_pairs("7,8", Some("codex")).unwrap();
assert_eq!(pairs, vec![(7, AgentKind::Codex), (8, AgentKind::Claude)]);
}
#[test]
fn pairs_match_agents_positionally() {
let pairs = parse_phase_agent_pairs("7, 8", Some("claude, codex")).unwrap();
assert_eq!(pairs, vec![(7, AgentKind::Claude), (8, AgentKind::Codex)]);
}
#[test]
fn pairs_default_all_to_claude_without_agents() {
let pairs = parse_phase_agent_pairs("3,4", None).unwrap();
assert_eq!(pairs, vec![(3, AgentKind::Claude), (4, AgentKind::Claude)]);
}
#[test]
fn pairs_reject_more_agents_than_phases() {
let err = parse_phase_agent_pairs("7", Some("claude,codex")).unwrap_err();
assert!(matches!(err, CliError::Message(_)));
}
#[test]
fn pairs_reject_invalid_phase() {
assert!(parse_phase_agent_pairs("7,x", None).is_err());
assert!(parse_phase_agent_pairs("", None).is_err());
}
#[test]
fn describe_worktree_dir_infers_phase_and_agent() {
assert_eq!(
describe_worktree_dir("phase-07-claude"),
" — phase 7, agent claude"
);
assert_eq!(describe_worktree_dir("phase-08"), " — phase 8");
assert_eq!(describe_worktree_dir("reference"), "");
}
#[test]
fn split_two_agents_requires_exactly_two() {
assert_eq!(
split_two_agents("claude, codex").unwrap(),
(AgentKind::Claude, AgentKind::Codex)
);
assert!(split_two_agents("claude").is_err());
assert!(split_two_agents("claude,codex,opencode").is_err());
assert!(split_two_agents("claude,bogus").is_err());
}
#[test]
fn retry_after_from_reason_strips_prefix() {
assert_eq!(
retry_after_from_reason(Some("rate limited until 2026-06-18T15:45:30Z")),
"2026-06-18T15:45:30Z"
);
assert_eq!(retry_after_from_reason(Some("usage limit")), "unknown");
assert_eq!(retry_after_from_reason(None), "unknown");
}
#[test]
fn cron_instruction_hints_include_hermes_command_per_phase() {
let dir = tempfile::tempdir().unwrap();
for phase in [7, 9] {
let instructions = devflow_core::ship::build_cron_instructions(
dir.path(),
phase,
"2026-06-18T15:45:30Z",
"claude,codex",
);
devflow_core::ship::write_cron_instructions(dir.path(), &instructions).unwrap();
}
let hints = cron_instruction_hints(dir.path());
assert_eq!(hints.len(), 2);
assert_eq!(
hints[0],
format!(
"Cron instruction pending (phase 7): hermes cron create --from-devflow {}",
dir.path().display()
)
);
assert!(hints[1].contains("(phase 9)"));
}
#[test]
fn parse_checkout_lock_timeout_defaults_and_parses() {
assert_eq!(
parse_checkout_lock_timeout(None),
std::time::Duration::from_secs(120)
);
assert_eq!(
parse_checkout_lock_timeout(Some("5".into())),
std::time::Duration::from_secs(5)
);
assert_eq!(
parse_checkout_lock_timeout(Some("nope".into())),
std::time::Duration::from_secs(120)
);
}
#[test]
fn checkout_hooks_skip_instead_of_running_unserialized_on_lock_timeout() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let _held = lock::acquire_project(root).expect("hold checkout lock");
unsafe {
std::env::set_var("DEVFLOW_CHECKOUT_LOCK_TIMEOUT_SECS", "0");
}
let state = State::new(33, AgentKind::Claude, Mode::Auto, root.to_path_buf());
run_checkout_hooks(root, &state, &hooks::hooks_after_ship(), Stage::Ship);
unsafe {
std::env::remove_var("DEVFLOW_CHECKOUT_LOCK_TIMEOUT_SECS");
}
assert!(
!root.join("CHANGELOG.md").exists(),
"hooks must not run while the checkout lock is held elsewhere"
);
let last = devflow_core::events::last_event_for_phase(root, 33)
.expect("skip must be recorded in events.jsonl");
assert_eq!(last["event"], "hook_run");
assert_eq!(last["ok"], false);
assert_eq!(last["skipped"], "checkout lock timeout");
}
#[test]
fn terminal_hook_failure_stops_before_branch_cleanup() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo(root);
let phase = 34;
let branch = "feature/phase-34";
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(&["branch", branch, "develop"]);
std::fs::remove_file(root.join("Cargo.toml")).unwrap();
std::fs::create_dir(root.join("Cargo.toml")).unwrap();
let state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
let succeeded = run_checkout_hooks(root, &state, &hooks::hooks_after_ship(), Stage::Ship);
assert!(!succeeded);
assert!(
GitFlow::new(root).branch_exists(branch),
"a failed terminal batch must preserve the branch for retry"
);
}
fn init_repo_no_version_file(root: &Path) {
let git = |args: &[&str]| {
let ok = std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success();
assert!(ok, "git {args:?} failed");
};
git(&["init", "-q"]);
git(&["config", "user.email", "devflow@example.com"]);
git(&["config", "user.name", "DevFlow Tests"]);
git(&["config", "commit.gpgsign", "false"]);
git(&["config", "tag.gpgsign", "false"]);
git(&["config", "core.hooksPath", "/dev/null"]);
std::fs::write(root.join("README.md"), "no version file in this repo\n").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "init"]);
git(&["branch", "-M", "main"]);
git(&["checkout", "-q", "-b", "develop"]);
}
#[test]
fn run_checkout_hooks_keeps_changelog_in_sync_with_tag_when_no_version_file() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
init_repo_no_version_file(root);
let phase = 47;
let branch = format!("feature/phase-{phase:02}");
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(&["branch", &branch, "develop"]);
std::fs::write(root.join(".gitignore"), ".devflow/\n").unwrap();
git(&["checkout", &branch]);
std::fs::write(root.join("feature.txt"), "phase work\n").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "phase work"]);
git(&["checkout", "develop"]);
let state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
let succeeded = run_checkout_hooks(root, &state, &hooks::hooks_after_ship(), Stage::Ship);
assert!(
succeeded,
"after-ship batch must succeed against a clean repo"
);
let all_tags = std::process::Command::new("git")
.arg("tag")
.current_dir(root)
.output()
.unwrap();
let all_tags = String::from_utf8_lossy(&all_tags.stdout);
assert_eq!(all_tags.lines().count(), 1, "expected exactly one tag");
let tag = all_tags.trim().to_string();
let tag_version = tag
.strip_prefix('v')
.expect("tag should be prefixed with v")
.to_string();
let changelog = std::fs::read_to_string(root.join("CHANGELOG.md")).unwrap();
let changelog_version = changelog
.lines()
.find(|l| l.starts_with("## "))
.and_then(|l| l.trim_start_matches("## ").split(' ').next())
.unwrap()
.to_string();
assert_ne!(
changelog_version, "unreleased",
"changelog heading must name the tagged version, not fall back to the literal"
);
assert_eq!(
changelog_version, tag_version,
"changelog heading must match the git tag ({tag}) produced by the same \
run_checkout_hooks call, even with no version file"
);
}
#[test]
fn default_logs_phase_prefers_single_active_state() {
let dir = tempfile::tempdir().unwrap();
let state = State::new(6, AgentKind::Claude, Mode::Auto, dir.path().to_path_buf());
workflow::save_state(&state).unwrap();
assert_eq!(default_logs_phase(dir.path()).unwrap(), 6);
}
#[test]
fn default_logs_phase_is_ambiguous_with_two_active_states() {
let dir = tempfile::tempdir().unwrap();
for phase in [6, 7] {
let state = State::new(
phase,
AgentKind::Claude,
Mode::Auto,
dir.path().to_path_buf(),
);
workflow::save_state(&state).unwrap();
}
let err = default_logs_phase(dir.path()).unwrap_err();
assert!(err.to_string().contains("--phase"));
}
#[test]
fn default_logs_phase_falls_back_to_newest_capture_file() {
let dir = tempfile::tempdir().unwrap();
std::fs::create_dir_all(dir.path().join(".devflow")).unwrap();
std::fs::write(agent_result::stdout_path(dir.path(), 3), "old").unwrap();
std::thread::sleep(std::time::Duration::from_millis(20));
std::fs::write(agent_result::stdout_path(dir.path(), 5), "new").unwrap();
assert_eq!(default_logs_phase(dir.path()).unwrap(), 5);
}
#[test]
fn default_logs_phase_errors_with_nothing_to_show() {
let dir = tempfile::tempdir().unwrap();
assert!(default_logs_phase(dir.path()).is_err());
}
#[test]
fn liveness_unknown_when_no_monitor_recorded() {
assert_eq!(liveness(None, false, false), Liveness::Unknown);
assert_eq!(liveness(None, false, true), Liveness::Unknown);
assert_eq!(liveness(None, true, false), Liveness::Unknown);
assert_eq!(liveness(None, true, true), Liveness::Unknown);
}
#[test]
fn liveness_matrix_covers_all_four_rows() {
let pid = Some(4242);
assert_eq!(liveness(pid, true, true), Liveness::Healthy);
assert_eq!(liveness(pid, true, false), Liveness::BetweenStages);
assert_eq!(liveness(pid, false, false), Liveness::Stuck);
assert_eq!(liveness(pid, false, true), Liveness::Stuck);
}
#[test]
fn liveness_treats_zero_and_overflow_pids_as_dead() {
assert!(!agent::agent_running(0));
assert!(!agent::agent_running(u32::MAX));
}
#[test]
fn monitor_pid_persisted_for_one_phase_does_not_disturb_a_sibling() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let mut phase7 = State::new(7, AgentKind::Claude, Mode::Auto, root.to_path_buf());
phase7.monitor_pid = Some(111);
workflow::save_state(&phase7).unwrap();
let mut phase8 = State::new(8, AgentKind::Claude, Mode::Auto, root.to_path_buf());
phase8.monitor_pid = Some(222);
workflow::save_state(&phase8).unwrap();
let reloaded7 = workflow::load_state(root, 7).unwrap();
let reloaded8 = workflow::load_state(root, 8).unwrap();
assert_eq!(reloaded7.monitor_pid, Some(111));
assert_eq!(reloaded8.monitor_pid, Some(222));
}
#[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 status_reading_monitor_liveness_writes_no_state_and_no_event() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 66;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.monitor_pid = Some(u32::MAX);
workflow::save_state(&state).unwrap();
let state_path = workflow::state_path(root, phase);
let before_len = std::fs::metadata(&state_path).unwrap().len();
let before_modified = std::fs::metadata(&state_path).unwrap().modified().unwrap();
let events_log = events::events_path(root);
let before_lines = std::fs::read_to_string(&events_log)
.unwrap_or_default()
.lines()
.count();
status(root).unwrap();
status(root).unwrap();
let after_len = std::fs::metadata(&state_path).unwrap().len();
let after_modified = std::fs::metadata(&state_path).unwrap().modified().unwrap();
let after_lines = std::fs::read_to_string(&events_log)
.unwrap_or_default()
.lines()
.count();
assert_eq!(
before_len, after_len,
"status must not rewrite the state file"
);
assert_eq!(
before_modified, after_modified,
"status must not touch the state file's mtime"
);
assert_eq!(
before_lines, after_lines,
"status must not append to events.jsonl"
);
}
#[test]
fn gate_respond_auto_resolves_single_open_gate() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
Gates::write_gate(root, 15, Stage::Ship, "approve merge?").unwrap();
gate_respond(root, 15, None, true, Some("lgtm".into())).unwrap();
let polled = Gates::poll_response(root, 15, Stage::Ship, 1).expect("response readable");
assert!(polled.approved);
assert_eq!(polled.note.as_deref(), Some("lgtm"));
let event = devflow_core::events::last_event_for_phase(root, 15).unwrap();
assert_eq!(event["event"], "gate_response_written");
assert_eq!(event["stage"], "ship");
}
#[test]
fn gate_respond_requires_stage_when_ambiguous_and_errors_when_none_open() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let err = gate_respond(root, 15, None, true, None).unwrap_err();
assert!(err.to_string().contains("no open gate"), "{err}");
Gates::write_gate(root, 15, Stage::Validate, "a").unwrap();
Gates::write_gate(root, 15, Stage::Ship, "b").unwrap();
let err = gate_respond(root, 15, None, false, Some("nope".into())).unwrap_err();
assert!(err.to_string().contains("--stage"), "{err}");
gate_respond(root, 15, Some(Stage::Validate), false, Some("gaps".into())).unwrap();
assert!(
Gates::response_path(root, 15, Stage::Validate).exists(),
"explicit-stage rejection must land"
);
assert!(!Gates::response_path(root, 15, Stage::Ship).exists());
}
#[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 rollover_offset_resets_on_shrunken_capture() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("capture");
std::fs::write(&path, "abc").unwrap();
assert_eq!(rollover_offset(&path, 10), 0);
assert_eq!(rollover_offset(&path, 3), 3);
assert_eq!(rollover_offset(&path, 2), 2);
assert_eq!(rollover_offset(&dir.path().join("gone"), 7), 7);
}
#[test]
fn print_capture_from_tracks_offsets_across_appends() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("capture");
std::fs::write(&path, "hello ").unwrap();
let offset = print_capture_from(&path, 0).unwrap();
assert_eq!(offset, 6);
assert_eq!(print_capture_from(&path, offset).unwrap(), 6);
use std::io::Write as _;
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(&path)
.unwrap();
f.write_all(b"world").unwrap();
drop(f);
assert_eq!(print_capture_from(&path, offset).unwrap(), 11);
assert_eq!(
print_capture_from(Path::new("/nonexistent/x"), 4).unwrap(),
4
);
}
fn init_repo(root: &Path) {
let git = |args: &[&str]| {
let ok = std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success();
assert!(ok, "git {args:?} failed");
};
git(&["init", "-q"]);
git(&["config", "user.email", "devflow@example.com"]);
git(&["config", "user.name", "DevFlow Tests"]);
git(&["config", "commit.gpgsign", "false"]);
git(&["config", "tag.gpgsign", "false"]);
git(&["config", "core.hooksPath", "/dev/null"]);
std::fs::write(root.join("Cargo.toml"), "[package]\nversion = \"2.0.0\"\n").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "init"]);
git(&["branch", "-M", "main"]);
git(&["checkout", "-q", "-b", "develop"]);
}
#[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 validate_failure_threshold_forces_gate_then_aborts() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 22;
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_eq!(state.consecutive_failures, mode::MAX_CONSECUTIVE_FAILURES);
assert!(
!Gates::gate_path(root, phase, Stage::Validate).exists(),
"forced gate's files must be cleaned up once it resolves to Abort"
);
let err = workflow::load_state(root, phase).unwrap_err();
assert!(matches!(err, workflow::WorkflowError::MissingState(_)));
}
fn drive_validate_advance_and_read_gate_context(
root: &Path,
phase: u32,
consecutive_failures: u32,
verdict_json: Option<&str>,
) -> String {
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
state.consecutive_failures = consecutive_failures;
workflow::save_state(&state).unwrap();
std::fs::create_dir_all(root.join(".devflow")).unwrap();
let marker = match verdict_json {
Some(verdict) => {
format!(r#"DEVFLOW_RESULT: {{"status":"success","verdict":"{verdict}"}}"#)
}
None => r#"DEVFLOW_RESULT: {"status":"success"}"#.to_string(),
};
std::fs::write(agent_result::stdout_path(root, phase), marker).unwrap();
let gate_path = Gates::gate_path(root, phase, Stage::Validate);
let response_path = Gates::response_path(root, phase, Stage::Validate);
let mut context = String::new();
std::thread::scope(|scope| {
scope.spawn(|| {
advance(root, Some(phase)).unwrap();
});
let mut seen = false;
for _ in 0..150 {
if gate_path.exists() {
seen = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
assert!(
seen,
"advance() must force a Validate gate, not advance silently"
);
context = std::fs::read_to_string(&gate_path).unwrap();
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();
});
context
}
#[test]
fn validate_gaps_does_not_advance_to_ship() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let context = drive_validate_advance_and_read_gate_context(
root,
60,
mode::MAX_CONSECUTIVE_FAILURES - 1,
Some("gaps"),
);
assert!(
context.contains("Validation failed"),
"a gaps verdict must be treated as a failed validation, not a pass: {context}"
);
}
#[test]
fn validate_missing_verdict_does_not_advance() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let context = drive_validate_advance_and_read_gate_context(
root,
61,
mode::MAX_CONSECUTIVE_FAILURES - 1,
None,
);
assert!(
context.contains("Validation failed"),
"a missing verdict must be treated as a failed validation, not a pass: {context}"
);
}
#[test]
fn validate_pass_advances() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let context = drive_validate_advance_and_read_gate_context(
root,
62,
mode::MAX_CONSECUTIVE_FAILURES,
Some("pass"),
);
assert!(
context.contains("Validation passed"),
"an explicit pass verdict must advance to Ship: {context}"
);
}
#[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 external_verify_agreement_advances_to_ship() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 90;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
workflow::save_state(&state).unwrap();
let result = agent_result::AgentResult {
status: AgentStatus::Success,
exit_code: None,
reason: None,
commits: None,
summary: None,
verdict: Some(Verdict::Pass),
decided_by_layer: Some(0),
};
let outcome = classify_validate_outcome(&result);
assert_eq!(outcome, ValidateOutcome::Passed);
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 _ = handle_validate_outcome(root, &mut state, outcome);
unsafe {
match &original_path {
Some(path) => std::env::set_var("PATH", path),
None => std::env::remove_var("PATH"),
}
}
assert_eq!(state.stage, Stage::Ship);
assert_eq!(
state.consecutive_failures, 0,
"an agreeing outcome must never touch the failure counter"
);
}
#[test]
fn external_verify_disagreement_gates_immediately() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 91;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
workflow::save_state(&state).unwrap();
let result = agent_result::AgentResult {
status: AgentStatus::Success,
exit_code: None,
reason: None,
commits: None,
summary: None,
verdict: Some(Verdict::Gaps),
decided_by_layer: Some(0),
};
let outcome = classify_validate_outcome(&result);
assert!(matches!(outcome, ValidateOutcome::Ambiguous(_)));
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_validate_outcome(root, &mut state, outcome).unwrap();
assert_eq!(
state.consecutive_failures, 0,
"an ambiguous outcome must gate on cycle one without touching the counter"
);
assert!(
!Gates::gate_path(root, phase, Stage::Validate).exists(),
"the immediate gate must resolve (and clean up) via the same abort path as any other gate"
);
}
#[test]
fn external_verify_no_verdict_gates_immediately() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 92;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
workflow::save_state(&state).unwrap();
let result = agent_result::AgentResult {
status: AgentStatus::Success,
exit_code: None,
reason: None,
commits: None,
summary: None,
verdict: None,
decided_by_layer: Some(0),
};
let outcome = classify_validate_outcome(&result);
assert!(matches!(outcome, ValidateOutcome::Ambiguous(_)));
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_validate_outcome(root, &mut state, outcome).unwrap();
assert_eq!(
state.consecutive_failures, 0,
"an ambiguous outcome must gate on cycle one without touching the counter"
);
}
#[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!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.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 resource_killed_on_code_bumps_infra_failures_not_consecutive_failures() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 73;
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;
state.consecutive_failures = 1;
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 err = workflow::load_state(root, phase).unwrap_err();
assert!(matches!(err, workflow::WorkflowError::MissingState(_)));
assert!(!Gates::gate_path(root, phase, Stage::Validate).exists());
}
#[test]
fn resource_killed_on_validate_bumps_infra_not_consecutive_failures() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 74;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
state.consecutive_failures = 2;
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: test cleanup","responded_by":"test"}"#,
)
.unwrap();
handle_infra_outcome(
root,
&mut state,
Stage::Validate,
Some("agent process was killed (exit code 137, likely OOM)".into()),
)
.unwrap();
assert_eq!(state.infra_failures, 1);
assert_eq!(
state.consecutive_failures, 2,
"consecutive_failures must be untouched by the infra path"
);
}
#[test]
fn infra_ceiling_aborts_instead_of_gating() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 75;
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();
handle_infra_outcome(root, &mut state, Stage::Code, Some("killed".into())).unwrap();
assert_eq!(state.infra_failures, mode::MAX_INFRA_FAILURES);
assert!(
!Gates::gate_path(root, phase, Stage::Code).exists(),
"at the ceiling, the run must abort rather than gate again"
);
let err = workflow::load_state(root, phase).unwrap_err();
assert!(matches!(err, workflow::WorkflowError::MissingState(_)));
}
#[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 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 consecutive_failures_reaches_ceiling_across_cycles() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 81;
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::Validate);
std::fs::create_dir_all(response_path.parent().unwrap()).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());
}
for _ in 0..mode::MAX_CONSECUTIVE_FAILURES {
std::fs::write(
&response_path,
r#"{"approved":false,"note":"abort: test cleanup","responded_by":"test"}"#,
)
.unwrap();
let _ = handle_validate_outcome(root, &mut state, ValidateOutcome::Failed);
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, mode::MAX_CONSECUTIVE_FAILURES);
assert!(
state
.mode
.should_gate(Stage::Validate, state.consecutive_failures),
"reaching the ceiling must force the Auto-mode Validate gate"
);
assert_eq!(
state.infra_failures, 0,
"infra_failures must still reset unconditionally on the same hop the consecutive reset now skips"
);
}
#[test]
fn external_verify_cycles_reach_ceiling_without_unbounded_loop() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
arm_a_ambiguous_outcome_gates_on_cycle_one(root, 93);
arm_b_genuine_failures_reach_the_ceiling(root, 94);
}
fn arm_a_ambiguous_outcome_gates_on_cycle_one(root: &Path, phase: u32) {
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
workflow::save_state(&state).unwrap();
let result = agent_result::AgentResult {
status: AgentStatus::Success,
exit_code: None,
reason: None,
commits: None,
summary: None,
verdict: Some(Verdict::Gaps),
decided_by_layer: Some(0),
};
let outcome = classify_validate_outcome(&result);
assert!(matches!(outcome, ValidateOutcome::Ambiguous(_)));
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_validate_outcome(root, &mut state, outcome).unwrap();
assert_eq!(
state.consecutive_failures, 0,
"18e's ambiguous gate must fire on cycle one, never touching 18d's counter"
);
}
fn arm_b_genuine_failures_reach_the_ceiling(root: &Path, phase: u32) {
let _guard = ENV_MUTEX.lock().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::Validate);
std::fs::create_dir_all(response_path.parent().unwrap()).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());
}
for _ in 0..mode::MAX_CONSECUTIVE_FAILURES {
std::fs::write(
&response_path,
r#"{"approved":false,"note":"abort: test cleanup","responded_by":"test"}"#,
)
.unwrap();
let _ = handle_validate_outcome(root, &mut state, ValidateOutcome::Failed);
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, mode::MAX_CONSECUTIVE_FAILURES);
assert!(
state
.mode
.should_gate(Stage::Validate, state.consecutive_failures),
"a genuine repeated failure must still reach the reachable ceiling (18d)"
);
}
#[test]
fn consecutive_failures_increment_saturates() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 82;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
state.consecutive_failures = u32::MAX;
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: test cleanup","responded_by":"test"}"#,
)
.unwrap();
handle_validate_outcome(root, &mut state, ValidateOutcome::Failed).unwrap();
assert_eq!(state.consecutive_failures, u32::MAX);
}
#[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"
);
}
#[test]
fn primary_loop_rate_limited_writes_single_agent_cron_instructions() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 76;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
workflow::save_state(&state).unwrap();
std::fs::create_dir_all(root.join(".devflow")).unwrap();
std::fs::write(
agent_result::stdout_path(root, phase),
r#"{"type":"result","subtype":"error_rate_limit","retry_after":"2026-06-18T15:45:30Z"}"#,
)
.unwrap();
advance(root, Some(phase)).unwrap();
let instructions = devflow_core::ship::load_cron_instructions(root, phase).unwrap();
assert_eq!(instructions.resume.command, "devflow");
assert_eq!(
instructions.resume.args,
["resume", "--phase", &phase.to_string()]
);
assert!(
instructions
.hermes_cron
.command
.contains(&format!("devflow resume --phase {phase}"))
);
let reloaded = workflow::load_state(root, phase).unwrap();
assert_eq!(reloaded.stage, Stage::Code);
assert!(!reloaded.gate_pending);
assert_eq!(reloaded.infra_failures, 1);
assert_eq!(reloaded.consecutive_failures, 0);
assert!(!Gates::gate_path(root, phase, Stage::Code).exists());
}
#[test]
fn rate_limited_at_infra_ceiling_stops_resuming_and_aborts() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 77;
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();
std::fs::create_dir_all(root.join(".devflow")).unwrap();
std::fs::write(
agent_result::stdout_path(root, phase),
r#"{"type":"result","subtype":"error_rate_limit","retry_after":"2026-06-18T15:45:30Z"}"#,
)
.unwrap();
advance(root, Some(phase)).unwrap();
let err = workflow::load_state(root, phase).unwrap_err();
assert!(
matches!(err, workflow::WorkflowError::MissingState(_)),
"the infra ceiling must abort, clearing state"
);
assert!(
devflow_core::ship::load_cron_instructions(root, phase).is_err(),
"must not schedule an auto-resume once the infra ceiling stops resumption"
);
}
#[test]
fn rate_limited_with_unparseable_retry_hint_gates_instead_of_stalling_silently() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 81;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
workflow::save_state(&state).unwrap();
std::fs::create_dir_all(root.join(".devflow")).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();
handle_rate_limited_outcome(
root,
&mut state,
phase,
Stage::Code,
Some("rate limited until usage limit".into()),
)
.unwrap();
let events =
std::fs::read_to_string(devflow_core::events::events_path(root)).unwrap_or_default();
assert!(
events.contains("gate_fired"),
"an unparseable retry hint must raise a gate, not stall the phase silently: {events}"
);
assert!(
events.contains("notify_fired"),
"the operator must be notified that a manual resume is needed: {events}"
);
assert!(
!events.contains("rate_limit_resume_scheduled"),
"nothing was scheduled — emitting a resume-scheduled event would be a false signal: {events}"
);
let instructions = devflow_core::ship::load_cron_instructions(root, phase).unwrap();
assert!(instructions.hermes_cron.schedule.is_empty());
}
#[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);
}
#[test]
fn parse_gate_timeout_env_override() {
const SEVEN_DAYS: u64 = 7 * 24 * 60 * 60;
assert_eq!(parse_gate_timeout(Some("42".into())), 42);
assert_eq!(parse_gate_timeout(Some("bad".into())), SEVEN_DAYS);
assert_eq!(parse_gate_timeout(None), SEVEN_DAYS);
}
#[test]
fn phase_artifact_on_develop_detects_context_and_fails_open() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let run = |args: &[&str]| {
let out = std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.expect("spawn git");
assert!(out.status.success(), "git {args:?} failed");
};
run(&["init", "-q", "-b", "main"]);
run(&["config", "user.email", "t@e.st"]);
run(&["config", "user.name", "t"]);
run(&["config", "commit.gpgsign", "false"]);
run(&["config", "core.hooksPath", "/dev/null"]);
std::fs::create_dir_all(root.join(".planning/phases/03-widget")).unwrap();
std::fs::write(root.join(".planning/phases/03-widget/03-CONTEXT.md"), "ctx").unwrap();
run(&["add", "-A"]);
run(&["commit", "-q", "-m", "init"]);
run(&["branch", "develop"]);
assert!(phase_artifact_on_develop(root, 3, "-CONTEXT.md"));
assert!(!phase_artifact_on_develop(root, 3, "-PLAN.md"));
assert!(!phase_artifact_on_develop(root, 4, "-CONTEXT.md"));
let empty = tempfile::tempdir().unwrap();
assert!(phase_artifact_on_develop(empty.path(), 3, "-CONTEXT.md"));
}
#[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!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.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");
}
struct AlwaysFailAdapter;
impl agents::AgentAdapter for AlwaysFailAdapter {
fn name(&self) -> &'static str {
"test-always-fail"
}
fn exec_command(
&self,
_phase: u32,
_prompt: &str,
_roots: &[PathBuf],
) -> (&'static str, Vec<String>) {
("true", Vec::new())
}
fn completion_signal_detected(&self, _output: &str) -> bool {
false
}
fn preflight(&self, _state: &State) -> Result<(), String> {
Err("test adapter always rejects".to_string())
}
}
struct FailOnceAdapter {
failed_once: std::cell::Cell<bool>,
}
impl FailOnceAdapter {
fn new() -> Self {
Self {
failed_once: std::cell::Cell::new(false),
}
}
}
impl agents::AgentAdapter for FailOnceAdapter {
fn name(&self) -> &'static str {
"test-fail-once"
}
fn exec_command(
&self,
_phase: u32,
_prompt: &str,
_roots: &[PathBuf],
) -> (&'static str, Vec<String>) {
("true", Vec::new())
}
fn completion_signal_detected(&self, _output: &str) -> bool {
false
}
fn preflight(&self, _state: &State) -> Result<(), String> {
if self.failed_once.get() {
Ok(())
} else {
self.failed_once.set(true);
Err("test adapter fails on the first preflight call only".to_string())
}
}
}
fn agent_free_git_only_path_dir() -> tempfile::TempDir {
let real_git = std::env::var_os("PATH")
.and_then(|paths| {
std::env::split_paths(&paths).find_map(|dir| {
let candidate = dir.join("git");
candidate.is_file().then_some(candidate)
})
})
.expect("git must be resolvable on PATH to run this test");
let dir = tempfile::tempdir().unwrap();
std::os::unix::fs::symlink(&real_git, dir.path().join("git")).unwrap();
dir
}
fn agent_free_dir_with_agent_stub(program: &str) -> tempfile::TempDir {
use std::os::unix::fs::PermissionsExt;
let dir = agent_free_git_only_path_dir();
let real_sh = std::env::var_os("PATH")
.and_then(|paths| {
std::env::split_paths(&paths).find_map(|d| {
let candidate = d.join("sh");
candidate.is_file().then_some(candidate)
})
})
.expect("sh must be resolvable on PATH to run this test");
std::os::unix::fs::symlink(&real_sh, dir.path().join("sh")).unwrap();
let path = dir.path().join(program);
std::fs::write(&path, "#!/bin/sh\nexit 0\n").unwrap();
let mut perms = std::fs::metadata(&path).unwrap().permissions();
perms.set_mode(0o755);
std::fs::set_permissions(&path, perms).unwrap();
dir
}
fn stub_agent_binary(name: &str) -> tempfile::TempDir {
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join(name);
std::fs::write(&path, "#!/bin/sh\nexit 0\n").unwrap();
let mut perms = std::fs::metadata(&path).unwrap().permissions();
perms.set_mode(0o755);
std::fs::set_permissions(&path, perms).unwrap();
dir
}
fn prepend_path(
stub_dir: &tempfile::TempDir,
original: &Option<std::ffi::OsString>,
) -> std::ffi::OsString {
let mut dirs = vec![stub_dir.path().to_path_buf()];
if let Some(original) = original {
dirs.extend(std::env::split_paths(original));
}
std::env::join_paths(dirs).unwrap()
}
fn stage_launched_count(root: &Path, phase: u32) -> usize {
std::fs::read_to_string(devflow_core::events::events_path(root))
.unwrap_or_default()
.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|event| {
event.get("phase").and_then(serde_json::Value::as_u64) == Some(u64::from(phase))
&& event.get("event").and_then(serde_json::Value::as_str)
== Some("stage_launched")
})
.count()
}
#[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"
);
}
#[test]
fn workflow_started_payload_carries_build_provenance() {
let state = State::new(66, AgentKind::Claude, Mode::Auto, PathBuf::from("/repo"));
let payload = workflow_started_payload(&state);
assert_eq!(payload["agent"], "claude");
assert_eq!(payload["mode"], "auto");
assert!(payload["version"].as_str().is_some());
assert!(payload["commit"].is_string());
assert!(payload["dirty"].is_string());
assert!(
payload.get("build_timestamp").is_none(),
"build_timestamp was removed (CR-02) and must not reappear"
);
assert!(payload["exe_path"].is_string() || payload["exe_path"].is_null());
}
#[test]
fn is_self_dogfood_workspace_matches_both_member_paths_only() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
std::fs::write(
root.join("Cargo.toml"),
"[workspace]\nmembers = [\n \"crates/devflow-core\",\n \"crates/devflow-cli\",\n]\n",
)
.unwrap();
assert!(is_self_dogfood_workspace(root));
let name_only = tempfile::tempdir().unwrap();
std::fs::write(
name_only.path().join("Cargo.toml"),
"[package]\nname = \"devflow-cli\"\n",
)
.unwrap();
assert!(
!is_self_dogfood_workspace(name_only.path()),
"a package NAME match must never fire — the CLI package is named `devflow`"
);
let partial = tempfile::tempdir().unwrap();
std::fs::write(
partial.path().join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\"]\n",
)
.unwrap();
assert!(!is_self_dogfood_workspace(partial.path()));
let missing = tempfile::tempdir().unwrap();
assert!(!is_self_dogfood_workspace(missing.path()));
}
#[test]
fn is_self_dogfood_workspace_requires_exact_member_paths_not_substrings() {
let lookalike = tempfile::tempdir().unwrap();
std::fs::write(
lookalike.path().join("Cargo.toml"),
"[workspace]\nmembers = [\n \"crates/devflow-core-extras\",\n \"crates/devflow-cli-plugin\",\n]\n",
)
.unwrap();
assert!(
!is_self_dogfood_workspace(lookalike.path()),
"`devflow-core-extras`/`devflow-cli-plugin` are not the real members — \
a substring match here would hard-block an unrelated project"
);
let prefixed = tempfile::tempdir().unwrap();
std::fs::write(
prefixed.path().join("Cargo.toml"),
"[workspace]\nmembers = [\n \"vendor/crates/devflow-core\",\n \"vendor/crates/devflow-cli\",\n]\n",
)
.unwrap();
assert!(
!is_self_dogfood_workspace(prefixed.path()),
"vendored copies at a different path are not DevFlow's own workspace"
);
}
#[test]
fn is_self_dogfood_workspace_anchors_on_members_not_default_members() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(
dir.path().join("Cargo.toml"),
"[workspace]\n\
default-members = [\"crates/devflow-cli\"]\n\
members = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
assert!(
is_self_dogfood_workspace(dir.path()),
"a `default-members` key ahead of `members` must not hide the real \
member list — that turns the D-18 hard block into a warning"
);
}
fn worktree_staleness_fixture() -> (tempfile::TempDir, PathBuf, String) {
let outer = tempfile::tempdir().unwrap();
let project_root = outer.path().join("project");
std::fs::create_dir_all(&project_root).unwrap();
let worktree_path = outer.path().join("worktree");
let git = |args: &[&str], cwd: &Path| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(cwd)
.output()
.unwrap()
.status
.success(),
"git {args:?} in {cwd:?} failed"
);
};
git(&["init", "-q", "-b", "develop"], &project_root);
git(&["config", "user.email", "t@e.st"], &project_root);
git(&["config", "user.name", "t"], &project_root);
git(&["config", "commit.gpgsign", "false"], &project_root);
std::fs::create_dir_all(project_root.join("src")).unwrap();
std::fs::write(project_root.join("src/lib.rs"), "// base\n").unwrap();
git(&["add", "."], &project_root);
git(&["commit", "-q", "-m", "base"], &project_root);
let embedded_commit = run_git_stdout(&project_root, &["rev-parse", "HEAD"])
.expect("rev-parse HEAD")
.trim()
.to_string();
git(
&[
"worktree",
"add",
"-b",
"feature/phase-90",
worktree_path.to_str().unwrap(),
"develop",
],
&project_root,
);
std::fs::write(worktree_path.join("src/lib.rs"), "// wt commit 1\n").unwrap();
git(&["add", "."], &worktree_path);
git(&["commit", "-q", "-m", "wt commit 1"], &worktree_path);
std::fs::write(worktree_path.join("src/lib.rs"), "// wt commit 2\n").unwrap();
git(&["add", "."], &worktree_path);
git(&["commit", "-q", "-m", "wt commit 2"], &worktree_path);
(outer, worktree_path, embedded_commit)
}
#[test]
fn embedded_commit_is_stale_uses_worktree_head() {
let _guard = ENV_MUTEX.lock().unwrap();
let (outer, worktree_path, embedded_commit) = worktree_staleness_fixture();
let project_root = outer.path().join("project");
assert_eq!(
embedded_commit_is_stale(&project_root, &embedded_commit),
Staleness::Fresh,
"project_root's HEAD never moved, so the embedded commit is still an exact match"
);
assert_eq!(
embedded_commit_is_stale(&worktree_path, &embedded_commit),
Staleness::Stale,
"the worktree branch advanced two commits past the embedded commit — Round 4 \
CR-01's mechanism: evaluated against the wrong tree, this same commit reads Fresh"
);
}
#[test]
fn enforce_build_staleness_blocks_self_dogfood_behind_worktree_head() {
let _guard = ENV_MUTEX.lock().unwrap();
let (outer, worktree_path, embedded_commit) = worktree_staleness_fixture();
let project_root = outer.path().join("project");
std::fs::write(
project_root.join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
assert!(is_self_dogfood_workspace(&project_root));
let phase = 90;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, project_root.clone());
state.stage = Stage::Code;
state.worktree_path = Some(worktree_path.clone());
let err =
enforce_build_staleness(&project_root, &state, &embedded_commit, false).unwrap_err();
let message = err.to_string();
assert!(
message.contains(&worktree_path.display().to_string()),
"block message must name the worktree that was actually evaluated: {message}"
);
assert!(
!message.contains(&project_root.display().to_string()),
"block message must not name project_root when a worktree was evaluated: {message}"
);
let last = devflow_core::events::last_event_for_phase(&project_root, phase)
.expect("staleness block must record an event before returning the error");
assert_eq!(last["reason"], "stale_build_blocked");
assert_eq!(last["worktree"], true);
}
#[test]
fn staleness_without_worktree_is_unchanged() {
let _guard = ENV_MUTEX.lock().unwrap();
let (outer, _worktree_path, embedded_commit) = worktree_staleness_fixture();
let project_root = outer.path().join("project");
std::fs::write(
project_root.join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
let phase = 91;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, project_root.clone());
state.stage = Stage::Code;
assert!(
state.worktree_path.is_none(),
"fixture precondition: no worktree recorded on this state"
);
assert!(
enforce_build_staleness(&project_root, &state, &embedded_commit, false).is_ok(),
"no worktree recorded must fall back to project_root, which the fixture never \
advances past embedded_commit"
);
}
fn init_repo_with_diverged_commit(root: &Path) -> (String, String) {
let git = |args: &[&str]| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
let rev_parse = || {
let out = std::process::Command::new("git")
.args(["rev-parse", "HEAD"])
.current_dir(root)
.output()
.unwrap();
String::from_utf8_lossy(&out.stdout).trim().to_string()
};
git(&["init", "-q", "-b", "trunk"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
std::fs::write(root.join("a.txt"), "one").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "base"]);
let base = rev_parse();
git(&["checkout", "-q", "-b", "side"]);
std::fs::write(root.join("side.txt"), "s").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "side"]);
let side = rev_parse();
git(&["checkout", "-q", "trunk"]);
std::fs::write(root.join("trunk2.txt"), "t2").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "trunk2"]);
(base, side)
}
#[test]
fn embedded_commit_is_stale_maps_ancestry_exit_codes() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let (base, side) = init_repo_with_diverged_commit(root);
let head = run_git_stdout(root, &["rev-parse", "HEAD"])
.expect("rev-parse HEAD")
.trim()
.to_string();
assert_eq!(embedded_commit_is_stale(root, &base), Staleness::Stale);
assert_eq!(embedded_commit_is_stale(root, &head), Staleness::Fresh);
assert_eq!(embedded_commit_is_stale(root, &side), Staleness::Stale);
assert_eq!(embedded_commit_is_stale(root, ""), Staleness::Indeterminate);
assert_eq!(
embedded_commit_is_stale(root, "deadbeefdeadbeefdeadbeefdeadbeefdeadbeef"),
Staleness::Indeterminate
);
}
#[test]
fn wr01_clean_tree_strict_ancestor_build_is_stale_and_hard_blocks() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let git = |args: &[&str]| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git(&["init", "-q"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
std::fs::write(
root.join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
std::fs::write(root.join("a.txt"), "one").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "workspace init"]);
let embedded_commit = run_git_stdout(root, &["rev-parse", "HEAD"])
.expect("rev-parse HEAD")
.trim()
.to_string();
std::fs::write(root.join("b.txt"), "two").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "unrelated follow-up"]);
let status = run_git_stdout(root, &["status", "--porcelain"]).unwrap();
assert!(
status.trim().is_empty(),
"fixture must have a clean working tree"
);
assert_eq!(
embedded_commit_is_stale(root, &embedded_commit),
Staleness::Stale
);
assert_eq!(
combined_staleness(root, &embedded_commit, false),
Staleness::Stale
);
assert!(is_self_dogfood_workspace(root));
let phase = 66;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
let err = enforce_build_staleness(root, &state, &embedded_commit, false).unwrap_err();
assert!(
err.to_string().contains("self-dogfood stale build blocked"),
"{err}"
);
let last = devflow_core::events::last_event_for_phase(root, phase)
.expect("staleness block must record an event before returning the error");
assert_eq!(last["event"], "self_dogfood_stale_blocked");
}
#[test]
fn ahead_build_from_descendant_commit_warns_instead_of_blocking() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let git = |args: &[&str]| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git(&["init", "-q"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
std::fs::write(
root.join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
std::fs::write(root.join("a.txt"), "one").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "workspace init"]);
let base_commit = run_git_stdout(root, &["rev-parse", "HEAD"])
.expect("rev-parse HEAD")
.trim()
.to_string();
std::fs::write(root.join("b.txt"), "two").unwrap();
git(&["add", "."]);
git(&[
"commit",
"-q",
"-m",
"newer work the checkout does not have",
]);
let embedded_commit = run_git_stdout(root, &["rev-parse", "HEAD"])
.expect("rev-parse HEAD")
.trim()
.to_string();
git(&["reset", "--hard", "-q", &base_commit]);
let status = run_git_stdout(root, &["status", "--porcelain"]).unwrap();
assert!(
status.trim().is_empty(),
"fixture must have a clean working tree"
);
assert_eq!(
embedded_commit_is_stale(root, &embedded_commit),
Staleness::Ahead,
"a descendant embedded commit is newer than HEAD, not stale"
);
assert_eq!(
staleness_outcome(true, Staleness::Ahead),
StalenessOutcome::Warn,
"an ahead build must warn, never hard-block, even for self-dogfood"
);
let phase = 67;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
assert!(
enforce_build_staleness(root, &state, &embedded_commit, false).is_ok(),
"ahead build must not block a self-dogfood workspace"
);
}
#[test]
fn dirty_flag_arm_ignores_non_build_files_but_still_flags_sources() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let git = |args: &[&str]| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git(&["init", "-q"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
std::fs::write(
root.join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
std::fs::write(root.join("CHANGELOG.md"), "# Changelog\n").unwrap();
std::fs::create_dir_all(root.join("crates/devflow-cli/src")).unwrap();
std::fs::write(
root.join("crates/devflow-cli/src/main.rs"),
"fn main() {}\n",
)
.unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "workspace init"]);
let embedded_commit = run_git_stdout(root, &["rev-parse", "HEAD"])
.expect("rev-parse HEAD")
.trim()
.to_string();
let build_dirty = false;
std::fs::write(root.join("CHANGELOG.md"), "# Changelog\n\n## 1.4.26\n").unwrap();
assert_eq!(
run_git_stdout(root, &["ls-files", "-m"]).unwrap().trim(),
"CHANGELOG.md",
"fixture must have exactly one dirty tracked file"
);
assert_eq!(
tree_has_modified_build_inputs(root),
Some(false),
"a dirty CHANGELOG.md cannot change the compiled binary"
);
assert_eq!(
combined_staleness(root, &embedded_commit, build_dirty),
Staleness::Fresh,
"a doc-only dirty tree must not be Stale"
);
let phase = 68;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
assert!(
enforce_build_staleness(root, &state, &embedded_commit, build_dirty).is_ok(),
"a doc-only dirty tree must not block Ship"
);
std::fs::write(
root.join("crates/devflow-cli/src/main.rs"),
"fn main() { /* edited after build */ }\n",
)
.unwrap();
assert_eq!(
tree_has_modified_build_inputs(root),
Some(true),
"a modified .rs file is genuine staleness input"
);
git(&["add", "crates/devflow-cli/src/main.rs"]);
assert!(
!run_git_stdout(root, &["ls-files", "-m"])
.unwrap()
.lines()
.any(|line| line.ends_with(".rs")),
"fixture precondition: `ls-files -m` is blind to the staged .rs edit"
);
assert_eq!(
tree_has_modified_build_inputs(root),
Some(true),
"a STAGED source edit is just as much a staleness input as an unstaged one"
);
assert_eq!(
combined_staleness(root, &embedded_commit, build_dirty),
Staleness::Stale,
"a staged, uncommitted source edit on a clean build is Stale"
);
git(&["reset", "-q"]);
assert_eq!(
combined_staleness(root, &embedded_commit, build_dirty),
Staleness::Stale
);
assert!(
enforce_build_staleness(root, &state, &embedded_commit, build_dirty).is_err(),
"a stale source build must still hard-block a self-dogfood workspace"
);
}
#[test]
fn content_hooks_target_the_worktree_while_terminal_hooks_stay_on_project_root() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let worktree = root.join(".worktrees/phase-70");
std::fs::create_dir_all(&worktree).unwrap();
let mut state = State::new(70, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.worktree_path = Some(worktree.clone());
assert_eq!(
hook_context_root(root, &state, false),
worktree,
"content hooks must write into the phase's worktree"
);
assert_eq!(
hook_context_root(root, &state, true),
root.to_path_buf(),
"terminal hooks merge/tag/delete against the primary checkout"
);
let mut no_worktree = state.clone();
no_worktree.worktree_path = None;
assert_eq!(hook_context_root(root, &no_worktree, false), root);
let mut missing = state.clone();
missing.worktree_path = Some(root.join(".worktrees/gone"));
assert_eq!(hook_context_root(root, &missing, false), root);
}
#[test]
fn combined_staleness_dirty_flag_arm_flags_modified_tree_when_build_was_clean() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let git = |args: &[&str]| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git(&["init", "-q"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
std::fs::create_dir_all(root.join("src")).unwrap();
std::fs::write(root.join("src/lib.rs"), "// one\n").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "init"]);
let head = {
let out = std::process::Command::new("git")
.args(["rev-parse", "HEAD"])
.current_dir(root)
.output()
.unwrap();
String::from_utf8_lossy(&out.stdout).trim().to_string()
};
assert_eq!(embedded_commit_is_stale(root, &head), Staleness::Fresh);
assert_eq!(combined_staleness(root, &head, false), Staleness::Fresh);
std::fs::write(root.join("src/lib.rs"), "// modified after build\n").unwrap();
assert_eq!(combined_staleness(root, &head, false), Staleness::Stale);
}
#[test]
fn combined_staleness_dirty_flag_arm_is_indeterminate_when_build_was_already_dirty() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let git = |args: &[&str]| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git(&["init", "-q"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
std::fs::write(
root.join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
std::fs::create_dir_all(root.join("src")).unwrap();
std::fs::write(root.join("src/lib.rs"), "// one\n").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "init"]);
let head = run_git_stdout(root, &["rev-parse", "HEAD"])
.expect("rev-parse HEAD")
.trim()
.to_string();
assert!(is_self_dogfood_workspace(root));
std::fs::write(root.join("src/lib.rs"), "// modified\n").unwrap();
assert_eq!(embedded_commit_is_stale(root, &head), Staleness::Fresh);
assert_eq!(
tree_has_modified_build_inputs(root),
Some(true),
"fixture must have a dirty, build-affecting tree"
);
let build_was_dirty = true;
assert_eq!(
combined_staleness(root, &head, build_was_dirty),
Staleness::Indeterminate,
"cannot distinguish \"same dirt\" from \"more dirt\" without a timestamp"
);
let phase = 71;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
assert!(
enforce_build_staleness(root, &state, &head, build_was_dirty).is_ok(),
"Indeterminate must never hard-block, even for a self-dogfood workspace (Pitfall 4)"
);
}
#[test]
fn enforce_build_staleness_blocks_self_dogfood_and_records_event_before_erroring() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let (_base, side) = init_repo_with_diverged_commit(root);
let git = |args: &[&str]| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
std::fs::write(
root.join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "add workspace cargo toml"]);
assert!(is_self_dogfood_workspace(root));
let phase = 63;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
let err = enforce_build_staleness(root, &state, &side, false).unwrap_err();
let message = err.to_string();
assert!(
message.contains("self-dogfood stale build blocked"),
"{message}"
);
assert!(
message.contains(&root.display().to_string()),
"the returned CliError (terminal-only) must still name the path: {message}"
);
let last = devflow_core::events::last_event_for_phase(root, phase)
.expect("staleness block must record an event before returning the error");
assert_eq!(last["event"], "self_dogfood_stale_blocked");
assert_eq!(last["reason"], "stale_build_blocked");
assert_eq!(last["worktree"], false);
let reason_str = last["reason"].as_str().unwrap();
assert!(
!reason_str.contains(&root.display().to_string()),
"persisted reason must never carry the project root path: {reason_str}"
);
}
#[test]
fn enforce_build_staleness_warns_for_ordinary_project_with_stale_commit() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let (_base, side) = init_repo_with_diverged_commit(root);
assert!(!is_self_dogfood_workspace(root));
let phase = 64;
let state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
let result = enforce_build_staleness(root, &state, &side, false);
assert!(
result.is_ok(),
"an ordinary project's stale build must only warn, never block"
);
assert!(
devflow_core::events::last_event_for_phase(root, phase).is_none(),
"a warn-only path must not fire the self_dogfood_stale_blocked event"
);
}
#[test]
fn enforce_build_staleness_never_blocks_on_indeterminate() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let git = |args: &[&str]| {
assert!(
std::process::Command::new("git")
.args(args)
.current_dir(root)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git(&["init", "-q"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
std::fs::write(
root.join("Cargo.toml"),
"[workspace]\nmembers = [\"crates/devflow-core\", \"crates/devflow-cli\"]\n",
)
.unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "init"]);
assert!(is_self_dogfood_workspace(root));
let phase = 65;
let state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
let result = enforce_build_staleness(
root,
&state,
"deadbeefdeadbeefdeadbeefdeadbeefdeadbeef",
false,
);
assert!(
result.is_ok(),
"an Indeterminate result must never hard-block"
);
}
#[test]
fn truncate_reason_caps_long_reasons_and_keeps_short_ones() {
assert_eq!(truncate_reason("short reason"), "short reason");
let long = "x".repeat(5000);
let capped = truncate_reason(&long);
assert!(capped.chars().count() <= 300);
assert!(capped.ends_with("[truncated; full output in .devflow/]"));
}
#[test]
fn gate_context_rendering_neutralizes_all_controls_and_obeys_limit() {
let rendered = render_gate_context("line 1\n\u{1b}[2J\tline 2\u{7}", 100);
assert!(!rendered.chars().any(char::is_control));
assert_eq!(rendered, "line 1 [2J line 2 ");
let bounded = render_gate_context(&"x".repeat(500), 100);
assert_eq!(bounded.chars().count(), 100);
assert!(bounded.ends_with("[truncated; full output in .devflow/]"));
}
#[test]
fn status_shows_pending_gate_prominently() {
let dir = tempfile::tempdir().unwrap();
let context = format!("first line\n\u{1b}[2J{}", "sensitive detail ".repeat(80));
Gates::write_gate(dir.path(), 16, Stage::Ship, &context).unwrap();
let open = Gates::list_open(dir.path());
let banner = render_pending_gate_banner(&open, u64::MAX).unwrap();
assert!(banner.contains("PENDING GATE"));
assert!(banner.contains("phase 16"));
assert!(banner.contains("ship"));
assert!(banner.contains("devflow gate approve 16 --stage ship"));
assert!(banner.contains("devflow gate reject 16 --stage ship"));
assert!(banner.contains("[truncated; full output in .devflow/]"));
assert!(!banner.contains(&context));
assert!(!banner.contains('\u{1b}'));
assert!(banner.contains("ESCALATED"));
}
#[test]
fn ship_agent_failed_fires_gate() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 40;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
let gate_path = Gates::gate_path(root, phase, Stage::Ship);
let response_path = Gates::response_path(root, phase, Stage::Ship);
std::thread::scope(|scope| {
scope.spawn(|| {
handle_ship_failure(root, &mut state, Some("agent crashed".into())).unwrap();
});
let mut seen = false;
for _ in 0..150 {
if gate_path.exists() {
seen = true;
break;
}
std::thread::sleep(std::time::Duration::from_millis(20));
}
assert!(
seen,
"handle_ship_failure must write a gate file, not silently return an Err"
);
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 ship_review_failed_loops_to_code() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 41;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Ship;
workflow::save_state(&state).unwrap();
let reason = Some("review: please fix naming".to_string());
assert!(is_ship_review_failure(&reason));
prepare_loop_back_to_code(root, &mut state, FixType::AuditFix).unwrap();
assert_eq!(state.stage, Stage::Code);
assert!(!Gates::gate_path(root, phase, Stage::Ship).exists());
assert!(workflow::load_state(root, phase).is_ok());
}
#[test]
fn ship_review_failed_uses_audit_fix() {
assert!(is_ship_review_failure(&Some(
"review: needs changes".into()
)));
assert!(is_ship_review_failure(&Some(" Review: nitpick".into())));
assert!(!is_ship_review_failure(&Some("agent crashed".into())));
assert!(!is_ship_review_failure(&None));
let prompt = prompt::fix_prompt(FixType::AuditFix, 11);
assert!(prompt.contains("/gsd-audit-fix"));
assert!(!prompt.contains("--gaps-only"));
}
#[test]
fn non_validate_failure_fires_gate_and_hook() {
let _guard = ENV_MUTEX.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let sentinel = root.join("notify-sentinel");
unsafe {
std::env::set_var(
"DEVFLOW_GATE_NOTIFY_CMD",
format!("touch {}", sentinel.display()),
);
}
let phase = 42;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Code;
workflow::save_state(&state).unwrap();
assert!(
!state
.mode
.should_gate(Stage::Code, state.consecutive_failures)
);
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();
let result =
handle_stage_failure(root, &mut state, Stage::Code, Some("build failed".into()));
unsafe {
std::env::remove_var("DEVFLOW_GATE_NOTIFY_CMD");
}
result.unwrap();
assert!(
sentinel.exists(),
"handle_stage_failure must fire the configured notify hook, not silently skip it"
);
}
#[test]
fn stage_failure_retry_cleans_stale_response() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 43;
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();
handle_stage_failure(root, &mut state, Stage::Code, Some("first failure".into())).unwrap();
assert!(!Gates::gate_path(root, phase, Stage::Code).exists());
assert!(!Gates::response_path(root, phase, Stage::Code).exists());
assert!(!Gates::ack_path(root, phase, Stage::Code).exists());
Gates::write_gate(root, phase, Stage::Code, "re-fired gate").unwrap();
let started = std::time::Instant::now();
let got = Gates::poll_response(root, phase, Stage::Code, 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));
}
#[cfg(test)]
mod doctor_reconciliation {
use super::*;
fn agreeing_facts(phase: u32) -> PhaseFacts {
PhaseFacts {
phase,
stage: Stage::Code,
gate_pending: false,
agent_pid: Some(4242),
agent_alive: true,
monitor_pid: Some(4343),
monitor_alive: true,
last_event: Some("stage_launched".into()),
last_launched_stage: Some(Stage::Code),
open_gate_stages: Vec::new(),
feature_branch_exists: true,
}
}
#[test]
fn reconcile_phase_returns_no_findings_when_all_agree() {
let facts = agreeing_facts(1);
assert!(reconcile_phase(&facts).is_empty());
}
#[test]
fn reconcile_phase_flags_gate_pending_without_open_gate() {
let facts = PhaseFacts {
gate_pending: true,
open_gate_stages: Vec::new(),
..agreeing_facts(2)
};
let findings = reconcile_phase(&facts);
assert_eq!(findings.len(), 1);
assert_eq!(findings[0].severity, Severity::Problem);
assert!(findings[0].detail.contains("gate_pending is true"));
assert_eq!(
findings[0].repair.as_deref(),
Some("devflow resume --phase 2")
);
}
#[test]
fn reconcile_phase_flags_orphan_open_gate() {
let facts = PhaseFacts {
gate_pending: false,
open_gate_stages: vec![Stage::Validate],
..agreeing_facts(3)
};
let findings = reconcile_phase(&facts);
assert_eq!(findings.len(), 1);
assert_eq!(findings[0].severity, Severity::Problem);
assert!(findings[0].detail.contains("gate open for stage validate"));
assert_eq!(
findings[0].repair.as_deref(),
Some("devflow gate approve 3 --stage validate")
);
}
#[test]
fn reconcile_phase_flags_dead_agent_at_agent_stage() {
let facts = PhaseFacts {
agent_pid: Some(999_999),
agent_alive: false,
..agreeing_facts(4)
};
let findings = reconcile_phase(&facts);
assert_eq!(findings.len(), 1);
assert_eq!(findings[0].severity, Severity::Problem);
assert!(findings[0].detail.contains("agent pid 999999"));
assert_eq!(
findings[0].repair.as_deref(),
Some("devflow resume --phase 4")
);
}
#[test]
fn reconcile_phase_flags_stage_event_drift() {
let facts = PhaseFacts {
stage: Stage::Validate,
last_launched_stage: Some(Stage::Code),
..agreeing_facts(5)
};
let findings = reconcile_phase(&facts);
assert_eq!(findings.len(), 1);
assert_eq!(findings[0].severity, Severity::Warn);
assert!(
findings[0]
.detail
.contains("last stage_launched event named code")
);
assert!(findings[0].repair.is_none());
}
#[test]
fn reconcile_phase_flags_missing_feature_branch() {
let facts = PhaseFacts {
stage: Stage::Plan,
last_launched_stage: Some(Stage::Plan),
feature_branch_exists: false,
..agreeing_facts(6)
};
let findings = reconcile_phase(&facts);
assert_eq!(findings.len(), 1);
assert_eq!(findings[0].severity, Severity::Warn);
assert!(findings[0].detail.contains("feature/phase-06"));
assert!(findings[0].repair.is_none());
}
#[test]
fn reconcile_reports_stuck_when_monitor_and_agent_are_both_dead() {
let facts = PhaseFacts {
monitor_pid: Some(5150),
monitor_alive: false,
agent_pid: Some(4242),
agent_alive: false,
..agreeing_facts(8)
};
let findings = reconcile_phase(&facts);
let monitor_finding = findings
.iter()
.find(|f| f.detail.contains("monitor pid"))
.expect("expected a monitor finding when monitor and agent are both dead");
assert_eq!(monitor_finding.severity, Severity::Problem);
assert!(monitor_finding.detail.contains("monitor pid 5150"));
assert_eq!(
monitor_finding.repair.as_deref(),
Some("devflow resume --phase 8")
);
}
#[test]
fn reconcile_is_silent_when_monitor_pid_is_unrecorded() {
let facts = PhaseFacts {
monitor_pid: None,
monitor_alive: false,
..agreeing_facts(9)
};
assert!(
reconcile_phase(&facts).is_empty(),
"an unrecorded monitor must never produce a finding"
);
}
#[test]
fn reconcile_is_silent_when_monitor_alive_and_agent_dead() {
let facts = PhaseFacts {
monitor_pid: Some(5150),
monitor_alive: true,
agent_alive: false,
..agreeing_facts(10)
};
let findings = reconcile_phase(&facts);
assert!(
findings.iter().all(|f| !f.detail.contains("monitor pid")),
"a live monitor with a dead agent must not produce a monitor finding"
);
}
#[test]
fn reconcile_phase_ordering_is_input_order_independent() {
let facts = PhaseFacts {
gate_pending: true,
agent_pid: Some(999_999),
agent_alive: false,
monitor_pid: Some(999_998),
monitor_alive: false,
last_launched_stage: Some(Stage::Validate),
open_gate_stages: Vec::new(),
feature_branch_exists: false,
..agreeing_facts(7)
};
let findings = reconcile_phase(&facts);
let severities: Vec<Severity> = findings.iter().map(|f| f.severity).collect();
assert_eq!(
severities,
vec![
Severity::Problem, Severity::Problem, Severity::Problem, Severity::Warn, Severity::Warn, ]
);
assert!(findings[0].detail.contains("gate_pending is true"));
assert!(findings[1].detail.contains("agent pid 999999"));
assert!(findings[2].detail.contains("monitor pid 999998"));
assert!(
findings[3]
.detail
.contains("last stage_launched event named validate")
);
assert!(findings[4].detail.contains("feature/phase-07"));
}
#[test]
fn doctor_reports_no_active_phases_when_idle() {
let dir = tempfile::tempdir().unwrap();
let facts = collect_phase_facts(dir.path());
assert!(facts.is_empty());
assert!(render_reconciliation_text(&facts).contains("no active phases"));
}
#[test]
fn doctor_reports_gate_pending_without_gate_file() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 90;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
state.gate_pending = true;
workflow::save_state(&state).unwrap();
let facts = collect_phase_facts(root);
assert_eq!(facts.len(), 1);
let text = render_reconciliation_text(&facts);
assert!(text.contains(&format!("phase {phase}: gate_pending is true")));
assert!(text.contains(&format!("repair: devflow resume --phase {phase}")));
}
#[test]
fn doctor_json_is_a_single_object_with_environment_and_reconciliation() {
let checks = vec![Check {
name: "git".into(),
status: "ok".into(),
version: Some("2.40.0".into()),
install_hint: None,
}];
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 92;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
state.gate_pending = true; workflow::save_state(&state).unwrap();
let facts = collect_phase_facts(root);
let body = doctor_json_body(&checks, &facts);
let serialized = serde_json::to_string(&body).unwrap();
let reparsed: serde_json::Value = serde_json::from_str(&serialized)
.expect("doctor --json must be single-document JSON, not two concatenated arrays");
assert!(
reparsed.get("environment").is_some(),
"must carry the tool checks under \"environment\": {reparsed}"
);
assert!(
reparsed.get("reconciliation").is_some(),
"must carry the reconciliation findings under \"reconciliation\": {reparsed}"
);
assert!(reparsed["environment"].is_array());
assert!(reparsed["reconciliation"].is_array());
let reconciliation = reparsed["reconciliation"].as_array().unwrap();
assert!(
!reconciliation.is_empty(),
"the mismatched gate_pending fixture must produce at least one finding"
);
assert!(
reconciliation.iter().any(|f| f["detail"]
.as_str()
.unwrap_or("")
.contains("gate_pending is true")),
"must carry the gate_pending finding: {reconciliation:?}"
);
}
#[test]
fn doctor_is_read_only_on_a_mismatched_project() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = 91;
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.stage = Stage::Validate;
state.gate_pending = true; workflow::save_state(&state).unwrap();
events::emit(
root,
phase,
"stage_launched",
serde_json::json!({"stage": "code"}),
);
let state_path = workflow::state_path(root, phase);
let before_len = std::fs::metadata(&state_path).unwrap().len();
let before_modified = std::fs::metadata(&state_path).unwrap().modified().unwrap();
let events_log = events::events_path(root);
let before_lines = std::fs::read_to_string(&events_log)
.unwrap()
.lines()
.count();
doctor(root, false).unwrap();
doctor(root, false).unwrap();
let after_len = std::fs::metadata(&state_path).unwrap().len();
let after_modified = std::fs::metadata(&state_path).unwrap().modified().unwrap();
let after_lines = std::fs::read_to_string(&events_log)
.unwrap()
.lines()
.count();
assert_eq!(
before_len, after_len,
"doctor must not rewrite the state file"
);
assert_eq!(
before_modified, after_modified,
"doctor must not touch the state file's mtime"
);
assert_eq!(
before_lines, after_lines,
"doctor must not append to events.jsonl"
);
}
}
}