use crate::events::EventEmitter;
use crate::paths::{self, AgentsHome};
use crate::protocol::{
read_request, write_request, write_response, ErrorCode, Namespace, Request, Response,
};
use crate::state::{self, RegistryEntry};
use crate::AgentStatus;
use serde_json::{json, Map, Value};
use std::os::unix::process::CommandExt; use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::net::{UnixListener, UnixStream};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DaemonState {
ColdStart,
Recovering,
Serving,
IdlePendingExit,
ShuttingDown,
Exited,
}
impl DaemonState {
pub fn as_str(&self) -> &'static str {
match self {
DaemonState::ColdStart => "cold_start",
DaemonState::Recovering => "recovering",
DaemonState::Serving => "serving",
DaemonState::IdlePendingExit => "idle_pending_exit",
DaemonState::ShuttingDown => "shutting_down",
DaemonState::Exited => "exited",
}
}
}
#[derive(Debug, Clone)]
pub struct DaemonOptions {
pub idle_exit: Duration,
pub worker_bin: PathBuf,
pub reconcile_on_start: bool,
pub dead_row_grace: Duration,
pub notify_on_blocked: bool,
pub notify_on_done: bool,
}
impl Default for DaemonOptions {
fn default() -> Self {
DaemonOptions {
idle_exit: Duration::from_secs(1800),
worker_bin: resolve_worker_bin(),
reconcile_on_start: true,
dead_row_grace: Duration::from_secs(crate::agents_config::DEFAULT_DEAD_ROW_GRACE_SECS),
notify_on_blocked: true,
notify_on_done: false,
}
}
}
fn resolve_worker_bin() -> PathBuf {
if let Some(v) = std::env::var_os("FNO_AGENTS_WORKER_BIN") {
return PathBuf::from(v);
}
std::env::current_exe()
.ok()
.and_then(|p| p.parent().map(|d| d.join("fno-agents-worker")))
.unwrap_or_else(|| PathBuf::from("fno-agents-worker"))
}
#[derive(Debug, thiserror::Error)]
pub enum DaemonError {
#[error("io: {0}")]
Io(#[from] std::io::Error),
#[error("another daemon is already serving on {0}")]
AlreadyRunning(PathBuf),
#[error("socket permission invariant failed: {0}")]
Permission(String),
#[error("filesystem does not support advisory locking at {0}: {1}")]
FlockUnsupported(PathBuf, String),
#[error("state: {0}")]
State(#[from] state::StateError),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum InconsistencyReason {
MissingStateJson,
UnreadableStateJson,
}
impl InconsistencyReason {
pub fn as_str(&self) -> &'static str {
match self {
InconsistencyReason::MissingStateJson => "missing_state_json",
InconsistencyReason::UnreadableStateJson => "unreadable_state_json",
}
}
}
#[derive(Debug, Default, PartialEq)]
pub struct RecoveryReport {
pub inconsistent: Vec<(String, InconsistencyReason)>,
pub archived_orphans: Vec<String>,
pub reaped_pids: Vec<u32>,
pub recovered_drives: Vec<String>,
}
pub fn recover(home: &AgentsHome, emitter: &EventEmitter) -> RecoveryReport {
let mut report = RecoveryReport::default();
let registry = state::load_registry(&home.registry_json()).unwrap_or_default();
let registered: std::collections::BTreeSet<String> = registry
.entries
.iter()
.map(|e| e.short_id.clone())
.collect();
for entry in ®istry.entries {
let is_claude_shellout = entry.harness_name() == "claude"
&& entry.host_mode_or_default() != crate::state::HOST_MODE_INTERACTIVE;
if entry.short_id.is_empty() || is_claude_shellout {
continue;
}
let state_path = home.state_json(&entry.short_id);
match state::load_state(&state_path) {
Ok(Some(mut st)) => {
let taken = st.pty.as_mut().and_then(|p| p.take_active_drive());
if let Some(drive) = taken {
let mut fields = Map::new();
if let Some(sid) = &drive.session_id {
fields.insert("session_id".into(), Value::String(sid.clone()));
}
fields.insert("reason".into(), Value::String("daemon_restart".into()));
let _ = emitter.emit_fields("drive_crashed", fields);
let _ = state::write_state_atomic(&state_path, &st);
report.recovered_drives.push(entry.short_id.clone());
}
}
Ok(None) => {
let reason = InconsistencyReason::MissingStateJson;
let _ = emitter.emit_fields(
"agent_inconsistent",
json_obj(&[
("short_id", Value::String(entry.short_id.clone())),
("reason", Value::String(reason.as_str().into())),
]),
);
report.inconsistent.push((entry.short_id.clone(), reason));
}
Err(_) => {
let reason = InconsistencyReason::UnreadableStateJson;
let _ = emitter.emit_fields(
"agent_inconsistent",
json_obj(&[
("short_id", Value::String(entry.short_id.clone())),
("reason", Value::String(reason.as_str().into())),
]),
);
report.inconsistent.push((entry.short_id.clone(), reason));
}
}
}
if let Ok(read) = std::fs::read_dir(home.root()) {
for entry in read.flatten() {
if !entry.file_type().map(|t| t.is_dir()).unwrap_or(false) {
continue;
}
let name = match entry.file_name().into_string() {
Ok(n) if !n.starts_with('.') => n,
_ => continue,
};
if registered.contains(&name) {
continue;
}
if home.state_json(&name).exists() {
let ts = now_compact();
let dest = home.orphan_archive_dest(&name, &ts);
let _ = std::fs::create_dir_all(home.orphaned_dir());
if std::fs::rename(home.agent_dir(&name), &dest).is_ok() {
let _ = emitter.emit_fields(
"agent_orphan_state_archived",
json_obj(&[
("short_id", Value::String(name.clone())),
(
"archived_to",
Value::String(dest.to_string_lossy().into_owned()),
),
]),
);
report.archived_orphans.push(name);
}
}
}
}
let live_workers = home.scan_worker_sockets();
let mut to_reap: Vec<(String, u32)> = Vec::new();
for entry in ®istry.entries {
if live_workers.contains(&entry.short_id) {
continue; }
if let Some(pid) = entry.pid {
if !pid_is_ours(pid, entry.pid_start_time) {
to_reap.push((entry.short_id.clone(), pid));
}
}
}
if !to_reap.is_empty() {
let reaped: std::collections::BTreeSet<String> =
to_reap.iter().map(|(s, _)| s.clone()).collect();
for e in ®istry.entries {
if reaped.contains(&e.short_id) {
emit_inside_leg_completion(emitter, e);
}
}
if let Err(e) = state::update_registry(&home.registry_json(), |r| {
for e in r.entries.iter_mut() {
if reaped.contains(&e.short_id) {
e.status = AgentStatus::Exited;
e.inside_leg = None;
e.screen_state = None;
}
}
}) {
let _ = emitter.emit(
"daemon_recovery_error",
&json!({"op": "reap_orphans", "error": e.to_string()}),
);
}
for (short_id, pid) in to_reap {
let _ = emitter.emit_fields(
"agent_orphan_reaped",
json_obj(&[
("short_id", Value::String(short_id)),
("pid", Value::Number(pid.into())),
]),
);
report.reaped_pids.push(pid);
}
}
report
}
#[cfg(target_os = "linux")]
pub fn process_start_time(pid: u32) -> Option<u64> {
let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
let after = stat.rsplit_once(')')?.1;
after.split_whitespace().nth(19)?.parse::<u64>().ok()
}
#[cfg(target_os = "macos")]
pub fn process_start_time(pid: u32) -> Option<u64> {
use std::mem;
let mut info: libc::proc_bsdinfo = unsafe { mem::zeroed() };
let size = mem::size_of::<libc::proc_bsdinfo>() as libc::c_int;
let written = unsafe {
libc::proc_pidinfo(
pid as libc::c_int,
libc::PROC_PIDTBSDINFO,
0,
&mut info as *mut _ as *mut libc::c_void,
size,
)
};
if written != size {
return None;
}
Some(info.pbi_start_tvsec * 1_000_000 + info.pbi_start_tvusec)
}
#[cfg(not(any(target_os = "linux", target_os = "macos")))]
pub fn process_start_time(_pid: u32) -> Option<u64> {
None
}
#[derive(Debug, Default, PartialEq)]
pub struct GcSummary {
pub reaped: Vec<String>,
pub kept_dirty: Vec<(String, String)>,
}
fn worktree_clean_probe(cwd: &str) -> Option<bool> {
let out = std::process::Command::new("git")
.current_dir(cwd)
.args(["status", "--porcelain"])
.output()
.ok()?;
if !out.status.success() {
return None;
}
Some(out.stdout.iter().all(u8::is_ascii_whitespace))
}
fn now_epoch_secs() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
fn dispatch_node_id(name: &str) -> Option<String> {
let mut parts = name.split('-');
match parts.next()? {
"target" | "reconcile" => {}
_ => return None,
}
let prefix = parts.next()?;
let hex = parts.next()?;
if prefix.is_empty()
|| !prefix
.chars()
.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit())
|| hex.is_empty()
|| !hex.chars().all(|c| c.is_ascii_hexdigit())
{
return None;
}
Some(format!("{prefix}-{hex}"))
}
fn global_events_path(home: &AgentsHome) -> PathBuf {
home.root()
.parent()
.unwrap_or_else(|| home.root())
.join("events.jsonl")
}
#[derive(Debug, PartialEq, Eq)]
enum DispatchTermination {
Found(String),
Absent(Option<String>),
Unknown(String),
}
fn dispatch_target_session_id(
entry: &RegistryEntry,
node_id: &str,
) -> Result<Option<String>, String> {
let manifest = PathBuf::from(&entry.cwd).join(".fno/target-state.md");
let content = match std::fs::read_to_string(&manifest) {
Ok(content) => content,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(err) => return Err(format!("read {}: {err}", manifest.display())),
};
let parsed = crate::loop_target::parse_target_manifest(&content)
.ok_or_else(|| format!("parse target session from {}", manifest.display()))?;
if parsed.input != node_id {
return Err(format!(
"target manifest {} belongs to input {}, not registry node {node_id}",
manifest.display(),
parsed.input
));
}
if let (Some(row_session), Some(manifest_session)) = (
entry.harness_session_id.as_deref(),
parsed.harness_session_id.as_deref(),
) {
if row_session != manifest_session {
return Err(format!(
"target manifest {} harness session {manifest_session} does not match registry {row_session}",
manifest.display()
));
}
}
Ok(Some(parsed.session_id))
}
fn dispatch_termination(
home: &AgentsHome,
entry: &RegistryEntry,
node_id: &str,
) -> DispatchTermination {
let session_id = match dispatch_target_session_id(entry, node_id) {
Ok(session_id) => session_id,
Err(err) => return DispatchTermination::Unknown(err),
};
let Some(session_id) = session_id else {
return DispatchTermination::Absent(None);
};
let journal = crate::loop_runtime::Journal::new(
crate::loop_runtime::ProjectJournalPath(
PathBuf::from(&entry.cwd).join(".fno/events.jsonl"),
),
crate::loop_runtime::GlobalJournalPath(global_events_path(home)),
);
match journal.find_termination_strict(&session_id) {
Ok(Some(_)) => DispatchTermination::Found(session_id),
Ok(None) => DispatchTermination::Absent(Some(session_id)),
Err(err) => DispatchTermination::Unknown(err.to_string()),
}
}
fn record_dead_dispatch(
home: &AgentsHome,
entry: &RegistryEntry,
node_id: &str,
target_session_id: Option<&str>,
) -> Result<(), String> {
EventEmitter::new(global_events_path(home), "daemon")
.emit(
"node_failed",
&json!({
"unit_id": node_id,
"session_id": target_session_id.unwrap_or(&entry.short_id),
"iteration": 0,
"exit_code": 1,
"short_id": entry.short_id,
"reason": "agent-row-reaped-no-termination",
}),
)
.map_err(|err| err.to_string())
}
fn restore_unaccounted_row(home: &AgentsHome, entry: &RegistryEntry) -> Result<(), String> {
let mut restored = false;
state::update_registry(&home.registry_json(), |registry| {
if !registry.entries.iter().any(|row| row.name == entry.name) {
registry.entries.push(entry.clone());
restored = true;
}
})
.map_err(|err| err.to_string())?;
if restored {
Ok(())
} else {
Err(format!(
"could not restore {}: a replacement row now owns that name",
entry.name
))
}
}
pub fn gc_sweep(home: &AgentsHome, emitter: &EventEmitter, grace: Duration) -> GcSummary {
let mut summary = GcSummary::default();
let registry = state::load_registry(&home.registry_json()).unwrap_or_default();
if registry.entries.is_empty() {
return summary; }
let live_workers = home.scan_worker_sockets();
let now = now_epoch_secs();
let grace_secs = grace.as_secs() as i64;
let mut to_reap: std::collections::BTreeMap<String, String> = std::collections::BTreeMap::new();
let mut to_stamp: std::collections::BTreeMap<String, String> =
std::collections::BTreeMap::new();
let mut to_clear: std::collections::BTreeMap<String, String> =
std::collections::BTreeMap::new();
for e in ®istry.entries {
let is_live = live_workers.contains(&e.short_id)
|| e.pid
.map(|p| pid_is_ours(p, e.pid_start_time))
.unwrap_or(false);
let pid_confirmed_dead = e
.pid
.map(|p| !pid_is_ours(p, e.pid_start_time))
.unwrap_or(false);
let is_ask = e.is_one_shot_ask();
let exited_at = e
.exited_at
.as_deref()
.and_then(state::rfc3339_like_to_secs)
.map(|s| s as i64);
let terminal_or_dead = matches!(e.status, AgentStatus::Exited | AgentStatus::PermanentDead)
|| pid_confirmed_dead;
let past_grace = matches!(exited_at, Some(t) if now.saturating_sub(t) > grace_secs);
let needs_probe = !is_live && terminal_or_dead && past_grace && !is_ask;
let worktree_clean = if needs_probe {
worktree_clean_probe(&e.cwd)
} else {
None
};
let row = crate::gc::GcRow {
status: e.status,
is_live,
pid_confirmed_dead,
is_ask,
exited_at,
worktree_clean,
};
let id = if e.short_id.is_empty() {
e.name.clone()
} else {
e.short_id.clone()
};
match crate::gc::gc_action(&row, now, grace_secs) {
crate::gc::GcAction::Reap => {
to_reap.insert(e.name.clone(), e.created_at.clone());
}
crate::gc::GcAction::StampExit => {
to_stamp.insert(e.name.clone(), e.created_at.clone());
}
crate::gc::GcAction::Keep => {
if is_live && e.exited_at.is_some() {
to_clear.insert(e.name.clone(), e.created_at.clone());
} else if needs_probe && matches!(worktree_clean, Some(false) | None) {
summary.kept_dirty.push((id, e.cwd.clone()));
}
}
}
}
if to_reap.is_empty() && to_stamp.is_empty() && to_clear.is_empty() {
return summary;
}
let now_stamp = now_rfc3339_like();
let mut reaped_names: std::collections::BTreeSet<String> = std::collections::BTreeSet::new();
let write = state::update_registry(&home.registry_json(), |r| {
for e in r.entries.iter_mut() {
if to_stamp.get(&e.name) == Some(&e.created_at) {
e.exited_at = Some(now_stamp.clone());
}
if to_clear.get(&e.name) == Some(&e.created_at) {
e.exited_at = None;
}
}
r.entries.retain(|e| {
if to_reap.get(&e.name) == Some(&e.created_at) {
reaped_names.insert(e.name.clone());
false
} else {
true
}
});
});
match write {
Ok(()) => {
for e in ®istry.entries {
if reaped_names.contains(&e.name) {
let node_id = dispatch_node_id(&e.name);
let mut target_session_id = None;
let mut termination_event = false;
let mut accounted = true;
if let Some(node_id) = node_id.as_deref() {
match dispatch_termination(home, e, node_id) {
DispatchTermination::Found(session_id) => {
target_session_id = Some(session_id);
termination_event = true;
}
DispatchTermination::Absent(session_id) => {
target_session_id = session_id;
if let Err(err) = record_dead_dispatch(
home,
e,
node_id,
target_session_id.as_deref(),
) {
accounted = false;
let restore = restore_unaccounted_row(home, e);
let _ = emitter.emit(
"daemon_recovery_error",
&json!({
"op": "record_dead_dispatch",
"short_id": e.short_id,
"error": err,
"restore_error": restore.err(),
}),
);
}
}
DispatchTermination::Unknown(err) => {
accounted = false;
let restore = restore_unaccounted_row(home, e);
let _ = emitter.emit(
"daemon_recovery_error",
&json!({
"op": "observe_dead_dispatch_termination",
"short_id": e.short_id,
"error": err,
"restore_error": restore.err(),
}),
);
}
}
}
if !accounted {
continue;
}
let _ = emitter.emit_fields(
"agent_row_reaped",
json_obj(&[
("short_id", Value::String(e.short_id.clone())),
("name", Value::String(e.name.clone())),
(
"node_id",
node_id.clone().map_or(Value::Null, Value::String),
),
(
"session_id",
target_session_id.map_or(Value::Null, Value::String),
),
("termination_event", Value::Bool(termination_event)),
]),
);
summary.reaped.push(if e.short_id.is_empty() {
e.name.clone()
} else {
e.short_id.clone()
});
}
}
}
Err(err) => {
let _ = emitter.emit(
"daemon_recovery_error",
&json!({"op": "gc_sweep", "error": err.to_string()}),
);
summary.reaped.clear();
}
}
summary
}
async fn terminal_stop_sweep(home: &AgentsHome, emitter: &EventEmitter) {
let home_read = home.clone();
let loaded = tokio::task::spawn_blocking(move || {
let markers = crate::terminal_stop::read_markers(&home_read);
if markers.is_empty() {
return (markers, None);
}
let roster = crate::claude_roster::ClaudeRoster::load_default();
(markers, Some(roster))
})
.await;
let (markers, roster) = match loaded {
Ok(v) => v,
Err(e) => {
eprintln!("daemon: terminal-stop sweep: read task failed: {e}");
return;
}
};
if markers.is_empty() {
return;
}
let roster = match roster {
Some(Ok(r)) => r,
Some(Err(e)) => {
eprintln!("daemon: terminal-stop sweep: roster load failed: {e} (retry next tick)");
return;
}
None => return,
};
for marker in markers {
let short = roster.find(&marker.uuid).map(|w| w.short_id().to_string());
match crate::terminal_stop::stop_decision(short) {
crate::terminal_stop::StopAction::Stop(short) => {
let stop = tokio::process::Command::new("claude")
.arg("stop")
.arg(&short)
.kill_on_drop(true)
.output();
let stopped = tokio::time::timeout(Duration::from_secs(15), stop).await;
match stopped {
Err(_) => eprintln!("daemon: claude stop {short} timed out (retry next tick)"),
Ok(Ok(o)) if o.status.success() => {
let _ = emitter.emit(
"bg_worker_terminal_stopped",
&json!({
"short_id": short,
"session_id": marker.uuid,
"reason": marker.reason,
}),
);
crate::terminal_stop::remove_marker(home, &marker.uuid);
}
Ok(Ok(o)) => eprintln!(
"daemon: claude stop {short} failed: {}",
String::from_utf8_lossy(&o.stderr).trim()
),
Ok(Err(e)) => eprintln!("daemon: could not exec `claude stop`: {e}"),
}
}
crate::terminal_stop::StopAction::RemoveStale => {
crate::terminal_stop::remove_marker(home, &marker.uuid);
}
}
}
}
pub fn pid_is_ours(pid: u32, recorded: Option<u64>) -> bool {
if pid <= 1 || pid > i32::MAX as u32 {
return false;
}
if unsafe { libc::kill(pid as libc::pid_t, 0) } != 0 {
return false;
}
match (recorded, process_start_time(pid)) {
(Some(rec), Some(now)) => rec == now,
_ => true,
}
}
pub async fn bind_supervisor_socket(home: &AgentsHome) -> Result<UnixListener, DaemonError> {
home.ensure_root()?;
flock_self_test(home)?;
let sock = home.supervisor_sock();
if sock.exists() {
if UnixStream::connect(&sock).await.is_ok() {
return Err(DaemonError::AlreadyRunning(sock));
}
let _ = std::fs::remove_file(&sock);
}
let listener = match UnixListener::bind(&sock) {
Ok(l) => l,
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => {
return Err(DaemonError::AlreadyRunning(sock));
}
Err(e) => return Err(e.into()),
};
paths::set_file_mode_0600(&sock)?;
#[cfg(unix)]
{
if !paths::is_dir_mode_0700(home.root()) {
return Err(DaemonError::Permission(format!(
"{} is not mode 0700",
home.root().display()
)));
}
if !paths::is_file_mode_0600(&sock) {
return Err(DaemonError::Permission(format!(
"{} is not mode 0600",
sock.display()
)));
}
}
Ok(listener)
}
fn flock_self_test(home: &AgentsHome) -> Result<(), DaemonError> {
let probe = home.root().join(".flock-probe");
let file = std::fs::OpenOptions::new()
.create(true)
.read(true)
.write(true)
.truncate(false)
.open(&probe)
.map_err(|e| DaemonError::FlockUnsupported(probe.clone(), e.to_string()))?;
let lock_res = file.lock();
if lock_res.is_ok() {
let _ = file.unlock();
}
let _ = std::fs::remove_file(&probe);
lock_res.map_err(|e| DaemonError::FlockUnsupported(probe.clone(), e.to_string()))?;
Ok(())
}
pub async fn run(home: AgentsHome, opts: DaemonOptions) -> Result<(), DaemonError> {
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let listener = match bind_supervisor_socket(&home).await {
Ok(l) => l,
Err(DaemonError::AlreadyRunning(_)) => {
return Ok(());
}
Err(e) => return Err(e),
};
emit_state(&emitter, DaemonState::Recovering);
let report = recover(&home, &emitter);
if opts.reconcile_on_start {
let swept: Result<ReconcileSweepResult, String> =
if std::env::var("FNO_AGENTS_FAIL_STARTUP_RECONCILE").is_ok() {
Err(
"forced startup-reconcile failure (FNO_AGENTS_FAIL_STARTUP_RECONCILE)"
.to_string(),
)
} else {
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
run_reconcile_sweep(&home, &emitter)
}))
.unwrap_or_else(|_| {
Err(
"startup reconcile sweep panicked; serving last-recorded status"
.to_string(),
)
})
};
match swept {
Ok(result) => {
let _ = emitter.emit(
"startup_reconcile_done",
&json!({
"updated": result.outcome.updated.len(),
"deferred": result.outcome.deferred,
}),
);
}
Err(msg) => {
let _ = emitter.emit("startup_reconcile_failed", &json!({"error": msg}));
}
}
}
let started_at = Instant::now();
let exe_fingerprint = crate::drift::ExeFingerprint::current();
if exe_fingerprint.is_none() {
let _ = emitter.emit("daemon_exe_fingerprint_unavailable", &json!({}));
}
let pid_start_time = process_start_time(std::process::id());
let _ = emitter.emit(
"daemon_started",
&json!({
"pid": std::process::id(),
"version": env!("CARGO_PKG_VERSION"),
"recovered_drives": report.recovered_drives.len(),
}),
);
emit_state(&emitter, DaemonState::Serving);
let ctx = Arc::new(Ctx {
home,
emitter,
opts,
started_at,
exe_fingerprint,
pid_start_time,
pending_inside_leg: std::sync::Mutex::new(std::collections::HashMap::new()),
});
let ab_live = Arc::new(std::sync::atomic::AtomicBool::new(false));
let ab_shutdown = Arc::new(std::sync::atomic::AtomicBool::new(false));
let ab_handle = {
let fno_bin = std::env::var("FNO_BIN").unwrap_or_else(|_| "fno".to_string());
let ab_emitter = EventEmitter::new(ctx.home.events_jsonl(), "active-backlog");
let live = Arc::clone(&ab_live);
let shutdown = Arc::clone(&ab_shutdown);
tokio::spawn(crate::active_backlog::run_supervisor(
fno_bin, ab_emitter, live, shutdown,
))
};
let mut sigterm = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())?;
let mut idle_check = tokio::time::interval(Duration::from_secs(5));
idle_check.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut last_activity = Instant::now();
let scrape_in_flight = Arc::new(std::sync::atomic::AtomicBool::new(false));
let terminal_stop_in_flight = Arc::new(std::sync::atomic::AtomicBool::new(false));
loop {
tokio::select! {
accepted = listener.accept() => {
if let Ok((stream, _)) = accepted {
last_activity = Instant::now();
let ctx = Arc::clone(&ctx);
tokio::spawn(async move {
serve_connection(ctx, stream).await;
});
}
}
_ = sigterm.recv() => {
emit_state(&ctx.emitter, DaemonState::ShuttingDown);
let _ = ctx.emitter.emit("daemon_shutting_down", &json!({"reason": "sigterm"}));
break;
}
_ = idle_check.tick() => {
reap_zombies();
if !scrape_in_flight.swap(true, std::sync::atomic::Ordering::SeqCst) {
let flag = Arc::clone(&scrape_in_flight);
let home = ctx.home.clone();
let emitter = EventEmitter::new(ctx.home.events_jsonl(), "daemon");
let notify_on_blocked = ctx.opts.notify_on_blocked;
tokio::task::spawn_blocking(move || {
crate::scrape::scrape_sweep(&home, &emitter, notify_on_blocked);
flag.store(false, std::sync::atomic::Ordering::SeqCst);
});
}
let _ = gc_sweep(&ctx.home, &ctx.emitter, ctx.opts.dead_row_grace);
if !terminal_stop_in_flight.swap(true, std::sync::atomic::Ordering::SeqCst) {
let flag = Arc::clone(&terminal_stop_in_flight);
let home = ctx.home.clone();
let emitter = EventEmitter::new(ctx.home.events_jsonl(), "daemon");
tokio::spawn(async move {
terminal_stop_sweep(&home, &emitter).await;
flag.store(false, std::sync::atomic::Ordering::SeqCst);
});
}
let empty = state::load_registry(&ctx.home.registry_json())
.map(|r| r.entries.is_empty())
.unwrap_or(true);
let ab_active = ab_live.load(std::sync::atomic::Ordering::SeqCst);
if empty && !ab_active && last_activity.elapsed() >= ctx.opts.idle_exit {
emit_state(&ctx.emitter, DaemonState::IdlePendingExit);
let _ = ctx.emitter.emit("daemon_idle_pending_exit", &json!({}));
emit_state(&ctx.emitter, DaemonState::ShuttingDown);
let _ = ctx.emitter.emit(
"daemon_shutting_down",
&json!({"reason": "idle"}),
);
break;
}
}
}
}
ab_shutdown.store(true, std::sync::atomic::Ordering::SeqCst);
ab_handle.abort();
let _ = std::fs::remove_file(ctx.home.supervisor_sock());
emit_state(&ctx.emitter, DaemonState::Exited);
let _ = ctx.emitter.emit("daemon_exited", &json!({"clean": true}));
Ok(())
}
struct Ctx {
home: AgentsHome,
emitter: EventEmitter,
opts: DaemonOptions,
started_at: Instant,
exe_fingerprint: Option<crate::drift::ExeFingerprint>,
pid_start_time: Option<u64>,
pending_inside_leg: std::sync::Mutex<std::collections::HashMap<String, state::InsideLegReport>>,
}
const PENDING_INSIDE_LEG_CAP: usize = 64;
fn emit_state(emitter: &EventEmitter, state: DaemonState) {
let _ = emitter.emit("daemon_state", &json!({"state": state.as_str()}));
}
const CONN_READ_TIMEOUT: Duration = Duration::from_secs(30);
async fn serve_connection(ctx: Arc<Ctx>, mut stream: UnixStream) {
let req = match tokio::time::timeout(CONN_READ_TIMEOUT, read_request(&mut stream)).await {
Err(_elapsed) => return, Ok(Ok(r)) => r,
Ok(Err(crate::protocol::ProtocolError::UnexpectedEof)) => return, Ok(Err(e)) => {
let resp = Response::err(0, ErrorCode::MalformedFrame, format!("{e}"));
let _ = write_response(&mut stream, &resp).await;
return;
}
};
if req.method == "agent.logs" {
crate::logs::handle_logs(&ctx.home, &req, stream).await;
return;
}
let resp = dispatch(&ctx, &req).await;
let _ = write_response(&mut stream, &resp).await;
}
async fn run_blocking<F>(ctx: &Arc<Ctx>, req: &Request, f: F) -> Response
where
F: FnOnce(&Ctx, &Request) -> Response + Send + 'static,
{
let ctx = Arc::clone(ctx);
let req = req.clone();
let id = req.id;
match tokio::task::spawn_blocking(move || f(&ctx, &req)).await {
Ok(resp) => resp,
Err(_) => Response::err(id, ErrorCode::Internal, "handler task panicked"),
}
}
async fn load_registry_offloaded(path: PathBuf) -> state::Registry {
tokio::task::spawn_blocking(move || state::load_registry(&path))
.await
.ok()
.and_then(|r| r.ok())
.unwrap_or_default()
}
async fn update_registry_offloaded<F, T>(path: PathBuf, f: F) -> Result<T, state::StateError>
where
F: FnOnce(&mut state::Registry) -> T + Send + 'static,
T: Send + 'static,
{
match tokio::task::spawn_blocking(move || state::update_registry(&path, f)).await {
Ok(result) => result,
Err(e) => Err(state::StateError::Io(std::io::Error::other(format!(
"update_registry task panicked: {e}"
)))),
}
}
async fn dispatch(ctx: &Arc<Ctx>, req: &Request) -> Response {
match Namespace::of(&req.method) {
Namespace::Agent => dispatch_agent(ctx, req).await,
Namespace::Channel => dispatch_channel(ctx, req).await,
Namespace::Unknown => Response::err(
req.id,
ErrorCode::UnknownMethod,
format!("unknown namespace for method `{}`", req.method),
),
}
}
async fn dispatch_agent(ctx: &Arc<Ctx>, req: &Request) -> Response {
match Namespace::verb(&req.method) {
Some("spawn") => handle_spawn(ctx, req).await,
Some("ask") => handle_ask(ctx, req).await,
Some("switchboard") | Some("switchboard_v2") => handle_switchboard(ctx, req).await,
Some("stop") => handle_stop(ctx, req).await,
Some("rm") => handle_rm(ctx, req).await,
Some("list") => run_blocking(ctx, req, handle_list).await,
Some("status") => handle_status(ctx, req).await,
Some("reconcile") => run_blocking(ctx, req, handle_reconcile).await,
Some("report") => run_blocking(ctx, req, handle_report).await,
_ => Response::err(
req.id,
ErrorCode::UnknownMethod,
format!("unknown agent verb in `{}`", req.method),
),
}
}
fn valid_agent_name(name: &str) -> bool {
!name.is_empty()
&& name.len() <= 64
&& name
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
}
fn derive_short_id(name: &str, registry: &state::Registry) -> String {
let base: String = name
.chars()
.filter(|c| c.is_ascii_alphanumeric())
.take(8)
.collect();
let base = if base.is_empty() {
"agent".into()
} else {
base
};
if registry.entries.iter().all(|e| e.short_id != base) {
return base;
}
for n in 1..10_000 {
let cand = format!("{base}{n}");
if registry.entries.iter().all(|e| e.short_id != cand) {
return cand;
}
}
format!("{base}-{}", now_compact())
}
pub fn entry_holds_session(e: &RegistryEntry, uuid: &str) -> bool {
e.codex_session_id.as_deref() == Some(uuid)
|| e.gemini_session_id.as_deref() == Some(uuid)
|| e.session_id.as_deref() == Some(uuid)
|| e.claude_session_uuid.as_deref() == Some(uuid)
}
fn is_non_terminal(s: AgentStatus) -> bool {
!matches!(s, AgentStatus::Exited | AgentStatus::PermanentDead)
}
async fn handle_spawn(ctx: &Ctx, req: &Request) -> Response {
let p = &req.params;
let name = match p.get("name").and_then(|v| v.as_str()) {
Some(n) if valid_agent_name(n) => n.to_string(),
Some(_) => {
return Response::err(
req.id,
ErrorCode::InvalidParams,
"name must be 1-64 chars of [A-Za-z0-9_-]",
)
}
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `name`"),
};
let provider = p
.get("provider")
.and_then(|v| v.as_str())
.unwrap_or("codex")
.to_string();
let cwd = match p.get("cwd").and_then(|v| v.as_str()) {
Some(c) => PathBuf::from(c),
None => {
let fallback = std::env::temp_dir();
let _ = ctx.emitter.emit(
"agent_spawn_cwd_fallback",
&json!({"name": name, "fallback": fallback.to_string_lossy()}),
);
fallback
}
};
let host_mode = p
.get("host_mode")
.and_then(|v| v.as_str())
.unwrap_or(crate::state::HOST_MODE_EXEC);
let resume_id = p
.get("resume_id")
.and_then(|v| v.as_str())
.map(|s| s.to_string());
if host_mode == crate::state::HOST_MODE_INTERACTIVE && provider == "claude" {
let claude_mode = p
.get("mode")
.and_then(|v| v.as_str())
.unwrap_or(crate::state::CLAUDE_MODE_STREAM_JSON);
if claude_mode != crate::state::CLAUDE_MODE_INTERACTIVE {
let explicit_argv = p.get("argv").and_then(|v| v.as_array()).map(|a| {
a.iter()
.filter_map(|v| v.as_str().map(String::from))
.collect::<Vec<String>>()
});
return spawn_claude_stream_lane(
ctx,
req,
&name,
&cwd,
resume_id.as_deref(),
explicit_argv,
)
.await;
}
}
let _ = ctx.emitter.emit(
"agent_spawn_failed",
&json!({"name": name, "reason": "daemon_pty_hosting_retired", "provider": provider}),
);
Response::err(
req.id,
ErrorCode::InvalidParams,
"daemon PTY hosting was retired at G4 (x-f54c): spawn a mux-hosted agent pane with \
`fno agents spawn --substrate pane`, or use `--substrate bg|headless`. The daemon \
serves only claude stream-json adoption (host_mode=interactive, mode=stream_json).",
)
}
fn stream_claim_holder(short_id: &str) -> String {
format!("stream:{short_id}")
}
fn is_live_writer(status: AgentStatus) -> bool {
matches!(
status,
AgentStatus::Live
| AgentStatus::Ready
| AgentStatus::Idle
| AgentStatus::Busy
| AgentStatus::Spawning
| AgentStatus::Restarting
)
}
fn claude_stream_worker_args(
short_id: &str,
home: &std::path::Path,
cwd: &std::path::Path,
uuid: &str,
holder: &str,
child_argv: &[String],
) -> Vec<String> {
let mut args = vec![
"--stream".into(),
"--short-id".into(),
short_id.into(),
"--home".into(),
home.to_string_lossy().into_owned(),
"--cwd".into(),
cwd.to_string_lossy().into_owned(),
"--session-uuid".into(),
uuid.into(),
"--holder".into(),
holder.into(),
"--".into(),
];
args.extend(child_argv.iter().cloned());
args
}
fn build_claude_stream_entry(
name: &str,
short_id: &str,
cwd: &std::path::Path,
uuid: &str,
pid: u32,
pid_start_time: Option<u64>,
log_path: PathBuf,
) -> RegistryEntry {
let cwd_s = cwd.to_string_lossy().into_owned();
RegistryEntry {
name: name.into(),
short_id: short_id.into(),
legacy_provider: String::new(),
harness: Some("claude".into()),
harness_session_id: Some(uuid.into()),
cwd: cwd_s.clone(),
project_root: cwd_s,
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: Some(uuid.into()),
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
cc_session_id: None,
host_mode: Some(crate::state::HOST_MODE_INTERACTIVE.into()),
status: AgentStatus::Live,
last_message_at: Some(now_rfc3339_like()),
created_at: now_rfc3339_like(),
pid: Some(pid),
pid_start_time,
log_path: Some(log_path.to_string_lossy().into_owned()),
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
}
}
#[derive(Debug)]
enum ClaimOutcome {
Acquired,
HeldByOther(String),
Unavailable(String),
}
fn acquire_session_claim(uuid: &str, holder: &str) -> ClaimOutcome {
match crate::claims::acquire(
&format!("session:{uuid}"),
holder,
crate::claims::AcquireOpts::default(),
) {
crate::claims::AcquireOutcome::Acquired(_) => ClaimOutcome::Acquired,
crate::claims::AcquireOutcome::HeldByOther { holder, .. } => {
ClaimOutcome::HeldByOther(holder)
}
crate::claims::AcquireOutcome::Error(e) => ClaimOutcome::Unavailable(e),
}
}
struct DaemonClaimGuard {
session_uuid: String,
holder: String,
armed: bool,
}
impl DaemonClaimGuard {
fn disarm(mut self) {
self.armed = false;
}
}
impl Drop for DaemonClaimGuard {
fn drop(&mut self) {
if !self.armed {
return;
}
let _ = crate::claims::release(
&format!("session:{}", self.session_uuid),
&self.holder,
None,
None,
);
}
}
async fn stream_worker_reports_child_alive(sock: &std::path::Path) -> bool {
let probe = async {
let mut conn = UnixStream::connect(sock).await.ok()?;
write_request(&mut conn, &Request::new(1, "stream.status", json!({})))
.await
.ok()?;
let resp = crate::protocol::read_response(&mut conn).await.ok()?;
Some(
resp.result()
.and_then(|r| r.get("child_alive"))
.and_then(Value::as_bool)
.unwrap_or(false),
)
};
matches!(
tokio::time::timeout(Duration::from_secs(STREAM_PROBE_TIMEOUT_S), probe).await,
Ok(Some(true))
)
}
async fn spawn_claude_stream_lane(
ctx: &Ctx,
req: &Request,
name: &str,
cwd: &std::path::Path,
resume_id: Option<&str>,
explicit_argv: Option<Vec<String>>,
) -> Response {
let uuid = match resume_id {
Some(u) if !u.trim().is_empty() => u,
_ => {
let _ = ctx.emitter.emit(
"agent_spawn_failed",
&json!({"name": name, "reason": "claude_host_needs_from"}),
);
return Response::err(
req.id,
ErrorCode::InvalidParams,
"claude has no fresh interactive host; adopt an idle session: `fno agents promote <name> --from <session-uuid> --provider claude`",
);
}
};
let registry = load_registry_offloaded(ctx.home.registry_json()).await;
if let Some(existing) = registry.find(name) {
return Response::err(
req.id,
ErrorCode::AgentExists,
format!(
"agent {name} already exists (short_id={}); use `fno agents rm` first",
existing.short_id
),
);
}
if let Some(h) = registry.entries.iter().rev().find(|e| {
e.harness_name() == "claude"
&& e.claude_session_uuid.as_deref() == Some(uuid)
&& is_live_writer(e.status)
}) {
return Response::err(
req.id,
ErrorCode::InvalidParams,
format!(
"session '{uuid}' is already hosted by live stream thread '{}'; one writer per session",
h.name
),
);
}
let short_id = derive_short_id(name, ®istry);
let holder = stream_claim_holder(&short_id);
let uuid_owned = uuid.to_string();
let holder_for_acq = holder.clone();
let claim_outcome =
tokio::task::spawn_blocking(move || acquire_session_claim(&uuid_owned, &holder_for_acq))
.await
.unwrap_or_else(|e| ClaimOutcome::Unavailable(format!("claim task panicked: {e}")));
let claim_guard = match claim_outcome {
ClaimOutcome::Acquired => DaemonClaimGuard {
session_uuid: uuid.to_string(),
holder: holder.clone(),
armed: true,
},
ClaimOutcome::HeldByOther(who) => {
let _ = ctx.emitter.emit(
"agent_spawn_failed",
&json!({"name": name, "reason": "session_claimed", "detail": who}),
);
return Response::err(
req.id,
ErrorCode::InvalidParams,
format!(
"session '{uuid}' is held by another writer ({who}); refusing to double-adopt"
),
);
}
ClaimOutcome::Unavailable(why) => {
let _ = ctx.emitter.emit(
"agent_stream_claim_unavailable",
&json!({"name": name, "session_uuid": uuid, "detail": why}),
);
DaemonClaimGuard {
session_uuid: uuid.to_string(),
holder: holder.clone(),
armed: false,
}
}
};
let child_argv =
explicit_argv.unwrap_or_else(|| crate::provider::claude_stream_json_resume_argv(uuid));
let worker_args =
claude_stream_worker_args(&short_id, ctx.home.root(), cwd, uuid, &holder, &child_argv);
let mut cmd = std::process::Command::new(&ctx.opts.worker_bin);
cmd.args(&worker_args);
cmd.process_group(0);
let child = match cmd.spawn() {
Ok(c) => c,
Err(e) => {
let _ = ctx.emitter.emit(
"agent_spawn_failed",
&json!({"name": name, "reason": "binary_not_found", "detail": e.to_string()}),
);
return Response::err(
req.id,
ErrorCode::SpawnFailed,
format!("could not launch stream worker: {e}"),
);
}
};
let worker_pid = child.id();
let worker_pid_start_time = process_start_time(worker_pid);
drop(child);
let sock = ctx.home.worker_sock(&short_id);
let start = Instant::now();
while !sock.exists() && start.elapsed() < Duration::from_secs(10) {
tokio::time::sleep(Duration::from_millis(25)).await;
}
if !sock.exists() {
let _ = ctx.emitter.emit(
"agent_create_no_session",
&json!({"name": name, "short_id": short_id, "lane": "stream"}),
);
return Response::err(
req.id,
ErrorCode::SpawnFailed,
"stream worker did not come up within 10s",
);
}
if !is_live_stream_thread(&sock).await {
best_effort_worker_shutdown(&sock).await;
let _ = ctx.emitter.emit(
"agent_create_no_session",
&json!({"name": name, "short_id": short_id, "reason": "not_a_stream_thread"}),
);
return Response::err(
req.id,
ErrorCode::SpawnFailed,
"stream worker came up but does not serve the stream protocol",
);
}
if !stream_worker_reports_child_alive(&sock).await {
best_effort_worker_shutdown(&sock).await;
let _ = ctx.emitter.emit(
"agent_create_no_session",
&json!({"name": name, "short_id": short_id, "reason": "resume_child_exited"}),
);
return Response::err(
req.id,
ErrorCode::SpawnFailed,
"claude --resume child exited before adoption (bad/expired session id, auth failure, or dead cwd)",
);
}
let entry = build_claude_stream_entry(
name,
&short_id,
cwd,
uuid,
worker_pid,
worker_pid_start_time,
ctx.home.timeline_jsonl(&short_id),
);
let uuid_for_lock = uuid.to_string();
let insert = update_registry_offloaded(ctx.home.registry_json(), move |r| {
if r.entries.iter().any(|e| e.name == entry.name) {
return false;
}
if r.entries.iter().any(|e| {
e.harness_name() == "claude"
&& e.claude_session_uuid.as_deref() == Some(&uuid_for_lock)
&& is_live_writer(e.status)
}) {
return false;
}
r.entries.push(entry);
true
})
.await;
match insert {
Ok(true) => flush_buffered_inside_leg(ctx, uuid, name),
Ok(false) => {
best_effort_worker_shutdown(&sock).await;
let _ = ctx.emitter.emit(
"agent_spawn_failed",
&json!({"name": name, "short_id": short_id, "reason": "session_taken_concurrent"}),
);
return Response::err(
req.id,
ErrorCode::AgentExists,
format!("session '{uuid}' was adopted by a concurrent call; this one refused"),
);
}
Err(e) => {
best_effort_worker_shutdown(&sock).await;
let _ = ctx.emitter.emit(
"agent_spawn_failed",
&json!({"name": name, "short_id": short_id, "reason": "registry_write_failed"}),
);
return Response::err(req.id, ErrorCode::Internal, format!("registry write: {e}"));
}
}
claim_guard.disarm();
let _ = ctx.emitter.emit(
"agent_spawned",
&json!({"name": name, "provider": "claude", "short_id": short_id, "lane": "stream", "session_uuid": uuid}),
);
Response::ok(
req.id,
json!({"short_id": short_id, "provider": "claude", "status": "live", "lane": "stream"}),
)
}
fn provider_readiness_detector(provider: &str) -> Box<dyn crate::readiness::ReadinessDetector> {
use crate::provider::ProviderWithPty as _;
match provider {
"codex" => crate::provider::CodexProvider.readiness_detector(),
"gemini" => crate::provider::GeminiProvider.readiness_detector(),
"agy" => crate::provider::AgyProvider.readiness_detector(),
"opencode" => crate::provider::OpencodeProvider.readiness_detector(),
"claude" => crate::provider::ClaudeInteractiveProvider.readiness_detector(),
other => Box::new(crate::readiness::NoSignalDetector {
provider: other.to_string(),
}),
}
}
#[derive(Debug, PartialEq, Eq)]
enum PollError {
Timeout { secs: u64 },
WorkerUnresponsive { secs: u64 },
}
impl std::fmt::Display for PollError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
PollError::Timeout { secs } => {
write!(f, "ask timed out after {secs}s before reply settled")
}
PollError::WorkerUnresponsive { secs } => write!(
f,
"ask timed out after {secs}s before reply settled (worker snapshot read did not return)"
),
}
}
}
async fn poll_until_ready<F, Fut>(
fetcher: F,
detector: Box<dyn crate::readiness::ReadinessDetector>,
poll_interval: Duration,
timeout: Duration,
) -> Result<String, PollError>
where
F: Fn() -> Fut,
Fut: std::future::Future<Output = Option<String>>,
{
let deadline = tokio::time::Instant::now() + timeout;
loop {
let now = tokio::time::Instant::now();
if now >= deadline {
return Err(PollError::Timeout {
secs: timeout.as_secs(),
});
}
let remaining = deadline.saturating_duration_since(now);
let fetched = match tokio::time::timeout(remaining, fetcher()).await {
Ok(opt) => opt,
Err(_) => {
return Err(PollError::WorkerUnresponsive {
secs: timeout.as_secs(),
})
}
};
if let Some(text) = fetched {
let mut grid = crate::screen::TerminalGrid::with_default_size();
grid.feed(text.as_bytes());
let owned = grid.snapshot();
let view = owned.view();
match detector.is_ready(&view) {
Ok(true) => return Ok(owned.text),
Ok(false) | Err(_) => {} }
}
tokio::time::sleep(poll_interval).await;
}
}
async fn handle_ask(ctx: &Ctx, req: &Request) -> Response {
let name = match req.params.get("name").and_then(|v| v.as_str()) {
Some(n) => n.to_string(),
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `name`"),
};
let message = req
.params
.get("message")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let provider_param = req
.params
.get("provider")
.and_then(|v| v.as_str())
.map(String::from);
let cwd_param = req
.params
.get("cwd")
.and_then(|v| v.as_str())
.map(PathBuf::from);
let from_name_param = req
.params
.get("from_name")
.and_then(|v| v.as_str())
.map(String::from);
let yolo_param = req
.params
.get("yolo")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let registry = load_registry_offloaded(ctx.home.registry_json()).await;
let entry = match registry.find(&name) {
Some(e) => e.clone(),
None => {
let provider = match provider_param {
Some(p) => p,
None => {
return Response::err(
req.id,
ErrorCode::InvalidParams,
format!(
"agent '{name}' not found; pass --provider to create it on first contact"
),
)
}
};
let spawn_cwd = match cwd_param {
Some(c) => c,
None => {
let fallback = std::env::temp_dir();
let _ = ctx.emitter.emit(
"agent_spawn_cwd_fallback",
&json!({
"name": name,
"fallback": fallback.to_string_lossy(),
"via": "ask_first_contact",
}),
);
fallback
}
};
let mut spawn_params = serde_json::Map::new();
spawn_params.insert("name".into(), serde_json::Value::String(name.clone()));
spawn_params.insert("provider".into(), serde_json::Value::String(provider));
spawn_params.insert(
"cwd".into(),
serde_json::Value::String(spawn_cwd.to_str().unwrap_or(".").to_string()),
);
spawn_params.insert("message".into(), serde_json::Value::String(message.clone()));
if let Some(ref fn_val) = from_name_param {
spawn_params.insert(
"from_name".into(),
serde_json::Value::String(fn_val.clone()),
);
}
if yolo_param {
spawn_params.insert("yolo".into(), serde_json::Value::Bool(true));
}
let spawn_req = Request::new(
req.id,
"agent.spawn",
serde_json::Value::Object(spawn_params),
);
let spawn_resp = handle_spawn(ctx, &spawn_req).await;
return match spawn_resp.payload {
crate::protocol::ResponsePayload::Ok(ref result) => {
let short_id = result
.get("short_id")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
Response::ok(req.id, json!({"created": true, "short_id": short_id}))
}
crate::protocol::ResponsePayload::Err(_) => spawn_resp,
};
}
};
if entry.status == AgentStatus::Orphaned {
return Response::err(
req.id,
ErrorCode::InvalidStatus,
format!("agent {name} is orphaned; use `fno agents reconcile` or `rm`"),
);
}
let sock = ctx.home.worker_sock(&entry.short_id);
let mut conn = match UnixStream::connect(&sock).await {
Ok(c) => c,
Err(_) => {
return Response::err(
req.id,
ErrorCode::InvalidStatus,
format!("worker for {name} is not reachable"),
)
}
};
let mut payload = message.clone();
if !payload.ends_with('\n') {
payload.push('\n');
}
if write_request(
&mut conn,
&Request::new(1, "worker.write", json!({"data": payload})),
)
.await
.is_err()
{
return Response::err(req.id, ErrorCode::Internal, "worker write failed");
}
match crate::protocol::read_response(&mut conn).await {
Ok(ack) if ack.is_err() => {
let msg = ack
.error()
.map(|e| e.message.clone())
.unwrap_or_else(|| "worker rejected the write".into());
return Response::err(req.id, ErrorCode::Internal, msg);
}
Ok(_) => {}
Err(_) => {
return Response::err(req.id, ErrorCode::Internal, "no write-ack from worker");
}
}
let timeout_secs = req
.params
.get("timeout")
.and_then(|v| v.as_u64())
.unwrap_or(600);
let detector = provider_readiness_detector(entry.harness_name());
let sock_path = sock.clone();
let fetcher = move || {
let p = sock_path.clone();
async move { read_worker_snapshot(&p).await }
};
let reply = match poll_until_ready(
fetcher,
detector,
Duration::from_millis(200),
Duration::from_secs(timeout_secs),
)
.await
{
Ok(text) => text,
Err(e) => {
return Response::err(req.id, ErrorCode::Internal, e.to_string());
}
};
let ask_name = name.clone();
let _ = update_registry_offloaded(ctx.home.registry_json(), move |r| {
if let Some(e) = r.find_mut(&ask_name) {
e.last_message_at = Some(now_rfc3339_like());
}
})
.await;
let _ = ctx
.emitter
.emit("agent_ask_done", &json!({"name": name, "backend": "pty"}));
Response::ok(req.id, json!({"reply": reply, "backend": "pty"}))
}
const MAX_INJECT_BODY_BYTES: usize = 16 * 1024 * 1024;
const SWITCHBOARD_TURN_TIMEOUT_MS: u64 = 120_000;
const SWITCHBOARD_POLL_MS: u64 = 50;
const STREAM_PROBE_TIMEOUT_S: u64 = 2;
const SWITCHBOARD_MIRROR_TIMEOUT_S: u64 = 5;
const SWITCHBOARD_DRIVE_GRACE_S: u64 = 5;
struct SwitchboardTurn {
reply: String,
is_error: bool,
saw_receipt: bool,
}
async fn is_live_stream_thread(sock: &std::path::Path) -> bool {
let probe = async {
let mut conn = UnixStream::connect(sock).await.ok()?;
write_request(&mut conn, &Request::new(1, "stream.ping", json!({})))
.await
.ok()?;
let resp = crate::protocol::read_response(&mut conn).await.ok()?;
Some(!resp.is_err())
};
matches!(
tokio::time::timeout(Duration::from_secs(STREAM_PROBE_TIMEOUT_S), probe).await,
Ok(Some(true))
)
}
async fn drive_stream_turn(
worker_sock: &std::path::Path,
text: &str,
deadline: Duration,
) -> Result<SwitchboardTurn, String> {
let mut conn = UnixStream::connect(worker_sock)
.await
.map_err(|e| format!("target not live (worker unreachable): {e}"))?;
write_request(
&mut conn,
&Request::new(0, "stream.read_frames", json!({ "cursor": u64::MAX })),
)
.await
.map_err(|e| format!("cursor probe send failed: {e}"))?;
let probe = crate::protocol::read_response(&mut conn)
.await
.map_err(|e| format!("cursor probe recv failed: {e}"))?;
let mut cursor = probe
.result()
.and_then(|r| r.get("next"))
.and_then(|v| v.as_u64())
.ok_or_else(|| "cursor probe returned no result".to_string())?;
write_request(
&mut conn,
&Request::new(1, "stream.write_turn", json!({ "text": text })),
)
.await
.map_err(|e| format!("write_turn send failed: {e}"))?;
match crate::protocol::read_response(&mut conn).await {
Ok(ack) if ack.is_err() => {
return Err(format!(
"write_turn rejected: {}",
ack.error().map(|e| e.message.as_str()).unwrap_or("?")
))
}
Ok(_) => {}
Err(e) => return Err(format!("no write_turn ack: {e}")),
}
let start = Instant::now();
let mut reply = String::new();
let mut saw_receipt = false;
let mut req_id = 100u64;
loop {
if start.elapsed() > deadline {
return Err("turn timed out before result".into());
}
write_request(
&mut conn,
&Request::new(req_id, "stream.read_frames", json!({ "cursor": cursor })),
)
.await
.map_err(|e| format!("read_frames send failed: {e}"))?;
req_id += 1;
let resp = crate::protocol::read_response(&mut conn)
.await
.map_err(|e| format!("read_frames recv failed: {e}"))?;
let res = resp
.result()
.ok_or_else(|| "read_frames returned no result".to_string())?;
if let Some(next) = res.get("next").and_then(|v| v.as_u64()) {
cursor = next;
}
let child_alive = res
.get("child_alive")
.and_then(|v| v.as_bool())
.unwrap_or(true);
if let Some(frames) = res.get("frames").and_then(|v| v.as_array()) {
for fr in frames {
match fr.get("kind").and_then(|k| k.as_str()) {
Some("user_echo") => saw_receipt = true,
Some("assistant") => {
if let Some(t) = fr.get("text").and_then(|t| t.as_str()) {
reply.push_str(t);
}
}
Some("result") => {
let is_error = fr
.get("is_error")
.and_then(|e| e.as_bool())
.unwrap_or(false);
if reply.is_empty() {
if let Some(r) = fr.get("result").and_then(|r| r.as_str()) {
reply.push_str(r);
}
}
return Ok(SwitchboardTurn {
reply,
is_error,
saw_receipt,
});
}
_ => {}
}
}
}
if !child_alive {
return Err("target child exited before result (orphaned)".into());
}
tokio::time::sleep(Duration::from_millis(SWITCHBOARD_POLL_MS)).await;
}
}
async fn mirror_into(worker_sock: &std::path::Path, text: &str) -> Result<(), String> {
let inner = async {
let mut conn = UnixStream::connect(worker_sock)
.await
.map_err(|e| format!("mirror target unreachable: {e}"))?;
write_request(
&mut conn,
&Request::new(1, "stream.write_turn", json!({ "text": text })),
)
.await
.map_err(|e| format!("mirror write failed: {e}"))?;
match crate::protocol::read_response(&mut conn).await {
Ok(ack) if ack.is_err() => Err(format!(
"mirror rejected: {}",
ack.error().map(|e| e.message.as_str()).unwrap_or("?")
)),
Ok(_) => Ok(()),
Err(e) => Err(format!("no mirror ack: {e}")),
}
};
match tokio::time::timeout(Duration::from_secs(SWITCHBOARD_MIRROR_TIMEOUT_S), inner).await {
Ok(r) => r,
Err(_) => Err("mirror timed out".into()),
}
}
async fn stamp_orphaned(
home: &AgentsHome,
name: String,
identity: Value,
) -> Result<bool, state::StateError> {
update_registry_offloaded(home.registry_json(), move |registry| {
let Some(entry) = registry.find_mut(&name) else {
return false;
};
if !switchboard_identity_matches(entry, &identity) {
return false;
}
if entry.status == AgentStatus::Live {
entry.status = AgentStatus::Orphaned;
}
true
})
.await
}
fn switchboard_identity_matches(entry: &RegistryEntry, identity: &Value) -> bool {
let Some(expected) = identity.as_object() else {
return false;
};
let Some(harness) = expected.get("harness").and_then(Value::as_str) else {
return false;
};
let Some(short_id) = expected.get("short_id").and_then(Value::as_str) else {
return false;
};
let Some(created_at) = expected.get("created_at").and_then(Value::as_str) else {
return false;
};
let session_id = match expected.get("session_id") {
Some(Value::Null) => None,
Some(Value::String(value)) => Some(value.as_str()),
_ => return false,
};
entry.harness_name() == harness
&& entry.harness_session_id.as_deref() == session_id
&& entry.short_id == short_id
&& entry.created_at == created_at
}
async fn handle_switchboard(ctx: &Ctx, req: &Request) -> Response {
let to = match req.params.get("to").and_then(|v| v.as_str()) {
Some(s) => s.to_string(),
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `to`"),
};
let from = req
.params
.get("from")
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string();
let body = match req.params.get("body").and_then(|v| v.as_str()) {
Some(b) => b.to_string(),
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `body`"),
};
let mirror = req
.params
.get("mirror")
.and_then(|v| v.as_bool())
.unwrap_or(true);
let recipient_identity = match req.params.get("recipient_identity") {
Some(value) if value.is_object() => value,
_ => {
return Response::err(
req.id,
ErrorCode::InvalidParams,
"missing `recipient_identity`",
)
}
};
let from_identity = req.params.get("from_identity");
if mirror && !from_identity.is_some_and(Value::is_object) {
return Response::err(
req.id,
ErrorCode::InvalidParams,
"missing `from_identity` for mirrored switchboard turn",
);
}
let timeout_ms = req
.params
.get("timeout_ms")
.and_then(|v| v.as_u64())
.unwrap_or(SWITCHBOARD_TURN_TIMEOUT_MS);
if body.len() > MAX_INJECT_BODY_BYTES {
return Response::err(
req.id,
ErrorCode::InvalidParams,
format!(
"body too large: {} bytes > {MAX_INJECT_BODY_BYTES}",
body.len()
),
);
}
let registry = load_registry_offloaded(ctx.home.registry_json()).await;
let to_entry = match registry.find(&to) {
Some(e) => e.clone(),
None => {
return Response::err(
req.id,
ErrorCode::AgentNotFound,
format!("agent '{to}' not found"),
)
}
};
if !switchboard_identity_matches(&to_entry, recipient_identity) {
return Response::ok(
req.id,
json!({"delivered": false, "reason": "recipient-identity-changed"}),
);
}
let to_sock = ctx.home.worker_sock(&to_entry.short_id);
if to_entry.harness_name() != "claude" || !is_live_stream_thread(&to_sock).await {
return Response::ok(
req.id,
json!({"delivered": false, "reason": "not-a-live-stream-thread"}),
);
}
let drive_deadline = Duration::from_millis(timeout_ms);
let outer = drive_deadline + Duration::from_secs(SWITCHBOARD_DRIVE_GRACE_S);
let drive_result =
match tokio::time::timeout(outer, drive_stream_turn(&to_sock, &body, drive_deadline)).await
{
Ok(inner) => inner,
Err(_) => Err("drive hung past the turn budget (timed out)".to_string()),
};
let outcome = match drive_result {
Ok(o) => o,
Err(reason) => {
match stamp_orphaned(&ctx.home, to.clone(), recipient_identity.clone()).await {
Ok(true) => {}
Ok(false) => {
let _ = ctx.emitter.emit(
"agent_deliver_status_write_failed",
&json!({
"name": to,
"from_name": from,
"provider": "claude",
"transport": "switchboard",
"reason": "recipient-identity-changed",
}),
);
}
Err(error) => {
let _ = ctx.emitter.emit(
"agent_deliver_status_write_failed",
&json!({
"name": to,
"from_name": from,
"provider": "claude",
"transport": "switchboard",
"reason": "registry-write-failed",
"error": error.to_string(),
}),
);
}
}
let _ = ctx.emitter.emit(
"agent_deliver_demoted",
&json!({
"name": to,
"from_name": from,
"provider": "claude",
"transport": "switchboard",
"reason": reason,
}),
);
return Response::ok(req.id, json!({"delivered": false, "reason": reason}));
}
};
let mut mirrored = false;
if mirror && from != to {
let fresh = load_registry_offloaded(ctx.home.registry_json()).await;
if let Some(from_entry) = fresh.find(&from).filter(|entry| {
from_identity.is_some_and(|identity| switchboard_identity_matches(entry, identity))
}) {
let from_sock = ctx.home.worker_sock(&from_entry.short_id);
if from_entry.harness_name() == "claude" && is_live_stream_thread(&from_sock).await {
match mirror_into(&from_sock, &outcome.reply).await {
Ok(()) => mirrored = true,
Err(e) => {
let _ = ctx.emitter.emit(
"agent_deliver_demoted",
&json!({
"name": from,
"from_name": to,
"provider": "claude",
"transport": "switchboard-mirror",
"reason": e,
}),
);
}
}
}
}
}
let _ = ctx.emitter.emit(
"agent_deliver_injected",
&json!({
"name": to,
"from_name": from,
"provider": "claude",
"transport": "switchboard",
"mirrored": mirrored,
"is_error": outcome.is_error,
}),
);
Response::ok(
req.id,
json!({
"delivered": true,
"identity_verified": true,
"transport": "switchboard",
"reply": outcome.reply,
"is_error": outcome.is_error,
"mirrored": mirrored,
"receipt": outcome.saw_receipt,
}),
)
}
async fn read_worker_snapshot(sock: &std::path::Path) -> Option<String> {
let mut conn = UnixStream::connect(sock).await.ok()?;
write_request(&mut conn, &Request::new(2, "worker.snapshot", json!({})))
.await
.ok()?;
let resp = crate::protocol::read_response(&mut conn).await.ok()?;
resp.result()
.and_then(|r| r.get("text").and_then(|t| t.as_str()).map(String::from))
}
fn reap_zombies() {
loop {
let mut status: libc::c_int = 0;
let pid = unsafe { libc::waitpid(-1, &mut status, libc::WNOHANG) };
if pid <= 0 {
break;
}
}
}
fn rendered_status_from_truth(truth: Option<&str>) -> &'static str {
match truth {
Some("working" | "watching" | "your-move") => "live",
Some("done" | "stalled") => "orphaned",
_ => "unknown",
}
}
fn registry_truth_handle(entry: &RegistryEntry) -> String {
if let Some(session_id) = entry.harness_session_id.as_deref() {
return session_id.to_string();
}
if !entry.short_id.is_empty() {
entry.short_id.clone()
} else {
entry.name.clone()
}
}
fn handle_list(ctx: &Ctx, req: &Request) -> Response {
handle_list_with_truth(ctx, req, crate::claude_ask::family1_truth_state)
}
fn handle_list_with_truth<F>(ctx: &Ctx, req: &Request, truth_fn: F) -> Response
where
F: Fn(&str) -> Option<String>,
{
let all = req
.params
.get("all")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let filter_cwd = req
.params
.get("cwd")
.and_then(|v| v.as_str())
.map(String::from);
let filter_provider = req
.params
.get("provider")
.and_then(|v| v.as_str())
.map(String::from);
let filter_status = req
.params
.get("status")
.and_then(|v| v.as_str())
.map(String::from);
let cwd_project = req
.params
.get("project_root")
.and_then(|v| v.as_str())
.map(String::from);
if let Some(ref st) = filter_status {
if st != "live" && st != "orphaned" && st != "unknown" {
return Response::err(
req.id,
ErrorCode::InvalidStatus,
format!("invalid --status '{st}' (expected: live | orphaned | unknown)"),
);
}
}
let norm_path = |p: &str| -> String {
std::fs::canonicalize(p)
.ok()
.and_then(|pb| pb.to_str().map(String::from))
.unwrap_or_else(|| p.to_string())
};
let filter_cwd_norm = filter_cwd.as_deref().map(&norm_path);
let registry = state::load_registry(&ctx.home.registry_json()).unwrap_or_default();
let classified: Vec<_> = registry
.entries
.iter()
.filter(|e| {
if !all {
if let Some(ref p) = cwd_project {
if &e.project_root != p {
return false;
}
}
}
if let Some(ref cwd) = filter_cwd_norm {
if &norm_path(&e.cwd) != cwd {
return false;
}
}
if let Some(ref prov) = filter_provider {
if e.harness_name() != prov.as_str() {
return false;
}
}
true
})
.map(|e| {
let truth_handle = registry_truth_handle(e);
let truth = truth_fn(&truth_handle);
(e, rendered_status_from_truth(truth.as_deref()))
})
.collect();
let entries: Vec<Value> = classified
.into_iter()
.filter(|(_e, rendered_status)| {
if let Some(ref st) = filter_status {
if rendered_status != &st.as_str() {
return false;
}
}
true
})
.map(|(e, rendered_status)| {
let resume_id: Option<String> = match e.harness_name() {
"claude" => e
.transport_short()
.map(str::to_string)
.or_else(|| e.session_id.clone()),
"codex" => e.codex_session_id.clone().or_else(|| e.session_id.clone()),
"gemini" => e.gemini_session_id.clone().or_else(|| e.session_id.clone()),
"opencode" => e
.harness_session_id
.clone()
.filter(|s| !s.is_empty())
.or_else(|| e.session_id.clone()),
_ => e.session_id.clone(),
};
let session_id: Value = resume_id.map(Value::String).unwrap_or(Value::Null);
let short_id: Value = e
.transport_short()
.map(|s| Value::String(s.to_string()))
.unwrap_or(Value::Null);
let log_path: Value = e
.log_path
.as_deref()
.map(|s| Value::String(s.to_string()))
.unwrap_or(Value::Null);
let crown: Value = match e.crown_level {
Some(level) => Value::String(format!(
"L{level} {}",
e.crown_scope
.as_deref()
.filter(|s| !s.is_empty())
.unwrap_or("?")
)),
None => Value::Null,
};
json!({
"name": e.name,
"harness": e.harness_name(),
"provider": e.harness_name(),
"harness_session_id": e.harness_session_id,
"short_id": short_id,
"session_id": session_id,
"cwd": e.cwd,
"created_at": e.created_at,
"last_message_at": e.last_message_at,
"status": rendered_status,
"live_status": null,
"pid": e.pid,
"last_reconciled_at": e.last_reconciled_at,
"log_path": log_path,
"mux": e.mux,
"crown": crown,
"crown_level": e.crown_level,
"crown_scope": e.crown_scope,
"crown_grantor": e.crown_grantor,
"project_root": e.project_root,
})
})
.collect();
let filters_applied = json!({
"cwd": filter_cwd_norm,
"provider": filter_provider,
"status": filter_status,
});
Response::ok(
req.id,
json!({"agents": entries, "filters_applied": filters_applied}),
)
}
async fn handle_status(ctx: &Ctx, req: &Request) -> Response {
let reg_path = ctx.home.registry_json();
let registry = tokio::task::spawn_blocking(move || state::load_registry(®_path))
.await
.ok()
.and_then(|r| r.ok())
.unwrap_or_default();
let mut by_status: Map<String, Value> = Map::new();
let mut restarting: u64 = 0;
let mut channels_registered: u64 = 0;
for e in ®istry.entries {
let key = format!("{:?}", e.status).to_lowercase();
let n = by_status.get(&key).and_then(|v| v.as_u64()).unwrap_or(0) + 1;
by_status.insert(key, Value::Number(n.into()));
if e.status == AgentStatus::Restarting {
restarting += 1;
}
if e.mcp_channel_id.is_some() {
channels_registered += 1;
}
}
Response::ok(
req.id,
json!({
"schema_version": 1,
"daemon": {
"state": DaemonState::Serving.as_str(),
"pid": std::process::id(),
"uptime_secs": ctx.started_at.elapsed().as_secs(),
"version": env!("CARGO_PKG_VERSION"),
"exe_path": ctx
.exe_fingerprint
.as_ref()
.map(|f| f.path.to_string_lossy().into_owned()),
"exe_mtime": ctx.exe_fingerprint.as_ref().map(|f| f.mtime_nanos),
"exe_size": ctx.exe_fingerprint.as_ref().map(|f| f.size),
"pid_start_time": ctx.pid_start_time,
},
"agents": {
"total": registry.entries.len(),
"by_status": by_status,
},
"restarts": {
"queue_depth": restarting,
"consecutive_failures_max_seen": 0,
},
"channels": { "registered": channels_registered },
}),
)
}
async fn entry_for_lifecycle(
registry: &state::Registry,
token: &str,
registry_path: &std::path::Path,
) -> Result<Option<RegistryEntry>, String> {
let Value::Array(rows) = serde_json::to_value(®istry.entries)
.map_err(|exc| format!("could not inspect registry identities: {exc}"))?
else {
return Err("could not inspect registry identities".to_string());
};
let worker_token = token.to_string();
let path = registry_path.to_path_buf();
let resolved = tokio::task::spawn_blocking(move || {
crate::client_verbs::resolve_entry_with_heal(&rows, &worker_token, &path)
})
.await
.map_err(|exc| format!("identity resolution task failed: {exc}"))?;
match resolved {
Ok(entry) => {
let mut entry: RegistryEntry = serde_json::from_value(entry)
.map_err(|exc| format!("resolved identity row is unreadable: {exc}"))?;
entry.backfill_harness_aliases();
if let Some(legacy) = entry.backfill_short_id() {
return Err(format!(
"resolved identity row {:?} has conflicting transport ids (legacy={legacy:?})",
entry.name
));
}
Ok(Some(entry))
}
Err(crate::client_verbs::ResolveError::NotFound(_)) => Ok(None),
Err(err) => Err(err.message()),
}
}
async fn handle_stop(ctx: &Ctx, req: &Request) -> Response {
let requested_name = match req.params.get("name").and_then(|v| v.as_str()) {
Some(n) => n.to_string(),
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `name`"),
};
let registry = load_registry_offloaded(ctx.home.registry_json()).await;
let entry =
match entry_for_lifecycle(®istry, &requested_name, &ctx.home.registry_json()).await {
Ok(Some(entry)) => entry,
Ok(None) => {
return Response::err(
req.id,
ErrorCode::AgentNotFound,
format!("agent {requested_name} not found"),
)
}
Err(message) => return Response::err(req.id, ErrorCode::InvalidParams, message),
};
let name = entry.name.clone();
if entry.status == AgentStatus::Exited {
return Response::ok(
req.id,
json!({"already_exited": true, "short_id": entry.short_id}),
);
}
if entry.harness_name() == "claude" {
return stop_claude(ctx, req, &name, &entry).await;
}
if entry.short_id.is_empty() {
let _ = ctx.emitter.emit(
"agent_stopped",
&json!({"name": name, "provider": entry.harness_name(), "claude_exit": Value::Null}),
);
return Response::ok(
req.id,
json!({"stopped": true, "provider": entry.harness_name(), "no_op": true}),
);
}
if !stop_worker_confirmed(ctx, &entry).await {
return Response::err(
req.id,
ErrorCode::Internal,
format!("agent {name}: worker did not confirm shutdown; it may still be running"),
);
}
let stop_name = name.clone();
if let Err(e) = update_registry_offloaded(ctx.home.registry_json(), move |r| {
if let Some(e) = r.find_mut(&stop_name) {
e.status = AgentStatus::Exited;
}
})
.await
{
let _ = ctx.emitter.emit(
"agent_stop_error",
&json!({"name": name, "error": e.to_string()}),
);
return Response::err(
req.id,
ErrorCode::Internal,
format!("agent {name}: worker stopped but registry write failed: {e}"),
);
}
let _ = ctx.emitter.emit("agent_stopped", &json!({"name": name}));
Response::ok(req.id, json!({"stopped": true, "short_id": entry.short_id}))
}
async fn best_effort_worker_shutdown(sock: &std::path::Path) {
if let Ok(mut conn) = UnixStream::connect(sock).await {
let _ = write_request(&mut conn, &Request::new(1, "worker.shutdown", json!({}))).await;
let _ = crate::protocol::read_response(&mut conn).await;
}
}
async fn stop_worker_confirmed(ctx: &Ctx, entry: &RegistryEntry) -> bool {
let sock = ctx.home.worker_sock(&entry.short_id);
if let Ok(mut conn) = UnixStream::connect(&sock).await {
let _ = write_request(&mut conn, &Request::new(1, "worker.shutdown", json!({}))).await;
let _ = crate::protocol::read_response(&mut conn).await;
}
let mut down = worker_down_within(&sock, Duration::from_secs(5)).await;
if !down {
if let Some(pid) = entry.pid {
if pid_is_ours(pid, entry.pid_start_time) {
unsafe {
libc::kill(pid as libc::pid_t, libc::SIGTERM);
}
down = worker_down_within(&sock, Duration::from_secs(5)).await;
if !down && pid_is_ours(pid, entry.pid_start_time) {
unsafe {
libc::kill(pid as libc::pid_t, libc::SIGKILL);
}
down = worker_down_within(&sock, Duration::from_secs(2)).await;
}
}
}
}
if down {
let _ = std::fs::remove_file(&sock);
}
down
}
async fn worker_socket_reachable(sock: &std::path::Path) -> bool {
UnixStream::connect(sock).await.is_ok()
}
async fn worker_down_within(sock: &std::path::Path, budget: Duration) -> bool {
let start = Instant::now();
loop {
if !worker_socket_reachable(sock).await {
return true;
}
if start.elapsed() >= budget {
return false;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
}
fn pid_confirmed_dead(pid: u32) -> bool {
if pid <= 1 || pid > i32::MAX as u32 {
return true;
}
if unsafe { libc::kill(pid as libc::pid_t, 0) } == 0 {
return false; }
std::io::Error::last_os_error().raw_os_error() == Some(libc::ESRCH)
}
fn pid_recycled(pid: u32, recorded_start: Option<u64>) -> bool {
let Some(recorded) = recorded_start else {
return false; };
if pid <= 1 || pid > i32::MAX as u32 {
return false;
}
if unsafe { libc::kill(pid as libc::pid_t, 0) } != 0 {
return false; }
match process_start_time(pid) {
Some(now) => now != recorded,
None => false, }
}
async fn pid_gone_within(pid: u32, recorded_start: Option<u64>, budget: Duration) -> bool {
let start = Instant::now();
loop {
if pid_confirmed_dead(pid) || pid_recycled(pid, recorded_start) {
return true;
}
if start.elapsed() >= budget {
return false;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
}
async fn stop_claude_pid_confirmed(entry: &RegistryEntry) -> bool {
let Some(pid) = entry.pid else {
return false;
};
if entry.pid_start_time.is_none() {
return false;
}
if !pid_is_ours(pid, entry.pid_start_time) {
return false;
}
unsafe {
libc::kill(pid as libc::pid_t, libc::SIGTERM);
}
if pid_gone_within(pid, entry.pid_start_time, Duration::from_secs(5)).await {
return true;
}
if pid_is_ours(pid, entry.pid_start_time) {
unsafe {
libc::kill(pid as libc::pid_t, libc::SIGKILL);
}
}
pid_gone_within(pid, entry.pid_start_time, Duration::from_secs(2)).await
}
async fn stop_claude(ctx: &Ctx, req: &Request, name: &str, entry: &RegistryEntry) -> Response {
let short = match entry
.transport_short()
.or(entry.session_id.as_deref())
.filter(|s| !s.is_empty())
{
Some(s) => s.to_string(),
None => {
if stop_claude_pid_confirmed(entry).await {
let claude_name = name.to_string();
if let Err(e) = update_registry_offloaded(ctx.home.registry_json(), move |r| {
if let Some(e) = r.find_mut(&claude_name) {
e.status = AgentStatus::Exited;
}
})
.await
{
return Response::err(
req.id,
ErrorCode::Internal,
format!("claude {name} stopped but registry write failed: {e}"),
);
}
let _ = ctx.emitter.emit(
"agent_stopped",
&json!({"name": name, "backend": "claude", "stopped_by": "pid"}),
);
return Response::ok(
req.id,
json!({"stopped": true, "backend": "claude", "pid": entry.pid}),
);
}
return Response::err(
req.id,
ErrorCode::InvalidStatus,
format!(
"agent {name} is claude but has no short id and no live process \
to stop; note that `rm` clears the row but does not stop a session"
),
);
}
};
match tokio::process::Command::new("claude")
.arg("stop")
.arg(&short)
.output()
.await
{
Ok(o) if o.status.success() => {
let claude_name = name.to_string();
if let Err(e) = update_registry_offloaded(ctx.home.registry_json(), move |r| {
if let Some(e) = r.find_mut(&claude_name) {
e.status = AgentStatus::Exited;
}
})
.await
{
return Response::err(
req.id,
ErrorCode::Internal,
format!("claude {name} stopped but registry write failed: {e}"),
);
}
let _ = ctx
.emitter
.emit("agent_stopped", &json!({"name": name, "backend": "claude"}));
Response::ok(
req.id,
json!({"stopped": true, "backend": "claude", "short_id": short}),
)
}
Ok(o) => Response::err(
req.id,
ErrorCode::Internal,
format!(
"claude stop {short} failed: {}",
String::from_utf8_lossy(&o.stderr).trim()
),
),
Err(e) => Response::err(
req.id,
ErrorCode::Internal,
format!("could not exec `claude stop`: {e}"),
),
}
}
async fn handle_rm(ctx: &Ctx, req: &Request) -> Response {
let requested_name = match req.params.get("name").and_then(|v| v.as_str()) {
Some(n) => n.to_string(),
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `name`"),
};
let force = req
.params
.get("force")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let registry = load_registry_offloaded(ctx.home.registry_json()).await;
let entry =
match entry_for_lifecycle(®istry, &requested_name, &ctx.home.registry_json()).await {
Ok(Some(entry)) => entry,
Ok(None) => {
return Response::err(
req.id,
ErrorCode::AgentNotFound,
format!("agent {requested_name} not found"),
)
}
Err(message) => return Response::err(req.id, ErrorCode::InvalidParams, message),
};
let name = entry.name.clone();
if entry.status == AgentStatus::Live && !force {
return Response::err(
req.id,
ErrorCode::Busy,
format!("agent {name} is still live; use `stop` first or pass --force"),
);
}
if entry.status == AgentStatus::Live && force && !stop_worker_confirmed(ctx, &entry).await {
return Response::err(
req.id,
ErrorCode::Internal,
format!("agent {name}: could not stop the worker before force-remove; refusing to orphan a live PTY"),
);
}
let was_orphaned = entry.status == AgentStatus::Orphaned;
let rm_name = name.clone();
if let Err(e) = update_registry_offloaded(ctx.home.registry_json(), move |r| {
r.entries.retain(|e| e.name != rm_name);
})
.await
{
return Response::err(
req.id,
ErrorCode::Internal,
format!("agent {name}: removal did not persist: {e}"),
);
}
let _ = ctx.emitter.emit(
"agent_removed",
&json!({"name": name, "was_orphaned": was_orphaned}),
);
Response::ok(
req.id,
json!({"removed": true, "was_orphaned": was_orphaned}),
)
}
const RECONCILE_PROBE_TIMEOUT: Duration = Duration::from_millis(250);
const RECONCILE_SWEEP_BUDGET: Duration = Duration::from_secs(5);
struct ReconcileChange {
name: String,
new_status: Option<AgentStatus>,
}
#[derive(Default, PartialEq, Debug)]
struct ReconcileOutcome {
updated: Vec<String>,
orphans: Vec<String>,
recovered: Vec<String>,
inconsistent: Vec<(String, String)>,
deferred: usize,
}
fn plan_reconcile<P, D, L>(
entries: &[RegistryEntry],
mut probe: P,
mut budget_exhausted: D,
mut pid_live: L,
) -> (Vec<ReconcileChange>, ReconcileOutcome)
where
P: FnMut(&RegistryEntry) -> Result<bool, crate::provider::ReachabilityProbeError>,
D: FnMut() -> bool,
L: FnMut(&RegistryEntry) -> bool,
{
let mut changes = Vec::new();
let mut out = ReconcileOutcome::default();
for (i, entry) in entries.iter().enumerate() {
if budget_exhausted() {
out.deferred = entries.len() - i;
break;
}
if entry.is_one_shot_ask() {
let new_status = if is_non_terminal(entry.status) {
out.updated.push(entry.name.clone());
Some(AgentStatus::Exited)
} else {
None
};
changes.push(ReconcileChange {
name: entry.name.clone(),
new_status,
});
continue;
}
let new_status = match probe(entry) {
Ok(true) => {
if entry.status == AgentStatus::Orphaned && pid_live(entry) {
out.recovered.push(entry.name.clone());
out.updated.push(entry.name.clone());
Some(AgentStatus::Live)
} else {
None
}
}
Ok(false) if entry.is_interactive() => {
if pid_live(entry) {
None
} else {
out.updated.push(entry.name.clone());
Some(AgentStatus::Exited)
}
}
Ok(false) if entry.mux.is_some() => {
if entry.pid.is_some() && pid_live(entry) {
None
} else if entry.pid.is_some() {
out.updated.push(entry.name.clone());
Some(AgentStatus::Exited)
} else {
let live_ish = matches!(
entry.status,
AgentStatus::Live
| AgentStatus::Ready
| AgentStatus::Idle
| AgentStatus::Busy
| AgentStatus::Spawning
);
if live_ish {
out.orphans.push(entry.name.clone());
out.updated.push(entry.name.clone());
Some(AgentStatus::Orphaned)
} else {
None
}
}
}
Ok(false) => {
let live_ish = matches!(
entry.status,
AgentStatus::Live
| AgentStatus::Ready
| AgentStatus::Idle
| AgentStatus::Busy
| AgentStatus::Spawning
);
if live_ish {
out.orphans.push(entry.name.clone());
out.updated.push(entry.name.clone());
Some(AgentStatus::Orphaned)
} else {
None
}
}
Err(e) => {
out.inconsistent
.push((entry.name.clone(), e.reason.clone()));
None
}
};
changes.push(ReconcileChange {
name: entry.name.clone(),
new_status,
});
}
(changes, out)
}
fn apply_reconcile_change(e: &mut RegistryEntry, new_status: Option<AgentStatus>, now: &str) {
e.last_reconciled_at = Some(now.to_string());
if let Some(s) = new_status {
e.status = s;
if matches!(s, AgentStatus::Exited) {
e.pid = None;
e.pid_start_time = None;
e.inside_leg = None;
e.screen_state = None;
}
}
}
fn emit_inside_leg_completion(emitter: &EventEmitter, e: &RegistryEntry) {
if let Some(rep) = &e.inside_leg {
let _ = emitter.emit(
"inside_leg_completed",
&json!({
"name": e.name,
"session_id": e.session_id,
"final_state": inside_leg_state_str(rep.state),
"seq": rep.seq,
}),
);
}
}
fn inside_leg_state_str(state: state::InsideLegState) -> &'static str {
match state {
state::InsideLegState::Working => "working",
state::InsideLegState::Blocked => "blocked",
state::InsideLegState::Done => "done",
}
}
fn to_agent_entry(e: &RegistryEntry) -> crate::provider::AgentEntry {
let session_id = match e.harness_name() {
"codex" => e.codex_session_id.clone().or_else(|| e.session_id.clone()),
"gemini" => e.gemini_session_id.clone().or_else(|| e.session_id.clone()),
"claude" => e
.transport_short()
.map(str::to_string)
.or_else(|| e.session_id.clone()),
"opencode" => e
.harness_session_id
.clone()
.or_else(|| e.session_id.clone()),
_ => e.session_id.clone(),
};
crate::provider::AgentEntry {
name: e.name.clone(),
provider: e.harness_name().to_string(),
session_id,
cwd: PathBuf::from(&e.cwd),
}
}
struct ReconcileSweepResult {
registry: crate::state::Registry,
entries: Vec<RegistryEntry>,
outcome: ReconcileOutcome,
}
fn run_reconcile_sweep(
home: &AgentsHome,
emitter: &EventEmitter,
) -> Result<ReconcileSweepResult, String> {
use crate::provider::ReachabilityProbeError;
let registry = state::load_registry(&home.registry_json()).unwrap_or_default();
let mut entries = registry.entries.clone();
entries.sort_by(|a, b| a.last_reconciled_at.cmp(&b.last_reconciled_at));
let start = Instant::now();
let probe = |e: &RegistryEntry| -> Result<bool, ReachabilityProbeError> {
if std::os::unix::net::UnixStream::connect(home.worker_sock(&e.short_id)).is_ok() {
return Ok(true);
}
match crate::provider::for_name(e.harness_name()) {
Some(p) => p.reachability(&to_agent_entry(e), RECONCILE_PROBE_TIMEOUT),
None => Err(ReachabilityProbeError::new(
e.harness_name(),
"unknown provider; cannot probe reachability",
)),
}
};
let pid_live = |e: &RegistryEntry| -> bool {
e.pid.map_or(true, |pid| pid_is_ours(pid, e.pid_start_time))
};
let (changes, outcome) = plan_reconcile(
&entries,
probe,
|| start.elapsed() >= RECONCILE_SWEEP_BUDGET,
pid_live,
);
for ch in &changes {
if matches!(ch.new_status, Some(AgentStatus::Exited)) {
if let Some(e) = registry.entries.iter().find(|e| e.name == ch.name) {
emit_inside_leg_completion(emitter, e);
}
}
}
let now = now_rfc3339_like();
if let Err(err) = state::update_registry(&home.registry_json(), |r| {
for ch in &changes {
if let Some(e) = r.find_mut(&ch.name) {
apply_reconcile_change(e, ch.new_status, &now);
}
}
}) {
let _ = emitter.emit("reconcile_error", &json!({"error": err.to_string()}));
return Err(format!(
"reconcile computed {} change(s) but the registry write failed: {err}",
changes.len()
));
}
for (name, reason) in &outcome.inconsistent {
let _ = emitter.emit(
"agent_inconsistent",
&json!({"name": name, "reason": reason}),
);
}
if outcome.deferred > 0 {
let _ = emitter.emit(
"reconcile_deferred",
&json!({"remaining_count": outcome.deferred}),
);
}
let _ = emitter.emit(
"reconcile_done",
&json!({
"updated": outcome.updated.len(),
"orphans": outcome.orphans.len(),
"recovered": outcome.recovered.len(),
}),
);
Ok(ReconcileSweepResult {
registry,
entries,
outcome,
})
}
fn handle_reconcile(ctx: &Ctx, req: &Request) -> Response {
let ReconcileSweepResult {
registry,
entries,
outcome,
} = match run_reconcile_sweep(&ctx.home, &ctx.emitter) {
Ok(r) => r,
Err(msg) => return Response::err(req.id, ErrorCode::Internal, msg),
};
let scanned = entries.len();
let probed = entries.len() - outcome.deferred;
let skipped_py: Vec<Value> = entries
.iter()
.skip(probed)
.map(|e| json!({"name": e.name, "provider": e.harness_name()}))
.collect();
let orphaned_py: Vec<Value> = outcome
.orphans
.iter()
.map(|n| {
let prov = registry
.entries
.iter()
.find(|e| &e.name == n)
.map(|e| e.harness_name())
.unwrap_or("unknown");
json!({"name": n, "provider": prov})
})
.collect();
let recovered_py: Vec<Value> = outcome
.recovered
.iter()
.map(|n| {
let prov = registry
.entries
.iter()
.find(|e| &e.name == n)
.map(|e| e.harness_name())
.unwrap_or("unknown");
json!({"name": n, "provider": prov})
})
.collect();
let errors_py: Vec<Value> = outcome
.inconsistent
.iter()
.map(|(n, reason)| json!({"name": n, "reason": reason}))
.collect();
Response::ok(
req.id,
json!({
"scanned": scanned,
"orphaned": orphaned_py,
"recovered": recovered_py,
"skipped": skipped_py,
"errors": errors_py,
"updated": outcome.updated,
"orphans": outcome.orphans,
"inconsistent": outcome.inconsistent.iter().map(|(n, _)| n.clone()).collect::<Vec<_>>(),
"deferred": outcome.deferred,
}),
)
}
enum BufferOutcome {
Buffered,
StaleSeq { last: u64 },
Full,
}
fn buffer_pending_report(
map: &mut std::collections::HashMap<String, state::InsideLegReport>,
session_id: &str,
report: state::InsideLegReport,
) -> BufferOutcome {
if let Some(prev) = map.get(session_id) {
if report.seq <= prev.seq {
return BufferOutcome::StaleSeq { last: prev.seq };
}
map.insert(session_id.to_string(), report);
return BufferOutcome::Buffered;
}
if map.len() >= PENDING_INSIDE_LEG_CAP {
return BufferOutcome::Full;
}
map.insert(session_id.to_string(), report);
BufferOutcome::Buffered
}
fn flush_buffered_inside_leg(ctx: &Ctx, session_uuid: &str, name: &str) {
let rep = match ctx.pending_inside_leg.lock() {
Ok(mut buf) => buf.remove(session_uuid),
Err(_) => None,
};
let Some(rep) = rep else {
return;
};
let (seq, state_str) = (rep.seq, inside_leg_state_str(rep.state));
let (rep_state, rep_reason) = (rep.state, rep.reason.clone());
let mut notify: Option<(String, String, bool)> = None;
let _ = state::update_registry(&ctx.home.registry_json(), |r| {
if let Some(e) = r
.entries
.iter_mut()
.find(|e| entry_holds_session(e, session_uuid))
{
let newer = e.inside_leg.as_ref().is_none_or(|cur| rep.seq > cur.seq);
if newer {
let prev_state = e.inside_leg.as_ref().map(|r| r.state);
let body = rep_reason.clone().unwrap_or_else(|| state_str.to_string());
if state::enters(prev_state, rep_state, state::InsideLegState::Blocked) {
notify = Some((name.to_string(), body, false));
} else if state::enters(prev_state, rep_state, state::InsideLegState::Done) {
notify = Some((name.to_string(), body, true));
}
e.inside_leg = Some(rep);
e.screen_state = None;
}
}
});
if let Some((title, body, is_done)) = notify {
let want = if is_done {
ctx.opts.notify_on_done
} else {
ctx.opts.notify_on_blocked
};
if want {
notify_transition(title, body);
}
}
let _ = ctx.emitter.emit(
"inside_leg_buffer_flushed",
&json!({"name": name, "session_id": session_uuid, "state": state_str, "seq": seq}),
);
}
pub(crate) fn notify_transition(title: String, body: String) {
let fno = std::env::var_os("FNO_BIN").unwrap_or_else(|| std::ffi::OsString::from("fno"));
std::thread::spawn(move || {
match std::process::Command::new(&fno)
.args(["notify", &title, &body])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
{
Ok(mut child) => {
let _ = child.wait();
}
Err(e) => eprintln!(
"fno-agents-daemon: badge notify skipped ({} notify): {e}",
fno.to_string_lossy()
),
}
});
}
enum UuidBackfill {
None,
One(usize),
Ambiguous,
}
fn find_uuid_backfill_row(entries: &[RegistryEntry], full_uuid: &str) -> UuidBackfill {
let mut found = None;
for (i, e) in entries.iter().enumerate() {
if e.harness_name() != "claude" || e.claude_session_uuid.is_some() {
continue;
}
let Some(short) = e.transport_short() else {
continue;
};
if short.is_empty()
|| !full_uuid
.strip_prefix(short)
.is_some_and(|rest| rest.starts_with('-'))
{
continue;
}
if found.is_some() {
return UuidBackfill::Ambiguous;
}
found = Some(i);
}
found.map_or(UuidBackfill::None, UuidBackfill::One)
}
fn handle_report(ctx: &Ctx, req: &Request) -> Response {
let session_id = match req.params.get("session_id").and_then(|v| v.as_str()) {
Some(s) if !s.is_empty() => s.to_string(),
_ => return Response::err(req.id, ErrorCode::InvalidParams, "missing `session_id`"),
};
let seq = match req.params.get("seq").and_then(|v| v.as_u64()) {
Some(n) => n,
None => {
return Response::err(
req.id,
ErrorCode::InvalidParams,
"missing or non-integer `seq`",
)
}
};
let state_label = match req.params.get("state").and_then(|v| v.as_str()) {
Some(s @ ("working" | "blocked" | "done")) => s.to_string(),
_ => {
return Response::err(
req.id,
ErrorCode::InvalidParams,
"`state` must be working|blocked|done",
)
}
};
let state = match state_label.as_str() {
"working" => state::InsideLegState::Working,
"blocked" => state::InsideLegState::Blocked,
_ => state::InsideLegState::Done,
};
let reason = req
.params
.get("reason")
.and_then(|v| v.as_str())
.map(String::from);
let ttl_ms = req.params.get("ttl_ms").and_then(|v| v.as_u64());
let report = state::InsideLegReport {
state,
seq,
reason,
received_at: now_rfc3339_like(),
ttl_ms,
};
let report_for_store = report.clone();
enum Outcome {
Stored,
StaleSeq { last: u64 },
Unknown,
}
let mut outcome = Outcome::Unknown;
let mut notify: Option<(String, String, bool)> = None;
if let Err(e) = state::update_registry(&ctx.home.registry_json(), |r| {
let idx = match r
.entries
.iter()
.position(|e| entry_holds_session(e, &session_id))
{
Some(i) => Some(i),
None => match find_uuid_backfill_row(&r.entries, &session_id) {
UuidBackfill::One(i) => {
r.entries[i].claude_session_uuid = Some(session_id.clone());
Some(i)
}
UuidBackfill::None | UuidBackfill::Ambiguous => None,
},
};
let Some(idx) = idx else {
outcome = Outcome::Unknown;
return;
};
let entry = &mut r.entries[idx];
if let Some(prev) = &entry.inside_leg {
if seq <= prev.seq {
outcome = Outcome::StaleSeq { last: prev.seq };
return;
}
}
let prev_state = entry.inside_leg.as_ref().map(|r| r.state);
if state::enters(prev_state, state, state::InsideLegState::Blocked) {
let body = report_for_store
.reason
.clone()
.unwrap_or_else(|| state_label.clone());
notify = Some((entry.name.clone(), body, false));
} else if state::enters(prev_state, state, state::InsideLegState::Done) {
let body = report_for_store
.reason
.clone()
.unwrap_or_else(|| state_label.clone());
notify = Some((entry.name.clone(), body, true));
}
entry.inside_leg = Some(report_for_store);
entry.screen_state = None;
outcome = Outcome::Stored;
}) {
return Response::err(
req.id,
ErrorCode::Internal,
format!("registry write failed during inside-leg report: {e}"),
);
}
match outcome {
Outcome::Stored => {
let _ = ctx.emitter.emit(
"inside_leg_report",
&json!({"session_id": session_id, "seq": seq, "state": state_label}),
);
if let Some((title, body, is_done)) = notify {
let want = if is_done {
ctx.opts.notify_on_done
} else {
ctx.opts.notify_on_blocked
};
if want {
notify_transition(title, body);
}
}
Response::ok(req.id, json!({"stored": true, "seq": seq}))
}
Outcome::StaleSeq { last } => {
let _ = ctx.emitter.emit(
"inside_leg_report_dropped",
&json!({"session_id": session_id, "seq": seq, "last_seq": last, "reason": "stale_seq"}),
);
Response::ok(
req.id,
json!({"stored": false, "dropped": "stale_seq", "last_seq": last}),
)
}
Outcome::Unknown => {
let buffered = ctx
.pending_inside_leg
.lock()
.map(|mut buf| buffer_pending_report(&mut buf, &session_id, report))
.ok();
match buffered {
Some(BufferOutcome::Buffered) => {
let _ = ctx.emitter.emit(
"inside_leg_report_buffered",
&json!({"session_id": session_id, "seq": seq, "state": state_label}),
);
Response::ok(
req.id,
json!({"stored": false, "buffered": true, "seq": seq}),
)
}
Some(BufferOutcome::StaleSeq { last }) => {
let _ = ctx.emitter.emit(
"inside_leg_report_dropped",
&json!({"session_id": session_id, "seq": seq, "last_seq": last, "reason": "stale_seq"}),
);
Response::ok(
req.id,
json!({"stored": false, "dropped": "stale_seq", "last_seq": last}),
)
}
Some(BufferOutcome::Full) => {
let _ = ctx.emitter.emit(
"inside_leg_report_dropped",
&json!({"session_id": session_id, "seq": seq, "reason": "buffer_full"}),
);
Response::ok(req.id, json!({"stored": false, "dropped": "buffer_full"}))
}
None => {
let _ = ctx.emitter.emit(
"inside_leg_report_dropped",
&json!({"session_id": session_id, "seq": seq, "reason": "unknown_session"}),
);
Response::ok(
req.id,
json!({"stored": false, "dropped": "unknown_session"}),
)
}
}
}
}
}
async fn dispatch_channel(ctx: &Arc<Ctx>, req: &Request) -> Response {
match Namespace::verb(&req.method) {
Some("register_channel") => run_blocking(ctx, req, handle_register_channel).await,
Some("unregister_channel") => run_blocking(ctx, req, handle_unregister_channel).await,
Some("push_to_channel") => run_blocking(ctx, req, handle_push_to_channel).await,
_ => Response::err(
req.id,
ErrorCode::UnknownMethod,
format!("unknown channel verb in `{}`", req.method),
),
}
}
fn handle_register_channel(ctx: &Ctx, req: &Request) -> Response {
let cc_session_id = match req.params.get("cc_session_id").and_then(|v| v.as_str()) {
Some(s) => s.to_string(),
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `cc_session_id`"),
};
let name = req
.params
.get("name")
.and_then(|v| v.as_str())
.map(String::from);
let channel_id = uuid_v4();
let mut matched = false;
if let Err(e) = state::update_registry(&ctx.home.registry_json(), |r| {
let target = match &name {
Some(n) => r.find_mut(n),
None => r
.entries
.iter_mut()
.find(|e| e.cc_session_id.as_deref() == Some(&cc_session_id)),
};
if let Some(e) = target {
e.cc_session_id = Some(cc_session_id.clone());
e.mcp_channel_id = Some(channel_id.clone());
matched = true;
}
}) {
return Response::err(
req.id,
ErrorCode::Internal,
format!("registry write failed during channel registration: {e}"),
);
}
if !matched {
return Response::err(
req.id,
ErrorCode::ChannelUnknown,
"no agent matched cc_session_id/name for registration",
);
}
let _ = ctx
.emitter
.emit("channel_registered", &json!({"mcp_channel_id": channel_id}));
Response::ok(req.id, json!({"mcp_channel_id": channel_id}))
}
fn handle_unregister_channel(ctx: &Ctx, req: &Request) -> Response {
let channel_id = match req.params.get("mcp_channel_id").and_then(|v| v.as_str()) {
Some(s) => s.to_string(),
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `mcp_channel_id`"),
};
let mut cleared = false;
let _ = state::update_registry(&ctx.home.registry_json(), |r| {
for e in r.entries.iter_mut() {
if e.mcp_channel_id.as_deref() == Some(&channel_id) {
e.mcp_channel_id = None;
cleared = true;
}
}
});
if !cleared {
return Response::err(req.id, ErrorCode::ChannelUnknown, "unknown channel id");
}
Response::ok(req.id, json!({"unregistered": true}))
}
fn handle_push_to_channel(ctx: &Ctx, req: &Request) -> Response {
let channel_id = match req.params.get("mcp_channel_id").and_then(|v| v.as_str()) {
Some(s) => s.to_string(),
None => return Response::err(req.id, ErrorCode::InvalidParams, "missing `mcp_channel_id`"),
};
let envelope = match req.params.get("envelope") {
None => None,
Some(v @ Value::Object(_)) => Some(v.clone()),
Some(_) => {
return Response::err(
req.id,
ErrorCode::InvalidParams,
"`envelope` must be a JSON object",
)
}
};
let registry = state::load_registry(&ctx.home.registry_json()).unwrap_or_default();
let found = registry
.entries
.iter()
.any(|e| e.mcp_channel_id.as_deref() == Some(&channel_id));
if !found {
return Response::err(
req.id,
ErrorCode::ChannelUnknown,
"channel id not registered (channel server should re-register)",
);
}
let envelope = match envelope {
Some(e) => e,
None => {
return Response::ok(req.id, json!({"routed": true}));
}
};
match deliver_envelope(&channel_id, &envelope) {
Ok(()) => Response::ok(req.id, json!({"routed": true, "delivered": true})),
Err(reason) => Response::ok(
req.id,
json!({"routed": true, "delivered": false, "reason": reason}),
),
}
}
fn deliver_envelope(channel_id: &str, envelope: &Value) -> Result<(), String> {
use std::io::Write;
use std::process::Stdio;
let mut child = crate::loop_dispatch::fno_cmd("fno")
.args(["mcp", "send", "--session-id", channel_id])
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.map_err(|e| format!("spawn `fno mcp send` failed: {e}"))?;
{
let mut stdin = child
.stdin
.take()
.ok_or_else(|| "child stdin unavailable".to_string())?;
let bytes = serde_json::to_vec(envelope).map_err(|e| format!("serialize envelope: {e}"))?;
stdin
.write_all(&bytes)
.map_err(|e| format!("write envelope to `fno mcp send`: {e}"))?;
}
let out = child
.wait_with_output()
.map_err(|e| format!("wait for `fno mcp send`: {e}"))?;
if out.status.success() {
return Ok(());
}
let stderr = String::from_utf8_lossy(&out.stderr);
let tail = stderr.trim().rsplit('\n').next().unwrap_or("").trim();
Err(if tail.is_empty() {
format!("`fno mcp send` exited {}", out.status)
} else {
tail.to_string()
})
}
fn json_obj(pairs: &[(&str, Value)]) -> Map<String, Value> {
let mut m = Map::new();
for (k, v) in pairs {
m.insert((*k).to_string(), v.clone());
}
m
}
fn now_compact() -> String {
let secs = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let (y, mo, d, h, mi, s) = civil(secs);
format!("{y:04}{mo:02}{d:02}T{h:02}{mi:02}{s:02}Z")
}
pub(crate) fn now_rfc3339_like() -> String {
let secs = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let (y, mo, d, h, mi, s) = civil(secs);
format!("{y:04}-{mo:02}-{d:02}T{h:02}:{mi:02}:{s:02}Z")
}
fn civil(secs: u64) -> (i64, u32, u32, u32, u32, u32) {
let days = (secs / 86_400) as i64;
let rem = secs % 86_400;
let (hh, mm, ss) = (
(rem / 3600) as u32,
((rem % 3600) / 60) as u32,
(rem % 60) as u32,
);
let z = days + 719_468;
let era = if z >= 0 { z } else { z - 146_096 } / 146_097;
let doe = z - era * 146_097;
let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
let y = yoe + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
let mp = (5 * doy + 2) / 153;
let d = (doy - (153 * mp + 2) / 5 + 1) as u32;
let m = if mp < 10 { mp + 3 } else { mp - 9 } as u32;
(if m <= 2 { y + 1 } else { y }, m, d, hh, mm, ss)
}
fn uuid_v4() -> String {
let mut b = [0u8; 16];
fill_random(&mut b);
b[6] = (b[6] & 0x0f) | 0x40; b[8] = (b[8] & 0x3f) | 0x80; format!(
"{:02x}{:02x}{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}{:02x}{:02x}{:02x}{:02x}",
b[0], b[1], b[2], b[3], b[4], b[5], b[6], b[7], b[8], b[9], b[10], b[11], b[12], b[13],
b[14], b[15]
)
}
fn fill_random(buf: &mut [u8]) {
if let Ok(mut f) = std::fs::File::open("/dev/urandom") {
use std::io::Read;
if f.read_exact(buf).is_ok() {
return;
}
}
let seed = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos() as u64
^ (std::process::id() as u64).rotate_left(17);
let mut x = seed | 1;
for byte in buf.iter_mut() {
x ^= x << 13;
x ^= x >> 7;
x ^= x << 17;
*byte = (x & 0xff) as u8;
}
}
#[cfg(test)]
mod tests {
use super::*;
fn canonical_name_in(registry: &state::Registry, token: &str) -> String {
let Ok(Value::Array(rows)) = serde_json::to_value(®istry.entries) else {
return token.to_string();
};
match crate::client_verbs::find_agent_entry(&rows, token) {
Ok(entry) => entry
.get("name")
.and_then(Value::as_str)
.unwrap_or(token)
.to_string(),
Err(_) => token.to_string(),
}
}
use crate::state::{AgentState, DriveWindow, PtyState};
fn tmp_home(tag: &str) -> AgentsHome {
let mut p = std::env::temp_dir();
p.push(format!(
"fno-agents-daemon-{}-{}-{}",
tag,
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let home = AgentsHome::at(&p);
home.ensure_root().unwrap();
home
}
fn read_events(home: &AgentsHome) -> Vec<Value> {
std::fs::read_to_string(home.events_jsonl())
.unwrap_or_default()
.lines()
.filter_map(|l| serde_json::from_str::<Value>(l).ok())
.collect()
}
fn ask_row(name: &str, exited_at: Option<&str>) -> RegistryEntry {
RegistryEntry {
name: name.into(),
short_id: String::new(),
legacy_provider: "claude".into(),
harness: None,
harness_session_id: None,
cwd: "/tmp".into(),
project_root: String::new(),
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: None,
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
cc_session_id: None,
host_mode: None,
status: AgentStatus::Exited,
last_message_at: None,
created_at: "2020-01-01T00:00:00Z".into(),
pid: None,
pid_start_time: None,
log_path: None,
last_reconciled_at: None,
inside_leg: None,
exited_at: exited_at.map(str::to_string),
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
}
}
#[test]
fn gc_sweep_reaps_stamped_stamps_unstamped_keeps_live() {
let home = tmp_home("gc-sweep");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
state::update_registry(&home.registry_json(), |r| {
r.entries
.push(ask_row("ask-old", Some("2020-01-01T00:00:00Z")));
r.entries.push(ask_row("ask-new", None));
let mut live = ask_row("live", None);
live.name = "live".into();
live.short_id = "wkL".into();
live.status = AgentStatus::Live;
live.pid = Some(std::process::id());
r.entries.push(live);
})
.unwrap();
let summary = gc_sweep(&home, &emitter, Duration::from_secs(3600));
assert_eq!(summary.reaped, vec!["ask-old".to_string()]);
let reg = state::load_registry(&home.registry_json()).unwrap();
let names: Vec<&str> = reg.entries.iter().map(|e| e.name.as_str()).collect();
assert!(!names.contains(&"ask-old"), "ask-old should be reaped");
assert!(
names.contains(&"ask-new"),
"ask-new should be kept (in grace)"
);
assert!(names.contains(&"live"), "live row must never be reaped");
let new = reg.entries.iter().find(|e| e.name == "ask-new").unwrap();
assert!(
new.exited_at.is_some(),
"ask-new should be stamped this pass"
);
let live = reg.entries.iter().find(|e| e.name == "live").unwrap();
assert!(live.exited_at.is_none());
let events = read_events(&home);
let reaped: Vec<&Value> = events
.iter()
.filter(|e| e.get("type").and_then(Value::as_str) == Some("agent_row_reaped"))
.collect();
assert_eq!(reaped.len(), 1);
assert_eq!(
reaped[0]
.get("data")
.and_then(|d| d.get("name"))
.and_then(Value::as_str),
Some("ask-old")
);
}
#[test]
fn gc_sweep_empty_registry_is_noop() {
let home = tmp_home("gc-empty");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let summary = gc_sweep(&home, &emitter, Duration::from_secs(3600));
assert!(summary.reaped.is_empty());
assert!(summary.kept_dirty.is_empty());
}
#[test]
fn gc_sweep_turns_unterminated_node_reap_into_durable_failure() {
let sandbox = tmp_home("gc-dead-dispatch");
let home = AgentsHome::at(sandbox.root().join("agents"));
home.ensure_root().unwrap();
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let dead_repo = home.root().join("dead-repo");
let done_repo = home.root().join("done-repo");
for repo in [&dead_repo, &done_repo] {
std::fs::create_dir_all(repo.join(".fno")).unwrap();
assert!(std::process::Command::new("git")
.args(["init", "-q"])
.current_dir(repo)
.status()
.unwrap()
.success());
}
let dead_session = "target-run-dead";
let done_session = "target-run-done";
std::fs::write(
dead_repo.join(".fno/target-state.md"),
format!("---\nfno_id: {dead_session}\ninput: x-a35a\nplan_path: \"\"\n---\n"),
)
.unwrap();
std::fs::write(
done_repo.join(".fno/target-state.md"),
format!("---\nfno_id: {done_session}\ninput: x-b44e\nplan_path: \"\"\n---\n"),
)
.unwrap();
state::update_registry(&home.registry_json(), |r| {
let mut dead = bg_claude_row("target-x-a35a-route-atomicity", "dead0001");
dead.status = AgentStatus::Exited;
dead.cwd = dead_repo.to_string_lossy().into_owned();
dead.exited_at = Some("2020-01-01T00:00:00Z".into());
dead.harness_session_id = Some("harness-dead-uuid".into());
r.entries.push(dead);
let mut done = bg_claude_row("target-x-b44e-finished", "done0002");
done.status = AgentStatus::Exited;
done.cwd = done_repo.to_string_lossy().into_owned();
done.exited_at = Some("2020-01-01T00:00:00Z".into());
done.harness_session_id = Some("harness-done-uuid".into());
r.entries.push(done);
})
.unwrap();
let global_events = home.root().parent().unwrap().join("events.jsonl");
std::fs::write(
done_repo.join(".fno/events.jsonl.1"),
format!(
"{{\"ts\":\"2026-07-24T00:00:00Z\",\"type\":\"termination\",\"source\":\"loop\",\"data\":{{\"session_id\":\"{done_session}\",\"reason\":\"DonePRGreen\",\"message\":\"done\"}}}}\n"
),
)
.unwrap();
let summary = gc_sweep(&home, &emitter, Duration::from_secs(0));
assert_eq!(summary.reaped.len(), 2);
let reaps = read_events(&home);
let dead_reap = reaps
.iter()
.find(|e| e["data"]["short_id"] == "dead0001")
.expect("dead dispatch reap event");
assert_eq!(dead_reap["data"]["node_id"], "x-a35a");
assert_eq!(dead_reap["data"]["termination_event"], false);
let done_reap = reaps
.iter()
.find(|e| e["data"]["short_id"] == "done0002")
.expect("completed dispatch reap event");
assert_eq!(done_reap["data"]["node_id"], "x-b44e");
assert_eq!(done_reap["data"]["termination_event"], true);
let global = std::fs::read_to_string(&global_events).unwrap();
let failures: Vec<Value> = global
.lines()
.filter_map(|line| serde_json::from_str(line).ok())
.filter(|e: &Value| e["type"] == "node_failed")
.collect();
assert_eq!(failures.len(), 1);
assert_eq!(failures[0]["data"]["unit_id"], "x-a35a");
assert_eq!(failures[0]["data"]["session_id"], dead_session);
assert_eq!(
failures[0]["data"]["reason"],
"agent-row-reaped-no-termination"
);
}
#[test]
fn gc_sweep_restores_row_when_termination_evidence_is_unknown() {
let sandbox = tmp_home("gc-unknown-termination");
let home = AgentsHome::at(sandbox.root().join("agents"));
home.ensure_root().unwrap();
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let repo = home.root().join("repo");
std::fs::create_dir_all(repo.join(".fno")).unwrap();
assert!(std::process::Command::new("git")
.args(["init", "-q"])
.current_dir(&repo)
.status()
.unwrap()
.success());
std::fs::write(
repo.join(".fno/target-state.md"),
"---\nfno_id: reused-run\ninput: x-other\nplan_path: \"\"\n---\n",
)
.unwrap();
state::update_registry(&home.registry_json(), |registry| {
let mut row = bg_claude_row("target-x-a35a-route-atomicity", "dead0001");
row.status = AgentStatus::Exited;
row.cwd = repo.to_string_lossy().into_owned();
row.exited_at = Some("2020-01-01T00:00:00Z".into());
registry.entries.push(row);
})
.unwrap();
let summary = gc_sweep(&home, &emitter, Duration::from_secs(0));
assert!(summary.reaped.is_empty());
let registry = state::load_registry(&home.registry_json()).unwrap();
assert!(registry
.entries
.iter()
.any(|row| row.name == "target-x-a35a-route-atomicity"));
let events = read_events(&home);
assert!(events.iter().any(|event| {
event["type"] == "daemon_recovery_error"
&& event["data"]["op"] == "observe_dead_dispatch_termination"
}));
}
#[test]
fn gc_sweep_restores_row_when_dead_dispatch_receipt_cannot_persist() {
let sandbox = tmp_home("gc-dead-dispatch-write-failure");
let home = AgentsHome::at(sandbox.root().join("agents"));
home.ensure_root().unwrap();
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let repo = home.root().join("repo");
std::fs::create_dir_all(&repo).unwrap();
assert!(std::process::Command::new("git")
.args(["init", "-q"])
.current_dir(&repo)
.status()
.unwrap()
.success());
state::update_registry(&home.registry_json(), |registry| {
let mut row = bg_claude_row("target-x-a35a-route-atomicity", "dead0001");
row.status = AgentStatus::Exited;
row.cwd = repo.to_string_lossy().into_owned();
row.exited_at = Some("2020-01-01T00:00:00Z".into());
registry.entries.push(row);
})
.unwrap();
std::fs::create_dir_all(global_events_path(&home)).unwrap();
let summary = gc_sweep(&home, &emitter, Duration::from_secs(0));
assert!(summary.reaped.is_empty());
let registry = state::load_registry(&home.registry_json()).unwrap();
assert!(registry
.entries
.iter()
.any(|row| row.name == "target-x-a35a-route-atomicity"));
let events = read_events(&home);
assert!(events.iter().any(|event| {
event["type"] == "daemon_recovery_error"
&& event["data"]["op"] == "record_dead_dispatch"
}));
}
#[test]
fn recovery_emits_drive_crashed_before_clearing_window() {
let home = tmp_home("recover-drive");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
state::update_registry(&home.registry_json(), |r| {
r.entries.push(RegistryEntry {
name: "worker-A".into(),
short_id: "wkA".into(),
legacy_provider: "codex".into(),
harness: None,
harness_session_id: None,
cwd: "/tmp".into(),
project_root: "/tmp".into(),
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: None,
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
cc_session_id: None,
host_mode: None,
status: AgentStatus::Live,
last_message_at: None,
created_at: "2026-05-24T00:00:00Z".into(),
pid: Some(std::process::id()), pid_start_time: None,
log_path: None,
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
});
})
.unwrap();
let mut st = AgentState::new_pty("wkA");
st.status = AgentStatus::Live;
st.pty = Some(PtyState {
active: true,
drive: Some(DriveWindow {
session_id: Some("drive-xyz".into()),
mode: Some("interactive".into()),
last_heartbeat_at_monotonic_ns: Some(123),
}),
});
state::write_state_atomic(&home.state_json("wkA"), &st).unwrap();
let report = recover(&home, &emitter);
assert_eq!(report.recovered_drives, vec!["wkA".to_string()]);
let events = read_events(&home);
let crashed = events
.iter()
.find(|e| e["type"] == "drive_crashed")
.expect("drive_crashed emitted");
assert_eq!(crashed["data"]["session_id"], "drive-xyz");
assert_eq!(crashed["data"]["reason"], "daemon_restart");
let after = state::load_state(&home.state_json("wkA")).unwrap().unwrap();
let pty = after.pty.unwrap();
assert!(pty.drive.is_none());
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn recovery_marks_missing_state_inconsistent() {
let home = tmp_home("recover-missing");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
state::update_registry(&home.registry_json(), |r| {
r.entries.push(RegistryEntry {
name: "ghost".into(),
short_id: "ghost".into(),
legacy_provider: "codex".into(),
harness: None,
harness_session_id: None,
cwd: "/tmp".into(),
project_root: "/tmp".into(),
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: None,
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
cc_session_id: None,
host_mode: None,
status: AgentStatus::Live,
last_message_at: None,
created_at: "2026-05-24T00:00:00Z".into(),
pid: None,
pid_start_time: None,
log_path: None,
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
});
})
.unwrap();
let report = recover(&home, &emitter);
assert_eq!(
report.inconsistent,
vec![("ghost".to_string(), InconsistencyReason::MissingStateJson)]
);
let events = read_events(&home);
assert!(events
.iter()
.any(|e| e["type"] == "agent_inconsistent"
&& e["data"]["reason"] == "missing_state_json"));
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn recovery_skips_claude_shellout_rows_no_spurious_inconsistent() {
let home = tmp_home("recover-claude-shellout");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
state::update_registry(&home.registry_json(), |r| {
let mut bg = bg_claude_row("bg-ask", "7c5dcf5d");
bg.host_mode = None;
r.entries.push(bg);
let mut adopted = bg_claude_row("cc-adopt", "deadbeef");
adopted.host_mode = Some(crate::state::HOST_MODE_ATTACHED.into());
adopted.pid = Some(4242);
r.entries.push(adopted);
})
.unwrap();
let report = recover(&home, &emitter);
assert!(
report.inconsistent.is_empty(),
"claude shellout/adopted rows must not be flagged inconsistent: {:?}",
report.inconsistent
);
let events = read_events(&home);
assert!(
!events.iter().any(|e| e["type"] == "agent_inconsistent"),
"no agent_inconsistent event for claude shellout rows"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn canonical_name_in_resolves_all_three_address_forms() {
let full = "aabbccdd-1111-2222-3333-444455556666";
let mut row = rentry("billing", AgentStatus::Live, None);
row.short_id = "a1b2c3d4".into();
row.harness_session_id = Some(full.into());
let reg = crate::state::Registry {
schema_version: crate::state::REGISTRY_SCHEMA_VERSION,
entries: vec![row],
};
assert_eq!(canonical_name_in(®, "billing"), "billing"); assert_eq!(canonical_name_in(®, "a1b2c3d4"), "billing"); assert_eq!(canonical_name_in(®, full), "billing"); assert_eq!(
canonical_name_in(®, "AABBCCDD-1111-2222-3333-444455556666"),
"billing"
); assert_eq!(canonical_name_in(®, "nope"), "nope");
}
#[tokio::test]
async fn lifecycle_name_resolution_never_falls_back_on_ambiguity() {
let mut named = rentry("deadbeef", AgentStatus::Live, None);
named.short_id = "transport-a".into();
named.harness_session_id = Some("aaaaaaaa-1111-2222-3333-444455556666".into());
let mut short = rentry("other", AgentStatus::Live, None);
short.short_id = "deadbeef".into();
short.harness_session_id = Some("bbbbbbbb-1111-2222-3333-000000000002".into());
let reg = crate::state::Registry {
schema_version: crate::state::REGISTRY_SCHEMA_VERSION,
entries: vec![named, short],
};
let error = entry_for_lifecycle(
®,
"deadbeef",
std::path::Path::new("/nonexistent/registry.json"),
)
.await
.expect_err("ambiguous token must not fall back to the matching row name");
assert!(error.contains("ambiguous across 2 agents"));
}
#[test]
fn recovery_reaps_dead_pid() {
let home = tmp_home("recover-reap");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
state::update_registry(&home.registry_json(), |r| {
r.entries.push(RegistryEntry {
name: "dead".into(),
short_id: "dead".into(),
legacy_provider: "codex".into(),
harness: None,
harness_session_id: None,
cwd: "/tmp".into(),
project_root: "/tmp".into(),
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: None,
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
cc_session_id: None,
host_mode: None,
status: AgentStatus::Live,
last_message_at: None,
created_at: "2026-05-24T00:00:00Z".into(),
pid: Some(0x7fff_fff0),
pid_start_time: None,
log_path: None,
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
});
})
.unwrap();
let mut st = AgentState::new_pty("dead");
st.status = AgentStatus::Live;
state::write_state_atomic(&home.state_json("dead"), &st).unwrap();
let report = recover(&home, &emitter);
assert_eq!(report.reaped_pids, vec![0x7fff_fff0]);
let reg = state::load_registry(&home.registry_json()).unwrap();
assert_eq!(reg.find("dead").unwrap().status, AgentStatus::Exited);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn recovery_marks_dead_interactive_exited_and_preserves_host_mode() {
let home = tmp_home("recover-interactive");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
state::update_registry(&home.registry_json(), |r| {
let mut e = rentry("hosted", AgentStatus::Live, None);
e.host_mode = Some(crate::state::HOST_MODE_INTERACTIVE.to_string());
e.pid = Some(0x7fff_fff0); r.entries.push(e);
})
.unwrap();
let mut st = AgentState::new_pty("hosted");
st.status = AgentStatus::Live;
state::write_state_atomic(&home.state_json("hosted"), &st).unwrap();
let _ = recover(&home, &emitter);
let reg = state::load_registry(&home.registry_json()).unwrap();
let row = reg.find("hosted").unwrap();
assert_eq!(
row.status,
AgentStatus::Exited,
"a dead interactive worker is exited, never orphaned"
);
assert_eq!(
row.host_mode_or_default(),
crate::state::HOST_MODE_INTERACTIVE,
"host_mode must survive recovery"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn pid_is_ours_distinguishes_recycled_pid() {
let me = std::process::id();
let Some(st) = process_start_time(me) else {
return; };
assert!(pid_is_ours(me, Some(st)), "correct start time -> ours");
assert!(
!pid_is_ours(me, Some(st.wrapping_add(1))),
"alive but mismatched start time -> recycled, not ours"
);
assert!(
!pid_is_ours(0x7fff_fff0, Some(st)),
"dead pid is never ours"
);
assert!(
pid_is_ours(me, None),
"no recorded start time -> fall back to bare liveness (legacy)"
);
}
#[tokio::test]
async fn stop_claude_pid_kills_a_real_child_and_spares_a_recycled_pid() {
let mut entry = ask_row("orphan", None);
assert!(
!stop_claude_pid_confirmed(&entry).await,
"no pid -> nothing to stop"
);
let out = std::process::Command::new("sh")
.arg("-c")
.arg("sleep 60 >/dev/null 2>&1 & echo $!")
.output()
.expect("spawn detached sleeper");
let pid: u32 = String::from_utf8_lossy(&out.stdout)
.trim()
.parse()
.expect("sleeper pid");
let start = process_start_time(pid);
let ps_says_alive = |pid: u32| {
std::process::Command::new("ps")
.args(["-p", &pid.to_string()])
.output()
.map(|o| {
String::from_utf8_lossy(&o.stdout)
.lines()
.filter(|l| l.split_whitespace().next() == Some(&pid.to_string()))
.count()
> 0
})
.unwrap_or(false)
};
entry.pid = Some(pid);
entry.pid_start_time = None;
assert!(
!stop_claude_pid_confirmed(&entry).await,
"no start token -> refuse"
);
assert!(ps_says_alive(pid), "a refused row must not be signalled");
if let Some(st) = start {
entry.pid_start_time = Some(st.wrapping_add(1));
assert!(
!stop_claude_pid_confirmed(&entry).await,
"recycled pid -> refuse"
);
assert!(
ps_says_alive(pid),
"an unrelated process must not be signalled"
);
}
entry.pid_start_time = start;
if start.is_some() {
assert!(
stop_claude_pid_confirmed(&entry).await,
"owned live pid -> stopped"
);
assert!(!ps_says_alive(pid), "process is gone");
} else {
unsafe {
libc::kill(pid as libc::pid_t, libc::SIGKILL);
}
}
}
#[test]
fn pid_is_ours_rejects_an_out_of_range_pid() {
assert!(!pid_is_ours(u32::MAX, None), "u32::MAX must never be ours");
assert!(
!pid_is_ours(i32::MAX as u32 + 1, Some(123)),
"anything past i32::MAX wraps negative"
);
assert!(
pid_confirmed_dead(u32::MAX),
"out-of-range is never running"
);
}
#[test]
fn recycle_and_death_each_demand_positive_evidence() {
let me = std::process::id();
let Some(st) = process_start_time(me) else {
return; };
assert!(!pid_confirmed_dead(me), "a live pid is not dead");
assert!(
!pid_recycled(me, Some(st)),
"matching token is not a recycle"
);
assert!(
pid_recycled(me, Some(st.wrapping_add(1))),
"reachable + differing token is a recycle"
);
assert!(!pid_recycled(me, None), "no token -> no recycle verdict");
let dead = 0x7fff_fff0u32;
assert!(pid_confirmed_dead(dead), "unused high pid reads as dead");
assert!(
!pid_recycled(dead, Some(st)),
"dead is not a recycle finding"
);
}
#[test]
fn recovery_reaps_recycled_pid() {
let home = tmp_home("recover-recycled");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let me = std::process::id();
if process_start_time(me).is_none() {
std::fs::remove_dir_all(home.root()).ok();
return; }
state::update_registry(&home.registry_json(), |r| {
r.entries.push(RegistryEntry {
name: "recycled".into(),
short_id: "recycled".into(),
legacy_provider: "codex".into(),
harness: None,
harness_session_id: None,
cwd: "/tmp".into(),
project_root: "/tmp".into(),
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: None,
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
cc_session_id: None,
host_mode: None,
status: AgentStatus::Live,
last_message_at: None,
created_at: "2026-05-24T00:00:00Z".into(),
pid: Some(me),
pid_start_time: Some(1),
log_path: None,
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
});
})
.unwrap();
let mut st = AgentState::new_pty("recycled");
st.status = AgentStatus::Live;
state::write_state_atomic(&home.state_json("recycled"), &st).unwrap();
let report = recover(&home, &emitter);
assert_eq!(report.reaped_pids, vec![me]);
let reg = state::load_registry(&home.registry_json()).unwrap();
assert_eq!(reg.find("recycled").unwrap().status, AgentStatus::Exited);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn recovery_archives_orphan_state_dir() {
let home = tmp_home("recover-orphan");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let mut st = AgentState::new_pty("loner");
st.status = AgentStatus::Live;
state::write_state_atomic(&home.state_json("loner"), &st).unwrap();
let report = recover(&home, &emitter);
assert_eq!(report.archived_orphans, vec!["loner".to_string()]);
assert!(!home.agent_dir("loner").exists(), "orphan dir moved aside");
assert!(home.orphaned_dir().exists());
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn agent_name_validation() {
assert!(valid_agent_name("worker-A_1"));
assert!(!valid_agent_name(""));
assert!(!valid_agent_name(&"x".repeat(65)));
assert!(!valid_agent_name("has space"));
assert!(!valid_agent_name("inject;rm"));
}
#[test]
fn uuid_v4_shape_and_uniqueness() {
let a = uuid_v4();
let b = uuid_v4();
assert_ne!(a, b);
assert_eq!(a.len(), 36);
let parts: Vec<&str> = a.split('-').collect();
assert_eq!(
parts.iter().map(|p| p.len()).collect::<Vec<_>>(),
vec![8, 4, 4, 4, 12]
);
assert_eq!(&a[14..15], "4");
assert!(matches!(&a[19..20], "8" | "9" | "a" | "b"));
}
#[test]
fn short_id_derivation_dedups() {
let mut reg = state::Registry::default();
assert_eq!(derive_short_id("worker-A", ®), "workerA");
reg.entries.push(RegistryEntry {
name: "x".into(),
short_id: "workerA".into(),
legacy_provider: "codex".into(),
harness: None,
harness_session_id: None,
cwd: "/".into(),
project_root: "/".into(),
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: None,
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
cc_session_id: None,
host_mode: None,
status: AgentStatus::Live,
last_message_at: None,
created_at: "t".into(),
pid: None,
pid_start_time: None,
log_path: None,
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
});
assert_eq!(derive_short_id("worker-A", ®), "workerA1");
}
fn rentry(name: &str, status: AgentStatus, last_reconciled: Option<&str>) -> RegistryEntry {
RegistryEntry {
name: name.into(),
short_id: name.into(),
legacy_provider: "codex".into(),
harness: None,
harness_session_id: None,
cwd: "/tmp".into(),
project_root: "/tmp".into(),
session_id: Some("sid".into()),
legacy_claude_short_id: None,
claude_session_uuid: None,
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
host_mode: None,
cc_session_id: None,
status,
last_message_at: None,
created_at: "t".into(),
pid: None,
pid_start_time: None,
log_path: None,
last_reconciled_at: last_reconciled.map(String::from),
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
}
}
fn probe_err() -> crate::provider::ReachabilityProbeError {
crate::provider::ReachabilityProbeError::new("codex", "store unavailable")
}
fn bg_claude_row(name: &str, short_id: &str) -> RegistryEntry {
let mut e = rentry(name, AgentStatus::Live, None);
e.legacy_provider = "claude".into();
e.short_id = short_id.into();
e.claude_session_uuid = None;
e
}
#[test]
fn find_uuid_backfill_row_matches_null_uuid_by_short_prefix() {
let rows = vec![bg_claude_row("w", "3228ccad")];
assert!(matches!(
find_uuid_backfill_row(&rows, "3228ccad-c078-4b53-a8c9-7199b831eae4"),
UuidBackfill::One(0)
));
}
#[test]
fn find_uuid_backfill_row_refuses_ambiguous_short_collision() {
let rows = vec![
bg_claude_row("w1", "3228ccad"),
bg_claude_row("w2", "3228ccad"),
];
assert!(matches!(
find_uuid_backfill_row(&rows, "3228ccad-c078-4b53-a8c9-7199b831eae4"),
UuidBackfill::Ambiguous
));
}
#[test]
fn find_uuid_backfill_row_skips_rows_that_already_have_a_uuid() {
let mut row = bg_claude_row("w", "3228ccad");
row.claude_session_uuid = Some("3228ccad-c078-4b53-a8c9-7199b831eae4".into());
assert!(matches!(
find_uuid_backfill_row(&[row], "3228ccad-c078-4b53-a8c9-7199b831eae4"),
UuidBackfill::None
));
}
#[test]
fn find_uuid_backfill_row_skips_non_claude_rows() {
let mut row = bg_claude_row("w", "3228ccad");
row.legacy_provider = "codex".into();
assert!(matches!(
find_uuid_backfill_row(&[row], "3228ccad-c078-4b53-a8c9-7199b831eae4"),
UuidBackfill::None
));
}
#[test]
fn find_uuid_backfill_row_requires_group_boundary() {
let rows = vec![bg_claude_row("w", "3228ccad")];
assert!(matches!(
find_uuid_backfill_row(&rows, "3228ccadd-c078-4b53-a8c9-7199b831eae4"),
UuidBackfill::None
));
}
#[test]
fn concurrent_spawn_name_reservation_inserts_once() {
let home = tmp_home("spawn-reserve");
let path = home.registry_json();
let reserve = |entry: RegistryEntry| -> bool {
state::update_registry(&path, move |r| {
if r.entries.iter().any(|e| e.name == entry.name) {
return false;
}
r.entries.push(entry);
true
})
.unwrap()
};
assert!(
reserve(rentry("dup", AgentStatus::Live, None)),
"first wins"
);
assert!(
!reserve(rentry("dup", AgentStatus::Live, None)),
"second loses the race -> no insert"
);
let reg = state::load_registry(&path).unwrap();
assert_eq!(
reg.entries.iter().filter(|e| e.name == "dup").count(),
1,
"exactly one row for the contended name"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn reconcile_flips_unreachable_live_to_orphaned_and_recovers_orphaned() {
let entries = vec![
rentry("live-but-gone", AgentStatus::Live, None),
rentry("back-from-dead", AgentStatus::Orphaned, None),
];
let (changes, out) = plan_reconcile(
&entries,
|e| match e.name.as_str() {
"live-but-gone" => Ok(false), _ => Ok(true), },
|| false,
|_| true,
);
assert_eq!(out.orphans, vec!["live-but-gone".to_string()]);
assert_eq!(out.recovered, vec!["back-from-dead".to_string()]);
assert_eq!(out.updated.len(), 2);
assert_eq!(
changes[0].new_status,
Some(AgentStatus::Orphaned),
"unreachable live agent should orphan"
);
assert_eq!(changes[1].new_status, Some(AgentStatus::Live));
}
#[test]
fn reconcile_does_not_orphan_a_live_interactive_host_on_store_miss() {
let mut interactive = rentry("hosted-tui", AgentStatus::Live, None);
interactive.host_mode = Some(crate::state::HOST_MODE_INTERACTIVE.to_string());
let exec = rentry("one-shot", AgentStatus::Live, None);
let entries = vec![interactive, exec];
let (changes, out) = plan_reconcile(&entries, |_| Ok(false), || false, |_| true);
assert_eq!(
changes[0].new_status, None,
"a live interactive host must not be orphaned on a session-store miss"
);
assert_eq!(
changes[1].new_status,
Some(AgentStatus::Orphaned),
"an exec sibling with the same probe result is still orphaned"
);
assert_eq!(out.orphans, vec!["one-shot".to_string()]);
}
#[test]
fn reconcile_reaps_a_dead_interactive_host_to_exited() {
let mut dead = rentry("dead-tui", AgentStatus::Live, None);
dead.host_mode = Some(crate::state::HOST_MODE_INTERACTIVE.to_string());
let mut live = rentry("live-tui", AgentStatus::Live, None);
live.host_mode = Some(crate::state::HOST_MODE_INTERACTIVE.to_string());
let entries = vec![dead, live];
let (changes, out) = plan_reconcile(
&entries,
|_| Ok(false), || false,
|e| e.name == "live-tui", );
assert_eq!(
changes[0].new_status,
Some(AgentStatus::Exited),
"a dead interactive host is reaped to Exited during reconcile"
);
assert_eq!(
changes[1].new_status, None,
"a live interactive host is left untouched"
);
assert!(out.orphans.is_empty());
assert_eq!(out.updated, vec!["dead-tui".to_string()]);
}
#[test]
fn reconcile_mux_pane_liveness_follows_the_pid_not_the_store() {
let mk = |name: &str, pid: Option<u32>| {
let mut e = rentry(name, AgentStatus::Live, None);
e.mux = Some(crate::state::MuxRef {
session: "main".into(),
pane_id: 7,
});
e.pid = pid;
e
};
let entries = vec![
mk("live-pane", Some(4242)), mk("dead-pane", Some(4243)), mk("pidless-pane", None), ];
let (changes, out) = plan_reconcile(
&entries,
|_| Ok(false), || false,
|e| e.name == "live-pane", );
assert_eq!(
changes[0].new_status, None,
"a live-pid mux pane is preserved"
);
assert_eq!(
changes[1].new_status,
Some(AgentStatus::Exited),
"a dead-pid mux pane is reaped to Exited"
);
assert_eq!(
changes[2].new_status,
Some(AgentStatus::Orphaned),
"a pid-less mux pane defers to store liveness (orphan), not immortal"
);
assert_eq!(out.orphans, vec!["pidless-pane".to_string()]);
}
#[test]
fn reconcile_store_hit_does_not_resurrect_a_pid_dead_row() {
let entries = vec![
rentry("dead-orphan", AgentStatus::Orphaned, None),
rentry("live-orphan", AgentStatus::Orphaned, None),
];
let (changes, out) = plan_reconcile(
&entries,
|_| Ok(true), || false,
|e| e.name == "live-orphan",
);
assert_eq!(
changes[0].new_status, None,
"a store hit must not recover a row whose pid is dead"
);
assert_eq!(
changes[1].new_status,
Some(AgentStatus::Live),
"a store hit on a live pid still recovers"
);
assert_eq!(out.recovered, vec!["live-orphan".to_string()]);
}
#[test]
fn reconcile_pidless_orphan_still_recovers_on_store_hit() {
let entries = vec![rentry("pidless", AgentStatus::Orphaned, None)];
let (changes, out) = plan_reconcile(&entries, |_| Ok(true), || false, |_| true);
assert_eq!(changes[0].new_status, Some(AgentStatus::Live));
assert_eq!(out.recovered, vec!["pidless".to_string()]);
}
#[test]
fn to_agent_entry_projects_the_opencode_session_id() {
let mut e = rentry("oc", AgentStatus::Live, None);
e.legacy_provider = "opencode".into();
e.harness = Some("opencode".into());
e.harness_session_id = Some("ses_09679f284ffeJv7NdBAoLQLnLZ".into());
e.session_id = None;
assert_eq!(
to_agent_entry(&e).session_id.as_deref(),
Some("ses_09679f284ffeJv7NdBAoLQLnLZ")
);
}
#[test]
fn reconcile_inconclusive_preserves_status() {
let entries = vec![rentry("flaky", AgentStatus::Live, None)];
let (changes, out) = plan_reconcile(&entries, |_| Err(probe_err()), || false, |_| true);
assert_eq!(changes[0].new_status, None, "must NOT flip on inconclusive");
assert!(out.orphans.is_empty());
assert_eq!(out.inconsistent.len(), 1);
assert_eq!(out.inconsistent[0].0, "flaky");
}
#[test]
fn reconcile_leaves_terminal_states_untouched() {
let entries = vec![
rentry("done", AgentStatus::Exited, None),
rentry("dead", AgentStatus::PermanentDead, None),
];
let (changes, out) = plan_reconcile(&entries, |_| Ok(false), || false, |_| true);
assert!(changes.iter().all(|c| c.new_status.is_none()));
assert!(out.orphans.is_empty() && out.updated.is_empty());
}
fn ask_entry(name: &str, status: AgentStatus) -> RegistryEntry {
let mut e = rentry(name, status, None);
e.short_id = String::new();
e.pid = None;
e.codex_session_id = Some("resume-uuid".into());
e.session_id = None;
e
}
#[test]
fn reconcile_one_shot_ask_settles_to_exited_even_when_reachable() {
let entries = vec![ask_entry("codex-ask", AgentStatus::Live)];
let (changes, out) = plan_reconcile(
&entries,
|_| Ok(true), || false,
|_| true,
);
assert_eq!(
changes[0].new_status,
Some(AgentStatus::Exited),
"a finished ask settles to exited even when its session file is reachable"
);
assert_eq!(out.updated, vec!["codex-ask".to_string()]);
assert!(out.orphans.is_empty(), "an ask is exited, never orphaned");
assert_eq!(entries[0].codex_session_id.as_deref(), Some("resume-uuid"));
}
#[test]
fn reconcile_one_shot_ask_already_terminal_is_untouched() {
let entries = vec![ask_entry("done-ask", AgentStatus::Exited)];
let (changes, out) = plan_reconcile(&entries, |_| Ok(true), || false, |_| true);
assert_eq!(changes[0].new_status, None);
assert!(out.updated.is_empty());
}
#[test]
fn apply_reconcile_change_clears_pid_only_on_exited() {
let mut to_exited = rentry("x", AgentStatus::Live, None);
to_exited.pid = Some(4242);
to_exited.pid_start_time = Some(99);
to_exited.inside_leg = Some(state::InsideLegReport {
state: state::InsideLegState::Working,
seq: 3,
reason: None,
received_at: "2026-06-27T00:00:00Z".into(),
ttl_ms: None,
});
apply_reconcile_change(&mut to_exited, Some(AgentStatus::Exited), "T1");
assert_eq!(to_exited.status, AgentStatus::Exited);
assert_eq!(to_exited.pid, None, "exited row must drop its pid");
assert_eq!(to_exited.pid_start_time, None);
assert_eq!(
to_exited.inside_leg, None,
"exited row must clear the inside-leg authority (E3.3 / AC-X2-4)"
);
assert_eq!(to_exited.last_reconciled_at.as_deref(), Some("T1"));
let mut to_orphaned = rentry("y", AgentStatus::Live, None);
to_orphaned.pid = Some(4242);
to_orphaned.inside_leg = Some(state::InsideLegReport {
state: state::InsideLegState::Working,
seq: 1,
reason: None,
received_at: "2026-06-27T00:00:00Z".into(),
ttl_ms: None,
});
apply_reconcile_change(&mut to_orphaned, Some(AgentStatus::Orphaned), "T2");
assert_eq!(to_orphaned.status, AgentStatus::Orphaned);
assert_eq!(
to_orphaned.pid,
Some(4242),
"non-exited transition keeps pid"
);
assert!(
to_orphaned.inside_leg.is_some(),
"a non-exit transition keeps the inside-leg report (only exit tears it down)"
);
let mut no_change = rentry("z", AgentStatus::Live, Some("OLD"));
no_change.pid = Some(4242);
apply_reconcile_change(&mut no_change, None, "T3");
assert_eq!(no_change.status, AgentStatus::Live);
assert_eq!(no_change.pid, Some(4242));
assert_eq!(no_change.last_reconciled_at.as_deref(), Some("T3"));
}
#[test]
fn emit_inside_leg_completion_publishes_only_for_report_bearing_rows() {
let home = tmp_home("inside-leg-completion");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let mut with_report = rentry("pane", AgentStatus::Live, None);
with_report.session_id = Some("sess-uuid".into());
with_report.inside_leg = Some(state::InsideLegReport {
state: state::InsideLegState::Working,
seq: 9,
reason: Some("running tests".into()),
received_at: "2026-06-27T00:00:00Z".into(),
ttl_ms: Some(5000),
});
emit_inside_leg_completion(&emitter, &with_report);
emit_inside_leg_completion(&emitter, &rentry("plain", AgentStatus::Live, None));
let log = std::fs::read_to_string(home.events_jsonl()).unwrap_or_default();
let events: Vec<serde_json::Value> = log
.lines()
.filter_map(|l| serde_json::from_str(l).ok())
.filter(|v: &serde_json::Value| v["type"] == "inside_leg_completed")
.collect();
assert_eq!(
events.len(),
1,
"exactly one completion, only for the report-bearing row"
);
let ev = &events[0];
assert_eq!(ev["data"]["name"], "pane");
assert_eq!(ev["data"]["session_id"], "sess-uuid");
assert_eq!(ev["data"]["final_state"], "working");
assert_eq!(ev["data"]["seq"], 9);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn buffer_pending_report_highest_seq_wins_and_is_bounded() {
use std::collections::HashMap;
let rep = |seq| state::InsideLegReport {
state: state::InsideLegState::Working,
seq,
reason: None,
received_at: "2026-06-27T00:00:00Z".into(),
ttl_ms: None,
};
let mut map: HashMap<String, state::InsideLegReport> = HashMap::new();
assert!(matches!(
buffer_pending_report(&mut map, "s1", rep(2)),
BufferOutcome::Buffered
));
assert_eq!(map["s1"].seq, 2);
assert!(matches!(
buffer_pending_report(&mut map, "s1", rep(1)),
BufferOutcome::StaleSeq { last: 2 }
));
assert_eq!(
map["s1"].seq, 2,
"stale early push must not regress the buffer"
);
assert!(matches!(
buffer_pending_report(&mut map, "s1", rep(5)),
BufferOutcome::Buffered
));
assert_eq!(map["s1"].seq, 5);
for i in 0..PENDING_INSIDE_LEG_CAP {
buffer_pending_report(&mut map, &format!("fill{i}"), rep(1));
}
assert!(map.len() >= PENDING_INSIDE_LEG_CAP);
assert!(matches!(
buffer_pending_report(&mut map, "brand-new", rep(1)),
BufferOutcome::Full
));
assert!(!map.contains_key("brand-new"));
assert!(
matches!(
buffer_pending_report(&mut map, "s1", rep(9)),
BufferOutcome::Buffered
),
"an already-buffered session advances even at cap (no new key)"
);
}
#[test]
fn flush_buffered_inside_leg_drains_onto_row_under_seq_gate() {
let home = tmp_home("inside-leg-flush");
let ctx = test_ctx_with_events(home.clone(), PathBuf::from("fno-agents-worker"));
let report = |seq| state::InsideLegReport {
state: state::InsideLegState::Working,
seq,
reason: None,
received_at: "2026-06-27T00:00:00Z".into(),
ttl_ms: Some(5000),
};
let mut row = rentry("pane", AgentStatus::Live, None);
row.legacy_provider = "claude".into();
row.claude_session_uuid = Some("uuid-x".into());
state::update_registry(&home.registry_json(), |r| r.entries.push(row)).unwrap();
ctx.pending_inside_leg
.lock()
.unwrap()
.insert("uuid-x".into(), report(4));
flush_buffered_inside_leg(&ctx, "uuid-x", "pane");
assert!(!ctx
.pending_inside_leg
.lock()
.unwrap()
.contains_key("uuid-x"));
let reg = state::load_registry(&home.registry_json()).unwrap();
assert_eq!(reg.entries[0].inside_leg.as_ref().map(|r| r.seq), Some(4));
let events = read_events(&home);
assert!(events
.iter()
.any(|e| e["type"] == "inside_leg_buffer_flushed"
&& e["data"]["name"] == "pane"
&& e["data"]["session_id"] == "uuid-x"
&& e["data"]["seq"] == 4));
state::update_registry(&home.registry_json(), |r| {
r.entries[0].inside_leg = Some(report(10));
})
.unwrap();
ctx.pending_inside_leg
.lock()
.unwrap()
.insert("uuid-x".into(), report(7));
flush_buffered_inside_leg(&ctx, "uuid-x", "pane");
let reg = state::load_registry(&home.registry_json()).unwrap();
assert_eq!(
reg.entries[0].inside_leg.as_ref().map(|r| r.seq),
Some(10),
"a stale buffered report must not regress a newer row state"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn reconcile_defers_remaining_when_budget_exhausted() {
let entries = vec![
rentry("a", AgentStatus::Live, None),
rentry("b", AgentStatus::Live, None),
rentry("c", AgentStatus::Live, None),
];
let mut probes = 0;
let (changes, out) = plan_reconcile(
&entries,
|_| {
probes += 1;
Ok(true)
},
{
let mut checked = 0;
move || {
let exhausted = checked >= 1;
checked += 1;
exhausted
}
},
|_| true,
);
assert_eq!(out.deferred, 2, "two trailing entries should defer");
assert_eq!(changes.len(), 1, "only one entry probed before budget");
}
#[test]
fn run_reconcile_sweep_empty_registry_is_noop() {
let home = tmp_home("sweep-empty");
let emitter = EventEmitter::new(home.events_jsonl(), "daemon");
let result = run_reconcile_sweep(&home, &emitter).expect("empty sweep ok");
assert!(result.entries.is_empty());
assert_eq!(result.outcome, ReconcileOutcome::default());
std::fs::remove_dir_all(home.root()).ok();
}
struct PromptDetector;
impl crate::readiness::ReadinessDetector for PromptDetector {
fn provider_name(&self) -> &str {
"test-cli"
}
fn is_ready(
&self,
screen: &crate::readiness::ScreenView,
) -> Result<bool, crate::readiness::ReadinessError> {
Ok(screen.visible_text.trim_end().ends_with('\u{276f}'))
}
}
struct NeverReadyDetector;
impl crate::readiness::ReadinessDetector for NeverReadyDetector {
fn provider_name(&self) -> &str {
"never"
}
fn is_ready(
&self,
_screen: &crate::readiness::ScreenView,
) -> Result<bool, crate::readiness::ReadinessError> {
Ok(false)
}
}
#[tokio::test(flavor = "current_thread")]
async fn poll_until_ready_returns_settled_reply_on_ready_prompt() {
let snapshots: &[&str] = &["loading...", "still loading...", "done \u{276f}"];
let idx = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let idx2 = idx.clone();
let fetcher = move || {
let i = idx2.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let text = snapshots[i.min(snapshots.len() - 1)].to_string();
std::future::ready(Some(text))
};
let result = poll_until_ready(
fetcher,
Box::new(PromptDetector),
Duration::from_millis(1),
Duration::from_secs(5),
)
.await;
assert!(result.is_ok(), "expected Ok, got {result:?}");
let reply = result.unwrap();
assert_eq!(
reply, "done \u{276f}",
"reply must be the settled snapshot text, got {reply:?}"
);
}
#[tokio::test(flavor = "current_thread")]
async fn poll_until_ready_returns_error_on_timeout() {
let fetcher = || std::future::ready(Some("still thinking...".to_string()));
let result = poll_until_ready(
fetcher,
Box::new(NeverReadyDetector),
Duration::from_millis(10),
Duration::from_millis(40), )
.await;
assert!(
result.is_err(),
"expected Err on timeout, got Ok({:?})",
result.ok()
);
}
#[tokio::test(flavor = "current_thread")]
async fn poll_until_ready_empty_settled_screen_returns_empty_string() {
let fetcher = || std::future::ready(Some("\u{276f}".to_string()));
let result = poll_until_ready(
fetcher,
Box::new(PromptDetector),
Duration::from_millis(1),
Duration::from_secs(5),
)
.await;
assert!(result.is_ok(), "expected Ok, got {result:?}");
let reply = result.unwrap();
assert!(!reply.contains("fabricated"), "must not fabricate content");
}
#[test]
fn entry_holds_session_matches_claude_session_uuid() {
let row = build_claude_stream_entry(
"peer",
"ab12cd34",
std::path::Path::new("/work"),
"sess-uuid-9",
4242,
None,
PathBuf::from("/tmp/log.jsonl"),
);
assert!(
entry_holds_session(&row, "sess-uuid-9"),
"a claude row must be matched by its claude_session_uuid"
);
assert!(!entry_holds_session(&row, "other-uuid"));
}
fn test_ctx(home: AgentsHome, worker_bin: PathBuf) -> Ctx {
Ctx {
home,
emitter: EventEmitter::new(std::path::PathBuf::from("/dev/null"), "daemon"),
opts: DaemonOptions {
idle_exit: Duration::from_secs(1800),
worker_bin,
reconcile_on_start: true,
dead_row_grace: Duration::from_secs(3600),
notify_on_blocked: false,
notify_on_done: false,
},
started_at: std::time::Instant::now(),
exe_fingerprint: crate::drift::ExeFingerprint::current(),
pid_start_time: process_start_time(std::process::id()),
pending_inside_leg: std::sync::Mutex::new(std::collections::HashMap::new()),
}
}
fn test_ctx_with_events(home: AgentsHome, worker_bin: PathBuf) -> Ctx {
let events_path = home.events_jsonl();
Ctx {
home,
emitter: EventEmitter::new(events_path, "daemon"),
opts: DaemonOptions {
idle_exit: Duration::from_secs(1800),
worker_bin,
reconcile_on_start: true,
dead_row_grace: Duration::from_secs(3600),
notify_on_blocked: false,
notify_on_done: false,
},
started_at: std::time::Instant::now(),
exe_fingerprint: crate::drift::ExeFingerprint::current(),
pid_start_time: process_start_time(std::process::id()),
pending_inside_leg: std::sync::Mutex::new(std::collections::HashMap::new()),
}
}
const FAKE_STREAM_EMITTER: &str = r#"
printf '%s\n' '{"type":"system","subtype":"init","session_id":"s1"}'
while IFS= read -r line; do
printf '%s\n' '{"type":"user","message":{"role":"user"}}'
printf '%s\n' '{"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"par"}}}'
printf '%s\n' '{"type":"assistant","message":{"content":[{"type":"text","text":"reply-text"}]}}'
printf '%s\n' '{"type":"result","subtype":"success","is_error":false,"result":"reply-text"}'
done
"#;
fn short_home(tag: &str) -> AgentsHome {
use std::sync::atomic::{AtomicU32, Ordering};
static C: AtomicU32 = AtomicU32::new(0);
let n = C.fetch_add(1, Ordering::Relaxed);
let p = PathBuf::from(format!("/tmp/fnosb{tag}{}_{n}", std::process::id()));
let home = AgentsHome::at(&p);
home.ensure_root().unwrap();
home
}
fn seed_stream_row(home: &AgentsHome, name: &str, short_id: &str) {
state::update_registry(&home.registry_json(), |r| {
r.entries.push(RegistryEntry {
name: name.into(),
short_id: short_id.into(),
legacy_provider: "claude".into(),
harness: None,
harness_session_id: None,
cwd: "/tmp".into(),
project_root: "/tmp".into(),
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: Some(format!("uuid-{short_id}")),
messaging_socket_path: None,
codex_session_id: None,
gemini_session_id: None,
mcp_channel_id: None,
cc_session_id: None,
host_mode: None,
status: AgentStatus::Live,
last_message_at: None,
created_at: "2026-06-09T00:00:00Z".into(),
pid: None,
pid_start_time: None,
log_path: None,
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
});
})
.unwrap();
}
#[test]
fn list_row_key_set_matches_shared_contract() {
const CONTRACT: &str = include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/../../schemas/agents-list-row.json"
));
let contract: Value = serde_json::from_str(CONTRACT).expect("contract is valid JSON");
let mut expected: std::collections::BTreeSet<String> = contract["required"]
.as_array()
.expect("required is an array")
.iter()
.map(|k| k.as_str().unwrap().to_string())
.collect();
expected.extend(
contract["rust_only"]["keys"]
.as_array()
.expect("rust_only.keys is an array")
.iter()
.map(|k| k.as_str().unwrap().to_string()),
);
let home = short_home("listcontract");
seed_stream_row(&home, "worker-contract", "abc12345");
state::update_registry(&home.registry_json(), |r| {
let e = &mut r.entries[0];
e.short_id = String::new();
e.harness = Some("claude".into());
e.harness_session_id = Some("e6f78b98-e594-47ed-ad81-84f8a78b8bb7".into());
e.claude_session_uuid = Some("e6f78b98-e594-47ed-ad81-84f8a78b8bb7".into());
e.mux = Some(crate::state::MuxRef {
session: "main".into(),
pane_id: 10,
});
e.crown_level = Some(1);
e.crown_scope = Some("epic-x".into());
e.crown_grantor = Some("king".into());
})
.unwrap();
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({}));
let response = handle_list_with_truth(&ctx, &req, |_handle| Some("working".into()));
let result = response.result().unwrap();
let row = &result["agents"][0];
let actual: std::collections::BTreeSet<String> =
row.as_object().unwrap().keys().cloned().collect();
assert_eq!(actual, expected, "list row key set drifted from contract");
assert_eq!(row["harness"], "claude");
assert_eq!(row["provider"], "claude", "legacy alias still emitted");
assert_eq!(
row["harness_session_id"],
"e6f78b98-e594-47ed-ad81-84f8a78b8bb7"
);
assert!(row["session_id"].is_null());
assert_eq!(row["mux"]["session"], "main");
assert_eq!(row["mux"]["pane_id"], 10);
assert_eq!(
row["crown"], "L1 epic-x",
"same formatter as Python crown_label"
);
assert_eq!(row["crown_level"], 1);
assert_eq!(row["crown_scope"], "epic-x");
assert_eq!(row["crown_grantor"], "king");
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn list_row_resolves_opencode_session_id_from_harness_session_id() {
let home = short_home("listopencode");
seed_stream_row(&home, "worker-opencode", "abc12345");
state::update_registry(&home.registry_json(), |r| {
let e = &mut r.entries[0];
e.harness = Some("opencode".into());
e.harness_session_id = Some("oc-sess-9f2".into());
e.session_id = None;
})
.unwrap();
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({}));
let response = handle_list_with_truth(&ctx, &req, |_handle| Some("working".into()));
let result = response.result().unwrap();
let row = &result["agents"][0];
assert_eq!(row["harness"], "opencode");
assert_eq!(row["session_id"], "oc-sess-9f2");
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn list_row_crown_label_falls_back_on_an_empty_scope() {
let home = short_home("listcrownempty");
seed_stream_row(&home, "worker-crown", "abc12345");
state::update_registry(&home.registry_json(), |r| {
let e = &mut r.entries[0];
e.crown_level = Some(1);
e.crown_scope = Some(String::new());
})
.unwrap();
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({}));
let response = handle_list_with_truth(&ctx, &req, |_handle| Some("working".into()));
let result = response.result().unwrap();
assert_eq!(result["agents"][0]["crown"], "L1 ?");
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn list_row_emits_absent_optional_fields_as_null() {
let home = short_home("listnulls");
seed_stream_row(&home, "worker-bare", "abc12345");
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({}));
let response = handle_list_with_truth(&ctx, &req, |_handle| Some("working".into()));
let result = response.result().unwrap();
let row = &result["agents"][0];
let obj = row.as_object().unwrap();
for key in ["mux", "crown", "crown_level"] {
assert!(obj.contains_key(key), "row omits key: {key}");
assert!(obj[key].is_null(), "key {key} should be null on a bare row");
}
std::fs::remove_dir_all(home.root()).ok();
}
fn stream_identity(short_id: &str) -> Value {
json!({
"harness": "claude",
"session_id": format!("uuid-{short_id}"),
"short_id": short_id,
"created_at": "2026-06-09T00:00:00Z",
})
}
fn switchboard_params(
to: &str,
to_short: &str,
from: &str,
from_short: Option<&str>,
body: &str,
) -> Value {
let mut params = json!({
"to": to,
"from": from,
"body": body,
"mirror": from_short.is_some(),
"recipient_identity": stream_identity(to_short),
});
if let Some(short_id) = from_short {
params["from_identity"] = stream_identity(short_id);
}
params
}
#[test]
fn list_renders_family1_truth_instead_of_stored_registry_status() {
let home = short_home("listtruth");
seed_stream_row(&home, "worker-list", "abc12345");
state::update_registry(&home.registry_json(), |registry| {
registry.entries[0].status = AgentStatus::Orphaned;
})
.unwrap();
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({"status": "live"}));
let response = handle_list_with_truth(&ctx, &req, |_handle| Some("working".into()));
let result = response.result().unwrap();
let agents = result["agents"].as_array().unwrap();
assert_eq!(agents.len(), 1);
assert_eq!(agents[0]["status"], "live");
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn list_queries_family1_by_session_identity_not_custom_name() {
let home = short_home("listidentity");
seed_stream_row(&home, "custom-worker-name", "abc12345");
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({"status": "live"}));
let seen = std::cell::RefCell::new(Vec::new());
let response = handle_list_with_truth(&ctx, &req, |handle| {
seen.borrow_mut().push(handle.to_string());
Some("working".into())
});
assert!(response.result().is_some());
assert_eq!(seen.into_inner(), vec!["uuid-abc12345"]);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn list_queries_pidless_row_by_bare_canonical_handle() {
let home = short_home("listpidless");
seed_stream_row(&home, "custom-worker-name", "unused");
state::update_registry(&home.registry_json(), |registry| {
registry.entries[0].short_id.clear();
registry.entries[0].harness = Some("codex".into());
registry.entries[0].harness_session_id =
Some("019f8ff2-1111-2222-3333-444444444444".into());
})
.unwrap();
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({"status": "live"}));
let seen = std::cell::RefCell::new(Vec::new());
let response = handle_list_with_truth(&ctx, &req, |handle| {
seen.borrow_mut().push(handle.to_string());
Some("working".into())
});
assert!(response.result().is_some());
assert_eq!(
seen.into_inner(),
vec!["019f8ff2-1111-2222-3333-444444444444"]
);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn list_queries_non_claude_row_by_transcript_identity() {
let home = short_home("listnonclaude");
seed_stream_row(&home, "custom-worker-name", "transport");
state::update_registry(&home.registry_json(), |registry| {
registry.entries[0].harness = Some("codex".into());
registry.entries[0].harness_session_id =
Some("019f8ff2-1111-2222-3333-444444444444".into());
})
.unwrap();
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({"status": "live"}));
let seen = std::cell::RefCell::new(Vec::new());
let response = handle_list_with_truth(&ctx, &req, |handle| {
seen.borrow_mut().push(handle.to_string());
Some("working".into())
});
assert!(response.result().is_some());
assert_eq!(
seen.into_inner(),
vec!["019f8ff2-1111-2222-3333-444444444444"]
);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn list_applies_cheap_filters_before_family1_subprocesses() {
let home = short_home("listprefilter");
seed_stream_row(&home, "claude-worker", "aaaaaaaa");
seed_stream_row(&home, "codex-worker", "bbbbbbbb");
state::update_registry(&home.registry_json(), |registry| {
registry.entries[1].harness = Some("codex".into());
registry.entries[1].harness_session_id =
Some("bbbbbbbb-1111-2222-3333-444444444444".into());
})
.unwrap();
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(1, "agent.list", json!({"provider": "codex"}));
let seen = std::cell::RefCell::new(Vec::new());
let response = handle_list_with_truth(&ctx, &req, |handle| {
seen.borrow_mut().push(handle.to_string());
Some("working".into())
});
assert!(response.result().is_some());
assert_eq!(
seen.into_inner(),
vec!["bbbbbbbb-1111-2222-3333-444444444444"]
);
std::fs::remove_dir_all(home.root()).ok();
}
async fn start_stream_worker(home: &AgentsHome, short_id: &str, script: &str) -> PathBuf {
let cfg = crate::stream_worker::StreamWorkerConfig::new(
short_id,
home.root().to_path_buf(),
std::env::temp_dir(),
vec!["bash".into(), "-c".into(), script.into()],
);
let short_id_dbg = short_id.to_string();
std::thread::spawn(move || {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
if let Err(e) = crate::stream_worker::run(cfg).await {
eprintln!("STREAM WORKER RUN ERROR ({short_id_dbg}): {e}");
}
});
});
let sock = home.worker_sock(short_id);
let bind_start = std::time::Instant::now();
while !sock.exists() && bind_start.elapsed() < Duration::from_secs(20) {
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(
sock.exists(),
"stream worker socket never appeared for {short_id}"
);
let live_start = std::time::Instant::now();
let mut live = false;
while live_start.elapsed() < Duration::from_secs(30) {
if is_live_stream_thread(&sock).await {
live = true;
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
assert!(
live,
"stream worker socket appeared but never answered a ping for {short_id}"
);
sock
}
fn built_worker_bin() -> Option<PathBuf> {
let exe = std::env::current_exe().ok()?;
let dir = exe.parent()?.parent()?; let cand = dir.join("fno-agents-worker");
cand.exists().then_some(cand)
}
#[test]
fn stream_claim_holder_is_short_id_scoped() {
assert_eq!(stream_claim_holder("sw7"), "stream:sw7");
}
#[test]
fn provider_readiness_detector_handles_claude() {
let d = provider_readiness_detector("claude");
assert_eq!(d.provider_name(), "claude");
assert_eq!(
provider_readiness_detector("aider").provider_name(),
"aider"
);
}
#[test]
fn is_live_writer_excludes_orphaned_and_terminal() {
for s in [
AgentStatus::Live,
AgentStatus::Ready,
AgentStatus::Idle,
AgentStatus::Busy,
AgentStatus::Spawning,
AgentStatus::Restarting,
] {
assert!(is_live_writer(s), "{s:?} should count as a live writer");
}
for s in [
AgentStatus::Orphaned,
AgentStatus::Failed,
AgentStatus::Exited,
AgentStatus::PermanentDead,
] {
assert!(!is_live_writer(s), "{s:?} must NOT block re-adoption");
}
}
#[test]
fn acquire_session_claim_maps_native_outcomes() {
let td = tempfile::tempdir().unwrap();
let _guard = crate::claims::test_env_lock()
.lock()
.unwrap_or_else(|e| e.into_inner());
std::env::set_var("FNO_CLAIMS_ROOT", td.path());
assert!(matches!(
acquire_session_claim("U-1", "stream:sw1"),
ClaimOutcome::Acquired
));
assert!(matches!(
acquire_session_claim("U-1", "stream:sw1"),
ClaimOutcome::Acquired
));
match acquire_session_claim("U-1", "stream:other") {
ClaimOutcome::HeldByOther(who) => assert_eq!(who, "stream:sw1"),
other => panic!("expected HeldByOther, got {other:?}"),
}
std::env::remove_var("FNO_CLAIMS_ROOT");
}
#[test]
fn claude_stream_worker_args_carry_stream_flags_and_child_argv() {
let child = crate::provider::claude_stream_json_resume_argv("U-9");
let args = claude_stream_worker_args(
"sw9",
std::path::Path::new("/home/agents"),
std::path::Path::new("/work"),
"U-9",
"stream:sw9",
&child,
);
assert!(args.contains(&"--stream".to_string()));
assert_eq!(
args.iter()
.position(|a| a == "--session-uuid")
.map(|i| &args[i + 1]),
Some(&"U-9".to_string())
);
assert_eq!(
args.iter()
.position(|a| a == "--holder")
.map(|i| &args[i + 1]),
Some(&"stream:sw9".to_string())
);
let sep = args
.iter()
.position(|a| a == "--")
.expect("missing -- separator");
assert_eq!(&args[sep + 1..], child.as_slice());
assert_eq!(child[0], "claude");
assert!(child.contains(&"--resume".to_string()) && child.contains(&"U-9".to_string()));
}
#[test]
fn build_claude_stream_entry_marks_interactive_claude_with_full_uuid() {
let e = build_claude_stream_entry(
"adopted",
"sw3",
std::path::Path::new("/proj"),
"FULL-UUID-3",
4242,
Some(99),
PathBuf::from("/proj/.fno/agents/sw3/timeline.jsonl"),
);
assert_eq!(e.harness_name(), "claude");
assert_eq!(
e.host_mode.as_deref(),
Some(crate::state::HOST_MODE_INTERACTIVE)
);
assert!(
e.is_interactive(),
"stream thread must read as interactive for reconcile"
);
assert_eq!(e.claude_session_uuid.as_deref(), Some("FULL-UUID-3"));
assert_eq!(e.status, AgentStatus::Live);
assert_eq!(e.pid, Some(4242));
assert_eq!(e.short_id, "sw3");
}
#[tokio::test(flavor = "current_thread")]
async fn host_claude_without_from_rejected_with_adopt_pointer() {
let home = short_home("clnofrom");
let ctx = test_ctx(home.clone(), PathBuf::from("/nonexistent-worker"));
let req = Request::new(
1,
"agent.spawn",
json!({"name": "cl", "provider": "claude", "host_mode": "interactive"}),
);
let resp = handle_spawn(&ctx, &req).await;
match &resp.payload {
crate::protocol::ResponsePayload::Err(e) => {
assert_eq!(e.code, ErrorCode::InvalidParams);
assert!(
e.message.contains("promote") && e.message.contains("--from"),
"claude host without --from must point at the adopt verb; got: {}",
e.message
);
}
_ => panic!("expected error for claude host without --from"),
}
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn promote_claude_duplicate_session_refused() {
let home = short_home("cldup");
seed_stream_row(&home, "first", "swDup"); let ctx = test_ctx(home.clone(), PathBuf::from("/nonexistent-worker"));
let req = Request::new(
1,
"agent.spawn",
json!({
"name": "second", "provider": "claude", "host_mode": "interactive",
"resume_id": "uuid-swDup"
}),
);
let resp = handle_spawn(&ctx, &req).await;
match &resp.payload {
crate::protocol::ResponsePayload::Err(e) => {
assert_eq!(e.code, ErrorCode::InvalidParams);
assert!(
e.message.contains("already hosted") && e.message.contains("first"),
"duplicate adopt must name the existing host; got: {}",
e.message
);
}
_ => panic!("expected single-writer refusal for duplicate adopt"),
}
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn promote_claude_spawns_live_stream_thread() {
let Some(worker_bin) = built_worker_bin() else {
eprintln!("skip promote_claude_spawns_live_stream_thread: worker bin not built");
return;
};
let _guard = crate::claims::test_env_lock()
.lock()
.unwrap_or_else(|e| e.into_inner());
let home = short_home("cle2e");
std::env::set_var("FNO_CLAIMS_ROOT", home.root());
let ctx = test_ctx(home.clone(), worker_bin);
let req = Request::new(
1,
"agent.spawn",
json!({
"name": "cl", "provider": "claude", "host_mode": "interactive",
"resume_id": "uuid-e2e", "cwd": "/tmp",
"argv": ["bash", "-c", FAKE_STREAM_EMITTER]
}),
);
let resp = handle_spawn(&ctx, &req).await;
let res = resp.result().expect("claude adopt errored");
assert_eq!(res["provider"], "claude");
assert_eq!(res["status"], "live");
assert_eq!(res["lane"], "stream");
let reg = load_registry_offloaded(home.registry_json()).await;
let row = reg.find("cl").expect("adopted row missing");
assert_eq!(row.harness_name(), "claude");
assert_eq!(row.host_mode.as_deref(), Some("interactive"));
assert_eq!(row.claude_session_uuid.as_deref(), Some("uuid-e2e"));
assert_eq!(row.status, AgentStatus::Live);
let sock = home.worker_sock(&row.short_id);
assert!(
is_live_stream_thread(&sock).await,
"adopted thread must serve the stream protocol"
);
best_effort_worker_shutdown(&sock).await;
std::fs::remove_dir_all(home.root()).ok();
std::env::remove_var("FNO_CLAIMS_ROOT");
}
#[tokio::test(flavor = "current_thread")]
async fn promote_claude_dead_on_arrival_resume_rejected() {
let Some(worker_bin) = built_worker_bin() else {
eprintln!("skip promote_claude_dead_on_arrival_resume_rejected: worker bin not built");
return;
};
let _guard = crate::claims::test_env_lock()
.lock()
.unwrap_or_else(|e| e.into_inner());
let home = short_home("cldoa");
std::env::set_var("FNO_CLAIMS_ROOT", home.root());
let ctx = test_ctx(home.clone(), worker_bin);
let req = Request::new(
1,
"agent.spawn",
json!({
"name": "cl", "provider": "claude", "host_mode": "interactive",
"resume_id": "uuid-doa", "cwd": "/tmp",
"argv": ["bash", "-c", "exit 1"]
}),
);
let resp = handle_spawn(&ctx, &req).await;
assert!(
resp.is_err(),
"DOA resume child must be rejected, not registered"
);
assert_eq!(resp.error().unwrap().code, ErrorCode::SpawnFailed);
let reg = load_registry_offloaded(home.registry_json()).await;
assert!(
reg.find("cl").is_none(),
"no row may be registered for a DOA adopt"
);
std::fs::remove_dir_all(home.root()).ok();
std::env::remove_var("FNO_CLAIMS_ROOT");
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_drives_b_and_mirrors_into_a() {
let home = short_home("hp");
seed_stream_row(&home, "A", "swA");
seed_stream_row(&home, "B", "swB");
let _a = start_stream_worker(&home, "swA", FAKE_STREAM_EMITTER).await;
let _b = start_stream_worker(&home, "swB", FAKE_STREAM_EMITTER).await;
let ctx = test_ctx_with_events(home.clone(), PathBuf::from("/nonexistent-worker"));
let req = Request::new(
1,
"agent.switchboard",
switchboard_params("B", "swB", "A", Some("swA"), "hello"),
);
let resp = handle_switchboard(&ctx, &req).await;
let res = resp.result().expect("switchboard errored");
assert_eq!(res["delivered"], true, "not delivered: {res:?}");
assert_eq!(res["reply"], "reply-text");
assert_eq!(res["is_error"], false);
assert_eq!(res["receipt"], true, "user-echo receipt not observed");
assert_eq!(res["mirrored"], true, "B's reply was not mirrored into A");
assert_eq!(res["identity_verified"], true);
let events = read_events(&home);
assert!(
events.iter().any(|e| e["type"] == "agent_deliver_injected"
&& e["data"]["transport"] == "switchboard"
&& e["data"]["mirrored"] == true),
"switchboard injected event missing: {events:?}"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_second_turn_returns_fresh_reply() {
const COUNTING_EMITTER: &str = r#"
printf '%s\n' '{"type":"system","subtype":"init","session_id":"s1"}'
n=0
while IFS= read -r line; do
n=$((n+1))
printf '%s\n' '{"type":"user","message":{"role":"user"}}'
printf '%s\n' "{\"type\":\"assistant\",\"message\":{\"content\":[{\"type\":\"text\",\"text\":\"reply-$n\"}]}}"
printf '%s\n' "{\"type\":\"result\",\"subtype\":\"success\",\"is_error\":false,\"result\":\"reply-$n\"}"
done
"#;
let home = short_home("fresh");
seed_stream_row(&home, "B", "swB");
let _b = start_stream_worker(&home, "swB", COUNTING_EMITTER).await;
let ctx = test_ctx(home.clone(), PathBuf::from("/nonexistent-worker"));
let r1 = handle_switchboard(
&ctx,
&Request::new(
1,
"agent.switchboard",
switchboard_params("B", "swB", "ghost", None, "first"),
),
)
.await;
assert_eq!(r1.result().expect("hop1")["reply"], "reply-1");
let r2 = handle_switchboard(
&ctx,
&Request::new(
2,
"agent.switchboard",
switchboard_params("B", "swB", "ghost", None, "second"),
),
)
.await;
assert_eq!(
r2.result().expect("hop2")["reply"],
"reply-2",
"second drive returned a STALE reply (cursor not advanced past the prior turn)"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_demotes_when_b_not_a_live_stream_thread() {
let home = short_home("demote");
seed_stream_row(&home, "B", "swB"); let ctx = test_ctx(home.clone(), PathBuf::from("/nonexistent-worker"));
let req = Request::new(
1,
"agent.switchboard",
switchboard_params("B", "swB", "A", None, "hi"),
);
let resp = handle_switchboard(&ctx, &req).await;
let res = resp
.result()
.expect("should be Ok-demote, not an RPC error");
assert_eq!(res["delivered"], false);
assert_eq!(res["reason"], "not-a-live-stream-thread");
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_one_way_when_peer_absent() {
let home = short_home("oneway");
seed_stream_row(&home, "B", "swB");
let _b = start_stream_worker(&home, "swB", FAKE_STREAM_EMITTER).await;
let ctx = test_ctx(home.clone(), PathBuf::from("/nonexistent-worker"));
let req = Request::new(
1,
"agent.switchboard",
switchboard_params("B", "swB", "ghost", None, "hi"),
);
let resp = handle_switchboard(&ctx, &req).await;
let res = resp.result().expect("switchboard errored");
assert_eq!(res["delivered"], true);
assert_eq!(res["reply"], "reply-text");
assert_eq!(res["mirrored"], false, "no peer to mirror into");
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_unknown_target_is_not_found() {
let home = short_home("404");
let ctx = test_ctx(home.clone(), PathBuf::from("/nonexistent-worker"));
let req = Request::new(
1,
"agent.switchboard",
switchboard_params("nope", "swNope", "A", None, "hi"),
);
let resp = handle_switchboard(&ctx, &req).await;
if let crate::protocol::ResponsePayload::Err(ref e) = resp.payload {
assert_eq!(e.code, ErrorCode::AgentNotFound);
} else {
panic!("expected AgentNotFound, got {resp:?}");
}
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_refuses_replaced_recipient_identity() {
let home = short_home("replaced");
seed_stream_row(&home, "victim", "swB");
let ctx = test_ctx(home.clone(), PathBuf::from("/nonexistent-worker"));
let mut params = switchboard_params("victim", "swB", "ghost", None, "secret");
params["recipient_identity"]["session_id"] = json!("uuid-swA");
let response =
handle_switchboard(&ctx, &Request::new(1, "agent.switchboard", params)).await;
let result = response.result().expect("identity mismatch is a demotion");
assert_eq!(result["delivered"], false);
assert_eq!(result["reason"], "recipient-identity-changed");
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_failed_drive_does_not_orphan_restamped_recipient() {
let home = short_home("failedrestamp");
seed_stream_row(&home, "B", "swB");
let turn_started = home.root().join("turn-started");
let restamp_done = home.root().join("restamp-done");
let script = format!(
r#"
printf '%s\n' '{{"type":"system","subtype":"init","session_id":"s1"}}'
while IFS= read -r line; do
touch '{}'
while [ ! -f '{}' ]; do sleep 0.01; done
exit 1
done
"#,
turn_started.display(),
restamp_done.display()
);
let _b = start_stream_worker(&home, "swB", &script).await;
let ctx = test_ctx_with_events(home.clone(), PathBuf::from("/nonexistent-worker"));
let registry_path = home.registry_json();
let restamp_signal = turn_started.clone();
let restamp_complete = restamp_done.clone();
let restamp = tokio::spawn(async move {
let start = Instant::now();
while !restamp_signal.exists() && start.elapsed() < Duration::from_secs(5) {
tokio::time::sleep(Duration::from_millis(5)).await;
}
assert!(restamp_signal.exists(), "drive never reached the worker");
state::update_registry(®istry_path, |registry| {
let row = registry.find_mut("B").expect("recipient row missing");
row.short_id = "swC".into();
row.harness_session_id = Some("uuid-replacement".into());
row.claude_session_uuid = Some("uuid-replacement".into());
row.created_at = "2026-06-09T00:00:01Z".into();
row.status = AgentStatus::Live;
})
.unwrap();
std::fs::write(restamp_complete, b"done\n").unwrap();
});
let response = handle_switchboard(
&ctx,
&Request::new(
1,
"agent.switchboard",
switchboard_params("B", "swB", "ghost", None, "fail after receipt"),
),
)
.await;
restamp.await.unwrap();
let result = response.result().expect("failed drive is a demotion");
assert_eq!(result["delivered"], false);
let registry = state::load_registry(&home.registry_json()).unwrap();
let replacement = registry.find("B").expect("replacement row missing");
assert_eq!(replacement.status, AgentStatus::Live);
assert_eq!(
replacement.harness_session_id.as_deref(),
Some("uuid-replacement")
);
let events = read_events(&home);
assert!(
events.iter().any(|event| {
event["type"] == "agent_deliver_status_write_failed"
&& event["data"]["name"] == "B"
&& event["data"]["reason"] == "recipient-identity-changed"
}),
"identity-CAS failure event missing: {events:?}"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_requires_recipient_identity() {
let home = short_home("identity-required");
let ctx = test_ctx(home.clone(), PathBuf::from("/nonexistent-worker"));
let response = handle_switchboard(
&ctx,
&Request::new(
1,
"agent.switchboard",
json!({"to": "victim", "from": "ghost", "body": "secret", "mirror": false}),
),
)
.await;
assert_eq!(
response.error().expect("missing identity must fail").code,
ErrorCode::InvalidParams
);
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn switchboard_v2_routes_to_identity_guard() {
let home = short_home("v2route");
let ctx = Arc::new(test_ctx(home.clone(), PathBuf::from("/nonexistent-worker")));
let response = dispatch_agent(
&ctx,
&Request::new(
1,
"agent.switchboard_v2",
json!({"to": "victim", "from": "ghost", "body": "secret"}),
),
)
.await;
let error = response.error().expect("missing identity must fail");
assert_eq!(error.code, ErrorCode::InvalidParams);
assert!(error.message.contains("recipient_identity"));
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn handle_spawn_codex_pty_hosting_retired_returns_pointer() {
let home = tmp_home("spawn-provider-argv");
let ctx = test_ctx(
home.clone(),
PathBuf::from("/nonexistent/fno-agents-worker"),
);
let req = Request::new(
1,
"agent.spawn",
json!({"name": "test-agent", "provider": "codex"}),
);
let resp = handle_spawn(&ctx, &req).await;
match &resp.payload {
crate::protocol::ResponsePayload::Err(e) => {
assert_eq!(e.code, ErrorCode::InvalidParams);
assert!(
e.message.contains("retired at G4"),
"codex spawn must point at the mux; got: {}",
e.message
);
}
_ => panic!("expected the G4 retirement error for a codex PTY spawn"),
}
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn handle_spawn_unknown_provider_no_argv_returns_invalid_params() {
let home = tmp_home("spawn-unknown-provider");
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(
1,
"agent.spawn",
json!({"name": "test-agent", "provider": "nonexistent-provider"}),
);
let resp = handle_spawn(&ctx, &req).await;
match &resp.payload {
crate::protocol::ResponsePayload::Err(e) => {
assert_eq!(
e.code,
ErrorCode::InvalidParams,
"unknown provider without argv must return InvalidParams"
);
}
_ => panic!("expected error response for unknown provider"),
}
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn handle_ask_first_contact_with_provider_routes_into_spawn() {
let home = tmp_home("ask-first-contact");
let ctx = test_ctx(
home.clone(),
PathBuf::from("/nonexistent/fno-agents-worker"),
);
let req = Request::new(
1,
"agent.ask",
json!({"name": "new-agent", "message": "hello", "provider": "codex"}),
);
let resp = handle_ask(&ctx, &req).await;
match &resp.payload {
crate::protocol::ResponsePayload::Err(e) => {
assert_ne!(
e.code,
ErrorCode::AgentNotFound,
"first-contact ask with --provider must route into the spawn branch, not short-circuit AgentNotFound; got: {}",
e.message
);
assert!(
e.message.contains("retired at G4"),
"first-contact codex spawn must surface the G4 mux pointer; got: {}",
e.message
);
}
crate::protocol::ResponsePayload::Ok(v) => {
panic!("post-G4 a codex first-contact spawn must fail, got Ok: {v}")
}
}
std::fs::remove_dir_all(home.root()).ok();
}
#[tokio::test(flavor = "current_thread")]
async fn handle_ask_first_contact_without_provider_returns_invalid_params() {
let home = tmp_home("ask-no-provider");
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(
1,
"agent.ask",
json!({"name": "ghost-agent", "message": "hello"}),
);
let resp = handle_ask(&ctx, &req).await;
match &resp.payload {
crate::protocol::ResponsePayload::Err(e) => {
assert_eq!(
e.code,
ErrorCode::InvalidParams,
"first-contact ask without --provider must return InvalidParams; got: {}",
e.message
);
assert!(
e.message.contains("provider"),
"error message must mention 'provider', got: {}",
e.message
);
}
_ => panic!("expected error for first-contact ask without provider"),
}
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn handle_report_stores_on_matching_row() {
let home = tmp_home("report-store");
seed_stream_row(&home, "worker-A", "repA"); let ctx = test_ctx_with_events(home.clone(), PathBuf::from("fno-agents-worker"));
let req = Request::new(
1,
"agent.report",
json!({"session_id": "uuid-repA", "seq": 3, "state": "working", "reason": "running tests"}),
);
let resp = handle_report(&ctx, &req);
assert!(!resp.is_err(), "report must return Ok: {resp:?}");
assert_eq!(resp.result().unwrap()["stored"], true);
let reg = state::load_registry(&home.registry_json()).unwrap();
let rep = reg.entries[0]
.inside_leg
.as_ref()
.expect("inside_leg stored");
assert_eq!(rep.state, state::InsideLegState::Working);
assert_eq!(rep.seq, 3);
assert_eq!(rep.reason.as_deref(), Some("running tests"));
assert!(!rep.received_at.is_empty(), "daemon stamps received_at");
let events = read_events(&home);
assert!(
events.iter().any(|e| e["type"] == "inside_leg_report"),
"inside_leg_report not emitted: {events:?}"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn handle_report_capability_flip_clears_screen_state() {
let home = tmp_home("report-flip-clears-scrape");
seed_stream_row(&home, "worker-A", "repF");
state::update_registry(&home.registry_json(), |r| {
r.entries[0].screen_state = Some(state::ScreenStateReport {
state: "idle".into(),
rule: "idle_prompt".into(),
seq: 4,
at: "2026-07-02T00:00:00Z".into(),
ttl_ms: Some(120_000),
answerable: None,
});
})
.unwrap();
let ctx = test_ctx_with_events(home.clone(), PathBuf::from("fno-agents-worker"));
let resp = handle_report(
&ctx,
&Request::new(
1,
"agent.report",
json!({"session_id": "uuid-repF", "seq": 1, "state": "working"}),
),
);
assert_eq!(resp.result().unwrap()["stored"], true);
let reg = state::load_registry(&home.registry_json()).unwrap();
assert!(reg.entries[0].inside_leg.is_some());
assert_eq!(
reg.entries[0].screen_state, None,
"capability flip must clear the scrape verdict"
);
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn handle_report_drops_stale_seq() {
let home = tmp_home("report-stale");
seed_stream_row(&home, "worker-A", "repB");
let ctx = test_ctx_with_events(home.clone(), PathBuf::from("fno-agents-worker"));
let _ = handle_report(
&ctx,
&Request::new(
1,
"agent.report",
json!({"session_id": "uuid-repB", "seq": 2, "state": "working"}),
),
);
let resp = handle_report(
&ctx,
&Request::new(
2,
"agent.report",
json!({"session_id": "uuid-repB", "seq": 1, "state": "done"}),
),
);
assert!(!resp.is_err());
assert_eq!(resp.result().unwrap()["stored"], false);
assert_eq!(resp.result().unwrap()["dropped"], "stale_seq");
let reg = state::load_registry(&home.registry_json()).unwrap();
let rep = reg.entries[0].inside_leg.as_ref().unwrap();
assert_eq!(rep.seq, 2);
assert_eq!(rep.state, state::InsideLegState::Working);
let events = read_events(&home);
assert!(events.iter().any(
|e| e["type"] == "inside_leg_report_dropped" && e["data"]["reason"] == "stale_seq"
));
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn handle_report_buffers_early_push_for_unknown_session() {
let home = tmp_home("report-unknown");
let ctx = test_ctx_with_events(home.clone(), PathBuf::from("fno-agents-worker"));
let resp = handle_report(
&ctx,
&Request::new(
1,
"agent.report",
json!({"session_id": "uuid-nope", "seq": 1, "state": "working"}),
),
);
assert!(!resp.is_err());
assert_eq!(resp.result().unwrap()["stored"], false);
assert_eq!(
resp.result().unwrap()["buffered"],
true,
"an early push is held, not dropped (E3.3)"
);
let reg = state::load_registry(&home.registry_json()).unwrap();
assert!(reg.entries.is_empty(), "no phantom row created");
assert_eq!(
ctx.pending_inside_leg
.lock()
.unwrap()
.get("uuid-nope")
.map(|r| r.seq),
Some(1)
);
let events = read_events(&home);
assert!(events
.iter()
.any(|e| e["type"] == "inside_leg_report_buffered"
&& e["data"]["session_id"] == "uuid-nope"));
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn handle_report_rejects_bad_params() {
let home = tmp_home("report-bad");
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
for params in [
json!({"seq": 1, "state": "working"}), json!({"session_id": "x", "state": "working"}), json!({"session_id": "x", "seq": 1}), json!({"session_id": "x", "seq": 1, "state": "idle"}), ] {
let resp = handle_report(&ctx, &Request::new(1, "agent.report", params.clone()));
assert!(resp.is_err(), "expected InvalidParams for {params}");
}
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn push_to_channel_rejects_non_object_envelope() {
let home = tmp_home("push-badenv");
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let resp = handle_push_to_channel(
&ctx,
&Request::new(
1,
"channel.push_to_channel",
json!({"mcp_channel_id": "c1", "envelope": "not-an-object"}),
),
);
assert!(resp.is_err(), "non-object envelope must be InvalidParams");
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn push_to_channel_unknown_channel_errors() {
let home = tmp_home("push-unknown");
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let resp = handle_push_to_channel(
&ctx,
&Request::new(
1,
"channel.push_to_channel",
json!({"mcp_channel_id": "nope", "envelope": {"a": 1}}),
),
);
assert!(resp.is_err(), "unknown channel must error");
std::fs::remove_dir_all(home.root()).ok();
}
#[test]
fn push_to_channel_no_envelope_is_confirm_only() {
let home = tmp_home("push-confirm");
seed_stream_row(&home, "worker-c", "chA");
state::update_registry(&home.registry_json(), |r| {
r.entries[0].mcp_channel_id = Some("c1".into());
})
.unwrap();
let ctx = test_ctx(home.clone(), PathBuf::from("fno-agents-worker"));
let resp = handle_push_to_channel(
&ctx,
&Request::new(
1,
"channel.push_to_channel",
json!({"mcp_channel_id": "c1"}),
),
);
let result = resp.result().unwrap();
assert_eq!(result["routed"], true);
assert!(
result.get("delivered").is_none(),
"confirm-only must not claim delivery"
);
std::fs::remove_dir_all(home.root()).ok();
}
}