diff --git a/CHANGELOG.md b/CHANGELOG.md index 3f04fb3..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 @@ -30,6 +35,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..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: Proposed +- Status: Accepted (Claude, Codex, and OpenCode retained turns implemented; explicit prewarming deferred) - Date: 2026-09-23 ## Problem @@ -16,42 +16,39 @@ 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. `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. -## 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 @@ -77,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 8d4e6bf..2b07eb4 100644 --- a/docs/reference/retained-runtimes.md +++ b/docs/reference/retained-runtimes.md @@ -11,6 +11,52 @@ 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. +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::{Claude, Codex, OpenCode}; +use temps_agent_runtime::{AgentRuntime, ProviderProcessRetention}; + +let mut builder = AgentRuntime::builder() + .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 +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. +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, harness, launch-context, environment, and timeout changes live. Provider, diff --git a/src/adapter.rs b/src/adapter.rs index 622088a..73f1f85 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, @@ -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. @@ -254,6 +256,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 +394,36 @@ 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. + 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..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}; +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/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/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 564f2dd..eaa9ab2 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,7 @@ impl RuntimeTurnExecutor for AgentRuntime { fn capabilities(&self, provider: Provider) -> RuntimeDriverCapabilities { let permissions = self.permission_support(provider).ok(); RuntimeDriverCapabilities { - retained_process: false, + retained_process: self.process_retention_enabled(provider), session_resume: true, live_interactions: permissions .is_some_and(|support| support.live_approvals || support.live_questions), @@ -501,6 +521,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 +688,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 +1036,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 +1051,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( @@ -1109,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); @@ -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( @@ -1181,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 @@ -1459,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, @@ -1989,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"), diff --git a/src/runtime.rs b/src/runtime.rs index 5a1cadb..10b04c8 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,255 @@ 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, +} + +/// 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 { + 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: ProviderProcessRetention, + permits: Arc, + processes: AsyncMutex>>, +} + +impl CodexProcessSupervisor { + fn new(config: ProviderProcessRetention) -> 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, +} + +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>, +} + +struct RetainedCodexIo { + process: crate::TransportProcess, + stdin: crate::TransportWriter, + reader: BufReader, + stderr_task: tokio::task::JoinHandle>, + stdout_task: Option>>, +} + +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(); + if let Some(stdout_task) = io.stdout_task.take() { + stdout_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 +581,8 @@ pub struct AgentRuntimeBuilder { max_prompt_bytes: usize, max_event_line_bytes: usize, startup_observer: Option>, + codex_process_retention: Option, + provider_process_retention: Option, } impl AgentRuntimeBuilder { @@ -344,6 +596,8 @@ 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, + provider_process_retention: None, }; #[cfg(feature = "claude")] builder.register(crate::providers::Claude::default()); @@ -400,6 +654,22 @@ 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 + } + + /// 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 { @@ -414,6 +684,21 @@ 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()?; + } + 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, @@ -421,6 +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: retention.map(CodexProcessSupervisor::new), + provider_process_retention: self.provider_process_retention, }) } } @@ -738,6 +1025,8 @@ pub struct AgentRuntime { max_prompt_bytes: usize, max_event_line_bytes: usize, startup_observer: Option>, + codex_process_retention: Option, + provider_process_retention: Option, } impl AgentRuntime { @@ -1918,12 +2207,62 @@ impl AgentRuntime { result } + 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) + .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.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 +2323,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 +2660,778 @@ impl AgentRuntime { }) } + #[allow(clippy::too_many_arguments)] + async fn run_retained_codex_process( + &self, + supervisor: &CodexProcessSupervisor, + runtime_id: &crate::lifecycle::RuntimeId, + adapter: Arc, + request: &TurnRequest, + events: &dyn EventSink, + interactions: &dyn InteractionHandler, + trace: &StartupTrace, + ) -> Result { + let provider = request.provider; + let initial_process = supervisor + .inner + .processes + .lock() + .await + .get(runtime_id) + .cloned(); + let reserved_permit = if initial_process.is_some() { + None + } else { + match supervisor.inner.permits.clone().try_acquire_owned() { + Ok(permit) => Some(permit), + Err(_) => { + // Decide the cold fallback before adapter or sandbox + // preparation: preparation may own resources and is not + // required to be side-effect-free. + return self + .run_process(adapter, request, events, interactions, trace) + .await; + } + } + }; + let mut state = AdapterState::default(); + state.result.session_id.clone_from(&request.session_id); + 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 { + 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 provider != Provider::OpenCode + && (!spec.interactive_stdin || !capabilities.interactive_stdin) + { + return Err(RuntimeError::TransportCapabilityUnavailable { + provider, + transport: self.transport.name().to_string(), + capability: "interactive_stdin", + message: format!("retained {provider} requires writable provider stdin"), + }); + } + if !capabilities.process_tree_termination { + return Err(RuntimeError::TransportCapabilityUnavailable { + provider, + transport: self.transport.name().to_string(), + capability: "process_tree_termination", + message: format!("retained {provider} requires complete process-tree termination"), + }); + } + let fingerprint = RetainedCodexFingerprint { + command: retained_fingerprint_command(provider, spec.clone()), + working_directory: request.working_directory.clone(), + model: request.model.clone(), + reasoning: request.reasoning.clone(), + permission_mode: request.permission_mode.clone(), + harness_options: request.harness_options.clone(), + launch_context: request.launch_context.clone(), + auto_compaction: request.auto_compaction, + required_sandbox_capabilities: request.required_sandbox_capabilities, + }; + + let previous = { + let mut processes = supervisor.inner.processes.lock().await; + match processes.get(runtime_id) { + Some(process) + if process.fingerprint == fingerprint + && process.usable.load(Ordering::Acquire) => + { + None + } + Some(_) => processes.remove(runtime_id), + None => None, + } + }; + let mut available_permit = reserved_permit; + if let Some(previous) = previous { + // Transfer the slot before teardown. This makes a configuration + // replacement atomic with respect to pool capacity: another + // runtime cannot steal the released slot after sandbox preparation. + if available_permit.is_none() { + available_permit = previous.permit.lock().await.take(); + } + previous.terminate().await?; + } + + let 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() { + // 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, + 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 preflight_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 { + state = preflight_state; + 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 { + if available_permit.is_none() { + if let Some(initial) = &initial_process { + available_permit = initial.permit.lock().await.take(); + initial.terminate().await?; + } + } + let permit = if let Some(permit) = available_permit { + permit + } else { + supervisor + .inner + .permits + .clone() + .try_acquire_owned() + .map_err(|_| RuntimeError::Transport { + provider, + source: TransportError::new( + TransportErrorKind::SpawnFailed, + self.transport.name(), + "retain_process", + format!( + "retained {provider} 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); + // 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)) + } + .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, + 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 + .inner + .processes + .lock() + .await + .insert(runtime_id.clone(), Arc::clone(&retained)); + retained + }; + + let generation = claimed_generation + .unwrap_or_else(|| 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: format!("retained {provider} process is no longer available"), + })?; + let reused = generation > 1; + if reused { + 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(), + } + })?; + 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, + 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); + 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; + 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; + } + } + }); + } + 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), + } + } + + #[allow(clippy::too_many_arguments)] + async fn drive_retained_codex_turn( + &self, + adapter: &dyn AgentAdapter, + request: &TurnRequest, + events: &dyn EventSink, + interactions: &dyn InteractionHandler, + trace: &StartupTrace, + io: &mut RetainedCodexIo, + retention: ProviderProcessRetention, + mut state: AdapterState, + mut prefetched_line: Option, + ) -> Result { + let provider = request.provider; + 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 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(), + }); + } + } + }; + let Some(line) = line else { + return Err(RuntimeError::ProcessFailed { + provider, + kind: ProviderProcessErrorKind::Unknown, + exit_code: None, + stderr: format!( + "retained {provider} process exited before completing the turn" + ), + 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)?; + 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()) + { + 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, @@ -3900,6 +5031,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/codex_app_server.rs b/tests/codex_app_server.rs index b65b0b8..3d436a7 100644 --- a/tests/codex_app_server.rs +++ b/tests/codex_app_server.rs @@ -9,18 +9,23 @@ use std::collections::BTreeMap; use std::path::Path; +use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::time::Duration; use async_trait::async_trait; use serde_json::{json, Value}; +use temps_agent_runtime::lifecycle::{InvocationId, RuntimeId}; use temps_agent_runtime::providers::{Codex, CodexTurnMode}; use temps_agent_runtime::retained::TurnAttachment; +use temps_agent_runtime::retained::{ + InProcessRuntimeClient, RuntimeClient, RuntimeSpec, TurnInput, +}; use temps_agent_runtime::{ - AgentRuntime, ApprovalDecision, ApprovalRequest, EventSink, ExecutionTransport, - InteractionHandler, McpServerConfig, PermissionMode, Provider, ProviderReadiness, - QuestionAnswer, QuestionRequest, Result, RuntimeError, SandboxCapabilities, SecretString, - TransportCapabilities, TransportError, TransportErrorKind, TransportExitStatus, + AgentRuntime, ApprovalDecision, ApprovalRequest, CodexProcessRetention, EventSink, + ExecutionTransport, InteractionHandler, McpServerConfig, PermissionMode, Provider, + ProviderReadiness, QuestionAnswer, QuestionRequest, Result, RuntimeError, SandboxCapabilities, + SecretString, TransportCapabilities, TransportError, TransportErrorKind, TransportExitStatus, TransportProcess, TransportProcessControl, TransportProcessHandle, TransportReader, TransportReadinessRequest, TransportResult, TransportSpawnRequest, TransportWriter, TurnEvent, TurnRequest, @@ -39,6 +44,10 @@ enum Script { AsyncQuestion, /// Stream one delta and then wait for `turn/interrupt`. Interrupt, + /// Exit after accepting a turn, before a terminal notification. + Crash, + /// Complete, then emit an old-turn frame while the process is idle. + DelayedIdleFrame, } #[derive(Clone)] @@ -48,6 +57,8 @@ struct AppServer { frames: Arc>>, /// Arguments the SDK asked the transport to spawn `codex` with. arguments: Arc>>, + spawns: Arc, + terminations: Arc, } impl AppServer { @@ -56,6 +67,8 @@ impl AppServer { script, frames: Arc::new(Mutex::new(Vec::new())), arguments: Arc::new(Mutex::new(Vec::new())), + spawns: Arc::new(AtomicUsize::new(0)), + terminations: Arc::new(AtomicUsize::new(0)), } } @@ -67,6 +80,14 @@ impl AppServer { self.frames.lock().unwrap().clone() } + fn spawn_count(&self) -> usize { + self.spawns.load(Ordering::Acquire) + } + + fn termination_count(&self) -> usize { + self.terminations.load(Ordering::Acquire) + } + fn method_frame(&self, method: &str) -> Option { self.frames() .into_iter() @@ -120,6 +141,7 @@ impl ExecutionTransport for AppServer { } async fn spawn(&self, request: TransportSpawnRequest) -> TransportResult { + self.spawns.fetch_add(1, Ordering::AcqRel); assert_eq!( request .command @@ -143,6 +165,7 @@ impl ExecutionTransport for AppServer { drop(server_stderr); let script = self.script; let frames = Arc::clone(&self.frames); + let terminations = Arc::clone(&self.terminations); tokio::spawn(async move { serve(script, frames, server_input, server_output).await }); Ok(TransportProcess::new( TransportProcessHandle { @@ -153,7 +176,7 @@ impl ExecutionTransport for AppServer { Some(Box::new(sdk_stdin) as TransportWriter), Box::new(sdk_stdout) as TransportReader, Box::new(sdk_stderr) as TransportReader, - Control, + Control { terminations }, )) } @@ -172,7 +195,9 @@ impl ExecutionTransport for AppServer { } } -struct Control; +struct Control { + terminations: Arc, +} #[async_trait] impl TransportProcessControl for Control { @@ -184,6 +209,7 @@ impl TransportProcessControl for Control { } async fn terminate(&mut self) -> TransportResult<()> { + self.terminations.fetch_add(1, Ordering::AcqRel); Ok(()) } } @@ -236,10 +262,14 @@ async fn serve( json!({"jsonrpc":"2.0","id":id,"result":{"turn":{"id":"turn-1"}}}), ) .await; + if script == Script::Crash { + return; + } send( &mut output, json!({"jsonrpc":"2.0","method":"item/agentMessage/delta", - "params":{"itemId":"item-1","delta":"Working"}}), + "params":{"threadId":"thread-fixture","turnId":"turn-1", + "itemId":"item-1","delta":"Working"}}), ) .await; for frame in opening_frames(script) { @@ -251,6 +281,16 @@ async fn serve( for frame in completion_frames(script) { send(&mut output, frame).await; } + if script == Script::DelayedIdleFrame { + tokio::time::sleep(Duration::from_millis(10)).await; + send( + &mut output, + json!({"jsonrpc":"2.0","method":"item/agentMessage/delta", + "params":{"threadId":"thread-fixture","turnId":"turn-1", + "itemId":"late","delta":"STALE"}}), + ) + .await; + } } // `initialize` and `turn/interrupt` need only a bare acknowledgement. _ => { @@ -265,7 +305,7 @@ fn opening_frames(script: Script) -> Vec { Script::Approval => vec![ json!({"jsonrpc":"2.0","method":"item/started","params":{"item":{ "id":"item-2","type":"commandExecution","command":"cargo test","status":"inProgress" - }}}), + },"threadId":"thread-fixture","turnId":"turn-1"}}), json!({"jsonrpc":"2.0","id":"server-1","method":"item/commandExecution/requestApproval", "params":{"threadId":"thread-fixture","turnId":"turn-1","itemId":"item-2", "command":"cargo test","startedAtMs":1}}), @@ -278,14 +318,14 @@ fn opening_frames(script: Script) -> Vec { "options":[{"label":"Banana","description":"Yellow"}, {"label":"Plantain","description":"Also yellow"}]}]} })], - Script::AsyncQuestion => vec![json!({ + Script::AsyncQuestion | Script::DelayedIdleFrame => vec![json!({ "jsonrpc":"2.0","id":"server-3","method":"item/tool/requestUserInput", "params":{"threadId":"thread-fixture","turnId":"turn-1","itemId":"item-4", "isBlocking":false, "questions":[{"id":"q2","header":"Theme","question":"Dark or light?", "options":[{"label":"Dark","description":"Dim"}]}]} })], - Script::Interrupt => Vec::new(), + Script::Interrupt | Script::Crash => Vec::new(), } } @@ -296,12 +336,12 @@ fn completion_frames(script: Script) -> Vec { json!({"jsonrpc":"2.0","method":"item/completed","params":{"item":{ "id":"item-2","type":"commandExecution","command":"cargo test", "status":"completed","aggregatedOutput":"ok","exitCode":0 - }}}), + },"threadId":"thread-fixture","turnId":"turn-1"}}), ); } frames.extend([ json!({"jsonrpc":"2.0","method":"thread/tokenUsage/updated","params":{ - "threadId":"thread-fixture", + "threadId":"thread-fixture","turnId":"turn-1", "tokenUsage":{"last":{"inputTokens":120,"outputTokens":34,"totalTokens":154}, "modelContextWindow":272_000} }}), @@ -369,6 +409,35 @@ fn runtime(transport: AppServer) -> AgentRuntime { builder.build().unwrap() } +fn retained_runtime(transport: AppServer, idle_timeout: Duration) -> AgentRuntime { + let mut builder = AgentRuntime::builder() + .transport(transport) + .codex_process_retention(CodexProcessRetention { + max_processes: 2, + idle_timeout, + }); + builder.register(Codex::app_server()); + builder.build().unwrap() +} + +async fn retained_handle( + runtime: AgentRuntime, +) -> ( + InProcessRuntimeClient, + temps_agent_runtime::retained::RuntimeHandle, +) { + let client = InProcessRuntimeClient::new(runtime); + let handle = client + .acquire(RuntimeSpec::new( + RuntimeId::new("codex-fixture-runtime").unwrap(), + Provider::Codex, + ".", + )) + .await + .unwrap(); + (client, handle) +} + fn request() -> TurnRequest { let mut request = TurnRequest::new(Provider::Codex, ".", "review the workspace"); request.permission_mode = PermissionMode::Default; @@ -400,6 +469,344 @@ async fn the_app_server_mode_advertises_live_interactions() { assert!(turn.native_image_attachments); } +#[tokio::test] +async fn retained_client_reuses_one_app_server_for_two_turns() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + assert!(handle.driver_capabilities().retained_process); + + let first = handle + .start_turn(TurnInput::new( + InvocationId::new("turn-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(first.session_id.as_deref(), Some("thread-fixture")); + + let second = handle + .start_turn(TurnInput::new( + InvocationId::new("turn-two").unwrap(), + "second", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(second.session_id.as_deref(), Some("thread-fixture")); + assert_eq!(transport.spawn_count(), 1); + assert_eq!( + transport + .frames() + .iter() + .filter(|frame| frame.get("method").and_then(Value::as_str) == Some("initialize")) + .count(), + 1 + ); +} + +#[tokio::test] +async fn retained_processes_are_disabled_by_default() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = runtime(transport.clone()); + let (_client, handle) = retained_handle(runtime).await; + assert!(!handle.driver_capabilities().retained_process); + + for invocation in ["default-one", "default-two"] { + handle + .start_turn(TurnInput::new( + InvocationId::new(invocation).unwrap(), + invocation, + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + } + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn changed_permission_replaces_the_retained_process() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + handle + .start_turn(TurnInput::new( + InvocationId::new("permission-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + let mut changed = TurnInput::new(InvocationId::new("permission-two").unwrap(), "second"); + changed.permission_mode = Some(PermissionMode::FullAccess); + handle + .start_turn(changed) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn configuration_replacement_transfers_its_only_pool_slot() { + let transport = AppServer::new(Script::AsyncQuestion); + let mut builder = AgentRuntime::builder() + .transport(transport.clone()) + .codex_process_retention(CodexProcessRetention { + max_processes: 1, + idle_timeout: Duration::from_secs(30), + }); + builder.register(Codex::app_server()); + let (_client, handle) = retained_handle(builder.build().unwrap()).await; + handle + .start_turn(TurnInput::new( + InvocationId::new("replace-only-slot-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + let mut changed = TurnInput::new( + InvocationId::new("replace-only-slot-two").unwrap(), + "second", + ); + changed.permission_mode = Some(PermissionMode::FullAccess); + handle + .start_turn(changed) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(transport.spawn_count(), 2); + assert!(transport.termination_count() >= 1); +} + +#[tokio::test] +async fn idle_expiry_and_dispose_terminate_retained_processes() { + let transport = AppServer::new(Script::AsyncQuestion); + let runtime = retained_runtime(transport.clone(), Duration::from_millis(20)); + let (client, handle) = retained_handle(runtime).await; + handle + .start_turn(TurnInput::new( + InvocationId::new("idle-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + tokio::time::sleep(Duration::from_millis(60)).await; + handle + .start_turn(TurnInput::new( + InvocationId::new("idle-two").unwrap(), + "second", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(transport.spawn_count(), 2); + let runtime_id = handle.runtime_id().clone(); + client.dispose(&runtime_id).await.unwrap(); + client + .acquire(RuntimeSpec::new(runtime_id, Provider::Codex, ".")) + .await + .expect("a disposed runtime identifier can be acquired again"); + assert!(transport.termination_count() >= 2); +} + +#[tokio::test] +async fn process_capacity_does_not_starve_an_existing_retained_runtime() { + let transport = AppServer::new(Script::AsyncQuestion); + let mut builder = AgentRuntime::builder() + .transport(transport.clone()) + .codex_process_retention(CodexProcessRetention { + max_processes: 1, + idle_timeout: Duration::from_secs(30), + }); + builder.register(Codex::app_server()); + let client = InProcessRuntimeClient::new(builder.build().unwrap()); + let first = client + .acquire(RuntimeSpec::new( + RuntimeId::new("capacity-one").unwrap(), + Provider::Codex, + ".", + )) + .await + .unwrap(); + let second = client + .acquire(RuntimeSpec::new( + RuntimeId::new("capacity-two").unwrap(), + Provider::Codex, + ".", + )) + .await + .unwrap(); + first + .start_turn(TurnInput::new( + InvocationId::new("capacity-first").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + second + .start_turn(TurnInput::new( + InvocationId::new("capacity-rejected").unwrap(), + "second runtime", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + first + .start_turn(TurnInput::new( + InvocationId::new("capacity-reuse").unwrap(), + "reuse", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn a_timed_out_turn_is_never_reused() { + let transport = AppServer::new(Script::Interrupt); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let client = InProcessRuntimeClient::new(runtime); + let mut spec = RuntimeSpec::new( + RuntimeId::new("timeout-runtime").unwrap(), + Provider::Codex, + ".", + ); + spec.turn_timeout = Duration::from_millis(20); + let handle = client.acquire(spec).await.unwrap(); + for invocation in ["timeout-one", "timeout-two"] { + let failure = handle + .start_turn(TurnInput::new( + InvocationId::new(invocation).unwrap(), + invocation, + )) + .await + .unwrap() + .wait() + .await + .unwrap_err(); + assert_eq!( + failure.kind, + temps_agent_runtime::lifecycle::RuntimeFailureKind::Timeout + ); + } + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn a_crashed_process_is_replaced_for_the_next_turn() { + let transport = AppServer::new(Script::Crash); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + for invocation in ["crash-one", "crash-two"] { + assert!(handle + .start_turn(TurnInput::new( + InvocationId::new(invocation).unwrap(), + invocation, + )) + .await + .unwrap() + .wait() + .await + .is_err()); + } + assert_eq!(transport.spawn_count(), 2); +} + +#[tokio::test] +async fn an_idle_frame_retires_the_process_before_the_next_turn() { + let transport = AppServer::new(Script::DelayedIdleFrame); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + let first = handle + .start_turn(TurnInput::new( + InvocationId::new("late-one").unwrap(), + "first", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert!(!first.text.contains("STALE")); + tokio::time::sleep(Duration::from_millis(50)).await; + let second = handle + .start_turn(TurnInput::new( + InvocationId::new("late-two").unwrap(), + "second", + )) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert!(!second.text.contains("STALE")); + assert_eq!(transport.spawn_count(), 2); + assert!(transport.termination_count() >= 1); +} + +#[tokio::test] +async fn cancellation_retires_the_process_before_an_immediate_retry() { + let transport = AppServer::new(Script::Interrupt); + let runtime = retained_runtime(transport.clone(), Duration::from_secs(30)); + let (_client, handle) = retained_handle(runtime).await; + let first = handle + .start_turn(TurnInput::new( + InvocationId::new("cancel-one").unwrap(), + "first", + )) + .await + .unwrap(); + transport.wait_for("turn/start").await.unwrap(); + first.interrupt().await; + + let second = handle + .start_turn(TurnInput::new( + InvocationId::new("cancel-two").unwrap(), + "second", + )) + .await + .unwrap(); + for _ in 0..100 { + if transport.spawn_count() == 2 { + break; + } + tokio::time::sleep(Duration::from_millis(5)).await; + } + assert_eq!(transport.spawn_count(), 2); + second.interrupt().await; +} + #[tokio::test] async fn an_approval_is_accepted_through_the_interaction_handler() { let transport = AppServer::new(Script::Approval); diff --git a/tests/opencode_serve.rs b/tests/opencode_serve.rs index 8d7e4eb..e7d9444 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,19 @@ 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, + /// Close the second session lookup; the replacement process answers it. + SessionCloseOnce, } /// One request the SDK made, recorded for assertions. @@ -53,6 +61,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 +75,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 +151,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 +185,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 +212,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 +238,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 +265,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 +336,40 @@ 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; + } + 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() + .filter(|request| { + request.path.starts_with("/session") && !request.path.contains("/message") + }) + .count() + }; + 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" { "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 +391,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 +404,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 +436,12 @@ 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_replaces_after_preprompt_bridge_error_without_stale_failure() { + use temps_agent_runtime::{ + lifecycle::{InvocationId, RuntimeId}, + retained::{RuntimeClient, TurnInput}, + }; + let server = Server::new(Script::SessionCloseOnce); + let (client, handle) = retained_fixture(&server, Duration::from_secs(30)).await; + for id in ["first", "second"] { + let result = handle + .start_turn(TurnInput::new(InvocationId::new(id).unwrap(), id)) + .await + .unwrap() + .wait() + .await + .unwrap(); + assert_eq!(result.text, "fixture reply"); + } + 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; +}