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