From 0d3487d917dbe536c1dd8583887fd3747f835455 Mon Sep 17 00:00:00 2001 From: David Viejo Date: Wed, 23 Sep 2026 17:07:31 +0200 Subject: [PATCH 1/5] feat(codex): retain app-server processes across turns --- CHANGELOG.md | 3 + docs/adr/0004-prepared-provider-processes.md | 41 +- docs/reference/retained-runtimes.md | 32 + src/adapter.rs | 23 +- src/lib.rs | 2 +- src/providers/codex.rs | 18 + src/providers/codex_app_server.rs | 73 ++ src/retained.rs | 70 +- src/runtime.rs | 727 ++++++++++++++++++- tests/codex_app_server.rs | 394 +++++++++- 10 files changed, 1322 insertions(+), 61 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3f04fb3..e1a1542 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -30,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/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/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/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 928409b..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 @@ -219,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) @@ -230,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, 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..b3b9a62 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: OwnedSemaphorePermit, +} + +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,424 @@ 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 has_runtime_process = supervisor + .inner + .processes + .lock() + .await + .contains_key(runtime_id); + let reserved_permit = if has_runtime_process { + 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, + } + }; + if let Some(previous) = previous { + previous.terminate().await?; + } + + let existing = supervisor + .inner + .processes + .lock() + .await + .get(runtime_id) + .cloned(); + let retained = if let Some(existing) = existing { + existing + } else { + let permit = if let Some(permit) = reserved_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: 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..7702810 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,307 @@ 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 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); From 55f4b9ddc38cb5f9bbf7b1d75c7ed3cdf1dd3e37 Mon Sep 17 00:00:00 2001 From: David Viejo Date: Wed, 23 Sep 2026 17:12:20 +0200 Subject: [PATCH 2/5] fix(codex): transfer retained process capacity on replacement --- src/runtime.rs | 26 ++++++++++++++++++++------ tests/codex_app_server.rs | 37 +++++++++++++++++++++++++++++++++++++ 2 files changed, 57 insertions(+), 6 deletions(-) diff --git a/src/runtime.rs b/src/runtime.rs index b3b9a62..bcdcb69 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -149,7 +149,7 @@ struct RetainedCodexProcess { io: AsyncMutex>, generation: AtomicU64, usable: AtomicBool, - _permit: OwnedSemaphorePermit, + permit: AsyncMutex>, } struct RetainedCodexIo { @@ -2588,13 +2588,14 @@ impl AgentRuntime { trace: &StartupTrace, ) -> Result { let provider = request.provider; - let has_runtime_process = supervisor + let initial_process = supervisor .inner .processes .lock() .await - .contains_key(runtime_id); - let reserved_permit = if has_runtime_process { + .get(runtime_id) + .cloned(); + let reserved_permit = if initial_process.is_some() { None } else { match supervisor.inner.permits.clone().try_acquire_owned() { @@ -2671,7 +2672,14 @@ impl AgentRuntime { 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?; } @@ -2685,7 +2693,13 @@ impl AgentRuntime { let retained = if let Some(existing) = existing { existing } else { - let permit = if let Some(permit) = reserved_permit { + 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 @@ -2765,7 +2779,7 @@ impl AgentRuntime { })), generation: AtomicU64::new(0), usable: AtomicBool::new(true), - _permit: permit, + permit: AsyncMutex::new(Some(permit)), }); supervisor .inner diff --git a/tests/codex_app_server.rs b/tests/codex_app_server.rs index 7702810..3d436a7 100644 --- a/tests/codex_app_server.rs +++ b/tests/codex_app_server.rs @@ -559,6 +559,43 @@ async fn changed_permission_replaces_the_retained_process() { 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); From 5c8873bbc70dffd2113d5bd6bf0ac2d9f5ae1f7f Mon Sep 17 00:00:00 2001 From: David Viejo Date: Wed, 23 Sep 2026 17:57:50 +0200 Subject: [PATCH 3/5] feat: retain Claude and OpenCode processes --- CHANGELOG.md | 5 + docs/adr/0004-prepared-provider-processes.md | 18 +- docs/reference/retained-runtimes.md | 28 +- src/adapter.rs | 19 + src/lib.rs | 11 +- src/providers/claude.rs | 4 + src/providers/opencode.rs | 33 + src/providers/opencode_serve.rs | 28 +- src/retained.rs | 3 +- src/runtime.rs | 616 +++++++++++++++---- src/types.rs | 22 + tests/opencode_serve.rs | 368 ++++++++++- tests/retained_providers.rs | 237 +++++++ 13 files changed, 1258 insertions(+), 134 deletions(-) create mode 100644 tests/retained_providers.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index e1a1542..e88e16f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,11 @@ Versioning and Keep a Changelog conventions. ### Changed +- Add opt-in bounded native process reuse for Claude streaming and OpenCode serve, + alongside Codex app-server reuse. Retained processes use provider-specific + preflight checks, strict session/configuration isolation, idle expiry, and + no-replay delivery boundaries. Default one-process-per-turn behavior is unchanged. + - Preserve Windows system and profile environment variables when launching providers, without inheriting unrelated credentials. Hide provider and cleanup console windows. Add native process fixtures for arguments, stdin, failures, cancellation, and diff --git a/docs/adr/0004-prepared-provider-processes.md b/docs/adr/0004-prepared-provider-processes.md index e6a936d..2138c5c 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: Accepted (Codex retained turns implemented; explicit prewarming deferred) +- Status: Accepted (Claude, Codex, and OpenCode retained turns implemented; explicit prewarming deferred) - Date: 2026-09-23 ## Problem @@ -24,9 +24,11 @@ records, authorization, feature rollout, and presentation. Listing projects or conversations must never spawn provider processes. 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. +lazy, one-process-per-turn behavior. `provider_process_retention` opts the +in-process retained client into bounded reuse for the built-in Claude streaming, +Codex app-server, and OpenCode serve protocols. `codex_process_retention` remains +the compatibility opt-in for Codex alone. Custom adapters remain disabled unless +they implement the lifecycle contract. ## Lifecycle @@ -72,9 +74,11 @@ a prerequisite to Fleet enabling retention. Claude keeps a streaming input channel open and separates initialization from user messages. Codex keeps one app-server connection and thread alive, issues one -initialize handshake per process, and starts later turns on that connection. Both -need bounded background draining, crash detection, cancellation, late-event -isolation and cleanup. One-shot Codex exec and unsupported providers remain lazy. +initialize handshake per process, and starts later turns on that connection. +OpenCode retains its loopback serve process while creating a fresh health-checked +HTTP/SSE bridge for each turn. All three use bounded initialization and inactivity +deadlines, crash detection, cancellation, late-event isolation, and process-tree +cleanup. One-shot modes and unsupported providers remain lazy. ## Fleet adoption diff --git a/docs/reference/retained-runtimes.md b/docs/reference/retained-runtimes.md index ee9d4bc..2b07eb4 100644 --- a/docs/reference/retained-runtimes.md +++ b/docs/reference/retained-runtimes.md @@ -11,22 +11,32 @@ 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: +Provider process reuse is explicit and in-process only. `ProviderProcessRetention` +enables the built-in Claude streaming driver, Codex app-server driver, and +OpenCode serve driver with one shared process budget: ```rust use std::time::Duration; -use temps_agent_runtime::providers::Codex; -use temps_agent_runtime::{AgentRuntime, CodexProcessRetention}; +use temps_agent_runtime::providers::{Claude, Codex, OpenCode}; +use temps_agent_runtime::{AgentRuntime, ProviderProcessRetention}; let mut builder = AgentRuntime::builder() - .codex_process_retention(CodexProcessRetention { + .provider_process_retention(ProviderProcessRetention { max_processes: 4, idle_timeout: Duration::from_secs(120), + initialization_timeout: Duration::from_secs(30), + active_inactivity_timeout: Some(Duration::from_secs(30 * 60)), }); +builder.register(Claude::default()); builder.register(Codex::app_server()); +builder.register(OpenCode::serve()); let runtime = builder.build()?; ``` +`codex_process_retention` remains available for applications that only want +Codex app-server reuse. Both settings are disabled by default. One-shot modes, +custom adapters, and remote OpenCode transports retain their existing behavior. + 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 @@ -39,9 +49,13 @@ 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. +Crash, failed health checks, cancellation, timeout, dropped turn futures, and +disposal retire the process tree. Before submitting a later prompt, Claude uses +a bounded read-only control probe and OpenCode checks its loopback health endpoint +through a fresh per-turn bridge. A pre-submission failure may start a replacement; +the SDK never automatically replays a prompt after delivery is possible. Codex +correlates native turn IDs, OpenCode requires the active session ID, and Claude +retires a connection that emits ambiguous idle output. Use `RuntimeHandle::configuration_impact` before presenting a live setting change. The compatibility driver applies per-turn model, reasoning, permission, diff --git a/src/adapter.rs b/src/adapter.rs index 02e5b73..73f1f85 100644 --- a/src/adapter.rs +++ b/src/adapter.rs @@ -213,6 +213,8 @@ pub struct AdapterOutput { pub writes: Vec>, /// True after a provider terminal frame. The runtime then closes stdin. pub terminal: bool, + /// The provider acknowledged the submitted prompt for this turn. + pub turn_submitted: bool, } /// Protocol carrier an adapter supplies in place of the provider's own stdio. @@ -392,6 +394,23 @@ pub trait AgentAdapter: Send + Sync { Ok(()) } + /// Seed a retained turn, optionally reusing provider-specific process state. + fn prepare_retained_turn( + &self, + request: &TurnRequest, + state: &mut AdapterState, + process_hint: Option, + ) -> Result<()> { + let _ = process_hint; + self.prepare_turn(request, state) + } + + /// Opaque provider-specific state needed to address this retained process. + fn retained_process_hint(&self, state: &AdapterState) -> Option { + let _ = state; + None + } + /// Begin another turn on an already initialized retained process. /// /// Returning `None` means the adapter cannot safely reuse its process. diff --git a/src/lib.rs b/src/lib.rs index bda4fa9..003af08 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -72,7 +72,9 @@ pub use extensions::{ HarnessMcpServer, HarnessSkill, McpServerManagementRequest, SkillManagementRequest, }; pub use interactions::{InteractionBroker, InteractionBrokerError, InteractionResolution}; -pub use runtime::{AgentRuntime, AgentRuntimeBuilder, CodexProcessRetention}; +pub use runtime::{ + AgentRuntime, AgentRuntimeBuilder, CodexProcessRetention, ProviderProcessRetention, +}; pub use sandbox::{ ResolvedSandboxProfile, SandboxBackend, SandboxCapabilities, SandboxContext, SandboxError, SandboxPathAccess, SandboxProfileChange, SandboxProfileManager, SandboxProfileRef, @@ -101,9 +103,10 @@ pub use types::{ AgentTaskActivityKind, AgentTaskUsage, ApprovalDecision, ApprovalRequest, AutoCompactionPolicy, CompactionTrigger, ContextCompaction, ContextWindowUsage, DenyAll, EventSink, InteractionHandler, LaunchContext, LaunchContextCapabilities, McpServerConfig, NoopEventSink, - PermissionMode, PermissionSupport, Provider, ProviderReadiness, QuestionAnswer, QuestionOption, - QuestionPrompt, QuestionRequest, RunStatus, SecretString, ToolCallStatus, ToolProcessPolicy, - TurnCapabilities, TurnEvent, TurnProvenance, TurnRequest, TurnResult, Usage, + PermissionMode, PermissionSupport, Provider, ProviderProcessStatus, ProviderReadiness, + QuestionAnswer, QuestionOption, QuestionPrompt, QuestionRequest, RunStatus, SecretString, + ToolCallStatus, ToolProcessPolicy, TurnCapabilities, TurnEvent, TurnProvenance, TurnRequest, + TurnResult, Usage, }; pub use startup::{StartupObserver, StartupStage, StartupTiming}; diff --git a/src/providers/claude.rs b/src/providers/claude.rs index ee03158..3aa650a 100644 --- a/src/providers/claude.rs +++ b/src/providers/claude.rs @@ -488,6 +488,10 @@ impl AgentAdapter for Claude { Provider::Claude } + fn supports_retained_process(&self) -> bool { + true + } + fn executable(&self) -> PathBuf { self.configured_executable() } diff --git a/src/providers/opencode.rs b/src/providers/opencode.rs index 2648173..85bdb64 100644 --- a/src/providers/opencode.rs +++ b/src/providers/opencode.rs @@ -123,6 +123,10 @@ impl AgentAdapter for OpenCode { Provider::OpenCode } + fn supports_retained_process(&self) -> bool { + self.serve_mode() + } + fn executable(&self) -> PathBuf { self.configured_executable() } @@ -292,6 +296,35 @@ impl AgentAdapter for OpenCode { Ok(()) } + fn prepare_retained_turn( + &self, + request: &TurnRequest, + state: &mut AdapterState, + process_hint: Option, + ) -> Result<()> { + if !self.serve_mode() { + return self.prepare_turn(request, state); + } + let port = process_hint + .and_then(|port| u16::try_from(port).ok()) + .ok_or_else(|| RuntimeError::Protocol { + provider: Provider::OpenCode, + message: "retained OpenCode process did not preserve its loopback port".into(), + })?; + super::opencode_serve::prepare_turn(request, state, port); + Ok(()) + } + + fn retained_process_hint(&self, state: &AdapterState) -> Option { + super::opencode_serve::turn_port(state).map(u64::from) + } + + fn mark_retained_turn(&self, state: &mut AdapterState) { + if self.serve_mode() { + super::opencode_serve::mark_retained(state); + } + } + fn command_for_turn(&self, request: &TurnRequest, state: &AdapterState) -> Result { match super::opencode_serve::turn_port(state) { Some(port) if self.serve_mode() => self.serve_command(request, port), diff --git a/src/providers/opencode_serve.rs b/src/providers/opencode_serve.rs index 2da5abc..00f17a5 100644 --- a/src/providers/opencode_serve.rs +++ b/src/providers/opencode_serve.rs @@ -82,6 +82,8 @@ pub(super) struct TurnState { error_message: Option, /// Whether any assistant text or tool call was observed. saw_activity: bool, + /// Require every turn-affecting SSE event to identify this turn's session. + retained: bool, } fn load(state: &AdapterState) -> TurnState { @@ -105,6 +107,12 @@ pub(super) fn turn_port(state: &AdapterState) -> Option { (port != 0).then_some(port) } +pub(super) fn mark_retained(state: &mut AdapterState) { + let mut turn = load(state); + turn.retained = true; + store(state, &turn); +} + fn protocol(message: impl Into) -> RuntimeError { RuntimeError::Protocol { provider: Provider::OpenCode, @@ -513,6 +521,9 @@ fn response( "path": "/event", }))?); } + Some(ID_PROMPT) if (200..300).contains(&status) => { + output.turn_submitted = true; + } Some(ID_PROMPT) if status == 0 || status >= 400 => { let message = body .pointer("/data/message") @@ -541,11 +552,18 @@ fn event( .and_then(Value::as_str) .unwrap_or_default(); let properties = event.get("properties").unwrap_or(&Value::Null); - // Another session's activity must never be charged to this turn. - if let (Some(reported), Some(current)) = ( - properties.get("sessionID").and_then(Value::as_str), - turn.session_id.as_deref(), - ) { + // A retained server carries traffic for more than one turn. Fail closed: + // every event that can mutate or finish a turn must identify the current + // session, so delayed or unrelated traffic cannot leak across turns. + let reported = properties.get("sessionID").and_then(Value::as_str); + if turn.retained { + if reported + .zip(turn.session_id.as_deref()) + .is_none_or(|(a, b)| a != b) + { + return; + } + } else if let (Some(reported), Some(current)) = (reported, turn.session_id.as_deref()) { if reported != current { return; } diff --git a/src/retained.rs b/src/retained.rs index 83a4c87..02d52b2 100644 --- a/src/retained.rs +++ b/src/retained.rs @@ -481,8 +481,7 @@ impl RuntimeTurnExecutor for AgentRuntime { fn capabilities(&self, provider: Provider) -> RuntimeDriverCapabilities { let permissions = self.permission_support(provider).ok(); RuntimeDriverCapabilities { - retained_process: provider == Provider::Codex - && self.codex_process_retention_enabled(provider), + retained_process: self.process_retention_enabled(provider), session_resume: true, live_interactions: permissions .is_some_and(|support| support.live_approvals || support.live_questions), diff --git a/src/runtime.rs b/src/runtime.rs index bcdcb69..e7f9b39 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -68,6 +68,46 @@ pub struct CodexProcessRetention { pub idle_timeout: Duration, } +/// Limits and health deadlines for opt-in built-in provider process reuse. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ProviderProcessRetention { + /// Maximum retained Claude, Codex, and OpenCode processes in aggregate. + pub max_processes: usize, + /// How long an idle process may remain alive after a completed turn. + pub idle_timeout: Duration, + /// Maximum time for a retained provider to become ready for a turn. + pub initialization_timeout: Duration, + /// Optional maximum silence while a submitted turn is active. + pub active_inactivity_timeout: Option, +} + +impl ProviderProcessRetention { + fn validate(self) -> Result<()> { + if self.max_processes == 0 { + return Err(RuntimeError::InvalidRequest { + field: "provider_process_retention.max_processes", + message: "must be greater than zero".to_string(), + }); + } + if self.idle_timeout.is_zero() || self.initialization_timeout.is_zero() { + return Err(RuntimeError::InvalidRequest { + field: "provider_process_retention", + message: "idle and initialization timeouts must be greater than zero".to_string(), + }); + } + if self + .active_inactivity_timeout + .is_some_and(|timeout| timeout.is_zero()) + { + return Err(RuntimeError::InvalidRequest { + field: "provider_process_retention.active_inactivity_timeout", + message: "must be greater than zero when configured".to_string(), + }); + } + Ok(()) + } +} + impl CodexProcessRetention { fn validate(self) -> Result<()> { if self.max_processes == 0 { @@ -92,13 +132,13 @@ struct CodexProcessSupervisor { } struct CodexProcessSupervisorInner { - config: CodexProcessRetention, + config: ProviderProcessRetention, permits: Arc, processes: AsyncMutex>>, } impl CodexProcessSupervisor { - fn new(config: CodexProcessRetention) -> Self { + fn new(config: ProviderProcessRetention) -> Self { Self { inner: Arc::new(CodexProcessSupervisorInner { config, @@ -144,11 +184,30 @@ struct RetainedCodexFingerprint { required_sandbox_capabilities: crate::SandboxCapabilities, } +fn retained_fingerprint_command(provider: Provider, mut command: CommandSpec) -> CommandSpec { + command.initial_stdin = None; + if provider == Provider::Claude { + let mut filtered = Vec::with_capacity(command.args.len()); + let mut arguments = command.args.into_iter(); + while let Some(argument) = arguments.next() { + if argument == "--resume" { + let _ = arguments.next(); + } else { + filtered.push(argument); + } + } + command.args = filtered; + } + command +} + struct RetainedCodexProcess { fingerprint: RetainedCodexFingerprint, io: AsyncMutex>, generation: AtomicU64, usable: AtomicBool, + process_hint: Option, + session_id: AsyncMutex>, permit: AsyncMutex>, } @@ -157,6 +216,7 @@ struct RetainedCodexIo { stdin: crate::TransportWriter, reader: BufReader, stderr_task: tokio::task::JoinHandle>, + stdout_task: Option>>, } async fn read_bounded_retained_line( @@ -204,6 +264,9 @@ impl RetainedCodexProcess { self.generation.fetch_add(1, Ordering::AcqRel); if let Some(mut io) = self.io.lock().await.take() { io.stderr_task.abort(); + if let Some(stdout_task) = io.stdout_task.take() { + stdout_task.abort(); + } io.process .terminate() .await @@ -519,6 +582,7 @@ pub struct AgentRuntimeBuilder { max_event_line_bytes: usize, startup_observer: Option>, codex_process_retention: Option, + provider_process_retention: Option, } impl AgentRuntimeBuilder { @@ -533,6 +597,7 @@ impl AgentRuntimeBuilder { max_event_line_bytes: DEFAULT_MAX_EVENT_LINE_BYTES, startup_observer: None, codex_process_retention: None, + provider_process_retention: None, }; #[cfg(feature = "claude")] builder.register(crate::providers::Claude::default()); @@ -599,6 +664,12 @@ impl AgentRuntimeBuilder { self } + /// Keep supported built-in provider processes across retained-runtime turns. + pub fn provider_process_retention(mut self, config: ProviderProcessRetention) -> Self { + self.provider_process_retention = Some(config); + self + } + /// Validate limits and construct the runtime. pub fn build(self) -> Result { if self.concurrency_limit == 0 { @@ -616,6 +687,18 @@ impl AgentRuntimeBuilder { if let Some(config) = self.codex_process_retention { config.validate()?; } + if let Some(config) = self.provider_process_retention { + config.validate()?; + } + let retention = self.provider_process_retention.or_else(|| { + self.codex_process_retention + .map(|config| ProviderProcessRetention { + max_processes: config.max_processes, + idle_timeout: config.idle_timeout, + initialization_timeout: Duration::from_secs(30), + active_inactivity_timeout: None, + }) + }); Ok(AgentRuntime { adapters: self.adapters, transport: self.transport, @@ -623,9 +706,8 @@ 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), + codex_process_retention: retention.map(CodexProcessSupervisor::new), + provider_process_retention: self.provider_process_retention, }) } } @@ -944,6 +1026,7 @@ pub struct AgentRuntime { max_event_line_bytes: usize, startup_observer: Option>, codex_process_retention: Option, + provider_process_retention: Option, } impl AgentRuntime { @@ -2124,9 +2207,10 @@ impl AgentRuntime { result } - pub(crate) fn codex_process_retention_enabled(&self, provider: Provider) -> bool { - provider == Provider::Codex - && self.codex_process_retention.is_some() + pub(crate) fn process_retention_enabled(&self, provider: Provider) -> bool { + self.codex_process_retention.is_some() + && (provider == Provider::Codex || self.provider_process_retention.is_some()) + && !(provider == Provider::OpenCode && self.transport.capabilities().remote) && self .adapters .get(&provider) @@ -2140,7 +2224,7 @@ impl AgentRuntime { events: &dyn EventSink, interactions: Option<&dyn InteractionHandler>, ) -> Result { - if !self.codex_process_retention_enabled(request.provider) { + if !self.process_retention_enabled(request.provider) { return self.run(request, events, interactions).await; } let mut trace = StartupTrace::new(request.provider, self.startup_observer.clone()); @@ -2612,7 +2696,15 @@ impl AgentRuntime { }; let mut state = AdapterState::default(); state.result.session_id.clone_from(&request.session_id); - adapter.prepare_turn(request, &mut state)?; + if provider == Provider::OpenCode { + if let Some(initial) = &initial_process { + adapter.prepare_retained_turn(request, &mut state, initial.process_hint)?; + } else { + adapter.prepare_turn(request, &mut state)?; + } + } else { + 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 { @@ -2631,12 +2723,14 @@ impl AgentRuntime { trace.record(StartupStage::SandboxPrepared); } let capabilities = self.transport.capabilities(); - if !spec.interactive_stdin || !capabilities.interactive_stdin { + if provider != Provider::OpenCode + && (!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(), + message: format!("retained {provider} requires writable provider stdin"), }); } if !capabilities.process_tree_termination { @@ -2644,11 +2738,11 @@ impl AgentRuntime { provider, transport: self.transport.name().to_string(), capability: "process_tree_termination", - message: "retained Codex requires complete process-tree termination".to_string(), + message: format!("retained {provider} requires complete process-tree termination"), }); } let fingerprint = RetainedCodexFingerprint { - command: spec.clone(), + command: retained_fingerprint_command(provider, spec.clone()), working_directory: request.working_directory.clone(), model: request.model.clone(), reasoning: request.reasoning.clone(), @@ -2683,13 +2777,196 @@ impl AgentRuntime { previous.terminate().await?; } - let existing = supervisor + let mut existing = supervisor .inner .processes .lock() .await .get(runtime_id) .cloned(); + if let Some(candidate) = existing.clone() { + if *candidate.session_id.lock().await != request.session_id { + supervisor.remove_if_same(runtime_id, &candidate).await; + if available_permit.is_none() { + available_permit = candidate.permit.lock().await.take(); + } + candidate.terminate().await?; + existing = None; + } + } + let mut prepared_protocol_streams = None; + let mut prefetched_line = None; + let mut preflight_events = Vec::new(); + let mut claimed_generation = None; + if provider == Provider::Claude { + if let Some(candidate) = existing.clone() { + events + .emit(TurnEvent::ProviderProcessStatus { + status: crate::ProviderProcessStatus::Checking, + message: "Checking the retained Claude process before sending the prompt." + .to_string(), + }) + .await?; + let claim = candidate.generation.fetch_add(1, Ordering::AcqRel) + 1; + claimed_generation = Some(claim); + let request_id = format!("temps-agent-runtime-health-{claim}"); + let probe = serde_json::to_vec(&serde_json::json!({ + "type": "control_request", + "request_id": request_id, + "request": { "subtype": "mcp_status" } + })) + .map_err(|error| RuntimeError::Protocol { + provider, + message: format!("could not encode Claude health probe: {error}"), + })?; + let health_timeout = supervisor + .inner + .config + .initialization_timeout + .min(Duration::from_secs(3)); + let healthy = { + let mut guard = candidate.io.lock().await; + if let Some(io) = guard.as_mut() { + let exchange = async { + write_provider_frames( + provider, + Some(&mut io.stdin), + &[probe], + "health probe", + ) + .await?; + read_bounded_retained_line(&mut io.reader, self.max_event_line_bytes) + .await + .map_err(|source| RuntimeError::ProcessIo { + provider, + stream: "stdout read", + source, + }) + }; + tokio::time::timeout(health_timeout, exchange) + .await + .ok() + .and_then(std::result::Result::ok) + .flatten() + .and_then(|line| serde_json::from_str::(&line).ok()) + .is_some_and(|frame| { + frame + .pointer("/response/request_id") + .and_then(serde_json::Value::as_str) + == Some(request_id.as_str()) + && frame + .pointer("/response/subtype") + .and_then(serde_json::Value::as_str) + == Some("success") + }) + } else { + false + } + }; + if !healthy { + supervisor.remove_if_same(runtime_id, &candidate).await; + if available_permit.is_none() { + available_permit = candidate.permit.lock().await.take(); + } + candidate.terminate().await?; + existing = None; + claimed_generation = None; + events + .emit(TurnEvent::ProviderProcessStatus { + status: crate::ProviderProcessStatus::Replacing, + message: "The retained Claude process stopped responding before prompt submission; starting a replacement.".to_string(), + }) + .await?; + } + } + } else if provider == Provider::OpenCode { + if let Some(candidate) = existing.clone() { + events + .emit(TurnEvent::ProviderProcessStatus { + status: crate::ProviderProcessStatus::Checking, + message: "Checking the retained OpenCode server before sending the prompt." + .to_string(), + }) + .await?; + claimed_generation = Some(candidate.generation.fetch_add(1, Ordering::AcqRel) + 1); + let health_timeout = supervisor + .inner + .config + .initialization_timeout + .min(Duration::from_secs(3)); + let health = tokio::time::timeout(health_timeout, async { + let streams = adapter.attach(request, &state).await?.ok_or_else(|| { + RuntimeError::Protocol { + provider, + message: + "retained OpenCode adapter did not provide HTTP bridge streams" + .to_string(), + } + })?; + let mut writer = streams.writer; + let mut reader = BufReader::new(streams.reader); + let mut pending_events = Vec::new(); + loop { + let line = + read_bounded_retained_line(&mut reader, self.max_event_line_bytes) + .await + .map_err(|source| RuntimeError::ProcessIo { + provider, + stream: "HTTP bridge read", + source, + })? + .ok_or_else(|| RuntimeError::Protocol { + provider, + message: "OpenCode's pre-prompt bridge closed".to_string(), + })?; + let kind = serde_json::from_str::(&line) + .ok() + .and_then(|frame| { + frame + .get("type") + .and_then(serde_json::Value::as_str) + .map(str::to_owned) + }); + if kind.as_deref() == Some("@subscribed") { + break Ok::<_, RuntimeError>((writer, reader, line, pending_events)); + } + let output = adapter.parse_line(&line, &mut state)?; + if output.terminal || output.interaction.is_some() { + return Err(RuntimeError::Protocol { + provider, + message: "OpenCode failed before prompt submission".to_string(), + }); + } + pending_events.extend(output.events); + write_provider_frames( + provider, + Some(&mut writer), + &output.writes, + "OpenCode preflight", + ) + .await?; + } + }) + .await; + if let Ok(Ok((writer, reader, line, pending_events))) = health { + prepared_protocol_streams = Some((writer, reader)); + prefetched_line = Some(line); + preflight_events = pending_events; + } else { + supervisor.remove_if_same(runtime_id, &candidate).await; + if available_permit.is_none() { + available_permit = candidate.permit.lock().await.take(); + } + candidate.terminate().await?; + existing = None; + claimed_generation = None; + events.emit(TurnEvent::ProviderProcessStatus { + status: crate::ProviderProcessStatus::Replacing, + message: "The retained OpenCode server stopped responding before prompt submission; starting a replacement.".to_string(), + }).await?; + } + } + } let retained = if let Some(existing) = existing { existing } else { @@ -2713,7 +2990,9 @@ impl AgentRuntime { TransportErrorKind::SpawnFailed, self.transport.name(), "retain_process", - "retained Codex capacity changed while preparing the process", + format!( + "retained {provider} capacity changed while preparing the process" + ), true, ), })? @@ -2740,45 +3019,86 @@ impl AgentRuntime { } })?; 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); + // From this point until the process is registered, every fallible + // setup step must explicitly terminate the child. + let setup_result: Result<_> = async { + let mut stdin = if provider == Provider::OpenCode { + let (writer, _reader) = tokio::io::duplex(1); + Box::new(writer) as crate::TransportWriter + } else { + process.take_stdin().ok_or_else(|| RuntimeError::Protocol { + provider, + message: "retained provider 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: format!("retained {provider} process did not expose stdout"), + })?; + let stderr = process + .take_stderr() + .ok_or_else(|| RuntimeError::Protocol { + provider, + message: format!("retained {provider} process did not expose stderr"), + })?; + let (stdin, reader, stdout_task) = if provider == Provider::OpenCode { + let streams = adapter.attach(request, &state).await?.ok_or_else(|| { + RuntimeError::Protocol { + provider, + message: + "retained OpenCode adapter did not provide HTTP bridge streams" + .to_string(), + } + })?; + ( + streams.writer, + BufReader::new(streams.reader), + Some(tokio::spawn(crate::process::bounded_stderr( + stdout, + STDERR_TAIL_BYTES, + ))), + ) + } else { + (stdin, BufReader::new(stdout), None) + }; + Ok((stdin, reader, stdout_task, stderr)) } - 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(), - })?; + .await; + let (stdin, reader, stdout_task, stderr) = match setup_result { + Ok(parts) => parts, + Err(error) => { + let _ = process.terminate().await; + return Err(error); + } + }; let retained = Arc::new(RetainedCodexProcess { fingerprint, io: AsyncMutex::new(Some(RetainedCodexIo { process, stdin, - reader: BufReader::new(stdout), + reader, stderr_task: tokio::spawn(crate::process::bounded_stderr( stderr, STDERR_TAIL_BYTES, )), + stdout_task, })), generation: AtomicU64::new(0), usable: AtomicBool::new(true), + process_hint: adapter.retained_process_hint(&state), + session_id: AsyncMutex::new(request.session_id.clone()), permit: AsyncMutex::new(Some(permit)), }); supervisor @@ -2790,7 +3110,8 @@ impl AgentRuntime { retained }; - let generation = retained.generation.fetch_add(1, Ordering::AcqRel) + 1; + let generation = claimed_generation + .unwrap_or_else(|| retained.generation.fetch_add(1, Ordering::AcqRel) + 1); let mut cleanup = RetainedTurnCleanup { armed: true, supervisor: supervisor.clone(), @@ -2800,33 +3121,67 @@ impl AgentRuntime { 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(), + message: format!("retained {provider} process is no longer available"), })?; 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(), + if provider == Provider::OpenCode { + if let Some((writer, reader)) = prepared_protocol_streams.take() { + io.stdin = writer; + io.reader = reader; + } else { + let streams = adapter.attach(request, &state).await?.ok_or_else(|| { + RuntimeError::Protocol { + provider, + message: + "retained OpenCode adapter did not provide HTTP bridge streams" + .to_string(), + } })?; - write_provider_frames( - provider, - Some(&mut io.stdin), - std::slice::from_ref(&start), - "retained turn start", - ) - .await?; + io.stdin = streams.writer; + io.reader = BufReader::new(streams.reader); + } + } + if provider != Provider::OpenCode { + let start = if provider == Provider::Claude { + spec.initial_stdin.clone() + } else { + adapter.retained_turn_start(&state)? + } + .ok_or_else(|| RuntimeError::Protocol { + provider, + message: format!("configured {provider} adapter cannot start a retained turn"), + })?; + write_provider_frames( + provider, + Some(&mut io.stdin), + std::slice::from_ref(&start), + "retained turn start", + ) + .await?; + } } trace.record(StartupStage::StreamsAttached); + for event in preflight_events { + events.emit(event).await?; + } let result = self - .drive_retained_codex_turn(adapter.as_ref(), request, events, interactions, trace, io) + .drive_retained_codex_turn( + adapter.as_ref(), + request, + events, + interactions, + trace, + io, + supervisor.inner.config, + state, + prefetched_line, + ) .await; drop(io_guard); match result { Ok(result) => { + *retained.session_id.lock().await = result.session_id.clone(); cleanup.disarm(); let runtime_id = runtime_id.clone(); let retained_for_expiry = Arc::clone(&retained); @@ -2836,38 +3191,45 @@ impl AgentRuntime { 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 { + if provider != Provider::OpenCode { + 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; } - 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 { @@ -2883,6 +3245,7 @@ impl AgentRuntime { } } + #[allow(clippy::too_many_arguments)] async fn drive_retained_codex_turn( &self, adapter: &dyn AgentAdapter, @@ -2891,34 +3254,56 @@ impl AgentRuntime { interactions: &dyn InteractionHandler, trace: &StartupTrace, io: &mut RetainedCodexIo, + retention: ProviderProcessRetention, + mut state: AdapterState, + mut prefetched_line: Option, ) -> 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 saw_semantic_activity = false; + let mut turn_submitted = provider != Provider::OpenCode; + let mut last_semantic_activity = tokio::time::Instant::now(); let mut first_text = true; + let mut process_ready = false; 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; + let inactivity_timeout = if !saw_semantic_activity || !turn_submitted { + retention.initialization_timeout + } else { + retention + .active_inactivity_timeout + .unwrap_or(request.timeout) + }; + let line = if let Some(line) = prefetched_line.take() { + Some(line) + } else { + 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, + })?, + _ = tokio::time::sleep_until(last_semantic_activity + inactivity_timeout) => { + return Err(RuntimeError::Timeout { + provider, + seconds: inactivity_timeout.as_secs(), + }); } - 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(), + stderr: format!( + "retained {provider} process exited before completing the turn" + ), provider_code: None, delivery: crate::lifecycle::DeliveryState::PossiblySent, }); @@ -2934,6 +3319,30 @@ impl AgentRuntime { }); } let output = adapter.parse_line(&line, &mut state)?; + let semantic_activity = !output.events.is_empty() + || !output.writes.is_empty() + || output.interaction.is_some() + || output.terminal + || output.turn_submitted; + if output.turn_submitted { + turn_submitted = true; + } + if semantic_activity { + saw_semantic_activity = true; + last_semantic_activity = tokio::time::Instant::now(); + } + if !process_ready + && semantic_activity + && (provider != Provider::OpenCode || output.turn_submitted) + { + events + .emit(TurnEvent::ProviderProcessStatus { + status: crate::ProviderProcessStatus::Ready, + message: format!("The retained {provider} process is ready."), + }) + .await?; + process_ready = true; + } for event in output.events { if first_text && matches!(&event, TurnEvent::TextDelta { text } if !text.is_empty()) { @@ -4607,6 +5016,7 @@ mod tests { }), writes: Vec::new(), terminal: false, + turn_submitted: false, }) } diff --git a/src/types.rs b/src/types.rs index 6c8322f..2851240 100644 --- a/src/types.rs +++ b/src/types.rs @@ -828,6 +828,13 @@ pub enum TurnEvent { /// Human-readable diagnostic. message: String, }, + /// Retained provider process lifecycle progress. + ProviderProcessStatus { + /// Current process lifecycle state. + status: ProviderProcessStatus, + /// Bounded content-free explanation suitable for user feedback. + message: String, + }, /// A sandbox backend identified a denied provider tool step. SandboxAccessDenied { /// Exact managed profile revision, when the turn used one. @@ -851,6 +858,21 @@ pub enum TurnEvent { }, } +/// Lifecycle state for an opt-in retained provider process. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +#[non_exhaustive] +pub enum ProviderProcessStatus { + /// The runtime is checking whether an idle process is responsive. + Checking, + /// A stale process is being replaced before prompt submission. + Replacing, + /// A provider process is ready to receive a turn. + Ready, + /// A retained process was stopped after a failure or lifecycle boundary. + Stopped, +} + /// Human approval requested by an adapter. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct ApprovalRequest { diff --git a/tests/opencode_serve.rs b/tests/opencode_serve.rs index 8d7e4eb..f824d62 100644 --- a/tests/opencode_serve.rs +++ b/tests/opencode_serve.rs @@ -13,6 +13,7 @@ use std::collections::BTreeMap; use std::path::Path; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::time::Duration; @@ -35,12 +36,17 @@ use tokio_util::sync::CancellationToken; /// Turn shape the fixture server plays out once the prompt arrives. #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum Script { + Complete, /// Ask one `bash` permission, then reply according to the decision. Permission, /// Close the event stream mid-turn, the way a crashing server does. Crash, /// Stream one delta and then go quiet, so the turn must be cancelled. Interrupt, + /// Pass health, then never answer session setup. + SessionHang, + /// Hang the second session lookup; the replacement process answers it. + SessionHangOnce, } /// One request the SDK made, recorded for assertions. @@ -53,6 +59,9 @@ struct Recorded { #[derive(Clone)] struct Server { script: Script, + spawns: Arc, + health: Arc>>>, + stoppers: Arc>>, requests: Arc>>, /// Argument vector the SDK asked the transport to spawn. arguments: Arc>>, @@ -64,6 +73,9 @@ impl Server { fn new(script: Script) -> Self { Self { script, + spawns: Arc::new(AtomicUsize::new(0)), + health: Arc::default(), + stoppers: Arc::default(), requests: Arc::new(Mutex::new(Vec::new())), arguments: Arc::new(Mutex::new(Vec::new())), config: Arc::new(Mutex::new(None)), @@ -137,6 +149,7 @@ impl ExecutionTransport for Server { } async fn spawn(&self, request: TransportSpawnRequest) -> TransportResult { + self.spawns.fetch_add(1, Ordering::SeqCst); let arguments = request .command .args @@ -170,11 +183,15 @@ impl ExecutionTransport for Server { let listener = TcpListener::bind(("127.0.0.1", port)) .await .expect("the reserved port is free for the child to bind"); - tokio::spawn(accept_loop( + let health = Arc::new(AtomicBool::new(true)); + self.health.lock().unwrap().push(health.clone()); + let task = tokio::spawn(accept_loop( listener, self.script, Arc::clone(&self.requests), + health, )); + self.stoppers.lock().unwrap().push(task.abort_handle()); let (sdk_stdin, child_input) = duplex(1024); let (sdk_stdout, child_output) = duplex(1024); @@ -193,7 +210,10 @@ impl ExecutionTransport for Server { Some(Box::new(sdk_stdin) as TransportWriter), Box::new(sdk_stdout) as TransportReader, Box::new(sdk_stderr) as TransportReader, - Control::default(), + Control { + stopped: false, + task: Some(task), + }, )) } @@ -216,6 +236,14 @@ impl ExecutionTransport for Server { #[derive(Default)] struct Control { stopped: bool, + task: Option>, +} +impl Drop for Control { + fn drop(&mut self) { + if let Some(task) = self.task.take() { + task.abort(); + } + } } #[async_trait] @@ -235,21 +263,40 @@ impl TransportProcessControl for Control { async fn terminate(&mut self) -> TransportResult<()> { self.stopped = true; + if let Some(task) = self.task.take() { + task.abort(); + let _ = task.await; + } Ok(()) } } -async fn accept_loop(listener: TcpListener, script: Script, requests: Arc>>) { +async fn accept_loop( + listener: TcpListener, + script: Script, + requests: Arc>>, + health: Arc, +) { loop { let Ok((socket, _)) = listener.accept().await else { return; }; - tokio::spawn(handle(socket, script, Arc::clone(&requests))); + tokio::spawn(handle( + socket, + script, + Arc::clone(&requests), + health.clone(), + )); } } /// Read one request, record it, and answer it. -async fn handle(mut socket: TcpStream, script: Script, requests: Arc>>) { +async fn handle( + mut socket: TcpStream, + script: Script, + requests: Arc>>, + health: Arc, +) { let mut reader = BufReader::new(&mut socket); let mut request_line = String::new(); if reader.read_line(&mut request_line).await.unwrap_or(0) == 0 { @@ -287,11 +334,37 @@ async fn handle(mut socket: TcpStream, script: Script, requests: Arc().await; + } + if script == Script::SessionHang && (path.starts_with("/session?") || path == "/session") { + std::future::pending::<()>().await; + } + if script == Script::SessionHangOnce + && (path.starts_with("/session?") + || path == "/session" + || path.split('?').next() == Some("/session/session-fixture")) + && requests + .lock() + .unwrap() + .iter() + .filter(|request| { + request.path.starts_with("/session") && !request.path.contains("/message") + }) + .count() + == 2 + { + std::future::pending::<()>().await; + } let payload = if path == "/global/health" { json!({"healthy": true}).to_string() } else if path == "/app" { "OpenCode UI".to_string() - } else if path.starts_with("/session?") || path == "/session" { + } else if path.starts_with("/session?") + || path == "/session" + || (request_line.starts_with("GET ") + && path.split('?').next() == Some("/session/session-fixture")) + { json!({"id": "session-fixture", "title": "Fixture session"}).to_string() } else { json!({}).to_string() @@ -313,6 +386,12 @@ async fn send(socket: &mut TcpStream, event: &Value) -> bool { /// Play the scripted turn out over a chunked SSE stream. async fn stream_events(mut socket: TcpStream, script: Script, requests: Arc>>) { + let prompts_before = requests + .lock() + .unwrap() + .iter() + .filter(|r| r.path.contains("/message")) + .count(); let _ = socket .write_all( b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nTransfer-Encoding: chunked\r\n\r\n", @@ -320,6 +399,19 @@ async fn stream_events(mut socket: TcpStream, script: Script, requests: Arc prompts_before + { + break; + } + tokio::time::sleep(Duration::from_millis(5)).await; + } // Establish the assistant message every part below belongs to. if !send( &mut socket, @@ -339,6 +431,9 @@ async fn stream_events(mut socket: TcpStream, script: Script, requests: Arc ( + temps_agent_runtime::retained::InProcessRuntimeClient, + temps_agent_runtime::retained::RuntimeHandle, +) { + use temps_agent_runtime::retained::{InProcessRuntimeClient, RuntimeClient, RuntimeSpec}; + let mut builder = AgentRuntime::builder() + .transport(server.clone()) + .provider_process_retention(temps_agent_runtime::ProviderProcessRetention { + max_processes: 1, + idle_timeout: Duration::from_secs(30), + initialization_timeout: Duration::from_secs(3), + active_inactivity_timeout: Some(active), + }); + builder.register(OpenCode::serve()); + let client = InProcessRuntimeClient::new(builder.build().unwrap()); + let handle = client + .acquire(RuntimeSpec::new( + temps_agent_runtime::lifecycle::RuntimeId::new("retained-opencode").unwrap(), + Provider::OpenCode, + std::env::temp_dir(), + )) + .await + .unwrap(); + (client, handle) +} + +#[tokio::test] +async fn retained_server_reuses_port_and_session_for_two_turns() { + use temps_agent_runtime::{ + lifecycle::{InvocationId, RuntimeId}, + retained::{RuntimeClient, TurnInput}, + }; + let server = Server::new(Script::Complete); + let (client, handle) = retained_fixture(&server, Duration::from_secs(3)).await; + for id in ["first", "second"] { + let result = tokio::time::timeout( + Duration::from_secs(5), + handle + .start_turn(TurnInput::new(InvocationId::new(id).unwrap(), id)) + .await + .unwrap() + .wait(), + ) + .await + .unwrap() + .unwrap(); + assert_eq!(result.session_id.as_deref(), Some("session-fixture")); + } + assert_eq!(server.spawns.load(Ordering::SeqCst), 1); + assert_eq!( + server + .requests() + .iter() + .filter(|r| r.path.contains("/message")) + .count(), + 2 + ); + client + .dispose(&RuntimeId::new("retained-opencode").unwrap()) + .await + .unwrap(); +} + +#[tokio::test] +async fn retained_server_stall_is_bounded_and_does_not_replay() { + use temps_agent_runtime::{ + lifecycle::{InvocationId, RuntimeId}, + retained::{RuntimeClient, TurnInput}, + }; + let server = Server::new(Script::Interrupt); + let (client, handle) = retained_fixture(&server, Duration::from_millis(100)).await; + let result = tokio::time::timeout( + Duration::from_secs(5), + handle + .start_turn(TurnInput::new( + InvocationId::new("stalled").unwrap(), + "stalled", + )) + .await + .unwrap() + .wait(), + ) + .await + .unwrap(); + assert!(result.is_err()); + assert_eq!( + server + .requests() + .iter() + .filter(|r| r.path.contains("/message")) + .count(), + 1 + ); + client + .dispose(&RuntimeId::new("retained-opencode").unwrap()) + .await + .unwrap(); +} + +#[tokio::test] +async fn retained_session_setup_hang_uses_initialization_deadline() { + use temps_agent_runtime::{ + lifecycle::{InvocationId, RuntimeId}, + retained::{RuntimeClient, TurnInput}, + }; + let server = Server::new(Script::SessionHang); + // The active allowance is deliberately much longer than the fixture's + // three-second initialization deadline. + let (client, handle) = retained_fixture(&server, Duration::from_secs(30)).await; + let started = tokio::time::Instant::now(); + let result = handle + .start_turn(TurnInput::new( + InvocationId::new("session-hang").unwrap(), + "never submitted", + )) + .await + .unwrap() + .wait() + .await; + assert!(result.is_err()); + assert!(started.elapsed() < Duration::from_secs(5)); + assert_eq!( + server + .requests() + .iter() + .filter(|request| request.path.contains("/message")) + .count(), + 0 + ); + client + .dispose(&RuntimeId::new("retained-opencode").unwrap()) + .await + .unwrap(); +} + +#[tokio::test] +async fn retained_replaces_when_session_setup_hangs_before_prompt() { + use temps_agent_runtime::{ + lifecycle::{InvocationId, RuntimeId}, + retained::{RuntimeClient, TurnInput}, + }; + let server = Server::new(Script::SessionHangOnce); + let (client, handle) = retained_fixture(&server, Duration::from_secs(30)).await; + for id in ["first", "second"] { + handle + .start_turn(TurnInput::new(InvocationId::new(id).unwrap(), id)) + .await + .unwrap() + .wait() + .await + .unwrap(); + } + assert_eq!(server.spawns.load(Ordering::SeqCst), 2); + assert_eq!( + server + .requests() + .iter() + .filter(|request| request.path.contains("/message")) + .count(), + 2 + ); + client + .dispose(&RuntimeId::new("retained-opencode").unwrap()) + .await + .unwrap(); +} + +#[tokio::test] +async fn retained_unresponsive_server_recovers_before_prompt_delivery() { + use temps_agent_runtime::{ + lifecycle::{InvocationId, RuntimeId}, + retained::{RuntimeClient, TurnInput}, + }; + let server = Server::new(Script::Complete); + let (client, handle) = retained_fixture(&server, Duration::from_secs(3)).await; + handle + .start_turn(TurnInput::new(InvocationId::new("first").unwrap(), "first")) + .await + .unwrap() + .wait() + .await + .unwrap(); + server.health.lock().unwrap()[0].store(false, Ordering::SeqCst); + tokio::time::timeout( + Duration::from_secs(8), + handle + .start_turn(TurnInput::new( + InvocationId::new("recovered").unwrap(), + "recovered", + )) + .await + .unwrap() + .wait(), + ) + .await + .unwrap() + .unwrap(); + assert_eq!(server.spawns.load(Ordering::SeqCst), 2); + assert_eq!( + server + .requests() + .iter() + .filter(|r| r.path.contains("/message")) + .count(), + 2 + ); + client + .dispose(&RuntimeId::new("retained-opencode").unwrap()) + .await + .unwrap(); +} + +#[tokio::test] +async fn retained_dead_server_recovers_before_prompt_delivery() { + use temps_agent_runtime::{ + lifecycle::{InvocationId, RuntimeId}, + retained::{RuntimeClient, TurnInput}, + }; + let server = Server::new(Script::Complete); + let (client, handle) = retained_fixture(&server, Duration::from_secs(3)).await; + handle + .start_turn(TurnInput::new(InvocationId::new("first").unwrap(), "first")) + .await + .unwrap() + .wait() + .await + .unwrap(); + server.stoppers.lock().unwrap()[0].abort(); + tokio::time::sleep(Duration::from_millis(10)).await; + tokio::time::timeout( + Duration::from_secs(8), + handle + .start_turn(TurnInput::new( + InvocationId::new("recovered").unwrap(), + "recovered", + )) + .await + .unwrap() + .wait(), + ) + .await + .unwrap() + .unwrap(); + assert_eq!(server.spawns.load(Ordering::SeqCst), 2); + assert_eq!( + server + .requests() + .iter() + .filter(|r| r.path.contains("/message")) + .count(), + 2 + ); + client + .dispose(&RuntimeId::new("retained-opencode").unwrap()) + .await + .unwrap(); +} diff --git a/tests/retained_providers.rs b/tests/retained_providers.rs new file mode 100644 index 0000000..ed867fb --- /dev/null +++ b/tests/retained_providers.rs @@ -0,0 +1,237 @@ +//! Native-process fixtures: no provider account, network model, or real tools. +#![cfg(all(unix, feature = "claude"))] + +use std::{fs, os::unix::fs::PermissionsExt, time::Duration}; +use temps_agent_runtime::{ + lifecycle::{InvocationId, RuntimeFailure, RuntimeId}, + providers::Claude, + retained::{InProcessRuntimeClient, RuntimeClient, RuntimeHandle, RuntimeSpec, TurnInput}, + AgentRuntime, PermissionMode, Provider, ProviderProcessRetention, TurnResult, +}; + +const SCRIPT: &str = r"#!/usr/bin/env python3 +import json,os,sys,time +LOG=LOG_PATH +def log(kind, **data): + with open(LOG,'a') as f:f.write(json.dumps(dict(kind=kind,pid=os.getpid(),**data))+'\n') +def emit(frame): + print(json.dumps(frame),flush=True) +log('spawn') +for line in sys.stdin: + frame=json.loads(line) + if frame.get('type')=='control_request': + log('probe') + emit({'type':'control_response','response':{'request_id':frame['request_id'],'subtype':'success','response':{}}}) + elif frame.get('type')=='user': + prompt=frame['message']['content'][0]['text'];log('prompt',prompt=prompt) + if prompt=='crash':sys.exit(2) + if prompt=='initial-hang': + while True:time.sleep(1) + emit({'type':'system','subtype':'init','session_id':'fixture-session'}) + if prompt=='chatter': + while True: + emit({'type':'unknown_noop'});time.sleep(.01) + if prompt=='hang': + while True:time.sleep(1) + emit({'type':'result','subtype':'success','session_id':'fixture-session','result':'reply:'+prompt,'is_error':False,'num_turns':1,'total_cost_usd':0,'usage':{'input_tokens':1,'output_tokens':1}}) +"; + +struct Fixture { + _dir: tempfile::TempDir, + log: std::path::PathBuf, + client: InProcessRuntimeClient, + handle: RuntimeHandle, +} +impl Fixture { + async fn new(enabled: bool, idle_timeout: Duration) -> Self { + let dir = tempfile::tempdir().unwrap(); + let executable = dir.path().join("claude-fixture"); + let log = dir.path().join("events.jsonl"); + fs::write( + &executable, + SCRIPT.replace("LOG_PATH", &serde_json::to_string(&log).unwrap()), + ) + .unwrap(); + fs::set_permissions(&executable, fs::Permissions::from_mode(0o700)).unwrap(); + let mut builder = AgentRuntime::builder(); + builder.register(Claude::with_executable(executable)); + if enabled { + builder = builder.provider_process_retention(ProviderProcessRetention { + max_processes: 1, + idle_timeout, + initialization_timeout: Duration::from_secs(10), + active_inactivity_timeout: Some(Duration::from_secs(2)), + }); + } + let client = InProcessRuntimeClient::new(builder.build().unwrap()); + let handle = client + .acquire(RuntimeSpec::new( + RuntimeId::new("claude-native").unwrap(), + Provider::Claude, + dir.path(), + )) + .await + .unwrap(); + Self { + _dir: dir, + log, + client, + handle, + } + } + fn events(&self) -> Vec { + fs::read_to_string(&self.log) + .unwrap_or_default() + .lines() + .map(|line| serde_json::from_str(line).unwrap()) + .collect() + } + fn spawns(&self) -> Vec { + self.events() + .into_iter() + .filter(|v| v["kind"] == "spawn") + .map(|v| i32::try_from(v["pid"].as_i64().unwrap()).unwrap()) + .collect() + } + async fn turn(&self, id: &str, prompt: &str) -> Result { + tokio::time::timeout(Duration::from_secs(15), async { + self.handle + .start_turn(TurnInput::new(InvocationId::new(id).unwrap(), prompt)) + .await + .unwrap() + .wait() + .await + }) + .await + .expect("turn must be bounded") + } + async fn dispose(&self) { + self.client + .dispose(&RuntimeId::new("claude-native").unwrap()) + .await + .unwrap(); + } +} + +#[tokio::test] +async fn claude_reuses_first_session_and_health_checks_before_second_prompt() { + let f = Fixture::new(true, Duration::from_secs(30)).await; + assert_eq!( + f.turn("one", "first").await.unwrap().session_id.as_deref(), + Some("fixture-session") + ); + assert_eq!(f.turn("two", "second").await.unwrap().text, "reply:second"); + assert_eq!(f.spawns().len(), 1, "{:?}", f.events()); + assert!(f.events().iter().any(|v| v["kind"] == "probe")); + f.dispose().await; +} +#[tokio::test] +async fn claude_default_still_exits_each_turn() { + let f = Fixture::new(false, Duration::from_secs(30)).await; + assert!(!f.handle.driver_capabilities().retained_process); + f.turn("one", "first").await.unwrap(); + f.turn("two", "second").await.unwrap(); + assert_eq!(f.spawns().len(), 2); + f.dispose().await; +} +#[tokio::test] +async fn claude_permission_change_replaces_with_single_pool_slot() { + let f = Fixture::new(true, Duration::from_secs(30)).await; + f.turn("one", "first").await.unwrap(); + let mut input = TurnInput::new(InvocationId::new("changed").unwrap(), "changed"); + input.permission_mode = Some(PermissionMode::AcceptEdits); + f.handle + .start_turn(input) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(f.spawns().len(), 2); + f.dispose().await; +} +#[tokio::test] +async fn claude_idle_crash_recovers_before_delivery() { + let f = Fixture::new(true, Duration::from_secs(30)).await; + f.turn("one", "first").await.unwrap(); + nix::sys::signal::kill( + nix::unistd::Pid::from_raw(f.spawns()[0]), + nix::sys::signal::Signal::SIGKILL, + ) + .unwrap(); + assert_eq!(f.turn("two", "second").await.unwrap().text, "reply:second"); + assert_eq!(f.spawns().len(), 2); + f.dispose().await; +} +#[tokio::test] +async fn claude_frozen_idle_process_is_replaced_without_duplicate_prompt() { + let f = Fixture::new(true, Duration::from_secs(30)).await; + f.turn("one", "first").await.unwrap(); + nix::sys::signal::kill( + nix::unistd::Pid::from_raw(f.spawns()[0]), + nix::sys::signal::Signal::SIGSTOP, + ) + .unwrap(); + assert_eq!(f.turn("two", "second").await.unwrap().text, "reply:second"); + assert_eq!( + f.events() + .iter() + .filter(|v| v["prompt"] == "second") + .count(), + 1 + ); + f.dispose().await; +} +#[tokio::test] +async fn claude_active_failure_is_not_replayed_and_next_turn_recovers() { + let f = Fixture::new(true, Duration::from_secs(30)).await; + assert!(f.turn("one", "crash").await.is_err()); + assert_eq!(f.turn("two", "second").await.unwrap().text, "reply:second"); + assert_eq!( + f.events().iter().filter(|v| v["prompt"] == "crash").count(), + 1 + ); + f.dispose().await; +} +#[tokio::test] +async fn claude_active_silence_is_bounded_and_runtime_is_reusable() { + let f = Fixture::new(true, Duration::from_secs(30)).await; + assert!(f.turn("one", "hang").await.is_err()); + assert_eq!(f.turn("two", "second").await.unwrap().text, "reply:second"); + f.dispose().await; +} +#[tokio::test] +async fn claude_idle_expiry_resumes_in_a_new_process() { + let f = Fixture::new(true, Duration::from_millis(50)).await; + f.turn("one", "first").await.unwrap(); + tokio::time::sleep(Duration::from_millis(150)).await; + assert_eq!( + f.turn("two", "second").await.unwrap().session_id.as_deref(), + Some("fixture-session") + ); + assert_eq!(f.spawns().len(), 2); + f.dispose().await; +} + +#[tokio::test] +async fn claude_initial_silence_is_bounded_without_replaying_prompt() { + let f = Fixture::new(true, Duration::from_secs(30)).await; + assert!(f.turn("one", "initial-hang").await.is_err()); + assert_eq!( + f.events() + .iter() + .filter(|v| v["prompt"] == "initial-hang") + .count(), + 1 + ); + assert_eq!(f.turn("two", "second").await.unwrap().text, "reply:second"); + f.dispose().await; +} + +#[tokio::test] +async fn claude_noop_frames_do_not_hide_a_stalled_turn() { + let f = Fixture::new(true, Duration::from_secs(30)).await; + assert!(f.turn("one", "chatter").await.is_err()); + assert_eq!(f.turn("two", "second").await.unwrap().text, "reply:second"); + f.dispose().await; +} From 194639350b1fe6e501c197a273da37005a6f8a9e Mon Sep 17 00:00:00 2001 From: David Viejo Date: Wed, 23 Sep 2026 17:59:54 +0200 Subject: [PATCH 4/5] fix: isolate OpenCode retention preflight state --- src/runtime.rs | 17 +++++++++++- tests/opencode_serve.rs | 57 +++++++++++++++++++++++++++++++++++------ 2 files changed, 65 insertions(+), 9 deletions(-) diff --git a/src/runtime.rs b/src/runtime.rs index e7f9b39..10b04c8 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -2881,6 +2881,20 @@ impl AgentRuntime { } } else if provider == Provider::OpenCode { if let Some(candidate) = existing.clone() { + // Handshake parsing can record a terminal failure. Keep it in + // disposable state until every pre-prompt step succeeds so a + // replacement never inherits failure from the stale server. + let mut preflight_state = AdapterState::default(); + preflight_state + .result + .session_id + .clone_from(&request.session_id); + adapter.prepare_retained_turn( + request, + &mut preflight_state, + candidate.process_hint, + )?; + adapter.mark_retained_turn(&mut preflight_state); events .emit(TurnEvent::ProviderProcessStatus { status: crate::ProviderProcessStatus::Checking, @@ -2930,7 +2944,7 @@ impl AgentRuntime { if kind.as_deref() == Some("@subscribed") { break Ok::<_, RuntimeError>((writer, reader, line, pending_events)); } - let output = adapter.parse_line(&line, &mut state)?; + let output = adapter.parse_line(&line, &mut preflight_state)?; if output.terminal || output.interaction.is_some() { return Err(RuntimeError::Protocol { provider, @@ -2949,6 +2963,7 @@ impl AgentRuntime { }) .await; if let Ok(Ok((writer, reader, line, pending_events))) = health { + state = preflight_state; prepared_protocol_streams = Some((writer, reader)); prefetched_line = Some(line); preflight_events = pending_events; diff --git a/tests/opencode_serve.rs b/tests/opencode_serve.rs index f824d62..e7d9444 100644 --- a/tests/opencode_serve.rs +++ b/tests/opencode_serve.rs @@ -47,6 +47,8 @@ enum Script { SessionHang, /// Hang the second session lookup; the replacement process answers it. SessionHangOnce, + /// Close the second session lookup; the replacement process answers it. + SessionCloseOnce, } /// One request the SDK made, recorded for assertions. @@ -340,11 +342,11 @@ async fn handle( if script == Script::SessionHang && (path.starts_with("/session?") || path == "/session") { std::future::pending::<()>().await; } - if script == Script::SessionHangOnce - && (path.starts_with("/session?") - || path == "/session" - || path.split('?').next() == Some("/session/session-fixture")) - && requests + let is_session_lookup = path.starts_with("/session?") + || path == "/session" + || path.split('?').next() == Some("/session/session-fixture"); + let session_lookup_count = || { + requests .lock() .unwrap() .iter() @@ -352,10 +354,13 @@ async fn handle( request.path.starts_with("/session") && !request.path.contains("/message") }) .count() - == 2 - { + }; + if script == Script::SessionHangOnce && is_session_lookup && session_lookup_count() == 2 { std::future::pending::<()>().await; } + if script == Script::SessionCloseOnce && is_session_lookup && session_lookup_count() == 2 { + return; + } let payload = if path == "/global/health" { json!({"healthy": true}).to_string() } else if path == "/app" { @@ -431,7 +436,10 @@ async fn stream_events(mut socket: TcpStream, script: Script, requests: Arc Date: Thu, 24 Sep 2026 18:48:54 +0200 Subject: [PATCH 5/5] fix(retained): keep provider session after interrupted turns --- src/retained.rs | 151 +++++++++++++++++++++++++++++++++++++++++++++--- 1 file changed, 144 insertions(+), 7 deletions(-) diff --git a/src/retained.rs b/src/retained.rs index 02d52b2..eaa9ab2 100644 --- a/src/retained.rs +++ b/src/retained.rs @@ -1159,6 +1159,7 @@ impl RuntimeEntry { invocation_id: invocation_id.clone(), sequence: AtomicU64::new(0), sender: event_sender, + announced_session_id: std::sync::Mutex::new(None), }; let entry = Arc::clone(self); let task_terminal = Arc::clone(&terminal); @@ -1236,13 +1237,22 @@ impl RuntimeEntry { { state.active = None; } - if let Ok(turn) = &result { - if let Some(session_id) = &turn.session_id { - state.provider_session_id = Some(session_id.clone()); - } - state.detail = None; - } else if let Err(failure) = &result { - state.detail = Some(failure.message.clone()); + // A provider session exists as soon as the provider announces + // it, not only once the turn succeeds. An interrupted or + // failed turn still leaves that session on disk with the + // conversation so far; forgetting it would make the next + // turn start a fresh session with no history. + let session_id = result + .as_ref() + .ok() + .and_then(|turn| turn.session_id.clone()) + .or_else(|| sink.announced_session_id()); + if let Some(session_id) = session_id { + state.provider_session_id = Some(session_id); + } + match &result { + Ok(_) => state.detail = None, + Err(failure) => state.detail = Some(failure.message.clone()), } state.status = if entry.disposed.load(Ordering::Acquire) { RuntimeStatus::Stopped @@ -1514,17 +1524,32 @@ struct ChannelEventSink { invocation_id: InvocationId, sequence: AtomicU64, sender: mpsc::Sender, + /// Latest `SessionStarted` the provider emitted during this invocation. + announced_session_id: std::sync::Mutex>, } #[async_trait] impl EventSink for ChannelEventSink { async fn emit(&self, event: TurnEvent) -> crate::Result<()> { + if let TurnEvent::SessionStarted { session_id, .. } = &event { + *self + .announced_session_id + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(session_id.clone()); + } self.emit_runtime(RuntimeEvent::ProviderEvent { event }) .await } } impl ChannelEventSink { + fn announced_session_id(&self) -> Option { + self.announced_session_id + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + } + async fn emit_runtime(&self, event: RuntimeEvent) -> crate::Result<()> { let envelope = EventEnvelope { schema_version: 1, @@ -2044,6 +2069,118 @@ mod tests { } } + /// Announces a provider session on the first turn, then blocks until the + /// turn is interrupted -- a user sending a follow-up mid-turn. Later + /// turns finish immediately and echo the session they were asked to + /// resume. + struct SessionThenBlockExecutor { + requested_sessions: StdMutex>>, + } + + #[async_trait] + impl RuntimeTurnExecutor for SessionThenBlockExecutor { + fn capabilities(&self, _provider: Provider) -> RuntimeDriverCapabilities { + RuntimeDriverCapabilities { + retained_process: true, + session_resume: true, + ..RuntimeDriverCapabilities::default() + } + } + + fn configuration_impact(&self, _key: RuntimeConfigurationKey) -> ConfigurationImpact { + ConfigurationImpact::Live + } + + async fn execute( + &self, + request: TurnRequest, + events: &dyn EventSink, + _interactions: Option<&dyn InteractionHandler>, + ) -> crate::Result { + let first = { + let mut sessions = self + .requested_sessions + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + sessions.push(request.session_id.clone()); + sessions.len() == 1 + }; + if first { + events + .emit(TurnEvent::SessionStarted { + session_id: "session-interrupted".to_owned(), + title: None, + }) + .await?; + request.cancellation.cancelled().await; + return Err(RuntimeError::Cancelled { + provider: request.provider, + }); + } + Ok(TurnResult { + status: RunStatus::Succeeded, + text: "ok".to_owned(), + reasoning: None, + session_id: request.session_id.clone(), + session_title: None, + model: None, + usage: Usage::default(), + }) + } + } + + #[tokio::test] + async fn interrupted_turn_keeps_the_session_the_provider_already_started() { + let executor = Arc::new(SessionThenBlockExecutor { + requested_sessions: StdMutex::new(Vec::new()), + }); + let client = InProcessRuntimeClient::from_executor(executor.clone()); + let handle = client + .acquire(runtime_spec("runtime-interrupted-session")) + .await + .expect("acquire runtime"); + + let turn = handle + .start_turn(turn_input("turn-interrupted", "investigate the alerts")) + .await + .expect("start first turn"); + let interrupt = turn.interrupt_handle(); + let (mut events, completion) = turn.into_parts(); + loop { + let envelope = events.next().await.expect("session announcement"); + if matches!( + envelope.event, + RuntimeEvent::ProviderEvent { + event: TurnEvent::SessionStarted { .. } + } + ) { + break; + } + } + assert_eq!(interrupt.interrupt().await, InterruptOutcome::Interrupted); + while events.next().await.is_some() {} + let failure = completion.wait().await.expect_err("interrupted turn"); + assert_eq!(failure.kind, RuntimeFailureKind::Cancelled); + + handle + .start_turn(turn_input("turn-follow-up", "keep going")) + .await + .expect("start follow-up turn") + .wait() + .await + .expect("follow-up result"); + + let sessions = executor + .requested_sessions + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + assert_eq!( + sessions.as_slice(), + [None, Some("session-interrupted".to_owned())], + "a follow-up after an interrupted turn must resume the provider session, not start a fresh one" + ); + } + fn runtime_spec(id: &str) -> RuntimeSpec { RuntimeSpec::new( RuntimeId::new(id).expect("valid runtime identifier"),