use std::time::{Duration, Instant};
use serde::Serialize;
use serde_json::Value;
use octl_core::{read_manifest_opt, read_node_opt, NodeId, RunLock, RunPaths, Status};
use crate::error::CliError;
use crate::output::{self, OutputFormat, OutputSpec};
use crate::run::dto::SupervisorView;
use crate::run::stalled::{stall_kind, StallKind};
use crate::run::{from_core, run_paths_from_cli_arg, status_kebab};
const DEFAULT_NODE_ID: &str = "n-0001";
const BACKOFF_START: Duration = Duration::from_millis(100);
const BACKOFF_CAP: Duration = Duration::from_secs(2);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Condition {
All,
Any,
}
impl Condition {
fn wire(self) -> &'static str {
match self {
Condition::All => "all",
Condition::Any => "any",
}
}
}
pub struct Args<'a> {
pub run_ids: Vec<String>,
pub any: bool,
pub timeout: Option<Duration>,
pub fail_on_error: bool,
pub progress: bool,
pub poll_interval: Option<Duration>,
pub spec: &'a OutputSpec,
pub warnings: &'a [String],
}
#[derive(Serialize)]
struct RunOutcome {
run_id: String,
status: &'static str,
merged: bool,
landed: bool,
landed_method: &'static str,
stalled: bool,
#[serde(skip_serializing_if = "Option::is_none")]
summary: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
error: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
recoverable_work: Option<Value>,
}
#[derive(Serialize)]
struct WaitData {
waited_ms: u64,
condition: &'static str,
runs: Vec<RunOutcome>,
}
enum Stop {
Met,
TimedOut,
}
pub fn run(args: Args<'_>) -> Result<(), CliError> {
let condition = if args.any {
Condition::Any
} else {
Condition::All
};
let root = crate::home::root_dir()?;
let mut runs: Vec<(String, RunPaths)> = Vec::with_capacity(args.run_ids.len());
for run_id in &args.run_ids {
let paths = run_paths_from_cli_arg(&root, run_id)?;
if current_status(&paths)?.is_none() {
return Err(
CliError::user("unknown_run", format!("no run with id {run_id}"))
.with_invalid_value(run_id),
);
}
runs.push((run_id.clone(), paths));
}
let start = Instant::now();
let (stop, latched_stall) = wait_loop(
&runs,
condition,
args.timeout,
args.poll_interval,
args.progress,
)?;
let waited_ms = start.elapsed().as_millis() as u64;
let mut outcomes = Vec::with_capacity(runs.len());
for (i, (run_id, paths)) in runs.iter().enumerate() {
outcomes.push(read_outcome(run_id, paths, latched_stall[i])?);
}
let exit_code: u8 = match stop {
Stop::TimedOut => 2,
Stop::Met if args.fail_on_error && any_settled_error(&outcomes) => 3,
Stop::Met => 0,
};
let data = WaitData {
waited_ms,
condition: condition.wire(),
runs: outcomes,
};
emit(&data, args.spec, args.warnings)?;
if exit_code == 0 {
return Ok(());
}
crate::cli::flush_logs();
std::process::exit(i32::from(exit_code));
}
fn wait_loop(
runs: &[(String, RunPaths)],
condition: Condition,
timeout: Option<Duration>,
poll_interval: Option<Duration>,
progress: bool,
) -> Result<(Stop, Vec<Option<StallKind>>), CliError> {
let start = Instant::now();
let mut backoff = poll_interval.unwrap_or(BACKOFF_START);
let mut prev: Vec<Option<(Status, Option<StallKind>)>> = vec![None; runs.len()];
loop {
let mut settled = 0usize;
let mut stall_now: Vec<Option<StallKind>> = vec![None; runs.len()];
let now = chrono::Utc::now();
for (i, (run_id, paths)) in runs.iter().enumerate() {
let settle = current_settle(paths, now)?.ok_or_else(|| {
CliError::system(
"io_error",
format!("run {run_id} manifest disappeared while waiting"),
)
})?;
stall_now[i] = settle.stall;
if settle.status.is_terminal() || settle.stall.is_some() {
settled += 1;
}
let key = (settle.status, settle.stall);
if progress && prev[i] != Some(key) {
emit_progress(run_id, settle.status, settle.stall.is_some());
}
prev[i] = Some(key);
}
let met = match condition {
Condition::All => settled == runs.len(),
Condition::Any => settled >= 1,
};
if met {
return Ok((Stop::Met, stall_now));
}
if let Some(t) = timeout {
if start.elapsed() >= t {
return Ok((Stop::TimedOut, stall_now));
}
}
let sleep_dur = match timeout {
Some(t) => backoff.min(t.saturating_sub(start.elapsed())),
None => backoff,
};
if sleep_dur.is_zero() {
std::thread::sleep(Duration::from_millis(1));
} else {
std::thread::sleep(sleep_dur);
}
if poll_interval.is_none() {
backoff = (backoff * 2).min(BACKOFF_CAP);
}
}
}
fn current_status(paths: &RunPaths) -> Result<Option<Status>, CliError> {
RunLock::with_shared_lock(&paths.lock(), || {
Ok(read_manifest_opt(paths)?.map(|m| m.status))
})
.map_err(from_core)
}
struct Settle {
status: Status,
stall: Option<StallKind>,
}
fn current_settle(
paths: &RunPaths,
now: chrono::DateTime<chrono::Utc>,
) -> Result<Option<Settle>, CliError> {
RunLock::with_shared_lock(&paths.lock(), || {
let Some(m) = read_manifest_opt(paths)? else {
return Ok(None);
};
let supervisor = SupervisorView::probe(paths);
Ok(Some(Settle {
status: m.status,
stall: stall_kind(
m.status,
supervisor.presumed_working(),
m.node_count,
m.created_at,
m.updated_at,
now,
),
}))
})
.map_err(from_core)
}
fn read_outcome(
run_id: &str,
paths: &RunPaths,
latched_stall: Option<StallKind>,
) -> Result<RunOutcome, CliError> {
let node_id = NodeId::parse_str(DEFAULT_NODE_ID).expect("DEFAULT_NODE_ID is a valid node id");
let git_inputs = RunLock::with_shared_lock(&paths.lock(), || {
let manifest = read_manifest_opt(paths)?;
let node = read_node_opt(paths, &node_id)?;
Ok(GitInputs {
status: manifest.as_ref().map(|m| m.status),
source_repo: manifest.as_ref().and_then(|m| m.source_repo.clone()),
source_branch: manifest.as_ref().and_then(|m| m.source_branch.clone()),
worktree_path: node.as_ref().and_then(|n| n.worktree_path.clone()),
branch: node.as_ref().and_then(|n| n.branch.clone()),
base_sha: node.as_ref().and_then(|n| n.base_sha.clone()),
report: node.and_then(|n| n.last_report),
})
})
.map_err(from_core)?;
let status = git_inputs.status.ok_or_else(|| {
CliError::system(
"io_error",
format!("run {run_id} manifest disappeared while waiting"),
)
})?;
let report = git_inputs.report.clone();
let signal = crate::run::landed::landing_signal(
&crate::run::landed::LandingInputs {
source_repo: git_inputs.source_repo.as_deref(),
source_branch: git_inputs.source_branch.as_deref(),
worktree_path: git_inputs.worktree_path.as_deref(),
branch: git_inputs.branch.as_deref(),
base_sha: git_inputs.base_sha.as_deref(),
report: report.as_ref(),
},
&crate::supervise::cleanup::git_bin(),
);
let merged = report
.as_ref()
.and_then(|r| r.get("via"))
.and_then(Value::as_str)
== Some("explicit-merge");
let summary = report
.as_ref()
.and_then(|r| r.get("summary"))
.and_then(Value::as_str)
.map(str::to_string);
let stall = if status.is_terminal() {
None
} else {
latched_stall
};
let stalled = stall.is_some();
let error = if let Some(kind) = stall {
Some(stall_reason(kind).to_string())
} else if matches!(status, Status::Failed | Status::Cancelled) {
report
.as_ref()
.and_then(|r| r.get("reason"))
.and_then(Value::as_str)
.map(str::to_string)
} else {
None
};
let recoverable_work = if matches!(status, Status::Failed) {
report
.as_ref()
.and_then(|r| r.get("recoverable_work"))
.filter(|v| v.is_object())
.cloned()
} else {
None
};
Ok(RunOutcome {
run_id: run_id.to_string(),
status: status_kebab(status),
merged,
landed: signal.landed,
landed_method: signal.method.wire(),
stalled,
summary,
error,
recoverable_work,
})
}
fn stall_reason(kind: StallKind) -> &'static str {
match kind {
StallKind::Stillborn => "supervisor died before creating any worker node",
StallKind::Orphaned => "supervisor died mid-run; work is stranded and cannot be rolled up",
}
}
struct GitInputs {
status: Option<Status>,
source_repo: Option<String>,
source_branch: Option<String>,
worktree_path: Option<String>,
branch: Option<String>,
base_sha: Option<String>,
report: Option<Value>,
}
fn any_settled_error(outcomes: &[RunOutcome]) -> bool {
outcomes
.iter()
.any(|o| matches!(o.status, "failed" | "cancelled") || o.stalled)
}
fn emit_progress(run_id: &str, status: Status, stalled: bool) {
if let Ok(line) = serde_json::to_string(&serde_json::json!({
"run_id": run_id,
"status": status_kebab(status),
"stalled": stalled,
})) {
eprintln!("{line}");
}
}
pub(crate) fn recoverable_summary(block: Option<&Value>) -> Option<String> {
let obj = block?.as_object()?;
let unmerged = obj.get("unmerged_commits").and_then(Value::as_u64)?;
let recoverable = obj
.get("recoverable")
.and_then(Value::as_bool)
.unwrap_or(false);
let branch = obj.get("branch").and_then(Value::as_str).unwrap_or("?");
let (noun, merge_verb, not_merge_verb) = if unmerged == 1 {
("commit", "merges", "does NOT merge")
} else {
("commits", "merge", "do NOT merge")
};
if recoverable {
Some(format!(
"recoverable={unmerged} unmerged {noun} {merge_verb} cleanly on {branch}"
))
} else {
Some(format!(
"recoverable=false ({unmerged} unmerged {noun} on {branch} {not_merge_verb} cleanly)"
))
}
}
fn emit(data: &WaitData, spec: &OutputSpec, warnings: &[String]) -> Result<(), CliError> {
match spec.format {
OutputFormat::Json | OutputFormat::Jsonl => {
output::emit_envelope(data, spec, warnings)?;
}
OutputFormat::Text => {
println!("condition: {}", data.condition);
println!("waited_ms: {}", data.waited_ms);
for r in &data.runs {
print!(
"{} status={} landed={} ({}) merged={}",
r.run_id, r.status, r.landed, r.landed_method, r.merged
);
if r.stalled {
print!(
" stalled=true (`run reattach {id}` or `run cancel {id}`)",
id = r.run_id
);
}
if let Some(s) = &r.summary {
print!(" summary={}", output::escape_one_line(s));
}
if let Some(e) = &r.error {
print!(" error={}", output::escape_one_line(e));
}
if let Some(line) = recoverable_summary(r.recoverable_work.as_ref()) {
print!(" {line}");
}
println!();
}
output::emit_text_warnings(warnings);
}
}
Ok(())
}
pub fn parse_duration(s: &str) -> Result<Duration, String> {
let trimmed = s.trim();
if !trimmed.is_empty() && trimmed.bytes().all(|b| b.is_ascii_digit()) {
let secs = trimmed
.parse::<u64>()
.map_err(|_| format!("invalid duration '{s}': bare seconds exceed {}", u64::MAX))?;
return Ok(Duration::from_secs(secs));
}
humantime::parse_duration(trimmed).map_err(|e| format!("invalid duration '{s}': {e}"))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_duration_accepts_human_units() {
assert_eq!(parse_duration("30s").unwrap(), Duration::from_secs(30));
assert_eq!(parse_duration("5m").unwrap(), Duration::from_secs(300));
assert_eq!(parse_duration("1h").unwrap(), Duration::from_secs(3600));
assert_eq!(
parse_duration("2400sec").unwrap(),
Duration::from_secs(2400)
);
assert_eq!(parse_duration("40min").unwrap(), Duration::from_secs(2400));
assert_eq!(parse_duration("500ms").unwrap(), Duration::from_millis(500));
}
#[test]
fn parse_duration_bare_integer_is_seconds() {
assert_eq!(parse_duration("2400").unwrap(), Duration::from_secs(2400));
assert_eq!(parse_duration("0").unwrap(), Duration::from_secs(0));
assert_eq!(parse_duration("00030").unwrap(), Duration::from_secs(30));
assert_eq!(
parse_duration("30").unwrap(),
parse_duration("30sec").unwrap()
);
assert_eq!(parse_duration(" 30 ").unwrap(), Duration::from_secs(30));
assert_eq!(parse_duration(" 30s ").unwrap(), Duration::from_secs(30));
}
#[test]
fn parse_duration_bare_integer_overflow_is_a_clear_error() {
let err = parse_duration("18446744073709551616").unwrap_err();
assert!(err.contains("bare seconds exceed"), "got: {err}");
}
#[test]
fn parse_duration_rejects_garbage() {
assert!(parse_duration("soon").is_err());
assert!(parse_duration("").is_err());
assert!(parse_duration(" ").is_err());
assert!(parse_duration("-5s").is_err());
assert!(parse_duration("-5").is_err());
assert!(parse_duration("+30").is_err());
assert!(parse_duration("2 400").is_err());
}
#[test]
fn any_settled_error_only_counts_terminal_failures() {
let mk = |status: &'static str| RunOutcome {
run_id: "r".into(),
status,
merged: false,
landed: false,
landed_method: "unverified",
stalled: false,
summary: None,
error: None,
recoverable_work: None,
};
let mk_stalled = || RunOutcome {
stalled: true,
..mk("pending")
};
assert!(!any_settled_error(&[mk("done"), mk("done")]));
assert!(any_settled_error(&[mk("done"), mk("failed")]));
assert!(any_settled_error(&[mk("cancelled")]));
assert!(!any_settled_error(&[mk("done"), mk("running")]));
assert!(any_settled_error(&[mk_stalled()]));
assert!(any_settled_error(&[mk("done"), mk_stalled()]));
}
#[test]
fn recoverable_summary_phrasing() {
assert_eq!(recoverable_summary(None), None);
assert_eq!(recoverable_summary(Some(&serde_json::json!("x"))), None);
let clean = serde_json::json!({
"recoverable": true,
"unmerged_commits": 1,
"merges_cleanly": true,
"branch": "wt/foo",
});
assert_eq!(
recoverable_summary(Some(&clean)).as_deref(),
Some("recoverable=1 unmerged commit merges cleanly on wt/foo"),
);
let many = serde_json::json!({
"recoverable": true, "unmerged_commits": 3, "merges_cleanly": true, "branch": "wt/bar",
});
assert_eq!(
recoverable_summary(Some(&many)).as_deref(),
Some("recoverable=3 unmerged commits merge cleanly on wt/bar"),
);
let dirty = serde_json::json!({
"recoverable": false, "unmerged_commits": 2, "merges_cleanly": false, "branch": "wt/baz",
});
assert_eq!(
recoverable_summary(Some(&dirty)).as_deref(),
Some("recoverable=false (2 unmerged commits on wt/baz do NOT merge cleanly)"),
);
let dirty1 = serde_json::json!({
"recoverable": false, "unmerged_commits": 1, "merges_cleanly": false, "branch": "wt/q",
});
assert_eq!(
recoverable_summary(Some(&dirty1)).as_deref(),
Some("recoverable=false (1 unmerged commit on wt/q does NOT merge cleanly)"),
);
}
}