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, recover, worktree};
use devflow_core::{
agent_result::{AgentStatus, Verdict},
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>,
},
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::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",
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()),
}),
);
if let Err(err) = launch_stage(&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 launch_stage(
state: &State,
prompt_override: Option<String>,
archived_stage: Option<Stage>,
) -> Result<(), CliError> {
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)?;
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}")))?;
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 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": format!("{:?}", result.status).to_ascii_lowercase(),
"verdict": result.verdict.map(|v| format!("{v:?}").to_ascii_lowercase()),
"reason": result.reason.as_deref().map(truncate_reason),
}),
);
let failed = matches!(
result.status,
AgentStatus::Failed | AgentStatus::RateLimited
);
if failed {
return match stage {
Stage::Validate => handle_validate_outcome(project_root, &mut state, false),
Stage::Ship => handle_ship_failure(project_root, &mut state, result.reason),
_ => handle_stage_failure(project_root, &mut state, stage, result.reason),
};
}
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 => {
let passed = matches!(result.verdict, Some(Verdict::Pass));
handle_validate_outcome(project_root, &mut state, passed)
}
Stage::Ship => handle_ship_outcome(project_root, &mut state),
}
}
fn handle_validate_outcome(
project_root: &Path,
state: &mut State,
passed: bool,
) -> Result<(), CliError> {
if !passed {
state.consecutive_failures += 1;
workflow::save_state(state)?;
}
if state
.mode
.should_gate(Stage::Validate, state.consecutive_failures)
{
let context = if passed {
"Validation passed — approve to ship?".to_string()
} else {
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),
};
}
if passed {
transition(project_root, state, Stage::Ship)
} else {
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 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();
for hook in batch {
let ctx = HookContext {
phase: state.phase,
project_root: project_root.to_path_buf(),
stage,
git_flow: git_flow.clone(),
};
let outcome = hook.run(&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;
state.consecutive_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,
}));
}
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(())
}
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());
match agent_pid_from_file(project_root, state.phase) {
Some(pid) => {
println!(
" agent_pid: {pid} (running: {})",
agent::agent_running(pid)
);
}
None => println!(" agent_pid: none"),
}
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 -- -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(", ")
)))
}
}
fn doctor(_project_root: &Path, json: bool) -> Result<(), CliError> {
use std::process::Command;
struct Check {
name: String,
status: String,
version: Option<String>,
install_hint: Option<String>,
}
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,
},
];
if json {
let mut out = String::from("[\n");
for (i, c) in checks.iter().enumerate() {
out.push_str(" {\n");
out.push_str(&format!(" \"name\": {:?},\n", c.name));
out.push_str(&format!(" \"status\": {:?},\n", c.status));
if let Some(v) = &c.version {
out.push_str(&format!(" \"version\": {:?},\n", v));
} else {
out.push_str(" \"version\": null,\n");
}
if let Some(h) = &c.install_hint {
out.push_str(&format!(" \"install_hint\": {:?}\n", h));
} else {
out.push_str(" \"install_hint\": null\n");
}
out.push('}');
if i + 1 < checks.len() {
out.push(',');
}
out.push('\n');
}
out.push_str("]\n");
print!("{out}");
} 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!();
}
}
Ok(())
}
#[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_for_transition(Stage::Validate, Stage::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"
);
}
#[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 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 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();
}
std::thread::scope(|scope| {
let handles: Vec<_> = phases
.iter()
.map(|&phase| scope.spawn(move || advance(root, Some(phase))))
.collect();
for handle in handles {
handle.join().expect("advance thread").expect("advance ok");
}
});
for &phase in &phases {
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"
);
}
}
#[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, false).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, false).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 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 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));
}
}