use crate::CliError;
use crate::config_parse;
use crate::config_parse::GATE_ESCALATION_THRESHOLD_SECS;
use crate::parallel::ensure_phase_worktree;
use crate::pipeline_gate::print_dry_run;
use crate::pipeline_launch::{launch_stage, single_active_phase};
use crate::pipeline_outcomes::render_gate_context;
use crate::preflight::{
agent_program, ensure_agent_binary, ensure_base_ref_current, ensure_phase_reachable_on_base,
};
use crate::staleness::{enforce_build_staleness, run_git_stdout};
use devflow_core::agent;
use devflow_core::agent_result;
use devflow_core::agents;
use devflow_core::config::{self, DEVELOP, FEATURE_PREFIX, MAIN};
use devflow_core::events;
use devflow_core::gates::{GateAction, GateError, GateResponse, Gates, OpenGate};
use devflow_core::git::{GitFlow, git_command, hermetic_command};
use devflow_core::history;
use devflow_core::lock;
use devflow_core::mode::Mode;
use devflow_core::phase_id::PhaseId;
use devflow_core::recover;
use devflow_core::registry;
use devflow_core::ship_evidence;
use devflow_core::stage::Stage;
use devflow_core::state::{AgentKind, State};
use devflow_core::version;
use devflow_core::workflow;
use devflow_core::worktree;
use std::collections::HashSet;
use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
pub(crate) 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)]
pub(crate) fn phase_artifact_on_develop(project_root: &Path, phase: PhaseId, suffix: &str) -> bool {
let prefix = format!(".planning/phases/{padded}-", padded = phase.padded());
let output = git_command(project_root)
.args([
"ls-tree",
"-r",
"--name-only",
"develop",
"--",
".planning/phases/",
])
.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))
})
}
pub(crate) fn fresh_state_carrying_phase_failures(
project_root: &Path,
phase: PhaseId,
agent: AgentKind,
mode: Mode,
) -> State {
let mut state = State::new(phase, agent, mode, project_root.to_path_buf());
let (carried, warning) =
carried_phase_failures(phase, workflow::load_state(project_root, phase));
state.phase_validate_failures = carried;
if let Some(warning) = warning {
println!("{warning}");
}
state
}
fn carried_phase_failures(
phase: PhaseId,
loaded: Result<State, workflow::WorkflowError>,
) -> (u32, Option<String>) {
match loaded {
Ok(persisted) => (persisted.phase_validate_failures, None),
Err(workflow::WorkflowError::MissingState(_)) => (0, None),
Err(err) => (
0,
Some(format!(
"warning: phase {phase} state could not be read ({err}) — the per-phase \
Validate-failure budget restarts at zero"
)),
),
}
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn start(
project_root: &Path,
phase: PhaseId,
agent: AgentKind,
mode: Mode,
force: bool,
worktree: bool,
dry_run: bool,
until: Option<Stage>,
yes_ship: bool,
legacy_claude_launch: bool,
) -> Result<(), CliError> {
let mut state = fresh_state_carrying_phase_failures(project_root, phase, agent, mode);
state.stop_until = until;
let config_yes_ship = config::yes_ship(project_root);
if config_yes_ship && !yes_ship {
println!(
"note: Ship gate pre-authorized by devflow.toml (yes_ship = true) — see D-12, 28-CONTEXT.md"
);
}
state.yes_ship = yes_ship || config_yes_ship;
if crate::pipeline_launch::apply_legacy_launch_opt_out(&mut state, legacy_claude_launch) {
println!(
"note: legacy Claude launch forced by DEVFLOW_CLAUDE_LEGACY_LAUNCH \
(D-11, 31-CONTEXT.md) — a persisted default is never a silent one"
);
}
if dry_run {
print_dry_run(&state);
return Ok(());
}
ensure_agent_binary(agent_program(agent))?;
ensure_base_ref_current(project_root, DEVELOP)?;
ensure_phase_reachable_on_base(project_root, phase, DEVELOP)?;
let driver = agents::driver_for(agent);
if driver.interactivity_mode(Stage::Define)
== agents::InteractivityMode::RequiresExistingArtifact
&& !phase_artifact_on_develop(project_root, phase, "-CONTEXT.md")
{
return Err(CliError::Message(format!(
"phase {phase} has no CONTEXT.md on develop, and {} cannot run an \
interactive discussion headless. Run /gsd-discuss-phase {phase} \
interactively first (any agent), or use --agent claude.",
driver.name()
)));
}
if driver.interactivity_mode(Stage::Plan) == agents::InteractivityMode::RequiresExistingArtifact
&& !phase_artifact_on_develop(project_root, phase, "-PLAN.md")
{
println!(
"warning: phase {phase} has no PLAN.md on develop — headless {} \
planning is untested and may need input; pre-writing plans is safer",
driver.name()
);
}
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-{padded})",
wt.display(),
padded = phase.padded(),
);
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());
}
}
}
let launch_root: PathBuf = state
.worktree_path
.clone()
.unwrap_or_else(|| project_root.to_path_buf());
crate::pipeline_launch::repair_leaked_auto_chain_flag(
project_root,
&launch_root,
phase,
crate::pipeline_launch::AUTO_CHAIN_REPAIR_FROM_START,
);
let evidence_root: PathBuf = state
.worktree_path
.clone()
.unwrap_or_else(|| project_root.to_path_buf());
state.last_verification_fingerprint =
agent_result::phase_verification_fingerprint(&evidence_root, phase);
state.last_verification_mtime_nanos =
agent_result::phase_verification_mtime_nanos(&evidence_root, phase);
state.verification_baseline_captured = true;
enforce_build_staleness(
project_root,
&state,
env!("DEVFLOW_BUILD_COMMIT"),
env!("DEVFLOW_BUILD_DIRTY") == "true",
)?;
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 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()
.and_then(|p| p.file_name().map(|n| n.to_string_lossy().into_owned())),
})
}
pub(crate) 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 phase_from_worktree_path(worktrees_dir: &Path, path: &Path) -> Option<PhaseId> {
let name = path.strip_prefix(worktrees_dir).ok()?.to_str()?;
let rest = name.strip_prefix("phase-")?;
let label: String = rest
.chars()
.take_while(|c| c.is_ascii_digit() || *c == '.')
.collect();
label.parse().ok()
}
fn state_for_worktree<'a>(
states: &'a [State],
worktrees_dir: &Path,
wt: &worktree::WorktreeInfo,
) -> Option<&'a State> {
if let Some(state) = states
.iter()
.find(|s| s.worktree_path.as_deref() == Some(wt.path.as_path()))
{
return Some(state);
}
if let Some(phase) = phase_from_worktree_path(worktrees_dir, &wt.path)
&& let Some(state) = states.iter().find(|s| s.phase == phase)
{
return Some(state);
}
if let Some(branch) = &wt.branch {
return states
.iter()
.find(|s| *branch == format!("{FEATURE_PREFIX}phase-{}", s.phase.padded()));
}
None
}
fn remove_worktree_with_retry(
project_root: &Path,
path: &Path,
force: bool,
) -> Result<(), worktree::WorktreeError> {
const ATTEMPTS: u32 = 3;
const BASE_DELAY_MS: u64 = 50;
let mut last_err = None;
for attempt in 0..ATTEMPTS {
match worktree::remove(project_root, path, force) {
Ok(()) => return Ok(()),
Err(err) => {
last_err = Some(err);
if attempt + 1 < ATTEMPTS {
std::thread::sleep(std::time::Duration::from_millis(
BASE_DELAY_MS * 2u64.pow(attempt),
));
}
}
}
}
Err(last_err.expect("loop runs ATTEMPTS >= 1 times"))
}
pub(crate) 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 states = workflow::list_states(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;
}
let matched_state = state_for_worktree(&states, &worktrees_dir, wt);
let phase = matched_state
.map(|s| s.phase)
.or_else(|| phase_from_worktree_path(&worktrees_dir, &wt.path));
let agent_alive = phase
.and_then(|p| agent_pid_from_file(project_root, p))
.is_some_and(agent::agent_running);
let monitor_pid = matched_state.and_then(|s| s.monitor_pid);
let monitor_alive = monitor_pid.is_some_and(agent::agent_running);
let phase_liveness = liveness(monitor_pid, monitor_alive, agent_alive);
let stopped = matched_state.is_some_and(|s| s.stopped);
if stopped && !force {
let phase_label = phase
.map(|p| p.to_string())
.unwrap_or_else(|| "?".to_string());
println!(
"keeping worktree {} for phase {phase_label} — halted via --until; run `devflow resume --phase {phase_label}` first, or pass --force to discard it",
wt.path.display()
);
continue;
}
if agent_alive || matches!(phase_liveness, Liveness::Healthy | Liveness::BetweenStages) {
let phase_label = phase
.map(|p| p.to_string())
.unwrap_or_else(|| "?".to_string());
return Err(CliError::Message(format!(
"refusing to remove worktree {} for phase {phase_label} ({}) — run `devflow resume --phase {phase_label}` or wait for it to finish",
wt.path.display(),
phase_liveness.describe(),
)));
}
match remove_worktree_with_retry(project_root, &wt.path, force) {
Ok(()) => {
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;
}
Err(err) => {
println!(
"warning: could not remove worktree {} after retrying — manually delete this directory: {err}",
wt.path.display()
);
}
}
}
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)]
pub(crate) enum Liveness {
Healthy,
BetweenStages,
Stuck,
Unknown,
}
impl Liveness {
pub(crate) 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 recovery_hints(state: &State, liveness: Liveness) -> Vec<String> {
if liveness != Liveness::Stuck {
return Vec::new();
}
let mut hints = vec![format!("devflow resume --phase {}", state.phase)];
if state.gate_pending {
hints.push(format!("devflow advance --phase {}", state.phase));
}
hints
}
fn render_stage_progress_line(stage: Stage, stage_launched_ts: Option<u64>) -> String {
match stage_launched_ts {
Some(ts) => format!(
" in stage {stage}: {}",
recover::format_age(&ts.to_string())
),
None => format!(" in stage {stage}"),
}
}
pub(crate) 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::driver_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());
let phase_summary = last_events.remove(&state.phase);
println!(
"{}",
render_stage_progress_line(
state.stage,
phase_summary.as_ref().and_then(|s| s.stage_launched_ts)
)
);
for hint in recovery_hints(state, phase_liveness) {
println!(" → {hint}");
}
if let Some(summary) = phase_summary {
let ago = summary
.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(&summary.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(())
}
pub(crate) fn evidence(
project_root: &Path,
phase: PhaseId,
json: bool,
require_shipped: bool,
) -> Result<(), CliError> {
let evidence = ship_evidence::collect(project_root, phase);
if json {
println!(
"{}",
serde_json::to_string_pretty(&evidence).expect("ShipEvidence must serialize")
);
} else {
println!("phase: {}", evidence.phase);
println!("shipped: {}", evidence.shipped);
println!(
"workflow_finished_seen: {}",
evidence.workflow_finished_seen
);
println!(
"finished_reason: {}",
evidence.finished_reason.as_deref().unwrap_or("none")
);
println!(
"stage: {}",
evidence
.stage
.map(|s| s.to_string())
.unwrap_or_else(|| "none".into())
);
println!("state_present: {}", evidence.state_present);
println!("feature_branch_exists: {}", evidence.feature_branch_exists);
println!("merged_into_develop: {}", evidence.merged_into_develop);
println!("has_remote: {}", evidence.has_remote);
}
if require_shipped && !evidence.shipped {
let detail = if ship_evidence::is_stopped_at(&evidence) {
format!(
"phase {phase} has not shipped — DevFlow's own record shows it stopped after \
one stage (--until) rather than reaching a finalized Ship"
)
} else {
format!("phase {phase} has not shipped — DevFlow has no record of a completed Ship")
};
return Err(CliError::Message(detail));
}
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)
}
pub(crate) fn gate_list(project_root: &Path, all_roots: bool) -> Result<(), CliError> {
if all_roots {
return gate_list_all_roots();
}
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_list_all_roots() -> Result<(), CliError> {
let roots = registry::load_roots();
let mut rows: Vec<(PathBuf, OpenGate)> = Vec::new();
for root in &roots {
for gate in Gates::list_open(&root.project_root) {
rows.push((root.project_root.clone(), gate));
}
}
if rows.is_empty() {
println!("no open gates across {} registered root(s)", roots.len());
return Ok(());
}
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
println!("{:<6} {:<9} {:<9} ROOT / CONTEXT", "PHASE", "STAGE", "AGE");
for (root, gate) in &rows {
println!("{}", render_all_roots_gate_row(root, gate, now));
}
Ok(())
}
fn render_all_roots_gate_row(root: &Path, gate: &OpenGate, now: u64) -> String {
let context = render_gate_context(&gate.context, 100);
format!(
"{:<6} {:<9} {:<9} {}\n {context}",
gate.phase,
gate.stage.to_string(),
render_gate_age(&gate.timestamp, now),
root.display(),
)
}
fn render_gate_age(timestamp: &str, now: u64) -> String {
let Ok(ts) = timestamp.parse::<u64>() else {
return "?".to_string();
};
let Some(age) = now.checked_sub(ts) else {
return "?".to_string();
};
let compact = match age {
s if s < 60 => format!("{s}s"),
s if s < 3600 => format!("{}m", s / 60),
s if s < 86400 => format!("{}h", s / 3600),
s => format!("{}d", s / 86400),
};
if age >= GATE_ESCALATION_THRESHOLD_SECS {
format!("{compact}!")
} else {
compact
}
}
fn resolve_single_open_gate_stage(open: &[OpenGate], phase: PhaseId) -> Result<Stage, CliError> {
let matching: Vec<&OpenGate> = open.iter().filter(|g| g.phase == phase).collect();
match matching.as_slice() {
[] => Err(CliError::Message(format!(
"no open gate for phase {phase} — see `devflow gate list`"
))),
[one] => Ok(one.stage),
many => Err(CliError::Message(format!(
"phase {phase} has several open gates ({}) — pass --stage",
many.iter()
.map(|g| g.stage.to_string())
.collect::<Vec<_>>()
.join(", ")
))),
}
}
pub(crate) fn gate_respond(
project_root: &Path,
phase: PhaseId,
stage: Option<Stage>,
approved: bool,
note: Option<String>,
) -> Result<(), CliError> {
let stage = match stage {
Some(stage) => stage,
None => resolve_single_open_gate_stage(&Gates::list_open(project_root), phase)?,
};
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(())
}
pub(crate) fn gate_sweep(
max_age_secs: Option<u64>,
dry_run: bool,
root: Option<PathBuf>,
reap_strays: bool,
) -> Result<(), CliError> {
let threshold = max_age_secs.unwrap_or_else(config_parse::gate_max_unattended_age_secs);
let explicit_root: Vec<PathBuf> = root.clone().into_iter().collect();
let roots: Vec<PathBuf> = match root {
Some(root) => vec![root],
None => {
registry::prune_missing();
registry::load_roots()
.into_iter()
.map(|r| r.project_root)
.collect()
}
};
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let mut reaped = 0u32;
let mut skipped = 0u32;
let mut left_alone = 0u32;
for project_root in &roots {
for gate in Gates::list_open(project_root) {
let Ok(ts) = gate.timestamp.parse::<u64>() else {
left_alone += 1;
continue;
};
let age = now.saturating_sub(ts);
if age < threshold {
left_alone += 1;
continue;
}
if dry_run {
reaped += 1;
println!(
"would reap phase {} {} (age {age}s) at {}",
gate.phase,
gate.stage,
project_root.display()
);
continue;
}
match Gates::reap(
project_root,
gate.phase,
gate.stage,
"abort: reaped by devflow gate sweep (unattended gate exceeded max age)",
"devflow-reap",
) {
Ok(_) => {
reaped += 1;
events::emit(
project_root,
gate.phase,
"gate_reaped",
serde_json::json!({
"stage": gate.stage.to_string(),
"age_secs": age,
"max_age_secs": threshold,
}),
);
println!(
"reaped phase {} {} (age {age}s) at {}",
gate.phase,
gate.stage,
project_root.display()
);
}
Err(err) => {
skipped += 1;
println!(
"skipped phase {} {} at {} — already answered ({err})",
gate.phase,
gate.stage,
project_root.display()
);
}
}
}
}
if reap_strays {
if !explicit_root.is_empty() {
println!(
"note: --root does not scope this stray pass -- a stray has no project root \
for any registry entry, lock file, or state file to name, so discovery and \
reaping are always machine-wide; the reachability safety filter is also \
computed across every registered root regardless of --root, deliberately, \
because narrowing it would leave other projects' live processes unprotected"
);
}
let candidates = unreachable_stray_candidates(&explicit_root);
let results = reap_stray_candidates(&candidates, dry_run, agent::STRAY_MIN_AGE);
let event_root = roots.first();
for result in &results {
let layer = stray_layer_label(result.layer);
match result.outcome {
StrayReapOutcome::Reaped if dry_run => {
reaped += 1;
println!("would reap stray pid {} ({layer})", result.pid);
}
StrayReapOutcome::Reaped => {
reaped += 1;
if let Some(event_root) = event_root {
events::emit(
event_root,
PhaseId::new(0),
"stray_reaped",
serde_json::json!({
"pid": result.pid,
"layer": layer,
}),
);
}
println!("reaped stray pid {} ({layer})", result.pid);
}
StrayReapOutcome::IdentityMismatch => {
skipped += 1;
println!(
"skipped stray pid {} ({layer}) — identity could not be re-confirmed \
immediately before signalling (the pid may have been recycled since \
discovery); inspect it manually (e.g. `ps -p {}`) before assuming it \
is safe",
result.pid, result.pid
);
}
StrayReapOutcome::ReapFailed => {
skipped += 1;
println!(
"failed to verify death for stray pid {} ({layer}) even after SIGKILL \
escalation — inspect it manually (e.g. `ps -p {}`)",
result.pid, result.pid
);
}
StrayReapOutcome::TooYoung => {
skipped += 1;
println!(
"skipped stray pid {} ({layer}) — younger than the minimum age \
({:?}); it may be a process that has not finished starting. A \
genuine stray will still be there on the next invocation",
result.pid,
agent::STRAY_MIN_AGE
);
}
}
}
if !dry_run && !results.is_empty() {
let remaining = agent::discover_stray_devflow_processes();
if !remaining.is_empty() {
println!(
"note: {} stray process(es) still discoverable after this pass — re-run \
`devflow gate sweep --reap-strays` to clear them",
remaining.len()
);
}
}
}
if dry_run {
println!(
"sweep complete (dry run): {reaped} would be reaped, {skipped} skipped, {left_alone} left alone"
);
} else {
println!("sweep complete: {reaped} reaped, {skipped} skipped, {left_alone} left alone");
}
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum StrayReapOutcome {
Reaped,
IdentityMismatch,
ReapFailed,
TooYoung,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct StrayReapResult {
pid: u32,
layer: agent::StrayLayer,
outcome: StrayReapOutcome,
}
fn reap_stray_candidates(
candidates: &[agent::StrayProcess],
dry_run: bool,
min_age: std::time::Duration,
) -> Vec<StrayReapResult> {
candidates
.iter()
.map(|candidate| {
if !agent::is_same_process(candidate.pid, candidate.start_time) {
return StrayReapResult {
pid: candidate.pid,
layer: candidate.layer,
outcome: StrayReapOutcome::IdentityMismatch,
};
}
let old_enough = agent::process_age(candidate.pid).is_some_and(|age| age >= min_age);
if !old_enough {
return StrayReapResult {
pid: candidate.pid,
layer: candidate.layer,
outcome: StrayReapOutcome::TooYoung,
};
}
if dry_run {
return StrayReapResult {
pid: candidate.pid,
layer: candidate.layer,
outcome: StrayReapOutcome::Reaped,
};
}
let cleared = agent::terminate_and_verify(
candidate.pid,
agent::TERMINATE_VERIFY_WAIT,
agent::TERMINATE_VERIFY_POLL,
);
StrayReapResult {
pid: candidate.pid,
layer: candidate.layer,
outcome: if cleared {
StrayReapOutcome::Reaped
} else {
StrayReapOutcome::ReapFailed
},
}
})
.collect()
}
pub(crate) fn stop(project_root: &Path, phase: PhaseId) -> Result<(), CliError> {
if !stop_via_gate(project_root, phase)? {
stop_via_lock(project_root, phase)?;
}
persist_stopped_state(project_root, phase)
}
fn stop_via_gate(project_root: &Path, phase: PhaseId) -> Result<bool, CliError> {
let Some(gate) = Gates::list_open(project_root)
.into_iter()
.find(|g| g.phase == phase)
else {
return Ok(false);
};
match Gates::reap(
project_root,
phase,
gate.stage,
"abort: stopped by `devflow stop`",
"devflow-stop",
) {
Ok(path) => {
println!(
"stop: wrote a rejection for phase {phase} {} at {} — the process waiting \
on it will pick this up on its next poll, within the 60s backoff cap",
gate.stage,
path.display()
);
Ok(true)
}
Err(GateError::AlreadyResponded { .. }) => {
println!(
"stop: phase {phase} {} already has a response awaiting pickup — the phase \
is already ending",
gate.stage
);
Ok(true)
}
Err(GateError::NoOpenGate { .. }) => Ok(false),
Err(err) => Err(err.into()),
}
}
fn stop_via_lock(project_root: &Path, phase: PhaseId) -> Result<(), CliError> {
let Some((pid_str, _path)) = lock::holder(project_root, phase) else {
println!("stop: no lock held for phase {phase} — nothing is running `advance()`");
return Ok(());
};
let Ok(pid) = pid_str.parse::<u32>() else {
println!(
"stop: phase {phase}'s lock file holds a corrupt pid ({pid_str}) — treating it \
as stale"
);
return Ok(());
};
if !agent::agent_running(pid) {
println!(
"stop: phase {phase}'s lock names pid {pid}, which is not alive — stale lock, \
nothing to signal"
);
return Ok(());
}
match lock::holder_identity(project_root, phase) {
Some((recorded_pid, Some(recorded_start))) if recorded_pid == pid => {
if !agent::is_same_process(pid, recorded_start) {
return Err(CliError::Message(format!(
"refusing to signal pid {pid} for phase {phase} — it is not the \
process that took the lock. The lock recorded start time \
{recorded_start}, but pid {pid} now reports {:?}, so the pid has \
been recycled and belongs to something else. Inspect it manually \
(e.g. `ps -p {pid}`) before proceeding.",
agent::process_start_time(pid)
)));
}
}
Some((_, None)) => {
return Err(CliError::Message(format!(
"refusing to signal pid {pid} for phase {phase} — the lock file records \
no start time, so this process's identity cannot be confirmed. The lock \
predates identity recording; if the run is genuinely stuck, inspect the \
pid manually (e.g. `ps -p {pid}`) and remove \
.devflow/lock-{padded} once you are satisfied.",
padded = phase.padded()
)));
}
_ => {
return Err(CliError::Message(format!(
"refusing to signal pid {pid} for phase {phase} — the lock file's holder \
could not be read back for identity confirmation. Inspect it manually \
(e.g. `ps -p {pid}`) before proceeding."
)));
}
}
if agent::terminate(pid) {
println!("stop: signalled pid {pid}, phase {phase}'s lock holder");
} else {
println!("stop: pid {pid} could not be signalled (it may have just exited)");
}
Ok(())
}
fn persist_stopped_state(project_root: &Path, phase: PhaseId) -> Result<(), CliError> {
let mut state = match workflow::load_state(project_root, phase) {
Ok(state) => state,
Err(workflow::WorkflowError::MissingState(_)) => {
println!("stop: no persisted state for phase {phase} — already stopped");
return Ok(());
}
Err(err) => return Err(err.into()),
};
state.stopped = true;
let reason = "stopped via `devflow stop`".to_string();
state.stop_reason = Some(match state.stop_reason.take() {
Some(existing) if !existing.is_empty() => format!("{existing}; {reason}"),
_ => reason,
});
workflow::save_state(&state)?;
Ok(())
}
pub(crate) fn gate_show(
project_root: &Path,
phase: PhaseId,
stage: Option<Stage>,
) -> Result<(), CliError> {
let open = Gates::list_open(project_root);
let stage = match stage {
Some(stage) => stage,
None => resolve_single_open_gate_stage(&open, phase)?,
};
let gate = open
.into_iter()
.find(|g| g.phase == phase && g.stage == stage)
.ok_or_else(|| {
CliError::Message(format!(
"no open gate for phase {phase} stage {stage} — see `devflow gate list`"
))
})?;
println!("{}", render_gate_show(&gate));
Ok(())
}
fn render_gate_show(gate: &OpenGate) -> String {
format!(
"phase {} {} ({})\n{}",
gate.phase,
gate.stage,
recover::format_age(&gate.timestamp),
render_gate_context(&gate.context, usize::MAX),
)
}
pub(crate) fn logs(
project_root: &Path,
phase: Option<PhaseId>,
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;
}
}
pub(crate) fn history_cmd(project_root: &Path, phase: Option<PhaseId>) -> 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> {
let stdout = std::io::stdout();
write_capture_from(path, offset, &mut stdout.lock())
}
fn write_capture_from(
path: &Path,
offset: u64,
output: &mut impl std::io::Write,
) -> Result<u64, CliError> {
use std::io::{Read, Seek, SeekFrom};
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 _ = output.write_all(&buf);
let _ = output.flush();
}
Ok(offset + buf.len() as u64)
}
fn default_logs_phase(project_root: &Path) -> Result<PhaseId, 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, PhaseId)> = 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::<PhaseId>().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: PhaseId) -> 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| cron_hint_line(instructions, project_root))
.collect()
}
fn cron_hint_line(
instructions: &devflow_core::ship::CronInstructions,
project_root: &Path,
) -> String {
let base = format!(
"Cron instruction pending (phase {}): hermes cron create --from-devflow {}",
instructions.phase,
project_root.display()
);
let retry_after = instructions.retry_after.trim();
if retry_after.is_empty() {
base
} else {
let reset = render_gate_context(retry_after, 100);
format!("{base} (rate-limit resets: {reset})")
}
}
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')),
}
}
pub(crate) 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);
}
}
pub(crate) fn recover_cmd(
project_root: &Path,
do_clean: bool,
phase: Option<PhaseId>,
) -> 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::driver_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(())
}
pub(crate) 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 = hermetic_command("sh", project_root)
.arg("-c")
.arg(cmd)
.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(", ")
)))
}
}
pub(crate) struct Check {
pub(crate) name: String,
pub(crate) status: String,
pub(crate) version: Option<String>,
pub(crate) install_hint: Option<String>,
}
pub(crate) 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",
),
cmd_check(
"pi",
"pi",
"--version",
"Install Pi (see https://github.com/earendil-works/pi-mono)",
),
pi_subagent_dispatch_check(),
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);
let doc_findings = collect_planning_doc_findings(project_root);
let stray_findings = collect_stray_process_findings(project_root);
if json {
let body = doctor_json_body(&checks, &facts, &doc_findings, &stray_findings);
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));
print!("{}", render_planning_doc_text(&doc_findings));
print!("{}", render_stray_process_text(&stray_findings));
}
Ok(())
}
pub(crate) fn release_check(project_root: &Path) -> Result<(), CliError> {
let checks: Vec<Check> = vec![
check_self_pin(project_root),
check_divergence(project_root),
check_publish_order(project_root),
check_changelog_version(project_root),
];
let mut failed = false;
for c in &checks {
let icon = match c.status.as_str() {
"ok" => "✓",
"warn" => "⚠",
"fail" => "✗",
_ => "?",
};
let detail = c.version.as_deref().unwrap_or("-");
println!(" {:<32} {icon} {detail}", c.name);
if matches!(c.status.as_str(), "warn" | "fail")
&& let Some(hint) = &c.install_hint
{
println!(" — {hint}");
}
if c.status == "fail" {
failed = true;
}
}
if failed {
Err(CliError::Message(
"release preflight failed — see checks above".into(),
))
} else {
println!("\nrelease preflight passed");
Ok(())
}
}
fn pi_subagent_dispatch_check() -> Check {
let dispatch = agents::driver_for(AgentKind::Pi)
.capabilities()
.subagent_dispatch;
pi_subagent_dispatch_check_for(dispatch)
}
fn pi_subagent_dispatch_check_for(dispatch: bool) -> Check {
Check {
name: "pi subagent dispatch".into(),
status: if dispatch { "ok".into() } else { "warn".into() },
version: Some(if dispatch {
"available".into()
} else {
"not installed".into()
}),
install_hint: if dispatch {
None
} else {
Some(
"optional — `pi install npm:@bacnh85/pi-subagent` (user scope) enables subagent dispatch"
.into(),
)
},
}
}
fn check_self_pin(project_root: &Path) -> Check {
const NAME: &str = "self-pin (workspace member versions)";
let cargo_toml = project_root.join("Cargo.toml");
let contents = match std::fs::read_to_string(&cargo_toml) {
Ok(contents) => contents,
Err(err) => {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some(format!("could not read Cargo.toml: {err}")),
install_hint: None,
};
}
};
let (workspace_version, pins) = version::read_workspace_self_pins(&contents);
let Some(workspace_version) = workspace_version else {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some("not a workspace Cargo.toml (no [workspace.package] version)".into()),
install_hint: None,
};
};
let drifted: Vec<String> = pins
.iter()
.filter(|pin| pin.version != workspace_version)
.map(|pin| format!("{} pinned {} != {workspace_version}", pin.name, pin.version))
.collect();
if drifted.is_empty() {
Check {
name: NAME.into(),
status: "ok".into(),
version: Some(format!(
"{} member pin(s) match {workspace_version}",
pins.len()
)),
install_hint: None,
}
} else {
Check {
name: NAME.into(),
status: "fail".into(),
version: Some(drifted.join("; ")),
install_hint: Some(format!(
"every [workspace.dependencies] self-pin must equal [workspace.package] \
version = \"{workspace_version}\" — VersionBump should have rewritten this; \
see 20a/DEN-49"
)),
}
}
}
fn check_divergence(project_root: &Path) -> Check {
const NAME: &str = "develop/main divergence (origin/main ancestor)";
match devflow_core::git::origin_main_ancestor_status(project_root) {
devflow_core::git::AncestorStatus::Ancestor => Check {
name: NAME.into(),
status: "ok".into(),
version: Some("origin/main is an ancestor of HEAD — sync would be a no-op".into()),
install_hint: None,
},
devflow_core::git::AncestorStatus::Diverged => Check {
name: NAME.into(),
status: "fail".into(),
version: Some("origin/main is NOT an ancestor of HEAD — develop has diverged".into()),
install_hint: Some(
"run scripts/sync-main-to-develop.sh before cutting the next release PR".into(),
),
},
devflow_core::git::AncestorStatus::RefAbsent => Check {
name: NAME.into(),
status: "warn".into(),
version: Some("origin/main not fetched — cannot determine divergence".into()),
install_hint: Some("run `git fetch` first, then re-run this check".into()),
},
}
}
fn check_publish_order(project_root: &Path) -> Check {
const NAME: &str = "crates.io publish order";
let order = devflow_core::git::publish_order(project_root);
if order.is_empty() {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some("could not determine workspace publish order".into()),
install_hint: None,
};
}
Check {
name: NAME.into(),
status: "ok".into(),
version: Some(format!("publish in order: {}", order.join(" -> "))),
install_hint: None,
}
}
fn check_changelog_version(project_root: &Path) -> Check {
const NAME: &str = "changelog version (matches workspace)";
let cargo_toml = project_root.join("Cargo.toml");
let contents = match std::fs::read_to_string(&cargo_toml) {
Ok(contents) => contents,
Err(err) => {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some(format!("could not read Cargo.toml: {err}")),
install_hint: None,
};
}
};
let (workspace_version, _) = version::read_workspace_self_pins(&contents);
let Some(workspace_version) = workspace_version else {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some("not a workspace Cargo.toml (no [workspace.package] version)".into()),
install_hint: None,
};
};
let changelog_contents = match std::fs::read_to_string(project_root.join("CHANGELOG.md")) {
Ok(contents) => contents,
Err(err) => {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some(format!("could not read CHANGELOG.md: {err}")),
install_hint: None,
};
}
};
let changelog_version = changelog_contents.lines().find_map(|line| {
let rest = line.trim().strip_prefix("## ")?;
let version = rest.split_whitespace().next().map(|w| {
w.trim_start_matches('[')
.trim_end_matches(']')
.trim_start_matches('v')
});
version
.filter(|v| {
let Some(dot) = v.find('.') else {
return false;
};
let bytes = v.as_bytes();
bytes.first().is_some_and(|b| b.is_ascii_digit())
&& bytes.get(dot + 1).is_some_and(|b| b.is_ascii_digit())
})
.map(str::to_string)
});
let Some(changelog_version) = changelog_version else {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some("no `## <version>` heading found in CHANGELOG.md".into()),
install_hint: None,
};
};
if changelog_version == workspace_version {
Check {
name: NAME.into(),
status: "ok".into(),
version: Some(format!("changelog {changelog_version} matches workspace")),
install_hint: None,
}
} else {
let direction = match compare_versions(&changelog_version, &workspace_version) {
Some(std::cmp::Ordering::Greater) => {
"changelog ahead of Cargo.toml (version bump not yet applied)"
}
Some(std::cmp::Ordering::Less) => {
"Cargo.toml ahead of changelog (release notes missing)"
}
_ => "direction undetermined (equal-numeric or unparseable version)",
};
Check {
name: NAME.into(),
status: "fail".into(),
version: Some(format!(
"changelog {changelog_version} != workspace {workspace_version} — {direction}"
)),
install_hint: Some(direction.into()),
}
}
}
fn compare_versions(a: &str, b: &str) -> Option<std::cmp::Ordering> {
let components = |v: &str| -> Option<(Vec<u32>, bool)> {
let parts = v
.split(|c: char| !c.is_ascii_digit())
.filter(|s| !s.is_empty())
.take(3)
.map(|s| s.parse().ok())
.collect::<Option<Vec<_>>>()?;
if parts.is_empty() {
return None;
}
let has_prerelease = v.contains('-');
Some((parts, has_prerelease))
};
let (a, a_pre) = components(a)?;
let (b, b_pre) = components(b)?;
let mut cmp = a.cmp(&b);
if cmp == std::cmp::Ordering::Equal {
cmp = match (a_pre, b_pre) {
(false, true) => std::cmp::Ordering::Greater,
(true, false) => std::cmp::Ordering::Less,
_ => std::cmp::Ordering::Equal,
};
}
Some(cmp)
}
pub(crate) fn release_verify(project_root: &Path) -> Result<(), CliError> {
let checks: Vec<Check> = vec![
check_tag_on_main(project_root),
check_sync_done(project_root),
];
let mut failed = false;
for c in &checks {
let icon = match c.status.as_str() {
"ok" => "✓",
"warn" => "⚠",
"fail" => "✗",
_ => "?",
};
let detail = c.version.as_deref().unwrap_or("-");
println!(" {:<32} {icon} {detail}", c.name);
if matches!(c.status.as_str(), "warn" | "fail")
&& let Some(hint) = &c.install_hint
{
println!(" — {hint}");
}
if c.status == "fail" {
failed = true;
}
}
if failed {
Err(CliError::Message(
"release verification failed — see checks above".into(),
))
} else {
println!("\nrelease verification passed");
Ok(())
}
}
fn check_tag_on_main(project_root: &Path) -> Check {
const NAME: &str = "release tag on main";
let cargo_toml = project_root.join("Cargo.toml");
let contents = match std::fs::read_to_string(&cargo_toml) {
Ok(c) => c,
Err(err) => {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some(format!("could not read Cargo.toml: {err}")),
install_hint: None,
};
}
};
let (Some(workspace_version), _) = version::read_workspace_self_pins(&contents) else {
return Check {
name: NAME.into(),
status: "warn".into(),
version: Some("no [workspace.package] version to derive the tag from".into()),
install_hint: None,
};
};
let tag = format!("v{workspace_version}");
let tag_exists = devflow_core::git::git_command(project_root)
.args(["rev-parse", "--verify", "--quiet", &tag])
.output()
.map(|out| out.status.success())
.unwrap_or(false);
if !tag_exists {
return Check {
name: NAME.into(),
status: "fail".into(),
version: Some(format!("tag {tag} does not exist")),
install_hint: Some("tag the release commit on main before verifying".into()),
};
}
match devflow_core::git::ref_is_ancestor(project_root, &tag, "origin/main") {
devflow_core::git::AncestorStatus::Ancestor => Check {
name: NAME.into(),
status: "ok".into(),
version: Some(format!("{tag} is on origin/main")),
install_hint: None,
},
devflow_core::git::AncestorStatus::Diverged => Check {
name: NAME.into(),
status: "fail".into(),
version: Some(format!(
"{tag} is NOT an ancestor of origin/main — tagged the wrong branch"
)),
install_hint: Some(
"re-tag on main: git -c user.signingkey=\"$(git config --get \
devflow.releaseSigningKey)\" tag -s -f vX.Y.Z origin/main"
.into(),
),
},
devflow_core::git::AncestorStatus::RefAbsent => Check {
name: NAME.into(),
status: "warn".into(),
version: Some("origin/main not fetched — cannot compare tag placement".into()),
install_hint: Some("run `git fetch` first, then re-run".into()),
},
}
}
fn check_sync_done(project_root: &Path) -> Check {
const NAME: &str = "main→develop sync";
match devflow_core::git::ref_is_ancestor(project_root, "origin/main", "origin/develop") {
devflow_core::git::AncestorStatus::Ancestor => Check {
name: NAME.into(),
status: "ok".into(),
version: Some("origin/main is an ancestor of origin/develop".into()),
install_hint: None,
},
devflow_core::git::AncestorStatus::Diverged => Check {
name: NAME.into(),
status: "fail".into(),
version: Some(
"origin/main is NOT an ancestor of origin/develop — sync was skipped".into(),
),
install_hint: Some(
"run scripts/sync-main-to-develop.sh, then PR the merge commit (not squash)".into(),
),
},
devflow_core::git::AncestorStatus::RefAbsent => Check {
name: NAME.into(),
status: "warn".into(),
version: Some("origin/main or origin/develop not fetched — cannot check sync".into()),
install_hint: Some("run `git fetch` first, then re-run".into()),
},
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Severity {
Ok,
Warn,
Problem,
}
impl Severity {
pub(crate) fn label(self) -> &'static str {
match self {
Severity::Ok => "ok",
Severity::Warn => "warn",
Severity::Problem => "problem",
}
}
}
pub(crate) struct PhaseFacts {
pub(crate) phase: PhaseId,
pub(crate) stage: Stage,
pub(crate) gate_pending: bool,
pub(crate) agent_pid: Option<u32>,
pub(crate) agent_alive: bool,
pub(crate) monitor_pid: Option<u32>,
pub(crate) monitor_alive: bool,
pub(crate) last_event: Option<String>,
pub(crate) last_launched_stage: Option<Stage>,
pub(crate) open_gate_stages: Vec<Stage>,
pub(crate) feature_branch_exists: bool,
pub(crate) stopped: bool,
}
pub(crate) struct PhaseFinding {
pub(crate) phase: PhaseId,
pub(crate) severity: Severity,
pub(crate) detail: String,
pub(crate) 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.stopped || 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 facts.stopped
|| 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-{} does not exist but stage is {}",
facts.phase,
facts.phase.padded(),
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<PhaseId, events::PhaseEventSummary>,
open_gates: &[OpenGate],
) -> PhaseFacts {
let phase = state.phase;
let stopped = state.stopped;
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).map(|summary| summary.event);
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-{padded}", padded = phase.padded());
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,
stopped,
}
}
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],
doc_findings: &[PlanningDocFinding],
stray_findings: &[StrayProcessFinding],
) -> serde_json::Value {
serde_json::json!({
"environment": checks_json_value(checks),
"reconciliation": render_reconciliation_json(facts),
"planning_doc_staleness": render_planning_doc_findings_json(doc_findings),
"stray_processes": render_stray_process_findings_json(stray_findings),
})
}
const PLANNING_DOC_STALENESS_CUTOFF: (u32, u32, u32) = (1, 5, 0);
pub(crate) struct PlanningDocFinding {
pub(crate) source: String,
pub(crate) claim: String,
pub(crate) severity: Severity,
pub(crate) detail: String,
pub(crate) repair: Option<String>,
}
pub(crate) fn parse_semver(cell: &str) -> Option<(u32, u32, u32)> {
let cell = cell.strip_prefix('v').unwrap_or(cell);
let mut parts = cell.split('.');
let major = parts.next()?.trim().parse().ok()?;
let minor = parts.next()?.trim().parse().ok()?;
let patch = parts.next()?.trim().parse().ok()?;
if parts.next().is_some() {
return None; }
Some((major, minor, patch))
}
pub(crate) fn parse_planning_doc_versions(text: &str, source: &str) -> Vec<(String, String)> {
let mut rows = Vec::new();
for line in text.lines() {
let trimmed = line.trim();
if !trimmed.starts_with('|') {
continue;
}
let cells: Vec<&str> = trimmed
.trim_matches('|')
.split('|')
.map(str::trim)
.collect();
if cells.len() < 2 {
continue;
}
let label = cells[0];
if label.is_empty()
|| label.eq_ignore_ascii_case("phase")
|| label.chars().all(|c| c == '-')
{
continue;
}
for cell in &cells[1..] {
if parse_semver(cell).is_some() {
rows.push((format!("{source} phase {label}"), (*cell).to_string()));
}
}
}
rows
}
pub(crate) fn tag_exists_and_reachable(project_root: &Path, tag: &str, base_branch: &str) -> bool {
let exists = git_command(project_root)
.args(["rev-parse", "--verify", &format!("refs/tags/{tag}")])
.output()
.is_ok_and(|o| o.status.success());
exists
&& git_command(project_root)
.args(["merge-base", "--is-ancestor", tag, base_branch])
.output()
.is_ok_and(|o| o.status.success())
}
pub(crate) fn reconcile_planning_docs(
rows: &[(String, String)],
tag_lookup: &mut impl FnMut(&str) -> bool,
) -> Vec<PlanningDocFinding> {
let mut findings = Vec::new();
for (label, version_cell) in rows {
let Some(parsed) = parse_semver(version_cell) else {
continue;
};
let tag = if version_cell.starts_with('v') {
version_cell.clone()
} else {
format!("v{version_cell}")
};
if tag_lookup(&tag) {
continue;
}
let severity = if parsed >= PLANNING_DOC_STALENESS_CUTOFF {
Severity::Problem
} else {
Severity::Warn
};
findings.push(PlanningDocFinding {
source: label.clone(),
claim: format!("{label} claims {tag}"),
severity,
detail: format!(
"{label} claims {tag}, but no git tag `{tag}` exists (or it isn't reachable from the base branch)"
),
repair: None,
});
}
findings
}
fn collect_planning_doc_findings(project_root: &Path) -> Vec<PlanningDocFinding> {
let roadmap =
std::fs::read_to_string(project_root.join(".planning/ROADMAP.md")).unwrap_or_default();
let state =
std::fs::read_to_string(project_root.join(".planning/STATE.md")).unwrap_or_default();
let mut rows = parse_planning_doc_versions(&roadmap, "ROADMAP.md");
rows.extend(parse_planning_doc_versions(&state, "STATE.md"));
let mut lookup = |tag: &str| tag_exists_and_reachable(project_root, tag, MAIN);
reconcile_planning_docs(&rows, &mut lookup)
}
fn render_planning_doc_findings_json(findings: &[PlanningDocFinding]) -> serde_json::Value {
serde_json::Value::Array(
findings
.iter()
.map(|f| {
serde_json::json!({
"source": f.source,
"claim": f.claim,
"severity": f.severity.label(),
"detail": f.detail,
"repair": f.repair,
})
})
.collect(),
)
}
fn render_planning_doc_text(findings: &[PlanningDocFinding]) -> String {
if findings.is_empty() {
return "\nplanning docs: consistent with git tags\n".to_string();
}
let mut out = String::from("\nplanning docs:\n");
for finding in findings {
out.push_str(&format!(
" [{}] {}\n",
finding.severity.label(),
finding.detail
));
}
out
}
fn stray_layer_label(layer: agent::StrayLayer) -> &'static str {
match layer {
agent::StrayLayer::MonitorWrapper => "monitor wrapper",
agent::StrayLayer::AdvanceChild => "advance child",
}
}
pub(crate) fn registry_reachable_pids(roots: &[PathBuf]) -> HashSet<u32> {
let mut reachable = HashSet::new();
let mut scanned = HashSet::new();
for root in roots {
if !scanned.insert(root.clone()) {
continue;
}
for state in workflow::list_states(root) {
if let Some(pid) = state.monitor_pid {
reachable.insert(pid);
}
if let Some((pid, _start_time)) = lock::holder_identity(root, state.phase) {
reachable.insert(pid);
}
}
}
reachable
}
pub(crate) fn retain_unreachable_strays(
strays: &[agent::StrayProcess],
reachable: &HashSet<u32>,
) -> Vec<agent::StrayProcess> {
strays
.iter()
.filter(|stray| !reachable.contains(&stray.pid))
.copied()
.collect()
}
fn stray_safety_roots(extra: &[PathBuf]) -> Vec<PathBuf> {
let mut roots: Vec<PathBuf> = registry::load_roots()
.into_iter()
.map(|r| r.project_root)
.collect();
for root in extra {
if !roots.contains(root) {
roots.push(root.clone());
}
}
roots
}
fn unreachable_stray_candidates(extra_roots: &[PathBuf]) -> Vec<agent::StrayProcess> {
let reachable = registry_reachable_pids(&stray_safety_roots(extra_roots));
retain_unreachable_strays(&agent::discover_stray_devflow_processes(), &reachable)
}
pub(crate) struct StrayProcessFinding {
pub(crate) pid: u32,
pub(crate) layer: &'static str,
pub(crate) severity: Severity,
pub(crate) detail: String,
pub(crate) repair: Option<String>,
}
pub(crate) fn build_stray_process_findings(
strays: &[agent::StrayProcess],
) -> Vec<StrayProcessFinding> {
strays
.iter()
.map(|stray| {
let layer = stray_layer_label(stray.layer);
StrayProcessFinding {
pid: stray.pid,
layer,
severity: Severity::Problem,
detail: format!(
"pid {} ({layer}) matches DevFlow's monitor-wrapper or advance-child \
argv shape, is owned by the calling user, and is named by no \
registered project root's state file or lock file",
stray.pid
),
repair: Some(
"devflow gate sweep --reap-strays --dry-run (preview first; re-run \
without --dry-run to reap)"
.to_string(),
),
}
})
.collect()
}
fn collect_stray_process_findings(project_root: &Path) -> Vec<StrayProcessFinding> {
build_stray_process_findings(&unreachable_stray_candidates(&[project_root.to_path_buf()]))
}
fn render_stray_process_findings_json(findings: &[StrayProcessFinding]) -> serde_json::Value {
serde_json::Value::Array(
findings
.iter()
.map(|f| {
serde_json::json!({
"pid": f.pid,
"layer": f.layer,
"severity": f.severity.label(),
"detail": f.detail,
"repair": f.repair,
})
})
.collect(),
)
}
fn render_stray_process_text(findings: &[StrayProcessFinding]) -> String {
if findings.is_empty() {
return String::new();
}
let mut out = String::from(
"\nstray processes (state-orphaned -- no registry/lock/state file reaches them):\n",
);
for finding in findings {
out.push_str(&format!(" {}\n", finding.detail));
if let Some(repair) = &finding.repair {
out.push_str(&format!(" repair: {repair}\n"));
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{Cli, Command, GateCmd};
use clap::Parser;
#[test]
fn pi_subagent_dispatch_check_renders_both_arms() {
let available = pi_subagent_dispatch_check_for(true);
assert_eq!(available.name, "pi subagent dispatch");
assert_eq!(available.status, "ok");
assert_eq!(available.version.as_deref(), Some("available"));
assert_eq!(available.install_hint, None);
let missing = pi_subagent_dispatch_check_for(false);
assert_eq!(missing.name, "pi subagent dispatch");
assert_eq!(missing.status, "warn");
assert_eq!(missing.version.as_deref(), Some("not installed"));
assert!(missing.install_hint.is_some());
assert!(
missing
.install_hint
.as_deref()
.is_some_and(|h| h.contains("@bacnh85/pi-subagent")),
"the absent hint must name the vetted install command"
);
}
#[test]
fn phase_validate_failures_survive_a_forced_restart() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(71);
let mut persisted = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
persisted.phase_validate_failures = 6;
persisted.consecutive_failures = 2;
persisted.infra_failures = 4;
persisted.preflight_retries = 3;
persisted.checkpoint_resumes = 1;
persisted.last_validate_failure_commit_count = Some(9);
persisted.stage = Stage::Validate;
persisted.stop_until = Some(Stage::Code);
workflow::save_state(&persisted).unwrap();
let fresh = fresh_state_carrying_phase_failures(root, phase, AgentKind::Claude, Mode::Auto);
assert_eq!(
fresh.phase_validate_failures, 6,
"the per-phase total must be carried across a forced restart — a new process \
starting is not one of A-11's two reset events"
);
assert_eq!(
fresh.consecutive_failures, 0,
"the streak is per-run and must start at zero"
);
assert_eq!(fresh.infra_failures, 0);
assert_eq!(fresh.preflight_retries, 0);
assert_eq!(fresh.checkpoint_resumes, 0);
assert_eq!(
fresh.last_validate_failure_commit_count, None,
"the commit baseline is per-run; carrying it would compare a new run's count \
against an old run's observation"
);
assert_eq!(fresh.stage, Stage::Define);
assert_eq!(fresh.stop_until, None);
}
#[test]
fn phase_validate_failures_reset_when_the_phase_completes() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(72);
let mut persisted = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
persisted.phase_validate_failures = 8;
workflow::save_state(&persisted).unwrap();
assert_eq!(
fresh_state_carrying_phase_failures(root, phase, AgentKind::Claude, Mode::Auto)
.phase_validate_failures,
8,
"premise: while the state file exists, the total IS carried — otherwise the \
assertion below proves nothing about completion"
);
workflow::clear_state(root, phase).unwrap();
assert_eq!(
fresh_state_carrying_phase_failures(root, phase, AgentKind::Claude, Mode::Auto)
.phase_validate_failures,
0,
"phase completion cleared the state, so the next start begins with a full budget"
);
}
#[test]
fn a_corrupt_state_file_warns_while_an_absent_one_is_silent() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(73);
let (absent_total, absent_warning) =
carried_phase_failures(phase, workflow::load_state(root, phase));
assert_eq!(absent_total, 0);
assert_eq!(
absent_warning, None,
"a phase's first start is not an anomaly and must stay silent"
);
let mut persisted = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
persisted.phase_validate_failures = 6;
workflow::save_state(&persisted).unwrap();
let (ok_total, ok_warning) =
carried_phase_failures(phase, workflow::load_state(root, phase));
assert_eq!(ok_total, 6);
assert_eq!(ok_warning, None);
let state_path = workflow::state_path(root, phase);
assert!(
state_path.exists(),
"the fixture must corrupt an EXISTING file; a missing one is the case above"
);
std::fs::write(&state_path, "{\"phase\": 73, \"stage\":").unwrap();
assert!(
!matches!(
workflow::load_state(root, phase),
Err(workflow::WorkflowError::MissingState(_))
),
"premise: a corrupt file must be a DIFFERENT error from an absent one, or \
nothing downstream could tell them apart"
);
let (corrupt_total, corrupt_warning) =
carried_phase_failures(phase, workflow::load_state(root, phase));
assert_eq!(
corrupt_total, 0,
"there is no total to carry — that is the point of the warning"
);
let corrupt_warning =
corrupt_warning.expect("an unreadable budget must not restart at zero silently");
assert!(
corrupt_warning.contains("restarts at zero"),
"the operator must be told what the consequence is, got: {corrupt_warning:?}"
);
}
#[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 gate_show_arg_parsing_accepts_phase_and_optional_stage() {
let bare = Cli::try_parse_from(["devflow", "gate", "show", "15"]).unwrap();
let Command::Gate {
action: GateCmd::Show { phase, stage, .. },
} = bare.command
else {
panic!("expected gate show command");
};
assert_eq!(phase, PhaseId::new(15));
assert_eq!(stage, None);
let flagged =
Cli::try_parse_from(["devflow", "gate", "show", "15", "--stage", "ship"]).unwrap();
let Command::Gate {
action: GateCmd::Show { phase, stage, .. },
} = flagged.command
else {
panic!("expected gate show command with stage");
};
assert_eq!(phase, PhaseId::new(15));
assert_eq!(stage, Some(Stage::Ship));
}
#[test]
fn gate_show_renders_full_untruncated_sanitized_context() {
let dir = tempfile::tempdir().unwrap();
let context = format!("first line\n\u{1b}[2J{}", "x".repeat(150));
Gates::write_gate(dir.path(), PhaseId::new(15), Stage::Ship, &context).unwrap();
let gate = Gates::list_open(dir.path())
.into_iter()
.find(|g| g.phase == PhaseId::new(15))
.unwrap();
let rendered = render_gate_show(&gate);
assert!(rendered.contains(&"x".repeat(150)));
assert!(!rendered.contains("[truncated"));
assert!(!rendered.contains('\u{1b}'));
}
#[test]
fn gate_show_errors_naming_gate_list_when_no_open_gate() {
let dir = tempfile::tempdir().unwrap();
let err = gate_show(dir.path(), PhaseId::new(15), None).unwrap_err();
assert!(err.to_string().contains("devflow gate list"));
}
#[test]
fn gate_show_errors_asking_for_stage_with_several_open_gates() {
let dir = tempfile::tempdir().unwrap();
Gates::write_gate(dir.path(), PhaseId::new(15), Stage::Ship, "ctx1").unwrap();
Gates::write_gate(dir.path(), PhaseId::new(15), Stage::Validate, "ctx2").unwrap();
let err = gate_show(dir.path(), PhaseId::new(15), None).unwrap_err();
assert!(err.to_string().contains("--stage"));
}
#[test]
fn gate_show_auto_resolves_single_open_gate() {
let dir = tempfile::tempdir().unwrap();
Gates::write_gate(
dir.path(),
PhaseId::new(15),
Stage::Ship,
"the only open gate",
)
.unwrap();
assert!(gate_show(dir.path(), PhaseId::new(15), None).is_ok());
}
#[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 cron_instruction_hints_include_hermes_command_per_phase() {
let dir = tempfile::tempdir().unwrap();
for phase in [PhaseId::new(7), PhaseId::new(9)] {
let instructions =
devflow_core::ship::build_single_agent_cron_instructions(dir.path(), phase, "");
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 cron_hint_line_appends_sanitized_reset_when_retry_after_present() {
let dir = tempfile::tempdir().unwrap();
let instructions = devflow_core::ship::build_single_agent_cron_instructions(
dir.path(),
PhaseId::new(7),
"2026-06-18T15:45:30Z",
);
let hint = cron_hint_line(&instructions, dir.path());
assert!(hint.starts_with(&format!(
"Cron instruction pending (phase 7): hermes cron create --from-devflow {}",
dir.path().display()
)));
assert!(hint.contains("(rate-limit resets: 2026-06-18T15:45:30Z)"));
}
#[test]
fn cron_hint_line_omits_reset_fragment_when_retry_after_empty() {
let dir = tempfile::tempdir().unwrap();
let instructions = devflow_core::ship::build_single_agent_cron_instructions(
dir.path(),
PhaseId::new(7),
"",
);
let hint = cron_hint_line(&instructions, dir.path());
assert_eq!(
hint,
format!(
"Cron instruction pending (phase 7): hermes cron create --from-devflow {}",
dir.path().display()
)
);
assert!(!hint.contains("resets"));
}
#[test]
fn default_logs_phase_prefers_single_active_state() {
let dir = tempfile::tempdir().unwrap();
let state = State::new(
PhaseId::new(6),
AgentKind::Claude,
Mode::Auto,
dir.path().to_path_buf(),
);
workflow::save_state(&state).unwrap();
assert_eq!(default_logs_phase(dir.path()).unwrap(), PhaseId::new(6));
}
#[test]
fn default_logs_phase_is_ambiguous_with_two_active_states() {
let dir = tempfile::tempdir().unwrap();
for phase in [PhaseId::new(6), PhaseId::new(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(), PhaseId::new(3)),
"old",
)
.unwrap();
std::thread::sleep(std::time::Duration::from_millis(20));
std::fs::write(
agent_result::stdout_path(dir.path(), PhaseId::new(5)),
"new",
)
.unwrap();
assert_eq!(default_logs_phase(dir.path()).unwrap(), PhaseId::new(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(
PhaseId::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(
PhaseId::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, PhaseId::new(7)).unwrap();
let reloaded8 = workflow::load_state(root, PhaseId::new(8)).unwrap();
assert_eq!(reloaded7.monitor_pid, Some(111));
assert_eq!(reloaded8.monitor_pid, Some(222));
}
#[test]
fn status_reading_monitor_liveness_writes_no_state_and_no_event() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(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, PhaseId::new(15), Stage::Ship, "approve merge?").unwrap();
gate_respond(root, PhaseId::new(15), None, true, Some("lgtm".into())).unwrap();
let polled = Gates::poll_response(root, PhaseId::new(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, PhaseId::new(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, PhaseId::new(15), None, true, None).unwrap_err();
assert!(err.to_string().contains("no open gate"), "{err}");
Gates::write_gate(root, PhaseId::new(15), Stage::Validate, "a").unwrap();
Gates::write_gate(root, PhaseId::new(15), Stage::Ship, "b").unwrap();
let err =
gate_respond(root, PhaseId::new(15), None, false, Some("nope".into())).unwrap_err();
assert!(err.to_string().contains("--stage"), "{err}");
gate_respond(
root,
PhaseId::new(15),
Some(Stage::Validate),
false,
Some("gaps".into()),
)
.unwrap();
assert!(
Gates::response_path(root, PhaseId::new(15), Stage::Validate).exists(),
"explicit-stage rejection must land"
);
assert!(!Gates::response_path(root, PhaseId::new(15), Stage::Ship).exists());
}
fn aged_past_threshold() -> u64 {
config_parse::gate_max_unattended_age_secs() + 60 * 60
}
fn backdate_gate(root: &Path, phase: PhaseId, stage: Stage, age_secs: u64) {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs();
let gate = devflow_core::gates::GateFile {
phase,
stage,
context: "ctx".to_string(),
timestamp: now.saturating_sub(age_secs).to_string(),
};
std::fs::write(
Gates::gate_path(root, phase, stage),
serde_json::to_string_pretty(&gate).unwrap(),
)
.unwrap();
}
#[test]
fn gate_sweep_dry_run_does_not_write_a_response() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
Gates::write_gate(root, PhaseId::new(30), Stage::Ship, "ctx").unwrap();
backdate_gate(root, PhaseId::new(30), Stage::Ship, aged_past_threshold());
gate_sweep(None, true, Some(root.to_path_buf()), false).unwrap();
assert!(
!Gates::response_path(root, PhaseId::new(30), Stage::Ship).exists(),
"dry-run must never write a response file"
);
let open = Gates::list_open(root);
assert_eq!(open.len(), 1, "dry-run must leave the gate open");
}
#[test]
fn gate_sweep_skips_already_responded_gate_without_clobbering() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
Gates::write_gate(root, PhaseId::new(31), Stage::Ship, "ctx").unwrap();
backdate_gate(root, PhaseId::new(31), Stage::Ship, aged_past_threshold());
let response = GateResponse {
approved: true,
note: None,
responded_by: Some("human".into()),
};
Gates::respond(root, PhaseId::new(31), Stage::Ship, &response).unwrap();
let before =
std::fs::read_to_string(Gates::response_path(root, PhaseId::new(31), Stage::Ship))
.unwrap();
let result = gate_sweep(None, false, Some(root.to_path_buf()), false);
assert!(result.is_ok(), "an already-answered gate must not error");
let after =
std::fs::read_to_string(Gates::response_path(root, PhaseId::new(31), Stage::Ship))
.unwrap();
assert_eq!(
before, after,
"the pre-existing response must be byte-identical afterwards"
);
}
#[test]
fn gate_sweep_emits_gate_reaped_event_on_reap() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
Gates::write_gate(root, PhaseId::new(32), Stage::Ship, "ctx").unwrap();
backdate_gate(root, PhaseId::new(32), Stage::Ship, aged_past_threshold());
gate_sweep(None, false, Some(root.to_path_buf()), false).unwrap();
let event = devflow_core::events::last_event_for_phase(root, PhaseId::new(32)).unwrap();
assert_eq!(event["event"], "gate_reaped");
assert_eq!(event["stage"], "ship");
}
fn stray_candidate_for(pid: u32, layer: agent::StrayLayer) -> agent::StrayProcess {
agent::StrayProcess {
pid,
start_time: agent::process_start_time(pid)
.expect("must be able to read the fixture's own recorded start time"),
layer,
}
}
#[test]
fn reap_stray_candidates_dry_run_never_signals() {
let mut child = std::process::Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn sleep fixture");
let pid = child.id();
let candidate = stray_candidate_for(pid, agent::StrayLayer::MonitorWrapper);
let results = reap_stray_candidates(&[candidate], true, std::time::Duration::ZERO);
assert_eq!(results.len(), 1);
assert_eq!(results[0].outcome, StrayReapOutcome::Reaped);
assert!(
agent::agent_running(pid),
"dry-run must never actually signal the candidate"
);
let _ = child.kill();
let _ = child.wait();
}
#[test]
fn reap_stray_candidates_clears_a_real_child_with_verified_death() {
let mut child = std::process::Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn sleep fixture");
let pid = child.id();
let candidate = stray_candidate_for(pid, agent::StrayLayer::AdvanceChild);
let results = reap_stray_candidates(&[candidate], false, std::time::Duration::ZERO);
assert_eq!(results.len(), 1);
assert_eq!(results[0].outcome, StrayReapOutcome::Reaped);
assert!(
!agent::agent_running(pid),
"a successful reap must leave the candidate verified dead, not merely signalled"
);
let _ = child.wait();
}
#[test]
fn reap_stray_candidates_escalates_to_kill_for_a_term_ignoring_child() {
let mut child = std::process::Command::new("sh")
.arg("-c")
.arg("trap '' TERM; sleep 30")
.spawn()
.expect("spawn TERM-ignoring fixture");
let pid = child.id();
std::thread::sleep(std::time::Duration::from_millis(100));
let candidate = stray_candidate_for(pid, agent::StrayLayer::MonitorWrapper);
let results = reap_stray_candidates(&[candidate], false, std::time::Duration::ZERO);
assert_eq!(results.len(), 1);
assert_eq!(results[0].outcome, StrayReapOutcome::Reaped);
assert!(
!agent::agent_running(pid),
"a TERM-ignoring candidate must still be cleared via SIGKILL escalation"
);
let _ = child.wait();
}
#[test]
fn reap_stray_candidates_refuses_on_identity_mismatch_without_signalling() {
let mut child = std::process::Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn sleep fixture");
let pid = child.id();
let real_start = agent::process_start_time(pid).expect("read real start time");
let mismatched = agent::StrayProcess {
pid,
start_time: real_start.wrapping_add(1),
layer: agent::StrayLayer::MonitorWrapper,
};
let results = reap_stray_candidates(&[mismatched], false, std::time::Duration::ZERO);
assert_eq!(results.len(), 1);
assert_eq!(results[0].outcome, StrayReapOutcome::IdentityMismatch);
assert!(
agent::agent_running(pid),
"an identity mismatch must never be signalled — the whole point of the \
re-confirmation"
);
let _ = child.kill();
let _ = child.wait();
}
#[test]
fn reap_stray_candidates_refuses_a_candidate_younger_than_the_minimum_age() {
let mut child = std::process::Command::new("sh")
.arg("-c")
.arg("trap cleanup TERM INT; sleep 30")
.spawn()
.expect("spawn monitor-wrapper-shaped fixture");
let pid = child.id();
assert!(
devflow_core::test_support::wait_for_exec_visibility(
pid,
"sh",
devflow_core::test_support::EXEC_VISIBILITY_WAIT,
devflow_core::test_support::EXEC_VISIBILITY_POLL,
),
"pid {pid}: exec visibility timed out before the fixture became discoverable"
);
let candidate = stray_candidate_for(pid, agent::StrayLayer::MonitorWrapper);
let results = reap_stray_candidates(&[candidate], false, agent::STRAY_MIN_AGE);
assert_eq!(results.len(), 1);
assert_eq!(results[0].outcome, StrayReapOutcome::TooYoung);
assert!(
agent::agent_running(pid),
"a candidate refused for youth must never actually be signalled"
);
let _ = child.kill();
let _ = child.wait();
}
#[test]
fn reap_stray_candidates_reaps_when_the_floor_is_zero() {
let mut child = std::process::Command::new("sh")
.arg("-c")
.arg("trap cleanup TERM INT; sleep 30")
.spawn()
.expect("spawn monitor-wrapper-shaped fixture");
let pid = child.id();
assert!(
devflow_core::test_support::wait_for_exec_visibility(
pid,
"sh",
devflow_core::test_support::EXEC_VISIBILITY_WAIT,
devflow_core::test_support::EXEC_VISIBILITY_POLL,
),
"pid {pid}: exec visibility timed out before the fixture became discoverable"
);
let candidate = stray_candidate_for(pid, agent::StrayLayer::MonitorWrapper);
let results = reap_stray_candidates(&[candidate], false, std::time::Duration::ZERO);
assert_eq!(results.len(), 1);
assert_eq!(results[0].outcome, StrayReapOutcome::Reaped);
assert!(
!agent::agent_running(pid),
"with the floor disabled, a genuine candidate must still be reaped and verified dead"
);
let _ = child.wait();
}
#[test]
fn reap_stray_candidates_refuses_a_dead_pid_as_identity_mismatch_before_the_age_check_runs() {
let dead_pid = 9_999_999;
assert!(
agent::process_age(dead_pid).is_none(),
"precondition: a dead pid's age must be unresolvable"
);
let candidate = agent::StrayProcess {
pid: dead_pid,
start_time: 0,
layer: agent::StrayLayer::MonitorWrapper,
};
let results = reap_stray_candidates(&[candidate], false, agent::STRAY_MIN_AGE);
assert_eq!(results.len(), 1);
assert_eq!(
results[0].outcome,
StrayReapOutcome::IdentityMismatch,
"is_same_process is evaluated before the age check and fails first for a dead pid"
);
}
#[test]
fn gate_sweep_reap_strays_dry_run_discovers_a_real_stray_without_signalling() {
let mut child = std::process::Command::new("sh")
.arg("-c")
.arg("trap cleanup TERM INT; sleep 30")
.spawn()
.expect("spawn monitor-wrapper-shaped fixture");
let pid = child.id();
let dir = tempfile::tempdir().unwrap();
assert!(
devflow_core::test_support::wait_for_exec_visibility(
pid,
"sh",
devflow_core::test_support::EXEC_VISIBILITY_WAIT,
devflow_core::test_support::EXEC_VISIBILITY_POLL,
),
"pid {pid}: exec visibility timed out before the fixture became discoverable"
);
assert!(
agent::discover_stray_devflow_processes()
.iter()
.any(|p| p.pid == pid),
"the fixture must be part of the real discovery census gate_sweep would use"
);
gate_sweep(None, true, Some(dir.path().to_path_buf()), true).unwrap();
assert!(
agent::agent_running(pid),
"--dry-run must never signal a discovered stray, no matter what else the machine \
is running"
);
let _ = child.kill();
let _ = child.wait();
}
#[test]
fn gate_sweep_without_reap_strays_flag_ignores_a_live_stray() {
let mut child = std::process::Command::new("sh")
.arg("-c")
.arg("trap cleanup TERM INT; sleep 30")
.spawn()
.expect("spawn monitor-wrapper-shaped fixture");
let pid = child.id();
let dir = tempfile::tempdir().unwrap();
gate_sweep(None, false, Some(dir.path().to_path_buf()), false).unwrap();
assert!(
agent::agent_running(pid),
"gate_sweep without --reap-strays must never signal anything discoverable only \
via the process table"
);
let _ = child.kill();
let _ = child.wait();
}
#[test]
fn stop_refuses_to_signal_a_live_pid_that_fails_the_identity_check() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(200);
let pid = std::process::id();
let lock_path = root
.join(".devflow")
.join(format!("lock-{padded}", padded = phase.padded()));
std::fs::create_dir_all(lock_path.parent().unwrap()).unwrap();
std::fs::write(&lock_path, pid.to_string()).unwrap();
assert_eq!(
lock::holder_identity(root, phase),
Some((pid, None)),
"the fixture must be a legacy lock (pid recorded, no start time) — the shape \
stop_via_lock's identity guard treats as unconfirmable"
);
let err = stop(root, phase)
.expect_err("a legacy lock with no recorded start time must be refused");
let message = err.to_string();
assert!(
message.contains(&pid.to_string()),
"error must name the pid it refused to signal, got: {message}"
);
assert!(
message.contains("records no start time"),
"error must say identity cannot be confirmed for a legacy lock, got: {message}"
);
assert!(
lock_path.exists(),
"the lock file must be untouched — stop must not signal anything"
);
}
#[test]
fn stop_refuses_when_the_recorded_start_time_does_not_match() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(201);
let pid = std::process::id();
let real_start = devflow_core::agent::process_start_time(pid)
.expect("read our own start time from /proc");
let wrong_start = real_start + 1;
let lock_path = root
.join(".devflow")
.join(format!("lock-{padded}", padded = phase.padded()));
std::fs::create_dir_all(lock_path.parent().unwrap()).unwrap();
std::fs::write(&lock_path, format!("{pid}\n{wrong_start}")).unwrap();
let err = stop(root, phase).expect_err("a start-time mismatch must be refused");
let message = err.to_string();
assert!(
message.contains(&pid.to_string()),
"error must name the pid it refused to signal, got: {message}"
);
assert!(
message.contains("not the process that took the lock"),
"error must say the identity did not match, got: {message}"
);
assert!(
lock_path.exists(),
"the lock file must be untouched — stop must not signal anything"
);
}
#[test]
fn stop_signals_the_holder_when_the_recorded_identity_matches() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(202);
let mut child = std::process::Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn sleep");
let child_pid = child.id();
let mut child_start = None;
for _ in 0..100 {
child_start = devflow_core::agent::process_start_time(child_pid);
if child_start.is_some() {
break;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
let child_start = child_start.expect("read the child's start time");
let lock_path = root
.join(".devflow")
.join(format!("lock-{padded}", padded = phase.padded()));
std::fs::create_dir_all(lock_path.parent().unwrap()).unwrap();
std::fs::write(&lock_path, format!("{child_pid}\n{child_start}")).unwrap();
let result = stop(root, phase);
let _ = child.kill();
let _ = child.wait();
assert!(
result.is_ok(),
"a matching identity must be signalled, not refused: {result:?}"
);
}
#[test]
fn stop_is_a_success_no_op_when_the_lock_names_a_dead_pid() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(202);
let lock_path = root
.join(".devflow")
.join(format!("lock-{padded}", padded = phase.padded()));
std::fs::create_dir_all(lock_path.parent().unwrap()).unwrap();
std::fs::write(&lock_path, "9999999").unwrap();
let state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
workflow::save_state(&state).unwrap();
stop(root, phase).expect("stop against a stale lock must succeed, not error");
}
#[test]
fn stop_never_treats_monitor_pid_as_a_signalling_target() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let phase = PhaseId::new(203);
let mut state = State::new(phase, AgentKind::Claude, Mode::Auto, root.to_path_buf());
state.monitor_pid = Some(std::process::id());
workflow::save_state(&state).unwrap();
stop(root, phase).expect(
"stop must succeed — a recorded monitor_pid alone must never be signalled or \
treated as a blocker",
);
}
#[test]
fn stop_is_a_success_no_op_with_no_gate_and_no_lock() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
stop(root, PhaseId::new(201)).expect("stop against nothing present must succeed");
}
#[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 mut output = Vec::new();
let offset = write_capture_from(&path, 0, &mut output).unwrap();
assert_eq!(offset, 6);
assert_eq!(output, b"hello ");
output.clear();
assert_eq!(write_capture_from(&path, offset, &mut output).unwrap(), 6);
assert!(output.is_empty());
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!(write_capture_from(&path, offset, &mut output).unwrap(), 11);
assert_eq!(output, b"world");
output.clear();
assert_eq!(
write_capture_from(Path::new("/nonexistent/x"), 4, &mut output).unwrap(),
4
);
assert!(output.is_empty());
}
#[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 = devflow_core::test_support::git_command(root)
.args(args)
.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,
PhaseId::new(3),
"-CONTEXT.md"
));
assert!(!phase_artifact_on_develop(
root,
PhaseId::new(3),
"-PLAN.md"
));
assert!(!phase_artifact_on_develop(
root,
PhaseId::new(4),
"-CONTEXT.md"
));
let empty = tempfile::tempdir().unwrap();
assert!(phase_artifact_on_develop(
empty.path(),
PhaseId::new(3),
"-CONTEXT.md"
));
}
#[test]
fn workflow_started_payload_carries_build_provenance() {
let state = State::new(
PhaseId::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.get("exe_path").is_some(),
"WR-02: exe_path key must still exist — a future refactor must not \
satisfy the redaction assertion by deleting the field"
);
assert!(payload["exe_path"].is_string() || payload["exe_path"].is_null());
if let Some(exe_path) = payload["exe_path"].as_str() {
assert!(
!exe_path.contains('/') && !exe_path.contains('\\'),
"WR-02: exe_path must be a bare filename with no directory \
separator — OPERATIONS.md documents events.jsonl as safe to \
tail and paste, so a full absolute path here leaks the \
operator's home directory and OS username; got {exe_path:?}"
);
}
}
#[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(), PhaseId::new(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 render_gate_age_marks_escalated_gate_urgent() {
let now = 10_000u64;
let old_timestamp = (now - GATE_ESCALATION_THRESHOLD_SECS - 60).to_string();
let age = render_gate_age(&old_timestamp, now);
assert!(
age.ends_with('!'),
"an escalated gate's age must carry a trailing urgency marker, got {age:?}"
);
}
#[test]
fn render_gate_age_no_marker_for_fresh_gate() {
let now = 10_000u64;
let fresh_timestamp = (now - 30).to_string();
let age = render_gate_age(&fresh_timestamp, now);
assert!(
!age.ends_with('!'),
"a fresh gate must not carry an urgency marker, got {age:?}"
);
}
#[test]
fn render_gate_age_unknown_for_non_numeric_timestamp() {
let age = render_gate_age("not-a-number", 10_000);
assert_eq!(age, "?");
}
#[test]
fn all_roots_row_includes_gate_with_non_numeric_timestamp() {
let gate = OpenGate {
phase: PhaseId::new(42),
stage: Stage::Ship,
context: "ctx".to_string(),
timestamp: "not-a-number".to_string(),
};
let row = render_all_roots_gate_row(Path::new("/tmp/some-root"), &gate, 10_000);
assert!(row.contains("42"), "row must still name the phase: {row}");
assert!(row.contains('?'), "row must render the unknown age: {row}");
}
#[test]
fn recovery_hints_includes_resume_for_stuck() {
let dir = tempfile::tempdir().unwrap();
let state = State::new(
PhaseId::new(7),
AgentKind::Claude,
Mode::Auto,
dir.path().to_path_buf(),
);
let hints = recovery_hints(&state, Liveness::Stuck);
assert_eq!(hints, vec!["devflow resume --phase 7".to_string()]);
}
#[test]
fn recovery_hints_includes_advance_when_stuck_and_gate_pending() {
let dir = tempfile::tempdir().unwrap();
let mut state = State::new(
PhaseId::new(7),
AgentKind::Claude,
Mode::Auto,
dir.path().to_path_buf(),
);
state.gate_pending = true;
let hints = recovery_hints(&state, Liveness::Stuck);
assert_eq!(
hints,
vec![
"devflow resume --phase 7".to_string(),
"devflow advance --phase 7".to_string(),
]
);
}
#[test]
fn recovery_hints_empty_for_healthy() {
let dir = tempfile::tempdir().unwrap();
let state = State::new(
PhaseId::new(7),
AgentKind::Claude,
Mode::Auto,
dir.path().to_path_buf(),
);
assert!(recovery_hints(&state, Liveness::Healthy).is_empty());
}
#[test]
fn stage_launched_ts_none_without_event() {
let dir = tempfile::tempdir().unwrap();
assert_eq!(
events::last_events_by_phase(dir.path())
.get(&PhaseId::new(7))
.and_then(|s| s.stage_launched_ts),
None
);
}
#[test]
fn stage_launched_ts_reflects_event_age_not_phase_started_at() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs();
let stage_ts = now - 90;
let phase_started_at = now - 30 * 60;
let mut state = State::new(
PhaseId::new(7),
AgentKind::Claude,
Mode::Auto,
root.to_path_buf(),
);
state.started_at = phase_started_at.to_string();
workflow::save_state(&state).unwrap();
events::emit(
root,
PhaseId::new(7),
"stage_launched",
serde_json::json!({"stage": "code", "agent": "claude", "monitor_pid": 1}),
);
let events_path = devflow_core::events::events_path(root);
let rewritten: String = std::fs::read_to_string(&events_path)
.unwrap()
.lines()
.map(|line| {
let mut value: serde_json::Value = serde_json::from_str(line).unwrap();
value["ts"] = serde_json::json!(stage_ts);
value.to_string()
})
.collect::<Vec<_>>()
.join("\n")
+ "\n";
std::fs::write(&events_path, rewritten).unwrap();
let ts = events::last_events_by_phase(root)
.get(&PhaseId::new(7))
.and_then(|s| s.stage_launched_ts);
assert_eq!(ts, Some(stage_ts));
let line = render_stage_progress_line(Stage::Code, ts);
assert!(line.contains("1m ago"), "expected ~90s age, got: {line}");
assert!(
!line.contains("30m ago"),
"must not render phase-level started_at age: {line}"
);
}
#[test]
fn render_stage_progress_line_omits_age_without_stage_launched_event() {
assert_eq!(
render_stage_progress_line(Stage::Plan, None),
" in stage plan"
);
}
#[cfg(test)]
mod doctor_reconciliation {
use super::*;
fn agreeing_facts(phase: PhaseId) -> 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,
stopped: false,
}
}
#[test]
fn reconcile_phase_returns_no_findings_when_all_agree() {
let facts = agreeing_facts(PhaseId::new(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(PhaseId::new(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(PhaseId::new(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(PhaseId::new(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(PhaseId::new(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(PhaseId::new(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(PhaseId::new(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_phase_ignores_dead_agent_when_stopped() {
let facts = PhaseFacts {
stage: Stage::Plan,
agent_pid: Some(999_999),
agent_alive: false,
stopped: true,
..agreeing_facts(PhaseId::new(11))
};
let findings = reconcile_phase(&facts);
assert!(
findings.iter().all(|f| f.severity != Severity::Problem),
"a --until-stopped phase must yield zero Problem findings, got: \
{:?}",
findings.iter().map(|f| &f.detail).collect::<Vec<_>>()
);
}
#[test]
fn reconcile_phase_ignores_dead_monitor_when_stopped() {
let facts = PhaseFacts {
stage: Stage::Plan,
monitor_pid: Some(5150),
monitor_alive: false,
agent_pid: Some(4242),
agent_alive: false,
stopped: true,
..agreeing_facts(PhaseId::new(12))
};
let findings = reconcile_phase(&facts);
assert!(
findings.iter().all(|f| f.severity != Severity::Problem),
"a --until-stopped phase must yield zero Problem findings even with a \
stale monitor_pid, got: {:?}",
findings.iter().map(|f| &f.detail).collect::<Vec<_>>()
);
}
#[test]
fn reconcile_is_silent_when_monitor_pid_is_unrecorded() {
let facts = PhaseFacts {
monitor_pid: None,
monitor_alive: false,
..agreeing_facts(PhaseId::new(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(PhaseId::new(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(PhaseId::new(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 = PhaseId::new(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 = PhaseId::new(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.get("planning_doc_staleness").is_some(),
"21b: must carry the planning-doc findings under a THIRD key, \
never a second concatenated array: {reparsed}"
);
assert!(
reparsed.get("stray_processes").is_some(),
"999.44: must carry the stray-process findings under a FOURTH key, \
never a second concatenated array: {reparsed}"
);
assert_eq!(
reparsed.as_object().unwrap().len(),
4,
"doctor --json must have exactly four top-level keys: {reparsed}"
);
assert!(reparsed["environment"].is_array());
assert!(reparsed["reconciliation"].is_array());
assert!(reparsed["planning_doc_staleness"].is_array());
assert!(reparsed["stray_processes"].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 = PhaseId::new(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"
);
}
}
#[cfg(test)]
mod stray_process_finding {
use super::*;
#[test]
fn build_stray_process_findings_is_empty_for_no_strays() {
assert!(build_stray_process_findings(&[]).is_empty());
}
#[test]
fn build_stray_process_findings_names_pid_layer_and_repair() {
let strays = vec![agent::StrayProcess {
pid: 424242,
start_time: 0,
layer: agent::StrayLayer::MonitorWrapper,
}];
let findings = build_stray_process_findings(&strays);
assert_eq!(findings.len(), 1);
let finding = &findings[0];
assert_eq!(finding.severity, Severity::Problem);
assert!(finding.detail.contains("424242"));
assert!(finding.detail.contains("monitor wrapper"));
assert!(
finding
.repair
.as_deref()
.unwrap()
.contains("--reap-strays --dry-run"),
"the repair must name the preview form first, not the destructive form alone: \
{:?}",
finding.repair
);
assert!(
!finding.detail.contains('/'),
"no finding may embed a filesystem path (WR-02): {}",
finding.detail
);
}
#[test]
fn build_stray_process_findings_names_advance_child_layer() {
let strays = vec![agent::StrayProcess {
pid: 424243,
start_time: 0,
layer: agent::StrayLayer::AdvanceChild,
}];
let findings = build_stray_process_findings(&strays);
assert!(findings[0].detail.contains("advance child"));
}
#[test]
fn render_stray_process_text_is_empty_when_no_strays() {
assert_eq!(render_stray_process_text(&[]), "");
}
#[test]
fn render_stray_process_text_names_pid_and_repair_when_present() {
let strays = vec![agent::StrayProcess {
pid: 555555,
start_time: 0,
layer: agent::StrayLayer::MonitorWrapper,
}];
let text = render_stray_process_text(&build_stray_process_findings(&strays));
assert!(text.contains("555555"));
assert!(text.contains("repair: devflow gate sweep --reap-strays --dry-run"));
}
#[test]
fn doctor_json_body_carries_stray_processes_as_a_fourth_key() {
let checks: Vec<Check> = Vec::new();
let facts: Vec<PhaseFacts> = Vec::new();
let stray_findings = vec![StrayProcessFinding {
pid: 1,
layer: "monitor wrapper",
severity: Severity::Problem,
detail: "detail".to_string(),
repair: None,
}];
let body = doctor_json_body(&checks, &facts, &[], &stray_findings);
let obj = body.as_object().unwrap();
assert_eq!(obj.len(), 4, "must be exactly four top-level keys: {body}");
assert!(obj.contains_key("stray_processes"));
let strays = obj["stray_processes"].as_array().unwrap();
assert_eq!(strays.len(), 1);
assert_eq!(strays[0]["pid"], 1);
}
#[test]
fn doctor_finds_a_real_stray_and_never_signals_it_across_two_runs() {
let mut child = std::process::Command::new("sh")
.arg("-c")
.arg("trap cleanup TERM INT; sleep 30")
.spawn()
.expect("spawn monitor-wrapper-shaped fixture");
let pid = child.id();
let dir = tempfile::tempdir().unwrap();
assert!(
devflow_core::test_support::wait_for_exec_visibility(
pid,
"sh",
devflow_core::test_support::EXEC_VISIBILITY_WAIT,
devflow_core::test_support::EXEC_VISIBILITY_POLL,
),
"pid {pid}: exec visibility timed out before the fixture became discoverable"
);
let first = collect_stray_process_findings(dir.path());
assert!(
agent::agent_running(pid),
"fixture must still be alive after the first collection"
);
let second = collect_stray_process_findings(dir.path());
assert!(
agent::agent_running(pid),
"fixture must still be alive after the second collection — doctor's stray \
finding must never signal"
);
assert!(
first.iter().any(|f| f.pid == pid),
"the fixture must be reported by the first run"
);
assert!(
second.iter().any(|f| f.pid == pid),
"the fixture must be reported by the second run"
);
doctor(dir.path(), false).unwrap();
assert!(agent::agent_running(pid), "doctor() itself must not signal");
doctor(dir.path(), false).unwrap();
assert!(
agent::agent_running(pid),
"doctor() must remain read-only across repeated runs"
);
let _ = child.kill();
let _ = child.wait();
}
}
#[test]
fn stray_finding_detail_states_only_what_was_checked() {
let strays = vec![agent::StrayProcess {
pid: 909090,
start_time: 0,
layer: agent::StrayLayer::MonitorWrapper,
}];
let findings = build_stray_process_findings(&strays);
assert_eq!(findings.len(), 1);
let finding = &findings[0];
assert!(finding.detail.contains("909090"));
assert!(finding.detail.contains("monitor-wrapper"));
assert!(finding.detail.contains("owned by the calling user"));
assert!(
finding
.detail
.contains("named by no registered project root's state file or lock file"),
"the detail must name the specific checked absence, not an unqualified \
conclusion: {}",
finding.detail
);
let old_unverified_orphan_phrase = format!(
"{}{}",
"reachable through no registry entry, ", "lock file, or state file"
);
assert!(
!finding.detail.contains(&old_unverified_orphan_phrase),
"the old, unverified-orphan phrasing must be gone: {}",
finding.detail
);
assert!(
!finding.detail.contains('/'),
"no finding may embed a filesystem path (WR-02): {}",
finding.detail
);
assert!(
finding
.repair
.as_deref()
.unwrap()
.starts_with("devflow gate sweep --reap-strays --dry-run"),
"the repair must name the --dry-run preview form first: {:?}",
finding.repair
);
}
fn spawn_wrapper_shaped_fixture() -> std::process::Child {
let child = std::process::Command::new("sh")
.arg("-c")
.arg("trap cleanup TERM INT; sleep 30")
.spawn()
.expect("spawn monitor-wrapper-shaped fixture");
let pid = child.id();
assert!(
devflow_core::test_support::wait_for_exec_visibility(
pid,
"sh",
devflow_core::test_support::EXEC_VISIBILITY_WAIT,
devflow_core::test_support::EXEC_VISIBILITY_POLL,
),
"pid {pid}: exec visibility timed out before the fixture became discoverable"
);
child
}
#[test]
fn reachable_pids_are_excluded_from_both_the_findings_and_the_reap_candidates() {
let mut state_child = spawn_wrapper_shaped_fixture();
let mut lock_child = spawn_wrapper_shaped_fixture();
let mut orphan_child = spawn_wrapper_shaped_fixture();
let state_pid = state_child.id();
let lock_pid = lock_child.id();
let orphan_pid = orphan_child.id();
let cache_dir = tempfile::tempdir().unwrap();
let project_root_guard = tempfile::tempdir().unwrap();
let project_root = project_root_guard.path().to_path_buf();
let mut state = State::new(
PhaseId::new(1),
AgentKind::Claude,
Mode::Auto,
project_root.clone(),
);
state.monitor_pid = Some(state_pid);
workflow::save_state(&state).unwrap();
let state_2 = State::new(
PhaseId::new(2),
AgentKind::Claude,
Mode::Auto,
project_root.clone(),
);
workflow::save_state(&state_2).unwrap();
let lock_start_time = agent::process_start_time(lock_pid)
.expect("must read the lock fixture's own recorded start time");
let devflow_dir = project_root.join(".devflow");
std::fs::create_dir_all(&devflow_dir).unwrap();
std::fs::write(
devflow_dir.join("lock-02"),
format!("{lock_pid}\n{lock_start_time}"),
)
.unwrap();
assert_eq!(
lock::holder_identity(&project_root, PhaseId::new(2)),
Some((lock_pid, Some(lock_start_time))),
"the directly-written lock file must read back through holder_identity exactly \
as one lock::acquire itself wrote would"
);
registry::register_in(cache_dir.path(), &project_root, PhaseId::new(1)).unwrap();
let registered_roots: Vec<PathBuf> = registry::load_roots_in(cache_dir.path())
.into_iter()
.map(|r| r.project_root)
.collect();
let reachable = registry_reachable_pids(®istered_roots);
assert!(
reachable.contains(&state_pid),
"state_pid must be reachable via its recorded monitor_pid"
);
assert!(
reachable.contains(&lock_pid),
"lock_pid must be reachable via the lock file's holder_identity"
);
assert!(
!reachable.contains(&orphan_pid),
"orphan_pid is named by no state file and no lock file, so it must not be \
reachable"
);
let census = agent::discover_stray_devflow_processes();
for (pid, label) in [
(state_pid, "state"),
(lock_pid, "lock"),
(orphan_pid, "orphan"),
] {
assert!(
census.iter().any(|p| p.pid == pid),
"the {label} fixture (pid {pid}) must be part of the real /proc census"
);
}
let retained = retain_unreachable_strays(&census, &reachable);
assert!(
retained.iter().any(|p| p.pid == orphan_pid),
"the orphan must survive the filter"
);
assert!(
!retained.iter().any(|p| p.pid == state_pid),
"the state-named pid must be filtered out"
);
assert!(
!retained.iter().any(|p| p.pid == lock_pid),
"the lock-held pid must be filtered out"
);
let findings = build_stray_process_findings(&retained);
assert!(
findings.iter().any(|f| f.pid == orphan_pid),
"the orphan must produce a finding"
);
assert!(
!findings.iter().any(|f| f.pid == state_pid),
"the state-named pid must produce no finding"
);
assert!(
!findings.iter().any(|f| f.pid == lock_pid),
"the lock-held pid must produce no finding"
);
for pid in [state_pid, lock_pid, orphan_pid] {
agent::terminate_and_verify(
pid,
agent::TERMINATE_VERIFY_WAIT,
agent::TERMINATE_VERIFY_POLL,
);
}
let _ = state_child.wait();
let _ = lock_child.wait();
let _ = orphan_child.wait();
}
#[test]
fn a_deleted_root_contributes_nothing_to_the_reachable_set() {
let outer = tempfile::tempdir().unwrap();
let root = outer.path().join("project-root");
std::fs::create_dir_all(&root).unwrap();
let mut child = spawn_wrapper_shaped_fixture();
let pid = child.id();
let mut state = State::new(PhaseId::new(1), AgentKind::Claude, Mode::Auto, root.clone());
state.monitor_pid = Some(pid);
workflow::save_state(&state).unwrap();
assert!(
registry_reachable_pids(std::slice::from_ref(&root)).contains(&pid),
"the pid must be reachable while its root and state file still exist"
);
std::fs::remove_dir_all(&root).unwrap();
assert!(!root.exists(), "the root must actually be gone");
assert!(
agent::agent_running(pid),
"deleting the root must not touch the still-alive process"
);
let reachable_after_deletion = registry_reachable_pids(std::slice::from_ref(&root));
assert!(
reachable_after_deletion.is_empty(),
"a deleted root's state file and lock file are gone with it, so it must \
contribute nothing to the reachable set: got {reachable_after_deletion:?}"
);
let census = vec![agent::StrayProcess {
pid,
start_time: agent::process_start_time(pid)
.expect("fixture must still be alive and readable"),
layer: agent::StrayLayer::MonitorWrapper,
}];
let retained = retain_unreachable_strays(&census, &reachable_after_deletion);
assert!(
retained.iter().any(|p| p.pid == pid),
"with the reachable set empty, the pid must still be treated as a stray"
);
agent::terminate_and_verify(
pid,
agent::TERMINATE_VERIFY_WAIT,
agent::TERMINATE_VERIFY_POLL,
);
let _ = child.wait();
}
#[cfg(test)]
mod planning_doc_staleness {
use super::*;
const SAMPLE_TABLE: &str = "\
| Phase | Name | Version |
|---|---|---|
| 20 | Release Correctness | 1.7.0 |
| 10 | Logging | — |
| 1–5 | Core workflow | 0.1.0–0.6.0 |
| 9 | OSS Polish | 1.2.0 |
| 11 | GSD-Native | 1.2.0 |
";
#[test]
fn parse_planning_doc_versions_skips_non_semver_cells() {
let rows = parse_planning_doc_versions(SAMPLE_TABLE, "ROADMAP.md");
assert_eq!(
rows,
vec![
("ROADMAP.md phase 20".to_string(), "1.7.0".to_string()),
("ROADMAP.md phase 9".to_string(), "1.2.0".to_string()),
("ROADMAP.md phase 11".to_string(), "1.2.0".to_string()),
],
"em-dash and range cells must be skipped; duplicate versions across \
phases (9 and 11 both claim 1.2.0) must both still parse"
);
}
#[test]
fn parse_planning_doc_versions_accepts_v_prefixed_cells() {
let text = "| Phase | Description | Version | Date |\n\
|---|---|---|---|\n\
| 18 | Dogfood Hardening | v1.5.0 | 2026-07-21 |\n";
let rows = parse_planning_doc_versions(text, "STATE.md");
assert_eq!(
rows,
vec![("STATE.md phase 18".to_string(), "v1.5.0".to_string())]
);
}
#[test]
fn parse_semver_rejects_ranges_and_em_dash() {
assert_eq!(parse_semver("1.7.0"), Some((1, 7, 0)));
assert_eq!(parse_semver("v1.7.0"), Some((1, 7, 0)));
assert_eq!(parse_semver("0.1.0–0.6.0"), None);
assert_eq!(parse_semver("—"), None);
assert_eq!(parse_semver("1.7"), None);
assert_eq!(parse_semver("1.7.0.1"), None);
}
#[test]
fn reconcile_planning_docs_flags_problem_for_unreachable_post_cutoff_version() {
let rows = vec![("ROADMAP.md phase 20".to_string(), "1.7.0".to_string())];
let mut lookup = |_tag: &str| false; let findings = reconcile_planning_docs(&rows, &mut lookup);
assert_eq!(findings.len(), 1);
assert_eq!(findings[0].severity, Severity::Problem);
assert!(
findings[0].repair.is_none(),
"D-04: detection-only, no repair"
);
assert!(findings[0].detail.contains("v1.7.0"));
}
#[test]
fn reconcile_planning_docs_downgrades_pre_cutoff_mismatch_to_warn() {
let rows = vec![("ROADMAP.md phase 7".to_string(), "1.0.0".to_string())];
let mut lookup = |_tag: &str| false;
let findings = reconcile_planning_docs(&rows, &mut lookup);
assert_eq!(findings.len(), 1);
assert_eq!(
findings[0].severity,
Severity::Warn,
"pre-v1.5.0 mismatches must downgrade to Warn, never Problem"
);
assert!(findings[0].repair.is_none());
}
#[test]
fn reconcile_planning_docs_numeric_cutoff_is_not_lexicographic() {
let rows = vec![
("label A".to_string(), "1.10.0".to_string()),
("label B".to_string(), "1.4.0".to_string()),
];
let mut lookup = |_tag: &str| false;
let findings = reconcile_planning_docs(&rows, &mut lookup);
assert_eq!(findings.len(), 2);
assert_eq!(
findings[0].severity,
Severity::Problem,
"1.10.0 is numerically >= v1.5.0 (post-cutoff), even though \
\"1.10.0\" < \"1.5.0\" as a string"
);
assert_eq!(
findings[1].severity,
Severity::Warn,
"1.4.0 is numerically < v1.5.0 (pre-cutoff)"
);
}
#[test]
fn reconcile_planning_docs_produces_no_finding_when_tag_is_reachable() {
let rows = vec![("ROADMAP.md phase 20".to_string(), "1.7.0".to_string())];
let mut lookup = |_tag: &str| true; let findings = reconcile_planning_docs(&rows, &mut lookup);
assert!(findings.is_empty());
}
#[test]
fn reconcile_planning_docs_normalizes_bare_cell_to_v_prefixed_tag() {
let rows = vec![("ROADMAP.md phase 20".to_string(), "1.7.0".to_string())];
let mut seen_tag = None;
let mut lookup = |tag: &str| {
seen_tag = Some(tag.to_string());
true
};
reconcile_planning_docs(&rows, &mut lookup);
assert_eq!(seen_tag.as_deref(), Some("v1.7.0"));
}
#[test]
fn reconcile_planning_docs_skips_a_malformed_row_defensively() {
let rows = vec![("bad row".to_string(), "not-a-version".to_string())];
let mut lookup = |_tag: &str| false;
let findings = reconcile_planning_docs(&rows, &mut lookup);
assert!(findings.is_empty());
}
fn init_tagged_repo(root: &Path) {
let git = |args: &[&str]| {
assert!(
devflow_core::test_support::git_command(root)
.args(args)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git(&["init", "-q", "-b", "main"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
git(&["config", "tag.gpgsign", "false"]);
git(&["config", "core.hooksPath", "/dev/null"]);
std::fs::write(root.join("a.txt"), "one").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "base"]);
git(&["tag", "v1.7.0"]);
git(&["checkout", "-q", "-b", "side"]);
std::fs::write(root.join("side.txt"), "s").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "side"]);
git(&["tag", "v9.9.9"]);
git(&["checkout", "-q", "main"]);
}
#[test]
fn tag_exists_and_reachable_true_for_a_tagged_ancestor() {
let dir = tempfile::tempdir().unwrap();
init_tagged_repo(dir.path());
assert!(tag_exists_and_reachable(dir.path(), "v1.7.0", "main"));
}
#[test]
fn tag_exists_and_reachable_false_for_a_missing_tag() {
let dir = tempfile::tempdir().unwrap();
init_tagged_repo(dir.path());
assert!(!tag_exists_and_reachable(dir.path(), "v0.0.1", "main"));
}
#[test]
fn tag_exists_and_reachable_false_for_a_tag_unreachable_from_base() {
let dir = tempfile::tempdir().unwrap();
init_tagged_repo(dir.path());
assert!(!tag_exists_and_reachable(dir.path(), "v9.9.9", "main"));
}
#[test]
fn tag_exists_and_reachable_resolves_caller_root_under_a_hostile_git_dir() {
const INNER_ROOT: &str = "DEVFLOW_27_04_TAG_INNER_ROOT";
const INNER_TAG: &str = "DEVFLOW_27_04_TAG_INNER_TAG";
const INNER_BASE: &str = "DEVFLOW_27_04_TAG_INNER_BASE";
if let Ok(root) = std::env::var(INNER_ROOT) {
let tag = std::env::var(INNER_TAG).expect("inner tag env set by parent");
let base = std::env::var(INNER_BASE).expect("inner base env set by parent");
assert!(
!tag_exists_and_reachable(Path::new(&root), &tag, &base),
"a hostile GIT_DIR pointed at a foreign repository that DOES \
carry this tag must not cause tag_exists_and_reachable to \
report it as belonging to project_root"
);
return;
}
let real = tempfile::tempdir().unwrap();
let real_root = real.path();
let git = |args: &[&str]| {
assert!(
devflow_core::test_support::git_command(real_root)
.args(args)
.output()
.unwrap()
.status
.success(),
"git {args:?} failed"
);
};
git(&["init", "-q", "-b", "main"]);
git(&["config", "user.email", "t@e.st"]);
git(&["config", "user.name", "t"]);
git(&["config", "commit.gpgsign", "false"]);
git(&["config", "core.hooksPath", "/dev/null"]);
std::fs::write(real_root.join("a.txt"), "one").unwrap();
git(&["add", "."]);
git(&["commit", "-q", "-m", "base"]);
let foreign = tempfile::tempdir().unwrap();
init_tagged_repo(foreign.path());
let exe = std::env::current_exe().expect("current_exe for child re-invocation");
let status = std::process::Command::new(&exe)
.arg("tag_exists_and_reachable_resolves_caller_root_under_a_hostile_git_dir")
.arg("--test-threads=1")
.env(INNER_ROOT, real_root.to_str().unwrap())
.env(INNER_TAG, "v1.7.0")
.env(INNER_BASE, "main")
.env("GIT_DIR", foreign.path().join(".git"))
.status()
.expect("spawn hostile child test process");
assert!(
status.success(),
"child test process (hostile GIT_DIR pointed at a foreign repo \
that DOES carry v1.7.0) must still report \
tag_exists_and_reachable == false for the real repository; \
child exit status {status:?}"
);
}
#[test]
fn collect_planning_doc_findings_missing_files_yield_no_findings_not_error() {
let dir = tempfile::tempdir().unwrap();
let findings = collect_planning_doc_findings(dir.path());
assert!(
findings.is_empty(),
"a project with no .planning/ dir at all must yield zero findings, not an error"
);
}
#[test]
fn collect_planning_doc_findings_reconciles_against_main() {
let dir = tempfile::tempdir().unwrap();
init_tagged_repo(dir.path());
std::fs::create_dir_all(dir.path().join(".planning")).unwrap();
std::fs::write(
dir.path().join(".planning/ROADMAP.md"),
"| Phase | Name | Version |\n|---|---|---|\n| 99 | Fixture | 9.9.9 |\n",
)
.unwrap();
let findings = collect_planning_doc_findings(dir.path());
assert_eq!(findings.len(), 1);
assert_eq!(findings[0].severity, Severity::Problem);
}
#[test]
fn render_planning_doc_text_reports_consistent_when_no_findings() {
assert_eq!(
render_planning_doc_text(&[]),
"\nplanning docs: consistent with git tags\n"
);
}
#[test]
fn render_planning_doc_text_lists_each_finding_detail() {
let findings = vec![PlanningDocFinding {
source: "ROADMAP.md phase 20".to_string(),
claim: "ROADMAP.md phase 20 claims v1.7.0".to_string(),
severity: Severity::Problem,
detail: "ROADMAP.md phase 20 claims v1.7.0, but no git tag `v1.7.0` exists"
.to_string(),
repair: None,
}];
let text = render_planning_doc_text(&findings);
assert!(text.contains("[problem]"));
assert!(text.contains("ROADMAP.md phase 20 claims v1.7.0"));
}
#[test]
fn render_planning_doc_findings_json_is_an_array_of_objects() {
let findings = vec![PlanningDocFinding {
source: "ROADMAP.md phase 20".to_string(),
claim: "ROADMAP.md phase 20 claims v1.7.0".to_string(),
severity: Severity::Problem,
detail: "detail text".to_string(),
repair: None,
}];
let value = render_planning_doc_findings_json(&findings);
assert!(value.is_array());
let arr = value.as_array().unwrap();
assert_eq!(arr.len(), 1);
assert_eq!(arr[0]["severity"], "problem");
assert_eq!(arr[0]["source"], "ROADMAP.md phase 20");
assert_eq!(arr[0]["repair"], serde_json::Value::Null);
}
#[test]
fn doctor_json_body_carries_planning_doc_staleness_as_a_third_key() {
let checks: Vec<Check> = Vec::new();
let facts: Vec<PhaseFacts> = Vec::new();
let doc_findings = vec![PlanningDocFinding {
source: "ROADMAP.md phase 20".to_string(),
claim: "claim".to_string(),
severity: Severity::Problem,
detail: "detail".to_string(),
repair: None,
}];
let body = doctor_json_body(&checks, &facts, &doc_findings, &[]);
let obj = body.as_object().unwrap();
assert_eq!(
obj.len(),
4,
"must be exactly {{environment, reconciliation, planning_doc_staleness, \
stray_processes}}: {body}"
);
assert!(obj.contains_key("environment"));
assert!(obj.contains_key("reconciliation"));
let staleness = obj["planning_doc_staleness"].as_array().unwrap();
assert_eq!(staleness.len(), 1);
assert!(obj.contains_key("stray_processes"));
}
}
#[test]
fn evidence_require_shipped_exits_ok_iff_the_phase_has_shipped() {
let dir = tempfile::tempdir().unwrap();
events::emit(
dir.path(),
PhaseId::new(30),
"workflow_shipped",
serde_json::json!({"stage": "ship"}),
);
assert!(evidence(dir.path(), PhaseId::new(30), false, true).is_ok());
assert!(evidence(dir.path(), PhaseId::new(31), false, true).is_err());
}
#[test]
fn evidence_require_shipped_names_stopped_at_rather_than_generic_not_shipped() {
let dir = tempfile::tempdir().unwrap();
events::emit(
dir.path(),
PhaseId::new(32),
"workflow_finished",
serde_json::json!({"reason": "stopped_at", "stage": "plan"}),
);
let err = evidence(dir.path(), PhaseId::new(32), false, true).unwrap_err();
let message = err.to_string();
assert!(
message.contains("stopped"),
"message must name the stopped-at case, got: {message}"
);
}
#[test]
fn evidence_require_shipped_failure_message_is_single_line_and_names_phase() {
let dir = tempfile::tempdir().unwrap();
let err = evidence(dir.path(), PhaseId::new(33), false, true).unwrap_err();
let message = err.to_string();
assert!(
!message.contains('\n'),
"message must be one line: {message:?}"
);
assert!(
message.contains("33"),
"message must name the phase: {message}"
);
}
#[test]
fn changelog_version_check_flags_mismatch_and_passes_on_agreement() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path();
let write_workspace = |version: &str| {
std::fs::write(
root.join("Cargo.toml"),
format!("[workspace.package]\nversion = \"{version}\"\n"),
)
.unwrap();
};
write_workspace("2.5.0");
std::fs::write(root.join("CHANGELOG.md"), "## 2.4.0 — 2026-01-01\n").unwrap();
let mismatch = check_changelog_version(root);
assert_eq!(mismatch.status, "fail");
assert!(
mismatch.version.unwrap().contains("Cargo.toml ahead"),
"workspace 2.5.0 is newer than changelog 2.4.0"
);
write_workspace("2.5.0");
std::fs::write(root.join("CHANGELOG.md"), "## 2.6.0 — 2026-01-01\n").unwrap();
let reverse = check_changelog_version(root);
assert_eq!(reverse.status, "fail");
assert!(
reverse.version.unwrap().contains("changelog ahead"),
"changelog 2.6.0 is newer than workspace 2.5.0"
);
write_workspace("2.5.0");
std::fs::write(root.join("CHANGELOG.md"), "## 2.5.0 — 2026-08-15\n").unwrap();
assert_eq!(check_changelog_version(root).status, "ok");
std::fs::write(root.join("CHANGELOG.md"), "## [2.5.0] - 2026-08-15\n").unwrap();
assert_eq!(check_changelog_version(root).status, "ok");
std::fs::write(root.join("CHANGELOG.md"), "no heading here\n").unwrap();
assert_eq!(check_changelog_version(root).status, "warn");
}
}