From e50e3973151b068324118eeb38ad65f4789620da Mon Sep 17 00:00:00 2001 From: meh Date: Sun, 23 Aug 2026 22:47:34 +0700 Subject: [PATCH] fix(codex): interrupt and bound resumed turns --- Cargo.toml | 2 +- src/codex_app_server.rs | 170 +++++++++++++++++++++++++++++++++++++++- src/lib.rs | 2 +- src/request.rs | 13 +++ src/run.rs | 74 +++++++++++++++++ tests/live.rs | 49 +++++++++++- 6 files changed, 306 insertions(+), 4 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index efda928..c7f6dc6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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. diff --git a/src/codex_app_server.rs b/src/codex_app_server.rs index c21ea6a..8beaf2a 100644 --- a/src/codex_app_server.rs +++ b/src/codex_app_server.rs @@ -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 { @@ -64,6 +65,7 @@ pub(crate) struct Protocol { pub failure: Option, pending: HashMap, pending_steers: HashSet, + interrupt_requested: bool, next_id: u64, } @@ -78,6 +80,7 @@ impl Protocol { failure: None, pending: HashMap::new(), pending_steers: HashSet::new(), + interrupt_requested: false, next_id: 10, } } @@ -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", @@ -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; } @@ -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 @@ -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 { + 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 { let pending = self.pending.remove(id)?; @@ -399,6 +462,20 @@ impl Protocol { } } +fn active_turn_id(value: &Value) -> Option { + 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") } @@ -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()); @@ -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(); diff --git a/src/lib.rs b/src/lib.rs index 0bc3f5a..fe7ea0a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -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}; diff --git a/src/request.rs b/src/request.rs index e8bd5c4..960e650 100644 --- a/src/request.rs +++ b/src/request.rs @@ -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. @@ -95,6 +107,7 @@ impl Request { timeout: None, binding: None, is_command: false, + operation: Operation::Run, } } diff --git a/src/run.rs b/src/run.rs index 331b5d3..2a4ee89 100644 --- a/src/run.rs +++ b/src/run.rs @@ -468,6 +468,36 @@ pub async fn run(request: &Request) -> Result { 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 { + 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. @@ -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 { @@ -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 }); @@ -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, + 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::(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( diff --git a/tests/live.rs b/tests/live.rs index 3052aa7..5e55d90 100644 --- a/tests/live.rs +++ b/tests/live.rs @@ -15,7 +15,7 @@ use std::time::Duration; use agent_abstraction::{ Agent, AuthState, AuthStatus, EnvPolicy, Event, Format, Permission, Probe, Request, - SessionStore, VersionStatus, run, stream, + SessionStore, VersionStatus, interrupt, run, stream, }; /// A prompt with exactly one correct answer, so the assertion is about the @@ -311,6 +311,53 @@ async fn codex_reports_account_usage_without_a_terminal() { assert!(window.window_minutes.is_some(), "{window:?}"); } +/// Provider-level recovery without a model call. The caller supplies a Codex +/// thread that currently has an orphaned `inProgress` turn; the test proves the +/// one-shot control path reaches `turn/interrupt` without submitting a prompt. +#[tokio::test] +#[ignore = "requires AGENT_ABSTRACTION_INTERRUPT_SESSION naming an active Codex turn"] +async fn codex_interrupts_an_orphaned_active_turn() { + if !available(Agent::Codex) { + return; + } + let session = std::env::var("AGENT_ABSTRACTION_INTERRUPT_SESSION") + .expect("set AGENT_ABSTRACTION_INTERRUPT_SESSION to an active Codex thread id"); + let request = Request::new(Agent::Codex, "") + .resume(session) + .permission(Permission::ReadOnly) + .timeout(Duration::from_secs(20)); + + assert!( + interrupt(&request) + .await + .expect("Codex session interruption failed"), + "the supplied thread had no active turn" + ); +} + +/// Long-lived idle threads used to exceed the bounded JSON line on resume and +/// time out before the caller could learn that no recovery was needed. +#[tokio::test] +#[ignore = "requires AGENT_ABSTRACTION_IDLE_SESSION naming an idle Codex thread id"] +async fn codex_reads_a_long_idle_session_without_starting_a_turn() { + if !available(Agent::Codex) { + return; + } + let session = std::env::var("AGENT_ABSTRACTION_IDLE_SESSION") + .expect("set AGENT_ABSTRACTION_IDLE_SESSION to an idle Codex thread id"); + let request = Request::new(Agent::Codex, "") + .resume(session) + .permission(Permission::ReadOnly) + .timeout(Duration::from_secs(20)); + + assert!( + !interrupt(&request) + .await + .expect("Codex idle-session inspection failed"), + "the supplied thread unexpectedly had an active turn" + ); +} + /// Claude's context tracker, end to end: it is the one agent that reports the /// window alongside the tokens, so a share of it is knowable. #[tokio::test]