diff --git a/packages/paperclip-runner/README.md b/packages/paperclip-runner/README.md index 619485c0fa..2342a79dc6 100644 --- a/packages/paperclip-runner/README.md +++ b/packages/paperclip-runner/README.md @@ -9,13 +9,14 @@ reducer oracle. It also contains a package-local Rust runner, scripted fake harness, bounded process supervisor, cross-language replay oracle, and durable PRP transport. The transport authenticates and encrypts loopback WebSocket sessions, persists an ACK-driven outbox and command journal, and reconnects with -a short-lived lease. No server code starts or invokes it. The package does not -add a provider, server adapter, semantic-tool authorization, or production -Paperclip behavior. +a short-lived lease. The Rust runner now includes a Codex-only app-server +provider bridge with durable thread resume, cancellation, structured questions, +and provider-neutral event normalization. No server code starts or invokes it. +The package does not add a server adapter, semantic-tool authorization, or +production Paperclip behavior. -The first provider scope is Codex. The protocol contains provider-neutral event -and semantic receipt shapes, but their presence does not authorize a tool or -enable a provider. +The first and only installed provider is Codex. Dynamic semantic tools remain +undiscoverable because the catalog and authorization layers have not landed. The root export is intentionally narrow. The `./testing` entry point and package release boundary will arrive with the later package-boundary change. @@ -40,6 +41,8 @@ Durability and failure semantics are documented in [`runner/DURABLE_TRANSPORT.md`](runner/DURABLE_TRANSPORT.md). The fault suite drops a connection before its event ACK, reconnects with the bound lease, replays the same event, and proves the duplicated command effect ran once. +Codex launch, resume, cancellation, and normalization behavior is documented in +[`runner/CODEX_PROVIDER.md`](runner/CODEX_PROVIDER.md). Use `generate:protocol-manifest` after a schema or fixture change, `generate:protocol-types` after a schema change, and diff --git a/packages/paperclip-runner/package.json b/packages/paperclip-runner/package.json index 51686f9fc6..cdbbb6ad6a 100644 --- a/packages/paperclip-runner/package.json +++ b/packages/paperclip-runner/package.json @@ -28,6 +28,7 @@ "test": "pnpm run test:typescript && pnpm run test:rust", "test:typescript": "node --test test/protocol-contract.test.mjs && vitest run", "test:rust": "cargo test --manifest-path runner/Cargo.toml --locked --workspace", + "test:codex": "cargo test --manifest-path runner/Cargo.toml --locked -p paperclip-runner-core --test codex_provider", "test:durable": "cargo test --manifest-path runner/Cargo.toml --locked -p paperclip-runner-core durable::", "generate:protocol-manifest": "node scripts/generate-protocol-manifest.mjs", "check:protocol-manifest": "node scripts/generate-protocol-manifest.mjs --check", diff --git a/packages/paperclip-runner/runner/CODEX_PROVIDER.md b/packages/paperclip-runner/runner/CODEX_PROVIDER.md new file mode 100644 index 0000000000..0c9ca47a81 --- /dev/null +++ b/packages/paperclip-runner/runner/CODEX_PROVIDER.md @@ -0,0 +1,50 @@ +# Codex provider boundary + +`paperclip-runnerd` supports one provider in this layer: Codex app-server as a +local supervised process. The Paperclip server does not select or launch this +path yet. + +## Command lifecycle + +- `run.prepare` accepts a `provider` object containing `provider: "codex"`, + `driver: "codex_app_server"`, `providerVersion`, `command`, bounded `args`, an + existing absolute `cwd`, optional `model`, `instructions`, and + `approvalPolicy: "never"`. +- `session.open` initializes Codex and starts a thread. A recovered runner + resumes the recorded thread and reads it before accepting another turn. +- `turn.start` requires bounded non-empty `payload.text`. `turn.steer`, + `turn.interrupt`, `turn.stop`, and `run.cancel` use Codex's native turn IDs, + while PRP continues to use its own stable run and turn identities. +- `request.resolve` translates a validated `paperclip.question_response.v1` + response back to Codex's user-input response shape. +- `session.close` and `session.destroy` explicitly terminate the provider + process group. Runner suspend and shutdown also stop the process without + deleting the resumable thread identity. + +## Recovery and duplicate safety + +The provider descriptor, Codex thread ID, account session ID, active Codex turn +ID, and unacknowledged normalized-event prefix are written to a private, +bounded, atomically replaced sidecar in the runner state directory. On restart, +runnerd resumes that exact thread and reconciles the active turn from +`thread/read`. An unexpected provider exit retains the last active turn until +that reconciliation proves whether it is still running, so cancellation and a +later turn cannot diverge from Codex's native state. + +PRP journals every command before the provider effect. Exact command replay +returns the durable result without invoking Codex again. A crash in the effect +window remains indeterminate and is not retried. Codex JSON-RPC notifications +received before a synchronous response are buffered rather than lost. Reusing +a pending structured-input request ID with different content fails closed. +Normalized events remain in the provider sidecar until the durable PRP outbox +has committed and acknowledged them. + +## Normalization and authorization + +Codex-native envelopes do not cross PRP. The provider backend emits bounded, +redacted session, turn, item, plan, usage, tool-execution, notice, and structured +input events. Unknown notifications are ignored. + +Codex starts with an empty dynamic-tool inventory. Catalog presence is not +authorization, and semantic operations remain unavailable until the separate +catalog and run-scoped authorization layers land. diff --git a/packages/paperclip-runner/runner/DURABLE_TRANSPORT.md b/packages/paperclip-runner/runner/DURABLE_TRANSPORT.md index aae730d796..fbc6d9550a 100644 --- a/packages/paperclip-runner/runner/DURABLE_TRANSPORT.md +++ b/packages/paperclip-runner/runner/DURABLE_TRANSPORT.md @@ -1,8 +1,8 @@ # Durable PRP transport This layer gives `paperclip-runnerd` a provider-neutral, package-local PRP v1 -transport. Nothing in the Paperclip server invokes the durable mode yet, and no -provider is installed by this change. +transport. Nothing in the Paperclip server invokes the durable mode yet. Codex +is the only installed provider; other providers remain unavailable. ## Trust boundary @@ -29,6 +29,17 @@ source sequence the runner has produced; acknowledged prefixes are removed atomically from durable state. After disconnect, every remaining event is sent again with the same identity and source sequence. +Executors retain polled events until runnerd acknowledges each event after its +outbox commit. Batches commit one event at a time, so a later oversized event or +capacity failure cannot roll back the accepted prefix or discard the +unacknowledged suffix. Each retained executor event has a stable identity that +runnerd derives into its PRP `sourceEventId`. If the process stops after the +outbox commit but before the executor acknowledgement, a bounded durable +receipt journal recognizes and byte-validates the retained copy without +appending a second event. Receipts outlive transport ACK removal; because the +provider queue is ordered and bounded, a possibly retained front event cannot +be evicted while later events advance the journal. + Commands require a contiguous controller sequence. The runner journals a pending command before invoking its executor and persists its result afterward. An exact duplicate returns the stored result without repeating the effect. If @@ -49,7 +60,8 @@ P0 reserve is an explicit unrecoverable condition. ## Current boundary Durable mode is selected only when `paperclip-runnerd` receives -`--connect-url`. Its transport-only executor handles runner lifecycle commands -and rejects provider commands with `provider_not_installed`. The existing local -fake-runner mode remains unchanged. Codex execution, semantic tools, server -coordination, and the user-facing adapter belong to later layers. +`--connect-url`. Its executor accepts a Codex app-server descriptor through +`run.prepare`, owns the provider process group, resumes the persisted Codex +thread after runner restart, and translates provider notifications to PRP +events. The existing local fake-runner mode remains unchanged. Semantic tools, +server coordination, and the user-facing adapter belong to later layers. diff --git a/packages/paperclip-runner/runner/crates/runner-core/Cargo.toml b/packages/paperclip-runner/runner/crates/runner-core/Cargo.toml index bf26046eb3..6754c383b2 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/Cargo.toml +++ b/packages/paperclip-runner/runner/crates/runner-core/Cargo.toml @@ -25,3 +25,7 @@ path = "src/bin/paperclip-runnerd.rs" [[bin]] name = "fake-harness" path = "src/bin/fake-harness.rs" + +[[bin]] +name = "fake-codex-app-server" +path = "src/bin/fake-codex-app-server.rs" diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/bin/fake-codex-app-server.rs b/packages/paperclip-runner/runner/crates/runner-core/src/bin/fake-codex-app-server.rs new file mode 100644 index 0000000000..2284ea25a2 --- /dev/null +++ b/packages/paperclip-runner/runner/crates/runner-core/src/bin/fake-codex-app-server.rs @@ -0,0 +1,204 @@ +use std::fs::{self, OpenOptions}; +use std::io::{self, BufRead, Write}; +use std::path::{Path, PathBuf}; +use std::process::ExitCode; + +use serde::{Deserialize, Serialize}; +use serde_json::{json, Value}; + +#[derive(Clone, Debug, Default, Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +struct FakeState { + thread_id: String, + active_turn_id: Option, +} + +fn argument(args: &[String], name: &str) -> Option { + args.iter() + .position(|value| value == name) + .and_then(|index| args.get(index + 1)) + .cloned() +} + +fn send(value: Value) -> io::Result<()> { + let mut stdout = io::stdout().lock(); + serde_json::to_writer(&mut stdout, &value)?; + stdout.write_all(b"\n")?; + stdout.flush() +} + +fn load_state(path: &Path) -> FakeState { + fs::read(path) + .ok() + .and_then(|bytes| serde_json::from_slice(&bytes).ok()) + .unwrap_or_else(|| FakeState { + thread_id: "codex-thread-1".to_owned(), + active_turn_id: None, + }) +} + +fn save_state(path: &Path, state: &FakeState) -> io::Result<()> { + fs::write(path, serde_json::to_vec_pretty(state)?) +} + +fn log_call(path: Option<&Path>, method: &str) -> io::Result<()> { + let Some(path) = path else { return Ok(()) }; + let mut file = OpenOptions::new().create(true).append(true).open(path)?; + writeln!(file, "{method}") +} + +fn finish_turn(state_path: &Path, state: &mut FakeState, status: &str) -> io::Result<()> { + let turn_id = state + .active_turn_id + .clone() + .unwrap_or_else(|| "provider-turn-1".to_owned()); + send(json!({ + "method": "item/completed", + "params": {"item": { + "id": "message-1", + "type": "agentMessage", + "status": "completed", + "text": "Codex completed the fake turn." + }} + }))?; + send(json!({ + "method": "thread/tokenUsage/updated", + "params": { + "threadId": state.thread_id, + "tokenUsage": { + "total": {"inputTokens": 12, "outputTokens": 3}, + "last": {"inputTokens": 12, "outputTokens": 3, "requests": 1} + } + } + }))?; + send(json!({ + "method": "turn/completed", + "params": {"turn": {"id": turn_id, "status": status}} + }))?; + state.active_turn_id = None; + save_state(state_path, state) +} + +fn run() -> Result<(), Box> { + let args = std::env::args().skip(1).collect::>(); + let state_path = + PathBuf::from(argument(&args, "--state-file").ok_or("--state-file is required")?); + let call_log = argument(&args, "--call-log").map(PathBuf::from); + let emit_question = args.iter().any(|value| value == "--emit-question"); + let hold_turn = args.iter().any(|value| value == "--hold-turn"); + let exit_after_turn_start = args.iter().any(|value| value == "--exit-after-turn-start"); + let pre_response_notification = args + .iter() + .any(|value| value == "--notification-before-response"); + let mut state = load_state(&state_path); + + for line in io::stdin().lock().lines() { + let message: Value = serde_json::from_str(&line?)?; + if message.get("method").is_none() && message.get("id") == Some(&json!("runtime-request-1")) + { + finish_turn(&state_path, &mut state, "completed")?; + continue; + } + let Some(method) = message.get("method").and_then(Value::as_str) else { + continue; + }; + log_call(call_log.as_deref(), method)?; + let id = message.get("id").cloned(); + match method { + "initialize" => send(json!({ + "id": id, + "result": {"user": {"sessionId": "codex-account-session"}} + }))?, + "initialized" => {} + "thread/start" => { + state.thread_id = "codex-thread-1".to_owned(); + state.active_turn_id = None; + save_state(&state_path, &state)?; + if pre_response_notification { + send(json!({ + "method": "warning", + "params": {"message": "buffered before thread response"} + }))?; + } + send(json!({ + "id": id, + "result": {"thread": {"id": state.thread_id, "sessionId": "codex-account-session"}} + }))?; + } + "thread/resume" => send(json!({ + "id": id, + "result": {"thread": {"id": state.thread_id, "sessionId": "codex-account-session"}} + }))?, + "thread/read" => { + let turns = state + .active_turn_id + .as_ref() + .map(|turn_id| vec![json!({"id": turn_id, "status": "inProgress"})]) + .unwrap_or_default(); + send(json!({ + "id": id, + "result": {"thread": {"id": state.thread_id, "turns": turns}} + }))?; + } + "turn/start" => { + state.active_turn_id = Some("provider-turn-1".to_owned()); + save_state(&state_path, &state)?; + send(json!({ + "id": id, + "result": {"turn": {"id": "provider-turn-1", "status": "inProgress"}} + }))?; + send(json!({ + "method": "turn/started", + "params": {"turn": {"id": "provider-turn-1"}} + }))?; + if exit_after_turn_start { + return Ok(()); + } else if emit_question { + send(json!({ + "id": "runtime-request-1", + "method": "item/tool/requestUserInput", + "params": { + "threadId": state.thread_id, + "turnId": "provider-turn-1", + "itemId": "question-item-1", + "isBlocking": true, + "title": "Deployment input", + "questions": [{ + "id": "environment", + "header": "Environment", + "question": "Where should we deploy?", + "options": [ + {"label": "Staging", "description": "Deploy safely."}, + {"label": "Production", "description": "Deploy directly."} + ] + }] + } + }))?; + } else if !hold_turn { + finish_turn(&state_path, &mut state, "completed")?; + } + } + "turn/steer" => send(json!({"id": id, "result": {"accepted": true}}))?, + "turn/interrupt" => { + send(json!({"id": id, "result": {"accepted": true}}))?; + finish_turn(&state_path, &mut state, "interrupted")?; + } + _ if id.is_some() => send(json!({ + "id": id, + "error": {"code": -32601, "message": format!("unsupported fake method {method}")} + }))?, + _ => {} + } + } + Ok(()) +} + +fn main() -> ExitCode { + match run() { + Ok(()) => ExitCode::SUCCESS, + Err(error) => { + eprintln!("fake-codex-app-server: {error}"); + ExitCode::FAILURE + } + } +} diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/bin/paperclip-runnerd.rs b/packages/paperclip-runner/runner/crates/runner-core/src/bin/paperclip-runnerd.rs index 69c1ea0b50..4764ed1e88 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/bin/paperclip-runnerd.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/bin/paperclip-runnerd.rs @@ -3,11 +3,10 @@ use std::process::ExitCode; use std::time::Duration; use paperclip_runner_core::durable::{ - capture_bootstrap_ticket, run_durable_runner, Command, CommandExecution, CommandExecutor, - DurableRunnerConfig, DurableRunnerError, + capture_bootstrap_ticket, run_durable_runner, DurableRunnerConfig, }; use paperclip_runner_core::local_runner::{run_local_runner, LocalRunnerError, RunnerConfig}; -use serde_json::json; +use paperclip_runner_core::provider_backend::CodexCommandExecutor; fn value(args: &[String], name: &str) -> Result { let index = args @@ -39,27 +38,6 @@ fn usize_value(args: &[String], name: &str, default: usize) -> Result Result { - Ok(CommandExecution::result( - if matches!( - command.command_type.as_str(), - "runner.shutdown" | "runner.suspend" | "runner.drain" - ) { - json!({"status": "completed"}) - } else { - json!({ - "status": "rejected", - "code": "provider_not_installed", - "message": "the durable transport is active, but no provider is installed in this build", - }) - }, - )) - } -} - fn run_durable(args: &[String]) -> Result<(), LocalRunnerError> { let ticket = capture_bootstrap_ticket() .map_err(|error| LocalRunnerError::invalid(error.to_string()))? @@ -71,10 +49,11 @@ fn run_durable(args: &[String]) -> Result<(), LocalRunnerError> { let duration = |name: &str, default: u64| { optional_u64(args, name).map(|value| Duration::from_millis(value.unwrap_or(default))) }; + let state_dir = PathBuf::from(value(args, "--state-dir")?); run_durable_runner( DurableRunnerConfig { connect_url: value(args, "--connect-url")?, - state_dir: PathBuf::from(value(args, "--state-dir")?), + state_dir: state_dir.clone(), runner_instance_id: value(args, "--runner-id")?, environment_lease_id: value(args, "--environment-lease-id")?, run_id: value(args, "--run-id")?, @@ -90,7 +69,7 @@ fn run_durable(args: &[String]) -> Result<(), LocalRunnerError> { max_runtime: duration("--max-runtime-ms", 60 * 60 * 1000)?, }, ticket, - TransportOnlyExecutor, + CodexCommandExecutor::new(state_dir), ) .map_err(|error| LocalRunnerError::invalid(error.to_string())) } diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/codex_provider.rs b/packages/paperclip-runner/runner/crates/runner-core/src/codex_provider.rs new file mode 100644 index 0000000000..6195562247 --- /dev/null +++ b/packages/paperclip-runner/runner/crates/runner-core/src/codex_provider.rs @@ -0,0 +1,847 @@ +use std::collections::{BTreeMap, VecDeque}; +use std::path::{Path, PathBuf}; +use std::time::Duration; + +use serde::{Deserialize, Serialize}; +use serde_json::{json, Value}; + +use crate::durable::redact_text; +use crate::local_runner::LocalRunnerError; +use crate::process_supervisor::SupervisedProcess; + +pub const CODEX_APP_SERVER_MAX_FRAME_BYTES: usize = 4 * 1024 * 1024; +const MAX_BUFFERED_MESSAGES: usize = 1_024; +const MAX_INSTRUCTIONS_BYTES: usize = 1024 * 1024; +type QuestionOptionLabels = BTreeMap>; +type QuestionSetMapping = (String, Value, QuestionOptionLabels); + +fn default_approval_policy() -> String { + "never".to_owned() +} + +#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct CodexProviderConfig { + pub provider: String, + pub driver: String, + pub provider_version: String, + pub command: PathBuf, + #[serde(default)] + pub args: Vec, + pub cwd: String, + #[serde(default)] + pub model: Option, + #[serde(default)] + pub instructions: String, + #[serde(default = "default_approval_policy")] + pub approval_policy: String, +} + +impl CodexProviderConfig { + pub fn validate(&self) -> Result<(), LocalRunnerError> { + if self.provider != "codex" || self.driver != "codex_app_server" { + return Err(LocalRunnerError::invalid( + "the initial runner provider must be codex through codex_app_server", + )); + } + if self.provider_version.trim().is_empty() || self.provider_version.len() > 120 { + return Err(LocalRunnerError::invalid( + "Codex providerVersion is empty or oversized", + )); + } + if self.command.as_os_str().is_empty() { + return Err(LocalRunnerError::invalid("Codex command is required")); + } + let cwd = Path::new(&self.cwd); + if !cwd.is_absolute() || !cwd.is_dir() { + return Err(LocalRunnerError::invalid( + "Codex cwd must be an existing absolute directory", + )); + } + if self.args.len() > 64 + || self.args.iter().any(|argument| { + argument.len() > 4096 || argument.chars().any(|character| character == '\0') + }) + { + return Err(LocalRunnerError::invalid( + "Codex arguments exceed the bounded launch contract", + )); + } + if self + .model + .as_ref() + .is_some_and(|model| model.is_empty() || model.len() > 240) + { + return Err(LocalRunnerError::invalid("Codex model is invalid")); + } + if self.instructions.len() > MAX_INSTRUCTIONS_BYTES { + return Err(LocalRunnerError::invalid( + "Codex instructions exceed the 1 MiB limit", + )); + } + if self.approval_policy != "never" { + return Err(LocalRunnerError::invalid( + "the initial Codex runner requires approvalPolicy=never; governed actions use PRP", + )); + } + Ok(()) + } +} + +#[derive(Clone, Debug, PartialEq)] +pub enum CodexProviderEvent { + Notification { + method: String, + params: Value, + }, + RuntimeRequest { + request_id: String, + question_set: Value, + }, + Exited { + exit_code: Option, + success: bool, + }, +} + +#[derive(Clone, Debug, PartialEq)] +struct PendingRuntimeRequest { + rpc_id: Value, + method: String, + params: Value, + question_set: Value, + option_labels: QuestionOptionLabels, +} + +pub struct CodexProvider { + process: SupervisedProcess, + next_request_id: u64, + thread_id: String, + provider_session_id: Option, + active_provider_turn_id: Option, + pending_messages: VecDeque, + pending_runtime_requests: BTreeMap, + expected_shutdown: bool, +} + +impl CodexProvider { + pub fn start( + config: &CodexProviderConfig, + resume_thread_id: Option<&str>, + ) -> Result { + config.validate()?; + let mut provider = Self { + process: SupervisedProcess::spawn( + &config.command, + &config.args, + Duration::from_secs(2), + CODEX_APP_SERVER_MAX_FRAME_BYTES, + )?, + next_request_id: 1, + thread_id: String::new(), + provider_session_id: None, + active_provider_turn_id: None, + pending_messages: VecDeque::new(), + pending_runtime_requests: BTreeMap::new(), + expected_shutdown: false, + }; + let initialized = provider.request( + "initialize", + json!({ + "clientInfo": { + "name": "paperclip-runnerd", + "title": "Paperclip Runner", + "version": "1", + }, + "capabilities": { + "experimentalApi": true, + "requestAttestation": false, + }, + }), + )?; + provider.process.send(&json!({"method": "initialized"}))?; + + let mut params = json!({ + "cwd": config.cwd, + "model": config.model, + "approvalPolicy": config.approval_policy, + "permissions": "paperclip-runner-workspace-only", + "runtimeWorkspaceRoots": [config.cwd], + "baseInstructions": config.instructions, + }); + let params_object = params + .as_object_mut() + .expect("Codex thread parameters are an object"); + let method = if let Some(thread_id) = resume_thread_id { + params_object.insert("threadId".to_owned(), json!(thread_id)); + "thread/resume" + } else { + // This PR does not grant any semantic tools. A later catalog and + // authorization layer can project a run-scoped inventory here. + params_object.insert("dynamicTools".to_owned(), json!([])); + params_object.insert("experimentalRawEvents".to_owned(), json!(false)); + "thread/start" + }; + let opened = provider.request(method, params)?; + provider.thread_id = opened + .pointer("/thread/id") + .or_else(|| opened.get("threadId")) + .and_then(Value::as_str) + .filter(|value| !value.is_empty()) + .ok_or_else(|| LocalRunnerError::invalid(format!("Codex {method} omitted thread.id")))? + .to_owned(); + if resume_thread_id.is_some_and(|expected| expected != provider.thread_id) { + return Err(LocalRunnerError::invalid( + "Codex resumed a different provider thread", + )); + } + provider.provider_session_id = opened + .pointer("/thread/sessionId") + .or_else(|| initialized.pointer("/user/sessionId")) + .and_then(Value::as_str) + .filter(|value| !value.is_empty()) + .map(str::to_owned); + + if resume_thread_id.is_some() { + let snapshot = provider.read_thread()?; + provider.active_provider_turn_id = latest_active_turn_id(&snapshot); + } + Ok(provider) + } + + pub fn process_id(&self) -> u32 { + self.process.id() + } + + pub fn thread_id(&self) -> &str { + &self.thread_id + } + + pub fn provider_session_id(&self) -> Option<&str> { + self.provider_session_id.as_deref() + } + + pub fn active_provider_turn_id(&self) -> Option<&str> { + self.active_provider_turn_id.as_deref() + } + + pub fn start_turn(&mut self, message: &str, cwd: &str) -> Result { + if self.active_provider_turn_id.is_some() { + return Err(LocalRunnerError::invalid( + "Codex already has an active provider turn", + )); + } + if message.is_empty() || message.len() > MAX_INSTRUCTIONS_BYTES { + return Err(LocalRunnerError::invalid( + "Codex turn text is empty or exceeds the 1 MiB limit", + )); + } + let result = self.request( + "turn/start", + json!({ + "threadId": self.thread_id, + "cwd": cwd, + "permissions": "paperclip-runner-workspace-only", + "runtimeWorkspaceRoots": [cwd], + "input": [{"type": "text", "text": message, "text_elements": []}], + }), + )?; + let provider_turn_id = result + .pointer("/turn/id") + .or_else(|| result.get("turnId")) + .and_then(Value::as_str) + .filter(|value| !value.is_empty()) + .ok_or_else(|| LocalRunnerError::invalid("Codex turn/start omitted turn.id"))? + .to_owned(); + self.active_provider_turn_id = Some(provider_turn_id); + Ok(result) + } + + pub fn steer_turn(&mut self, message: &str) -> Result { + let turn_id = self + .active_provider_turn_id + .clone() + .ok_or_else(|| LocalRunnerError::invalid("Codex has no active provider turn"))?; + if message.is_empty() || message.len() > MAX_INSTRUCTIONS_BYTES { + return Err(LocalRunnerError::invalid( + "Codex steering text is empty or oversized", + )); + } + self.request( + "turn/steer", + json!({ + "threadId": self.thread_id, + "expectedTurnId": turn_id, + "input": [{"type": "text", "text": message, "text_elements": []}], + }), + ) + } + + pub fn interrupt_turn(&mut self) -> Result { + let turn_id = self + .active_provider_turn_id + .clone() + .ok_or_else(|| LocalRunnerError::invalid("Codex has no active provider turn"))?; + self.cancel_pending_runtime_requests()?; + self.request( + "turn/interrupt", + json!({"threadId": self.thread_id, "turnId": turn_id}), + ) + } + + pub fn read_thread(&mut self) -> Result { + self.request( + "thread/read", + json!({"threadId": self.thread_id, "includeTurns": true}), + ) + } + + pub fn resolve_runtime_request( + &mut self, + request_id: &str, + response: &Value, + ) -> Result<(), LocalRunnerError> { + let pending = self + .pending_runtime_requests + .get(request_id) + .cloned() + .ok_or_else(|| { + LocalRunnerError::invalid("runtime response has no pending Codex request") + })?; + let result = codex_question_response(&pending, response)?; + self.process + .send(&json!({"id": pending.rpc_id, "result": result}))?; + self.pending_runtime_requests.remove(request_id); + Ok(()) + } + + pub fn poll(&mut self) -> Result, LocalRunnerError> { + let message = if let Some(message) = self.pending_messages.pop_front() { + message + } else { + let Some(line) = self.process.receive_stdout_line(Duration::from_millis(1))? else { + return if let Some(exit) = self.process.try_wait()? { + Ok(Some(CodexProviderEvent::Exited { + exit_code: exit.exit_code, + success: exit.success && self.expected_shutdown, + })) + } else { + Ok(None) + }; + }; + parse_provider_message(&line)? + }; + + if let (Some(rpc_id), Some(method)) = ( + message.get("id").cloned(), + message.get("method").and_then(Value::as_str), + ) { + if method == "item/tool/requestUserInput" { + let params = message.get("params").cloned().unwrap_or(Value::Null); + if params.get("threadId").and_then(Value::as_str) != Some(self.thread_id.as_str()) { + return Err(LocalRunnerError::invalid( + "Codex runtime request named another thread", + )); + } + let active_turn_id = self.active_provider_turn_id.as_deref().ok_or_else(|| { + LocalRunnerError::invalid( + "Codex runtime request arrived outside an active turn", + ) + })?; + if params.get("turnId").and_then(Value::as_str) != Some(active_turn_id) { + return Err(LocalRunnerError::invalid( + "Codex runtime request named another turn", + )); + } + let (request_id, question_set, option_labels) = + codex_question_set(&rpc_id, ¶ms)?; + let pending = PendingRuntimeRequest { + rpc_id, + method: method.to_owned(), + params, + question_set: question_set.clone(), + option_labels, + }; + if let Some(existing) = self.pending_runtime_requests.get(&request_id) { + if existing != &pending { + return Err(LocalRunnerError::invalid( + "Codex reused a runtime request id with different input", + )); + } + return Ok(None); + } else { + self.pending_runtime_requests + .insert(request_id.clone(), pending); + } + return Ok(Some(CodexProviderEvent::RuntimeRequest { + request_id, + question_set, + })); + } + self.process.send(&json!({ + "id": rpc_id, + "error": {"code": -32601, "message": "provider request is unavailable in this runner layer"}, + }))?; + return Ok(Some(CodexProviderEvent::Notification { + method: "warning".to_owned(), + params: json!({"message": format!("unsupported Codex request {}", bounded_method(method))}), + })); + } + + if let Some(method) = message.get("method").and_then(Value::as_str) { + let params = message.get("params").cloned().unwrap_or(Value::Null); + validate_notification_binding( + &self.thread_id, + self.active_provider_turn_id.as_deref(), + ¶ms, + )?; + if method == "turn/completed" { + self.active_provider_turn_id = None; + } + return Ok(Some(CodexProviderEvent::Notification { + method: method.to_owned(), + params, + })); + } + Ok(None) + } + + pub fn shutdown(&mut self) -> Result<(), LocalRunnerError> { + self.expected_shutdown = true; + self.cancel_pending_runtime_requests()?; + self.process.terminate_group().map(|_| ()) + } + + fn cancel_pending_runtime_requests(&mut self) -> Result<(), LocalRunnerError> { + let pending = std::mem::take(&mut self.pending_runtime_requests); + for request in pending.into_values() { + self.process.send(&json!({ + "id": request.rpc_id, + "result": {"answers": {}}, + }))?; + } + Ok(()) + } + + fn request(&mut self, method: &str, params: Value) -> Result { + let request_id = self.next_request_id; + self.next_request_id = self + .next_request_id + .checked_add(1) + .ok_or_else(|| LocalRunnerError::invalid("Codex request id exhausted"))?; + self.process + .send(&json!({"id": request_id, "method": method, "params": params}))?; + loop { + let line = self + .process + .receive_stdout_line(Duration::from_secs(30))? + .ok_or_else(|| { + LocalRunnerError::invalid(format!("Codex {method} response timed out")) + })?; + let message = parse_provider_message(&line)?; + if message.get("id").and_then(Value::as_u64) == Some(request_id) + && message.get("method").is_none() + { + if let Some(error) = message.get("error") { + return Err(LocalRunnerError::invalid(format!( + "Codex {method} failed: {}", + redact_text(&error.to_string()) + ))); + } + return Ok(message.get("result").cloned().unwrap_or(Value::Null)); + } + if self.pending_messages.len() >= MAX_BUFFERED_MESSAGES { + return Err(LocalRunnerError::invalid( + "Codex emitted too many messages before a request response", + )); + } + self.pending_messages.push_back(message); + } + } +} + +fn parse_provider_message(line: &str) -> Result { + let value: Value = serde_json::from_str(line).map_err(|error| { + LocalRunnerError::invalid(format!("Codex emitted invalid JSON-RPC: {error}")) + })?; + if !value.is_object() { + return Err(LocalRunnerError::invalid( + "Codex emitted a non-object JSON-RPC frame", + )); + } + Ok(value) +} + +fn validate_notification_binding( + thread_id: &str, + active_turn_id: Option<&str>, + params: &Value, +) -> Result<(), LocalRunnerError> { + let notification_thread_id = params + .get("threadId") + .or_else(|| params.pointer("/thread/id")) + .and_then(Value::as_str) + .filter(|value| !value.is_empty()); + if notification_thread_id.is_some_and(|value| value != thread_id) { + return Err(LocalRunnerError::invalid( + "Codex notification named another thread", + )); + } + let notification_turn_id = params + .get("turnId") + .or_else(|| params.pointer("/turn/id")) + .and_then(Value::as_str) + .filter(|value| !value.is_empty()); + if let Some(active_turn_id) = active_turn_id { + if notification_turn_id.is_some_and(|value| value != active_turn_id) { + return Err(LocalRunnerError::invalid( + "Codex notification named another active turn", + )); + } + } + Ok(()) +} + +fn latest_active_turn_id(snapshot: &Value) -> Option { + snapshot + .pointer("/thread/turns") + .and_then(Value::as_array) + .into_iter() + .flatten() + .rev() + .find(|turn| { + matches!( + turn.get("status").and_then(Value::as_str), + Some("inProgress" | "running" | "pending") + ) + }) + .and_then(|turn| turn.get("id").and_then(Value::as_str)) + .filter(|value| !value.is_empty()) + .map(str::to_owned) +} + +fn bounded_method(method: &str) -> String { + method + .chars() + .filter(|character| character.is_ascii_alphanumeric() || "._/-".contains(*character)) + .take(160) + .collect() +} + +fn codex_question_set( + rpc_id: &Value, + params: &Value, +) -> Result { + let request_id = match rpc_id { + Value::String(value) if !value.is_empty() && value.len() <= 160 => value.clone(), + Value::Number(value) => value.to_string(), + _ => { + return Err(LocalRunnerError::invalid( + "Codex user-input request id is invalid", + )) + } + }; + let questions = params + .get("questions") + .and_then(Value::as_array) + .ok_or_else(|| LocalRunnerError::invalid("Codex user-input request omitted questions"))?; + if questions.is_empty() || questions.len() > 3 { + return Err(LocalRunnerError::invalid( + "Codex user-input request must contain one to three questions", + )); + } + let mut canonical = Vec::new(); + let mut option_labels = BTreeMap::new(); + for question in questions { + if question.get("isSecret").and_then(Value::as_bool) == Some(true) { + return Err(LocalRunnerError::invalid( + "Codex secret input cannot use the persisted question channel", + )); + } + let id = question + .get("id") + .and_then(Value::as_str) + .filter(|value| !value.is_empty() && value.len() <= 160) + .ok_or_else(|| LocalRunnerError::invalid("Codex question id is invalid"))?; + if option_labels.contains_key(id) { + return Err(LocalRunnerError::invalid( + "Codex question ids must be unique", + )); + } + let prompt = question + .get("question") + .and_then(Value::as_str) + .filter(|value| !value.is_empty()) + .ok_or_else(|| LocalRunnerError::invalid("Codex question prompt is required"))?; + let options = question + .get("options") + .and_then(Value::as_array) + .cloned() + .unwrap_or_default(); + let mut labels = BTreeMap::new(); + let canonical_options = options + .iter() + .take(64) + .enumerate() + .filter_map(|(index, option)| { + let label = option.get("label")?.as_str()?.trim(); + if label.is_empty() { + return None; + } + let option_id = format!("option-{}", index + 1); + labels.insert(option_id.clone(), label.chars().take(240).collect()); + Some(json!({ + "id": option_id, + "label": label.chars().take(240).collect::(), + "description": option.get("description").and_then(Value::as_str).map(|value| value.chars().take(1000).collect::()), + })) + }) + .collect::>(); + if !options.is_empty() && canonical_options.len() != options.len() { + return Err(LocalRunnerError::invalid( + "Codex question contains an invalid option", + )); + } + option_labels.insert(id.to_owned(), labels); + let mut canonical_question = json!({ + "id": id, + "header": question.get("header").and_then(Value::as_str).unwrap_or("Question").chars().take(80).collect::(), + "prompt": prompt.chars().take(4000).collect::(), + "required": true, + "answerMode": if canonical_options.is_empty() { "text" } else { "single_select" }, + "options": canonical_options, + }); + if question.get("isOther").and_then(Value::as_bool) == Some(true) { + canonical_question + .as_object_mut() + .expect("canonical question is an object") + .insert( + "customAnswer".to_owned(), + json!({ + "enabled": true, + "label": "Other", + "placeholder": "Enter another answer", + }), + ); + } + canonical.push(canonical_question); + } + Ok(( + request_id, + json!({ + "schema": "paperclip.question_set.v1", + "title": params.get("title").and_then(Value::as_str).unwrap_or("Codex input").chars().take(240).collect::(), + "submitLabel": "Submit answers", + "questions": canonical, + }), + option_labels, + )) +} + +fn codex_question_response( + pending: &PendingRuntimeRequest, + response: &Value, +) -> Result { + let response_object = response + .as_object() + .ok_or_else(|| LocalRunnerError::invalid("runtime response must be an object"))?; + if response_object + .keys() + .any(|key| !matches!(key.as_str(), "schema" | "answers")) + { + return Err(LocalRunnerError::invalid( + "runtime response contains an unknown top-level field", + )); + } + if response.get("schema").and_then(Value::as_str) != Some("paperclip.question_response.v1") { + return Err(LocalRunnerError::invalid( + "runtime response requires paperclip.question_response.v1", + )); + } + let answers = response + .get("answers") + .and_then(Value::as_object) + .ok_or_else(|| LocalRunnerError::invalid("runtime response answers are required"))?; + let questions = pending + .question_set + .get("questions") + .and_then(Value::as_array) + .ok_or_else(|| LocalRunnerError::invalid("pending question set is malformed"))?; + let mut native = serde_json::Map::new(); + for question in questions { + let id = question + .get("id") + .and_then(Value::as_str) + .ok_or_else(|| LocalRunnerError::invalid("pending question id is malformed"))?; + let answer = answers + .get(id) + .ok_or_else(|| LocalRunnerError::invalid(format!("missing answer for {id}")))?; + let answer = answer.as_object().ok_or_else(|| { + LocalRunnerError::invalid(format!("answer for {id} must be an object")) + })?; + if answer + .keys() + .any(|key| !matches!(key.as_str(), "selectedOptionIds" | "text" | "customText")) + { + return Err(LocalRunnerError::invalid(format!( + "answer for {id} contains an unknown field" + ))); + } + let selected = answer.get("selectedOptionIds"); + let text = answer.get("text"); + let custom = answer.get("customText"); + let values = if question.get("answerMode").and_then(Value::as_str) == Some("single_select") + { + if text.is_some() || (selected.is_some() && custom.is_some()) { + return Err(LocalRunnerError::invalid(format!( + "{id} must contain one selected option or one custom answer" + ))); + } + if let Some(custom_text) = custom { + if question + .pointer("/customAnswer/enabled") + .and_then(Value::as_bool) + != Some(true) + { + return Err(LocalRunnerError::invalid(format!( + "{id} does not allow a custom answer" + ))); + } + vec![custom_text + .as_str() + .filter(|value| !value.is_empty() && value.len() <= 4000) + .ok_or_else(|| { + LocalRunnerError::invalid(format!( + "{id} custom answer must be non-empty and bounded" + )) + })? + .to_owned()] + } else { + let selected = selected + .and_then(Value::as_array) + .filter(|values| values.len() == 1) + .ok_or_else(|| { + LocalRunnerError::invalid(format!("{id} requires one selected option")) + })?; + let option_id = selected[0] + .as_str() + .ok_or_else(|| LocalRunnerError::invalid("selected option id is invalid"))?; + vec![pending + .option_labels + .get(id) + .and_then(|labels| labels.get(option_id)) + .cloned() + .ok_or_else(|| LocalRunnerError::invalid("selected option is not available"))?] + } + } else { + if selected.is_some() || custom.is_some() { + return Err(LocalRunnerError::invalid(format!( + "{id} text answer cannot contain select or custom fields" + ))); + } + vec![text + .and_then(Value::as_str) + .filter(|value| !value.is_empty() && value.len() <= 4000) + .ok_or_else(|| LocalRunnerError::invalid(format!("{id} requires text")))? + .to_owned()] + }; + native.insert(id.to_owned(), json!({"answers": values})); + } + if answers.keys().any(|id| !native.contains_key(id)) { + return Err(LocalRunnerError::invalid( + "runtime response contains an unknown question id", + )); + } + Ok(json!({"answers": native})) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn converts_codex_questions_and_responses_without_provider_leakage() { + let (request_id, question_set, labels) = codex_question_set( + &json!(41), + &json!({ + "requestId": "request-1", + "questions": [{ + "id": "environment", + "header": "Environment", + "question": "Where should we deploy?", + "options": [{"label": "Staging", "description": "Deploy safely."}], + }], + }), + ) + .unwrap(); + assert_eq!(request_id, "41"); + assert_eq!(question_set["schema"], "paperclip.question_set.v1"); + let pending = PendingRuntimeRequest { + rpc_id: json!(41), + method: "item/tool/requestUserInput".to_owned(), + params: Value::Null, + question_set, + option_labels: labels, + }; + let native = codex_question_response( + &pending, + &json!({ + "schema": "paperclip.question_response.v1", + "answers": {"environment": {"selectedOptionIds": ["option-1"]}}, + }), + ) + .unwrap(); + assert_eq!(native["answers"]["environment"]["answers"][0], "Staging"); + assert!(codex_question_response( + &pending, + &json!({ + "schema": "paperclip.question_response.v1", + "answers": {"environment": { + "selectedOptionIds": ["option-1"], + "customText": "Production", + }}, + }), + ) + .is_err()); + assert!(codex_question_response( + &pending, + &json!({ + "schema": "paperclip.question_response.v1", + "answers": {"environment": {"selectedOptionIds": ["option-1"]}}, + "providerEnvelope": {}, + }), + ) + .is_err()); + } + + #[test] + fn finds_only_active_turns_during_resume() { + let snapshot = json!({"thread": {"turns": [ + {"id": "done", "status": "completed"}, + {"id": "active", "status": "inProgress"} + ]}}); + assert_eq!(latest_active_turn_id(&snapshot).as_deref(), Some("active")); + } + + #[test] + fn rejects_notifications_bound_to_another_thread_or_active_turn() { + assert!(validate_notification_binding( + "thread-1", + Some("turn-1"), + &json!({"threadId": "thread-2", "turnId": "turn-1"}), + ) + .is_err()); + assert!(validate_notification_binding( + "thread-1", + Some("turn-1"), + &json!({"threadId": "thread-1", "turnId": "turn-2"}), + ) + .is_err()); + assert!(validate_notification_binding( + "thread-1", + Some("turn-1"), + &json!({"threadId": "thread-1", "turnId": "turn-1"}), + ) + .is_ok()); + } +} diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/durable/mod.rs b/packages/paperclip-runner/runner/crates/runner-core/src/durable/mod.rs index cbb1e1ebbd..ef3ec61456 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/durable/mod.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/durable/mod.rs @@ -7,7 +7,10 @@ use std::fmt::{self, Display, Formatter}; use std::path::PathBuf; use std::time::Duration; -pub use runner::{run_durable_runner, CommandExecution, CommandExecutor}; +pub use runner::{run_durable_runner, CommandExecution, CommandExecutor, PolledEvent}; +pub(crate) use state::{ + create_private_temporary_file, open_private_regular_file, redact_text, verify_private_directory, +}; pub use state::{ Command, CommandDisposition, DurableState, DurableStateStore, EventPriority, StoredCommandResult, StoredOutboxEvent, diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/durable/runner.rs b/packages/paperclip-runner/runner/crates/runner-core/src/durable/runner.rs index dba5f6b57d..b4b2bc39a9 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/durable/runner.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/durable/runner.rs @@ -1,6 +1,7 @@ use std::thread; use std::time::Instant; +use serde::{Deserialize, Serialize}; use serde_json::{json, Value}; use super::state::{ @@ -19,6 +20,15 @@ pub struct CommandExecution { pub events: Vec<(String, EventPriority, Value)>, } +#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct PolledEvent { + pub executor_event_id: String, + pub event_type: String, + pub priority: EventPriority, + pub payload: Value, +} + impl CommandExecution { pub fn result(result: Value) -> Self { Self { @@ -30,6 +40,21 @@ impl CommandExecution { pub trait CommandExecutor { fn execute(&mut self, command: &Command) -> Result; + + fn poll_events(&mut self) -> Result, DurableRunnerError> { + Ok(Vec::new()) + } + + /// Removes the prefix returned by `poll_events` after each event is + /// durably committed to the PRP outbox. Implementations that retain + /// provider events must not remove them before this acknowledgement. + fn acknowledge_events(&mut self, _count: usize) -> Result<(), DurableRunnerError> { + Ok(()) + } + + fn shutdown(&mut self) -> Result<(), DurableRunnerError> { + Ok(()) + } } pub fn run_durable_runner( @@ -63,6 +88,7 @@ pub fn run_durable_runner( loop { if started.elapsed() >= config.max_runtime { + let _ = executor.shutdown(); state.lifecycle = "recoverable_failure".to_owned(); state.recoverable_failure = Some("transport_reconnect_deadline_exceeded".to_owned()); state.record_diagnostic( @@ -76,6 +102,7 @@ pub fn run_durable_runner( if lease.as_ref().is_some_and(|credential| { current_unix_ms().is_ok_and(|now| now >= credential.expires_at_unix_ms) }) { + let _ = executor.shutdown(); state.lifecycle = "recoverable_failure".to_owned(); state.recoverable_failure = Some("lease_expired_requires_bootstrap".to_owned()); state.record_diagnostic("connection lease expired; a fresh bootstrap is required"); @@ -138,6 +165,7 @@ pub fn run_durable_runner( disconnected = true; } if stop_after_reply && !disconnected { + executor.shutdown()?; state.lifecycle = "stopped".to_owned(); store.save(&state)?; return Ok(()); @@ -154,7 +182,15 @@ pub fn run_durable_runner( if started.elapsed() >= config.max_runtime { break; } + poll_executor_events(&mut state, &store, &config, &mut executor)?; + if let Err(error) = send_outbox(&mut transport, &state, &mut sent_source_seq) { + state.record_diagnostic(error.to_string()); + state.reconnect_count = state.reconnect_count.saturating_add(1); + store.save(&state)?; + break; + } if current_unix_ms()? >= connection.expires_at_unix_ms { + let _ = executor.shutdown(); state.lifecycle = "recoverable_failure".to_owned(); state.recoverable_failure = Some("lease_expired_requires_bootstrap".to_owned()); state.record_diagnostic("active connection lease expired"); @@ -210,6 +246,7 @@ pub fn run_durable_runner( break; } if stop { + executor.shutdown()?; state.lifecycle = "stopped".to_owned(); store.save(&state)?; return Ok(()); @@ -229,6 +266,7 @@ pub fn run_durable_runner( } state.lifecycle = "revoked".to_owned(); state.record_diagnostic("connection capability was revoked"); + executor.shutdown()?; store.save(&state)?; return Ok(()); } @@ -258,6 +296,42 @@ pub fn run_durable_runner( } } +fn poll_executor_events( + state: &mut DurableState, + store: &DurableStateStore, + config: &DurableRunnerConfig, + executor: &mut E, +) -> Result<(), DurableRunnerError> { + let events = executor.poll_events()?; + if events.is_empty() { + return Ok(()); + } + for event in events { + // Commit and acknowledge one event at a time. If a later event is + // oversized or the outbox is full, the accepted prefix is already + // durable and the unacknowledged suffix remains with the executor. + if state.has_executor_event_receipt( + &event.executor_event_id, + &event.event_type, + event.priority, + &event.payload, + )? { + executor.acknowledge_events(1)?; + continue; + } + state.enqueue_executor_event( + config, + event.executor_event_id, + event.event_type, + event.priority, + event.payload, + )?; + store.save(state)?; + executor.acknowledge_events(1)?; + } + Ok(()) +} + fn process_command( state: &mut DurableState, store: &DurableStateStore, @@ -266,7 +340,15 @@ fn process_command( command: &Command, ) -> Result<(StoredCommandResult, bool), DurableRunnerError> { match state.begin_command(command)? { - CommandDisposition::Replay(result) | CommandDisposition::Reject(result) => { + CommandDisposition::Replay(result) => { + let stop = result.status == "completed" + && matches!( + command.command_type.as_str(), + "runner.shutdown" | "runner.suspend" + ); + return Ok((result, stop)); + } + CommandDisposition::Reject(result) => { return Ok((result, false)); } CommandDisposition::Execute => {} @@ -342,6 +424,7 @@ fn control_envelope( #[cfg(test)] mod tests { + use std::collections::VecDeque; use std::fs; use std::path::PathBuf; use std::time::Duration; @@ -352,6 +435,11 @@ mod tests { calls: usize, } + struct RetainingEventExecutor { + events: VecDeque, + fail_acknowledgement: bool, + } + impl CommandExecutor for CountingExecutor { fn execute(&mut self, _command: &Command) -> Result { self.calls += 1; @@ -359,6 +447,31 @@ mod tests { } } + impl CommandExecutor for RetainingEventExecutor { + fn execute(&mut self, _command: &Command) -> Result { + Ok(CommandExecution::result(json!({"status": "completed"}))) + } + + fn poll_events(&mut self) -> Result, DurableRunnerError> { + Ok(self.events.iter().cloned().collect()) + } + + fn acknowledge_events(&mut self, count: usize) -> Result<(), DurableRunnerError> { + if self.fail_acknowledgement { + return Err(DurableRunnerError::invalid( + "simulated crash before provider acknowledgement", + )); + } + if count > self.events.len() { + return Err(DurableRunnerError::invalid( + "test acknowledgement exceeded pending events", + )); + } + self.events.drain(..count); + Ok(()) + } + } + fn config(directory: PathBuf) -> DurableRunnerConfig { DurableRunnerConfig { connect_url: "ws://127.0.0.1:3000/path".to_owned(), @@ -379,12 +492,12 @@ mod tests { } } - fn command() -> Command { + fn command(command_type: &str) -> Command { Command { schema: "paperclip.prp.command.v1".to_owned(), command_id: "command_1".to_owned(), controller_seq: 1, - command_type: "session.open".to_owned(), + command_type: command_type.to_owned(), issued_at: "2026-08-24T00:00:00.000Z".to_owned(), deadline_at: None, precondition: None, @@ -392,6 +505,101 @@ mod tests { } } + #[test] + fn event_batch_keeps_accepted_prefix_and_unacknowledged_suffix() { + let directory = std::env::temp_dir().join(format!( + "paperclip-runner-event-batch-failure-{}", + std::process::id() + )); + let _ = fs::remove_dir_all(&directory); + let mut config = config(directory.clone()); + config.max_frame_bytes = 1024; + let store = DurableStateStore::new(&directory).unwrap(); + let (mut state, _) = store.load_or_create(&config).unwrap(); + let mut executor = RetainingEventExecutor { + events: VecDeque::from([ + PolledEvent { + executor_event_id: "provider-event-1".to_owned(), + event_type: "provider.notice.recorded".to_owned(), + priority: EventPriority::P1, + payload: json!({"message": "durable prefix"}), + }, + PolledEvent { + executor_event_id: "provider-event-2".to_owned(), + event_type: "provider.notice.recorded".to_owned(), + priority: EventPriority::P1, + payload: json!({"message": "x".repeat(2048)}), + }, + ]), + fail_acknowledgement: false, + }; + + let error = poll_executor_events(&mut state, &store, &config, &mut executor) + .expect_err("the oversized suffix must fail closed"); + assert!(error.to_string().contains("transport frame limit")); + assert_eq!(state.outbox.len(), 1); + assert_eq!(state.outbox[0].event_type, "provider.notice.recorded"); + assert_eq!(executor.events.len(), 1); + assert_eq!( + executor.events[0].payload["message"], + Value::String("x".repeat(2048)) + ); + + let (reloaded, recovered) = store.load_or_create(&config).unwrap(); + assert!(recovered); + assert_eq!(reloaded.outbox.len(), 1); + fs::remove_dir_all(directory).unwrap(); + } + + #[test] + fn receipt_survives_outbox_ack_and_prevents_duplicate_delivery() { + let directory = std::env::temp_dir().join(format!( + "paperclip-runner-event-ack-crash-{}", + std::process::id() + )); + let _ = fs::remove_dir_all(&directory); + let config = config(directory.clone()); + let store = DurableStateStore::new(&directory).unwrap(); + let (mut state, _) = store.load_or_create(&config).unwrap(); + let mut executor = RetainingEventExecutor { + events: VecDeque::from([PolledEvent { + executor_event_id: "provider-event-before-ack-crash".to_owned(), + event_type: "provider.notice.recorded".to_owned(), + priority: EventPriority::P1, + payload: json!({"message": "deliver exactly once"}), + }]), + fail_acknowledgement: true, + }; + + let error = poll_executor_events(&mut state, &store, &config, &mut executor) + .expect_err("simulate a crash after outbox persistence"); + assert!(error + .to_string() + .contains("before provider acknowledgement")); + assert_eq!(state.outbox.len(), 1); + assert_eq!(executor.events.len(), 1); + state + .apply_ack(1) + .expect("controller ACK removes the durable outbox copy"); + store.save(&state).unwrap(); + + let (mut recovered_state, recovered) = store.load_or_create(&config).unwrap(); + assert!(recovered); + assert!(recovered_state.outbox.is_empty()); + executor.fail_acknowledgement = false; + executor.events[0].payload = json!({"message": "different data"}); + let mismatch = poll_executor_events(&mut recovered_state, &store, &config, &mut executor) + .expect_err("a retained identity cannot name different event data"); + assert!(mismatch.to_string().contains("reused with different")); + executor.events[0].payload = json!({"message": "deliver exactly once"}); + poll_executor_events(&mut recovered_state, &store, &config, &mut executor) + .expect("recovery acknowledges the retained provider copy"); + assert!(executor.events.is_empty()); + assert!(recovered_state.outbox.is_empty()); + assert_eq!(recovered_state.highest_source_seq(), 1); + fs::remove_dir_all(directory).unwrap(); + } + #[test] fn duplicate_delivery_replays_the_durable_result() { let directory = std::env::temp_dir().join(format!( @@ -403,14 +611,39 @@ mod tests { let store = DurableStateStore::new(&directory).unwrap(); let (mut state, _) = store.load_or_create(&config).unwrap(); let mut executor = CountingExecutor { calls: 0 }; - let first = process_command(&mut state, &store, &config, &mut executor, &command()) + let command = command("session.open"); + let first = process_command(&mut state, &store, &config, &mut executor, &command) .unwrap() .0; - let replay = process_command(&mut state, &store, &config, &mut executor, &command()) + let replay = process_command(&mut state, &store, &config, &mut executor, &command) .unwrap() .0; assert_eq!(executor.calls, 1); assert_eq!(first, replay); fs::remove_dir_all(directory).unwrap(); } + + #[test] + fn completed_shutdown_replay_still_stops_after_delivery() { + let directory = std::env::temp_dir().join(format!( + "paperclip-runner-shutdown-replay-{}", + std::process::id() + )); + let _ = fs::remove_dir_all(&directory); + let config = config(directory.clone()); + let store = DurableStateStore::new(&directory).unwrap(); + let (mut state, _) = store.load_or_create(&config).unwrap(); + let mut executor = CountingExecutor { calls: 0 }; + let command = command("runner.shutdown"); + + let (_, first_stop) = + process_command(&mut state, &store, &config, &mut executor, &command).unwrap(); + let (_, replay_stop) = + process_command(&mut state, &store, &config, &mut executor, &command).unwrap(); + + assert!(first_stop); + assert!(replay_stop); + assert_eq!(executor.calls, 1); + fs::remove_dir_all(directory).unwrap(); + } } diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/durable/state.rs b/packages/paperclip-runner/runner/crates/runner-core/src/durable/state.rs index 9dc4174229..50c6df265c 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/durable/state.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/durable/state.rs @@ -1,4 +1,4 @@ -use std::collections::BTreeMap; +use std::collections::{BTreeMap, HashSet}; use std::fs::{self, File, OpenOptions}; use std::io::{self, Read, Write}; use std::path::{Path, PathBuf}; @@ -18,10 +18,12 @@ const STATE_FILE: &str = "runner-state.json"; const MAX_RECENT_COMMANDS: usize = 128; const MAX_DIAGNOSTICS: usize = 32; const MAX_COMMAND_RESULT_BYTES: usize = 64 * 1024; +const MAX_EXECUTOR_EVENT_RECEIPTS: usize = 256; const STATE_OVERHEAD_BYTES: usize = 16 * 1024 * 1024; const TEMP_FILE_ATTEMPTS: usize = 32; -#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] pub enum EventPriority { P0, P1, @@ -149,6 +151,13 @@ pub struct StoredCommandResult { pub result: Value, } +#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +struct ExecutorEventReceipt { + fingerprint: String, + source_seq: u64, +} + #[derive(Clone, Debug, PartialEq)] pub enum CommandDisposition { Execute, @@ -179,6 +188,8 @@ pub struct DurableState { pub processed_commands: BTreeMap, #[serde(default)] pub processed_command_fingerprints: BTreeMap, + #[serde(default)] + executor_event_receipts: BTreeMap, pub diagnostics: Vec, pub backpressure: bool, pub recoverable_failure: Option, @@ -206,6 +217,7 @@ impl DurableState { outbox: Vec::new(), processed_commands: BTreeMap::new(), processed_command_fingerprints: BTreeMap::new(), + executor_event_receipts: BTreeMap::new(), diagnostics: Vec::new(), backpressure: false, recoverable_failure: None, @@ -226,8 +238,120 @@ impl DurableState { event_type: impl Into, priority: EventPriority, payload: Value, + ) -> Result { + let source_event_id = format!( + "event_{}_{:016}", + self.runner_instance_id, self.next_source_seq + ); + self.enqueue_event_with_source_event_id( + config, + source_event_id, + event_type, + priority, + payload, + ) + } + + fn source_event_id_for_executor( + &self, + executor_event_id: &str, + ) -> Result { + if executor_event_id.is_empty() + || executor_event_id.len() > 160 + || executor_event_id.chars().any(char::is_control) + { + return Err(DurableRunnerError::invalid( + "executor event identity is empty, oversized, or contains control characters", + )); + } + let mut hasher = Sha256::new(); + hasher.update(b"paperclip.executor-event.v1\0"); + hasher.update(self.runner_instance_id.as_bytes()); + hasher.update(b"\0"); + hasher.update(executor_event_id.as_bytes()); + Ok(format!("event_executor_{:x}", hasher.finalize())) + } + + fn has_source_event_id(&self, source_event_id: &str) -> bool { + self.outbox.iter().any(|event| { + event + .envelope + .pointer("/payload/sourceEventId") + .and_then(Value::as_str) + == Some(source_event_id) + }) + } + + pub(crate) fn has_executor_event_receipt( + &self, + executor_event_id: &str, + event_type: &str, + priority: EventPriority, + payload: &Value, + ) -> Result { + self.source_event_id_for_executor(executor_event_id)?; + let Some(existing) = self.executor_event_receipts.get(executor_event_id) else { + return Ok(false); + }; + if existing.fingerprint != executor_event_fingerprint(event_type, priority, payload) { + return Err(DurableRunnerError::invalid( + "executor event identity was reused with different event data", + )); + } + Ok(true) + } + + pub(crate) fn enqueue_executor_event( + &mut self, + config: &DurableRunnerConfig, + executor_event_id: String, + event_type: String, + priority: EventPriority, + payload: Value, + ) -> Result { + if self.has_executor_event_receipt(&executor_event_id, &event_type, priority, &payload)? { + return Err(DurableRunnerError::invalid( + "executor event identity is already committed", + )); + } + let source_event_id = self.source_event_id_for_executor(&executor_event_id)?; + let fingerprint = executor_event_fingerprint(&event_type, priority, &payload); + let source_seq = self.enqueue_event_with_source_event_id( + config, + source_event_id, + event_type, + priority, + payload, + )?; + self.executor_event_receipts.insert( + executor_event_id, + ExecutorEventReceipt { + fingerprint, + source_seq, + }, + ); + self.compact_executor_event_receipts(); + Ok(source_seq) + } + + pub(crate) fn enqueue_event_with_source_event_id( + &mut self, + config: &DurableRunnerConfig, + source_event_id: String, + event_type: impl Into, + priority: EventPriority, + payload: Value, ) -> Result { let event_type = event_type.into(); + if source_event_id.is_empty() + || source_event_id.len() > 160 + || source_event_id.chars().any(char::is_control) + || self.has_source_event_id(&source_event_id) + { + return Err(DurableRunnerError::invalid( + "source event identity is malformed or already queued", + )); + } if event_type.is_empty() || event_type.len() > 160 || event_type.chars().any(char::is_control) @@ -256,7 +380,7 @@ impl DurableState { "itemId": self.item_id, "payload": { "schema": "paperclip.prp.event.v1", - "sourceEventId": format!("event_{}_{source_seq:016}", self.runner_instance_id), + "sourceEventId": source_event_id, "sourceSeq": source_seq, "sourceInstanceId": self.runner_instance_id, "sourceKind": "runner", @@ -484,6 +608,20 @@ impl DurableState { } } } + + fn compact_executor_event_receipts(&mut self) { + while self.executor_event_receipts.len() > MAX_EXECUTOR_EVENT_RECEIPTS { + let Some(oldest_id) = self + .executor_event_receipts + .iter() + .min_by_key(|(_, receipt)| receipt.source_seq) + .map(|(event_id, _)| event_id.clone()) + else { + break; + }; + self.executor_event_receipts.remove(&oldest_id); + } + } } fn command_result(command: &Command, status: &str, result: Value) -> StoredCommandResult { @@ -510,6 +648,19 @@ fn command_fingerprint(command: &Command) -> Result Ok(fingerprint) } +fn executor_event_fingerprint( + event_type: &str, + priority: EventPriority, + payload: &Value, +) -> String { + let identity = json!({ + "eventType": event_type, + "priority": priority.number(), + "payload": sanitize_value(payload), + }); + format!("{:x}", Sha256::digest(canonical_json(&identity).as_bytes())) +} + fn canonical_json(value: &Value) -> String { match value { Value::Array(values) => format!( @@ -757,6 +908,23 @@ fn validate_binding( || (allow_legacy_command_journal && !state.processed_commands.is_empty() && state.processed_command_fingerprints.is_empty()); + let mut executor_receipt_sequences = HashSet::new(); + let executor_event_receipts_are_valid = state.executor_event_receipts.len() + <= MAX_EXECUTOR_EVENT_RECEIPTS + && state + .executor_event_receipts + .iter() + .all(|(event_id, receipt)| { + state.source_event_id_for_executor(event_id).is_ok() + && receipt.fingerprint.len() == 64 + && receipt + .fingerprint + .bytes() + .all(|byte| byte.is_ascii_hexdigit()) + && receipt.source_seq > 0 + && receipt.source_seq <= state.highest_source_seq() + && executor_receipt_sequences.insert(receipt.source_seq) + }); command_sequences.sort_unstable(); let command_cursors_are_valid = match (command_sequences.first(), command_sequences.last()) { (None, None) => state.compacted_through_controller_seq == state.last_controller_command_seq, @@ -782,6 +950,7 @@ fn validate_binding( || state.compacted_through_controller_seq > state.last_controller_command_seq || !command_cursors_are_valid || !command_fingerprints_are_valid + || !executor_event_receipts_are_valid { return Err(DurableRunnerError::invalid( "durable state cursors, bounds, or journals are inconsistent", @@ -790,7 +959,7 @@ fn validate_binding( Ok(()) } -fn verify_private_directory(path: &Path) -> Result<(), DurableRunnerError> { +pub(crate) fn verify_private_directory(path: &Path) -> Result<(), DurableRunnerError> { let metadata = fs::symlink_metadata(path) .map_err(|error| DurableRunnerError::invalid(error.to_string()))?; if metadata.file_type().is_symlink() || !metadata.is_dir() { @@ -807,7 +976,7 @@ fn verify_private_directory(path: &Path) -> Result<(), DurableRunnerError> { Ok(()) } -fn open_private_regular_file(path: &Path) -> io::Result { +pub(crate) fn open_private_regular_file(path: &Path) -> io::Result { let mut options = OpenOptions::new(); options.read(true); #[cfg(unix)] @@ -854,7 +1023,9 @@ const fn no_follow_flag() -> i32 { 0 } -fn create_private_temporary_file(path: &Path) -> Result<(PathBuf, File), DurableRunnerError> { +pub(crate) fn create_private_temporary_file( + path: &Path, +) -> Result<(PathBuf, File), DurableRunnerError> { let parent = path .parent() .ok_or_else(|| DurableRunnerError::invalid("durable state path has no parent"))?; @@ -882,6 +1053,17 @@ fn create_private_temporary_file(path: &Path) -> Result<(PathBuf, File), Durable fn sensitive_key(key: &str) -> bool { let normalized = key.to_ascii_lowercase().replace(['-', '_'], ""); + if matches!( + normalized.as_str(), + "inputtokens" + | "outputtokens" + | "cachereadtokens" + | "cachewritetokens" + | "pretokens" + | "posttokens" + ) { + return false; + } [ "authorization", "cookie", @@ -919,7 +1101,7 @@ fn sanitize_value(value: &Value) -> Value { } } -fn redact_text(input: &str) -> String { +pub(crate) fn redact_text(input: &str) -> String { let (bounded, truncated) = if input.len() > 4096 { let boundary = input .char_indices() @@ -1124,7 +1306,7 @@ mod tests { &config, "runner.diagnostic", EventPriority::P1, - json!({"nested": {"api_token": "secret-value"}}), + json!({"nested": {"api_token": "secret-value", "inputTokens": 42}}), ) .unwrap(); assert_eq!( @@ -1133,6 +1315,12 @@ mod tests { .pointer("/payload/payload/nested/api_token"), Some(&Value::String("[REDACTED]".to_owned())) ); + assert_eq!( + state.outbox[0] + .envelope + .pointer("/payload/payload/nested/inputTokens"), + Some(&json!(42)) + ); } #[test] diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/durable/transport.rs b/packages/paperclip-runner/runner/crates/runner-core/src/durable/transport.rs index 8bbab5733b..8386624a59 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/durable/transport.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/durable/transport.rs @@ -1402,6 +1402,7 @@ mod tests { fn reconnect_replays_unacked_events_and_not_command_effects() { struct EventExecutor { session_open_calls: Arc, + shutdown_calls: Arc, } impl super::super::CommandExecutor for EventExecutor { @@ -1424,6 +1425,11 @@ mod tests { json!({"status": "completed"}), )) } + + fn shutdown(&mut self) -> Result<(), DurableRunnerError> { + self.shutdown_calls.fetch_add(1, Ordering::SeqCst); + Ok(()) + } } let listener = TcpListener::bind("127.0.0.1:0").unwrap(); @@ -1544,16 +1550,19 @@ mod tests { }); let session_open_calls = Arc::new(AtomicUsize::new(0)); + let shutdown_calls = Arc::new(AtomicUsize::new(0)); super::super::run_durable_runner( config, BootstrapTicket::new("bootstrap-secret".to_owned()).unwrap(), EventExecutor { session_open_calls: session_open_calls.clone(), + shutdown_calls: shutdown_calls.clone(), }, ) .unwrap(); server.join().unwrap(); assert_eq!(session_open_calls.load(Ordering::SeqCst), 1); + assert_eq!(shutdown_calls.load(Ordering::SeqCst), 1); let store = super::super::DurableStateStore::new(&directory).unwrap(); let state_bytes = std::fs::read(store.path()).unwrap(); let final_state: DurableState = serde_json::from_slice(&state_bytes).unwrap(); diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/lib.rs b/packages/paperclip-runner/runner/crates/runner-core/src/lib.rs index 926aa0c540..9960607bab 100644 --- a/packages/paperclip-runner/runner/crates/runner-core/src/lib.rs +++ b/packages/paperclip-runner/runner/crates/runner-core/src/lib.rs @@ -1,9 +1,12 @@ #![forbid(unsafe_code)] +pub mod codex_provider; pub mod durable; pub mod fake_harness; pub mod local_runner; pub mod process_supervisor; +pub mod provider_backend; +pub mod provider_events; pub mod replay; use std::error::Error; diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs b/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs new file mode 100644 index 0000000000..77186cf892 --- /dev/null +++ b/packages/paperclip-runner/runner/crates/runner-core/src/provider_backend.rs @@ -0,0 +1,781 @@ +use std::collections::{HashSet, VecDeque}; +use std::fs; +use std::io::{Read, Write}; +use std::path::PathBuf; + +#[cfg(unix)] +use std::fs::File; +#[cfg(unix)] +use std::os::unix::fs::PermissionsExt; + +use serde::{Deserialize, Serialize}; +use serde_json::{json, Value}; + +use crate::codex_provider::{CodexProvider, CodexProviderConfig, CodexProviderEvent}; +use crate::durable::{ + create_private_temporary_file, open_private_regular_file, verify_private_directory, Command, + CommandExecution, CommandExecutor, DurableRunnerError, EventPriority, PolledEvent, +}; +use crate::provider_events::{normalize_codex_notification, NormalizedProviderEvent}; + +const PROVIDER_STATE_SCHEMA: &str = "paperclip.runner.codex-provider-state.v1"; +const PROVIDER_STATE_FILE: &str = "codex-provider-state.json"; +const MAX_PROVIDER_STATE_BYTES: u64 = 2 * 1024 * 1024; +const MAX_EVENTS_PER_POLL: usize = 128; + +fn initial_provider_event_seq() -> u64 { + 1 +} + +fn provider_event_id(sequence: u64) -> String { + format!("codex_provider_{sequence:016}") +} + +fn provider_event_sequence(event_id: &str) -> Option { + let sequence = event_id.strip_prefix("codex_provider_")?.parse().ok()?; + (provider_event_id(sequence) == event_id).then_some(sequence) +} + +#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)] +#[serde(rename_all = "camelCase")] +struct CodexProviderState { + schema: String, + lifecycle: String, + config: CodexProviderConfig, + #[serde(default)] + thread_id: Option, + #[serde(default)] + provider_session_id: Option, + #[serde(default)] + active_provider_turn_id: Option, + #[serde(default)] + pending_events: VecDeque, + #[serde(default = "initial_provider_event_seq")] + next_provider_event_seq: u64, +} + +impl CodexProviderState { + fn new(config: CodexProviderConfig) -> Self { + Self { + schema: PROVIDER_STATE_SCHEMA.to_owned(), + lifecycle: "prepared".to_owned(), + config, + thread_id: None, + provider_session_id: None, + active_provider_turn_id: None, + pending_events: VecDeque::new(), + next_provider_event_seq: initial_provider_event_seq(), + } + } + + fn validate(&self) -> Result<(), DurableRunnerError> { + self.config + .validate() + .map_err(|error| DurableRunnerError::invalid(error.to_string()))?; + let mut pending_event_ids = HashSet::new(); + if self.schema != PROVIDER_STATE_SCHEMA + || !matches!( + self.lifecycle.as_str(), + "prepared" | "session_open" | "turn_active" | "provider_exited" | "closed" + ) + || self + .thread_id + .as_ref() + .is_some_and(|value| value.is_empty() || value.len() > 240) + || self + .provider_session_id + .as_ref() + .is_some_and(|value| value.is_empty() || value.len() > 240) + || self + .active_provider_turn_id + .as_ref() + .is_some_and(|value| value.is_empty() || value.len() > 240) + || (self.thread_id.is_none() + && (self.provider_session_id.is_some() + || self.active_provider_turn_id.is_some() + || matches!(self.lifecycle.as_str(), "session_open" | "turn_active"))) + || (self.lifecycle == "turn_active" && self.active_provider_turn_id.is_none()) + || (matches!( + self.lifecycle.as_str(), + "prepared" | "session_open" | "closed" + ) && self.active_provider_turn_id.is_some()) + || self.next_provider_event_seq == 0 + || self.pending_events.len() > MAX_EVENTS_PER_POLL + 1 + || self.pending_events.iter().any(|event| { + provider_event_sequence(&event.executor_event_id) + .is_none_or(|sequence| sequence >= self.next_provider_event_seq) + || !pending_event_ids.insert(event.executor_event_id.as_str()) + || event.event_type.is_empty() + || event.event_type.len() > 160 + || event.event_type.chars().any(char::is_control) + || !event.payload.is_object() + }) + { + return Err(DurableRunnerError::invalid( + "Codex provider state is malformed or inconsistent", + )); + } + Ok(()) + } + + fn push_event(&mut self, event: NormalizedProviderEvent) -> Result<(), DurableRunnerError> { + let sequence = self.next_provider_event_seq; + self.next_provider_event_seq = sequence + .checked_add(1) + .ok_or_else(|| DurableRunnerError::invalid("provider event sequence exhausted"))?; + self.pending_events.push_back(PolledEvent { + executor_event_id: provider_event_id(sequence), + event_type: event.event_type, + priority: event.priority, + payload: event.payload, + }); + Ok(()) + } + + fn extend_events( + &mut self, + events: impl IntoIterator, + ) -> Result<(), DurableRunnerError> { + for event in events { + self.push_event(event)?; + } + Ok(()) + } +} + +pub struct CodexCommandExecutor { + state_dir: PathBuf, + state: Option, + provider: Option, + restore_checked: bool, +} + +impl CodexCommandExecutor { + pub fn new(state_dir: impl Into) -> Self { + Self { + state_dir: state_dir.into(), + state: None, + provider: None, + restore_checked: false, + } + } + + fn state_path(&self) -> PathBuf { + self.state_dir.join(PROVIDER_STATE_FILE) + } + + fn restore(&mut self) -> Result<(), DurableRunnerError> { + if self.restore_checked { + return Ok(()); + } + self.restore_checked = true; + let path = self.state_path(); + let mut file = match open_private_regular_file(&path) { + Ok(file) => file, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()), + Err(error) => { + return Err(DurableRunnerError::invalid(format!( + "failed to open private Codex provider state: {error}" + ))) + } + }; + let metadata = file.metadata().map_err(|error| { + DurableRunnerError::invalid(format!("failed to inspect Codex provider state: {error}")) + })?; + if metadata.len() > MAX_PROVIDER_STATE_BYTES { + return Err(DurableRunnerError::invalid( + "Codex provider state must be a bounded regular file", + )); + } + let mut input = Vec::with_capacity(metadata.len() as usize); + file.read_to_end(&mut input).map_err(|error| { + DurableRunnerError::invalid(format!("failed to read Codex provider state: {error}")) + })?; + let state: CodexProviderState = serde_json::from_slice(&input).map_err(|error| { + DurableRunnerError::invalid(format!("Codex provider state is malformed: {error}")) + })?; + state.validate()?; + self.state = Some(state); + self.restore_provider_if_needed() + } + + fn restore_provider_if_needed(&mut self) -> Result<(), DurableRunnerError> { + let Some(state) = self.state.as_ref() else { + return Ok(()); + }; + if self.provider.is_some() + || !matches!( + state.lifecycle.as_str(), + "session_open" | "turn_active" | "provider_exited" + ) + { + return Ok(()); + } + let provider_had_exited = state.lifecycle == "provider_exited"; + let thread_id = state.thread_id.clone().ok_or_else(|| { + DurableRunnerError::invalid("recoverable Codex state omitted its thread id") + })?; + let previous_active_turn_id = state.active_provider_turn_id.clone(); + let provider = CodexProvider::start(&state.config, Some(&thread_id)).map_err(|error| { + DurableRunnerError::invalid(format!("failed to resume Codex provider: {error}")) + })?; + let recovered_active_turn_id = provider.active_provider_turn_id().map(str::to_owned); + self.provider = Some(provider); + if provider_had_exited || recovered_active_turn_id != previous_active_turn_id { + let state = self + .state + .as_mut() + .expect("Codex state remains available during recovery"); + state.active_provider_turn_id = recovered_active_turn_id.clone(); + state.lifecycle = if recovered_active_turn_id.is_some() { + "turn_active".to_owned() + } else { + "session_open".to_owned() + }; + state.push_event(NormalizedProviderEvent { + event_type: "session.reconciled".to_owned(), + priority: EventPriority::P0, + payload: json!({ + "provider": "codex", + "providerSessionId": thread_id, + "previousProviderTurnId": previous_active_turn_id, + "activeProviderTurnId": recovered_active_turn_id, + }), + })?; + self.save_state()?; + } + Ok(()) + } + + fn save_state(&self) -> Result<(), DurableRunnerError> { + let state = self + .state + .as_ref() + .ok_or_else(|| DurableRunnerError::invalid("Codex provider state is unavailable"))?; + self.persist_state(state) + } + + fn persist_state(&self, state: &CodexProviderState) -> Result<(), DurableRunnerError> { + state.validate()?; + fs::create_dir_all(&self.state_dir).map_err(|error| { + DurableRunnerError::invalid(format!( + "failed to create provider state directory: {error}" + )) + })?; + #[cfg(unix)] + fs::set_permissions(&self.state_dir, fs::Permissions::from_mode(0o700)).map_err( + |error| { + DurableRunnerError::invalid(format!( + "failed to protect provider state directory: {error}" + )) + }, + )?; + verify_private_directory(&self.state_dir)?; + let path = self.state_path(); + let bytes = serde_json::to_vec_pretty(state).map_err(|error| { + DurableRunnerError::invalid(format!( + "failed to serialize Codex provider state: {error}" + )) + })?; + if bytes.len() as u64 > MAX_PROVIDER_STATE_BYTES { + return Err(DurableRunnerError::invalid( + "Codex provider state exceeds the 2 MiB limit", + )); + } + let (temporary, mut file) = create_private_temporary_file(&path)?; + let result = (|| -> std::io::Result<()> { + file.write_all(&bytes)?; + file.sync_all()?; + drop(file); + fs::rename(&temporary, &path)?; + #[cfg(unix)] + File::open(&self.state_dir)?.sync_all()?; + Ok(()) + })(); + if let Err(error) = result { + let _ = fs::remove_file(&temporary); + return Err(DurableRunnerError::invalid(format!( + "failed to replace provider state atomically: {error}" + ))); + } + #[cfg(unix)] + fs::set_permissions(&path, fs::Permissions::from_mode(0o600)).map_err(|error| { + DurableRunnerError::invalid(format!("failed to protect provider state: {error}")) + })?; + Ok(()) + } + + fn prepare(&mut self, payload: &Value) -> Result { + let config: CodexProviderConfig = serde_json::from_value( + payload + .get("provider") + .cloned() + .ok_or_else(|| DurableRunnerError::invalid("run.prepare requires provider"))?, + ) + .map_err(|error| { + DurableRunnerError::invalid(format!("run.prepare provider is invalid: {error}")) + })?; + config + .validate() + .map_err(|error| DurableRunnerError::invalid(error.to_string()))?; + if let Some(state) = &self.state { + if state.config != config { + return Err(DurableRunnerError::invalid( + "Codex provider configuration changed across the durable run", + )); + } + if state.lifecycle == "closed" { + return Err(DurableRunnerError::invalid( + "Codex provider session is already closed", + )); + } + } else { + self.state = Some(CodexProviderState::new(config)); + self.save_state()?; + } + Ok(CommandExecution::result(json!({ + "status": "prepared", + "provider": "codex", + "driver": "codex_app_server", + }))) + } + + fn ensure_provider(&mut self) -> Result<&mut CodexProvider, DurableRunnerError> { + self.restore_provider_if_needed()?; + if self.provider.is_none() { + let state = self.state.as_ref().ok_or_else(|| { + DurableRunnerError::invalid("Codex provider has not been prepared") + })?; + if state.lifecycle == "closed" { + return Err(DurableRunnerError::invalid( + "Codex provider session is closed", + )); + } + let provider = CodexProvider::start(&state.config, state.thread_id.as_deref()) + .map_err(|error| { + DurableRunnerError::invalid(format!("failed to start Codex provider: {error}")) + })?; + self.provider = Some(provider); + } + self.provider + .as_mut() + .ok_or_else(|| DurableRunnerError::invalid("Codex provider is unavailable")) + } + + fn open_session(&mut self) -> Result { + self.restore_provider_if_needed()?; + if self + .state + .as_ref() + .is_some_and(|state| state.lifecycle == "turn_active") + { + return Err(DurableRunnerError::invalid( + "cannot open a new Codex session while a provider turn is active", + )); + } + let resumed = self + .state + .as_ref() + .and_then(|state| state.thread_id.as_ref()) + .is_some(); + let (thread_id, provider_session_id, process_id) = { + let provider = self.ensure_provider()?; + ( + provider.thread_id().to_owned(), + provider.provider_session_id().map(str::to_owned), + provider.process_id(), + ) + }; + let provider_version = { + let state = self + .state + .as_mut() + .expect("Codex state exists after provider start"); + state.thread_id = Some(thread_id.clone()); + state.provider_session_id = provider_session_id.clone(); + state.active_provider_turn_id = None; + state.lifecycle = "session_open".to_owned(); + state.config.provider_version.clone() + }; + self.save_state()?; + Ok(CommandExecution { + result: json!({ + "status": if resumed { "resumed" } else { "started" }, + "provider": "codex", + "driver": "codex_app_server", + "providerVersion": provider_version, + "providerSessionId": thread_id, + "processId": process_id, + }), + events: vec![( + if resumed { + "session.resumed" + } else { + "session.started" + } + .to_owned(), + EventPriority::P0, + json!({ + "provider": "codex", + "providerSessionId": thread_id, + "providerAccountSessionId": provider_session_id, + "processId": process_id, + }), + )], + }) + } + + fn start_turn(&mut self, payload: &Value) -> Result { + self.restore_provider_if_needed()?; + if self + .state + .as_ref() + .is_some_and(|state| state.active_provider_turn_id.is_some()) + { + return Err(DurableRunnerError::invalid( + "Codex already has an active provider turn", + )); + } + let text = payload + .get("text") + .and_then(Value::as_str) + .ok_or_else(|| DurableRunnerError::invalid("turn.start payload.text is required"))?; + let cwd = self + .state + .as_ref() + .ok_or_else(|| DurableRunnerError::invalid("Codex provider is not prepared"))? + .config + .cwd + .clone(); + let (provider_turn_id, thread_id) = { + let provider = self.ensure_provider()?; + provider.start_turn(text, &cwd).map_err(|error| { + DurableRunnerError::invalid(format!("Codex turn/start failed: {error}")) + })?; + ( + provider + .active_provider_turn_id() + .ok_or_else(|| { + DurableRunnerError::invalid("Codex turn/start omitted its turn identity") + })? + .to_owned(), + provider.thread_id().to_owned(), + ) + }; + let state = self + .state + .as_mut() + .expect("Codex state exists after turn start"); + state.active_provider_turn_id = Some(provider_turn_id.clone()); + state.lifecycle = "turn_active".to_owned(); + self.save_state()?; + Ok(CommandExecution { + result: json!({"status": "accepted", "providerTurnId": provider_turn_id}), + events: vec![( + "turn.accepted".to_owned(), + EventPriority::P0, + json!({"provider": "codex", "providerSessionId": thread_id, "providerTurnId": provider_turn_id}), + )], + }) + } + + fn interrupt_turn(&mut self, reason: &str) -> Result { + self.restore_provider_if_needed()?; + let provider_turn_id = self + .state + .as_ref() + .and_then(|state| state.active_provider_turn_id.clone()); + if provider_turn_id.is_none() { + return Ok(CommandExecution::result(json!({ + "status": "already_settled", + "reason": reason, + }))); + } + self.ensure_provider()?.interrupt_turn().map_err(|error| { + DurableRunnerError::invalid(format!("Codex turn interrupt failed: {error}")) + })?; + Ok(CommandExecution::result(json!({ + "status": "interrupt_requested", + "reason": reason, + "providerTurnId": provider_turn_id, + }))) + } + + fn steer_turn(&mut self, payload: &Value) -> Result { + let text = payload + .get("text") + .and_then(Value::as_str) + .ok_or_else(|| DurableRunnerError::invalid("turn.steer payload.text is required"))?; + self.ensure_provider()?.steer_turn(text).map_err(|error| { + DurableRunnerError::invalid(format!("Codex turn steer failed: {error}")) + })?; + Ok(CommandExecution::result(json!({"status": "steered"}))) + } + + fn resolve_request(&mut self, payload: &Value) -> Result { + let request_id = payload + .get("requestId") + .and_then(Value::as_str) + .ok_or_else(|| DurableRunnerError::invalid("request.resolve requires requestId"))?; + let response = payload + .get("response") + .ok_or_else(|| DurableRunnerError::invalid("request.resolve requires response"))?; + self.ensure_provider()? + .resolve_runtime_request(request_id, response) + .map_err(|error| { + DurableRunnerError::invalid(format!("Codex runtime response failed: {error}")) + })?; + Ok(CommandExecution { + result: json!({"status": "delivered", "requestId": request_id}), + events: vec![( + "runtime_request.resolved".to_owned(), + EventPriority::P0, + json!({"provider": "codex", "requestId": request_id, "status": "delivered"}), + )], + }) + } + + fn close_session(&mut self) -> Result { + if let Some(provider) = self.provider.as_mut() { + provider.shutdown().map_err(|error| { + DurableRunnerError::invalid(format!("failed to stop Codex provider: {error}")) + })?; + } + self.provider = None; + let state = self + .state + .as_mut() + .ok_or_else(|| DurableRunnerError::invalid("Codex provider is not prepared"))?; + state.active_provider_turn_id = None; + state.lifecycle = "closed".to_owned(); + let thread_id = state.thread_id.clone(); + self.save_state()?; + Ok(CommandExecution { + result: json!({"status": "closed", "providerSessionId": thread_id}), + events: vec![( + "session.closed".to_owned(), + EventPriority::P0, + json!({"provider": "codex", "providerSessionId": thread_id}), + )], + }) + } + + fn snapshot(&mut self) -> Result { + self.restore_provider_if_needed()?; + let state = self + .state + .as_ref() + .ok_or_else(|| DurableRunnerError::invalid("Codex provider is not prepared"))?; + Ok(CommandExecution::result(json!({ + "status": state.lifecycle, + "provider": "codex", + "driver": "codex_app_server", + "providerSessionId": state.thread_id, + "activeProviderTurnId": state.active_provider_turn_id, + }))) + } + + fn poll_provider(&mut self) -> Result<(), DurableRunnerError> { + self.restore()?; + if self + .state + .as_ref() + .is_some_and(|state| !state.pending_events.is_empty()) + { + return Ok(()); + } + if self.provider.is_none() { + return Ok(()); + } + for _ in 0..MAX_EVENTS_PER_POLL { + let event = self + .provider + .as_mut() + .expect("provider remains present while polling") + .poll() + .map_err(|error| { + DurableRunnerError::invalid(format!("Codex provider failed: {error}")) + })?; + let Some(event) = event else { break }; + match event { + CodexProviderEvent::Notification { method, params } => { + let normalized = normalize_codex_notification(&method, ¶ms); + let state = self + .state + .as_mut() + .expect("Codex state remains available while polling"); + if method == "turn/completed" { + state.active_provider_turn_id = None; + state.lifecycle = "session_open".to_owned(); + } + state.extend_events(normalized)?; + self.save_state()?; + } + CodexProviderEvent::RuntimeRequest { + request_id, + question_set, + } => { + let prompt = question_set + .get("title") + .or_else(|| question_set.pointer("/questions/0/prompt")) + .and_then(Value::as_str) + .unwrap_or("Codex needs your input"); + self.state + .as_mut() + .expect("Codex state remains available while polling") + .push_event(NormalizedProviderEvent { + event_type: "runtime_request.created".to_owned(), + priority: EventPriority::P0, + payload: json!({ + "request": { + "schema": "paperclip.runtime_request.v2", + "requestKind": "runtime", + "requestId": request_id, + "type": "input", + "status": "pending", + "prompt": prompt, + "input": question_set, + "origin": { + "adapter": "codex-app-server", + "provider": "codex", + "method": "item/tool/requestUserInput", + }, + }, + }), + })?; + self.save_state()?; + } + CodexProviderEvent::Exited { exit_code, success } => { + self.provider = None; + if let Some(state) = self.state.as_mut() { + state.lifecycle = "provider_exited".to_owned(); + state.push_event(NormalizedProviderEvent { + event_type: "session.failed".to_owned(), + priority: EventPriority::P0, + payload: json!({ + "provider": "codex", + "code": "provider_exited", + "exitCode": exit_code, + "expected": success, + }), + })?; + } + self.save_state()?; + break; + } + } + } + Ok(()) + } +} + +impl CommandExecutor for CodexCommandExecutor { + fn execute(&mut self, command: &Command) -> Result { + self.restore()?; + match command.command_type.as_str() { + "run.prepare" => self.prepare(&command.payload), + "run.attach" => { + if self.state.is_none() && command.payload.get("provider").is_some() { + self.prepare(&command.payload)?; + } + let mut execution = self.open_session()?; + execution.events.push(( + "run.attached".to_owned(), + EventPriority::P0, + json!({"provider": "codex"}), + )); + Ok(execution) + } + "session.open" => self.open_session(), + "turn.start" => self.start_turn(&command.payload), + "turn.steer" => self.steer_turn(&command.payload), + "turn.interrupt" | "turn.stop" | "run.cancel" => { + self.interrupt_turn(&command.command_type) + } + "request.resolve" => self.resolve_request(&command.payload), + "session.snapshot" => self.snapshot(), + "session.close" | "session.destroy" => self.close_session(), + "runner.drain" | "runner.suspend" | "runner.shutdown" => { + Ok(CommandExecution::result(json!({"status": "completed"}))) + } + _ => Ok(CommandExecution::result(json!({ + "status": "rejected", + "code": "provider_command_unavailable", + "message": "the Codex provider does not implement this command in the current layer", + }))), + } + } + + fn poll_events(&mut self) -> Result, DurableRunnerError> { + self.poll_provider()?; + Ok(self + .state + .as_ref() + .into_iter() + .flat_map(|state| state.pending_events.iter()) + .cloned() + .collect()) + } + + fn acknowledge_events(&mut self, count: usize) -> Result<(), DurableRunnerError> { + if count == 0 { + return Ok(()); + } + let mut next_state = self + .state + .clone() + .ok_or_else(|| DurableRunnerError::invalid("Codex provider state is unavailable"))?; + if count > next_state.pending_events.len() { + return Err(DurableRunnerError::invalid( + "provider event acknowledgement exceeded the pending prefix", + )); + } + next_state.pending_events.drain(..count); + self.persist_state(&next_state)?; + self.state = Some(next_state); + Ok(()) + } + + fn shutdown(&mut self) -> Result<(), DurableRunnerError> { + if let Some(provider) = self.provider.as_mut() { + provider.shutdown().map_err(|error| { + DurableRunnerError::invalid(format!("failed to stop Codex provider: {error}")) + })?; + } + self.provider = None; + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn rejects_inconsistent_provider_state() { + let state = CodexProviderState { + schema: PROVIDER_STATE_SCHEMA.to_owned(), + lifecycle: "turn_active".to_owned(), + config: CodexProviderConfig { + provider: "codex".to_owned(), + driver: "codex_app_server".to_owned(), + provider_version: "test".to_owned(), + command: PathBuf::from("codex"), + args: Vec::new(), + cwd: std::env::current_dir() + .unwrap() + .to_string_lossy() + .into_owned(), + model: None, + instructions: String::new(), + approval_policy: "never".to_owned(), + }, + thread_id: Some("thread-1".to_owned()), + provider_session_id: None, + active_provider_turn_id: None, + pending_events: VecDeque::new(), + next_provider_event_seq: initial_provider_event_seq(), + }; + assert!(state.validate().is_err()); + } +} diff --git a/packages/paperclip-runner/runner/crates/runner-core/src/provider_events.rs b/packages/paperclip-runner/runner/crates/runner-core/src/provider_events.rs new file mode 100644 index 0000000000..f9e7b0d8fe --- /dev/null +++ b/packages/paperclip-runner/runner/crates/runner-core/src/provider_events.rs @@ -0,0 +1,375 @@ +use serde_json::{json, Value}; +use sha2::{Digest, Sha256}; + +use crate::durable::{redact_text, EventPriority}; + +const MAX_TEXT_CHARS: usize = 4_000; + +#[derive(Clone, Debug, PartialEq)] +pub struct NormalizedProviderEvent { + pub event_type: String, + pub priority: EventPriority, + pub payload: Value, +} + +fn bounded_text(value: &str, max_chars: usize) -> String { + redact_text(value).chars().take(max_chars).collect() +} + +fn string(value: Option<&Value>) -> &str { + value.and_then(Value::as_str).unwrap_or("") +} + +fn stable_id(value: &str, fallback: &str) -> String { + let value: String = value + .chars() + .map(|character| { + if character.is_ascii_alphanumeric() || "._:-".contains(character) { + character + } else { + '-' + } + }) + .take(160) + .collect(); + if value + .chars() + .next() + .is_some_and(|character| character.is_ascii_alphanumeric()) + { + value + } else { + fallback.to_owned() + } +} + +fn item(params: &Value) -> &Value { + params.get("item").unwrap_or(&Value::Null) +} + +fn provider_status(value: &str, completed: bool) -> &'static str { + match value { + "failed" | "error" => "failed", + "cancelled" | "canceled" => "cancelled", + "interrupted" | "aborted" => "interrupted", + _ if completed => "completed", + _ => "running", + } +} + +fn bounded_output(value: &str) -> Value { + let output = redact_text(value); + let output_truncated = output != value; + json!({ + "output": output, + "outputBytes": value.len(), + "outputTruncated": output_truncated, + "outputDigest": format!("sha256:{:x}", Sha256::digest(value.as_bytes())), + }) +} + +fn measurement(value: &Value) -> Value { + json!({ + "inputTokens": value.get("inputTokens").and_then(Value::as_u64).unwrap_or(0), + "outputTokens": value.get("outputTokens").and_then(Value::as_u64).unwrap_or(0), + "cacheReadTokens": value.get("cachedInputTokens").or_else(|| value.get("cacheReadTokens")).and_then(Value::as_u64).unwrap_or(0), + "cacheWriteTokens": value.get("cacheWriteTokens").and_then(Value::as_u64).unwrap_or(0), + "activeSeconds": value.get("activeSeconds").and_then(Value::as_f64).filter(|value| *value >= 0.0).unwrap_or(0.0), + "requests": value.get("requests").and_then(Value::as_u64).unwrap_or(0), + "providerCostUsd": value.get("providerCostUsd").and_then(Value::as_f64).filter(|value| *value >= 0.0).unwrap_or(0.0), + }) +} + +/// Converts Codex app-server notifications into provider-neutral PRP events. +/// Provider-native envelopes are consumed here and never cross the PRP boundary. +pub fn normalize_codex_notification(method: &str, params: &Value) -> Vec { + let mut events = Vec::new(); + let push = |events: &mut Vec, + event_type: &str, + priority: EventPriority, + payload: Value| { + events.push(NormalizedProviderEvent { + event_type: event_type.to_owned(), + priority, + payload, + }); + }; + + match method { + "thread/compacted" => push( + &mut events, + "context.compacted", + EventPriority::P1, + json!({ + "schema": "paperclip.context.compacted.v1", + "compactionId": stable_id(string(params.get("threadId")), "codex-compaction"), + "reason": "provider", + "preTokens": Value::Null, + "postTokens": Value::Null, + "sameSession": true, + }), + ), + "turn/started" => push( + &mut events, + "turn.started", + EventPriority::P0, + json!({ + "provider": "codex", + "providerTurnId": params.pointer("/turn/id").or_else(|| params.get("turnId")).and_then(Value::as_str), + }), + ), + "turn/completed" => { + let status = string( + params + .pointer("/turn/status") + .or_else(|| params.get("status")), + ); + let event_type = match status { + "failed" | "error" => "turn.failed", + "cancelled" | "canceled" => "turn.cancelled", + "interrupted" | "aborted" => "turn.interrupted", + _ => "turn.completed", + }; + push( + &mut events, + event_type, + EventPriority::P0, + json!({ + "provider": "codex", + "providerTurnId": params.pointer("/turn/id").or_else(|| params.get("turnId")).and_then(Value::as_str), + "status": provider_status(status, true), + }), + ); + } + "turn/plan/updated" => { + let plan_id = stable_id(string(params.get("turnId")), "codex-plan"); + let steps = params + .get("plan") + .and_then(Value::as_array) + .map(|steps| { + steps + .iter() + .take(256) + .enumerate() + .filter_map(|(index, step)| { + let body = bounded_text(string(step.get("step")), MAX_TEXT_CHARS); + if body.trim().is_empty() { + return None; + } + Some(json!({ + "stepId": format!("step-{}", index + 1), + "body": body, + "status": match string(step.get("status")) { + "inProgress" | "in_progress" => "in_progress", + "completed" => "completed", + "blocked" | "failed" | "error" => "blocked", + _ => "pending", + }, + })) + }) + .collect::>() + }) + .unwrap_or_default(); + let complete = !steps.is_empty() + && steps + .iter() + .all(|step| step.get("status").and_then(Value::as_str) == Some("completed")); + push( + &mut events, + "plan.updated", + EventPriority::P1, + json!({ + "schema": "paperclip.plan.updated.v1", + "planId": plan_id, + "revision": params.get("revision").and_then(Value::as_u64).filter(|value| *value > 0).unwrap_or(1), + "explanation": params.get("explanation").and_then(Value::as_str).map(|value| bounded_text(value, MAX_TEXT_CHARS)), + "steps": steps, + "complete": complete, + "syncStatus": "not_applicable", + "documentRevision": Value::Null, + }), + ); + } + "thread/tokenUsage/updated" => { + let cumulative = params + .get("tokenUsage") + .and_then(|value| value.get("total")) + .or_else(|| params.get("total")) + .unwrap_or(&Value::Null); + let run_delta = params + .get("tokenUsage") + .and_then(|value| value.get("last")) + .or_else(|| params.get("last")) + .unwrap_or(cumulative); + push( + &mut events, + "usage.reported", + EventPriority::P0, + json!({ + "provider": "codex", + "model": params.get("model").and_then(Value::as_str).map(|value| bounded_text(value, 240)), + "providerSessionId": params.get("threadId").and_then(Value::as_str).map(|value| bounded_text(value, 240)), + "providerRequestId": Value::Null, + "cumulative": measurement(cumulative), + "runDelta": measurement(run_delta), + }), + ); + } + "error" | "warning" | "deprecationNotice" | "configWarning" => push( + &mut events, + "provider.notice.recorded", + EventPriority::P0, + json!({ + "schema": "paperclip.provider.notice.v1", + "noticeId": stable_id(&format!("codex-{method}"), "codex-notice"), + "severity": if method == "error" { "error" } else { "warning" }, + "category": method.replace('/', "_"), + "scope": if method.contains("config") { "environment" } else { "turn" }, + "recoverable": method != "error", + "userActionable": true, + "summary": bounded_text(string(params.get("message")), MAX_TEXT_CHARS), + }), + ), + "item/agentMessage/delta" => push( + &mut events, + "item.delta", + EventPriority::P2, + json!({ + "provider": "codex", + "itemId": stable_id(string(params.get("itemId")), "codex-message"), + "kind": "agentMessage", + "channel": "progress", + "providerMethod": method, + "text": bounded_text(string(params.get("delta")), MAX_TEXT_CHARS), + }), + ), + "item/started" | "item/completed" => { + let provider_item = item(params); + let item_id = stable_id(string(provider_item.get("id")), "codex-item"); + let item_type = string(provider_item.get("type")); + let completed = method == "item/completed"; + if matches!(item_type, "commandExecution" | "mcpToolCall") { + let mut payload = json!({ + "schema": "paperclip.tool.execution.v1", + "executionId": item_id, + "transport": if item_type == "mcpToolCall" { "mcp" } else { "process" }, + "operation": if item_type == "commandExecution" { "execute" } else { "unknown" }, + "name": provider_item.get("tool").or_else(|| provider_item.get("command")).and_then(Value::as_str).map(|value| bounded_text(value, 240)), + "target": Value::Null, + "namespace": provider_item.get("server").and_then(Value::as_str).map(|value| bounded_text(value, 240)), + "readOnly": provider_item.get("readOnlyHint").and_then(Value::as_bool), + "status": provider_status(string(provider_item.get("status")), completed), + "durationMs": provider_item.get("durationMs").and_then(Value::as_u64), + "exitCode": provider_item.get("exitCode").and_then(Value::as_i64), + "progress": Value::Null, + }); + if let (Some(object), Value::Object(output)) = ( + payload.as_object_mut(), + bounded_output(string( + provider_item + .get("aggregatedOutput") + .or_else(|| provider_item.get("output")), + )), + ) { + object.extend(output); + } + push( + &mut events, + if completed { + "tool.execution.completed" + } else { + "tool.execution.started" + }, + if completed { + EventPriority::P1 + } else { + EventPriority::P2 + }, + payload, + ); + } else { + push( + &mut events, + if completed { + "item.completed" + } else { + "item.started" + }, + if completed { + EventPriority::P1 + } else { + EventPriority::P2 + }, + json!({ + "provider": "codex", + "itemId": item_id, + "kind": bounded_text(item_type, 160), + "status": provider_status(string(provider_item.get("status")), completed), + "channel": if item_type == "agentMessage" { "progress" } else { "detail" }, + "text": provider_item.get("text").and_then(Value::as_str).map(|value| bounded_text(value, MAX_TEXT_CHARS)), + }), + ); + } + } + _ => {} + } + + events +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn maps_codex_plan_without_retaining_native_envelope() { + let events = normalize_codex_notification( + "turn/plan/updated", + &json!({ + "turnId": "turn-1", + "revision": 2, + "plan": [{"step": "Inspect", "status": "inProgress"}], + "accessToken": "secret-value", + }), + ); + assert_eq!(events[0].event_type, "plan.updated"); + assert_eq!(events[0].payload["steps"][0]["status"], "in_progress"); + assert!(!events[0].payload.to_string().contains("secret-value")); + } + + #[test] + fn bounds_and_redacts_command_output() { + let events = normalize_codex_notification( + "item/completed", + &json!({"item": { + "id": "exec-1", + "type": "commandExecution", + "status": "completed", + "command": "printenv", + "aggregatedOutput": "Authorization: Bearer top-secret", + }}), + ); + assert_eq!(events[0].event_type, "tool.execution.completed"); + assert_eq!(events[0].payload["outputTruncated"], true); + assert_eq!(events[0].payload["outputBytes"], 32); + assert!(!events[0].payload.to_string().contains("top-secret")); + } + + #[test] + fn maps_terminal_and_usage_events_at_priority_zero() { + let terminal = normalize_codex_notification( + "turn/completed", + &json!({"turn": {"id": "provider-turn", "status": "failed"}}), + ); + assert_eq!(terminal[0].event_type, "turn.failed"); + assert_eq!(terminal[0].priority, EventPriority::P0); + + let usage = normalize_codex_notification( + "thread/tokenUsage/updated", + &json!({"tokenUsage": {"total": {"inputTokens": 12, "outputTokens": 3}}}), + ); + assert_eq!(usage[0].event_type, "usage.reported"); + assert_eq!(usage[0].payload["cumulative"]["inputTokens"], 12); + assert_eq!(usage[0].priority, EventPriority::P0); + } +} diff --git a/packages/paperclip-runner/runner/crates/runner-core/tests/codex_provider.rs b/packages/paperclip-runner/runner/crates/runner-core/tests/codex_provider.rs new file mode 100644 index 0000000000..22859843e6 --- /dev/null +++ b/packages/paperclip-runner/runner/crates/runner-core/tests/codex_provider.rs @@ -0,0 +1,385 @@ +use std::fs; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicU64, Ordering}; + +use paperclip_runner_core::codex_provider::{ + CodexProvider, CodexProviderConfig, CodexProviderEvent, +}; +use paperclip_runner_core::durable::{Command, CommandExecutor, DurableRunnerError, PolledEvent}; +use paperclip_runner_core::provider_backend::CodexCommandExecutor; +use paperclip_runner_core::provider_events::normalize_codex_notification; +use serde_json::{json, Value}; + +static NEXT_DIRECTORY: AtomicU64 = AtomicU64::new(1); + +fn temporary_directory(label: &str) -> PathBuf { + let directory = std::env::temp_dir().join(format!( + "paperclip-runner-codex-{label}-{}-{}", + std::process::id(), + NEXT_DIRECTORY.fetch_add(1, Ordering::Relaxed) + )); + let _ = fs::remove_dir_all(&directory); + fs::create_dir_all(&directory).expect("create Codex integration-test directory"); + directory +} + +fn provider_config(directory: &Path, switches: &[&str]) -> CodexProviderConfig { + let mut args = vec![ + "--state-file".to_owned(), + directory + .join("fake-state.json") + .to_string_lossy() + .into_owned(), + "--call-log".to_owned(), + directory.join("calls.log").to_string_lossy().into_owned(), + ]; + args.extend(switches.iter().map(|value| (*value).to_owned())); + CodexProviderConfig { + provider: "codex".to_owned(), + driver: "codex_app_server".to_owned(), + provider_version: "fake-1".to_owned(), + command: PathBuf::from(env!("CARGO_BIN_EXE_fake-codex-app-server")), + args, + cwd: std::env::current_dir() + .expect("resolve test cwd") + .to_string_lossy() + .into_owned(), + model: Some("test-model".to_owned()), + instructions: "Stay inside the test workspace.".to_owned(), + approval_policy: "never".to_owned(), + } +} + +fn command(id: &str, sequence: u64, command_type: &str, payload: Value) -> Command { + Command { + schema: "paperclip.prp.command.v1".to_owned(), + command_id: id.to_owned(), + controller_seq: sequence, + command_type: command_type.to_owned(), + issued_at: "2026-08-24T00:00:00.000Z".to_owned(), + deadline_at: None, + precondition: None, + payload, + } +} + +fn call_count(directory: &Path, method: &str) -> usize { + fs::read_to_string(directory.join("calls.log")) + .unwrap_or_default() + .lines() + .filter(|line| *line == method) + .count() +} + +fn poll_and_ack( + executor: &mut CodexCommandExecutor, +) -> Result, DurableRunnerError> { + let events = executor.poll_events()?; + executor.acknowledge_events(events.len())?; + Ok(events) +} + +#[test] +fn codex_transport_buffers_notifications_while_waiting_for_responses() { + let directory = temporary_directory("buffering"); + let config = provider_config(&directory, &["--notification-before-response"]); + let mut provider = CodexProvider::start(&config, None).expect("start fake Codex provider"); + let event = provider + .poll() + .expect("poll buffered notification") + .expect("buffered notification is available"); + let CodexProviderEvent::Notification { method, params } = event else { + panic!("expected the pre-response warning notification"); + }; + assert_eq!(method, "warning"); + let normalized = normalize_codex_notification(&method, ¶ms); + assert_eq!(normalized[0].event_type, "provider.notice.recorded"); + + provider + .start_turn("Complete the fake task.", &config.cwd) + .expect("start provider turn"); + let mut event_types = Vec::new(); + for _ in 0..16 { + if let Some(CodexProviderEvent::Notification { method, params }) = + provider.poll().expect("poll provider event") + { + event_types.extend( + normalize_codex_notification(&method, ¶ms) + .into_iter() + .map(|event| event.event_type), + ); + } + if event_types.iter().any(|event| event == "turn.completed") { + break; + } + } + assert!(event_types.iter().any(|event| event == "turn.started")); + assert!(event_types.iter().any(|event| event == "item.completed")); + assert!(event_types.iter().any(|event| event == "usage.reported")); + assert!(event_types.iter().any(|event| event == "turn.completed")); + provider.shutdown().expect("stop provider"); + fs::remove_dir_all(directory).expect("remove Codex integration-test directory"); +} + +#[test] +fn durable_backend_resumes_the_active_thread_without_restarting_the_turn() { + let directory = temporary_directory("resume"); + let config = provider_config(&directory, &["--hold-turn"]); + let mut first = CodexCommandExecutor::new(&directory); + first + .execute(&command( + "prepare", + 1, + "run.prepare", + json!({"provider": config}), + )) + .expect("prepare Codex provider"); + first + .execute(&command("open", 2, "session.open", json!({}))) + .expect("open Codex session"); + first + .execute(&command( + "turn", + 3, + "turn.start", + json!({"text": "Hold this turn for recovery."}), + )) + .expect("start held provider turn"); + assert_eq!(call_count(&directory, "turn/start"), 1); + first.shutdown().expect("stop first provider process"); + drop(first); + + let mut recovered = CodexCommandExecutor::new(&directory); + let snapshot = recovered + .execute(&command("snapshot", 4, "session.snapshot", json!({}))) + .expect("restore provider session"); + assert_eq!(snapshot.result["status"], "turn_active"); + assert_eq!(snapshot.result["activeProviderTurnId"], "provider-turn-1"); + assert_eq!(call_count(&directory, "turn/start"), 1); + assert_eq!(call_count(&directory, "thread/resume"), 1); + assert_eq!(call_count(&directory, "thread/read"), 1); + + recovered + .execute(&command("interrupt", 5, "turn.interrupt", json!({}))) + .expect("interrupt recovered provider turn"); + let mut terminal_seen = false; + for _ in 0..16 { + let events = poll_and_ack(&mut recovered).expect("poll interrupted turn"); + terminal_seen |= events + .iter() + .any(|event| event.event_type == "turn.interrupted"); + if terminal_seen { + break; + } + } + assert!(terminal_seen); + recovered + .shutdown() + .expect("stop recovered provider process"); + fs::remove_dir_all(directory).expect("remove Codex integration-test directory"); +} + +#[test] +fn provider_exit_preserves_and_reconciles_the_active_turn() { + let directory = temporary_directory("exit-active-turn"); + let config = provider_config(&directory, &["--exit-after-turn-start"]); + let mut executor = CodexCommandExecutor::new(&directory); + executor + .execute(&command( + "prepare", + 1, + "run.prepare", + json!({"provider": config}), + )) + .expect("prepare Codex provider"); + executor + .execute(&command("open", 2, "session.open", json!({}))) + .expect("open Codex session"); + executor + .execute(&command( + "turn", + 3, + "turn.start", + json!({"text": "Keep the native turn active while the provider exits."}), + )) + .expect("start provider turn"); + + let mut provider_exit_seen = false; + for _ in 0..32 { + provider_exit_seen |= poll_and_ack(&mut executor) + .expect("poll provider exit") + .iter() + .any(|event| event.event_type == "session.failed"); + if provider_exit_seen { + break; + } + } + assert!(provider_exit_seen); + let persisted: Value = serde_json::from_slice( + &fs::read(directory.join("codex-provider-state.json")) + .expect("read provider state after exit"), + ) + .expect("parse provider state after exit"); + assert_eq!(persisted["lifecycle"], "provider_exited"); + assert_eq!(persisted["activeProviderTurnId"], "provider-turn-1"); + + let interrupted = executor + .execute(&command("interrupt", 4, "turn.interrupt", json!({}))) + .expect("interrupt reconciled provider turn"); + assert_eq!(interrupted.result["status"], "interrupt_requested"); + assert_eq!(call_count(&directory, "thread/resume"), 1); + assert_eq!(call_count(&directory, "thread/read"), 1); + assert_eq!(call_count(&directory, "turn/interrupt"), 1); + + let mut terminal_seen = false; + for _ in 0..32 { + terminal_seen |= poll_and_ack(&mut executor) + .expect("poll reconciled interruption") + .iter() + .any(|event| event.event_type == "turn.interrupted"); + if terminal_seen { + break; + } + } + assert!(terminal_seen); + executor.shutdown().expect("stop resumed provider process"); + fs::remove_dir_all(directory).expect("remove Codex integration-test directory"); +} + +#[test] +fn unacknowledged_provider_events_survive_executor_restart() { + let directory = temporary_directory("pending-event-recovery"); + let config = provider_config(&directory, &["--emit-question"]); + let mut first = CodexCommandExecutor::new(&directory); + first + .execute(&command( + "prepare", + 1, + "run.prepare", + json!({"provider": config}), + )) + .expect("prepare Codex provider"); + first + .execute(&command("open", 2, "session.open", json!({}))) + .expect("open Codex session"); + first + .execute(&command( + "turn", + 3, + "turn.start", + json!({"text": "Emit a durable question."}), + )) + .expect("start provider turn"); + + let mut retained = None; + for _ in 0..32 { + let events = first.poll_events().expect("poll provider events"); + if events + .iter() + .any(|event| event.event_type == "runtime_request.created") + { + retained = Some(events); + break; + } + first + .acknowledge_events(events.len()) + .expect("acknowledge events before the question"); + } + let retained = retained.expect("observe a durable runtime request"); + first.shutdown().expect("stop first provider process"); + drop(first); + + let mut recovered = CodexCommandExecutor::new(&directory); + let replayed = recovered + .poll_events() + .expect("reload unacknowledged provider events"); + assert_eq!(&replayed[..retained.len()], retained.as_slice()); + recovered + .acknowledge_events(replayed.len()) + .expect("acknowledge reloaded provider events"); + recovered + .shutdown() + .expect("stop recovered provider process"); + fs::remove_dir_all(directory).expect("remove Codex integration-test directory"); +} + +#[test] +fn structured_question_round_trips_through_the_normalized_backend() { + let directory = temporary_directory("questions"); + let config = provider_config(&directory, &["--emit-question"]); + let mut executor = CodexCommandExecutor::new(&directory); + executor + .execute(&command( + "prepare", + 1, + "run.prepare", + json!({"provider": config}), + )) + .expect("prepare provider"); + executor + .execute(&command("open", 2, "session.open", json!({}))) + .expect("open provider session"); + let started = executor + .execute(&command( + "turn", + 3, + "turn.start", + json!({"text": "Ask for deployment input."}), + )) + .expect("start provider turn"); + assert_eq!(started.events.len(), 1); + assert_eq!(started.events[0].0, "turn.accepted"); + + let mut question_set = None; + let mut provider_started_events = 0; + for _ in 0..16 { + for event in poll_and_ack(&mut executor).expect("poll question") { + provider_started_events += usize::from(event.event_type == "turn.started"); + if event.event_type == "runtime_request.created" { + assert_eq!( + event.payload["request"]["schema"], + "paperclip.runtime_request.v2" + ); + question_set = event.payload.pointer("/request/input").cloned(); + } + } + if question_set.is_some() { + break; + } + } + let question_set = question_set.expect("normalized question set is emitted"); + assert_eq!(provider_started_events, 1); + assert_eq!(question_set["schema"], "paperclip.question_set.v1"); + assert_eq!( + question_set["questions"][0]["options"][0]["label"], + "Staging" + ); + + executor + .execute(&command( + "resolve", + 4, + "request.resolve", + json!({ + "requestId": "runtime-request-1", + "response": { + "schema": "paperclip.question_response.v1", + "answers": {"environment": {"selectedOptionIds": ["option-1"]}} + } + }), + )) + .expect("deliver normalized response"); + let mut completed = false; + for _ in 0..16 { + completed |= poll_and_ack(&mut executor) + .expect("poll completed question turn") + .iter() + .any(|event| event.event_type == "turn.completed"); + if completed { + break; + } + } + assert!(completed); + executor.shutdown().expect("stop provider process"); + fs::remove_dir_all(directory).expect("remove Codex integration-test directory"); +}