use std::io::{BufRead, BufReader, Write};
use std::os::unix::process::CommandExt;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use crate::claude_ask::emit_event;
use crate::paths::AgentsHome;
use crate::provider::normalize_codex_command;
use crate::state::{load_registry, update_registry};
use crate::AgentStatus;
const EV_SESSION: &str = "thread.started";
const EV_COMPLETE: &str = "turn.completed";
const EV_ITEM: &str = "item.completed";
const ITEM_MESSAGE: &str = "agent_message";
const ITEM_ERROR: &str = "error";
const LOCK_ACQUIRE_TIMEOUT: Duration = Duration::from_secs(30);
const DEFAULT_FOLLOWUP_TIMEOUT: Duration = Duration::from_secs(600);
pub fn inject_from_name(prompt: &str, from_name: &str) -> String {
format!("[from: {}]\n\n{}", from_name, prompt)
}
pub fn sandbox_flag(yolo: bool) -> Vec<String> {
if yolo {
vec!["--dangerously-bypass-approvals-and-sandbox".to_string()]
} else {
vec!["--sandbox".to_string(), "workspace-write".to_string()]
}
}
pub fn approval_flag(yolo: bool) -> Vec<String> {
if yolo {
vec![]
} else {
vec!["--ask-for-approval".to_string(), "never".to_string()]
}
}
pub fn sandbox_flag_resume(yolo: bool) -> Vec<String> {
if yolo {
vec!["--dangerously-bypass-approvals-and-sandbox".to_string()]
} else {
vec![]
}
}
pub fn build_argv_create(
cwd: &Path,
full_prompt: &str,
yolo: bool,
model: Option<&str>,
reasoning_effort: Option<&str>,
add_dir: Option<&str>,
) -> Vec<String> {
let mut argv = vec!["codex".to_string()];
argv.extend(approval_flag(yolo));
argv.extend([
"exec".to_string(),
"--json".to_string(),
"-C".to_string(),
cwd.to_string_lossy().to_string(),
"--skip-git-repo-check".to_string(),
]);
if let Some(d) = add_dir.filter(|d| !d.is_empty()) {
argv.push("--add-dir".to_string());
argv.push(d.to_string());
}
if let Some(m) = model.filter(|m| !m.is_empty()) {
argv.push("--model".to_string());
argv.push(m.to_string());
}
if let Some(effort) = reasoning_effort.filter(|e| !e.is_empty()) {
argv.push("-c".to_string());
argv.push(format!("model_reasoning_effort={effort}"));
}
argv.extend(sandbox_flag(yolo));
argv.push(full_prompt.to_string());
argv
}
pub fn build_argv_resume(session_id: &str, full_prompt: &str, yolo: bool) -> Vec<String> {
let mut argv = vec![
"codex".to_string(),
"exec".to_string(),
"resume".to_string(),
session_id.to_string(),
"--json".to_string(),
"--skip-git-repo-check".to_string(),
];
argv.extend(sandbox_flag_resume(yolo));
argv.push(full_prompt.to_string());
argv
}
#[derive(Debug)]
pub enum JsonlEvent {
ThreadStarted { thread_id: String },
AgentMessage { text: String },
SoftError { message: String },
TurnCompleted,
Other { type_name: Option<String> },
}
pub fn parse_jsonl_line(line: &str) -> Option<JsonlEvent> {
let line = line.trim_end_matches('\n');
if line.is_empty() || !line.starts_with('{') {
return None;
}
let v: serde_json::Value = serde_json::from_str(line).ok()?;
if !v.is_object() {
return None;
}
let ev_type = v.get("type").and_then(|t| t.as_str());
match ev_type {
Some(t) if t == EV_SESSION => {
let thread_id = v.get("thread_id").and_then(|x| x.as_str())?;
Some(JsonlEvent::ThreadStarted {
thread_id: thread_id.to_string(),
})
}
Some(t) if t == EV_COMPLETE => Some(JsonlEvent::TurnCompleted),
Some(t) if t == EV_ITEM => {
let item = v.get("item").and_then(|x| x.as_object());
match item {
Some(item) => {
let item_type = item.get("type").and_then(|x| x.as_str());
match item_type {
Some(t) if t == ITEM_MESSAGE => {
let text = item.get("text").and_then(|x| x.as_str()).unwrap_or("");
Some(JsonlEvent::AgentMessage {
text: text.to_string(),
})
}
Some(t) if t == ITEM_ERROR => {
let msg = item.get("message").and_then(|x| x.as_str()).unwrap_or("");
Some(JsonlEvent::SoftError {
message: msg.to_string(),
})
}
_ => Some(JsonlEvent::Other {
type_name: ev_type.map(String::from),
}),
}
}
None => Some(JsonlEvent::Other {
type_name: ev_type.map(String::from),
}),
}
}
Some(t) => Some(JsonlEvent::Other {
type_name: Some(t.to_string()),
}),
None => Some(JsonlEvent::Other { type_name: None }),
}
}
#[derive(Debug)]
pub enum CodexAskError {
NotFound,
NoSessionId { types_seen: Vec<String> },
TeeOpen { message: String },
Timeout { timeout_sec: f64 },
Invocation { exit_code: i32, message: String },
SigkillEscalated { partial_exit_code: i32 },
OsError { message: String },
Interrupted,
}
impl CodexAskError {
pub fn exit_code(&self) -> i32 {
match self {
CodexAskError::NotFound => 14,
CodexAskError::NoSessionId { .. } => 11,
CodexAskError::TeeOpen { .. } => 12,
CodexAskError::Timeout { .. } => 15,
CodexAskError::Invocation { exit_code, .. } => {
if *exit_code != 0 {
*exit_code
} else {
1
}
}
CodexAskError::SigkillEscalated { partial_exit_code } => {
if *partial_exit_code != 0 {
*partial_exit_code
} else {
1
}
}
CodexAskError::OsError { .. } => 1,
CodexAskError::Interrupted => 130,
}
}
}
impl std::fmt::Display for CodexAskError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
CodexAskError::NotFound => write!(f, "codex binary not found on PATH"),
CodexAskError::NoSessionId { types_seen } => write!(
f,
"codex did not emit session id; saw events: {:?}; expected one of: [\"thread.started\"]",
types_seen
),
CodexAskError::TeeOpen { message } => {
write!(f, "codex provider: cannot open output tee: {}", message)
}
CodexAskError::Timeout { timeout_sec } => {
write!(f, "codex timed out after {}s", timeout_sec)
}
CodexAskError::Invocation { exit_code, message } => {
write!(f, "codex exited {} ({})", exit_code, message)
}
CodexAskError::SigkillEscalated { partial_exit_code } => write!(
f,
"codex was SIGKILL'd during reap (exit {}); partial reply discarded",
partial_exit_code
),
CodexAskError::OsError { message } => {
write!(f, "codex provider: OSError invoking codex: {}", message)
}
CodexAskError::Interrupted => {
write!(f, "codex interrupted by SIGINT (Ctrl-C)")
}
}
}
}
impl std::error::Error for CodexAskError {}
#[derive(Debug, Clone)]
pub struct CodexResult {
pub exit_code: i32,
pub session_id: Option<String>,
pub last_msg: String,
pub duration_ms: u64,
}
fn open_tee(log_path: &Path) -> Result<std::fs::File, CodexAskError> {
crate::subprocess_ask::open_tee(log_path).map_err(|e| CodexAskError::TeeOpen {
message: e.to_string(),
})
}
fn run_codex(
argv: &[String],
output_path: &Path,
timeout: Option<Duration>,
expect_session: bool,
popen_cwd: Option<&Path>,
agent_self: Option<&str>,
) -> Result<CodexResult, CodexAskError> {
use std::process::{Command, Stdio};
let started = Instant::now();
let tee_fh = open_tee(output_path)?;
let argv =
crate::spawn_gate::qos_wrap(popen_cwd.unwrap_or_else(|| Path::new(".")), argv.to_vec());
let mut cmd = Command::new(&argv[0]);
cmd.args(&argv[1..]);
cmd.stdin(Stdio::null()); cmd.stdout(Stdio::piped());
cmd.stderr(Stdio::piped());
if let Some(cwd) = popen_cwd {
cmd.current_dir(cwd);
}
if let Some(name) = agent_self {
cmd.env("FNO_AGENT_SELF", name);
cmd.env("FNO_AGENT_PROVIDER", "codex");
}
unsafe {
cmd.pre_exec(|| {
libc::setpgid(0, 0);
Ok(())
});
}
let mut child = match cmd.spawn() {
Ok(c) => c,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
return Err(CodexAskError::NotFound);
}
Err(e) => {
return Err(CodexAskError::OsError {
message: e.to_string(),
});
}
};
let pid = child.id();
let _sigint_guard = crate::subprocess_ask::SigintForwarder::install(pid);
let stdout_pipe = child.stdout.take().expect("stdout piped");
let stderr_pipe = child.stderr.take().expect("stderr piped");
let tee = std::sync::Arc::new(std::sync::Mutex::new(tee_fh));
let tee_stderr = tee.clone();
let stderr_handle = std::thread::spawn(move || {
let mut stderr_tee_warned: std::collections::HashSet<String> =
std::collections::HashSet::new();
for line in BufReader::new(stderr_pipe).lines() {
match line {
Ok(l) => {
let raw = format!("{}\n", l);
if let Ok(mut guard) = tee_stderr.lock() {
if let Err(e) = guard.write_all(raw.as_bytes()) {
let key = e.to_string();
if stderr_tee_warned.insert(key) {
eprintln!("codex provider: stderr tee write failed: {}", e);
}
} else {
let _ = guard.flush();
}
}
}
Err(_) => break,
}
}
});
let mut watchdog = crate::subprocess_ask::AskWatchdog::spawn(pid, timeout);
let mut session_id: Option<String> = None;
let mut last_msg = String::new();
let mut last_error_msg = String::new();
let mut types_seen: Vec<String> = Vec::new();
let mut tee_warned: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut stream_read_error: Option<String> = None;
let mut broke_on_complete = false;
let stdout_reader = BufReader::new(stdout_pipe);
for raw_line in stdout_reader.lines() {
let raw = match raw_line {
Ok(l) => l,
Err(e) => {
eprintln!("codex provider: stdout stream read error: {}", e);
stream_read_error = Some(e.to_string());
break;
}
};
let tee_line = format!("{}\n", raw);
if let Ok(mut guard) = tee.lock() {
if let Err(e) = guard.write_all(tee_line.as_bytes()) {
let key = e.to_string();
if !tee_warned.contains(&key) {
tee_warned.insert(key.clone());
eprintln!("codex provider: tee write failed: {}", e);
}
} else {
let _ = guard.flush();
}
}
match parse_jsonl_line(&raw) {
Some(JsonlEvent::ThreadStarted { thread_id }) => {
types_seen.push(EV_SESSION.to_string());
if session_id.is_none() && !thread_id.is_empty() {
session_id = Some(thread_id);
}
}
Some(JsonlEvent::AgentMessage { text }) => {
types_seen.push(EV_ITEM.to_string());
last_msg = text;
}
Some(JsonlEvent::SoftError { message }) => {
types_seen.push(EV_ITEM.to_string());
last_error_msg = message;
}
Some(JsonlEvent::TurnCompleted) => {
types_seen.push(EV_COMPLETE.to_string());
broke_on_complete = true;
break;
}
Some(JsonlEvent::Other { type_name }) => {
if let Some(t) = type_name {
if !types_seen.contains(&t) {
types_seen.push(t);
}
}
}
None => {} }
}
watchdog.cancel();
let (exit_code, sigkill_escalated) =
crate::subprocess_ask::wait_with_grace(pid, &mut child, 5.0);
watchdog.join();
if stderr_handle.join().is_err() {
eprintln!("codex provider: stderr drain thread panicked");
}
let duration_ms = started.elapsed().as_millis() as u64;
let was_timed_out = watchdog.timed_out();
if crate::subprocess_ask::ask_interrupted() {
return Err(CodexAskError::Interrupted);
}
if was_timed_out {
return Err(CodexAskError::Timeout {
timeout_sec: timeout.map(|d| d.as_secs_f64()).unwrap_or(0.0),
});
}
if expect_session && session_id.is_none() {
types_seen.sort();
types_seen.dedup();
return Err(CodexAskError::NoSessionId { types_seen });
}
if sigkill_escalated {
return Err(CodexAskError::SigkillEscalated {
partial_exit_code: exit_code,
});
}
if let Some(err) = stream_read_error {
if !broke_on_complete {
return Err(CodexAskError::Invocation {
exit_code,
message: format!(
"stream read error before turn.completed: {} (see output.jsonl)",
err
),
});
}
}
if exit_code != 0 && last_msg.is_empty() {
return Err(CodexAskError::Invocation {
exit_code,
message: format!("see output.jsonl for details"),
});
}
let effective_last_msg = if !last_msg.is_empty() {
last_msg
} else {
last_error_msg
};
Ok(CodexResult {
exit_code,
session_id,
last_msg: effective_last_msg,
duration_ms,
})
}
pub fn codex_create(
cwd: &Path,
prompt: &str,
from_name: &str,
yolo: bool,
output_path: &Path,
timeout: Option<Duration>,
agent_self: Option<&str>,
model: Option<&str>,
reasoning_effort: Option<&str>,
add_dir: Option<&str>,
) -> Result<CodexResult, CodexAskError> {
let effective_prompt = normalize_codex_command(prompt);
let full_prompt = inject_from_name(&effective_prompt, from_name);
let eff = crate::agents_config::effective_yolo(
yolo,
crate::agents_config::headless_yolo_enabled("codex", cwd),
);
let argv = build_argv_create(cwd, &full_prompt, eff, model, reasoning_effort, add_dir);
run_codex(&argv, output_path, timeout, true, None, agent_self)
}
pub fn codex_resume(
session_id: &str,
cwd: &Path,
prompt: &str,
from_name: &str,
yolo: bool,
output_path: &Path,
timeout: Option<Duration>,
) -> Result<CodexResult, CodexAskError> {
let effective_prompt = normalize_codex_command(prompt);
let full_prompt = inject_from_name(&effective_prompt, from_name);
let eff = crate::agents_config::effective_yolo(
yolo,
crate::agents_config::headless_yolo_enabled("codex", cwd),
);
let argv = build_argv_resume(session_id, &full_prompt, eff);
run_codex(&argv, output_path, timeout, false, Some(cwd), None)
}
fn derive_log_path(home: &AgentsHome, name: &str) -> PathBuf {
home.root()
.join("agents")
.join("logs")
.join(format!("{}.jsonl", name))
}
fn now_iso() -> String {
chrono::Utc::now().format("%Y-%m-%dT%H:%M:%SZ").to_string()
}
struct AgentLock {
_file: std::fs::File,
}
impl AgentLock {
fn acquire(home: &AgentsHome, name: &str, timeout: Duration) -> Result<Self, ()> {
let locks_dir = home.root().join("locks");
let _ = std::fs::create_dir_all(&locks_dir);
let path = locks_dir.join(format!("{}.lock", name));
let file = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.write(true)
.open(&path)
.map_err(|_| ())?;
let deadline = Instant::now() + timeout;
loop {
match file.try_lock() {
Ok(()) => return Ok(Self { _file: file }),
Err(_) => {
if Instant::now() >= deadline {
return Err(());
}
std::thread::sleep(Duration::from_millis(25));
}
}
}
}
}
impl Drop for AgentLock {
fn drop(&mut self) {
let _ = self._file.unlock();
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AskOutcome {
pub stdout: String,
pub stderr: String,
pub exit_code: i32,
}
impl AskOutcome {
fn ok_reply(reply: String) -> Self {
Self {
stdout: reply,
stderr: String::new(),
exit_code: 0,
}
}
fn err(msg: impl Into<String>, code: i32) -> Self {
Self {
stdout: String::new(),
stderr: format!("{}\n", msg.into()),
exit_code: code,
}
}
}
#[allow(clippy::too_many_arguments)]
pub fn dispatch_codex_ask(
home: &AgentsHome,
name: &str,
message: &str,
from_name: &str,
_cwd: &Path,
yolo: bool,
timeout: Option<Duration>,
) -> AskOutcome {
if let Err(msg) = crate::claude_ask::validate_inputs(name, message, from_name) {
return AskOutcome::err(msg, 2);
}
let events = home.events_jsonl();
let registry_path = home.registry_json();
let _lock = match AgentLock::acquire(home, name, LOCK_ACQUIRE_TIMEOUT) {
Ok(l) => l,
Err(()) => {
emit_event(
&events,
"agent_ask_failed",
&[
("stage", "lock-timeout".into()),
("name", name.into()),
("provider", "codex".into()),
],
);
return AskOutcome::err(
format!(
"lock timeout for agent {:?} after {:.1}s",
name,
LOCK_ACQUIRE_TIMEOUT.as_secs_f64()
),
11,
);
}
};
let registry = match load_registry(®istry_path) {
Ok(r) => r,
Err(e) => {
emit_event(
&events,
"agent_ask_failed",
&[
("stage", "registry-read".into()),
("name", name.into()),
("provider", "codex".into()),
("error", e.to_string().into()),
],
);
return AskOutcome::err(format!("registry read failed: {}", e), 12);
}
};
let existing = registry.find(name).cloned();
match existing {
None => {
emit_event(
&events,
"agent_ask_failed",
&[
("stage", "unknown-name".into()),
("name", name.into()),
("provider", "codex".into()),
],
);
AskOutcome::err(
format!(
"unknown agent {}; spawn it first: fno agents spawn {} --harness <harness>",
crate::claude_ask::py_repr(name),
name
),
16,
)
}
Some(entry) => dispatch_resume(
home,
&events,
®istry_path,
name,
&entry,
message,
from_name,
yolo,
timeout,
),
}
}
#[allow(clippy::too_many_arguments)]
pub fn dispatch_codex_once(
home: &AgentsHome,
name: &str,
message: &str,
from_name: &str,
cwd: &Path,
yolo: bool,
timeout: Option<Duration>,
model: Option<&str>,
reasoning_effort: Option<&str>,
add_dir: Option<&str>,
) -> AskOutcome {
use crate::claude_ask::py_repr;
if let Err(msg) = crate::claude_ask::validate_spawn_inputs(name, from_name) {
return AskOutcome::err(msg, 2);
}
let events = home.events_jsonl();
let registry_path = home.registry_json();
let _lock = match AgentLock::acquire(home, name, LOCK_ACQUIRE_TIMEOUT) {
Ok(l) => l,
Err(()) => {
emit_event(
&events,
"agent_ask_failed",
&[
("stage", "lock-timeout".into()),
("name", name.into()),
("provider", "codex".into()),
],
);
return AskOutcome::err(
format!(
"lock timeout for agent {} after {:.1}s",
py_repr(name),
LOCK_ACQUIRE_TIMEOUT.as_secs_f64()
),
11,
);
}
};
let registry = match load_registry(®istry_path) {
Ok(r) => r,
Err(e) => {
return AskOutcome::err(format!("registry read failed: {}", e), 12);
}
};
if registry.find(name).is_some() {
return AskOutcome::err(
format!(
"agent {} already exists; use 'fno agents rm {}' first or pick another name",
py_repr(name),
name
),
2,
);
}
let effective_message = if message.is_empty() { "hello" } else { message };
let inner = dispatch_create(
home,
&events,
®istry_path,
name,
effective_message,
from_name,
cwd,
yolo,
timeout,
model,
reasoning_effort,
add_dir,
);
if inner.exit_code != 0 {
return inner;
}
let session_or_short_id = load_registry(®istry_path)
.ok()
.and_then(|r| r.find(name).and_then(|e| e.codex_session_id.clone()))
.unwrap_or_default();
let teardown_err = update_registry(®istry_path, |reg| {
reg.entries.retain(|e| e.name != name);
true
})
.err();
let teardown_receipt = if let Some(e) = teardown_err {
format!(
"fno agents spawn: warning: teardown failed for {} (codex/{}): {}. Peer leaked -- clean up via 'fno agents rm {}'\n",
py_repr(name),
session_or_short_id,
e,
name
)
} else {
format!("once: {} (codex/{}) torn down\n", name, session_or_short_id)
};
AskOutcome {
stdout: inner.stdout,
stderr: teardown_receipt,
exit_code: 0,
}
}
fn dispatch_create(
home: &AgentsHome,
events: &Path,
registry_path: &Path,
name: &str,
message: &str,
from_name: &str,
cwd: &Path,
yolo: bool,
timeout: Option<Duration>,
model: Option<&str>,
reasoning_effort: Option<&str>,
add_dir: Option<&str>,
) -> AskOutcome {
let output_path = derive_log_path(home, name);
let timeout_sec = timeout.unwrap_or(DEFAULT_FOLLOWUP_TIMEOUT);
let result = match codex_create(
cwd,
message,
from_name,
yolo,
&output_path,
Some(timeout_sec),
Some(name),
model,
reasoning_effort,
add_dir,
) {
Ok(r) => r,
Err(e) => {
let stage = match &e {
CodexAskError::NoSessionId { .. } => "codex-no-session",
CodexAskError::Timeout { .. } => "codex-timeout",
CodexAskError::Interrupted => "codex-interrupted",
_ => "codex-subprocess",
};
let exit_code = e.exit_code();
let msg = format!("{} (see {} for details)", e, output_path.display());
emit_event(
events,
"agent_ask_failed",
&[
("stage", stage.into()),
("name", name.into()),
("provider", "codex".into()),
("returncode", exit_code.into()),
],
);
return AskOutcome::err(msg, exit_code);
}
};
let session_id = result
.session_id
.expect("codex_create guarantees session_id on success (expect_session=true)");
use crate::state::RegistryEntry;
let new_entry = RegistryEntry {
name: name.to_string(),
short_id: String::new(),
legacy_provider: String::new(),
harness: Some("codex".to_string()),
harness_session_id: Some(session_id.clone()),
cwd: cwd.to_string_lossy().to_string(),
project_root: String::new(),
session_id: None,
legacy_claude_short_id: None,
claude_session_uuid: None,
messaging_socket_path: None,
codex_session_id: Some(session_id.clone()),
gemini_session_id: None,
mcp_channel_id: None,
host_mode: None, cc_session_id: None,
status: AgentStatus::Live,
last_message_at: None,
created_at: now_iso(),
pid: None,
pid_start_time: None,
log_path: Some(output_path.to_string_lossy().to_string()),
last_reconciled_at: None,
inside_leg: None,
exited_at: None,
mux: None,
screen_state: None,
crown_level: None,
crown_scope: None,
crown_grantor: None,
};
match update_registry(registry_path, |reg| {
if reg.find(name).is_some() {
false
} else {
reg.entries.push(new_entry.clone());
true
}
}) {
Ok(true) => {}
Ok(false) => {
emit_event(
events,
"agent_ask_failed",
&[
("stage", "name-collision".into()),
("name", name.into()),
("provider", "codex".into()),
("codex_session_id", session_id.clone().into()),
],
);
return AskOutcome::err(
format!(
"agent {:?} already exists (registered concurrently); orphaned codex session: {:?}",
name, session_id
),
12,
);
}
Err(e) => {
emit_event(
events,
"agent_ask_failed",
&[
("stage", "registry-write".into()),
("name", name.into()),
("provider", "codex".into()),
("codex_session_id", session_id.clone().into()),
],
);
return AskOutcome::err(
format!(
"registry write failed: {}. orphaned codex session: {:?} (see output.jsonl)",
e, session_id
),
12,
);
}
}
emit_event(
events,
"agent_ask_done",
&[
("stage", "dispatch".into()),
("name", name.into()),
("provider", "codex".into()),
("codex_session_id", session_id.clone().into()),
("duration_ms", (result.duration_ms as u64).into()),
("yolo", yolo.into()),
],
);
AskOutcome::ok_reply(result.last_msg)
}
fn dispatch_resume(
_home: &AgentsHome,
events: &Path,
registry_path: &Path,
name: &str,
entry: &crate::state::RegistryEntry,
message: &str,
from_name: &str,
yolo: bool,
timeout: Option<Duration>,
) -> AskOutcome {
let session_id = match entry.codex_session_id.as_deref() {
Some(s) if !s.is_empty() => s.to_string(),
_ => {
return AskOutcome::err(
format!(
"registry entry {:?} has no codex_session_id; cannot follow up. \
Remove with 'fno agents rm {}' and recreate.",
name, name
),
11,
);
}
};
let log_path = match entry.log_path.as_deref() {
Some(p) if !p.is_empty() => PathBuf::from(p),
_ => {
return AskOutcome::err(
format!(
"registry entry {:?} has empty log_path; run 'fno agents rm {}' and recreate.",
name, name
),
11,
);
}
};
let registered_cwd = match entry.cwd.as_str() {
"" => {
return AskOutcome::err(
format!(
"registry entry {:?} has empty cwd; codex sessions are cwd-pinned. \
Run 'fno agents rm {}' and recreate.",
name, name
),
11,
);
}
c => PathBuf::from(c),
};
emit_event(
events,
"agent_followup_started",
&[
("name", name.into()),
("provider", "codex".into()),
("codex_session_id", session_id.clone().into()),
("yolo", yolo.into()),
],
);
let timeout_sec = timeout.unwrap_or(DEFAULT_FOLLOWUP_TIMEOUT);
let result = match codex_resume(
&session_id,
®istered_cwd,
message,
from_name,
yolo,
&log_path,
Some(timeout_sec),
) {
Ok(r) => r,
Err(e) => {
let stage = match &e {
CodexAskError::Timeout { .. } => "codex-timeout",
CodexAskError::Interrupted => "codex-interrupted",
_ => "codex-subprocess",
};
let exit_code = e.exit_code();
emit_event(
events,
"agent_followup_failed",
&[
("stage", stage.into()),
("name", name.into()),
("provider", "codex".into()),
("codex_session_id", session_id.clone().into()),
("returncode", exit_code.into()),
],
);
let msg = format!(
"{} (see {} for details). If the session was lost, run 'fno agents rm {}' then re-ask.",
e, log_path.display(), name
);
return AskOutcome::err(msg, exit_code);
}
};
if let Err(e) = update_registry(registry_path, |reg| {
if let Some(en) = reg.find_mut(name) {
en.status = AgentStatus::Live;
en.last_message_at = Some(now_iso());
}
}) {
emit_event(
events,
"agent_followup_failed",
&[
("stage", "registry-write".into()),
("name", name.into()),
("provider", "codex".into()),
("codex_session_id", session_id.clone().into()),
("error", e.to_string().into()),
],
);
return AskOutcome::err(
format!(
"registry write failed: {}. NOTE: message was already delivered; do not retry. \
(agent={:?} session={:?})",
e, name, session_id
),
12,
);
}
emit_event(
events,
"agent_followup_done",
&[
("stage", "followup".into()),
("name", name.into()),
("provider", "codex".into()),
("codex_session_id", session_id.clone().into()),
(
"reply_chars",
(result.last_msg.chars().count() as u64).into(),
),
("yolo", yolo.into()),
],
);
AskOutcome::ok_reply(result.last_msg)
}
pub fn maybe_run_codex_ask(
home: &AgentsHome,
params: &serde_json::Value,
name: &str,
) -> Option<i32> {
let provider_param = params.get("provider").and_then(|v| v.as_str());
let registry = match load_registry(&home.registry_json()) {
Ok(r) => r,
Err(e) => {
eprintln!(
"fno-agents: cannot read agents registry at {:?}: {}",
home.registry_json(),
e
);
return Some(12);
}
};
let existing_provider = registry.find(name).map(|e| e.harness_name().to_string());
if let (Some(ep), Some(pp)) = (existing_provider.as_deref(), provider_param) {
if ep == "codex" && pp != "codex" {
eprintln!(
"fno-agents: agent {:?} already exists with provider 'codex'; \
refusing to override with --provider {}",
name, pp
);
return Some(2);
}
}
let resolved = existing_provider.as_deref().or(provider_param);
if resolved != Some("codex") {
return None; }
let message = params.get("message").and_then(|v| v.as_str()).unwrap_or("");
let from_name = params
.get("from_name")
.and_then(|v| v.as_str())
.unwrap_or("fno");
let cwd = crate::subprocess_ask::resolve_ask_cwd(params.get("cwd").and_then(|v| v.as_str()));
let timeout = params
.get("timeout")
.and_then(|v| v.as_u64())
.map(std::time::Duration::from_secs);
let yolo = params
.get("yolo")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let outcome = dispatch_codex_ask(home, name, message, from_name, &cwd, yolo, timeout);
if !outcome.stderr.is_empty() {
eprint!("{}", outcome.stderr);
}
if !outcome.stdout.is_empty() {
print!("{}", outcome.stdout);
}
Some(outcome.exit_code)
}