diff --git a/CHANGELOG.md b/CHANGELOG.md index 1226bf6..e1a1542 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,12 @@ Versioning and Keep a Changelog conventions. ## [Unreleased] +### Fixed + +- Codex resumes and active-writer forks omit historical turns from their replies, + so long conversations can continue without exceeding the protocol frame limit. + Saved provider context is preserved. + ### Changed - Preserve Windows system and profile environment variables when launching providers, @@ -24,6 +30,9 @@ Versioning and Keep a Changelog conventions. ### Added +- Add opt-in, bounded Codex app-server process reuse for in-process retained + runtimes, with strict runtime/configuration isolation, cold fallback at pool + capacity, idle expiry, and process-tree cleanup on cancellation or failure. - Contain observer panic-payload cleanup failures and disable failed observers across runtime clones. Report event-delivery wait separately from observed first-text latency, preserving bounded backpressure. diff --git a/README.md b/README.md index 6520a6e..231df69 100644 --- a/README.md +++ b/README.md @@ -146,6 +146,8 @@ and questions are left unanswered. Production applications should implement `InteractionHandler` and connect it to their durable approval workflow. For a complete walkthrough, read the [quickstart](docs/tutorials/quickstart.md). +For long Codex threads, read [Resume large Codex conversations](docs/how-to/resume-large-codex-conversations.md). + For sandbox setup, read [Manage Nono sandboxes](docs/how-to/nono.md) or [Implement a sandbox backend](docs/how-to/custom-sandbox.md). For durable profile updates and retry, read diff --git a/docs/adr/0004-prepared-provider-processes.md b/docs/adr/0004-prepared-provider-processes.md index 71f2370..e6a936d 100644 --- a/docs/adr/0004-prepared-provider-processes.md +++ b/docs/adr/0004-prepared-provider-processes.md @@ -1,6 +1,6 @@ # ADR 0004: Prepared provider processes -- Status: Proposed +- Status: Accepted (Codex retained turns implemented; explicit prewarming deferred) - Date: 2026-09-23 ## Problem @@ -16,42 +16,37 @@ first assistant text. A first output frame is not evidence of provider readiness or model request submission. Provider-specific readiness requires an explicit handshake acknowledgement, not a sleep or an empty model turn. -## Proposed ownership +## Ownership The SDK owns a bounded process supervisor and provider protocol state. Fleet owns when a user has selected enough configuration to prepare, durable conversation records, authorization, feature rollout, and presentation. Listing projects or conversations must never spawn provider processes. -Existing `AgentRuntime::run` and retained-client constructors preserve lazy, -one-process-per-turn behavior. A separate opt-in client configuration enables -retained processes. Unsupported adapters and remote protocol versions report a -typed capability error rather than silently claiming preparation succeeded. +Existing `AgentRuntime::run` and default retained-client constructors preserve +lazy, one-process-per-turn behavior. `codex_process_retention` opts the in-process +retained client into bounded app-server reuse. Custom adapters remain disabled +unless they implement the lifecycle contract. -## Proposed lifecycle +## Lifecycle 1. Acquire a logical runtime with project, provider, sandbox and launch settings. -2. Explicitly prepare it, without a prompt or invocation identifier. This starts - the process and completes the provider handshake; it does not call a model, - create a synthetic transcript message, or grant tool execution. -3. Send a turn. Lazy sending performs preparation automatically. Sending during - preparation joins the same bounded operation and submits exactly once. -4. On a successful turn, keep the provider connection and continue draining its +2. Send the first real turn. This lazily starts and initializes the app server, + then submits the prompt exactly once. +3. On a successful turn, keep the provider connection and continue draining its bounded event stream. Idle tool/background events belong to the runtime and must not be attached to the next invocation. -5. Dispose or expire the idle process, confirming process-tree teardown. Preserve +4. Dispose or expire the idle process, confirming process-tree teardown. Preserve session identity so later work can explicitly resume from provider persistence. -The process lifecycle distinguishes unprepared, queued, preparing, ready, busy, -failed, stopping and stopped. A logical runtime being acquired is not provider -readiness. Preparation errors expose delivery=not_sent. Failure after submission -preserves the existing accepted/possibly_sent semantics and must not replay a -prompt automatically. +Acquiring a logical runtime does not start a provider process or claim provider +readiness. Failure before submission remains `delivery=not_sent`; failure after +submission preserves accepted/possibly-sent semantics and never replays a prompt. -Preparation has a configurable deadline, bounded concurrency, global process -capacity and idle expiration. Active turns cannot be evicted to admit speculative -preparation. Abandoned preparations release their capacity. Disposing while -preparing prevents late successful readiness from resurrecting the runtime. +Retention has bounded turn concurrency, separate global process capacity and idle +expiration. Capacity exhaustion falls back to the ordinary cold-turn path. +Cancellation, timeout, dropped futures, disposal and ambiguous idle output poison +the connection and terminate its process tree before it can be reused. ## Configuration and credentials diff --git a/docs/how-to/resume-large-codex-conversations.md b/docs/how-to/resume-large-codex-conversations.md new file mode 100644 index 0000000..8150fa5 --- /dev/null +++ b/docs/how-to/resume-large-codex-conversations.md @@ -0,0 +1,11 @@ +# Resume large Codex conversations + +The app-server adapter requests `excludeTurns: true` when resuming a thread or +forking one after an active-writer conflict. This returns the metadata needed to +start the next turn without replaying the historical turns in one JSON frame. +Codex still uses the saved conversation context; this does not clear or truncate +history. New threads use the normal start request. + +This prevents growing history replies from hitting the default 2 MiB event-line +limit. The limit remains in place for other protocol events. Clients displaying +history should fetch it separately using the provider's paginated history APIs. diff --git a/docs/reference/retained-runtimes.md b/docs/reference/retained-runtimes.md index 8d4e6bf..ee9d4bc 100644 --- a/docs/reference/retained-runtimes.md +++ b/docs/reference/retained-runtimes.md @@ -11,6 +11,38 @@ reports `session_resume: true` and `retained_process: false`: Claude, Codex, and OpenCode can continue their provider-native sessions even though the CLI is currently relaunched for each turn. +Codex app-server reuse is explicit and in-process only: + +```rust +use std::time::Duration; +use temps_agent_runtime::providers::Codex; +use temps_agent_runtime::{AgentRuntime, CodexProcessRetention}; + +let mut builder = AgentRuntime::builder() + .codex_process_retention(CodexProcessRetention { + max_processes: 4, + idle_timeout: Duration::from_secs(120), + }); +builder.register(Codex::app_server()); +let runtime = builder.build()?; +``` + +Pass that runtime to `InProcessRuntimeClient::new`. Each `RuntimeId` owns at most +one process. The SDK compares the complete sandbox-wrapped command plus working +directory, model, reasoning, permission, harness, launch context, compaction and +sandbox requirements before reuse. Changed explicit credentials or environment +replace the process. Ambient inherited environment is read when a process starts; +applications that rotate ambient credentials must dispose the logical runtime or +rebuild the `AgentRuntime`. + +The process pool is bounded independently from acquired logical runtimes. A new +runtime that reaches capacity runs through the ordinary one-process turn path; +existing retained runtimes remain warm and usable without waiting for an idle +slot. +Any unsolicited idle frame, crash, cancellation, timeout, dropped turn future or +disposal retires the process tree. Late frames are correlated by native turn ID +and cannot enter a later invocation. + Use `RuntimeHandle::configuration_impact` before presenting a live setting change. The compatibility driver applies per-turn model, reasoning, permission, harness, launch-context, environment, and timeout changes live. Provider, diff --git a/site/src/content/docs.ts b/site/src/content/docs.ts index 4a353b2..d542d8f 100644 --- a/site/src/content/docs.ts +++ b/site/src/content/docs.ts @@ -1,6 +1,7 @@ import architecture from "../../../docs/explanation/architecture.md?raw"; import approvals from "../../../docs/how-to/persist-approvals.md?raw"; import commandExecution from "../../../docs/how-to/persist-command-execution.md?raw"; +import codexResume from "../../../docs/how-to/resume-large-codex-conversations.md?raw"; import persistence from "../../../docs/how-to/persist-conversations.md?raw"; import toolProcesses from "../../../docs/how-to/keep-tool-processes-running.md?raw"; import managedProcesses from "../../../docs/how-to/manage-background-processes.md?raw"; @@ -36,6 +37,14 @@ export const categories: DocCategory[] = [ ]; export const docs: DocPage[] = [ + { + slug: "resume-large-codex-conversations", + sourcePath: "docs/how-to/resume-large-codex-conversations.md", + title: "Resume large Codex conversations", + description: "Continue long threads without returning oversized history frames or losing provider context.", + category: "How-to guides", + body: codexResume, + }, { slug: "quickstart", sourcePath: "docs/tutorials/quickstart.md", diff --git a/src/adapter.rs b/src/adapter.rs index 622088a..02e5b73 100644 --- a/src/adapter.rs +++ b/src/adapter.rs @@ -18,7 +18,7 @@ use crate::{ /// /// Programs and arguments remain separate values throughout execution; this /// crate never constructs a shell command string. -#[derive(Clone)] +#[derive(Clone, PartialEq, Eq)] pub struct CommandSpec { /// Executable path. pub program: PathBuf, @@ -254,6 +254,14 @@ pub trait AgentAdapter: Send + Sync { /// Provider implemented by this adapter. fn provider(&self) -> Provider; + /// Whether this exact adapter supports retaining one native process across turns. + /// + /// Custom adapters remain disabled unless they explicitly implement the + /// complete lifecycle contract. + fn supports_retained_process(&self) -> bool { + false + } + /// Executable name or path meaningful inside the selected execution transport. fn executable(&self) -> PathBuf { PathBuf::from(match self.provider() { @@ -384,6 +392,19 @@ pub trait AgentAdapter: Send + Sync { Ok(()) } + /// Begin another turn on an already initialized retained process. + /// + /// Returning `None` means the adapter cannot safely reuse its process. + fn retained_turn_start(&self, state: &AdapterState) -> Result>> { + let _ = state; + Ok(None) + } + + /// Mark parser state as belonging to a retained native process. + fn mark_retained_turn(&self, state: &mut AdapterState) { + let _ = state; + } + /// Supply a protocol carrier to use instead of the child's stdout and stdin. /// /// Called once, after the provider process is spawned and before the first diff --git a/src/lib.rs b/src/lib.rs index 44aadaf..bda4fa9 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -72,7 +72,7 @@ pub use extensions::{ HarnessMcpServer, HarnessSkill, McpServerManagementRequest, SkillManagementRequest, }; pub use interactions::{InteractionBroker, InteractionBrokerError, InteractionResolution}; -pub use runtime::{AgentRuntime, AgentRuntimeBuilder}; +pub use runtime::{AgentRuntime, AgentRuntimeBuilder, CodexProcessRetention}; pub use sandbox::{ ResolvedSandboxProfile, SandboxBackend, SandboxCapabilities, SandboxContext, SandboxError, SandboxPathAccess, SandboxProfileChange, SandboxProfileManager, SandboxProfileRef, diff --git a/src/providers/codex.rs b/src/providers/codex.rs index 5145a09..5474d86 100644 --- a/src/providers/codex.rs +++ b/src/providers/codex.rs @@ -638,6 +638,10 @@ impl AgentAdapter for Codex { Provider::Codex } + fn supports_retained_process(&self) -> bool { + self.app_server_mode() + } + fn executable(&self) -> PathBuf { self.configured_executable() } @@ -1016,6 +1020,20 @@ impl AgentAdapter for Codex { Ok(()) } + fn retained_turn_start(&self, state: &AdapterState) -> Result>> { + if self.app_server_mode() { + codex_app_server::retained_turn_start(state).map(Some) + } else { + Ok(None) + } + } + + fn mark_retained_turn(&self, state: &mut AdapterState) { + if self.app_server_mode() { + codex_app_server::mark_retained(state); + } + } + fn interrupt_request(&self, state: &AdapterState) -> Option> { self.app_server_mode() .then(|| codex_app_server::interrupt(state)) diff --git a/src/providers/codex_app_server.rs b/src/providers/codex_app_server.rs index d31a35c..9252864 100644 --- a/src/providers/codex_app_server.rs +++ b/src/providers/codex_app_server.rs @@ -60,6 +60,7 @@ struct TurnState { turn_params: Value, /// Whether this turn resumes an existing Codex thread. resume: bool, + retained: bool, /// Model requested for this turn, used to label context-window usage /// before the app server reports the thread's resolved model. model: Option, @@ -74,6 +75,12 @@ struct TurnState { error_message: Option, } +pub(super) fn mark_retained(state: &mut AdapterState) { + let mut turn = load(state); + turn.retained = true; + store(state, &turn); +} + fn load(state: &AdapterState) -> TurnState { state .extensions @@ -156,6 +163,10 @@ pub(super) fn prepare_turn(request: &TurnRequest, state: &mut AdapterState) -> R } if let Some(resumed) = request.session_id.as_deref() { thread_params["threadId"] = json!(resumed); + // Only metadata is needed to start the next turn. Returning the entire + // history can exceed the event frame limit on long conversations. + // The same parameters also cover the active-writer fork fallback. + thread_params["excludeTurns"] = json!(true); } // Image attachments become native `localImage` user inputs; every other @@ -215,6 +226,17 @@ pub(super) fn interrupt(state: &AdapterState) -> Option> { .ok() } +/// Start a turn after this app-server connection has already initialized. +pub(super) fn retained_turn_start(state: &AdapterState) -> Result> { + let turn = load(state); + let method = if turn.resume { + "thread/resume" + } else { + "thread/start" + }; + encode(&request(ID_THREAD, method, turn.thread_params)) +} + /// Translate one JSON-RPC message from the app server. pub(super) fn parse_line(line: &str, state: &mut AdapterState) -> Result { let value: Value = serde_json::from_str(line) @@ -226,6 +248,28 @@ pub(super) fn parse_line(line: &str, state: &mut AdapterState) -> Result Result bool { + matches!( + method, + "item/agentMessage/delta" + | "item/reasoning/textDelta" + | "item/reasoning/summaryTextDelta" + | "item/started" + | "item/completed" + | "item/commandExecution/requestApproval" + | "item/fileChange/requestApproval" + | "item/permissions/requestApproval" + | "item/tool/requestUserInput" + | "thread/tokenUsage/updated" + | "turn/completed" + | "turn/failed" + ) +} + +fn belongs_to_active_turn(value: &Value, turn: &TurnState) -> bool { + if !turn.retained { + return true; + } + let reported = reported_turn_id(value); + matches!((reported, turn.turn_id.as_deref()), (Some(reported), Some(current)) if reported == current) +} + +fn reported_turn_id(value: &Value) -> Option<&str> { + value + .pointer("/params/turnId") + .or_else(|| value.pointer("/params/turn/id")) + .and_then(Value::as_str) +} + /// Advance the handshake with the response to one of our own requests. fn parse_response( value: &Value, @@ -1218,6 +1295,32 @@ mod tests { ); } + #[test] + fn resume_and_writer_conflict_fork_exclude_historical_turns() { + let mut request = TurnRequest::new(Provider::Codex, ".", "continue"); + let mut state = AdapterState::default(); + prepare_turn(&request, &mut state).unwrap(); + let opened = parse_line(r#"{"id":1,"result":{}}"#, &mut state).unwrap(); + let start = decode(opened.writes[0].clone()); + assert_eq!(start["method"], "thread/start"); + assert!(start["params"].get("excludeTurns").is_none()); + + request.session_id = Some("large-thread".into()); + prepare_turn(&request, &mut state).unwrap(); + let opened = parse_line(r#"{"id":1,"result":{}}"#, &mut state).unwrap(); + let resume = decode(opened.writes[0].clone()); + assert_eq!(resume["params"]["excludeTurns"], true); + let conflict = parse_line( + r#"{"id":2,"error":{"message":"thread already has an active writer"}}"#, + &mut state, + ) + .unwrap(); + let fork = decode(conflict.writes[0].clone()); + assert_eq!(fork["method"], "thread/fork"); + assert_eq!(fork["params"]["threadId"], "large-thread"); + assert_eq!(fork["params"]["excludeTurns"], true); + } + #[test] fn the_handshake_response_opens_the_requested_thread_and_starts_the_turn() { let mut request = TurnRequest::new(Provider::Codex, ".", "continue"); @@ -1230,6 +1333,7 @@ mod tests { let resume = decode(opened.writes[0].clone()); assert_eq!(resume["method"], json!("thread/resume")); assert_eq!(resume["params"]["threadId"], json!("thread-7")); + assert_eq!(resume["params"]["excludeTurns"], json!(true)); let started = parse_line( r#"{"jsonrpc":"2.0","id":2,"result":{"thread":{"id":"thread-7"}}}"#, diff --git a/src/retained.rs b/src/retained.rs index 564f2dd..83a4c87 100644 --- a/src/retained.rs +++ b/src/retained.rs @@ -454,6 +454,26 @@ pub trait RuntimeTurnExecutor: Send + Sync { events: &dyn EventSink, interactions: Option<&dyn InteractionHandler>, ) -> crate::Result; + + /// Execute a turn associated with one logical retained runtime. + /// + /// Custom executors keep the legacy behavior by default. Built-in drivers + /// may use the stable runtime identity to isolate native process reuse. + async fn execute_retained( + &self, + runtime_id: &RuntimeId, + request: TurnRequest, + events: &dyn EventSink, + interactions: Option<&dyn InteractionHandler>, + ) -> crate::Result { + let _ = runtime_id; + self.execute(request, events, interactions).await + } + + /// Release resources owned by one logical retained runtime. + async fn dispose_retained(&self, _runtime_id: &RuntimeId) -> crate::Result<()> { + Ok(()) + } } #[async_trait] @@ -461,7 +481,8 @@ impl RuntimeTurnExecutor for AgentRuntime { fn capabilities(&self, provider: Provider) -> RuntimeDriverCapabilities { let permissions = self.permission_support(provider).ok(); RuntimeDriverCapabilities { - retained_process: false, + retained_process: provider == Provider::Codex + && self.codex_process_retention_enabled(provider), session_resume: true, live_interactions: permissions .is_some_and(|support| support.live_approvals || support.live_questions), @@ -501,6 +522,21 @@ impl RuntimeTurnExecutor for AgentRuntime { ) -> crate::Result { self.run(request, events, interactions).await } + + async fn execute_retained( + &self, + runtime_id: &RuntimeId, + request: TurnRequest, + events: &dyn EventSink, + interactions: Option<&dyn InteractionHandler>, + ) -> crate::Result { + self.run_retained(runtime_id, request, events, interactions) + .await + } + + async fn dispose_retained(&self, runtime_id: &RuntimeId) -> crate::Result<()> { + self.dispose_retained_process(runtime_id).await + } } /// Client contract implemented by in-process and future remote runtime hosts. @@ -653,8 +689,9 @@ impl RuntimeClient for InProcessRuntimeClient { let Some(entry) = entry else { return Ok(DisposeOutcome::NotFound); }; - entry.dispose().await?; + let result = entry.dispose().await; self.runtimes.write().await.remove(runtime_id); + result?; Ok(DisposeOutcome::Disposed) } } @@ -1000,7 +1037,7 @@ impl RuntimeEntry { None }; drop(state); - if let Some((invocation_id, terminal)) = terminal { + let confirmation = if let Some((invocation_id, terminal)) = terminal { tokio::time::timeout( INTERRUPT_CONFIRMATION_TIMEOUT, terminal.wait_for_completion(), @@ -1015,9 +1052,23 @@ impl RuntimeEntry { DeliveryState::PossiblySent, "retained provider termination was not confirmed during disposal", ) - })?; - } - Ok(()) + }) + .map(|_| ()) + } else { + Ok(()) + }; + let cleanup = self + .executor + .dispose_retained(&self.spec.runtime_id) + .await + .map_err(|error| { + runtime_error_to_failure( + error, + self.spec.runtime_id.clone(), + InvocationId::new("dispose").expect("static invocation id is valid"), + ) + }); + confirmation.and(cleanup) } async fn start_turn( @@ -1132,7 +1183,12 @@ impl RuntimeEntry { let mut result = match started { Ok(()) => entry .executor - .execute(request, &sink, interactions.as_deref()) + .execute_retained( + &entry.spec.runtime_id, + request, + &sink, + interactions.as_deref(), + ) .await .map_err(|error| { runtime_error_to_failure( diff --git a/src/runtime.rs b/src/runtime.rs index 5a1cadb..bcdcb69 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -1,10 +1,11 @@ use std::collections::{BTreeMap, HashMap}; use std::path::PathBuf; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use std::time::Duration; use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader}; -use tokio::sync::Semaphore; +use tokio::sync::{Mutex as AsyncMutex, OwnedSemaphorePermit, Semaphore}; use crate::adapter::{AdapterState, AgentAdapter, CommandSpec, InteractionRequest}; use crate::error::classify_provider_failure; @@ -58,6 +59,192 @@ const INTERRUPT_GRACE: Duration = Duration::from_secs(5); /// stopped; this only bounds how long the turn waits to observe it. const ATTACHED_SHUTDOWN_GRACE: Duration = Duration::from_secs(5); +/// Limits for opt-in Codex app-server process reuse. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct CodexProcessRetention { + /// Maximum native Codex processes retained by one [`AgentRuntime`]. + pub max_processes: usize, + /// How long an idle process may remain alive after a completed turn. + pub idle_timeout: Duration, +} + +impl CodexProcessRetention { + fn validate(self) -> Result<()> { + if self.max_processes == 0 { + return Err(RuntimeError::InvalidRequest { + field: "codex_process_retention.max_processes", + message: "must be greater than zero".to_string(), + }); + } + if self.idle_timeout.is_zero() { + return Err(RuntimeError::InvalidRequest { + field: "codex_process_retention.idle_timeout", + message: "must be greater than zero".to_string(), + }); + } + Ok(()) + } +} + +#[derive(Clone)] +struct CodexProcessSupervisor { + inner: Arc, +} + +struct CodexProcessSupervisorInner { + config: CodexProcessRetention, + permits: Arc, + processes: AsyncMutex>>, +} + +impl CodexProcessSupervisor { + fn new(config: CodexProcessRetention) -> Self { + Self { + inner: Arc::new(CodexProcessSupervisorInner { + config, + permits: Arc::new(Semaphore::new(config.max_processes)), + processes: AsyncMutex::new(HashMap::new()), + }), + } + } + + async fn dispose(&self, runtime_id: &crate::lifecycle::RuntimeId) -> Result<()> { + let process = self.inner.processes.lock().await.remove(runtime_id); + if let Some(process) = process { + process.terminate().await?; + } + Ok(()) + } + + async fn remove_if_same( + &self, + runtime_id: &crate::lifecycle::RuntimeId, + expected: &Arc, + ) { + let mut processes = self.inner.processes.lock().await; + if processes + .get(runtime_id) + .is_some_and(|current| Arc::ptr_eq(current, expected)) + { + processes.remove(runtime_id); + } + } +} + +#[derive(Clone, PartialEq, Eq)] +struct RetainedCodexFingerprint { + command: CommandSpec, + working_directory: PathBuf, + model: Option, + reasoning: Option, + permission_mode: crate::PermissionMode, + harness_options: BTreeMap, + launch_context: crate::LaunchContext, + auto_compaction: crate::AutoCompactionPolicy, + required_sandbox_capabilities: crate::SandboxCapabilities, +} + +struct RetainedCodexProcess { + fingerprint: RetainedCodexFingerprint, + io: AsyncMutex>, + generation: AtomicU64, + usable: AtomicBool, + permit: AsyncMutex>, +} + +struct RetainedCodexIo { + process: crate::TransportProcess, + stdin: crate::TransportWriter, + reader: BufReader, + stderr_task: tokio::task::JoinHandle>, +} + +async fn read_bounded_retained_line( + reader: &mut BufReader, + limit: usize, +) -> std::io::Result> { + let mut bytes = Vec::new(); + loop { + let available = tokio::io::AsyncBufReadExt::fill_buf(reader).await?; + if available.is_empty() { + if bytes.is_empty() { + return Ok(None); + } + break; + } + let take = available + .iter() + .position(|byte| *byte == b'\n') + .map_or(available.len(), |position| position + 1); + if bytes.len().saturating_add(take) > limit.saturating_add(1) { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "provider event line exceeded configured limit", + )); + } + let ended = available.get(take.saturating_sub(1)) == Some(&b'\n'); + bytes.extend_from_slice(&available[..take]); + tokio::io::AsyncBufReadExt::consume(reader, take); + if ended { + bytes.pop(); + if bytes.last() == Some(&b'\r') { + bytes.pop(); + } + break; + } + } + String::from_utf8(bytes) + .map(Some) + .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error.to_string())) +} + +impl RetainedCodexProcess { + async fn terminate(&self) -> Result<()> { + self.usable.store(false, Ordering::Release); + self.generation.fetch_add(1, Ordering::AcqRel); + if let Some(mut io) = self.io.lock().await.take() { + io.stderr_task.abort(); + io.process + .terminate() + .await + .map_err(|source| RuntimeError::Transport { + provider: Provider::Codex, + source, + })?; + } + Ok(()) + } +} + +struct RetainedTurnCleanup { + armed: bool, + supervisor: CodexProcessSupervisor, + runtime_id: crate::lifecycle::RuntimeId, + process: Arc, +} + +impl RetainedTurnCleanup { + fn disarm(&mut self) { + self.armed = false; + } +} + +impl Drop for RetainedTurnCleanup { + fn drop(&mut self) { + if !self.armed { + return; + } + self.process.usable.store(false, Ordering::Release); + let supervisor = self.supervisor.clone(); + let runtime_id = self.runtime_id.clone(); + let process = Arc::clone(&self.process); + tokio::spawn(async move { + supervisor.remove_if_same(&runtime_id, &process).await; + let _ = process.terminate().await; + }); + } +} + /// Write newline-terminated provider frames to an interactive stdin. async fn write_provider_frames( provider: Provider, @@ -331,6 +518,7 @@ pub struct AgentRuntimeBuilder { max_prompt_bytes: usize, max_event_line_bytes: usize, startup_observer: Option>, + codex_process_retention: Option, } impl AgentRuntimeBuilder { @@ -344,6 +532,7 @@ impl AgentRuntimeBuilder { max_prompt_bytes: DEFAULT_MAX_PROMPT_BYTES, max_event_line_bytes: DEFAULT_MAX_EVENT_LINE_BYTES, startup_observer: None, + codex_process_retention: None, }; #[cfg(feature = "claude")] builder.register(crate::providers::Claude::default()); @@ -400,6 +589,16 @@ impl AgentRuntimeBuilder { self } + /// Keep opted-in Codex app-server processes across retained-runtime turns. + /// + /// Ordinary [`AgentRuntime::run`] calls preserve their existing + /// one-process-per-turn behavior. This setting is used only by the + /// in-process retained-runtime client. + pub fn codex_process_retention(mut self, config: CodexProcessRetention) -> Self { + self.codex_process_retention = Some(config); + self + } + /// Validate limits and construct the runtime. pub fn build(self) -> Result { if self.concurrency_limit == 0 { @@ -414,6 +613,9 @@ impl AgentRuntimeBuilder { message: "prompt and event-line limits must be greater than zero".to_string(), }); } + if let Some(config) = self.codex_process_retention { + config.validate()?; + } Ok(AgentRuntime { adapters: self.adapters, transport: self.transport, @@ -421,6 +623,9 @@ impl AgentRuntimeBuilder { max_prompt_bytes: self.max_prompt_bytes, max_event_line_bytes: self.max_event_line_bytes, startup_observer: self.startup_observer, + codex_process_retention: self + .codex_process_retention + .map(CodexProcessSupervisor::new), }) } } @@ -738,6 +943,7 @@ pub struct AgentRuntime { max_prompt_bytes: usize, max_event_line_bytes: usize, startup_observer: Option>, + codex_process_retention: Option, } impl AgentRuntime { @@ -1918,12 +2124,61 @@ impl AgentRuntime { result } + pub(crate) fn codex_process_retention_enabled(&self, provider: Provider) -> bool { + provider == Provider::Codex + && self.codex_process_retention.is_some() + && self + .adapters + .get(&provider) + .is_some_and(|adapter| adapter.supports_retained_process()) + } + + pub(crate) async fn run_retained( + &self, + runtime_id: &crate::lifecycle::RuntimeId, + request: TurnRequest, + events: &dyn EventSink, + interactions: Option<&dyn InteractionHandler>, + ) -> Result { + if !self.codex_process_retention_enabled(request.provider) { + return self.run(request, events, interactions).await; + } + let mut trace = StartupTrace::new(request.provider, self.startup_observer.clone()); + let result = self + .run_inner_with_retention(Some(runtime_id), request, events, interactions, &trace) + .await; + trace.finish(&result); + result + } + + pub(crate) async fn dispose_retained_process( + &self, + runtime_id: &crate::lifecycle::RuntimeId, + ) -> Result<()> { + if let Some(supervisor) = &self.codex_process_retention { + supervisor.dispose(runtime_id).await?; + } + Ok(()) + } + async fn run_inner( &self, request: TurnRequest, events: &dyn EventSink, interactions: Option<&dyn InteractionHandler>, trace: &StartupTrace, + ) -> Result { + self.run_inner_with_retention(None, request, events, interactions, trace) + .await + } + + async fn run_inner_with_retention( + &self, + retained_runtime_id: Option<&crate::lifecycle::RuntimeId>, + request: TurnRequest, + events: &dyn EventSink, + interactions: Option<&dyn InteractionHandler>, + trace: &StartupTrace, ) -> Result { self.validate(&request)?; self.validate_working_directory(&request).await?; @@ -1984,25 +2239,45 @@ impl AgentRuntime { }; trace.record(StartupStage::PermitAcquired); let timeout = request.timeout; - let result = tokio::time::timeout( - timeout, - self.run_process( - adapter, - &request, - events, - interactions.unwrap_or(&DenyAll), - trace, - ), - ) - .await; + let process = async { + if let (Some(runtime_id), Some(supervisor)) = + (retained_runtime_id, &self.codex_process_retention) + { + self.run_retained_codex_process( + supervisor, + runtime_id, + adapter, + &request, + events, + interactions.unwrap_or(&DenyAll), + trace, + ) + .await + } else { + self.run_process( + adapter, + &request, + events, + interactions.unwrap_or(&DenyAll), + trace, + ) + .await + } + }; + let result = tokio::time::timeout(timeout, process).await; drop(permit); - match result { - Ok(result) => result, - Err(_) => Err(RuntimeError::Timeout { + let Ok(result) = result else { + if let (Some(runtime_id), Some(supervisor)) = + (retained_runtime_id, &self.codex_process_retention) + { + supervisor.dispose(runtime_id).await?; + } + return Err(RuntimeError::Timeout { provider, seconds: timeout.as_secs(), - }), - } + }); + }; + result } /// Run a turn with opt-in, bounded recovery for an application-managed @@ -2301,6 +2576,438 @@ impl AgentRuntime { }) } + #[allow(clippy::too_many_arguments)] + async fn run_retained_codex_process( + &self, + supervisor: &CodexProcessSupervisor, + runtime_id: &crate::lifecycle::RuntimeId, + adapter: Arc, + request: &TurnRequest, + events: &dyn EventSink, + interactions: &dyn InteractionHandler, + trace: &StartupTrace, + ) -> Result { + let provider = request.provider; + let initial_process = supervisor + .inner + .processes + .lock() + .await + .get(runtime_id) + .cloned(); + let reserved_permit = if initial_process.is_some() { + None + } else { + match supervisor.inner.permits.clone().try_acquire_owned() { + Ok(permit) => Some(permit), + Err(_) => { + // Decide the cold fallback before adapter or sandbox + // preparation: preparation may own resources and is not + // required to be side-effect-free. + return self + .run_process(adapter, request, events, interactions, trace) + .await; + } + } + }; + let mut state = AdapterState::default(); + state.result.session_id.clone_from(&request.session_id); + adapter.prepare_turn(request, &mut state)?; + adapter.mark_retained_turn(&mut state); + let mut spec = adapter.command_for_turn(request, &state)?; + for (name, value) in &request.environment { + spec.environment.insert(name.into(), value.expose().into()); + } + trace.record(StartupStage::CommandPrepared); + if let Some(sandbox) = &request.sandbox { + let context = SandboxContext { + provider, + working_directory: request.working_directory.clone(), + }; + spec = tokio::select! { + _ = request.cancellation.cancelled() => return Err(RuntimeError::Cancelled { provider }), + prepared = sandbox.prepare(context, spec) => prepared?, + }; + trace.record(StartupStage::SandboxPrepared); + } + let capabilities = self.transport.capabilities(); + if !spec.interactive_stdin || !capabilities.interactive_stdin { + return Err(RuntimeError::TransportCapabilityUnavailable { + provider, + transport: self.transport.name().to_string(), + capability: "interactive_stdin", + message: "retained Codex requires writable provider stdin".to_string(), + }); + } + if !capabilities.process_tree_termination { + return Err(RuntimeError::TransportCapabilityUnavailable { + provider, + transport: self.transport.name().to_string(), + capability: "process_tree_termination", + message: "retained Codex requires complete process-tree termination".to_string(), + }); + } + let fingerprint = RetainedCodexFingerprint { + command: spec.clone(), + working_directory: request.working_directory.clone(), + model: request.model.clone(), + reasoning: request.reasoning.clone(), + permission_mode: request.permission_mode.clone(), + harness_options: request.harness_options.clone(), + launch_context: request.launch_context.clone(), + auto_compaction: request.auto_compaction, + required_sandbox_capabilities: request.required_sandbox_capabilities, + }; + + let previous = { + let mut processes = supervisor.inner.processes.lock().await; + match processes.get(runtime_id) { + Some(process) + if process.fingerprint == fingerprint + && process.usable.load(Ordering::Acquire) => + { + None + } + Some(_) => processes.remove(runtime_id), + None => None, + } + }; + let mut available_permit = reserved_permit; + if let Some(previous) = previous { + // Transfer the slot before teardown. This makes a configuration + // replacement atomic with respect to pool capacity: another + // runtime cannot steal the released slot after sandbox preparation. + if available_permit.is_none() { + available_permit = previous.permit.lock().await.take(); + } + previous.terminate().await?; + } + + let existing = supervisor + .inner + .processes + .lock() + .await + .get(runtime_id) + .cloned(); + let retained = if let Some(existing) = existing { + existing + } else { + if available_permit.is_none() { + if let Some(initial) = &initial_process { + available_permit = initial.permit.lock().await.take(); + initial.terminate().await?; + } + } + let permit = if let Some(permit) = available_permit { + permit + } else { + supervisor + .inner + .permits + .clone() + .try_acquire_owned() + .map_err(|_| RuntimeError::Transport { + provider, + source: TransportError::new( + TransportErrorKind::SpawnFailed, + self.transport.name(), + "retain_process", + "retained Codex capacity changed while preparing the process", + true, + ), + })? + }; + let program = spec.program.clone(); + let mut process = self + .transport + .spawn(TransportSpawnRequest { + command: spec.clone(), + working_directory: request.working_directory.clone(), + }) + .await + .map_err(|source| { + if source.kind == TransportErrorKind::ExecutableNotFound { + RuntimeError::ExecutableNotFound { + provider, + executable: program.display().to_string(), + } + } else { + RuntimeError::Transport { + provider, + source: redact_transport_error(source, &request.environment), + } + } + })?; + trace.record(StartupStage::ProcessSpawned); + let mut stdin = process.take_stdin().ok_or_else(|| RuntimeError::Protocol { + provider, + message: "retained Codex process did not expose stdin".to_string(), + })?; + if let Some(initial) = &spec.initial_stdin { + write_provider_frames( + provider, + Some(&mut stdin), + std::slice::from_ref(initial), + "initial input", + ) + .await?; + trace.record(StartupStage::InitialInputWritten); + } + let stdout = process + .take_stdout() + .ok_or_else(|| RuntimeError::Protocol { + provider, + message: "retained Codex process did not expose stdout".to_string(), + })?; + let stderr = process + .take_stderr() + .ok_or_else(|| RuntimeError::Protocol { + provider, + message: "retained Codex process did not expose stderr".to_string(), + })?; + let retained = Arc::new(RetainedCodexProcess { + fingerprint, + io: AsyncMutex::new(Some(RetainedCodexIo { + process, + stdin, + reader: BufReader::new(stdout), + stderr_task: tokio::spawn(crate::process::bounded_stderr( + stderr, + STDERR_TAIL_BYTES, + )), + })), + generation: AtomicU64::new(0), + usable: AtomicBool::new(true), + permit: AsyncMutex::new(Some(permit)), + }); + supervisor + .inner + .processes + .lock() + .await + .insert(runtime_id.clone(), Arc::clone(&retained)); + retained + }; + + let generation = retained.generation.fetch_add(1, Ordering::AcqRel) + 1; + let mut cleanup = RetainedTurnCleanup { + armed: true, + supervisor: supervisor.clone(), + runtime_id: runtime_id.clone(), + process: Arc::clone(&retained), + }; + let mut io_guard = retained.io.lock().await; + let io = io_guard.as_mut().ok_or_else(|| RuntimeError::Protocol { + provider, + message: "retained Codex process is no longer available".to_string(), + })?; + let reused = generation > 1; + if reused { + let start = + adapter + .retained_turn_start(&state)? + .ok_or_else(|| RuntimeError::Protocol { + provider, + message: "configured Codex adapter cannot start a retained turn" + .to_string(), + })?; + write_provider_frames( + provider, + Some(&mut io.stdin), + std::slice::from_ref(&start), + "retained turn start", + ) + .await?; + } + trace.record(StartupStage::StreamsAttached); + let result = self + .drive_retained_codex_turn(adapter.as_ref(), request, events, interactions, trace, io) + .await; + drop(io_guard); + match result { + Ok(result) => { + cleanup.disarm(); + let runtime_id = runtime_id.clone(); + let retained_for_expiry = Arc::clone(&retained); + let supervisor_for_expiry = supervisor.clone(); + let idle_timeout = supervisor.inner.config.idle_timeout; + let retained_for_drain = Arc::clone(&retained); + let supervisor_for_drain = supervisor.clone(); + let runtime_id_for_drain = runtime_id.clone(); + let max_event_line_bytes = self.max_event_line_bytes; + tokio::spawn(async move { + loop { + if retained_for_drain.generation.load(Ordering::Acquire) != generation + || !retained_for_drain.usable.load(Ordering::Acquire) + { + return; + } + let observed = { + let mut io = retained_for_drain.io.lock().await; + if retained_for_drain.generation.load(Ordering::Acquire) != generation { + return; + } + let Some(io) = io.as_mut() else { return }; + tokio::time::timeout( + Duration::from_millis(25), + read_bounded_retained_line(&mut io.reader, max_event_line_bytes), + ) + .await + }; + if observed.is_ok() { + // Any frame after the terminal event makes the + // connection ambiguous. Retire it instead of + // assigning or replying under a later turn. + retained_for_drain.usable.store(false, Ordering::Release); + supervisor_for_drain + .remove_if_same(&runtime_id_for_drain, &retained_for_drain) + .await; + let _ = retained_for_drain.terminate().await; + return; + } + } + }); + tokio::spawn(async move { + tokio::time::sleep(idle_timeout).await; + if retained_for_expiry.generation.load(Ordering::Acquire) == generation { + supervisor_for_expiry + .remove_if_same(&runtime_id, &retained_for_expiry) + .await; + let _ = retained_for_expiry.terminate().await; + } + }); + Ok(result) + } + Err(error) => Err(error), + } + } + + async fn drive_retained_codex_turn( + &self, + adapter: &dyn AgentAdapter, + request: &TurnRequest, + events: &dyn EventSink, + interactions: &dyn InteractionHandler, + trace: &StartupTrace, + io: &mut RetainedCodexIo, + ) -> Result { + let provider = request.provider; + let mut state = AdapterState::default(); + state.result.session_id.clone_from(&request.session_id); + adapter.prepare_turn(request, &mut state)?; + adapter.mark_retained_turn(&mut state); + let mut first_output = true; + let mut first_text = true; + loop { + let line = tokio::select! { + _ = request.cancellation.cancelled() => { + if let Some(interrupt) = adapter.interrupt_request(&state) { + let _ = write_provider_frames(provider, Some(&mut io.stdin), &[interrupt], "interrupt").await; + } + return Err(RuntimeError::Cancelled { provider }); + } + line = read_bounded_retained_line(&mut io.reader, self.max_event_line_bytes) => line.map_err(|source| RuntimeError::ProcessIo { + provider, + stream: "stdout read", + source, + })?, + }; + let Some(line) = line else { + return Err(RuntimeError::ProcessFailed { + provider, + kind: ProviderProcessErrorKind::Unknown, + exit_code: None, + stderr: "retained Codex process exited before completing the turn".to_string(), + provider_code: None, + delivery: crate::lifecycle::DeliveryState::PossiblySent, + }); + }; + if first_output { + first_output = false; + trace.record(StartupStage::FirstOutput); + } + if line.len() > self.max_event_line_bytes { + return Err(RuntimeError::Protocol { + provider, + message: format!("event line exceeded {} bytes", self.max_event_line_bytes), + }); + } + let output = adapter.parse_line(&line, &mut state)?; + for event in output.events { + if first_text && matches!(&event, TurnEvent::TextDelta { text } if !text.is_empty()) + { + first_text = false; + trace.record(StartupStage::FirstText); + } + let _delivery = trace.event_delivery(); + events.emit(event).await?; + } + write_provider_frames( + provider, + Some(&mut io.stdin), + &output.writes, + "provider write", + ) + .await?; + if let Some(interaction) = output.interaction { + let response = match interaction { + InteractionRequest::Approval { + request: approval, + original, + } => { + let decision = tokio::select! { + _ = request.cancellation.cancelled() => return Err(RuntimeError::Cancelled { provider }), + decision = tokio::time::timeout(request.interaction_timeout, interactions.approve(approval.clone())) => { + decision.unwrap_or_else(|_| crate::ApprovalDecision::Deny { reason: Some("Approval timed out".to_string()) }) + } + }; + adapter.approval_response(&approval, &original, decision)? + } + InteractionRequest::Question { + request: question, + original, + } => { + let answer = tokio::select! { + _ = request.cancellation.cancelled() => return Err(RuntimeError::Cancelled { provider }), + answer = tokio::time::timeout(request.interaction_timeout, interactions.answer(question.clone())) => answer.ok().flatten(), + }; + adapter.question_response(&question, &original, answer)? + } + }; + if let Some(response) = response { + write_provider_frames( + provider, + Some(&mut io.stdin), + &[response], + "interaction response", + ) + .await?; + } + } + if output.terminal { + if let Some(failure) = state.terminal_failure.take() { + return Err(RuntimeError::ProcessFailed { + provider, + kind: failure.kind, + exit_code: None, + stderr: redact_secrets(&failure.diagnostic, &request.environment), + provider_code: failure.provider_code, + delivery: failure.delivery, + }); + } + if state.result.text.is_empty() { + events + .emit(TurnEvent::Warning { + message: format!("{provider} completed without a text response"), + }) + .await?; + } + return Ok(state.result); + } + } + } + async fn run_process( &self, adapter: Arc, diff --git a/tests/codex_app_server.rs b/tests/codex_app_server.rs index b65b0b8..3d436a7 100644 --- a/tests/codex_app_server.rs +++ b/tests/codex_app_server.rs @@ -9,18 +9,23 @@ use std::collections::BTreeMap; use std::path::Path; +use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::time::Duration; use async_trait::async_trait; use serde_json::{json, Value}; +use temps_agent_runtime::lifecycle::{InvocationId, RuntimeId}; use temps_agent_runtime::providers::{Codex, CodexTurnMode}; use temps_agent_runtime::retained::TurnAttachment; +use temps_agent_runtime::retained::{ + InProcessRuntimeClient, RuntimeClient, RuntimeSpec, TurnInput, +}; use temps_agent_runtime::{ - AgentRuntime, ApprovalDecision, ApprovalRequest, EventSink, ExecutionTransport, - InteractionHandler, McpServerConfig, PermissionMode, Provider, ProviderReadiness, - QuestionAnswer, QuestionRequest, Result, RuntimeError, SandboxCapabilities, SecretString, - TransportCapabilities, TransportError, TransportErrorKind, TransportExitStatus, + AgentRuntime, ApprovalDecision, ApprovalRequest, CodexProcessRetention, EventSink, + ExecutionTransport, InteractionHandler, McpServerConfig, PermissionMode, Provider, + ProviderReadiness, QuestionAnswer, QuestionRequest, Result, RuntimeError, SandboxCapabilities, + SecretString, TransportCapabilities, TransportError, TransportErrorKind, TransportExitStatus, TransportProcess, TransportProcessControl, TransportProcessHandle, TransportReader, TransportReadinessRequest, TransportResult, TransportSpawnRequest, TransportWriter, TurnEvent, TurnRequest, @@ -39,6 +44,10 @@ enum Script { AsyncQuestion, /// Stream one delta and then wait for `turn/interrupt`. Interrupt, + /// Exit after accepting a turn, before a terminal notification. + Crash, + /// Complete, then emit an old-turn frame while the process is idle. + DelayedIdleFrame, } #[derive(Clone)] @@ -48,6 +57,8 @@ struct AppServer { frames: Arc>>, /// Arguments the SDK asked the transport to spawn `codex` with. arguments: Arc>>, + spawns: Arc, + terminations: Arc, } impl AppServer { @@ -56,6 +67,8 @@ impl AppServer { script, frames: Arc::new(Mutex::new(Vec::new())), arguments: Arc::new(Mutex::new(Vec::new())), + spawns: Arc::new(AtomicUsize::new(0)), + terminations: Arc::new(AtomicUsize::new(0)), } } @@ -67,6 +80,14 @@ impl AppServer { self.frames.lock().unwrap().clone() } + fn spawn_count(&self) -> usize { + self.spawns.load(Ordering::Acquire) + } + + fn termination_count(&self) -> usize { + self.terminations.load(Ordering::Acquire) + } + fn method_frame(&self, method: &str) -> Option { self.frames() .into_iter() @@ -120,6 +141,7 @@ impl ExecutionTransport for AppServer { } async fn spawn(&self, request: TransportSpawnRequest) -> TransportResult { + self.spawns.fetch_add(1, Ordering::AcqRel); assert_eq!( request .command @@ -143,6 +165,7 @@ impl ExecutionTransport for AppServer { drop(server_stderr); let script = self.script; let frames = Arc::clone(&self.frames); + let terminations = Arc::clone(&self.terminations); tokio::spawn(async move { serve(script, frames, server_input, server_output).await }); Ok(TransportProcess::new( TransportProcessHandle { @@ -153,7 +176,7 @@ impl ExecutionTransport for AppServer { Some(Box::new(sdk_stdin) as TransportWriter), Box::new(sdk_stdout) as TransportReader, Box::new(sdk_stderr) as TransportReader, - Control, + Control { terminations }, )) } @@ -172,7 +195,9 @@ impl ExecutionTransport for AppServer { } } -struct Control; +struct Control { + terminations: Arc, +} #[async_trait] impl TransportProcessControl for Control { @@ -184,6 +209,7 @@ impl TransportProcessControl for Control { } async fn terminate(&mut self) -> TransportResult<()> { + self.terminations.fetch_add(1, Ordering::AcqRel); Ok(()) } } @@ -236,10 +262,14 @@ async fn serve( json!({"jsonrpc":"2.0","id":id,"result":{"turn":{"id":"turn-1"}}}), ) .await; + if script == Script::Crash { + return; + } send( &mut output, json!({"jsonrpc":"2.0","method":"item/agentMessage/delta", - "params":{"itemId":"item-1","delta":"Working"}}), + "params":{"threadId":"thread-fixture","turnId":"turn-1", + "itemId":"item-1","delta":"Working"}}), ) .await; for frame in opening_frames(script) { @@ -251,6 +281,16 @@ async fn serve( for frame in completion_frames(script) { send(&mut output, frame).await; } + if script == Script::DelayedIdleFrame { + tokio::time::sleep(Duration::from_millis(10)).await; + send( + &mut output, + json!({"jsonrpc":"2.0","method":"item/agentMessage/delta", + "params":{"threadId":"thread-fixture","turnId":"turn-1", + "itemId":"late","delta":"STALE"}}), + ) + .await; + } } // `initialize` and `turn/interrupt` need only a bare acknowledgement. _ => { @@ -265,7 +305,7 @@ fn opening_frames(script: Script) -> Vec { Script::Approval => vec![ json!({"jsonrpc":"2.0","method":"item/started","params":{"item":{ "id":"item-2","type":"commandExecution","command":"cargo test","status":"inProgress" - }}}), + },"threadId":"thread-fixture","turnId":"turn-1"}}), json!({"jsonrpc":"2.0","id":"server-1","method":"item/commandExecution/requestApproval", "params":{"threadId":"thread-fixture","turnId":"turn-1","itemId":"item-2", "command":"cargo test","startedAtMs":1}}), @@ -278,14 +318,14 @@ fn opening_frames(script: Script) -> Vec { "options":[{"label":"Banana","description":"Yellow"}, {"label":"Plantain","description":"Also yellow"}]}]} })], - Script::AsyncQuestion => vec![json!({ + Script::AsyncQuestion | Script::DelayedIdleFrame => vec![json!({ "jsonrpc":"2.0","id":"server-3","method":"item/tool/requestUserInput", "params":{"threadId":"thread-fixture","turnId":"turn-1","itemId":"item-4", "isBlocking":false, "questions":[{"id":"q2","header":"Theme","question":"Dark or light?", "options":[{"label":"Dark","description":"Dim"}]}]} })], - Script::Interrupt => Vec::new(), + Script::Interrupt | Script::Crash => Vec::new(), } } @@ -296,12 +336,12 @@ fn completion_frames(script: Script) -> Vec { json!({"jsonrpc":"2.0","method":"item/completed","params":{"item":{ "id":"item-2","type":"commandExecution","command":"cargo test", "status":"completed","aggregatedOutput":"ok","exitCode":0 - }}}), + },"threadId":"thread-fixture","turnId":"turn-1"}}), ); } frames.extend([ json!({"jsonrpc":"2.0","method":"thread/tokenUsage/updated","params":{ - "threadId":"thread-fixture", + "threadId":"thread-fixture","turnId":"turn-1", "tokenUsage":{"last":{"inputTokens":120,"outputTokens":34,"totalTokens":154}, "modelContextWindow":272_000} }}), @@ -369,6 +409,35 @@ fn runtime(transport: AppServer) -> AgentRuntime { builder.build().unwrap() } +fn retained_runtime(transport: AppServer, idle_timeout: Duration) -> AgentRuntime { + let mut builder = AgentRuntime::builder() + .transport(transport) + .codex_process_retention(CodexProcessRetention { + max_processes: 2, + idle_timeout, + }); + builder.register(Codex::app_server()); + builder.build().unwrap() +} + +async fn retained_handle( + runtime: AgentRuntime, +) -> ( + InProcessRuntimeClient, + temps_agent_runtime::retained::RuntimeHandle, +) { + let client = InProcessRuntimeClient::new(runtime); + let handle = client + .acquire(RuntimeSpec::new( + RuntimeId::new("codex-fixture-runtime").unwrap(), + Provider::Codex, + ".", + )) + .await + .unwrap(); + (client, handle) +} + fn request() -> TurnRequest { let mut request = TurnRequest::new(Provider::Codex, ".", "review the workspace"); request.permission_mode = PermissionMode::Default; @@ -400,6 +469,344 @@ async fn the_app_server_mode_advertises_live_interactions() { assert!(turn.native_image_attachments); } +#[tokio::test] +async fn retained_client_reuses_one_app_server_for_two_turns() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + assert!(handle.driver_capabilities().retained_process); + + let first = handle + .start_turn(TurnInput::new( + InvocationId::new("turn-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(first.session_id.as_deref(), Some("thread-fixture")); + + let second = handle + .start_turn(TurnInput::new( + InvocationId::new("turn-two").unwrap(), + "second", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(second.session_id.as_deref(), Some("thread-fixture")); + assert_eq!(transport.spawn_count(), 1); + assert_eq!( + transport + .frames() + .iter() + .filter(|frame| frame.get("method").and_then(Value::as_str) == Some("initialize")) + .count(), + 1 + ); +} + +#[tokio::test] +async fn retained_processes_are_disabled_by_default() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = runtime(transport.clone()); + let (_client, handle) = retained_handle(runtime).await; + assert!(!handle.driver_capabilities().retained_process); + + for invocation in ["default-one", "default-two"] { + handle + .start_turn(TurnInput::new( + InvocationId::new(invocation).unwrap(), + invocation, + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + } + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn changed_permission_replaces_the_retained_process() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + handle + .start_turn(TurnInput::new( + InvocationId::new("permission-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + let mut changed = TurnInput::new(InvocationId::new("permission-two").unwrap(), "second"); + changed.permission_mode = Some(PermissionMode::FullAccess); + handle + .start_turn(changed) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn configuration_replacement_transfers_its_only_pool_slot() { + let transport = AppServer::new(Script::AsyncQuestion); + let mut builder = AgentRuntime::builder() + .transport(transport.clone()) + .codex_process_retention(CodexProcessRetention { + max_processes: 1, + idle_timeout: Duration::from_secs(30), + }); + builder.register(Codex::app_server()); + let (_client, handle) = retained_handle(builder.build().unwrap()).await; + handle + .start_turn(TurnInput::new( + InvocationId::new("replace-only-slot-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + let mut changed = TurnInput::new( + InvocationId::new("replace-only-slot-two").unwrap(), + "second", + ); + changed.permission_mode = Some(PermissionMode::FullAccess); + handle + .start_turn(changed) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(transport.spawn_count(), 2); + assert!(transport.termination_count() >= 1); +} + +#[tokio::test] +async fn idle_expiry_and_dispose_terminate_retained_processes() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = retained_runtime(transport.clone(), Duration::from_millis(20)); + let (client, handle) = retained_handle(runtime).await; + handle + .start_turn(TurnInput::new( + InvocationId::new("idle-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(60)).await; + handle + .start_turn(TurnInput::new( + InvocationId::new("idle-two").unwrap(), + "second", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(transport.spawn_count(), 2); + let runtime_id = handle.runtime_id().clone(); + client.dispose(&runtime_id).await.unwrap(); + client + .acquire(RuntimeSpec::new(runtime_id, Provider::Codex, ".")) + .await + .expect("a disposed runtime identifier can be acquired again"); + assert!(transport.termination_count() >= 2); +} + +#[tokio::test] +async fn process_capacity_does_not_starve_an_existing_retained_runtime() { + let transport = AppServer::new(Script::AsyncQuestion); + let mut builder = AgentRuntime::builder() + .transport(transport.clone()) + .codex_process_retention(CodexProcessRetention { + max_processes: 1, + idle_timeout: Duration::from_secs(30), + }); + builder.register(Codex::app_server()); + let client = InProcessRuntimeClient::new(builder.build().unwrap()); + let first = client + .acquire(RuntimeSpec::new( + RuntimeId::new("capacity-one").unwrap(), + Provider::Codex, + ".", + )) + .await + .unwrap(); + let second = client + .acquire(RuntimeSpec::new( + RuntimeId::new("capacity-two").unwrap(), + Provider::Codex, + ".", + )) + .await + .unwrap(); + first + .start_turn(TurnInput::new( + InvocationId::new("capacity-first").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + second + .start_turn(TurnInput::new( + InvocationId::new("capacity-rejected").unwrap(), + "second runtime", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + first + .start_turn(TurnInput::new( + InvocationId::new("capacity-reuse").unwrap(), + "reuse", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn a_timed_out_turn_is_never_reused() { + let transport = AppServer::new(Script::Interrupt); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let client = InProcessRuntimeClient::new(runtime); + let mut spec = RuntimeSpec::new( + RuntimeId::new("timeout-runtime").unwrap(), + Provider::Codex, + ".", + ); + spec.turn_timeout = Duration::from_millis(20); + let handle = client.acquire(spec).await.unwrap(); + for invocation in ["timeout-one", "timeout-two"] { + let failure = handle + .start_turn(TurnInput::new( + InvocationId::new(invocation).unwrap(), + invocation, + )) + .await + .unwrap() + .wait() + .await + .unwrap_err(); + assert_eq!( + failure.kind, + temps_agent_runtime::lifecycle::RuntimeFailureKind::Timeout + ); + } + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn a_crashed_process_is_replaced_for_the_next_turn() { + let transport = AppServer::new(Script::Crash); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + for invocation in ["crash-one", "crash-two"] { + assert!(handle + .start_turn(TurnInput::new( + InvocationId::new(invocation).unwrap(), + invocation, + )) + .await + .unwrap() + .wait() + .await + .is_err()); + } + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn an_idle_frame_retires_the_process_before_the_next_turn() { + let transport = AppServer::new(Script::DelayedIdleFrame); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + let first = handle + .start_turn(TurnInput::new( + InvocationId::new("late-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert!(!first.text.contains("STALE")); + tokio::time::sleep(Duration::from_millis(50)).await; + let second = handle + .start_turn(TurnInput::new( + InvocationId::new("late-two").unwrap(), + "second", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert!(!second.text.contains("STALE")); + assert_eq!(transport.spawn_count(), 2); + assert!(transport.termination_count() >= 1); +} + +#[tokio::test] +async fn cancellation_retires_the_process_before_an_immediate_retry() { + let transport = AppServer::new(Script::Interrupt); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + let first = handle + .start_turn(TurnInput::new( + InvocationId::new("cancel-one").unwrap(), + "first", + )) + .await + .unwrap(); + transport.wait_for("turn/start").await.unwrap(); + first.interrupt().await; + + let second = handle + .start_turn(TurnInput::new( + InvocationId::new("cancel-two").unwrap(), + "second", + )) + .await + .unwrap(); + for _ in 0..100 { + if transport.spawn_count() == 2 { + break; + } + tokio::time::sleep(Duration::from_millis(5)).await; + } + assert_eq!(transport.spawn_count(), 2); + second.interrupt().await; +} + #[tokio::test] async fn an_approval_is_accepted_through_the_interaction_handler() { let transport = AppServer::new(Script::Approval);