Files
ai-app/server/src/session/codex.rs
T

1533 lines
53 KiB
Rust

//! Codex CLI sessions over its persistent app-server protocol.
//!
//! One app-server owns one conversation tree. Unlike `codex exec --json`,
//! this surface can interrupt a turn without killing the connection and can
//! steer an active turn through `turn/steer`. Child threads are multiplexed
//! onto the same stdout and routed into subagent transcripts. Its stdio is a
//! fifo plus logs so both the process and in-flight turns survive an
//! ai-server restart.
mod translate;
use std::collections::VecDeque;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use tokio::io::AsyncWriteExt;
use tokio::sync::mpsc;
use super::driver::{
AttachmentRef, Driver, Event, EventSink, SessionStatus, Unqueued, store_image,
};
use super::process;
use super::subagent::Subagents;
use super::transport::{Launch, Streams, Transport};
use crate::config::{ProviderConfig, SessionConfig};
use translate::Translator;
const STDIN_FIFO: &str = "codex-stdin.fifo";
const STDOUT_LOG: &str = "codex-stdout.log";
const STDERR_LOG: &str = "codex-stderr.log";
const THREAD_FILE: &str = "codex-thread.json";
const STATE_FILE: &str = "codex-state.json";
const POLL: std::time::Duration = std::time::Duration::from_millis(50);
#[derive(Clone, Default, Serialize, Deserialize)]
struct Waiting {
/// The phone's queued-bubble id. Empty for the message that starts a turn.
#[serde(default)]
id: String,
/// The id Codex echoes on its user-message item.
#[serde(default)]
client_id: String,
/// Typed while a turn was already running, so the CLI is told as much --
/// see `super::driver::message_body`. Recorded at the moment it was typed
/// rather than read off the turn state when it is dispatched, because a
/// steer Codex refuses as `activeTurnNotSteerable` is requeued and sent as
/// the start of the next turn, which is precisely the case the note exists
/// for.
#[serde(default)]
steering: bool,
text: String,
#[serde(default)]
attachments: Vec<AttachmentRef>,
}
#[derive(Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
enum RequestKind {
Start,
Steer,
}
#[derive(Clone, Serialize, Deserialize)]
struct PendingRequest {
id: String,
client_id: String,
kind: RequestKind,
}
#[derive(Default, Serialize, Deserialize)]
struct ProtocolState {
#[serde(default)]
waiting: VecDeque<Waiting>,
#[serde(default)]
sent: VecDeque<Waiting>,
#[serde(default)]
pending: Vec<PendingRequest>,
#[serde(default)]
initialize_request: Option<String>,
#[serde(default)]
thread_request: Option<String>,
#[serde(default)]
active_turn: Option<String>,
#[serde(default)]
interrupt_when_started: bool,
#[serde(skip)]
running: bool,
#[serde(skip)]
closed: bool,
}
#[derive(Clone)]
struct Settings {
model: Option<String>,
permission_mode: Option<String>,
effort: Option<String>,
}
struct Inner {
sink: EventSink,
state: Mutex<ProtocolState>,
settings: Mutex<Settings>,
to_child: mpsc::UnboundedSender<String>,
transport: Transport,
session_dir: PathBuf,
subagents: Arc<Subagents>,
reading: AtomicBool,
}
pub struct CodexDriver {
inner: Arc<Inner>,
}
impl CodexDriver {
pub fn launch(
meta: &SessionConfig,
provider: &ProviderConfig,
transport: Transport,
session_dir: &Path,
sink: EventSink,
subagents: Arc<Subagents>,
) -> Result<Self> {
let mut state = read_state(session_dir);
let recorded = process::recorded(session_dir);
let (record, started_here) = match recorded {
Some((record, process::Liveness::Alive | process::Liveness::Unknown)) => {
state.running = state.active_turn.is_some()
|| state
.pending
.iter()
.any(|request| request.kind == RequestKind::Start);
(record, false)
}
Some((_, process::Liveness::Dead)) | None => {
process::clear(session_dir);
// Requests written to a dead process have no recipient. Put their messages back
// in front of the unsent queue so restarting cannot silently lose them.
while let Some(message) = state.sent.pop_back() {
state.waiting.push_front(message);
}
state.pending.clear();
state.initialize_request = None;
state.thread_request = None;
state.active_turn = None;
state.interrupt_when_started = false;
state.running = !state.waiting.is_empty();
(
start_process(meta, provider, &transport, session_dir)?,
true,
)
}
};
let stdin = std::fs::OpenOptions::new()
.write(true)
.open(session_dir.join(STDIN_FIFO))
.with_context(|| format!("opening {STDIN_FIFO} for session {}", meta.id))?;
let (to_child, mut from_driver) = mpsc::unbounded_channel::<String>();
tokio::spawn(async move {
let mut stdin = tokio::fs::File::from_std(stdin);
while let Some(line) = from_driver.recv().await {
if stdin.write_all(line.as_bytes()).await.is_err()
|| stdin.write_all(b"\n").await.is_err()
|| stdin.flush().await.is_err()
{
break;
}
}
});
let inner = Arc::new(Inner {
sink,
state: Mutex::new(state),
settings: Mutex::new(Settings {
model: meta.model.clone(),
permission_mode: meta.permission_mode.clone(),
effort: meta.effort.clone(),
}),
to_child,
transport,
session_dir: session_dir.to_path_buf(),
subagents,
reading: AtomicBool::new(true),
});
if started_here {
begin_initialization(&inner);
let _ = inner.sink.send(Event::Status {
state: SessionStatus::Idle,
});
}
save_state(&inner);
spawn_follower(Arc::clone(&inner), record);
Ok(Self { inner })
}
}
impl Driver for CodexDriver {
fn send_user_message(&self, text: String, attachments: Vec<AttachmentRef>) {
let mut state = self.inner.state.lock().unwrap();
if state.closed {
drop(state);
let _ = self.inner.sink.send(Event::Error {
message: "this Codex session has been stopped".to_string(),
});
return;
}
// Initialization and missing-thread recovery can leave an otherwise idle session without
// a thread to receive this yet. Announce that wait just like one behind an active turn: the
// phone stops preserving its local bridge once this method returns, so the server event is
// the durable copy until Codex acknowledges the message.
let missing_thread = read_thread(&self.inner.session_dir).is_none();
let waiting_for_protocol =
state.initialize_request.is_some() || state.thread_request.is_some() || missing_thread;
let id = (state.running || waiting_for_protocol)
.then(super::random_hex)
.unwrap_or_default();
let steering = state.running;
state.running = true;
state.waiting.push_back(Waiting {
id: id.clone(),
client_id: format!("ai-app-{}", super::random_hex()),
steering,
text: text.clone(),
attachments: attachments.clone(),
});
let start_thread = (missing_thread
&& state.initialize_request.is_none()
&& state.thread_request.is_none())
.then(|| {
let request = request_id();
state.thread_request = Some(request.clone());
request
});
drop(state);
save_state(&self.inner);
if !id.is_empty() {
let _ = self.inner.sink.send(Event::MessageQueued {
id,
text,
attachments,
});
}
if let Some(request) = start_thread {
send_json(
&self.inner,
json!({
"id": request,
"method": "thread/start",
"params": thread_params(&self.inner)
}),
);
}
dispatch_waiting(&self.inner);
}
fn unqueue(&self, id: &str) -> Unqueued {
let mut state = self.inner.state.lock().unwrap();
if let Some(at) = state.waiting.iter().position(|message| message.id == id) {
state.waiting.remove(at);
drop(state);
save_state(&self.inner);
let _ = self
.inner
.sink
.send(Event::MessageDropped { id: id.to_string() });
return Unqueued::Dropped;
}
if state.sent.iter().any(|message| message.id == id) {
Unqueued::AlreadySent
} else {
Unqueued::Unknown
}
}
fn answer_question(&self, _id: &str, _answers: &[String]) {
let _ = self.inner.sink.send(Event::Error {
message: "Codex question answering is not available in this app yet".to_string(),
});
}
fn interrupt(&self) {
let mut state = self.inner.state.lock().unwrap();
let Some(thread_id) = read_thread(&self.inner.session_dir) else {
state.interrupt_when_started = state.running;
drop(state);
save_state(&self.inner);
return;
};
let Some(turn_id) = state.active_turn.clone() else {
state.interrupt_when_started = state.running;
drop(state);
save_state(&self.inner);
return;
};
state.interrupt_when_started = false;
drop(state);
save_state(&self.inner);
send_request(
&self.inner,
"turn/interrupt",
json!({"threadId": thread_id, "turnId": turn_id}),
);
}
fn set_model(&self, model: &str) {
self.inner.settings.lock().unwrap().model = Some(model.to_string());
let _ = self.inner.sink.send(Event::Settings {
model: Some(model.to_string()),
permission_mode: None,
});
}
fn set_permission_mode(&self, mode: &str) {
self.inner.settings.lock().unwrap().permission_mode = Some(mode.to_string());
let _ = self.inner.sink.send(Event::Settings {
model: None,
permission_mode: Some(mode.to_string()),
});
}
fn set_title(&self, title: &str) {
if let Some(thread_id) = read_thread(&self.inner.session_dir) {
send_request(
&self.inner,
"thread/name/set",
json!({"threadId": thread_id, "name": title}),
);
}
}
fn run_command(&self, text: &str) {
self.send_user_message(text.to_string(), Vec::new());
}
fn compact(&self) {
self.send_user_message("/compact".to_string(), Vec::new());
}
fn clear(&self) {
let path = self.inner.session_dir.join(THREAD_FILE);
if let Err(err) = std::fs::remove_file(&path)
&& err.kind() != std::io::ErrorKind::NotFound
{
let _ = self.inner.sink.send(Event::Error {
message: format!("couldn't clear the Codex thread id: {err}"),
});
return;
}
// A thread has no durable rollout until its first turn. Creating an empty replacement here
// leaves an id that a restarted app-server cannot resume; start it with the next message
// instead, when Codex can make the thread and its first turn together.
self.inner.state.lock().unwrap().thread_request = None;
save_state(&self.inner);
let _ = self.inner.sink.send(Event::Cleared);
}
fn between_turns(&self) -> bool {
let state = self.inner.state.lock().unwrap();
!state.running && !state.closed
}
fn detach(&self) {
self.inner.reading.store(false, Ordering::SeqCst);
}
fn stop(&self) {
self.inner.reading.store(false, Ordering::SeqCst);
let dropped = {
let mut state = self.inner.state.lock().unwrap();
state.closed = true;
let mut messages: Vec<_> = state.waiting.drain(..).collect();
messages.extend(state.sent.drain(..));
messages
.into_iter()
.filter_map(|message| (!message.id.is_empty()).then_some(message.id))
.collect::<Vec<_>>()
};
save_state(&self.inner);
for id in dropped {
let _ = self.inner.sink.send(Event::MessageDropped { id });
}
if let Some(record) = process::live(&self.inner.session_dir) {
process::stop(&record, process::STOP_GRACE);
}
process::clear(&self.inner.session_dir);
}
}
fn start_process(
meta: &SessionConfig,
provider: &ProviderConfig,
transport: &Transport,
session_dir: &Path,
) -> Result<process::Record> {
let stdin = process::make_fifo(&session_dir.join(STDIN_FIFO))?;
let stdout = process::create_log(&session_dir.join(STDOUT_LOG))?;
let stderr = process::create_log(&session_dir.join(STDERR_LOG))?;
let launch = Launch::new(
provider.program(),
vec!["app-server".to_string(), "--stdio".to_string()],
meta.cwd.as_deref(),
);
let mut child = transport.spawn(
&launch,
Streams::Detached {
stdin: stdin.into(),
stdout: stdout.into(),
stderr: stderr.into(),
},
)?;
let pid = child
.id()
.context("Codex app-server exited before it could be recorded")?;
tokio::spawn(async move {
let _ = child.wait().await;
});
let record = process::Record::of(pid, process::Detail::Stdio { stdout_read: 0 })
.context("Codex app-server was gone before its start time could be read")?;
process::write(session_dir, &record);
Ok(record)
}
fn begin_initialization(inner: &Arc<Inner>) {
let id = request_id();
inner.state.lock().unwrap().initialize_request = Some(id.clone());
save_state(inner);
send_json(
inner,
json!({
"id": id,
"method": "initialize",
"params": {"clientInfo": {
"name": "ai-app",
"title": "AI Sessions",
"version": env!("CARGO_PKG_VERSION")
}}
}),
);
}
fn send_json(inner: &Inner, value: Value) {
let _ = inner.to_child.send(value.to_string());
}
fn send_request(inner: &Inner, method: &str, params: Value) -> String {
let id = request_id();
send_json(inner, json!({"id": id, "method": method, "params": params}));
id
}
fn request_id() -> String {
format!("ai-app-{}", super::random_hex())
}
fn thread_params(inner: &Inner) -> Value {
let settings = inner.settings.lock().unwrap().clone();
let mut params = json!({"approvalPolicy": "never"});
if let Some(model) = settings.model {
params["model"] = Value::String(model);
}
if let Some(mode) = settings.permission_mode {
params["sandbox"] = Value::String(mode);
}
params
}
fn turn_params(inner: &Inner, thread_id: &str, message: &Waiting) -> Result<Value> {
let settings = inner.settings.lock().unwrap().clone();
let mut params = json!({
"threadId": thread_id,
"input": input_for(inner, message)?,
"clientUserMessageId": message.client_id,
});
if let Some(model) = settings.model {
params["model"] = Value::String(model);
}
if let Some(effort) = settings.effort {
params["effort"] = Value::String(effort);
}
if let Some(mode) = settings.permission_mode {
params["sandboxPolicy"] = sandbox_policy(&mode);
params["approvalPolicy"] = Value::String("never".to_string());
}
Ok(params)
}
fn sandbox_policy(mode: &str) -> Value {
match mode {
"danger-full-access" => json!({"type": "dangerFullAccess"}),
"read-only" => json!({"type": "readOnly"}),
_ => json!({"type": "workspaceWrite"}),
}
}
fn input_for(inner: &Inner, message: &Waiting) -> Result<Vec<Value>> {
let mut files = Vec::new();
let mut input = Vec::new();
for attachment in &message.attachments {
let path = attachment_path(&inner.session_dir, attachment)?;
if let Some(media_type) = crate::media::media_type_for(attachment) {
if matches!(inner.transport, Transport::Here) {
input.push(json!({"type": "localImage", "path": path}));
} else {
// The upload is on the server, not on the machine reached over ssh. Inline image
// input carries those bytes across the app-server connection just as Claude's
// image block does; naming the server path left Codex unable to see it.
input.push(inline_image(&path, media_type)?);
}
} else {
files.push(path);
}
}
let text = super::driver::message_body(&message.text, &files, message.steering);
if !text.is_empty() {
input.insert(0, json!({"type": "text", "text": text}));
}
Ok(input)
}
fn inline_image(path: &Path, media_type: &str) -> Result<Value> {
use base64::Engine;
let bytes =
std::fs::read(path).with_context(|| format!("read attachment {}", path.display()))?;
let data = base64::engine::general_purpose::STANDARD.encode(bytes);
Ok(json!({
"type": "image",
"url": format!("data:{media_type};base64,{data}")
}))
}
fn dispatch_waiting(inner: &Arc<Inner>) {
let Some(thread_id) = read_thread(&inner.session_dir) else {
return;
};
loop {
let (message, kind, active_turn) = {
let mut state = inner.state.lock().unwrap();
if state.closed || state.waiting.is_empty() {
return;
}
if state.active_turn.is_none()
&& state
.pending
.iter()
.any(|request| request.kind == RequestKind::Start)
{
return;
}
let kind = if state.active_turn.is_some() {
RequestKind::Steer
} else {
RequestKind::Start
};
let mut message = state.waiting.pop_front().unwrap();
// Queues written by the old per-turn driver predate app-server's
// client id. Give an adopted message one before Codex sees it so
// its eventual user item can still resolve the queued bubble.
if message.client_id.is_empty() {
message.client_id = format!("ai-app-{}", super::random_hex());
}
let active_turn = state.active_turn.clone();
(message, kind, active_turn)
};
let mut params = match turn_params(inner, &thread_id, &message) {
Ok(params) => params,
Err(err) => {
if !message.id.is_empty() {
let _ = inner.sink.send(Event::MessageDropped {
id: message.id.clone(),
});
}
let _ = inner.sink.send(Event::Error {
message: format!("couldn't send an attachment to Codex: {err:#}"),
});
continue;
}
};
let method = match kind {
RequestKind::Start => "turn/start",
RequestKind::Steer => {
params["expectedTurnId"] = Value::String(active_turn.unwrap());
"turn/steer"
}
};
let id = request_id();
{
let mut state = inner.state.lock().unwrap();
state.sent.push_back(message.clone());
state.pending.push(PendingRequest {
id: id.clone(),
client_id: message.client_id.clone(),
kind,
});
}
save_state(inner);
send_json(inner, json!({"id": id, "method": method, "params": params}));
if kind == RequestKind::Start {
return;
}
}
}
fn spawn_follower(inner: Arc<Inner>, record: process::Record) {
let offset = match record.detail {
process::Detail::Stdio { stdout_read } => stdout_read,
_ => 0,
};
tokio::spawn(follow(inner, record, offset));
}
async fn follow(inner: Arc<Inner>, mut record: process::Record, mut offset: u64) {
let stdout = inner.session_dir.join(STDOUT_LOG);
let stderr = inner.session_dir.join(STDERR_LOG);
let mut translator = Translator::new(
Arc::clone(&inner.subagents),
read_thread(&inner.session_dir),
inner.state.lock().unwrap().active_turn.is_some(),
);
// `offset` is the durable boundary after the last complete record. `read_at` may move beyond
// it while app-server is still writing one record. Image-bearing tool results can be several
// megabytes long, and rereading their incomplete prefix every 50 ms made arrival over ssh
// quadratic in time and allocation.
let mut read_at = offset;
let mut pending = Vec::new();
while inner.reading.load(Ordering::SeqCst) {
let (bytes, next) = match process::read_from(&stdout, read_at) {
Ok(read) => read,
Err(err) => {
let _ = inner.sink.send(Event::Error {
message: format!("couldn't read Codex output: {err:#}"),
});
return;
}
};
read_at = next;
pending.extend_from_slice(&bytes);
let complete = pending
.iter()
.rposition(|byte| *byte == b'\n')
.map(|at| at + 1)
.unwrap_or(0);
for line in String::from_utf8_lossy(&pending[..complete]).lines() {
let Ok(value) = serde_json::from_str::<Value>(line) else {
tracing::warn!(
"unparseable Codex JSONL line: {}",
line.chars().take(200).collect::<String>()
);
continue;
};
handle_line(&inner, &mut translator, &value);
}
if complete > 0 {
offset += complete as u64;
pending.drain(..complete);
record.detail = process::Detail::Stdio {
stdout_read: offset,
};
process::write(&inner.session_dir, &record);
}
match record.liveness() {
process::Liveness::Alive | process::Liveness::Unknown => {}
process::Liveness::Dead if complete > 0 => {}
process::Liveness::Dead => {
let stopped = process::stopping(&inner.session_dir);
process::clear(&inner.session_dir);
let dropped = {
let mut state = inner.state.lock().unwrap();
state.closed = true;
state.running = false;
let mut messages: Vec<_> = state.waiting.drain(..).collect();
messages.extend(state.sent.drain(..));
messages
.into_iter()
.filter_map(|message| (!message.id.is_empty()).then_some(message.id))
.collect::<Vec<_>>()
};
save_state(&inner);
for id in dropped {
let _ = inner.sink.send(Event::MessageDropped { id });
}
if !stopped {
let detail = stderr_tail(&stderr);
if !detail.is_empty() {
let _ = inner.sink.send(Event::Error {
message: format!("Codex exited:\n{detail}"),
});
}
}
let _ = inner.sink.send(Event::Status {
state: SessionStatus::Exited,
});
return;
}
}
tokio::time::sleep(POLL).await;
}
}
fn handle_line(inner: &Arc<Inner>, translator: &mut Translator, line: &Value) {
if line.get("id").is_some() {
handle_response(inner, line);
}
let method = line.get("method").and_then(Value::as_str);
let params = line.get("params").unwrap_or(&Value::Null);
let parent = translator.is_parent(line);
match method.filter(|_| parent) {
Some("thread/started") => {
if let Some(thread) = params.pointer("/thread/id").and_then(Value::as_str) {
write_thread(&inner.session_dir, thread);
}
}
Some("turn/started") => {
if let Some(turn) = params.pointer("/turn/id").and_then(Value::as_str) {
let interrupt = {
let mut state = inner.state.lock().unwrap();
state.active_turn = Some(turn.to_string());
state.running = true;
std::mem::take(&mut state.interrupt_when_started)
};
save_state(inner);
if interrupt {
if let Some(thread) = read_thread(&inner.session_dir) {
send_request(
inner,
"turn/interrupt",
json!({"threadId": thread, "turnId": turn}),
);
}
} else {
dispatch_waiting(inner);
}
}
}
Some("item/completed") => {
let item = &params["item"];
if item.get("type").and_then(Value::as_str) == Some("userMessage") {
announce_user(inner, item);
}
}
Some("turn/completed") => {
let mut state = inner.state.lock().unwrap();
state.active_turn = None;
state.running = !state.waiting.is_empty();
state.interrupt_when_started = false;
drop(state);
save_state(inner);
}
_ => {}
}
let images = if method == Some("item/completed") {
image_events(inner, &params["item"])
} else {
Vec::new()
};
for event in translator.translate_with_prefix(line, images) {
let _ = inner.sink.send(event);
}
if parent && method == Some("thread/started") {
// `thread/start` answers before this notification. Ordinarily its response releases the
// queue, but this notification is the authoritative point at which the replacement root
// exists and the translator knows its id. Retrying here makes that handoff level-triggered:
// a message cannot remain stuck until another send happens to call `dispatch_waiting`.
dispatch_waiting(inner);
}
if parent && method == Some("turn/completed") {
dispatch_waiting(inner);
}
}
/// Images embedded in a structured tool result, copied into the session before the translator's
/// `ToolEnd` is emitted so they stay attached to that call in transcript order.
fn image_events(inner: &Inner, item: &Value) -> Vec<Event> {
let Some(id) = item.get("id").and_then(Value::as_str) else {
return Vec::new();
};
let mut images = Vec::new();
match item.get("type").and_then(Value::as_str) {
Some("dynamicToolCall") => save_data_images(
&inner.session_dir,
item.get("contentItems").and_then(Value::as_array),
"inputImage",
"imageUrl",
&mut images,
),
Some("mcpToolCall") => {
if let Some(parts) = item.pointer("/result/content").and_then(Value::as_array) {
for part in parts {
if part.get("type").and_then(Value::as_str) != Some("image") {
continue;
}
if let (Some(data), Some(media_type)) = (
part.get("data").and_then(Value::as_str),
part.get("mimeType")
.or_else(|| part.get("mime_type"))
.and_then(Value::as_str),
) && let Some(image) =
save_base64_image(&inner.session_dir, media_type, data)
{
images.push(image);
}
}
}
}
Some("functionCallOutput") => save_data_images(
&inner.session_dir,
item.get("output").and_then(Value::as_array),
"input_image",
"image_url",
&mut images,
),
Some("imageView") => {
if let Some(path) = item.get("path").and_then(Value::as_str)
&& let Some(image) = save_viewed_image(inner, path)
{
images.push(image);
}
}
_ => {}
}
images
.into_iter()
.map(|image| Event::Image {
image,
about: Some(id.to_string()),
})
.collect()
}
fn save_data_images(
session_dir: &Path,
parts: Option<&Vec<Value>>,
image_kind: &str,
url_field: &str,
images: &mut Vec<String>,
) {
for part in parts.into_iter().flatten() {
if part.get("type").and_then(Value::as_str) != Some(image_kind) {
continue;
}
let Some(url) = part.get(url_field).and_then(Value::as_str) else {
continue;
};
let Some((header, data)) = url.split_once(',') else {
continue;
};
let Some(media_type) = header
.strip_prefix("data:")
.and_then(|header| header.strip_suffix(";base64"))
else {
continue;
};
if let Some(image) = save_base64_image(session_dir, media_type, data) {
images.push(image);
}
}
}
fn save_base64_image(session_dir: &Path, media_type: &str, data: &str) -> Option<String> {
use base64::Engine;
let bytes = base64::engine::general_purpose::STANDARD
.decode(data)
.ok()?;
store_image(session_dir, media_type, &bytes)
}
fn save_viewed_image(inner: &Inner, path: &str) -> Option<String> {
let media_type = crate::media::media_type_for(path).unwrap_or("image/png");
let bytes = match &inner.transport {
Transport::Here => std::fs::read(path).map_err(anyhow::Error::from),
Transport::Ssh { .. } => tokio::task::block_in_place(|| {
inner.transport.capture_bytes_blocking(&Launch::new(
"cat",
vec![path.to_string()],
None,
))
}),
};
match bytes {
Ok(bytes) => store_image(&inner.session_dir, media_type, &bytes),
Err(err) => {
tracing::error!("couldn't save image Codex read from {path}: {err:#}");
None
}
}
}
fn handle_response(inner: &Arc<Inner>, line: &Value) {
let Some(id) = line.get("id").and_then(Value::as_str) else {
return;
};
let mut state = inner.state.lock().unwrap();
if state.initialize_request.as_deref() == Some(id) {
state.initialize_request = None;
if let Some(error) = response_error(line) {
drop(state);
protocol_error(inner, error);
return;
}
let request = request_id();
state.thread_request = Some(request.clone());
drop(state);
save_state(inner);
send_json(inner, json!({"method": "initialized"}));
let mut params = thread_params(inner);
let method = match read_thread(&inner.session_dir) {
Some(thread) => {
params["threadId"] = Value::String(thread);
// This app already owns and pages its common transcript. Hydrating the complete
// Codex history here sends it a second time, including every base64 screenshot.
params["excludeTurns"] = Value::Bool(true);
"thread/resume"
}
None => "thread/start",
};
send_json(
inner,
json!({"id": request, "method": method, "params": params}),
);
return;
}
if state.thread_request.as_deref() == Some(id) {
state.thread_request = None;
if let Some(error) = response_error(line) {
// `thread/start` returns an id before the first turn creates its rollout. If the
// app-server is restarted in that window, the id we correctly persisted is not one
// Codex can resume. A rollout can also disappear independently of ai-app. In either
// case our common transcript still exists, so start a blank Codex thread and mark the
// point where the model's context stopped instead of stranding every queued message.
if missing_thread(&error) && read_thread(&inner.session_dir).is_some() {
protocol_error(inner, error.clone());
let path = inner.session_dir.join(THREAD_FILE);
match std::fs::remove_file(&path) {
Ok(()) => {}
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
Err(err) => {
drop(state);
protocol_error(
inner,
format!("couldn't forget the missing Codex thread: {err}"),
);
return;
}
}
let request = request_id();
state.thread_request = Some(request.clone());
drop(state);
save_state(inner);
send_json(
inner,
json!({
"id": request,
"method": "thread/start",
"params": thread_params(inner)
}),
);
let _ = inner.sink.send(Event::Cleared);
return;
}
drop(state);
protocol_error(inner, error);
return;
}
let thread = line.pointer("/result/thread/id").and_then(Value::as_str);
drop(state);
if let Some(thread) = thread {
write_thread(&inner.session_dir, thread);
}
save_state(inner);
dispatch_waiting(inner);
return;
}
let Some(at) = state.pending.iter().position(|request| request.id == id) else {
return;
};
let request = state.pending.remove(at);
if let Some(error) = response_error(line) {
let message = state
.sent
.iter()
.position(|message| message.client_id == request.client_id)
.and_then(|at| state.sent.remove(at));
if request.kind == RequestKind::Steer && active_turn_not_steerable(line) {
if let Some(message) = message {
state.waiting.push_back(message);
}
let between_turns = state.active_turn.is_none();
drop(state);
save_state(inner);
if between_turns {
dispatch_waiting(inner);
}
return;
}
state.running = state.active_turn.is_some()
|| state
.pending
.iter()
.any(|pending| pending.kind == RequestKind::Start)
|| !state.waiting.is_empty();
let idle = !state.running;
drop(state);
save_state(inner);
if let Some(message) = message
&& !message.id.is_empty()
{
let _ = inner.sink.send(Event::MessageDropped { id: message.id });
}
protocol_error(inner, error);
if idle {
let _ = inner.sink.send(Event::Status {
state: SessionStatus::Idle,
});
}
return;
}
if request.kind == RequestKind::Start
&& let Some(turn) = line.pointer("/result/turn/id").and_then(Value::as_str)
{
state.active_turn = Some(turn.to_string());
}
drop(state);
save_state(inner);
}
fn announce_user(inner: &Inner, item: &Value) {
let Some(client_id) = item.get("clientId").and_then(Value::as_str) else {
return;
};
let message = {
let mut state = inner.state.lock().unwrap();
state
.sent
.iter()
.position(|message| message.client_id == client_id)
.and_then(|at| state.sent.remove(at))
};
let Some(message) = message else {
return;
};
save_state(inner);
let _ = inner.sink.send(Event::MessageTaken {
id: (!message.id.is_empty()).then_some(message.id),
text: message.text,
attachments: message.attachments,
});
}
fn response_error(line: &Value) -> Option<String> {
line.pointer("/error/message")
.and_then(Value::as_str)
.map(str::to_string)
}
fn missing_thread(message: &str) -> bool {
let message = message.to_ascii_lowercase();
message.contains("thread not found") || message.contains("no rollout found for thread id")
}
fn active_turn_not_steerable(line: &Value) -> bool {
line.pointer("/error/data/codexErrorInfo/activeTurnNotSteerable")
.is_some()
|| line.to_string().contains("activeTurnNotSteerable")
}
fn protocol_error(inner: &Inner, message: String) {
let _ = inner.sink.send(Event::Error {
message: format!("Codex refused a request: {message}"),
});
}
fn attachment_path(session_dir: &Path, id: &str) -> Result<PathBuf> {
if id.contains("..")
|| !id
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '-' | '_'))
{
anyhow::bail!("invalid attachment id");
}
Ok(session_dir.join("attachments").join(id))
}
pub(super) fn read_thread(session_dir: &Path) -> Option<String> {
serde_json::from_str::<Value>(&std::fs::read_to_string(session_dir.join(THREAD_FILE)).ok()?)
.ok()?
.get("threadId")?
.as_str()
.map(str::to_string)
}
fn write_thread(session_dir: &Path, thread: &str) {
let path = session_dir.join(THREAD_FILE);
if let Err(err) = std::fs::write(&path, json!({"threadId": thread}).to_string()) {
tracing::error!(
"couldn't persist Codex thread id to {}: {err}",
path.display()
);
}
}
fn read_state(session_dir: &Path) -> ProtocolState {
if let Some(state) = std::fs::read_to_string(session_dir.join(STATE_FILE))
.ok()
.and_then(|text| serde_json::from_str(&text).ok())
{
return state;
}
// Upgrade messages left by the old per-turn exec driver.
std::fs::read_to_string(session_dir.join("codex-queue.json"))
.ok()
.and_then(|text| serde_json::from_str::<Value>(&text).ok())
.and_then(|value| serde_json::from_value(value["waiting"].clone()).ok())
.map(|waiting| ProtocolState {
waiting,
..ProtocolState::default()
})
.unwrap_or_default()
}
fn save_state(inner: &Inner) {
let path = inner.session_dir.join(STATE_FILE);
let text = match serde_json::to_string(&*inner.state.lock().unwrap()) {
Ok(text) => text,
Err(err) => {
tracing::error!("couldn't serialize the Codex protocol state: {err}");
return;
}
};
let temporary = path.with_extension("json.new");
let written =
std::fs::write(&temporary, text).and_then(|()| std::fs::rename(&temporary, &path));
if let Err(err) = written {
tracing::error!(
"couldn't persist Codex protocol state to {}: {err}",
path.display()
);
let _ = std::fs::remove_file(temporary);
}
}
fn stderr_tail(path: &Path) -> String {
std::fs::read_to_string(path)
.unwrap_or_default()
.lines()
.rev()
.take(20)
.collect::<Vec<_>>()
.into_iter()
.rev()
.collect::<Vec<_>>()
.join("\n")
}
const DELETE_TRANSCRIPT_SCRIPT: &str = r#"
state=missing
for f in "$HOME"/.codex/sessions/*/*/*/rollout-*-${1}.jsonl; do
[ -f "$f" ] || continue
if rm -f "$f"; then state=deleted; else state=failed; fi
done
printf '%s\n' "$state"
"#;
const READ_CONTEXT_SCRIPT: &str = r#"
for f in "$HOME"/.codex/sessions/*/*/*/rollout-*-${1}.jsonl; do
[ -f "$f" ] || continue
grep '"type":"token_count"' "$f" | tail -1
exit 0
done
"#;
/// The last measured prompt size from Codex's own rollout, for a session whose common transcript
/// predates context events. This is the same `last.inputTokens` app-server reports live, under the
/// rollout writer's snake-case names.
pub async fn context_of(transport: &Transport, id: &str) -> Option<u64> {
if !valid_thread_id(id) {
return None;
}
let launch = Launch::new(
"sh",
vec![
"-c".to_string(),
READ_CONTEXT_SCRIPT.to_string(),
"sh".to_string(),
id.to_string(),
],
None,
);
let line = transport.capture(&launch).await.ok()?;
context_from_rollout(&line)
}
fn context_from_rollout(line: &str) -> Option<u64> {
serde_json::from_str::<Value>(line)
.ok()?
.pointer("/payload/info/last_token_usage/input_tokens")?
.as_u64()
}
/// Removes the rollout whose suffix is this thread id.
pub async fn delete_transcript(transport: &Transport, id: &str) -> Result<()> {
if !valid_thread_id(id) {
anyhow::bail!("not a Codex thread id: {id}");
}
let launch = Launch::new(
"sh",
vec![
"-c".to_string(),
DELETE_TRANSCRIPT_SCRIPT.to_string(),
"sh".to_string(),
id.to_string(),
],
None,
);
match transport.capture(&launch).await?.trim() {
"deleted" => Ok(()),
"missing" => anyhow::bail!("no Codex thread {id} on that machine"),
_ => anyhow::bail!("couldn't remove Codex thread {id}"),
}
}
fn valid_thread_id(id: &str) -> bool {
!id.is_empty()
&& id
.chars()
.all(|character| character.is_ascii_hexdigit() || character == '-')
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn old_exec_queues_are_adopted() {
let dir = tempfile::tempdir().expect("tempdir");
std::fs::write(
dir.path().join("codex-queue.json"),
r#"{"waiting":[{"id":"q1","text":"next","attachments":[]}]}"#,
)
.expect("queue");
let state = read_state(dir.path());
assert_eq!(state.waiting.len(), 1);
assert_eq!(state.waiting[0].id, "q1");
}
#[test]
fn a_structured_image_is_saved_outside_the_transcript() {
let dir = tempfile::tempdir().expect("tempdir");
let parts = vec![json!({
"type": "inputImage",
"imageUrl": "data:image/png;base64,aGVsbG8="
})];
let mut images = Vec::new();
save_data_images(
dir.path(),
Some(&parts),
"inputImage",
"imageUrl",
&mut images,
);
assert_eq!(images.len(), 1);
assert_eq!(
std::fs::read(dir.path().join("files").join(&images[0])).expect("saved image"),
b"hello"
);
}
#[test]
fn an_inline_image_carries_its_bytes_to_a_remote_codex() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("shot.png");
std::fs::write(&path, b"hello").expect("image");
assert_eq!(
inline_image(&path, "image/png").expect("inline image"),
json!({"type": "image", "url": "data:image/png;base64,aGVsbG8="})
);
}
#[test]
fn transcript_delete_resolves_only_the_named_codex_rollout() {
use std::process::Command;
let home = tempfile::tempdir().expect("home");
let sessions = home.path().join(".codex/sessions/2026/09/09");
std::fs::create_dir_all(&sessions).expect("sessions");
let wanted = sessions.join("rollout-2026-09-09T12-00-00-abcd-1234.jsonl");
let other = sessions.join("rollout-2026-09-09T12-00-01-abcd-5678.jsonl");
std::fs::write(&wanted, "wanted").expect("wanted");
std::fs::write(&other, "other").expect("other");
let output = Command::new("sh")
.args(["-c", DELETE_TRANSCRIPT_SCRIPT, "sh", "abcd-1234"])
.env("HOME", home.path())
.output()
.expect("delete script");
assert_eq!(String::from_utf8_lossy(&output.stdout).trim(), "deleted");
assert!(!wanted.exists());
assert!(other.exists());
assert!(!valid_thread_id("../../something"));
}
#[test]
fn context_is_read_from_the_last_codex_model_call() {
assert_eq!(
context_from_rollout(
r#"{"type":"event_msg","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":118866,"cached_input_tokens":118144,"output_tokens":37,"total_tokens":118903}}}}"#
),
Some(118_866)
);
assert_eq!(context_from_rollout(""), None);
}
#[test]
fn both_missing_thread_failures_are_recognised() {
assert!(missing_thread(
"thread not found: 01a099bd-2e0f-7442-b750-a9f33998c8c5"
));
assert!(missing_thread(
"no rollout found for thread id 01a099bd-2e0f-7442-b750-a9f33998c8c5"
));
assert!(!missing_thread("thread is already running"));
}
#[test]
fn a_missing_thread_is_reported_then_replaced() {
let dir = tempfile::tempdir().expect("tempdir");
write_thread(dir.path(), "missing-thread");
let (sink, mut events) = mpsc::unbounded_channel();
let (to_child, mut requests) = mpsc::unbounded_channel();
let inner = Arc::new(Inner {
sink,
state: Mutex::new(ProtocolState {
thread_request: Some("resume-request".to_string()),
..ProtocolState::default()
}),
settings: Mutex::new(Settings {
model: None,
permission_mode: None,
effort: None,
}),
to_child,
transport: Transport::Here,
session_dir: dir.path().to_path_buf(),
subagents: Arc::new(Subagents::new(dir.path().to_path_buf())),
reading: AtomicBool::new(true),
});
handle_response(
&inner,
&json!({
"id": "resume-request",
"error": {"message": "thread not found: missing-thread"}
}),
);
assert!(matches!(
events.try_recv().expect("visible refusal"),
Event::Error { message } if message ==
"Codex refused a request: thread not found: missing-thread"
));
assert_eq!(events.try_recv().expect("context boundary"), Event::Cleared);
assert_eq!(read_thread(dir.path()), None);
let request: Value =
serde_json::from_str(&requests.try_recv().expect("replacement thread request"))
.expect("request json");
assert_eq!(request["method"], "thread/start");
let request_id = request["id"].as_str().expect("request id");
handle_response(
&inner,
&json!({"id": request_id, "result": {"thread": {"id": "replacement-thread"}}}),
);
assert_eq!(
read_thread(dir.path()).as_deref(),
Some("replacement-thread")
);
}
#[test]
fn a_message_waiting_for_a_thread_is_durable() {
let dir = tempfile::tempdir().expect("tempdir");
let (sink, mut events) = mpsc::unbounded_channel();
let (to_child, _requests) = mpsc::unbounded_channel();
let inner = Arc::new(Inner {
sink,
state: Mutex::new(ProtocolState {
thread_request: Some("replacement-request".to_string()),
..ProtocolState::default()
}),
settings: Mutex::new(Settings {
model: None,
permission_mode: None,
effort: None,
}),
to_child,
transport: Transport::Here,
session_dir: dir.path().to_path_buf(),
subagents: Arc::new(Subagents::new(dir.path().to_path_buf())),
reading: AtomicBool::new(true),
});
let driver = CodexDriver {
inner: Arc::clone(&inner),
};
driver.send_user_message("still send this".to_string(), Vec::new());
assert!(matches!(
events.try_recv().expect("durable queue event"),
Event::MessageQueued { text, .. } if text == "still send this"
));
let state = inner.state.lock().unwrap();
assert_eq!(state.waiting.len(), 1);
assert!(!state.waiting[0].id.is_empty());
}
#[test]
fn clear_defers_the_replacement_thread_until_its_first_message() {
let dir = tempfile::tempdir().expect("tempdir");
write_thread(dir.path(), "old-thread");
let (sink, mut events) = mpsc::unbounded_channel();
let (to_child, mut requests) = mpsc::unbounded_channel();
let inner = Arc::new(Inner {
sink,
state: Mutex::new(ProtocolState::default()),
settings: Mutex::new(Settings {
model: None,
permission_mode: None,
effort: None,
}),
to_child,
transport: Transport::Here,
session_dir: dir.path().to_path_buf(),
subagents: Arc::new(Subagents::new(dir.path().to_path_buf())),
reading: AtomicBool::new(true),
});
let driver = CodexDriver {
inner: Arc::clone(&inner),
};
driver.clear();
assert_eq!(events.try_recv().expect("clear boundary"), Event::Cleared);
assert_eq!(read_thread(dir.path()), None);
assert!(
requests.try_recv().is_err(),
"clear eagerly started a thread"
);
driver.send_user_message("first new message".to_string(), Vec::new());
assert!(matches!(
events.try_recv().expect("durable queue event"),
Event::MessageQueued { text, .. } if text == "first new message"
));
let start: Value = serde_json::from_str(&requests.try_recv().expect("thread start"))
.expect("request json");
assert_eq!(start["method"], "thread/start");
let request_id = start["id"].as_str().expect("request id");
handle_response(
&inner,
&json!({"id": request_id, "result": {"thread": {"id": "new-thread"}}}),
);
let turn: Value =
serde_json::from_str(&requests.try_recv().expect("first turn")).expect("request json");
assert_eq!(turn["method"], "turn/start");
assert_eq!(turn["params"]["threadId"], "new-thread");
}
#[test]
fn a_replacement_thread_notification_releases_its_waiting_message() {
let dir = tempfile::tempdir().expect("tempdir");
let (sink, _events) = mpsc::unbounded_channel();
let (to_child, mut requests) = mpsc::unbounded_channel();
let inner = Arc::new(Inner {
sink,
state: Mutex::new(ProtocolState {
waiting: VecDeque::from([Waiting {
id: "queued-message".to_string(),
client_id: "client-message".to_string(),
steering: false,
text: "deliver me".to_string(),
attachments: Vec::new(),
}]),
..ProtocolState::default()
}),
settings: Mutex::new(Settings {
model: None,
permission_mode: None,
effort: None,
}),
to_child,
transport: Transport::Here,
session_dir: dir.path().to_path_buf(),
subagents: Arc::new(Subagents::new(dir.path().to_path_buf())),
reading: AtomicBool::new(true),
});
let mut translator = Translator::new(
Arc::clone(&inner.subagents),
Some("missing-thread".to_string()),
false,
);
handle_line(
&inner,
&mut translator,
&json!({
"method": "thread/started",
"params": {"thread": {
"id": "replacement-thread",
"parentThreadId": null
}}
}),
);
assert_eq!(
read_thread(dir.path()).as_deref(),
Some("replacement-thread")
);
let request: Value = serde_json::from_str(&requests.try_recv().expect("turn request"))
.expect("request json");
assert_eq!(request["method"], "turn/start");
assert_eq!(request["params"]["threadId"], "replacement-thread");
assert_eq!(request["params"]["clientUserMessageId"], "client-message");
let state = inner.state.lock().unwrap();
assert!(state.waiting.is_empty());
assert_eq!(state.sent.len(), 1);
}
}