use std::borrow::Cow;
use std::io::{BufRead, BufReader, Read, Write};
use std::os::unix::net::UnixStream;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
pub const ARGV_OVERFLOW_THRESHOLD: usize = 200 * 1024;
pub const STDOUT_HEAD_LIMIT: usize = 200;
pub const TERMINAL_STATES: [&str; 4] = ["done", "completed", "failed", "needs-input"];
const LIVENESS_PROBE_TIMEOUT: Duration = Duration::from_millis(250);
const SEND_SOCKET_TIMEOUT: Duration = Duration::from_secs(5);
const RETRY_BACKOFF: Duration = Duration::from_millis(10);
const SPAWN_UUID_RETRY_ATTEMPTS: u32 = 6;
const SPAWN_UUID_RETRY_BACKOFF: Duration = Duration::from_millis(300);
fn is_terminal_state(state: &str) -> bool {
TERMINAL_STATES.contains(&state)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OrphanReason {
SocketNull,
NotFound,
LivenessFailed,
RosterLiveInjectFailed,
TruthLiveInjectFailed,
}
impl OrphanReason {
pub fn as_str(self) -> &'static str {
match self {
OrphanReason::SocketNull => "socket-null",
OrphanReason::NotFound => "not-found",
OrphanReason::LivenessFailed => "liveness-failed",
OrphanReason::RosterLiveInjectFailed => "roster-live-inject-failed",
OrphanReason::TruthLiveInjectFailed => "truth-live-inject-failed",
}
}
}
pub fn family1_truth_state(handle: &str) -> Option<String> {
let mut command = std::process::Command::new("fno");
command
.args(["agents", "truth", handle, "--json"])
.env("FNO_AGENTS_RUNTIME", "python");
family1_truth_state_with_command(command, Duration::from_secs(5), handle)
}
fn family1_truth_failure_detail(stdout: &[u8], stderr: &str) -> String {
let reason = serde_json::from_slice::<serde_json::Value>(stdout)
.ok()
.and_then(|value| value.get("reason")?.as_str().map(str::to_owned));
reason.unwrap_or_else(|| stderr.trim().to_owned())
}
const TRUTH_NOT_FOUND: &str = "not-found";
fn truth_failure_is_routine(detail: &str) -> bool {
detail.trim() == TRUTH_NOT_FOUND
}
fn family1_truth_state_with_command(
mut command: std::process::Command,
timeout: Duration,
handle: &str,
) -> Option<String> {
command
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped());
let mut child = match command.spawn() {
Ok(child) => child,
Err(error) => {
eprintln!("WARN: family-1 truth probe for {handle} failed to start: {error}");
return None;
}
};
let deadline = Instant::now() + timeout;
loop {
match child.try_wait() {
Ok(Some(_)) => break,
Ok(None) if Instant::now() < deadline => {
std::thread::sleep(Duration::from_millis(20));
}
Ok(None) => {
let _ = child.kill();
let _ = child.wait();
eprintln!("WARN: family-1 truth probe for {handle} timed out");
return None;
}
Err(error) => {
let _ = child.kill();
let _ = child.wait();
eprintln!("WARN: family-1 truth probe for {handle} wait failed: {error}");
return None;
}
}
}
let output = match child.wait_with_output() {
Ok(output) => output,
Err(error) => {
eprintln!("WARN: family-1 truth probe for {handle} output failed: {error}");
return None;
}
};
if !output.status.success() {
let detail =
family1_truth_failure_detail(&output.stdout, &String::from_utf8_lossy(&output.stderr));
if !truth_failure_is_routine(&detail) {
eprintln!(
"WARN: family-1 truth probe for {handle} exited {}: {}",
output.status, detail
);
}
return None;
}
let state = serde_json::from_slice::<serde_json::Value>(&output.stdout)
.ok()
.and_then(|value| value.get("state")?.as_str().map(str::to_owned));
match state.as_deref() {
Some("done" | "watching" | "your-move" | "working" | "stalled" | "unknown") => state,
_ => {
eprintln!("WARN: family-1 truth probe for {handle} returned malformed output");
None
}
}
}
fn family1_orphan_reason(handle: &str, confirmed: OrphanReason) -> OrphanReason {
match family1_truth_state(handle).as_deref() {
Some("done" | "stalled") => confirmed,
_ => OrphanReason::TruthLiveInjectFailed,
}
}
#[derive(Debug)]
pub enum AskError {
Parse { stdout_head: String },
Subprocess { exit_code: i32, stderr: String },
Orphan {
reason: OrphanReason,
short_id: String,
},
Socket { message: String },
Timeout { elapsed_sec: f64, short_id: String },
Io { message: String },
}
impl std::fmt::Display for AskError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
AskError::Parse { stdout_head } => write!(
f,
"unable to parse short-id from claude --bg output: first {} chars: {}",
stdout_head.chars().count(),
py_repr(stdout_head),
),
AskError::Subprocess { exit_code, stderr } => {
write!(f, "claude --bg exited {}: {}", exit_code, py_repr(stderr))
}
AskError::Orphan { reason, short_id } => write!(
f,
"agent short-id {} is not reachable (reason: {})",
py_repr(short_id),
reason.as_str()
),
AskError::Socket { message } => write!(f, "{}", message),
AskError::Io { message } => write!(f, "{}", message),
AskError::Timeout {
elapsed_sec,
short_id,
} => {
write!(f, "timed out waiting for reply after {:.1}s", elapsed_sec)?;
if !short_id.is_empty() {
write!(f, " (short_id={})", short_id)?;
}
Ok(())
}
}
}
}
impl std::error::Error for AskError {}
pub fn json_string_ascii(s: &str) -> String {
let mut out = String::with_capacity(s.len() + 2);
out.push('"');
for ch in s.chars() {
match ch {
'"' => out.push_str("\\\""),
'\\' => out.push_str("\\\\"),
'\n' => out.push_str("\\n"),
'\r' => out.push_str("\\r"),
'\t' => out.push_str("\\t"),
'\u{08}' => out.push_str("\\b"),
'\u{0c}' => out.push_str("\\f"),
c if (c as u32) < 0x20 => out.push_str(&format!("\\u{:04x}", c as u32)),
c if (c as u32) < 0x7f => out.push(c),
c => {
let cp = c as u32;
if cp <= 0xffff {
out.push_str(&format!("\\u{:04x}", cp));
} else {
let v = cp - 0x10000;
let hi = 0xd800 + (v >> 10);
let lo = 0xdc00 + (v & 0x3ff);
out.push_str(&format!("\\u{:04x}\\u{:04x}", hi, lo));
}
}
}
}
out.push('"');
out
}
pub fn html_escape_quote(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for ch in s.chars() {
match ch {
'&' => out.push_str("&"),
'<' => out.push_str("<"),
'>' => out.push_str(">"),
'"' => out.push_str("""),
'\'' => out.push_str("'"),
c => out.push(c),
}
}
out
}
pub fn py_repr(s: &str) -> String {
let use_double = s.contains('\'') && !s.contains('"');
let quote = if use_double { '"' } else { '\'' };
let mut out = String::with_capacity(s.len() + 2);
out.push(quote);
for ch in s.chars() {
match ch {
'\\' => out.push_str("\\\\"),
'\n' => out.push_str("\\n"),
'\r' => out.push_str("\\r"),
'\t' => out.push_str("\\t"),
c if c == quote => {
out.push('\\');
out.push(c);
}
c if (c as u32) < 0x20 || (c as u32) == 0x7f => {
out.push_str(&format!("\\x{:02x}", c as u32))
}
c => out.push(c),
}
}
out.push(quote);
out
}
#[derive(Debug, Clone)]
pub struct ClaudeHome {
home: PathBuf,
}
impl ClaudeHome {
pub fn from_env() -> Self {
let home = std::env::var_os("HOME")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("."));
Self { home }
}
pub fn at(home: impl Into<PathBuf>) -> Self {
Self { home: home.into() }
}
pub fn sessions_dir(&self) -> PathBuf {
self.home.join(".claude").join("sessions")
}
pub fn jobs_dir_for(&self, short_id: &str) -> PathBuf {
self.home.join(".claude").join("jobs").join(short_id)
}
pub fn daemon_roster_path(&self) -> PathBuf {
if let Some(dir) = std::env::var_os(crate::claude_roster::DAEMON_DIR_ENV) {
return PathBuf::from(dir).join("roster.json");
}
self.home.join(".claude").join("daemon").join("roster.json")
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SessionLocator {
pub pid: i64,
pub short_id: String,
pub messaging_socket_path: String,
pub jobs_dir: PathBuf,
pub session_id: Option<String>,
pub cwd: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct StateSnapshot {
pub state: String,
pub updated_at: Option<String>,
pub output_result: Option<String>,
pub intent: Option<String>,
}
pub fn parse_short_id(stdout: &str) -> Result<String, AskError> {
if stdout.is_empty() {
return Err(AskError::Parse {
stdout_head: String::new(),
});
}
let first_line = stdout.split('\n').next().unwrap_or("");
if let Some(id) = match_short_id(first_line) {
return Ok(id);
}
Err(AskError::Parse {
stdout_head: head_chars(stdout, STDOUT_HEAD_LIMIT),
})
}
fn match_short_id(line: &str) -> Option<String> {
const PREFIX: &str = "backgrounded \u{b7} ";
const SEP: &str = " \u{b7} ";
let cleaned = strip_ansi_csi(line);
let rest = cleaned.strip_prefix(PREFIX)?;
let rb = rest.as_bytes();
if rb.len() < 8 {
return None;
}
if !rb[..8]
.iter()
.all(|b| b.is_ascii_hexdigit() && !b.is_ascii_uppercase())
{
return None;
}
let (hex, after) = rest.split_at(8);
if after.starts_with(SEP) {
Some(hex.to_string())
} else {
None
}
}
fn strip_ansi_csi(s: &str) -> Cow<'_, str> {
if !s.contains('\u{1b}') {
return Cow::Borrowed(s);
}
let mut out = String::with_capacity(s.len());
let mut chars = s.chars().peekable();
while let Some(c) = chars.next() {
if c == '\u{1b}' && chars.peek() == Some(&'[') {
chars.next(); while let Some(&p) = chars.peek() {
if ('\u{20}'..='\u{3f}').contains(&p) {
chars.next();
} else {
break;
}
}
if let Some(&f) = chars.peek() {
if ('\u{40}'..='\u{7e}').contains(&f) {
chars.next();
}
}
continue;
}
out.push(c);
}
Cow::Owned(out)
}
fn head_chars(s: &str, n: usize) -> String {
s.chars().take(n).collect()
}
#[derive(Clone, Copy, Default)]
pub struct HarnessFlags<'a> {
pub add_dir: Option<&'a str>,
pub agent: Option<&'a str>,
pub allowed_tools: Option<&'a str>,
pub disallowed_tools: Option<&'a str>,
}
impl<'a> HarnessFlags<'a> {
fn push_onto(&self, argv: &mut Vec<String>) {
for (flag, value) in [
("--add-dir", self.add_dir),
("--agent", self.agent),
("--allowedTools", self.allowed_tools),
("--disallowedTools", self.disallowed_tools),
] {
if let Some(v) = value.filter(|v| !v.is_empty()) {
argv.push(flag.to_string());
argv.push(v.to_string());
}
}
}
}
pub fn build_argv(
name: &str,
message: &str,
use_stdin: bool,
model: Option<&str>,
permission_mode: Option<&str>,
effort: Option<&str>,
flags: HarnessFlags,
) -> Vec<String> {
let mut argv = vec![
"claude".to_string(),
"--bg".to_string(),
"--name".to_string(),
name.to_string(),
];
if let Some(m) = permission_mode.filter(|m| !m.is_empty()) {
argv.push("--permission-mode".to_string());
argv.push(m.to_string());
}
if let Some(value) = effort.filter(|v| !v.is_empty()) {
argv.push("--effort".to_string());
argv.push(value.to_string());
}
flags.push_onto(&mut argv);
if let Some(m) = model.filter(|m| !m.is_empty()) {
argv.push("--model".to_string());
argv.push(m.to_string());
}
if !use_stdin {
argv.push(message.to_string());
}
argv
}
pub fn use_stdin_for(message: &str) -> bool {
message.len() > ARGV_OVERFLOW_THRESHOLD
}
pub fn build_cross_session_container(message: &str, from_name: &str) -> String {
format!(
"<cross-session-message from-name=\"{}\">\n{}\n</cross-session-message>",
html_escape_quote(from_name),
message
)
}
pub fn build_envelope(message: &str, from_name: &str) -> Vec<u8> {
let wrapped = build_cross_session_container(message, from_name);
let content = json_string_ascii(&wrapped);
let line = format!(
"{{\"type\":\"user\",\"message\":{{\"role\":\"user\",\"content\":{}}},\"priority\":\"next\"}}\n",
content
);
line.into_bytes()
}
fn connect_unix_timeout(path: &str, timeout: Duration) -> std::io::Result<UnixStream> {
use std::os::unix::io::FromRawFd;
let c_path = std::ffi::CString::new(path).map_err(|_| {
std::io::Error::new(std::io::ErrorKind::InvalidInput, "socket path contains NUL")
})?;
unsafe {
let fd = libc::socket(libc::AF_UNIX, libc::SOCK_STREAM, 0);
if fd < 0 {
return Err(std::io::Error::last_os_error());
}
let mut addr: libc::sockaddr_un = std::mem::zeroed();
addr.sun_family = libc::AF_UNIX as libc::sa_family_t;
let bytes = c_path.as_bytes();
if bytes.len() >= std::mem::size_of_val(&addr.sun_path) {
libc::close(fd);
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"socket path too long",
));
}
std::ptr::copy_nonoverlapping(
bytes.as_ptr() as *const libc::c_char,
addr.sun_path.as_mut_ptr(),
bytes.len(),
);
let flags = libc::fcntl(fd, libc::F_GETFL, 0);
if flags < 0 || libc::fcntl(fd, libc::F_SETFL, flags | libc::O_NONBLOCK) < 0 {
let e = std::io::Error::last_os_error();
libc::close(fd);
return Err(e);
}
let addr_len = std::mem::size_of::<libc::sockaddr_un>() as libc::socklen_t;
let rc = libc::connect(fd, &addr as *const _ as *const libc::sockaddr, addr_len);
if rc != 0 {
let err = std::io::Error::last_os_error();
if err.raw_os_error() != Some(libc::EINPROGRESS) {
libc::close(fd);
return Err(err);
}
let mut pfd = libc::pollfd {
fd,
events: libc::POLLOUT,
revents: 0,
};
let ms = timeout.as_millis().min(i32::MAX as u128) as libc::c_int;
let pr = libc::poll(&mut pfd, 1, ms);
if pr < 0 {
let e = std::io::Error::last_os_error();
libc::close(fd);
return Err(e);
}
if pr == 0 {
libc::close(fd);
return Err(std::io::Error::new(
std::io::ErrorKind::TimedOut,
"connect timed out",
));
}
let mut soerr: libc::c_int = 0;
let mut len = std::mem::size_of::<libc::c_int>() as libc::socklen_t;
if libc::getsockopt(
fd,
libc::SOL_SOCKET,
libc::SO_ERROR,
&mut soerr as *mut _ as *mut libc::c_void,
&mut len,
) < 0
{
let e = std::io::Error::last_os_error();
libc::close(fd);
return Err(e);
}
if soerr != 0 {
libc::close(fd);
return Err(std::io::Error::from_raw_os_error(soerr));
}
}
if libc::fcntl(fd, libc::F_SETFL, flags) < 0 {
let e = std::io::Error::last_os_error();
libc::close(fd);
return Err(e);
}
Ok(UnixStream::from_raw_fd(fd))
}
}
pub fn send_to_session(sock_path: &str, content: &str, from_name: &str) -> Result<(), AskError> {
let payload = build_envelope(content, from_name);
let mut stream =
connect_unix_timeout(sock_path, SEND_SOCKET_TIMEOUT).map_err(|e| AskError::Socket {
message: e.to_string(),
})?;
let _ = stream.set_write_timeout(Some(SEND_SOCKET_TIMEOUT));
let _ = stream.set_read_timeout(Some(SEND_SOCKET_TIMEOUT));
let write_res = stream.write_all(&payload);
let close_res = stream
.flush()
.and_then(|_| stream.shutdown(std::net::Shutdown::Both));
if let Err(e) = write_res {
return Err(AskError::Socket {
message: e.to_string(),
});
}
if let Err(e) = close_res {
return Err(AskError::Socket {
message: format!("close after send failed: {}", e),
});
}
Ok(())
}
pub fn liveness_probe(sock_path: &str) -> bool {
connect_unix_timeout(sock_path, LIVENESS_PROBE_TIMEOUT).is_ok()
}
pub fn locate_session(home: &ClaudeHome, short_id: &str) -> Option<SessionLocator> {
let sessions = home.sessions_dir();
if !sessions.exists() {
return None;
}
let mut entries: Vec<PathBuf> = match std::fs::read_dir(&sessions) {
Ok(rd) => rd
.filter_map(|e| e.ok().map(|e| e.path()))
.filter(|p| p.extension().map(|x| x == "json").unwrap_or(false))
.collect(),
Err(_) => return None,
};
entries.sort();
for entry_path in entries {
let raw = match std::fs::read_to_string(&entry_path) {
Ok(t) => t,
Err(_) => continue,
};
let v: serde_json::Value = match serde_json::from_str(&raw) {
Ok(v) => v,
Err(_) => continue,
};
if !v.is_object() {
continue;
}
if v.get("jobId").and_then(|x| x.as_str()) != Some(short_id) {
continue;
}
if v.get("kind").and_then(|x| x.as_str()) != Some("bg") {
continue;
}
let sock = match v.get("messagingSocketPath").and_then(|x| x.as_str()) {
Some(s) if !s.is_empty() => s.to_string(),
_ => continue, };
let pid = match entry_path
.file_stem()
.and_then(|s| s.to_str())
.and_then(|s| s.parse::<i64>().ok())
{
Some(p) => p,
None => continue,
};
return Some(SessionLocator {
pid,
short_id: short_id.to_string(),
messaging_socket_path: sock,
jobs_dir: home.jobs_dir_for(short_id),
session_id: v
.get("sessionId")
.and_then(|x| x.as_str())
.map(String::from),
cwd: v.get("cwd").and_then(|x| x.as_str()).map(String::from),
});
}
None
}
pub fn resolve_session_uuid(home: &ClaudeHome, short_id: &str) -> Option<String> {
let sessions = home.sessions_dir();
if !sessions.exists() {
return None;
}
let mut entries: Vec<PathBuf> = match std::fs::read_dir(&sessions) {
Ok(rd) => rd
.filter_map(|e| e.ok().map(|e| e.path()))
.filter(|p| p.extension().map(|x| x == "json").unwrap_or(false))
.collect(),
Err(_) => return None,
};
entries.sort();
let mut fallback: Option<String> = None;
for entry_path in entries {
let raw = match std::fs::read_to_string(&entry_path) {
Ok(t) => t,
Err(_) => continue,
};
let v: serde_json::Value = match serde_json::from_str(&raw) {
Ok(v) => v,
Err(_) => continue,
};
if !v.is_object() {
continue;
}
if v.get("jobId").and_then(|x| x.as_str()) != Some(short_id) {
continue;
}
if v.get("kind").and_then(|x| x.as_str()) != Some("bg") {
continue;
}
let sid = match v.get("sessionId").and_then(|x| x.as_str()) {
Some(s) if !s.is_empty() => s.to_string(),
_ => continue,
};
match v.get("messagingSocketPath").and_then(|x| x.as_str()) {
Some(s) if !s.is_empty() => return Some(sid), _ => {
if fallback.is_none() {
fallback = Some(sid);
}
}
}
}
fallback
}
pub fn resolve_session_uuid_at_spawn(home: &ClaudeHome, short_id: &str) -> Option<String> {
if short_id.is_empty() || !home.sessions_dir().exists() {
return None;
}
for attempt in 0..SPAWN_UUID_RETRY_ATTEMPTS {
if let Some(uuid) = resolve_session_uuid(home, short_id) {
return Some(uuid);
}
if attempt + 1 < SPAWN_UUID_RETRY_ATTEMPTS {
std::thread::sleep(SPAWN_UUID_RETRY_BACKOFF);
}
}
None
}
pub fn classify_orphan_reason(home: &ClaudeHome, short_id: &str) -> OrphanReason {
let sessions = home.sessions_dir();
if let Ok(rd) = std::fs::read_dir(&sessions) {
for entry in rd.filter_map(|e| e.ok()) {
let p = entry.path();
if p.extension().map(|x| x != "json").unwrap_or(true) {
continue;
}
let raw = match std::fs::read_to_string(&p) {
Ok(t) => t,
Err(_) => continue,
};
let v: serde_json::Value = match serde_json::from_str(&raw) {
Ok(v) => v,
Err(_) => continue,
};
if v.get("jobId").and_then(|x| x.as_str()) == Some(short_id)
&& v.get("kind").and_then(|x| x.as_str()) == Some("bg")
{
let sock = v.get("messagingSocketPath").and_then(|x| x.as_str());
if sock.map(|s| s.is_empty()).unwrap_or(true) {
return OrphanReason::SocketNull;
}
}
}
}
OrphanReason::NotFound
}
#[derive(Debug)]
pub enum StateReadError {
NotFound,
Parse,
Io(std::io::Error),
}
pub fn read_state_json(jobs_dir: &Path) -> Result<StateSnapshot, StateReadError> {
let state_path = jobs_dir.join("state.json");
match parse_state(&state_path) {
Ok(s) => Ok(s),
Err(StateReadError::Parse) => {
std::thread::sleep(RETRY_BACKOFF);
parse_state(&state_path)
}
Err(e) => Err(e),
}
}
fn parse_state(state_path: &Path) -> Result<StateSnapshot, StateReadError> {
let raw_text = match std::fs::read_to_string(state_path) {
Ok(t) => t,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Err(StateReadError::NotFound),
Err(e) => return Err(StateReadError::Io(e)),
};
if raw_text.trim().is_empty() {
return Err(StateReadError::Parse);
}
let v: serde_json::Value =
serde_json::from_str(&raw_text).map_err(|_| StateReadError::Parse)?;
let output = v.get("output");
let output_result = output
.and_then(|o| if o.is_object() { o.get("result") } else { None })
.and_then(|r| r.as_str())
.map(String::from);
Ok(StateSnapshot {
state: v
.get("state")
.and_then(|x| x.as_str())
.unwrap_or("")
.to_string(),
updated_at: v
.get("updatedAt")
.and_then(|x| x.as_str())
.map(String::from),
output_result,
intent: v.get("intent").and_then(|x| x.as_str()).map(String::from),
})
}
pub fn read_timeline_tail(jobs_dir: &Path, offset: u64) -> String {
use std::io::{Seek, SeekFrom};
let timeline = jobs_dir.join("timeline.jsonl");
if !timeline.exists() {
return String::new();
}
let mut file = match std::fs::File::open(&timeline) {
Ok(f) => f,
Err(_) => return String::new(),
};
if file.seek(SeekFrom::Start(offset)).is_err() {
return String::new();
}
let mut tail = Vec::new();
if file.read_to_end(&mut tail).is_err() {
return String::new();
}
let text = match String::from_utf8(tail) {
Ok(t) => t,
Err(_) => return String::new(),
};
let mut chunks = String::new();
for line in text.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
let row: serde_json::Value = match serde_json::from_str(line) {
Ok(r) => r,
Err(_) => continue,
};
if !row.is_object() {
continue;
}
let state = row.get("state").and_then(|x| x.as_str()).unwrap_or("");
if !is_terminal_state(state) {
continue;
}
if let Some(piece) = row.get("text").and_then(|x| x.as_str()) {
if !piece.is_empty() {
chunks.push_str(piece);
}
}
}
chunks
}
pub fn timeline_offset(jobs_dir: &Path) -> u64 {
std::fs::metadata(jobs_dir.join("timeline.jsonl"))
.map(|m| m.len())
.unwrap_or(0)
}
pub const DEFAULT_POLL_INTERVAL: Duration = Duration::from_millis(500);
#[derive(Debug, Clone)]
pub struct CreateResult {
pub short_id: String,
pub stdout: String,
pub stderr: String,
pub duration_ms: u128,
}
enum ShortIdScan {
Found { short_id: String, consumed: String },
NoId { consumed: String },
}
fn scan_stdout_for_short_id<R: BufRead>(mut reader: R) -> ShortIdScan {
let mut consumed = String::new();
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line) {
Ok(0) => return ShortIdScan::NoId { consumed }, Ok(_) => {
consumed.push_str(&line);
if let Some(short_id) = match_short_id(line.trim_end()) {
return ShortIdScan::Found { short_id, consumed };
}
}
Err(_) => return ShortIdScan::NoId { consumed },
}
}
}
pub fn bg_create(
name: &str,
message: &str,
cwd: &Path,
timeout: Option<Duration>,
extra_env: &[(&str, &str)],
model: Option<&str>,
permission_mode: Option<&str>,
effort: Option<&str>,
flags: HarnessFlags,
) -> Result<CreateResult, AskError> {
use std::process::{Command, Stdio};
let use_stdin = use_stdin_for(message);
let argv = build_argv(
name,
message,
use_stdin,
model,
permission_mode,
effort,
flags,
);
let mut cmd = Command::new(&argv[0]);
cmd.args(&argv[1..]);
cmd.current_dir(cwd);
cmd.env("FNO_AGENT_SELF", name);
cmd.env("FNO_AGENT_PROVIDER", "claude");
for (k, v) in extra_env {
cmd.env(k, v);
}
cmd.stdout(Stdio::piped());
cmd.stderr(Stdio::piped());
cmd.stdin(if use_stdin {
Stdio::piped()
} else {
Stdio::null()
});
let start = std::time::Instant::now();
let mut child = match cmd.spawn() {
Ok(c) => c,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
return Err(AskError::Subprocess {
exit_code: 127,
stderr: format!("claude CLI not found: {}", e),
});
}
Err(e) => {
return Err(AskError::Subprocess {
exit_code: 127,
stderr: e.to_string(),
});
}
};
let stdin_handle = if use_stdin { child.stdin.take() } else { None };
let stdout_handle = match child.stdout.take() {
Some(h) => h,
None => {
return Err(AskError::Subprocess {
exit_code: 1,
stderr: "claude --bg exposed no stdout pipe".to_string(),
})
}
};
let stderr_handle = child.stderr.take();
let pid = child.id();
let stderr_join: Option<std::thread::JoinHandle<String>> = stderr_handle.map(|mut eh| {
std::thread::spawn(move || {
let mut s = String::new();
let _ = eh.read_to_string(&mut s);
s
})
});
if let Some(mut sin) = stdin_handle {
let msg = message.to_string();
std::thread::spawn(move || {
let _ = sin.write_all(msg.as_bytes());
});
}
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
let scan = scan_stdout_for_short_id(BufReader::new(stdout_handle));
let _ = tx.send(scan);
});
let scan = match timeout {
Some(d) => match rx.recv_timeout(d) {
Ok(s) => s,
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
unsafe {
libc::kill(pid as libc::pid_t, libc::SIGKILL);
}
let secs = d.as_secs_f64();
let secs_str = if secs.fract() == 0.0 {
format!("{}", secs as u64)
} else {
format!("{}", secs)
};
return Err(AskError::Subprocess {
exit_code: 124,
stderr: format!("claude --bg timed out after {}s", secs_str),
});
}
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
return Err(AskError::Subprocess {
exit_code: 1,
stderr: "claude --bg stdout scan thread disconnected".to_string(),
});
}
},
None => rx.recv().unwrap_or(ShortIdScan::NoId {
consumed: String::new(),
}),
};
let duration_ms = start.elapsed().as_millis();
match scan {
ShortIdScan::Found { short_id, consumed } => Ok(CreateResult {
short_id,
stdout: consumed,
stderr: String::new(),
duration_ms,
}),
ShortIdScan::NoId { consumed } => {
let exit_code = child.wait().ok().and_then(|s| s.code()).unwrap_or(1);
let stderr = stderr_join
.map(|h| h.join().unwrap_or_default())
.unwrap_or_default();
if exit_code != 0 {
Err(AskError::Subprocess { exit_code, stderr })
} else {
Err(AskError::Parse {
stdout_head: head_chars(&consumed, STDOUT_HEAD_LIMIT),
})
}
}
}
}
pub fn wait_for_reply(
jobs_dir: &Path,
baseline_updated_at: Option<&str>,
timeline_offset: u64,
timeout: Duration,
poll_interval: Duration,
short_id: &str,
) -> Result<String, AskError> {
let deadline = std::time::Instant::now() + timeout;
let final_snap = loop {
let snap = match read_state_json(jobs_dir) {
Ok(s) => Some(s),
Err(StateReadError::NotFound) | Err(StateReadError::Parse) => None,
Err(StateReadError::Io(e)) => {
return Err(AskError::Io {
message: e.to_string(),
});
}
};
if let Some(ref s) = snap {
if is_terminal_state(&s.state) {
let advanced = match (baseline_updated_at, &s.updated_at) {
(None, _) => true,
(Some(base), Some(cur)) => cur.as_str() > base,
(Some(_), None) => false,
};
if advanced {
break snap.unwrap();
}
}
}
if std::time::Instant::now() >= deadline {
return Err(AskError::Timeout {
elapsed_sec: timeout.as_secs_f64(),
short_id: short_id.to_string(),
});
}
std::thread::sleep(poll_interval);
};
match final_snap.output_result {
Some(ref r) if !r.is_empty() => Ok(r.clone()),
_ => Ok(read_timeline_tail(jobs_dir, timeline_offset)),
}
}
#[allow(clippy::too_many_arguments)]
fn roster_live(home: &ClaudeHome, short_id: &str) -> bool {
match crate::claude_roster::ClaudeRoster::load(&home.daemon_roster_path()) {
Ok(roster) => roster.find(short_id).is_some(),
Err(_) => false,
}
}
fn ask_via_control_sock(
short_id: &str,
message: &str,
from_name: &str,
timeout: Duration,
poll_interval: Duration,
target_jobs_dir: &Path,
) -> Result<String, AskError> {
let baseline_updated_at = read_state_json(target_jobs_dir)
.ok()
.and_then(|s| s.updated_at);
let offset = timeline_offset(target_jobs_dir);
let wrapped = build_cross_session_container(message, from_name);
if crate::mail_inject::deliver_via_control_sock(
short_id,
&wrapped,
crate::mail_inject::DEFAULT_ATTEMPTS,
crate::mail_inject::DEFAULT_INTERVAL_MS,
)
.is_err()
{
return Err(AskError::Orphan {
reason: OrphanReason::RosterLiveInjectFailed,
short_id: short_id.to_string(),
});
}
if !target_jobs_dir.exists() {
return Err(AskError::Timeout {
elapsed_sec: 0.0,
short_id: short_id.to_string(),
});
}
wait_for_reply(
target_jobs_dir,
baseline_updated_at.as_deref(),
offset,
timeout,
poll_interval,
short_id,
)
}
pub fn ask_followup(
home: &ClaudeHome,
claude_short_id: &str,
message: &str,
from_name: &str,
timeout: Duration,
poll_interval: Duration,
jobs_dir_override: Option<&Path>,
) -> Result<String, AskError> {
let locator = match locate_session(home, claude_short_id) {
Some(l) => l,
None => {
let reason = classify_orphan_reason(home, claude_short_id);
if reason == OrphanReason::SocketNull && roster_live(home, claude_short_id) {
let jd = jobs_dir_override
.map(|p| p.to_path_buf())
.unwrap_or_else(|| home.jobs_dir_for(claude_short_id));
return ask_via_control_sock(
claude_short_id,
message,
from_name,
timeout,
poll_interval,
&jd,
);
}
return Err(AskError::Orphan {
reason: family1_orphan_reason(claude_short_id, reason),
short_id: claude_short_id.to_string(),
});
}
};
if !liveness_probe(&locator.messaging_socket_path) {
if roster_live(home, claude_short_id) {
let jd = jobs_dir_override
.map(|p| p.to_path_buf())
.unwrap_or_else(|| locator.jobs_dir.clone());
return ask_via_control_sock(
claude_short_id,
message,
from_name,
timeout,
poll_interval,
&jd,
);
}
return Err(AskError::Orphan {
reason: family1_orphan_reason(claude_short_id, OrphanReason::LivenessFailed),
short_id: claude_short_id.to_string(),
});
}
let target_jobs_dir: PathBuf = jobs_dir_override
.map(|p| p.to_path_buf())
.unwrap_or_else(|| locator.jobs_dir.clone());
let baseline_updated_at = read_state_json(&target_jobs_dir)
.ok()
.and_then(|s| s.updated_at);
let offset = timeline_offset(&target_jobs_dir);
send_to_session(&locator.messaging_socket_path, message, from_name)?;
wait_for_reply(
&target_jobs_dir,
baseline_updated_at.as_deref(),
offset,
timeout,
poll_interval,
claude_short_id,
)
}
use crate::paths::AgentsHome;
use crate::state::{load_registry, update_registry, RegistryEntry};
use crate::AgentStatus;
const NAME_MAX_LEN: usize = 128;
const FROM_NAME_MAX_LEN: usize = 128;
const DEFAULT_FOLLOWUP_TIMEOUT: Duration = Duration::from_secs(600);
const LOCK_ACQUIRE_TIMEOUT: Duration = Duration::from_secs(30);
const DEFAULT_SPAWN_TIMEOUT: Duration = Duration::from_secs(120);
const PROVABLY_LIVE_WINDOW_SECS: u64 = 3600;
fn is_provably_live_report(
inside_leg: Option<&crate::state::InsideLegReport>,
now_secs: u64,
) -> bool {
inside_leg.is_some_and(|r| r.received_within(now_secs, PROVABLY_LIVE_WINDOW_SECS))
}
fn now_epoch_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
fn spawn_create_timeout(explicit: Option<Duration>) -> Duration {
explicit.unwrap_or(DEFAULT_SPAWN_TIMEOUT)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AskOutcome {
pub stdout: String,
pub stderr: String,
pub exit_code: i32,
}
impl AskOutcome {
fn ok_stdout(s: String) -> Self {
Self {
stdout: s,
stderr: String::new(),
exit_code: 0,
}
}
fn err(msg: impl Into<String>, code: i32) -> Self {
Self {
stdout: String::new(),
stderr: format!("{}\n", msg.into()),
exit_code: code,
}
}
}
fn now_iso() -> String {
chrono::Utc::now().format("%Y-%m-%dT%H:%M:%SZ").to_string()
}
pub fn emit_event(events_path: &Path, kind: &str, fields: &[(&str, serde_json::Value)]) {
fn enc_value(v: &serde_json::Value) -> String {
match v {
serde_json::Value::String(s) => json_string_ascii(s),
other => serde_json::to_string(other).unwrap_or_default(),
}
}
let ts = now_iso();
let mut parts: Vec<String> = fields
.iter()
.map(|(k, v)| format!("{}:{}", json_string_ascii(k), enc_value(v)))
.collect();
parts.push(format!("\"ts\":{}", json_string_ascii(&ts)));
parts.push(format!("\"kind\":{}", json_string_ascii(kind)));
let line = format!("{{{}}}\n", parts.join(","));
let res = (|| -> std::io::Result<()> {
if let Some(parent) = events_path.parent() {
std::fs::create_dir_all(parent)?;
}
let mut fh = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(events_path)?;
fh.write_all(line.as_bytes())
})();
if let Err(e) = res {
if !EMIT_EVENT_WRITE_WARNED.swap(true, std::sync::atomic::Ordering::Relaxed) {
eprintln!(
"fno-agents: failed to append {} event to {}: {} (further event-write failures this run suppressed)",
kind,
events_path.display(),
e
);
}
}
}
static EMIT_EVENT_WRITE_WARNED: std::sync::atomic::AtomicBool =
std::sync::atomic::AtomicBool::new(false);
pub fn validate_inputs(name: &str, message: &str, from_name: &str) -> Result<(), String> {
validate_spawn_inputs(name, from_name)?;
if message.is_empty() || message.trim().is_empty() {
return Err("message must be non-empty".into());
}
Ok(())
}
pub fn validate_spawn_inputs(name: &str, from_name: &str) -> Result<(), String> {
if name.is_empty() {
return Err("agent name must not be empty".into());
}
if name.contains('/') || name.contains('\\') || name.contains("..") {
return Err(format!(
"agent name must not contain path separators or '..': {}",
py_repr(name)
));
}
if name.chars().count() > NAME_MAX_LEN {
return Err(format!(
"name must be <={} chars (got {})",
NAME_MAX_LEN,
name.chars().count()
));
}
if name.len() == 8
&& name
.bytes()
.all(|b| b.is_ascii_hexdigit() && !b.is_ascii_uppercase())
{
return Err(format!(
"agent name {} must not match short-id shape ^[0-9a-f]{{8}}$ (prevents name/id collision)",
py_repr(name)
));
}
if let Some(bad) = ['\u{0}', '\n', '\r', '=']
.into_iter()
.find(|c| name.contains(*c))
{
return Err(format!(
"agent name {} contains a forbidden character ({} would corrupt subprocess env injection)",
py_repr(name),
py_repr(&bad.to_string())
));
}
if from_name.is_empty() {
return Err("from-name must not be empty".into());
}
if from_name.chars().count() > FROM_NAME_MAX_LEN {
return Err(format!(
"from-name must be <={} chars (got {})",
FROM_NAME_MAX_LEN,
from_name.chars().count()
));
}
if from_name
.chars()
.any(|c| matches!(c, '"' | '<' | '>' | '&'))
{
return Err("from-name must not contain XML-unsafe characters (\", <, >, &)".into());
}
Ok(())
}
struct AgentLock {
_file: std::fs::File,
}
impl AgentLock {
fn acquire(home: &AgentsHome, name: &str, timeout: Duration) -> Result<Self, ()> {
let locks_dir = home.root().join("locks");
let _ = std::fs::create_dir_all(&locks_dir);
let path = locks_dir.join(format!("{}.lock", name));
let file = match std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.write(true)
.open(&path)
{
Ok(f) => f,
Err(_) => return Err(()),
};
let deadline = std::time::Instant::now() + timeout;
loop {
match file.try_lock() {
Ok(()) => return Ok(Self { _file: file }),
Err(_) => {
if std::time::Instant::now() >= deadline {
return Err(());
}
std::thread::sleep(Duration::from_millis(25));
}
}
}
}
}
impl Drop for AgentLock {
fn drop(&mut self) {
let _ = self._file.unlock();
}
}
fn derive_log_path(home: &AgentsHome, name: &str) -> PathBuf {
home.root()
.join("agents")
.join("logs")
.join(format!("{}.log", name))
}
#[allow(clippy::too_many_arguments)]
pub fn dispatch_claude_ask(
home: &AgentsHome,
claude_home: &ClaudeHome,
name: &str,
message: &str,
from_name: &str,
_cwd: &Path,
_yolo: bool,
timeout: Option<Duration>,
_extra_env: &[(&str, &str)],
) -> AskOutcome {
if let Err(msg) = validate_inputs(name, message, from_name) {
return AskOutcome::err(msg, 2);
}
let events = home.events_jsonl();
let registry_path = home.registry_json();
let _lock = match AgentLock::acquire(home, name, LOCK_ACQUIRE_TIMEOUT) {
Ok(l) => l,
Err(()) => {
emit_event(
&events,
"agent_ask_failed",
&[("stage", "lock-timeout".into()), ("name", name.into())],
);
return AskOutcome::err(
format!(
"lock timeout for agent {} after {:.1}s",
py_repr(name),
LOCK_ACQUIRE_TIMEOUT.as_secs_f64()
),
11,
);
}
};
let registry = match load_registry(®istry_path) {
Ok(r) => r,
Err(e) => {
emit_event(
&events,
"agent_ask_failed",
&[
("stage", "registry-read".into()),
("name", name.into()),
("error", e.to_string().into()),
],
);
return AskOutcome::err(format!("registry read failed: {}", e), 12);
}
};
let existing = registry.find(name).cloned();
match existing {
Some(entry) => followup(
home,
claude_home,
&events,
®istry_path,
name,
&entry,
message,
from_name,
timeout,
),
None => {
emit_event(
&events,
"agent_ask_failed",
&[
("stage", "unknown-name".into()),
("name", name.into()),
("provider", "claude".into()),
],
);
AskOutcome::err(
format!(
"unknown agent {}; spawn it first: fno agents spawn {} --harness <harness>",
py_repr(name),
name
),
16,
)
}
}
}
#[allow(clippy::too_many_arguments)]
pub fn dispatch_claude_spawn(
home: &AgentsHome,
claude_home: &ClaudeHome,
name: &str,
message: &str,
from_name: &str,
cwd: &Path,
yolo: bool,
timeout: Option<Duration>,
extra_env: &[(&str, &str)],
model: Option<&str>,
permission_mode: Option<&str>,
effort: Option<&str>,
flags: HarnessFlags,
surface_cwd: bool,
) -> AskOutcome {
if let Err(msg) = validate_spawn_inputs(name, from_name) {
return AskOutcome::err(msg, 2);
}
let events = home.events_jsonl();
let registry_path = home.registry_json();
let _lock = match AgentLock::acquire(home, name, LOCK_ACQUIRE_TIMEOUT) {
Ok(l) => l,
Err(()) => {
emit_event(
&events,
"agent_ask_failed",
&[
("stage", "lock-timeout".into()),
("name", name.into()),
("provider", "claude".into()),
],
);
return AskOutcome::err(
format!(
"lock timeout for agent {} after {:.1}s",
py_repr(name),
LOCK_ACQUIRE_TIMEOUT.as_secs_f64()
),
11,
);
}
};
let registry = match load_registry(®istry_path) {
Ok(r) => r,
Err(e) => {
return AskOutcome::err(format!("registry read failed: {}", e), 12);
}
};
if registry.find(name).is_some() {
return AskOutcome::err(
format!(
"agent {} already exists; use 'fno agents rm {}' first or pick another name",
py_repr(name),
name
),
2,
);
}
let effective_mode: Option<&str> = match permission_mode {
Some(m) => Some(m),
None if yolo => Some("bypassPermissions"),
None => None,
};
let inner = create(
home,
claude_home,
&events,
®istry_path,
name,
message,
from_name,
cwd,
yolo,
Some(spawn_create_timeout(timeout)),
extra_env,
model,
effective_mode,
effort,
flags,
);
if inner.exit_code != 0 {
return inner;
}
let short_id = inner.stdout.trim_end_matches('\n').to_string();
assert!(
short_id.len() == 8 && short_id.bytes().all(|b| b.is_ascii_hexdigit()),
"create path produced non-8hex short_id: {short_id:?}"
);
let safe_name = name.replace('"', "\\\"");
let perm_field = match effective_mode.filter(|m| !m.is_empty()) {
Some(m) => format!(r#", "permission_mode": "{}""#, m.replace('"', "\\\"")),
None => String::new(),
};
let cwd_field = if surface_cwd {
format!(
", \"cwd\": {}",
json_string_ascii(&cwd.display().to_string())
)
} else {
String::new()
};
AskOutcome {
stdout: format!(
r#"{{"name": "{safe_name}", "short_id": "{short_id}", "provider": "claude", "status": "live"{perm_field}{cwd_field}}}"#
) + "\n",
stderr: inner.stderr,
exit_code: 0,
}
}
#[allow(clippy::too_many_arguments)]
pub fn dispatch_claude_headless(
_claude_home: &ClaudeHome,
name: &str,
message: &str,
from_name: &str,
cwd: &Path,
_yolo: bool,
timeout: Option<Duration>,
model: Option<&str>,
permission_mode: Option<&str>,
effort: Option<&str>,
flags: HarnessFlags,
) -> AskOutcome {
use std::io::Write;
use std::process::{Command, Stdio};
use std::time::Instant;
if let Err(msg) = validate_spawn_inputs(name, from_name) {
return AskOutcome::err(msg, 2);
}
let effective = if message.is_empty() { "hello" } else { message };
let use_stdin = use_stdin_for(effective);
let mut argv: Vec<String> = vec!["claude".into(), "-p".into()];
match permission_mode.filter(|m| !m.is_empty()) {
Some(m) => {
argv.push("--permission-mode".into());
argv.push(m.into());
}
None => argv.push("--dangerously-skip-permissions".into()),
}
if let Some(m) = model.filter(|m| !m.is_empty()) {
argv.push("--model".to_string());
argv.push(m.to_string());
}
if let Some(value) = effort.filter(|v| !v.is_empty()) {
argv.push("--effort".to_string());
argv.push(value.to_string());
}
flags.push_onto(&mut argv);
if !use_stdin {
argv.push(effective.to_string());
}
let argv = crate::spawn_gate::qos_wrap(cwd, argv);
let mut cmd = Command::new(&argv[0]);
cmd.args(&argv[1..]);
cmd.current_dir(cwd);
cmd.env("FNO_AGENT_SELF", name);
cmd.env("FNO_AGENT_PROVIDER", "claude");
cmd.env("FNO_AGENT_FROM", from_name);
cmd.stdout(Stdio::piped());
cmd.stderr(Stdio::piped());
cmd.stdin(if use_stdin {
Stdio::piped()
} else {
Stdio::null()
});
let mut child = match cmd.spawn() {
Ok(c) => c,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
return AskOutcome::err(format!("claude CLI not found: {}", e), 127);
}
Err(e) => return AskOutcome::err(e.to_string(), 127),
};
let stdout_join = child.stdout.take().map(|mut p| {
std::thread::spawn(move || {
let mut s = Vec::new();
let _ = p.read_to_end(&mut s);
s
})
});
let stderr_join = child.stderr.take().map(|mut p| {
std::thread::spawn(move || {
let mut s = Vec::new();
let _ = p.read_to_end(&mut s);
s
})
});
let stdin_writer = if use_stdin {
child.stdin.take().map(|mut s| {
let data = effective.to_string();
std::thread::spawn(move || {
let _ = s.write_all(data.as_bytes());
})
})
} else {
None
};
let wait_result: Result<Option<i32>, String> = if let Some(limit) = timeout {
let deadline = Instant::now() + limit;
loop {
match child.try_wait() {
Ok(Some(st)) => break Ok(Some(st.code().unwrap_or(1))),
Ok(None) => {
if Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
break Ok(None);
}
std::thread::sleep(Duration::from_millis(50));
}
Err(e) => break Err(format!("claude -p wait failed: {}", e)),
}
}
} else {
child
.wait()
.map(|st| Some(st.code().unwrap_or(1)))
.map_err(|e| format!("claude -p wait failed: {}", e))
};
match wait_result {
Err(msg) => AskOutcome::err(msg, 1),
Ok(None) => AskOutcome::err(
format!(
"claude -p timed out after {:.1}s",
timeout
.expect("timeout Some on the bounded path")
.as_secs_f64()
),
124,
),
Ok(Some(exit_code)) => {
let stdout = stdout_join
.map(|h| h.join().unwrap_or_default())
.unwrap_or_default();
let stderr = stderr_join
.map(|h| h.join().unwrap_or_default())
.unwrap_or_default();
if let Some(h) = stdin_writer {
let _ = h.join();
}
AskOutcome {
stdout: String::from_utf8_lossy(&stdout).into_owned(),
stderr: String::from_utf8_lossy(&stderr).into_owned(),
exit_code,
}
}
}
}
#[allow(clippy::too_many_arguments)]
fn followup(
_home: &AgentsHome,
claude_home: &ClaudeHome,
events: &Path,
registry_path: &Path,
name: &str,
entry: &RegistryEntry,
message: &str,
from_name: &str,
timeout: Option<Duration>,
) -> AskOutcome {
let job_id = if entry.host_mode_or_default() == crate::state::HOST_MODE_INTERACTIVE {
None
} else {
entry.transport_short()
};
let short_id = match job_id {
Some(s) => s.to_string(),
_ => {
return AskOutcome::err(
format!(
"registry entry {} has no short id on file; cannot follow up. Remove with 'fno agents rm {}' and recreate.",
py_repr(name), name
),
12,
);
}
};
emit_event(
events,
"agent_followup_started",
&[
("name", name.into()),
("provider", entry.harness_name().to_string().into()),
("short_id", short_id.clone().into()),
],
);
let wait = timeout.unwrap_or(DEFAULT_FOLLOWUP_TIMEOUT);
match ask_followup(
claude_home,
&short_id,
message,
from_name,
wait,
DEFAULT_POLL_INTERVAL,
None,
) {
Ok(reply) => {
if let Err(e) = update_registry(registry_path, |reg| {
if let Some(en) = reg.find_mut(name) {
en.status = AgentStatus::Live;
en.last_message_at = Some(now_iso());
}
}) {
emit_event(
events,
"agent_followup_failed",
&[
("stage", "registry-write".into()),
("name", name.into()),
("short_id", short_id.clone().into()),
("error", e.to_string().into()),
],
);
return AskOutcome::err(
format!(
"registry write failed: {}. NOTE: message was already delivered; do not retry.",
e
),
12,
);
}
emit_event(
events,
"agent_followup_done",
&[
("stage", "followup".into()),
("name", name.into()),
("provider", entry.harness_name().to_string().into()),
("short_id", short_id.clone().into()),
("reply_chars", (reply.chars().count() as u64).into()),
("backend", "socket".into()),
],
);
AskOutcome::ok_stdout(reply)
}
Err(AskError::Orphan { reason, .. }) => {
let now = now_epoch_secs();
let routing_gap = matches!(
reason,
OrphanReason::RosterLiveInjectFailed | OrphanReason::TruthLiveInjectFailed
);
let mut provably_live = false;
let mut stamp_warning = String::new();
if let Err(e) = update_registry(registry_path, |reg| {
if let Some(en) = reg.find_mut(name) {
if routing_gap || is_provably_live_report(en.inside_leg.as_ref(), now) {
provably_live = true;
} else {
en.status = AgentStatus::Orphaned;
}
}
}) {
stamp_warning = format!(
"fno agents: warning: failed to mark {} as orphaned: {}\n",
py_repr(name),
e
);
emit_event(
events,
"agent_status_stamp_failed",
&[
("name", name.into()),
("short_id", short_id.clone().into()),
("target_status", "orphaned".into()),
("error", e.to_string().into()),
],
);
}
if provably_live {
emit_event(
events,
"agent_followup_failed",
&[
("stage", "routing-gap".into()),
("name", name.into()),
("short_id", short_id.clone().into()),
("reason", reason.as_str().into()),
],
);
return AskOutcome {
stdout: String::new(),
stderr: format!(
"agent {} is live but not currently routable (reason: {}); message not delivered. Try 'claude attach {}'\n",
py_repr(name),
reason.as_str(),
short_id
),
exit_code: 13,
};
}
emit_event(
events,
"agent_followup_failed",
&[
("stage", "orphan".into()),
("name", name.into()),
("short_id", short_id.clone().into()),
("reason", reason.as_str().into()),
],
);
let hint = match reason {
OrphanReason::SocketNull => format!(
". Run 'claude attach {}' to wake the session, or 'fno agents rm {}' to remove",
short_id, name
),
OrphanReason::NotFound => format!(". Run 'fno agents rm {}' to clear the stale entry", name),
OrphanReason::LivenessFailed => format!(
". Socket exists but is unresponsive; try 'claude attach {}' or 'fno agents rm {}'",
short_id, name
),
OrphanReason::RosterLiveInjectFailed => format!(
". Inspect with 'fno agents logs {}' or remove via 'fno agents rm {}'",
name, name
),
OrphanReason::TruthLiveInjectFailed => format!(
". Inspect with 'fno agents logs {}' before removing",
name
),
};
let suspended = if reason == OrphanReason::SocketNull {
"; session is suspended"
} else {
""
};
AskOutcome {
stdout: String::new(),
stderr: format!(
"{}agent {} is not running (reason: {}{}){}\n",
stamp_warning,
py_repr(name),
reason.as_str(),
suspended,
hint
),
exit_code: 13,
}
}
Err(AskError::Socket { message }) => {
emit_event(
events,
"agent_followup_failed",
&[
("stage", "send".into()),
("name", name.into()),
("short_id", short_id.clone().into()),
("reason", "socket-error".into()),
],
);
AskOutcome::err(message, 1)
}
Err(AskError::Timeout { elapsed_sec, .. }) => {
emit_event(
events,
"agent_followup_failed",
&[
("stage", "poll-timeout".into()),
("name", name.into()),
("short_id", short_id.clone().into()),
("elapsed_sec", (elapsed_sec as u64).into()),
],
);
AskOutcome::err(
format!(
"message sent but no reply within {}s. Try 'fno agents logs {}' to read the transcript.",
elapsed_sec as u64, name
),
15,
)
}
Err(other) => AskOutcome::err(other.to_string(), 1),
}
}
#[allow(clippy::too_many_arguments)]
fn create(
home: &AgentsHome,
claude_home: &ClaudeHome,
events: &Path,
registry_path: &Path,
name: &str,
message: &str,
_from_name: &str,
cwd: &Path,
yolo: bool,
timeout: Option<Duration>,
extra_env: &[(&str, &str)],
model: Option<&str>,
permission_mode: Option<&str>,
effort: Option<&str>,
flags: HarnessFlags,
) -> AskOutcome {
let pre_stderr = String::new();
let result = match bg_create(
name,
message,
cwd,
timeout,
extra_env,
model,
permission_mode,
effort,
flags,
) {
Ok(r) => r,
Err(AskError::Subprocess { exit_code, stderr }) => {
emit_event(
events,
"agent_ask_failed",
&[
("stage", "subprocess".into()),
("name", name.into()),
("provider", "claude".into()),
("returncode", exit_code.into()),
],
);
let code = if exit_code == 127 { 14 } else { 1 };
return AskOutcome {
stdout: String::new(),
stderr: format!("{}{}\n", pre_stderr, stderr),
exit_code: code,
};
}
Err(AskError::Parse { stdout_head }) => {
emit_event(
events,
"agent_ask_failed",
&[
("stage", "parse".into()),
("name", name.into()),
("provider", "claude".into()),
("short_id_raw", stdout_head.clone().into()),
],
);
return AskOutcome {
stdout: String::new(),
stderr: format!(
"{}unable to parse short-id from claude --bg output: {}\n",
pre_stderr, stdout_head
),
exit_code: 1,
};
}
Err(other) => {
return AskOutcome {
stdout: String::new(),
stderr: format!("{}{}\n", pre_stderr, other),
exit_code: 1,
};
}
};
let short_id = result.short_id.clone();
let session_uuid = resolve_session_uuid_at_spawn(claude_home, &short_id);
let log_path = derive_log_path(home, name);
let new_entry = RegistryEntry {
name: name.to_string(),
short_id: short_id.clone(),
legacy_provider: String::new(),
harness: Some("claude".to_string()),
harness_session_id: session_uuid.clone(),
cwd: cwd.to_string_lossy().to_string(),
project_root: String::new(),
session_id: None,
claude_session_uuid: session_uuid,
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
host_mode: None, cc_session_id: None,
status: AgentStatus::Live,
last_message_at: None,
created_at: now_iso(),
pid: None,
pid_start_time: None,
log_path: Some(log_path.to_string_lossy().to_string()),
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
legacy_claude_short_id: None,
};
match update_registry(registry_path, |reg| {
if reg.find(name).is_some() {
false
} else {
reg.entries.push(new_entry.clone());
true
}
}) {
Ok(true) => {}
Ok(false) => {
emit_event(
events,
"agent_ask_failed",
&[
("stage", "name-collision".into()),
("name", name.into()),
("provider", "claude".into()),
("short_id", short_id.clone().into()),
],
);
return AskOutcome {
stdout: String::new(),
stderr: format!(
"{}agent {} already exists (registered concurrently); orphaned supervisor session: claude rm {} (registry not updated)\n",
pre_stderr,
py_repr(name),
short_id
),
exit_code: 12,
};
}
Err(_) => {
emit_event(
events,
"agent_ask_failed",
&[
("stage", "registry-write".into()),
("name", name.into()),
("provider", "claude".into()),
("short_id", short_id.clone().into()),
],
);
return AskOutcome {
stdout: String::new(),
stderr: format!(
"{}registry write failed. orphaned supervisor session: claude rm {} (registry not updated)\n",
pre_stderr, short_id
),
exit_code: 12,
};
}
}
emit_event(
events,
"agent_ask_done",
&[
("stage", "dispatch".into()),
("name", name.into()),
("provider", "claude".into()),
("short_id", short_id.clone().into()),
("duration_ms", (result.duration_ms as u64).into()),
("yolo", yolo.into()),
],
);
AskOutcome {
stdout: format!("{}\n", short_id),
stderr: pre_stderr,
exit_code: 0,
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
#[test]
fn spawn_create_timeout_defaults_when_unset() {
assert_eq!(spawn_create_timeout(None), DEFAULT_SPAWN_TIMEOUT);
}
#[test]
fn spawn_create_timeout_honors_explicit() {
let explicit = Duration::from_secs(5);
assert_eq!(spawn_create_timeout(Some(explicit)), explicit);
}
fn report_at(stamp: &str) -> crate::state::InsideLegReport {
crate::state::InsideLegReport {
state: crate::state::InsideLegState::Working,
seq: 1,
reason: None,
received_at: stamp.into(),
ttl_ms: None,
}
}
#[test]
fn provably_live_true_for_recent_inside_leg_report() {
let stamp = "2026-07-06T20:00:00Z";
let now = crate::state::rfc3339_like_to_secs(stamp).unwrap() + 30;
assert!(is_provably_live_report(Some(&report_at(stamp)), now));
}
#[test]
fn not_provably_live_without_inside_leg_report() {
assert!(!is_provably_live_report(None, 9_999_999_999));
}
#[test]
fn not_provably_live_when_inside_leg_is_stale() {
let stamp = "2026-07-06T20:00:00Z";
let now =
crate::state::rfc3339_like_to_secs(stamp).unwrap() + PROVABLY_LIVE_WINDOW_SECS + 60;
assert!(!is_provably_live_report(Some(&report_at(stamp)), now));
}
#[test]
fn not_provably_live_for_future_stamp() {
let stamp = "2026-07-06T20:00:00Z";
let now = crate::state::rfc3339_like_to_secs(stamp).unwrap() - 60;
assert!(!is_provably_live_report(Some(&report_at(stamp)), now));
}
#[test]
fn scan_returns_short_id_on_confirmation_line() {
let input = "backgrounded \u{b7} abcd1234 \u{b7} my-worker\n";
match scan_stdout_for_short_id(std::io::Cursor::new(input)) {
ShortIdScan::Found { short_id, .. } => assert_eq!(short_id, "abcd1234"),
ShortIdScan::NoId { .. } => panic!("expected Found"),
}
}
#[test]
fn scan_stops_at_confirmation_does_not_drain_to_eof() {
let input = "warming up\nbackgrounded \u{b7} 0011aabb \u{b7} w\nSENTINEL_AFTER_ID\n";
match scan_stdout_for_short_id(std::io::Cursor::new(input)) {
ShortIdScan::Found { short_id, consumed } => {
assert_eq!(short_id, "0011aabb");
assert!(consumed.contains("warming up"));
assert!(
!consumed.contains("SENTINEL_AFTER_ID"),
"scan read past the confirmation line: {consumed:?}"
);
}
ShortIdScan::NoId { .. } => panic!("expected Found"),
}
}
#[test]
fn scan_no_confirmation_is_noid_at_eof() {
let input = "error: could not start\ngiving up\n";
match scan_stdout_for_short_id(std::io::Cursor::new(input)) {
ShortIdScan::NoId { consumed } => assert!(consumed.contains("giving up")),
ShortIdScan::Found { .. } => panic!("expected NoId"),
}
}
#[test]
fn scan_matches_colorized_confirmation_line() {
let input = "backgrounded \u{b7} \u{1b}[36mdeadbeef\u{1b}[39m \u{b7} w\n";
match scan_stdout_for_short_id(std::io::Cursor::new(input)) {
ShortIdScan::Found { short_id, .. } => assert_eq!(short_id, "deadbeef"),
ShortIdScan::NoId { .. } => panic!("expected Found on colorized line"),
}
}
fn tmpdir() -> PathBuf {
use std::sync::atomic::{AtomicU64, Ordering};
static SEQ: AtomicU64 = AtomicU64::new(0);
let seq = SEQ.fetch_add(1, Ordering::Relaxed);
let p = std::env::temp_dir().join(format!(
"fno-claude-ask-{}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos(),
seq
));
fs::create_dir_all(&p).unwrap();
p
}
#[test]
fn parse_short_id_happy() {
let out = "backgrounded \u{b7} 7c5dcf5d \u{b7} alice\n";
assert_eq!(parse_short_id(out).unwrap(), "7c5dcf5d");
}
#[test]
fn parse_short_id_only_first_line() {
let out = "backgrounded \u{b7} 7c5dcf5d \u{b7} alice\nextra garbage\n";
assert_eq!(parse_short_id(out).unwrap(), "7c5dcf5d");
}
#[test]
fn parse_short_id_empty_is_error() {
assert!(matches!(parse_short_id(""), Err(AskError::Parse { .. })));
}
#[test]
fn parse_short_id_strips_ansi_color() {
let out = "backgrounded \u{b7} \u{1b}[36m441064a2\u{1b}[39m \u{b7} fnogates\n";
assert_eq!(parse_short_id(out).unwrap(), "441064a2");
}
#[test]
fn parse_short_id_strips_truecolor_sgr() {
let out = "backgrounded \u{b7} \u{1b}[38;2;215;119;87m7c5dcf5d\u{1b}[0m \u{b7} alice";
assert_eq!(parse_short_id(out).unwrap(), "7c5dcf5d");
}
#[test]
fn strip_ansi_csi_borrows_when_no_escape() {
assert!(matches!(
strip_ansi_csi("backgrounded \u{b7} 7c5dcf5d \u{b7} alice"),
std::borrow::Cow::Borrowed(_)
));
let owned = strip_ansi_csi("a\u{1b}[31mb\u{1b}[0mc");
assert!(matches!(owned, std::borrow::Cow::Owned(_)));
assert_eq!(owned.as_ref(), "abc");
}
#[test]
fn parse_short_id_no_panic_on_non_ascii_at_byte_8() {
let out = "backgrounded \u{b7} 1234567\u{e9} \u{b7} x";
assert!(parse_short_id(out).is_err());
}
#[test]
fn parse_short_id_rejects_non_hex_and_uppercase() {
assert!(parse_short_id("backgrounded \u{b7} 7C5DCF5D \u{b7} a").is_err());
assert!(parse_short_id("backgrounded \u{b7} zzzzzzzz \u{b7} a").is_err());
assert!(parse_short_id("nope \u{b7} 7c5dcf5d \u{b7} a").is_err());
assert!(parse_short_id("backgrounded \u{b7} 7c5dcf5d done").is_err());
}
#[test]
fn build_argv_inline_vs_stdin() {
assert_eq!(
build_argv("a", "hi", false, None, None, None, HarnessFlags::default()),
vec!["claude", "--bg", "--name", "a", "hi"]
);
assert_eq!(
build_argv("a", "hi", true, None, None, None, HarnessFlags::default()),
vec!["claude", "--bg", "--name", "a"]
);
}
#[test]
fn build_argv_appends_permission_mode() {
assert_eq!(
build_argv(
"a",
"hi",
false,
None,
Some("acceptEdits"),
None,
HarnessFlags::default()
),
vec![
"claude",
"--bg",
"--name",
"a",
"--permission-mode",
"acceptEdits",
"hi"
]
);
assert_eq!(
build_argv(
"a",
"hi",
false,
None,
Some(""),
None,
HarnessFlags::default()
),
build_argv("a", "hi", false, None, None, None, HarnessFlags::default())
);
assert!(!build_argv(
"a",
"hi",
false,
None,
Some("acceptEdits"),
None,
HarnessFlags::default()
)
.iter()
.any(|t| t == "--dangerously-skip-permissions"));
}
#[test]
fn build_argv_appends_model_pin() {
assert_eq!(
build_argv(
"a",
"hi",
false,
Some("fable"),
None,
None,
HarnessFlags::default()
),
vec!["claude", "--bg", "--name", "a", "--model", "fable", "hi"]
);
assert_eq!(
build_argv(
"a",
"hi",
true,
Some("fable"),
None,
None,
HarnessFlags::default()
),
vec!["claude", "--bg", "--name", "a", "--model", "fable"]
);
assert_eq!(
build_argv(
"a",
"hi",
false,
Some(""),
None,
None,
HarnessFlags::default()
),
build_argv("a", "hi", false, None, None, None, HarnessFlags::default())
);
}
#[test]
fn build_argv_appends_effort() {
assert_eq!(
build_argv(
"a",
"hi",
false,
None,
None,
Some("high"),
HarnessFlags::default()
),
vec!["claude", "--bg", "--name", "a", "--effort", "high", "hi"]
);
}
#[test]
fn build_argv_appends_harness_flags() {
let flags = HarnessFlags {
add_dir: Some("/work"),
agent: Some("reviewer"),
allowed_tools: Some("Read,Edit"),
disallowed_tools: Some("Bash"),
};
assert_eq!(
build_argv("a", "hi", false, None, None, None, flags),
vec![
"claude",
"--bg",
"--name",
"a",
"--add-dir",
"/work",
"--agent",
"reviewer",
"--allowedTools",
"Read,Edit",
"--disallowedTools",
"Bash",
"hi"
]
);
let only_dir = HarnessFlags {
add_dir: Some("/work"),
allowed_tools: Some(""),
..Default::default()
};
assert_eq!(
build_argv("a", "hi", false, None, None, None, only_dir),
vec!["claude", "--bg", "--name", "a", "--add-dir", "/work", "hi"]
);
}
#[test]
fn use_stdin_threshold() {
assert!(!use_stdin_for(&"x".repeat(ARGV_OVERFLOW_THRESHOLD)));
assert!(use_stdin_for(&"x".repeat(ARGV_OVERFLOW_THRESHOLD + 1)));
}
#[test]
fn envelope_exact_bytes_ascii() {
let env = build_envelope("hello", "bob");
let expected = "{\"type\":\"user\",\"message\":{\"role\":\"user\",\"content\":\"<cross-session-message from-name=\\\"bob\\\">\\nhello\\n</cross-session-message>\"},\"priority\":\"next\"}\n";
assert_eq!(String::from_utf8(env).unwrap(), expected);
}
#[test]
fn envelope_escapes_from_name_html() {
let env = build_envelope("hi", "a&b<c>\"d'e");
let s = String::from_utf8(env).unwrap();
assert!(
s.contains("from-name=\\\"a&b<c>"d'e\\\""),
"{}",
s
);
}
#[test]
fn envelope_ensure_ascii_non_ascii() {
let env = build_envelope("caf\u{e9}", "x");
let s = String::from_utf8(env).unwrap();
assert!(s.contains("caf\\u00e9"), "{}", s);
assert!(!s.contains('\u{e9}'), "raw non-ascii leaked: {}", s);
}
#[test]
fn envelope_astral_surrogate_pair() {
let env = build_envelope("\u{1F600}", "x");
let s = String::from_utf8(env).unwrap();
assert!(s.contains("\\ud83d\\ude00"), "{}", s);
}
#[test]
fn json_string_escapes_control_chars() {
assert_eq!(json_string_ascii("a\nb\tc"), "\"a\\nb\\tc\"");
assert_eq!(json_string_ascii("\u{01}"), "\"\\u0001\"");
assert_eq!(json_string_ascii("a\"b\\c"), "\"a\\\"b\\\\c\"");
}
fn write_session(dir: &Path, pid: &str, job: &str, kind: &str, sock: Option<&str>) {
let sock_field = match sock {
Some(s) => format!("\"{}\"", s),
None => "null".to_string(),
};
let body = format!(
"{{\"jobId\":\"{}\",\"kind\":\"{}\",\"messagingSocketPath\":{},\"sessionId\":\"sess-{}\",\"cwd\":\"/tmp\"}}",
job, kind, sock_field, job
);
fs::write(dir.join(format!("{}.json", pid)), body).unwrap();
}
#[test]
fn locate_session_happy() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(&sessions, "111", "7c5dcf5d", "bg", Some("/tmp/sock1"));
let ch = ClaudeHome::at(&home);
let loc = locate_session(&ch, "7c5dcf5d").unwrap();
assert_eq!(loc.pid, 111);
assert_eq!(loc.messaging_socket_path, "/tmp/sock1");
assert_eq!(loc.session_id.as_deref(), Some("sess-7c5dcf5d"));
assert_eq!(
loc.jobs_dir,
home.join(".claude").join("jobs").join("7c5dcf5d")
);
}
#[test]
fn locate_session_skips_null_socket_prefers_live() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(&sessions, "100", "abcd1234", "bg", None);
write_session(&sessions, "200", "abcd1234", "bg", Some("/tmp/live"));
let ch = ClaudeHome::at(&home);
let loc = locate_session(&ch, "abcd1234").unwrap();
assert_eq!(loc.messaging_socket_path, "/tmp/live");
assert_eq!(loc.pid, 200);
}
#[test]
fn locate_session_not_found_and_classify() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(&sessions, "100", "abcd1234", "bg", None);
let ch = ClaudeHome::at(&home);
assert!(locate_session(&ch, "abcd1234").is_none());
assert_eq!(
classify_orphan_reason(&ch, "abcd1234"),
OrphanReason::SocketNull
);
assert_eq!(
classify_orphan_reason(&ch, "ffffffff"),
OrphanReason::NotFound
);
}
#[test]
fn locate_session_skips_corrupt_and_non_bg() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
fs::write(sessions.join("1.json"), "{not json").unwrap();
write_session(&sessions, "2", "abcd1234", "interactive", Some("/tmp/x"));
write_session(&sessions, "3", "abcd1234", "bg", Some("/tmp/good"));
let ch = ClaudeHome::at(&home);
let loc = locate_session(&ch, "abcd1234").unwrap();
assert_eq!(loc.messaging_socket_path, "/tmp/good");
}
#[test]
fn locate_session_missing_dir_is_none() {
let home = tmpdir();
let ch = ClaudeHome::at(&home);
assert!(locate_session(&ch, "abcd1234").is_none());
}
#[test]
fn resolve_session_uuid_resolves_idle_bg() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(&sessions, "111", "7c5dcf5d", "bg", None); let ch = ClaudeHome::at(&home);
assert!(locate_session(&ch, "7c5dcf5d").is_none()); assert_eq!(
resolve_session_uuid(&ch, "7c5dcf5d").as_deref(),
Some("sess-7c5dcf5d") );
}
#[test]
fn resolve_session_uuid_skips_non_bg_and_unmatched() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(&sessions, "1", "7c5dcf5d", "interactive", Some("/tmp/x")); write_session(&sessions, "2", "deadbeef", "bg", Some("/tmp/y")); let ch = ClaudeHome::at(&home);
assert!(resolve_session_uuid(&ch, "7c5dcf5d").is_none());
assert!(resolve_session_uuid(&ch, "ffffffff").is_none());
}
#[test]
fn resolve_session_uuid_missing_dir_is_none() {
let home = tmpdir();
let ch = ClaudeHome::at(&home);
assert!(resolve_session_uuid(&ch, "7c5dcf5d").is_none());
}
#[test]
fn resolve_at_spawn_happy_empty_and_missing() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(&sessions, "111", "7c5dcf5d", "bg", Some("/tmp/sock"));
let ch = ClaudeHome::at(&home);
assert_eq!(
resolve_session_uuid_at_spawn(&ch, "7c5dcf5d").as_deref(),
Some("sess-7c5dcf5d")
);
assert!(resolve_session_uuid_at_spawn(&ch, "").is_none());
let empty = tmpdir();
let empty_ch = ClaudeHome::at(&empty);
assert!(resolve_session_uuid_at_spawn(&empty_ch, "7c5dcf5d").is_none());
}
#[test]
fn read_state_json_parses_fields() {
let jobs = tmpdir();
fs::write(
jobs.join("state.json"),
r#"{"state":"completed","updatedAt":"2026-05-27T10:00:00Z","output":{"result":"PONG"},"intent":"reply"}"#,
)
.unwrap();
let snap = read_state_json(&jobs).unwrap();
assert_eq!(snap.state, "completed");
assert_eq!(snap.updated_at.as_deref(), Some("2026-05-27T10:00:00Z"));
assert_eq!(snap.output_result.as_deref(), Some("PONG"));
}
#[test]
fn read_state_json_missing_is_notfound() {
let jobs = tmpdir();
assert!(matches!(
read_state_json(&jobs),
Err(StateReadError::NotFound)
));
}
#[test]
fn read_state_json_empty_is_parse_err() {
let jobs = tmpdir();
fs::write(jobs.join("state.json"), " ").unwrap();
assert!(matches!(read_state_json(&jobs), Err(StateReadError::Parse)));
}
#[test]
fn read_state_json_output_not_dict() {
let jobs = tmpdir();
fs::write(jobs.join("state.json"), r#"{"state":"done","output":null}"#).unwrap();
let snap = read_state_json(&jobs).unwrap();
assert_eq!(snap.output_result, None);
}
#[test]
fn timeline_tail_concats_terminal_text_from_offset() {
let jobs = tmpdir();
let tl = jobs.join("timeline.jsonl");
fs::write(&tl, "{\"state\":\"completed\",\"text\":\"OLD\"}\n").unwrap();
let offset = timeline_offset(&jobs);
let mut f = fs::OpenOptions::new().append(true).open(&tl).unwrap();
writeln!(f, "{{\"state\":\"running\",\"text\":\"tool call\"}}").unwrap();
writeln!(f, "{{\"state\":\"completed\",\"text\":\"AB\"}}").unwrap();
writeln!(f, "{{\"state\":\"done\",\"text\":\"CD\"}}").unwrap();
writeln!(f, "not json").unwrap();
assert_eq!(read_timeline_tail(&jobs, offset), "ABCD");
}
#[test]
fn timeline_tail_missing_is_empty() {
let jobs = tmpdir();
assert_eq!(read_timeline_tail(&jobs, 0), "");
}
#[test]
fn send_to_session_delivers_envelope_bytes() {
use std::os::unix::net::UnixListener;
let dir = tmpdir();
let sock = dir.join("s.sock");
let listener = UnixListener::bind(&sock).unwrap();
let sock_str = sock.to_str().unwrap().to_string();
let handle = std::thread::spawn(move || {
let (mut conn, _) = listener.accept().unwrap();
let mut buf = Vec::new();
conn.read_to_end(&mut buf).unwrap();
buf
});
send_to_session(&sock_str, "ping", "tester").unwrap();
let got = handle.join().unwrap();
assert_eq!(got, build_envelope("ping", "tester"));
}
#[test]
fn liveness_probe_true_when_listening_false_when_absent() {
use std::os::unix::net::UnixListener;
let dir = tmpdir();
let sock = dir.join("live.sock");
let _listener = UnixListener::bind(&sock).unwrap();
assert!(liveness_probe(sock.to_str().unwrap()));
assert!(!liveness_probe(dir.join("absent.sock").to_str().unwrap()));
}
#[test]
fn send_to_session_errors_on_missing_socket() {
let dir = tmpdir();
let res = send_to_session(dir.join("nope.sock").to_str().unwrap(), "x", "y");
assert!(matches!(res, Err(AskError::Socket { .. })));
}
fn write_state(jobs: &Path, state: &str, updated: &str, result: Option<&str>) {
let res = match result {
Some(r) => format!(",\"output\":{{\"result\":{}}}", json_string_ascii(r)),
None => String::new(),
};
fs::write(
jobs.join("state.json"),
format!(
"{{\"state\":\"{}\",\"updatedAt\":\"{}\"{}}}",
state, updated, res
),
)
.unwrap();
}
#[test]
fn wait_for_reply_prefers_output_result() {
let jobs = tmpdir();
write_state(&jobs, "completed", "2026-05-27T10:00:01Z", Some("PONG"));
let r = wait_for_reply(
&jobs,
Some("2026-05-27T10:00:00Z"),
0,
Duration::from_secs(2),
Duration::from_millis(10),
"sid",
)
.unwrap();
assert_eq!(r, "PONG");
}
#[test]
fn wait_for_reply_baseline_invariant_then_advance() {
let jobs = tmpdir();
write_state(&jobs, "completed", "2026-05-27T10:00:00Z", Some("STALE"));
let r = wait_for_reply(
&jobs,
Some("2026-05-27T10:00:00Z"),
0,
Duration::from_millis(200),
Duration::from_millis(10),
"sid",
);
assert!(
matches!(r, Err(AskError::Timeout { .. })),
"stale state equal to baseline must not satisfy wait_for_reply; got {:?}",
r
);
write_state(&jobs, "completed", "2026-05-27T10:00:05Z", Some("FRESH"));
let r = wait_for_reply(
&jobs,
Some("2026-05-27T10:00:00Z"),
0,
Duration::from_secs(2),
Duration::from_millis(10),
"sid",
)
.unwrap();
assert_eq!(r, "FRESH");
}
#[test]
fn wait_for_reply_falls_back_to_timeline_when_result_empty() {
let jobs = tmpdir();
let offset = timeline_offset(&jobs); write_state(&jobs, "done", "2026-05-27T10:00:01Z", None);
fs::write(
jobs.join("timeline.jsonl"),
"{\"state\":\"done\",\"text\":\"TAIL\"}\n",
)
.unwrap();
let r = wait_for_reply(
&jobs,
None,
offset,
Duration::from_secs(2),
Duration::from_millis(10),
"sid",
)
.unwrap();
assert_eq!(r, "TAIL");
}
#[test]
fn read_state_json_eacces_is_fatal_io_not_transient() {
if unsafe { libc::geteuid() } == 0 {
eprintln!("SKIP: running as root; permission bits not enforced");
return;
}
use std::os::unix::fs::PermissionsExt;
let jobs = tmpdir();
let sp = jobs.join("state.json");
fs::write(&sp, r#"{"state":"done","updatedAt":"t"}"#).unwrap();
fs::set_permissions(&sp, fs::Permissions::from_mode(0o000)).unwrap();
let got = read_state_json(&jobs);
let _ = fs::set_permissions(&sp, fs::Permissions::from_mode(0o644));
assert!(
matches!(got, Err(StateReadError::Io(_))),
"expected Io, got {:?}",
got
);
fs::set_permissions(&sp, fs::Permissions::from_mode(0o000)).unwrap();
let r = wait_for_reply(
&jobs,
None,
0,
Duration::from_secs(30),
Duration::from_millis(10),
"sid",
);
let _ = fs::set_permissions(&sp, fs::Permissions::from_mode(0o644));
assert!(
matches!(r, Err(AskError::Io { .. })),
"expected fatal Io, got {:?}",
r
);
}
#[test]
fn wait_for_reply_times_out() {
let jobs = tmpdir();
write_state(&jobs, "running", "2026-05-27T10:00:00Z", None);
let r = wait_for_reply(
&jobs,
None,
0,
Duration::from_millis(40),
Duration::from_millis(10),
"sid",
);
assert!(matches!(r, Err(AskError::Timeout { .. })));
}
#[test]
fn ask_followup_socket_to_reply() {
use std::os::unix::net::UnixListener;
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
let jobs = home.join(".claude").join("jobs").join("abcd1234");
fs::create_dir_all(&sessions).unwrap();
fs::create_dir_all(&jobs).unwrap();
let sock = home.join("msg.sock");
let listener = UnixListener::bind(&sock).unwrap();
write_session(
&sessions,
"999",
"abcd1234",
"bg",
Some(sock.to_str().unwrap()),
);
let jobs_for_thread = jobs.clone();
let handle = std::thread::spawn(move || {
loop {
let (mut conn, _) = listener.accept().unwrap();
let mut buf = Vec::new();
let _ = conn.read_to_end(&mut buf);
if buf.is_empty() {
continue; }
write_state(
&jobs_for_thread,
"completed",
"2026-05-27T10:00:09Z",
Some("REPLY!"),
);
break buf;
}
});
let ch = ClaudeHome::at(&home);
let reply = ask_followup(
&ch,
"abcd1234",
"ping",
"tester",
Duration::from_secs(10),
Duration::from_millis(10),
None,
)
.unwrap();
let envelope = handle.join().unwrap();
assert_eq!(reply, "REPLY!");
assert_eq!(envelope, build_envelope("ping", "tester"));
}
#[test]
fn ask_followup_orphan_socket_null() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(&sessions, "1", "abcd1234", "bg", None);
let ch = ClaudeHome::at(&home);
let err = ask_followup(
&ch,
"abcd1234",
"x",
"y",
Duration::from_secs(1),
Duration::from_millis(10),
None,
)
.unwrap_err();
match err {
AskError::Orphan { reason, .. } => {
assert_eq!(reason, OrphanReason::TruthLiveInjectFailed)
}
other => panic!("expected orphan, got {:?}", other),
}
}
#[test]
fn ask_followup_orphan_not_found() {
let home = tmpdir();
fs::create_dir_all(home.join(".claude").join("sessions")).unwrap();
let ch = ClaudeHome::at(&home);
let err = ask_followup(
&ch,
"ffffffff",
"x",
"y",
Duration::from_secs(1),
Duration::from_millis(10),
None,
)
.unwrap_err();
assert!(matches!(
err,
AskError::Orphan {
reason: OrphanReason::TruthLiveInjectFailed,
..
}
));
}
#[test]
fn ask_followup_liveness_failed_when_socket_dead() {
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(
&sessions,
"1",
"abcd1234",
"bg",
Some(home.join("dead.sock").to_str().unwrap()),
);
let ch = ClaudeHome::at(&home);
let err = ask_followup(
&ch,
"abcd1234",
"x",
"y",
Duration::from_secs(1),
Duration::from_millis(10),
None,
)
.unwrap_err();
assert!(matches!(
err,
AskError::Orphan {
reason: OrphanReason::TruthLiveInjectFailed,
..
}
));
}
static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
fn write_roster(home: &Path, session_uuid: &str) {
let daemon = home.join(".claude").join("daemon");
fs::create_dir_all(&daemon).unwrap();
let body = format!(
"{{\"workers\":{{\"w\":{{\"sessionId\":\"{}\",\"pid\":5}}}}}}",
session_uuid
);
fs::write(daemon.join("roster.json"), body).unwrap();
}
#[test]
fn build_cross_session_container_wraps_peer_turn() {
assert_eq!(
build_cross_session_container("hello", "fno"),
"<cross-session-message from-name=\"fno\">\nhello\n</cross-session-message>"
);
}
#[test]
fn orphan_reason_roster_live_inject_failed_token() {
assert_eq!(
OrphanReason::RosterLiveInjectFailed.as_str(),
"roster-live-inject-failed"
);
}
#[test]
fn daemon_roster_path_honors_env_override_first() {
let _g = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner());
let home = tmpdir();
let ch = ClaudeHome::at(&home);
std::env::remove_var(crate::claude_roster::DAEMON_DIR_ENV);
assert_eq!(
ch.daemon_roster_path(),
home.join(".claude").join("daemon").join("roster.json")
);
let alt = tmpdir();
std::env::set_var(crate::claude_roster::DAEMON_DIR_ENV, &alt);
assert_eq!(ch.daemon_roster_path(), alt.join("roster.json"));
std::env::remove_var(crate::claude_roster::DAEMON_DIR_ENV);
}
#[test]
fn ask_followup_socket_null_roster_live_falls_back_to_control_sock() {
let _g = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner());
std::env::remove_var(crate::claude_roster::DAEMON_DIR_ENV);
let home = tmpdir();
let sessions = home.join(".claude").join("sessions");
fs::create_dir_all(&sessions).unwrap();
write_session(&sessions, "1", "abcd1234", "bg", None);
write_roster(&home, "abcd1234-1111-2222-3333-444455556666");
let ch = ClaudeHome::at(&home);
let err = ask_followup(
&ch,
"abcd1234",
"x",
"y",
Duration::from_millis(200),
Duration::from_millis(10),
None,
)
.unwrap_err();
match err {
AskError::Orphan { reason, .. } => {
assert_eq!(reason, OrphanReason::RosterLiveInjectFailed)
}
other => panic!("expected roster-live-inject-failed orphan, got {:?}", other),
}
}
#[test]
fn ask_followup_not_found_never_falls_back_even_if_rostered() {
let home = tmpdir();
fs::create_dir_all(home.join(".claude").join("sessions")).unwrap();
write_roster(&home, "abcd1234-1111-2222-3333-444455556666");
let ch = ClaudeHome::at(&home);
let err = ask_followup(
&ch,
"abcd1234",
"x",
"y",
Duration::from_secs(1),
Duration::from_millis(10),
None,
)
.unwrap_err();
match err {
AskError::Orphan { reason, .. } => {
assert_eq!(reason, OrphanReason::TruthLiveInjectFailed)
}
other => panic!("expected routing gap, got {:?}", other),
}
}
#[test]
fn family1_truth_subprocess_is_bounded_and_validated() {
let mut valid = std::process::Command::new("sh");
valid.args(["-c", "printf '{\"state\":\"watching\"}'"]);
assert_eq!(
family1_truth_state_with_command(valid, Duration::from_secs(1), "h1").as_deref(),
Some("watching")
);
let mut invalid = std::process::Command::new("sh");
invalid.args(["-c", "printf '{\"state\":\"invented\"}'"]);
assert_eq!(
family1_truth_state_with_command(invalid, Duration::from_secs(1), "h1"),
None
);
let mut hung = std::process::Command::new("sh");
hung.args(["-c", "sleep 5"]);
let started = Instant::now();
assert_eq!(
family1_truth_state_with_command(hung, Duration::from_millis(50), "h1"),
None
);
assert!(started.elapsed() < Duration::from_secs(1));
}
#[test]
fn only_a_routine_not_found_is_silenced() {
assert!(truth_failure_is_routine("not-found"));
assert!(truth_failure_is_routine(" not-found\n"));
assert!(!truth_failure_is_routine("transcript-unreadable"));
assert!(!truth_failure_is_routine("resolver-error"));
assert!(!truth_failure_is_routine("ambiguous"));
assert!(!truth_failure_is_routine(""));
assert!(!truth_failure_is_routine("not-found-ish"));
}
#[test]
fn family1_truth_nonzero_exit_is_unresolved_either_way() {
let mut not_found = std::process::Command::new("sh");
not_found.args([
"-c",
"printf '{\"state\":\"unknown\",\"reason\":\"not-found\"}'; exit 13",
]);
assert_eq!(
family1_truth_state_with_command(not_found, Duration::from_secs(1), "ses_1d9e"),
None
);
let mut broken = std::process::Command::new("sh");
broken.args([
"-c",
"printf '{\"state\":\"unknown\",\"reason\":\"transcript-unreadable\"}'; exit 13",
]);
assert_eq!(
family1_truth_state_with_command(broken, Duration::from_secs(1), "abcd1234"),
None
);
}
#[test]
fn family1_truth_failure_detail_prefers_stdout_reason() {
let detail =
family1_truth_failure_detail(br#"{"state":"unknown","reason":"not-found"}"#, "");
assert_eq!(detail, "not-found");
}
#[test]
fn family1_truth_failure_detail_falls_back_to_stderr() {
let detail = family1_truth_failure_detail(b"not json", " banner ");
assert_eq!(detail, "banner");
}
}