#[derive(Debug, Clone)]
struct AgentSpawnOutcome {
status: std::process::ExitStatus,
timed_out: bool,
interrupted: bool,
timeout_secs: Option<u64>,
usage_capture_path: Option<PathBuf>,
}
#[cfg(not(test))]
const AGENT_OUTPUT_DRAIN_GRACE: Duration = Duration::from_millis(100);
#[cfg(test)]
const AGENT_OUTPUT_DRAIN_GRACE: Duration = Duration::from_millis(20);
impl InvocationOutcome for AgentSpawnOutcome {
fn was_interrupted(&self) -> bool {
self.interrupted
}
fn timed_out(&self) -> bool {
self.timed_out
}
fn status(&self) -> std::process::ExitStatus {
self.status
}
}
fn with_agent_log<T>(
log_file: &Arc<Mutex<fs::File>>,
write: impl FnOnce(&mut fs::File) -> std::io::Result<T>,
) -> std::io::Result<T> {
let mut guard = log_file
.lock()
.map_err(|_| std::io::Error::new(std::io::ErrorKind::Other, "agent log lock poisoned"))?;
write(&mut guard)
}
fn output_line(buf: &[u8]) -> String {
let line = buf.strip_suffix(b"\n").unwrap_or(buf);
let line = line.strip_suffix(b"\r").unwrap_or(line);
String::from_utf8_lossy(line).into_owned()
}
fn agent_stream_label(stream: rhei_tui::AgentStream) -> &'static str {
match stream {
rhei_tui::AgentStream::Stdout => "stdout",
rhei_tui::AgentStream::Stderr => "stderr",
}
}
fn spawn_agent_output_reader<R>(
reader: R,
stream: rhei_tui::AgentStream,
log_file: Arc<Mutex<fs::File>>,
sink: Arc<dyn rhei_tui::EventSink>,
slot: rhei_tui::Slot,
task_id: String,
usage_capture: Option<AgentUsageCapture>,
) -> std::thread::JoinHandle<std::io::Result<()>>
where
R: Read + Send + 'static,
{
std::thread::spawn(move || {
let mut reader = BufReader::new(reader);
let mut buf = Vec::new();
loop {
buf.clear();
let read = reader.read_until(b'\n', &mut buf)?;
if read == 0 {
break;
}
with_agent_log(&log_file, |f| {
f.write_all(&buf)?;
f.flush()
})?;
let raw_line = output_line(&buf);
capture_agent_output_usage(usage_capture.as_ref(), stream, &raw_line, &sink);
if let Some(line) =
display_agent_output_line(usage_capture.as_ref(), stream, &raw_line)
{
sink.emit(rhei_tui::RunEvent::AgentOutput {
slot,
task: task_id.clone(),
stream,
line,
wall_clock: std::time::SystemTime::now(),
});
}
}
Ok(())
})
}
fn drain_agent_output_reader(
handle: std::thread::JoinHandle<std::io::Result<()>>,
stream: rhei_tui::AgentStream,
) -> MietteResult<()> {
let deadline = Instant::now() + AGENT_OUTPUT_DRAIN_GRACE;
while !handle.is_finished() {
if Instant::now() >= deadline {
return Ok(());
}
std::thread::sleep(Duration::from_millis(1));
}
match handle.join() {
Ok(Ok(())) => Ok(()),
Ok(Err(err)) => {
Err(miette!(
help = agent_command_help(),
"failed to capture agent {}: {err}", agent_stream_label(stream)
))
}
Err(_) => Err(miette!(
help = internal_error_help(),
"agent {} capture thread panicked", agent_stream_label(stream)
)),
}
}
#[allow(clippy::too_many_arguments)]
fn spawn_and_wait_agent(
resolved: &ResolvedAgent,
prompt: &str,
rhei_root: &Path,
checkout_root: &Path,
worktree_root: Option<&Path>,
plan_path: &Path,
state_machine_path: Option<&Path>,
task_id: &str,
state_name: &str,
visit_count: u64,
tooling: &ResolvedTooling,
log_path: &Path,
runtime_dir: &Path,
snapshot_preload: Option<&SnapshotPreload>,
slot: rhei_tui::Slot,
sink: Arc<dyn rhei_tui::EventSink>,
intervene: Option<&Arc<RunInterveneSink>>,
result_identity: Option<&str>,
) -> MietteResult<AgentSpawnOutcome> {
if let Some(parent) = log_path.parent() {
fs::create_dir_all(parent)
.map_err(|e| miette!(
help = agent_log_help(),
"failed to create log directory '{}': {e}", parent.display()
))?;
}
let log_file = Arc::new(Mutex::new(
fs::File::create(log_path)
.map_err(|e| miette!(
help = agent_log_help(),
"failed to create log file '{}': {e}", log_path.display()
))?,
));
let started_wall = std::time::SystemTime::now();
with_agent_log(&log_file, |f| {
writeln!(f, "=== rhei agent log v1 ===")?;
writeln!(f, "agent: {}", resolved.agent.id())?;
if let Some(mode) = &resolved.mode {
writeln!(f, "mode: {mode}")?;
}
if let Some(target) = &resolved.target {
writeln!(f, "target: {}", target.selector())?;
}
if let Some(m) = &resolved.model {
writeln!(f, "model: {m}")?;
}
if let Some(provider) = &resolved.model_provider {
writeln!(f, "provider: {provider}")?;
}
if let Some(model_name) = &resolved.model_name {
writeln!(f, "model_name: {model_name}")?;
}
writeln!(f, "task: {task_id}")?;
writeln!(f, "state: {state_name}")?;
writeln!(f, "started: {}", format_iso8601_utc(started_wall))?;
if let Some(t) = resolved.timeout_secs {
writeln!(f, "timeout: {}", format_duration_human(t))?;
}
writeln!(f, "plan: {}", plan_path.display())?;
writeln!(f, "rhei_root: {}", rhei_root.display())?;
writeln!(f, "checkout_root: {}", checkout_root.display())?;
if let Some(path) = worktree_root {
writeln!(f, "worktree_root: {}", path.display())?;
}
let mcp_line = format_tooling_log_line(&tooling.mcp_servers, |e| {
(e.id.as_str(), e.optional, e.definition.is_some())
});
if let Some(line) = mcp_line {
writeln!(f, "mcp_servers: {line}")?;
}
let skill_line = format_tooling_log_line(&tooling.skills, |e| {
(e.id.as_str(), e.optional, e.definition.is_some())
});
if let Some(line) = skill_line {
writeln!(f, "skills: {line}")?;
}
writeln!(f, "===\n")?;
f.flush()
})
.map_err(|e| miette!(
help = agent_log_help(),
"failed to write log header '{}': {e}", log_path.display()
))?;
for warning in collect_unsupported_tooling_warnings(resolved, tooling) {
let _ = with_agent_log(&log_file, |f| writeln!(f, "{warning}"));
diag_warn!("{warning}");
}
let usage_capture_path =
accounting_capture_path_for_spawn(runtime_dir, task_id, state_name, resolved);
if let Some(parent) = usage_capture_path.as_ref().and_then(|path| path.parent()) {
let _ = fs::create_dir_all(parent);
}
let usage_capture = usage_capture_for_spawn(
resolved,
usage_capture_path.as_deref(),
task_id,
state_name,
visit_count,
slot,
);
let mut cmd = build_agent_command(
resolved,
prompt,
rhei_root,
checkout_root,
worktree_root,
plan_path,
state_machine_path,
task_id,
state_name,
visit_count,
tooling,
runtime_dir,
result_identity,
);
configure_accounting_capture(&mut cmd, usage_capture_path.as_deref());
if let Some(snapshot_preload) = snapshot_preload {
for arg in &snapshot_preload.extra_args {
cmd.arg(arg);
}
if let Some(session_dir) = snapshot_preload.session_dir.as_ref() {
cmd.env("RHEI_SNAPSHOT_SESSION_DIR", session_dir);
}
if let Some(parent_ref) = snapshot_preload.parent_ref.as_ref() {
cmd.env("RHEI_SNAPSHOT_PARENT_REF", parent_ref.to_string());
}
}
cmd.stdout(std::process::Stdio::piped()).stderr(std::process::Stdio::piped());
let mut supervised = match Supervised::spawn(&mut cmd, &format!("{task_id}@{state_name}")) {
Ok(supervised) => supervised,
Err(err) if spawn_was_interrupted(&err) => {
let _ = with_agent_log(&log_file, |f| {
writeln!(f, "\nagent not started: the run was interrupted first")?;
writeln!(f, "\n=== exit ===")?;
writeln!(f, "code: -")?;
writeln!(f, "ended: {}", format_iso8601_utc(std::time::SystemTime::now()))?;
writeln!(f, "interrupted: true")?;
writeln!(f, "===")?;
f.flush()
});
return Ok(AgentSpawnOutcome {
status: never_started_status(),
timed_out: false,
interrupted: true,
timeout_secs: resolved.timeout_secs,
usage_capture_path: None,
});
}
Err(e) => {
return Err(miette!(
help = "the agent command could not start. Check it exists on PATH and is executable: rhei diag",
"failed to spawn agent '{}': {e}", resolved.agent.id()
))
}
};
let child = &mut supervised.child;
let stdout_handle = child.stdout.take().map(|stdout| {
spawn_agent_output_reader(
stdout,
rhei_tui::AgentStream::Stdout,
log_file.clone(),
sink.clone(),
slot,
task_id.to_string(),
usage_capture.clone(),
)
});
let stderr_handle = child.stderr.take().map(|stderr| {
spawn_agent_output_reader(
stderr,
rhei_tui::AgentStream::Stderr,
log_file.clone(),
sink.clone(),
slot,
task_id.to_string(),
None,
)
});
let mut registered_intervene = false;
let stdin_format = agent_stdin_format(resolved);
if stdin_format == AgentStdinFormat::ClaudeCodeStreamJson {
if let Some(mut stdin) = child.stdin.take() {
use std::io::Write as _;
let _ = stdin.write_all(&stdin_message_bytes(stdin_format, prompt));
let _ = stdin.flush();
match (resolved.profile.intervene_stdin, intervene) {
(true, Some(registry)) => {
registry.register(
task_id,
slot,
state_name,
log_file.clone(),
stdin,
stdin_format,
);
registered_intervene = true;
}
_ => drop(stdin),
}
}
} else if resolved.profile.stdin_prompt {
if let Some(mut stdin) = child.stdin.take() {
use std::io::Write as _;
let _ = stdin.write_all(prompt.as_bytes());
let _ = stdin.flush();
match (resolved.profile.intervene_stdin, intervene) {
(true, Some(registry)) => {
registry.register(
task_id,
slot,
state_name,
log_file.clone(),
stdin,
stdin_format,
);
registered_intervene = true;
}
_ => drop(stdin),
}
}
} else if resolved.profile.intervene_stdin {
if let Some(stdin) = child.stdin.take() {
if let Some(registry) = intervene {
registry.register(task_id, slot, state_name, log_file.clone(), stdin, stdin_format);
registered_intervene = true;
} else {
drop(stdin);
}
}
}
let start = Instant::now();
let ended = supervised
.wait(
resolved.timeout_secs.map(Duration::from_secs),
&INTERRUPT,
¬ify_through_sink(&sink),
)
.map_err(|e| miette!(
help = agent_command_help(),
"error waiting for agent: {e}"
))?;
let status = ended.status;
let timed_out = ended.cause == EndCause::TimedOut;
let interrupted = ended.cause == EndCause::Interrupted;
if registered_intervene {
if let Some(registry) = intervene {
registry.unregister(task_id, slot);
}
}
if let Some(handle) = stdout_handle {
drain_agent_output_reader(handle, rhei_tui::AgentStream::Stdout)?;
}
if let Some(handle) = stderr_handle {
drain_agent_output_reader(handle, rhei_tui::AgentStream::Stderr)?;
}
let elapsed = start.elapsed();
let timeout_message =
if timed_out { resolved.timeout_secs.map(format_duration_human) } else { None };
with_agent_log(&log_file, |f| {
if let Some(duration) = &timeout_message {
writeln!(f, "\nagent timed out after {duration}")?;
} else if interrupted {
writeln!(
f,
"\nagent interrupted by run shutdown after {}",
format_duration_human(elapsed.as_secs())
)?;
}
writeln!(f, "\n=== exit ===")?;
writeln!(f, "code: {}", status.code().unwrap_or(-1))?;
writeln!(f, "duration: {}", format_duration_human(elapsed.as_secs()))?;
writeln!(f, "ended: {}", format_iso8601_utc(std::time::SystemTime::now()))?;
if timed_out {
writeln!(f, "timed_out: true")?;
}
if interrupted {
writeln!(f, "interrupted: true")?;
}
writeln!(f, "===")?;
f.flush()
})
.map_err(|e| miette!(
help = agent_log_help(),
"failed to append to log file '{}': {e}", log_path.display()
))?;
Ok(AgentSpawnOutcome {
status,
timed_out,
interrupted,
timeout_secs: resolved.timeout_secs,
usage_capture_path,
})
}