Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "agent-abstraction"
version = "0.4.18"
version = "0.4.19"
edition = "2024"
# The floor edition 2024 requires, and where the strictest dependencies (uuid,
# getrandom) sit. Derived from the dependency graph rather than compile-tested.
Expand Down
170 changes: 169 additions & 1 deletion src/codex_app_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ use crate::request::Request;

const OPEN_ID: u64 = 2;
const TURN_ID: u64 = 3;
const INTERRUPT_ID: u64 = 4;

#[derive(Debug, Clone)]
enum PendingApproval {
Expand Down Expand Up @@ -64,6 +65,7 @@ pub(crate) struct Protocol {
pub failure: Option<String>,
pending: HashMap<String, PendingApproval>,
pending_steers: HashSet<u64>,
interrupt_requested: bool,
next_id: u64,
}

Expand All @@ -78,6 +80,7 @@ impl Protocol {
failure: None,
pending: HashMap::new(),
pending_steers: HashSet::new(),
interrupt_requested: false,
next_id: 10,
}
}
Expand Down Expand Up @@ -121,6 +124,22 @@ impl Protocol {
}),
Continue::Resume(thread_id) => {
params["threadId"] = json!(thread_id);
// Verified against codex-cli 0.147.0. The default resume reply
// includes every turn and item. Long-lived threads exceed this
// crate's bounded JSON line and disappear as an unparsed line,
// leaving the run waiting forever for a response it received.
// Ordinary runs need no history. One-shot recovery needs only
// the newest turn id and status, never its items.
if self.request.operation == crate::request::Operation::Interrupt {
params["excludeTurns"] = json!(true);
params["initialTurnsPage"] = json!({
"limit": 1,
"sortDirection": "desc",
"itemsView": "notLoaded",
});
} else {
params["excludeTurns"] = json!(true);
}
json!({
"id": OPEN_ID,
"method": "thread/resume",
Expand Down Expand Up @@ -208,7 +227,21 @@ impl Protocol {
session: thread_id.clone(),
model,
});
step.writes.push(self.start_turn(&thread_id));
// Codex 0.147.0 returns the resumed thread's turns here. A
// process killed without `turn/interrupt` leaves its last turn
// `inProgress`; starting another then fails with "run already
// exists" forever. One-shot interruption targets that exact
// turn and deliberately starts no replacement prompt.
self.turn_id = active_turn_id(value);
if self.request.operation == crate::request::Operation::Interrupt {
if let Some(interrupt) = self.interrupt() {
step.writes.push(interrupt);
} else {
self.finished = true;
}
} else {
step.writes.push(self.start_turn(&thread_id));
}
}
return step;
}
Expand All @@ -223,6 +256,17 @@ impl Protocol {
return step;
}

if value.get("id").and_then(Value::as_u64) == Some(INTERRUPT_ID) && self.interrupt_requested
{
if let Some(error) = rpc_error(value) {
self.failure = Some(error);
} else {
self.terminal.stop = Stop::Other("interrupted".into());
}
self.finished = true;
return step;
}

// Unlike streamed notifications, `turn/steer` has a decisive JSON-RPC
// response: `{ result: { turnId } }` means Codex accepted the input,
// while an error means it did not enter the turn. Keep that receipt
Expand Down Expand Up @@ -377,6 +421,25 @@ impl Protocol {
})
}

/// Encode Codex's provider-level turn interruption once both ids are known.
///
/// Verified against the schema generated by codex-cli 0.147.0. Killing the
/// app-server process alone does not interrupt the server-owned turn and
/// leaves the thread permanently refusing a resumed `turn/start`.
pub fn interrupt(&mut self) -> Option<String> {
let thread_id = self.thread_id.as_ref()?;
let turn_id = self.turn_id.as_ref()?;
self.interrupt_requested = true;
Some(wire(&json!({
"id": INTERRUPT_ID,
"method": "turn/interrupt",
"params": {
"threadId": thread_id,
"turnId": turn_id,
},
})))
}

/// Encode the response to a server-side permission request.
pub fn respond(&mut self, id: &str, decision: &Decision) -> Option<String> {
let pending = self.pending.remove(id)?;
Expand All @@ -399,6 +462,20 @@ impl Protocol {
}
}

fn active_turn_id(value: &Value) -> Option<String> {
let turns = value
.pointer("/result/initialTurnsPage/data")
.or_else(|| value.pointer("/result/thread/turns"))?
.as_array()?;
turns
.iter()
.rev()
.find(|turn| turn.get("status").and_then(Value::as_str) == Some("inProgress"))?
.get("id")?
.as_str()
.map(str::to_string)
}

fn wire(value: &Value) -> String {
format!("{value}\n")
}
Expand Down Expand Up @@ -583,6 +660,32 @@ mod tests {
);
}

#[test]
fn ordinary_resume_excludes_unbounded_history() {
let protocol = Protocol::new(request().resume("thread-long"));
let opening = protocol.opening();
let open: Value = serde_json::from_str(opening.last().expect("open request")).unwrap();

assert_eq!(open["method"], "thread/resume");
assert_eq!(open["params"]["excludeTurns"], true);
assert!(open["params"].get("initialTurnsPage").is_none());
}

#[test]
fn interrupt_resume_requests_only_the_newest_turn_without_items() {
let mut request = request().resume("thread-long");
request.operation = crate::request::Operation::Interrupt;
let protocol = Protocol::new(request);
let opening = protocol.opening();
let open: Value = serde_json::from_str(opening.last().expect("open request")).unwrap();

assert_eq!(open["method"], "thread/resume");
assert_eq!(open["params"]["excludeTurns"], true);
assert_eq!(open["params"]["initialTurnsPage"]["limit"], 1);
assert_eq!(open["params"]["initialTurnsPage"]["sortDirection"], "desc");
assert_eq!(open["params"]["initialTurnsPage"]["itemsView"], "notLoaded");
}

#[test]
fn a_thread_response_starts_the_turn_and_exposes_the_session() {
let mut protocol = Protocol::new(request());
Expand Down Expand Up @@ -708,6 +811,71 @@ mod tests {
protocol
}

#[test]
fn interruption_targets_the_active_codex_turn() {
let mut protocol = running_protocol();
let encoded = protocol.interrupt().expect("interrupt request");
let wire: Value = serde_json::from_str(&encoded).unwrap();

assert_eq!(wire["method"], "turn/interrupt");
assert_eq!(wire["params"]["threadId"], "thread-7");
assert_eq!(wire["params"]["turnId"], "turn-9");

protocol.push(&json!({"id": INTERRUPT_ID, "result": {}}));
assert!(protocol.finished);
assert_eq!(protocol.terminal.stop, Stop::Other("interrupted".into()));
}

#[test]
fn interrupt_only_resume_repairs_an_orphan_without_starting_a_turn() {
let mut request = request().resume("thread-orphan");
request.operation = crate::request::Operation::Interrupt;
let mut protocol = Protocol::new(request);
let step = protocol.push(&json!({
"id": OPEN_ID,
"result": {
"thread": {
"id": "thread-orphan",
"turns": []
},
"initialTurnsPage": {"data": [
{"id": "turn-done", "status": "completed"},
{"id": "turn-stuck", "status": "inProgress"}
]},
"model": "gpt-5.6-sol"
}
}));

assert_eq!(step.writes.len(), 1);
let wire: Value = serde_json::from_str(&step.writes[0]).unwrap();
assert_eq!(wire["method"], "turn/interrupt");
assert_eq!(wire["params"]["turnId"], "turn-stuck");
}

#[test]
fn interrupt_only_resume_is_a_noop_without_an_active_turn() {
let mut request = request().resume("thread-idle");
request.operation = crate::request::Operation::Interrupt;
let mut protocol = Protocol::new(request);
let step = protocol.push(&json!({
"id": OPEN_ID,
"result": {
"thread": {
"id": "thread-idle",
"turns": []
},
"initialTurnsPage": {"data": [
{"id": "turn-done", "status": "completed"}
]},
"model": "gpt-5.6-sol"
}
}));

assert!(step.writes.is_empty());
assert!(protocol.finished);
assert_eq!(protocol.terminal.stop, Stop::Completed);
}

#[test]
fn a_steer_resolves_only_from_its_acceptance_response() {
let mut protocol = running_protocol();
Expand Down
2 changes: 1 addition & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -112,5 +112,5 @@ pub use model::{Kind, Model, Source, Verified};
pub use outcome::{Outcome, RateLimit, Stop, Usage};
pub use probe::{Probe, Version, VersionStatus};
pub use request::Request;
pub use run::{Run, RunControl, run, stream};
pub use run::{Run, RunControl, interrupt, run, stream};
pub use session::{Phase, SessionRecord, SessionStore};
13 changes: 13 additions & 0 deletions src/request.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,18 @@ pub struct Request {
/// Set by [`Request::command`]: the prompt is a slash command, so the
/// capability check refuses agents that have no command vocabulary.
pub(crate) is_command: bool,
/// Whether this request runs a prompt or only repairs a resumed session.
pub(crate) operation: Operation,
}

/// The lifecycle operation a request asks the provider transport to perform.
#[derive(Debug, Clone, Copy, Default, Eq, PartialEq)]
pub(crate) enum Operation {
/// Submit the request's prompt as a turn.
#[default]
Run,
/// Interrupt an orphaned active turn without starting a replacement.
Interrupt,
}

/// A named session this run is attached to.
Expand Down Expand Up @@ -95,6 +107,7 @@ impl Request {
timeout: None,
binding: None,
is_command: false,
operation: Operation::Run,
}
}

Expand Down
74 changes: 74 additions & 0 deletions src/run.rs
Original file line number Diff line number Diff line change
Expand Up @@ -468,6 +468,36 @@ pub async fn run(request: &Request) -> Result<Outcome> {
stream(request)?.finish().await
}

/// Interrupt an orphaned active Codex turn without starting a replacement.
///
/// The request must resume a Codex session. Its binary, environment policy,
/// environment overrides, and working directory are retained, while its prompt
/// is never submitted. Returns `true` when Codex reported an active turn and
/// accepted `turn/interrupt`, or `false` when the session had no active turn.
///
/// # Errors
/// Returns [`Error::Unsupported`] for another provider or a request that does
/// not resume a session, and the ordinary spawn/protocol errors otherwise.
pub async fn interrupt(request: &Request) -> Result<bool> {
if request.agent != crate::Agent::Codex {
return Err(Error::Unsupported {
agent: request.agent,
what: "provider-level session interruption",
});
}
if !matches!(request.cont, crate::agent::Continue::Resume(_)) {
return Err(Error::Unsupported {
agent: request.agent,
what: "interrupting without a resumed session",
});
}
let mut request = request.clone();
request.duplex = true;
request.operation = crate::request::Operation::Interrupt;
let outcome = run(&request).await?;
Ok(matches!(outcome.stop, Stop::Other(ref reason) if reason == "interrupted"))
}

/// Start `request`, returning a handle that streams its events.
///
/// Returns as soon as the child is spawned; the work proceeds on a task.
Expand Down Expand Up @@ -1236,6 +1266,8 @@ async fn drive_codex_app_server(
}
() = &mut deadline => {
let partial = protocol.terminal.text.clone();
interrupt_codex_turn(&mut protocol, &mut stdin, &mut reader, &mut line, &mut raw)
.await;
shut_down(&mut child, stderr_task).await;
reaped.store(true, std::sync::atomic::Ordering::SeqCst);
return Err(Error::Timeout {
Expand All @@ -1245,6 +1277,8 @@ async fn drive_codex_app_server(
});
}
_ = &mut cancel => {
interrupt_codex_turn(&mut protocol, &mut stdin, &mut reader, &mut line, &mut raw)
.await;
shut_down(&mut child, stderr_task).await;
reaped.store(true, std::sync::atomic::Ordering::SeqCst);
return Err(Error::Cancelled { bin });
Expand Down Expand Up @@ -1308,6 +1342,46 @@ async fn drive_codex_app_server(
})
}

/// Ask Codex to settle its server-owned turn before its local app-server dies.
///
/// Two seconds bounds a provider that is already wedged. Failure still falls
/// through to process-group teardown, but a healthy app-server gets the
/// `turn/interrupt` it needs to make this thread resumable again.
async fn interrupt_codex_turn(
protocol: &mut crate::codex_app_server::Protocol,
stdin: &mut tokio::process::ChildStdin,
reader: &mut BufReader<tokio::process::ChildStdout>,
line: &mut String,
raw: &mut String,
) {
let Some(encoded) = protocol.interrupt() else {
return;
};
if stdin.write_all(encoded.as_bytes()).await.is_err() || stdin.flush().await.is_err() {
return;
}
let settle = async {
while !protocol.finished {
let Ok(Some(_)) = read_bounded_line(reader, line).await else {
break;
};
append_capped(raw, line);
if let Ok(value) = serde_json::from_str::<serde_json::Value>(line) {
let step = protocol.push(&value);
for write in step.writes {
if stdin.write_all(write.as_bytes()).await.is_err() {
return;
}
}
if stdin.flush().await.is_err() {
return;
}
}
}
};
let _ = tokio::time::timeout(std::time::Duration::from_secs(2), settle).await;
}

/// Write every control whose protocol ids are available, preserving earlier
/// messages until thread and turn startup have both completed.
async fn flush_codex_controls(
Expand Down
Loading
Loading