From 4e3524819c0417cfd39b8e77c5f58900d8de78a4 Mon Sep 17 00:00:00 2001 From: guoxu1 Date: Tue, 8 Sep 2026 19:58:21 +0800 Subject: [PATCH 1/8] fix(replay): satisfy workspace clippy and pin OpenCode sampling --- .../persisting-replay/src/adapter/generic.rs | 164 +++--------------- 1 file changed, 28 insertions(+), 136 deletions(-) diff --git a/crates/persisting-replay/src/adapter/generic.rs b/crates/persisting-replay/src/adapter/generic.rs index a3b1359ba..38bdc5210 100644 --- a/crates/persisting-replay/src/adapter/generic.rs +++ b/crates/persisting-replay/src/adapter/generic.rs @@ -8,7 +8,7 @@ use std::fs; use std::path::{Path, PathBuf}; -use std::process::{Command, Stdio}; +use std::process::Command; use std::time::{Duration, Instant}; use serde_json::{Value, json}; @@ -25,7 +25,6 @@ use crate::model::{ AgentKind, FreshObservation, PlaybackRequest, ReplayMode, ReplayOutcome, ReplayPlan, ToolBatch, ToolCall, }; -use crate::opencode_bridge; use crate::process::{ProcessSpec, run_process}; #[derive(Debug, Clone, Copy)] @@ -994,8 +993,6 @@ fn continue_native_cli( let mut codex_bridge = None; let mut codex_transport_prompt = None; let mut codex_prompt_mode = None; - let mut opencode_bridge = None; - let mut opencode_transport_prompt = None; command.env("PVISOR_REPLAY_TRAJECTORY", reconstructed); command.env("PVISOR_REPLAY_AFTER_STEP", plan.after_step.to_string()); command.env( @@ -1010,32 +1007,11 @@ fn continue_native_cli( match agent { NativeJsonlAgent::Opencode => { let session_id = opencode_session_id(&session_id); - // `opencode run --session` refuses to start without a message. - // Pass a unique transport nonce as that message and strip it on - // the wire through the local bridge, so the first live request - // still ends exactly at the replayed boundary observation. - let explicit_prompt = context.request.boundary_user_prompt().map(str::to_owned); - let transport_prompt = explicit_prompt - .clone() - .unwrap_or_else(|| format!("pvisor-opencode-resume-{}", context.nonce)); - let temperature = env_f64("PVISOR_OPENCODE_TEMPERATURE"); - let top_p = env_f64("PVISOR_OPENCODE_TOP_P"); - let bridge = opencode_bridge::OpencodeBridgeHandle::start( - context.session_id, - explicit_prompt.is_none().then(|| transport_prompt.clone()), - temperature, - top_p, - context.request.disable_thinking, - )?; + command.env("PVISOR_REPLAY_SESSION_ID", &session_id); + let export_path = context.output_dir.join("native/opencode-session.json"); let opencode_config = context.state_dir.join("opencode-config"); let opencode_data = context.state_dir.join("opencode-data"); - write_opencode_provider_config( - &opencode_config, - Some(&bridge.base_url), - temperature, - top_p, - )?; - let export_path = context.output_dir.join("native/opencode-session.json"); + write_opencode_provider_config(&opencode_config)?; atomic_write_json( &export_path, &opencode_export(plan, prefix, &session_id, &context.request.workspace), @@ -1074,9 +1050,6 @@ fn continue_native_cli( command.env("XDG_CONFIG_HOME", &opencode_config); command.env("XDG_DATA_HOME", &opencode_data); command.env("OPENCODE_DISABLE_AUTOUPDATE", "1"); - for (name, value) in bridge.child_environment() { - command.env(name, value); - } if let Some(model) = configured_model_from_environment() { command.args(["--model", &model]); } @@ -1090,15 +1063,10 @@ fn continue_native_cli( if !context.request.disable_thinking { command.arg("--thinking"); } - command.arg("--"); - command.arg(&transport_prompt); - // OpenCode awaits stdin EOF whenever it is not a TTY; inheriting - // the controller's stdin would hang the continuation forever. - command.stdin(Stdio::null()); - if explicit_prompt.is_none() { - opencode_transport_prompt = Some(transport_prompt); + if let Some(prompt) = context.request.boundary_user_prompt() { + command.arg("--"); + command.arg(prompt); } - opencode_bridge = Some(bridge); } NativeJsonlAgent::Codex => { let explicit_prompt = context.request.boundary_user_prompt().map(str::to_owned); @@ -1183,10 +1151,7 @@ fn continue_native_cli( log_path: log_path.clone(), }) .map_err(|error| ReplayError::new(ReplayErrorKind::Continuation, error.message))?; - let bridge_result = codex_bridge - .take() - .map(|bridge| bridge.finish()) - .or_else(|| opencode_bridge.take().map(|bridge| bridge.finish())); + let bridge_result = codex_bridge.take().map(|bridge| bridge.finish()); let bridge_error = bridge_result.and_then(|result| result.err()); if !output.status.success() { let process_error = ReplayError::classify_continuation( @@ -1241,13 +1206,6 @@ fn continue_native_cli( NativeJsonlAgent::Opencode => { let raw = read_regular_file(&log_path)?; let events = parse_json_lines_from_log(&raw); - // The transport nonce was only a CLI wake-up signal; drop it if - // the native stream echoed it back as a user or text event. - let nonce = opencode_transport_prompt.as_deref().unwrap_or_default(); - let events: Vec = events - .into_iter() - .filter(|event| !opencode_event_is_nonce(event, nonce)) - .collect(); let steps = count_opencode_turns(&events); let mut combined = prefix.to_vec(); combined.extend(events); @@ -1282,8 +1240,6 @@ fn env_f64(name: &str) -> Option { /// have no environment channel, so a live continuation would silently fall /// back to provider defaults and diverge from the recorded sampling. The /// shape mirrors what a SweEval trial writes for the original run. -/// `effective_base` overrides the environment endpoint (used for the local -/// sampling-injection proxy). fn opencode_provider_config( model: &str, base_url: Option<&str>, @@ -1315,54 +1271,28 @@ fn opencode_provider_config( Some(json!({ "provider": { provider: Value::Object(provider_config) } })) } -fn write_opencode_provider_config( - config_root: &Path, - effective_base: Option<&str>, - temperature: Option, - top_p: Option, -) -> Result<(), ReplayError> { +fn write_opencode_provider_config(config_root: &Path) -> Result<(), ReplayError> { let Some(model) = configured_model_from_environment() else { return Ok(()); }; - let base_url = match effective_base { - Some(base) => Some(base.to_owned()), - None => std::env::var("OPENAI_BASE_URL") - .ok() - .or_else(|| std::env::var("OPENAI_API_BASE").ok()), - }; - let config = opencode_provider_config(&model, base_url.as_deref(), temperature, top_p); + let base_url = std::env::var("OPENAI_BASE_URL") + .ok() + .or_else(|| std::env::var("OPENAI_API_BASE").ok()); + let config = opencode_provider_config( + &model, + base_url.as_deref(), + env_f64("PVISOR_OPENCODE_TEMPERATURE"), + env_f64("PVISOR_OPENCODE_TOP_P"), + ); let Some(config) = config else { return Ok(()); }; let directory = config_root.join("opencode"); - fs::create_dir_all(&directory).replay_context( - ReplayErrorKind::Executor, - "create OpenCode config directory", - )?; + fs::create_dir_all(&directory) + .replay_context(ReplayErrorKind::Executor, "create OpenCode config directory")?; atomic_write_json(&directory.join("opencode.json"), &config) } -/// True when the native event echoes the transport nonce back as a user or -/// text part; such events are transport noise, not model input. -fn opencode_event_is_nonce(event: &Value, nonce: &str) -> bool { - if nonce.is_empty() { - return false; - } - match event.get("type").and_then(Value::as_str) { - Some("user") => event - .get("parts") - .and_then(Value::as_array) - .map(|parts| { - parts - .iter() - .any(|part| part.get("text") == Some(&json!(nonce))) - }) - .unwrap_or(false), - Some("text") => event.pointer("/part/text") == Some(&json!(nonce)), - _ => false, - } -} - fn continuation_session_id( agent: NativeJsonlAgent, plan: &ReplayPlan, @@ -1703,17 +1633,6 @@ fn opencode_export( .get("user_prompt") .and_then(Value::as_str) .unwrap_or_default(); - // OpenCode resolves a session's default model from the last user message - // metadata when a request does not pin one. The synthetic placeholder - // must therefore carry the configured model; "pvisor/replay" would poison - // that fallback with a provider that does not exist. - let (placeholder_provider, placeholder_model) = configured_model_from_environment() - .and_then(|model| { - model - .split_once('/') - .map(|(p, m)| (p.to_owned(), m.to_owned())) - }) - .unwrap_or_else(|| ("pvisor".to_owned(), "replay".to_owned())); let mut messages = vec![json!({ "info": { "id": user_id, @@ -1721,7 +1640,7 @@ fn opencode_export( "role": "user", "time": {"created": 0}, "agent": "build", - "model": {"providerID": placeholder_provider, "modelID": placeholder_model}, + "model": {"providerID": "pvisor", "modelID": "replay"}, }, "parts": [{ "id": "prt_pvisor_user", @@ -1849,8 +1768,8 @@ fn opencode_export( "role": "assistant", "time": {"created": batch.ordinal as u64, "completed": batch.ordinal as u64}, "parentID": user_id, - "modelID": placeholder_model, - "providerID": placeholder_provider, + "modelID": "replay", + "providerID": "pvisor", "mode": "build", "agent": "build", "path": {"cwd": workspace.display().to_string(), "root": workspace.display().to_string()}, @@ -1987,37 +1906,12 @@ mod tests { use super::{ CallRecord, NativeJsonlAgent, RunContext, TurnRecord, codex_native_session_id, - continuation_session_id, is_actionable_turn, opencode_event_is_nonce, - opencode_provider_config, parse_codex, parse_jsonl, parse_opencode, - redact_codex_transport_nonce, validate_codex_continuation, + continuation_session_id, is_actionable_turn, opencode_provider_config, parse_codex, + parse_jsonl, parse_opencode, redact_codex_transport_nonce, validate_codex_continuation, }; use crate::model::{AgentKind, PlaybackRequest, ReplayMode, ReplayPlan, ToolBatch, ToolCall}; use serde_json::{Value, json}; - #[test] - fn opencode_nonce_events_are_filtered_from_the_continued_stream() { - let nonce = "pvisor-opencode-resume-nonce"; - let events = vec![ - json!({"type": "step_start", "sessionID": "ses"}), - json!({"type": "user", "sessionID": "ses", "parts": [{"type": "text", "text": nonce}]}), - json!({"type": "text", "sessionID": "ses", "part": {"type": "text", "text": nonce}}), - json!({"type": "text", "sessionID": "ses", "part": {"type": "text", "text": "real text"}}), - json!({"type": "step_finish", "sessionID": "ses"}), - ]; - let kept: Vec = events - .iter() - .filter(|event| !opencode_event_is_nonce(event, nonce)) - .cloned() - .collect(); - let kinds: Vec<&str> = kept.iter().map(|e| e["type"].as_str().unwrap()).collect(); - assert_eq!(kinds, vec!["step_start", "text", "step_finish"]); - assert_eq!(kept[1]["part"]["text"], "real text"); - // An empty nonce (explicit boundary prompt mode) filters nothing. - for event in &events { - assert!(!opencode_event_is_nonce(event, "")); - } - } - #[test] fn opencode_provider_config_mirrors_recorded_sampling() { let config = opencode_provider_config( @@ -2041,8 +1935,8 @@ mod tests { // Without sampling overrides the endpoint still comes from the // environment, so only the baseURL section is written. - let base_only = - opencode_provider_config("openai/model-x", Some("http://m:1/v1"), None, None).unwrap(); + let base_only = opencode_provider_config("openai/model-x", Some("http://m:1/v1"), None, None) + .unwrap(); assert_eq!( base_only, json!({"provider": {"openai": {"options": {"baseURL": "http://m:1/v1"}}}}) @@ -2051,9 +1945,7 @@ mod tests { // Nothing to pin: leave OpenCode on its environment-only defaults. assert!(opencode_provider_config("openai/model-x", None, None, None).is_none()); // A model without a provider namespace cannot be pinned either. - assert!( - opencode_provider_config("model-x", Some("http://m:1/v1"), Some(0.0), None).is_none() - ); + assert!(opencode_provider_config("model-x", Some("http://m:1/v1"), Some(0.0), None).is_none()); // Blank endpoints are ignored rather than written. assert!(opencode_provider_config("openai/model-x", Some(" "), None, None).is_none()); } From 79e41eca6f31e9713a04e5a321391c6a38d0e0d3 Mon Sep 17 00:00:00 2001 From: guoxu1 Date: Tue, 8 Sep 2026 20:15:23 +0800 Subject: [PATCH 2/8] fix(replay): stabilize OpenCode continuation --- .../persisting-replay/src/adapter/generic.rs | 590 ++++++++++++++++-- 1 file changed, 521 insertions(+), 69 deletions(-) diff --git a/crates/persisting-replay/src/adapter/generic.rs b/crates/persisting-replay/src/adapter/generic.rs index 38bdc5210..0506a74e5 100644 --- a/crates/persisting-replay/src/adapter/generic.rs +++ b/crates/persisting-replay/src/adapter/generic.rs @@ -1007,66 +1007,11 @@ fn continue_native_cli( match agent { NativeJsonlAgent::Opencode => { let session_id = opencode_session_id(&session_id); - command.env("PVISOR_REPLAY_SESSION_ID", &session_id); - let export_path = context.output_dir.join("native/opencode-session.json"); - let opencode_config = context.state_dir.join("opencode-config"); - let opencode_data = context.state_dir.join("opencode-data"); - write_opencode_provider_config(&opencode_config)?; - atomic_write_json( - &export_path, - &opencode_export(plan, prefix, &session_id, &context.request.workspace), - )?; - let mut import = agent_command(&launch.entrypoint, context); - import.args([ - "import", - export_path.to_str().ok_or_else(|| { - ReplayError::configuration("OpenCode export path is not valid UTF-8") - })?, - ]); - import.env("XDG_CONFIG_HOME", &opencode_config); - import.env("XDG_DATA_HOME", &opencode_data); - import.env("OPENCODE_DISABLE_AUTOUPDATE", "1"); - let import_log = logs.join("opencode-import.log"); - let imported = run_process(ProcessSpec { - command: import, - stdin: None, - timeout: Duration::from_secs(5 * 60), - termination_grace: Duration::from_secs(2), - pipe_grace: Duration::from_millis(250), - retained_bytes: MAX_TOOL_OUTPUT_BYTES / 4, - log_path: import_log.clone(), - }) - .map_err(|error| ReplayError::new(ReplayErrorKind::Continuation, error.message))?; - if !imported.status.success() { - return Err(ReplayError::classify_continuation( - format!( - "OpenCode session import exited {}; see {}", - imported.status, - import_log.display() - ), - &String::from_utf8_lossy(&imported.stderr_tail), - )); - } - command.env("XDG_CONFIG_HOME", &opencode_config); - command.env("XDG_DATA_HOME", &opencode_data); - command.env("OPENCODE_DISABLE_AUTOUPDATE", "1"); - if let Some(model) = configured_model_from_environment() { - command.args(["--model", &model]); - } - command.args([ - "run", - "--format=json", - "--session", - &session_id, - "--dangerously-skip-permissions", - ]); - if !context.request.disable_thinking { - command.arg("--thinking"); - } - if let Some(prompt) = context.request.boundary_user_prompt() { - command.arg("--"); - command.arg(prompt); - } + // The CLI refuses `run --session` without a message, and any + // message would contaminate the boundary. Continue through the + // server API instead; see opencode_server_continuation. + let _ = &command; + return opencode_server_continuation(plan, context, journal, prefix, &session_id, &log_path); } NativeJsonlAgent::Codex => { let explicit_prompt = context.request.boundary_user_prompt().map(str::to_owned); @@ -1240,6 +1185,8 @@ fn env_f64(name: &str) -> Option { /// have no environment channel, so a live continuation would silently fall /// back to provider defaults and diverge from the recorded sampling. The /// shape mirrors what a SweEval trial writes for the original run. +/// `effective_base` overrides the environment endpoint (used for the local +/// sampling-injection proxy). fn opencode_provider_config( model: &str, base_url: Option<&str>, @@ -1271,18 +1218,26 @@ fn opencode_provider_config( Some(json!({ "provider": { provider: Value::Object(provider_config) } })) } -fn write_opencode_provider_config(config_root: &Path) -> Result<(), ReplayError> { +fn write_opencode_provider_config( + config_root: &Path, + effective_base: Option<&str>, + temperature: Option, + top_p: Option, +) -> Result<(), ReplayError> { let Some(model) = configured_model_from_environment() else { return Ok(()); }; - let base_url = std::env::var("OPENAI_BASE_URL") - .ok() - .or_else(|| std::env::var("OPENAI_API_BASE").ok()); + let base_url = match effective_base { + Some(base) => Some(base.to_owned()), + None => std::env::var("OPENAI_BASE_URL") + .ok() + .or_else(|| std::env::var("OPENAI_API_BASE").ok()), + }; let config = opencode_provider_config( &model, base_url.as_deref(), - env_f64("PVISOR_OPENCODE_TEMPERATURE"), - env_f64("PVISOR_OPENCODE_TOP_P"), + temperature, + top_p, ); let Some(config) = config else { return Ok(()); @@ -1293,6 +1248,496 @@ fn write_opencode_provider_config(config_root: &Path) -> Result<(), ReplayError> atomic_write_json(&directory.join("opencode.json"), &config) } +/// Node script that forwards requests to the model endpoint and injects +/// sampling parameters into JSON bodies. OpenCode 1.17.7 never forwards +/// `temperature`/`top_p` to the Responses API regardless of configuration, +/// so pinning sampling has to happen on the wire. Node comes from the same +/// runtime tree as the OpenCode entrypoint. +fn opencode_sampling_proxy_script() -> &'static str { + r#"const http = require("node:http") +const https = require("node:https") +const fs = require("node:fs") +const upstream = new URL(process.env.PVISOR_PROXY_UPSTREAM) +const temperature = Number(process.env.PVISOR_PROXY_TEMPERATURE) +const topP = process.env.PVISOR_PROXY_TOP_P === "" ? null : Number(process.env.PVISOR_PROXY_TOP_P) +const server = http.createServer((request, response) => { + const chunks = [] + request.on("data", (chunk) => chunks.push(chunk)) + request.on("end", () => { + let body = Buffer.concat(chunks) + const headers = {...request.headers} + delete headers["content-length"] + delete headers["transfer-encoding"] + if (String(headers["content-type"] || "").includes("application/json") && body.length > 0) { + try { + const document = JSON.parse(body.toString("utf8")) + if (document && typeof document === "object" && !Array.isArray(document)) { + document.temperature = temperature + if (topP !== null) document.top_p = topP + body = Buffer.from(JSON.stringify(document)) + } + } catch (error) { /* forward the original body */ } + } + const target = new URL(upstream) + target.pathname = request.url + const transport = target.protocol === "https:" ? https : http + const forwarded = transport.request(target, {method: request.method, headers}, (up) => { + response.writeHead(up.statusCode, up.headers) + up.pipe(response) + }) + forwarded.on("error", (error) => { + if (!response.headersSent) response.writeHead(502) + if (!response.writableEnded) response.end(String(error)) + }) + forwarded.end(body) + }) +}) +server.listen(Number(process.env.PVISOR_PROXY_PORT), "127.0.0.1", () => { + fs.writeFileSync(process.env.PVISOR_PROXY_READY, "ready") +}) +"# +} + +struct ChildGuard(std::process::Child); + +impl Drop for ChildGuard { + fn drop(&mut self) { + let _ = self.0.kill(); + let _ = self.0.wait(); + } +} + +fn loopback_free_port() -> Result { + let listener = std::net::TcpListener::bind(("127.0.0.1", 0)) + .replay_context(ReplayErrorKind::Executor, "allocate loopback port")?; + let port = listener + .local_addr() + .replay_context(ReplayErrorKind::Executor, "read loopback port")? + .port(); + Ok(port) +} + +fn wait_for_file(path: &Path, timeout: Duration) -> Result<(), ReplayError> { + let started = Instant::now(); + while !path.is_file() { + if started.elapsed() > timeout { + return Err(ReplayError::continuation(format!( + "OpenCode helper did not become ready: {}", + path.display() + ))); + } + std::thread::sleep(Duration::from_millis(100)); + } + Ok(()) +} + +/// Start the sampling-injection proxy when sampling overrides are requested. +/// Returns the OpenAI-style API base URL OpenCode should use instead of the +/// real endpoint. +fn start_opencode_sampling_proxy( + entrypoint: &Path, + state_dir: &Path, + logs: &Path, + upstream: &str, + temperature: Option, + top_p: Option, + guard: &mut Vec, +) -> Result, ReplayError> { + if temperature.is_none() && top_p.is_none() { + return Ok(None); + } + let node = { + let sibling = entrypoint + .parent() + .map(|parent| parent.join("node")) + .filter(|path| path.is_file()) + .unwrap_or_else(|| PathBuf::from("node")); + sibling + }; + let port = loopback_free_port()?; + let upstream = upstream.trim_end_matches('/').to_owned(); + let ready = state_dir.join("opencode-sampling-proxy.ready"); + let _ = fs::remove_file(&ready); + let script = state_dir.join("opencode-sampling-proxy.js"); + atomic_write(&script, opencode_sampling_proxy_script().as_bytes())?; + let log = fs::File::create(logs.join("opencode-sampling-proxy.log")) + .replay_context(ReplayErrorKind::Executor, "create sampling proxy log")?; + let child = Command::new(&node) + .arg(&script) + .env("PVISOR_PROXY_UPSTREAM", &upstream) + .env( + "PVISOR_PROXY_TEMPERATURE", + temperature.map(|v| v.to_string()).unwrap_or_default(), + ) + .env( + "PVISOR_PROXY_TOP_P", + top_p.map(|v| v.to_string()).unwrap_or_default(), + ) + .env("PVISOR_PROXY_PORT", port.to_string()) + .env("PVISOR_PROXY_READY", &ready) + .stdout(log.try_clone().replay_context( + ReplayErrorKind::Executor, + "clone sampling proxy log handle", + )?) + .stderr(log) + .spawn() + .replay_context(ReplayErrorKind::Executor, "start OpenCode sampling proxy")?; + guard.push(ChildGuard(child)); + wait_for_file(&ready, Duration::from_secs(30))?; + Ok(Some(format!("http://127.0.0.1:{port}/v1"))) +} + +/// Project the assistant messages created by the live continuation into the +/// native OpenCode JSONL event stream. One assistant message is one agent +/// step; its parts are stored unordered, so restore the canonical order: +/// step-start first, step-finish last, content parts by their start time. +fn opencode_project_messages(messages: &[Value], session_id: &str) -> Vec { + let mut ordered: Vec<&Value> = messages + .iter() + .filter(|message| message.get("info").and_then(|i| i.get("role")) == Some(&json!("assistant"))) + .collect(); + ordered.sort_by_key(|message| { + message + .pointer("/info/time/created") + .and_then(Value::as_u64) + .unwrap_or(0) + }); + let mut events = Vec::new(); + for message in ordered { + let mut parts: Vec<&Value> = message + .get("parts") + .and_then(Value::as_array) + .map(|parts| parts.iter().collect()) + .unwrap_or_default(); + parts.sort_by(|a, b| opencode_part_order(a).cmp(&opencode_part_order(b))); + for part in parts { + let event_type = match part.get("type").and_then(Value::as_str) { + Some("step-start") => Some("step_start"), + Some("text") => Some("text"), + Some("reasoning") => Some("reasoning"), + Some("tool") => Some("tool_use"), + Some("step-finish") => Some("step_finish"), + _ => None, + }; + let Some(event_type) = event_type else { + continue; + }; + events.push(json!({ + "type": event_type, + "sessionID": session_id, + "part": part, + })); + } + } + events +} + +fn opencode_part_order(part: &Value) -> (u8, u64) { + let rank = match part.get("type").and_then(Value::as_str) { + Some("step-start") => 0, + Some("step-finish") => 2, + _ => 1, + }; + let start = part.pointer("/time/start").and_then(Value::as_u64).unwrap_or(0); + (rank, start) +} + +/// Continue an imported OpenCode session without adding any user message. +/// +/// `opencode run --session ` refuses to start without a message +/// ("You must provide a message or a command"), so the CLI cannot express a +/// clean continuation. The server API accepts `parts: []` on +/// `POST /session/{id}/message`; the resulting model request ends exactly at +/// the replayed boundary observation `O'N` with no injected prompt. +fn opencode_server_continuation( + plan: &ReplayPlan, + context: &RunContext<'_>, + journal: &mut Journal, + prefix: &[Value], + session_id: &str, + log_path: &Path, +) -> Result<(PathBuf, usize), ReplayError> { + let launch = context + .launch + .ok_or_else(|| ReplayError::continuation("OpenCode continuation has no launch spec"))?; + let state_dir = context.state_dir; + let opencode_config = state_dir.join("opencode-config"); + let opencode_data = state_dir.join("opencode-data"); + fs::create_dir_all(&opencode_data) + .replay_context(ReplayErrorKind::Executor, "create OpenCode data directory")?; + let logs = log_path.parent().unwrap_or_else(|| Path::new(".")).to_path_buf(); + + let temperature = env_f64("PVISOR_OPENCODE_TEMPERATURE"); + let top_p = env_f64("PVISOR_OPENCODE_TOP_P"); + let mut children: Vec = Vec::new(); + let result = opencode_server_continuation_inner( + plan, + context, + journal, + prefix, + session_id, + log_path, + &launch.entrypoint, + &opencode_config, + &opencode_data, + &logs, + temperature, + top_p, + &mut children, + ); + drop(children); + result +} + +#[allow(clippy::too_many_arguments)] +fn opencode_server_continuation_inner( + plan: &ReplayPlan, + context: &RunContext<'_>, + journal: &mut Journal, + prefix: &[Value], + session_id: &str, + log_path: &Path, + entrypoint: &Path, + opencode_config: &Path, + opencode_data: &Path, + logs: &Path, + temperature: Option, + top_p: Option, + children: &mut Vec, +) -> Result<(PathBuf, usize), ReplayError> { + let workspace = &context.request.workspace; + let upstream = std::env::var("OPENAI_BASE_URL") + .ok() + .or_else(|| std::env::var("OPENAI_API_BASE").ok()) + .filter(|value| !value.trim().is_empty()); + let effective_base = start_opencode_sampling_proxy( + entrypoint, + context.state_dir, + logs, + upstream.as_deref().unwrap_or("http://127.0.0.1"), + temperature, + top_p, + children, + )?; + write_opencode_provider_config(opencode_config, effective_base.as_deref(), temperature, top_p)?; + + let export_path = context.output_dir.join("native/opencode-session.json"); + atomic_write_json( + &export_path, + &opencode_export(plan, prefix, session_id, workspace), + )?; + let mut import = Command::new(entrypoint); + import + .arg("import") + .arg( + export_path + .to_str() + .ok_or_else(|| ReplayError::configuration("OpenCode export path is not valid UTF-8"))?, + ) + .env("XDG_CONFIG_HOME", opencode_config) + .env("XDG_DATA_HOME", opencode_data) + .env("OPENCODE_DISABLE_AUTOUPDATE", "1") + .current_dir(workspace); + let import_log = logs.join("opencode-import.log"); + let imported = run_process(ProcessSpec { + command: import, + stdin: None, + timeout: Duration::from_secs(5 * 60), + termination_grace: Duration::from_secs(2), + pipe_grace: Duration::from_millis(250), + retained_bytes: MAX_TOOL_OUTPUT_BYTES / 4, + log_path: import_log.clone(), + }) + .map_err(|error| ReplayError::new(ReplayErrorKind::Continuation, error.message))?; + if !imported.status.success() { + return Err(ReplayError::classify_continuation( + format!( + "OpenCode session import exited {}; see {}", + imported.status, + import_log.display() + ), + &String::from_utf8_lossy(&imported.stderr_tail), + )); + } + + journal.append( + "continuation_started", + [("agent".into(), json!("opencode"))], + )?; + + let serve_port = loopback_free_port()?; + let serve_log = fs::OpenOptions::new() + .create(true) + .append(true) + .open(log_path) + .replay_context(ReplayErrorKind::Executor, "open OpenCode serve log")?; + let mut serve = Command::new(entrypoint); + serve + .args(["serve", "--port", &serve_port.to_string(), "--hostname", "127.0.0.1"]) + .env("XDG_CONFIG_HOME", opencode_config) + .env("XDG_DATA_HOME", opencode_data) + .env("OPENCODE_DISABLE_AUTOUPDATE", "1") + .current_dir(workspace) + .stdout(serve_log.try_clone().replay_context( + ReplayErrorKind::Executor, + "clone OpenCode serve log handle", + )?) + .stderr(serve_log); + let serve = serve + .spawn() + .replay_context(ReplayErrorKind::Executor, "start OpenCode server")?; + children.push(ChildGuard(serve)); + + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .replay_context(ReplayErrorKind::Executor, "build OpenCode continuation runtime")?; + let base = format!("http://127.0.0.1:{serve_port}"); + let client = reqwest::Client::builder() + .no_proxy() + .build() + .replay_context(ReplayErrorKind::Executor, "build OpenCode continuation client")?; + let directory = workspace.to_string_lossy().to_string(); + + let ready = runtime.block_on(async { + for _ in 0..120 { + if let Ok(response) = client + .get(format!("{base}/config")) + .query(&[("directory", directory.as_str())]) + .send() + .await + { + if response.status().is_success() { + return Ok(()); + } + } + tokio::time::sleep(Duration::from_millis(500)).await; + } + Err(ReplayError::continuation( + "OpenCode server did not become ready; see opencode.log", + )) + }); + ready?; + + let messages_url = format!("{base}/session/{session_id}/message"); + let baseline: Vec = runtime + .block_on(async { + let response = client + .get(&messages_url) + .query(&[("directory", directory.as_str())]) + .send() + .await + .replay_context(ReplayErrorKind::Continuation, "read imported session messages")? + .error_for_status() + .replay_context(ReplayErrorKind::Continuation, "imported session request failed")?; + let messages: Value = response + .json() + .await + .replay_context(ReplayErrorKind::Continuation, "decode imported session messages")?; + Ok(message_ids(&messages)) + })?; + + let mut prompt_body = json!({"parts": []}); + if let Some(model) = configured_model_from_environment() { + if let Some((provider, model_id)) = model.split_once('/') { + prompt_body["model"] = json!({"providerID": provider, "modelID": model_id}); + } + } + let prompt_response: Value = runtime.block_on(async { + let response = client + .post(&messages_url) + .query(&[("directory", directory.as_str())]) + .json(&prompt_body) + .send() + .await + .replay_context(ReplayErrorKind::Continuation, "send OpenCode continuation prompt")?; + let status = response.status(); + let payload: Value = response + .json() + .await + .replay_context(ReplayErrorKind::Continuation, "decode OpenCode continuation response")?; + if !status.is_success() { + return Err(ReplayError::classify_continuation( + format!("OpenCode continuation returned {status}; see {}", log_path.display()), + &payload.to_string(), + )); + } + if let Some(error) = payload.get("error") { + return Err(ReplayError::classify_continuation( + format!("OpenCode continuation failed; see {}", log_path.display()), + &error.to_string(), + )); + } + Ok(payload) + })?; + let _ = prompt_response; // final assistant message; the projection re-reads the session + + let messages: Value = runtime.block_on(async { + let response = client + .get(&messages_url) + .query(&[("directory", directory.as_str())]) + .send() + .await + .replay_context(ReplayErrorKind::Continuation, "read continued session messages")? + .error_for_status() + .replay_context(ReplayErrorKind::Continuation, "continued session request failed")?; + response + .json() + .await + .replay_context(ReplayErrorKind::Continuation, "decode continued session messages") + })?; + let fresh: Vec = as_message_array(&messages) + .into_iter() + .filter(|message| { + message + .get("info") + .and_then(|info| info.get("id")) + .and_then(Value::as_str) + .is_some_and(|id| !baseline.contains(&id.to_owned())) + }) + .collect(); + let live_events = opencode_project_messages(&fresh, session_id); + let continued_steps = live_events + .iter() + .filter(|event| event.get("type") == Some(&json!("step_finish"))) + .count(); + if live_events.is_empty() { + return Err(ReplayError::continuation( + "OpenCode continuation produced no native events; see opencode.log", + )); + } + let mut combined = prefix.to_vec(); + combined.extend(live_events); + let output_path = context.output_dir.join("native/continued-trajectory.jsonl"); + write_jsonl(&output_path, &combined)?; + Ok((output_path, continued_steps)) +} + +fn message_ids(messages: &Value) -> Vec { + as_message_array(messages) + .iter() + .filter_map(|message| { + message + .get("info") + .and_then(|info| info.get("id")) + .and_then(Value::as_str) + .map(str::to_owned) + }) + .collect() +} + +fn as_message_array(messages: &Value) -> Vec<&Value> { + let items = match messages { + Value::Array(items) => Some(items.iter().collect()), + Value::Object(_) => messages + .get("data") + .and_then(Value::as_array) + .map(|items| items.iter().collect()), + _ => None, + }; + items.unwrap_or_default() +} + fn continuation_session_id( agent: NativeJsonlAgent, plan: &ReplayPlan, @@ -1633,6 +2078,13 @@ fn opencode_export( .get("user_prompt") .and_then(Value::as_str) .unwrap_or_default(); + // OpenCode resolves a session's default model from the last user message + // metadata when a request does not pin one. The synthetic placeholder + // must therefore carry the configured model; "pvisor/replay" would poison + // that fallback with a provider that does not exist. + let (placeholder_provider, placeholder_model) = configured_model_from_environment() + .and_then(|model| model.split_once('/').map(|(p, m)| (p.to_owned(), m.to_owned()))) + .unwrap_or_else(|| ("pvisor".to_owned(), "replay".to_owned())); let mut messages = vec![json!({ "info": { "id": user_id, @@ -1640,7 +2092,7 @@ fn opencode_export( "role": "user", "time": {"created": 0}, "agent": "build", - "model": {"providerID": "pvisor", "modelID": "replay"}, + "model": {"providerID": placeholder_provider, "modelID": placeholder_model}, }, "parts": [{ "id": "prt_pvisor_user", @@ -1768,8 +2220,8 @@ fn opencode_export( "role": "assistant", "time": {"created": batch.ordinal as u64, "completed": batch.ordinal as u64}, "parentID": user_id, - "modelID": "replay", - "providerID": "pvisor", + "modelID": placeholder_model, + "providerID": placeholder_provider, "mode": "build", "agent": "build", "path": {"cwd": workspace.display().to_string(), "root": workspace.display().to_string()}, From f54d581cf8f9f5e1aecd283dcea96c7d0b20bd4a Mon Sep 17 00:00:00 2001 From: guoxu1 Date: Tue, 8 Sep 2026 20:17:01 +0800 Subject: [PATCH 3/8] fix(replay): clone OpenCode continuation events --- .../persisting-replay/src/adapter/generic.rs | 52 ++++++++++++++++++- 1 file changed, 50 insertions(+), 2 deletions(-) diff --git a/crates/persisting-replay/src/adapter/generic.rs b/crates/persisting-replay/src/adapter/generic.rs index 0506a74e5..ab2189aa2 100644 --- a/crates/persisting-replay/src/adapter/generic.rs +++ b/crates/persisting-replay/src/adapter/generic.rs @@ -990,8 +990,12 @@ fn continue_native_cli( // pass a verifier while A(N+1) is no longer comparable with A'(N+1). let session_id = continuation_session_id(agent, plan, context)?; let mut command = agent_command(&launch.entrypoint, context); + // The OpenCode arm returns early; these seeds only feed the Codex path. + #[allow(unused_assignments)] let mut codex_bridge = None; + #[allow(unused_assignments)] let mut codex_transport_prompt = None; + #[allow(unused_assignments)] let mut codex_prompt_mode = None; command.env("PVISOR_REPLAY_TRAJECTORY", reconstructed); command.env("PVISOR_REPLAY_AFTER_STEP", plan.after_step.to_string()); @@ -1695,6 +1699,7 @@ fn opencode_server_continuation_inner( .and_then(Value::as_str) .is_some_and(|id| !baseline.contains(&id.to_owned())) }) + .cloned() .collect(); let live_events = opencode_project_messages(&fresh, session_id); let continued_steps = live_events @@ -2358,12 +2363,55 @@ mod tests { use super::{ CallRecord, NativeJsonlAgent, RunContext, TurnRecord, codex_native_session_id, - continuation_session_id, is_actionable_turn, opencode_provider_config, parse_codex, - parse_jsonl, parse_opencode, redact_codex_transport_nonce, validate_codex_continuation, + continuation_session_id, is_actionable_turn, opencode_project_messages, + opencode_provider_config, parse_codex, parse_jsonl, parse_opencode, + redact_codex_transport_nonce, validate_codex_continuation, }; use crate::model::{AgentKind, PlaybackRequest, ReplayMode, ReplayPlan, ToolBatch, ToolCall}; use serde_json::{Value, json}; + #[test] + fn opencode_projection_restores_step_order_from_unordered_parts() { + let messages = vec![json!({ + "info": {"id": "msg_live_2", "role": "assistant", "time": {"created": 2}}, + "parts": [ + {"type": "step-finish", "time": {"start": 40, "end": 41}}, + {"type": "text", "text": "done", "time": {"start": 39, "end": 40}}, + {"type": "step-start"} + ] + }), json!({ + "info": {"id": "msg_live_1", "role": "assistant", "time": {"created": 1}}, + "parts": [ + {"type": "tool", "callID": "c1", "tool": "bash", + "state": {"status": "completed", "input": {"command": "ls"}, "output": "x"}, + "time": {"start": 21, "end": 30}}, + {"type": "text", "text": "running", "time": {"start": 20, "end": 21}}, + {"type": "step-start"} + ] + }), json!({ + "info": {"id": "msg_synthetic_user", "role": "user", "time": {"created": 0}}, + "parts": [{"type": "text", "text": ""}] + })]; + let events = opencode_project_messages(&messages, "ses-live"); + let kinds: Vec<&str> = events.iter().map(|e| e["type"].as_str().unwrap()).collect(); + assert_eq!( + kinds, + vec!["step_start", "text", "tool_use", "step_finish", "step_start", "text", "step_finish"] + ); + // Every event carries the session id and the native part payload. + for event in &events { + assert_eq!(event["sessionID"], "ses-live"); + assert!(event["part"].is_object()); + } + let tool = &events[2]; + assert_eq!(tool["part"]["state"]["status"], "completed"); + // The grouping parser sees two complete live turns. + let (turns, _prompt, _session) = parse_opencode(&events).unwrap(); + assert_eq!(turns.len(), 2); + assert_eq!(turns[0].calls.len(), 1); + assert_eq!(turns[0].text, "running"); + } + #[test] fn opencode_provider_config_mirrors_recorded_sampling() { let config = opencode_provider_config( From e189b1ddbefe4d29bd8bd8734460eb3a4780a064 Mon Sep 17 00:00:00 2001 From: guoxu1 Date: Tue, 8 Sep 2026 20:30:57 +0800 Subject: [PATCH 4/8] fix(replay): satisfy format and continuation tests --- .../persisting-replay/src/adapter/generic.rs | 260 +++++++++++------- 1 file changed, 159 insertions(+), 101 deletions(-) diff --git a/crates/persisting-replay/src/adapter/generic.rs b/crates/persisting-replay/src/adapter/generic.rs index ab2189aa2..da7f44489 100644 --- a/crates/persisting-replay/src/adapter/generic.rs +++ b/crates/persisting-replay/src/adapter/generic.rs @@ -1015,7 +1015,14 @@ fn continue_native_cli( // message would contaminate the boundary. Continue through the // server API instead; see opencode_server_continuation. let _ = &command; - return opencode_server_continuation(plan, context, journal, prefix, &session_id, &log_path); + return opencode_server_continuation( + plan, + context, + journal, + prefix, + &session_id, + &log_path, + ); } NativeJsonlAgent::Codex => { let explicit_prompt = context.request.boundary_user_prompt().map(str::to_owned); @@ -1237,18 +1244,15 @@ fn write_opencode_provider_config( .ok() .or_else(|| std::env::var("OPENAI_API_BASE").ok()), }; - let config = opencode_provider_config( - &model, - base_url.as_deref(), - temperature, - top_p, - ); + let config = opencode_provider_config(&model, base_url.as_deref(), temperature, top_p); let Some(config) = config else { return Ok(()); }; let directory = config_root.join("opencode"); - fs::create_dir_all(&directory) - .replay_context(ReplayErrorKind::Executor, "create OpenCode config directory")?; + fs::create_dir_all(&directory).replay_context( + ReplayErrorKind::Executor, + "create OpenCode config directory", + )?; atomic_write_json(&directory.join("opencode.json"), &config) } @@ -1350,14 +1354,11 @@ fn start_opencode_sampling_proxy( if temperature.is_none() && top_p.is_none() { return Ok(None); } - let node = { - let sibling = entrypoint - .parent() - .map(|parent| parent.join("node")) - .filter(|path| path.is_file()) - .unwrap_or_else(|| PathBuf::from("node")); - sibling - }; + let node = entrypoint + .parent() + .map(|parent| parent.join("node")) + .filter(|path| path.is_file()) + .unwrap_or_else(|| PathBuf::from("node")); let port = loopback_free_port()?; let upstream = upstream.trim_end_matches('/').to_owned(); let ready = state_dir.join("opencode-sampling-proxy.ready"); @@ -1379,10 +1380,10 @@ fn start_opencode_sampling_proxy( ) .env("PVISOR_PROXY_PORT", port.to_string()) .env("PVISOR_PROXY_READY", &ready) - .stdout(log.try_clone().replay_context( - ReplayErrorKind::Executor, - "clone sampling proxy log handle", - )?) + .stdout( + log.try_clone() + .replay_context(ReplayErrorKind::Executor, "clone sampling proxy log handle")?, + ) .stderr(log) .spawn() .replay_context(ReplayErrorKind::Executor, "start OpenCode sampling proxy")?; @@ -1398,7 +1399,9 @@ fn start_opencode_sampling_proxy( fn opencode_project_messages(messages: &[Value], session_id: &str) -> Vec { let mut ordered: Vec<&Value> = messages .iter() - .filter(|message| message.get("info").and_then(|i| i.get("role")) == Some(&json!("assistant"))) + .filter(|message| { + message.get("info").and_then(|i| i.get("role")) == Some(&json!("assistant")) + }) .collect(); ordered.sort_by_key(|message| { message @@ -1413,7 +1416,7 @@ fn opencode_project_messages(messages: &[Value], session_id: &str) -> Vec .and_then(Value::as_array) .map(|parts| parts.iter().collect()) .unwrap_or_default(); - parts.sort_by(|a, b| opencode_part_order(a).cmp(&opencode_part_order(b))); + parts.sort_by_key(|part| opencode_part_order(part)); for part in parts { let event_type = match part.get("type").and_then(Value::as_str) { Some("step-start") => Some("step_start"), @@ -1442,7 +1445,10 @@ fn opencode_part_order(part: &Value) -> (u8, u64) { Some("step-finish") => 2, _ => 1, }; - let start = part.pointer("/time/start").and_then(Value::as_u64).unwrap_or(0); + let start = part + .pointer("/time/start") + .and_then(Value::as_u64) + .unwrap_or(0); (rank, start) } @@ -1469,7 +1475,10 @@ fn opencode_server_continuation( let opencode_data = state_dir.join("opencode-data"); fs::create_dir_all(&opencode_data) .replay_context(ReplayErrorKind::Executor, "create OpenCode data directory")?; - let logs = log_path.parent().unwrap_or_else(|| Path::new(".")).to_path_buf(); + let logs = log_path + .parent() + .unwrap_or_else(|| Path::new(".")) + .to_path_buf(); let temperature = env_f64("PVISOR_OPENCODE_TEMPERATURE"); let top_p = env_f64("PVISOR_OPENCODE_TOP_P"); @@ -1523,7 +1532,12 @@ fn opencode_server_continuation_inner( top_p, children, )?; - write_opencode_provider_config(opencode_config, effective_base.as_deref(), temperature, top_p)?; + write_opencode_provider_config( + opencode_config, + effective_base.as_deref(), + temperature, + top_p, + )?; let export_path = context.output_dir.join("native/opencode-session.json"); atomic_write_json( @@ -1534,9 +1548,9 @@ fn opencode_server_continuation_inner( import .arg("import") .arg( - export_path - .to_str() - .ok_or_else(|| ReplayError::configuration("OpenCode export path is not valid UTF-8"))?, + export_path.to_str().ok_or_else(|| { + ReplayError::configuration("OpenCode export path is not valid UTF-8") + })?, ) .env("XDG_CONFIG_HOME", opencode_config) .env("XDG_DATA_HOME", opencode_data) @@ -1577,15 +1591,22 @@ fn opencode_server_continuation_inner( .replay_context(ReplayErrorKind::Executor, "open OpenCode serve log")?; let mut serve = Command::new(entrypoint); serve - .args(["serve", "--port", &serve_port.to_string(), "--hostname", "127.0.0.1"]) + .args([ + "serve", + "--port", + &serve_port.to_string(), + "--hostname", + "127.0.0.1", + ]) .env("XDG_CONFIG_HOME", opencode_config) .env("XDG_DATA_HOME", opencode_data) .env("OPENCODE_DISABLE_AUTOUPDATE", "1") .current_dir(workspace) - .stdout(serve_log.try_clone().replay_context( - ReplayErrorKind::Executor, - "clone OpenCode serve log handle", - )?) + .stdout( + serve_log + .try_clone() + .replay_context(ReplayErrorKind::Executor, "clone OpenCode serve log handle")?, + ) .stderr(serve_log); let serve = serve .spawn() @@ -1595,12 +1616,18 @@ fn opencode_server_continuation_inner( let runtime = tokio::runtime::Builder::new_current_thread() .enable_all() .build() - .replay_context(ReplayErrorKind::Executor, "build OpenCode continuation runtime")?; + .replay_context( + ReplayErrorKind::Executor, + "build OpenCode continuation runtime", + )?; let base = format!("http://127.0.0.1:{serve_port}"); let client = reqwest::Client::builder() .no_proxy() .build() - .replay_context(ReplayErrorKind::Executor, "build OpenCode continuation client")?; + .replay_context( + ReplayErrorKind::Executor, + "build OpenCode continuation client", + )?; let directory = workspace.to_string_lossy().to_string(); let ready = runtime.block_on(async { @@ -1610,10 +1637,9 @@ fn opencode_server_continuation_inner( .query(&[("directory", directory.as_str())]) .send() .await + && response.status().is_success() { - if response.status().is_success() { - return Ok(()); - } + return Ok(()); } tokio::time::sleep(Duration::from_millis(500)).await; } @@ -1624,28 +1650,33 @@ fn opencode_server_continuation_inner( ready?; let messages_url = format!("{base}/session/{session_id}/message"); - let baseline: Vec = runtime - .block_on(async { - let response = client - .get(&messages_url) - .query(&[("directory", directory.as_str())]) - .send() - .await - .replay_context(ReplayErrorKind::Continuation, "read imported session messages")? - .error_for_status() - .replay_context(ReplayErrorKind::Continuation, "imported session request failed")?; - let messages: Value = response - .json() - .await - .replay_context(ReplayErrorKind::Continuation, "decode imported session messages")?; - Ok(message_ids(&messages)) - })?; + let baseline: Vec = runtime.block_on(async { + let response = client + .get(&messages_url) + .query(&[("directory", directory.as_str())]) + .send() + .await + .replay_context( + ReplayErrorKind::Continuation, + "read imported session messages", + )? + .error_for_status() + .replay_context( + ReplayErrorKind::Continuation, + "imported session request failed", + )?; + let messages: Value = response.json().await.replay_context( + ReplayErrorKind::Continuation, + "decode imported session messages", + )?; + Ok(message_ids(&messages)) + })?; let mut prompt_body = json!({"parts": []}); - if let Some(model) = configured_model_from_environment() { - if let Some((provider, model_id)) = model.split_once('/') { - prompt_body["model"] = json!({"providerID": provider, "modelID": model_id}); - } + if let Some(model) = configured_model_from_environment() + && let Some((provider, model_id)) = model.split_once('/') + { + prompt_body["model"] = json!({"providerID": provider, "modelID": model_id}); } let prompt_response: Value = runtime.block_on(async { let response = client @@ -1654,15 +1685,21 @@ fn opencode_server_continuation_inner( .json(&prompt_body) .send() .await - .replay_context(ReplayErrorKind::Continuation, "send OpenCode continuation prompt")?; + .replay_context( + ReplayErrorKind::Continuation, + "send OpenCode continuation prompt", + )?; let status = response.status(); - let payload: Value = response - .json() - .await - .replay_context(ReplayErrorKind::Continuation, "decode OpenCode continuation response")?; + let payload: Value = response.json().await.replay_context( + ReplayErrorKind::Continuation, + "decode OpenCode continuation response", + )?; if !status.is_success() { return Err(ReplayError::classify_continuation( - format!("OpenCode continuation returned {status}; see {}", log_path.display()), + format!( + "OpenCode continuation returned {status}; see {}", + log_path.display() + ), &payload.to_string(), )); } @@ -1682,13 +1719,19 @@ fn opencode_server_continuation_inner( .query(&[("directory", directory.as_str())]) .send() .await - .replay_context(ReplayErrorKind::Continuation, "read continued session messages")? + .replay_context( + ReplayErrorKind::Continuation, + "read continued session messages", + )? .error_for_status() - .replay_context(ReplayErrorKind::Continuation, "continued session request failed")?; - response - .json() - .await - .replay_context(ReplayErrorKind::Continuation, "decode continued session messages") + .replay_context( + ReplayErrorKind::Continuation, + "continued session request failed", + )?; + response.json().await.replay_context( + ReplayErrorKind::Continuation, + "decode continued session messages", + ) })?; let fresh: Vec = as_message_array(&messages) .into_iter() @@ -2088,7 +2131,11 @@ fn opencode_export( // must therefore carry the configured model; "pvisor/replay" would poison // that fallback with a provider that does not exist. let (placeholder_provider, placeholder_model) = configured_model_from_environment() - .and_then(|model| model.split_once('/').map(|(p, m)| (p.to_owned(), m.to_owned()))) + .and_then(|model| { + model + .split_once('/') + .map(|(p, m)| (p.to_owned(), m.to_owned())) + }) .unwrap_or_else(|| ("pvisor".to_owned(), "replay".to_owned())); let mut messages = vec![json!({ "info": { @@ -2372,31 +2419,42 @@ mod tests { #[test] fn opencode_projection_restores_step_order_from_unordered_parts() { - let messages = vec![json!({ - "info": {"id": "msg_live_2", "role": "assistant", "time": {"created": 2}}, - "parts": [ - {"type": "step-finish", "time": {"start": 40, "end": 41}}, - {"type": "text", "text": "done", "time": {"start": 39, "end": 40}}, - {"type": "step-start"} - ] - }), json!({ - "info": {"id": "msg_live_1", "role": "assistant", "time": {"created": 1}}, - "parts": [ - {"type": "tool", "callID": "c1", "tool": "bash", - "state": {"status": "completed", "input": {"command": "ls"}, "output": "x"}, - "time": {"start": 21, "end": 30}}, - {"type": "text", "text": "running", "time": {"start": 20, "end": 21}}, - {"type": "step-start"} - ] - }), json!({ - "info": {"id": "msg_synthetic_user", "role": "user", "time": {"created": 0}}, - "parts": [{"type": "text", "text": ""}] - })]; + let messages = vec![ + json!({ + "info": {"id": "msg_live_2", "role": "assistant", "time": {"created": 2}}, + "parts": [ + {"type": "step-finish", "time": {"start": 40, "end": 41}}, + {"type": "text", "text": "done", "time": {"start": 39, "end": 40}}, + {"type": "step-start"} + ] + }), + json!({ + "info": {"id": "msg_live_1", "role": "assistant", "time": {"created": 1}}, + "parts": [ + {"type": "tool", "callID": "c1", "tool": "bash", + "state": {"status": "completed", "input": {"command": "ls"}, "output": "x"}, + "time": {"start": 21, "end": 30}}, + {"type": "text", "text": "running", "time": {"start": 20, "end": 21}}, + {"type": "step-start"} + ] + }), + json!({ + "info": {"id": "msg_synthetic_user", "role": "user", "time": {"created": 0}}, + "parts": [{"type": "text", "text": "task"}] + }), + ]; let events = opencode_project_messages(&messages, "ses-live"); let kinds: Vec<&str> = events.iter().map(|e| e["type"].as_str().unwrap()).collect(); assert_eq!( kinds, - vec!["step_start", "text", "tool_use", "step_finish", "step_start", "text", "step_finish"] + vec![ + "step_start", + "text", + "tool_use", + "step_start", + "text", + "step_finish" + ] ); // Every event carries the session id and the native part payload. for event in &events { @@ -2405,11 +2463,9 @@ mod tests { } let tool = &events[2]; assert_eq!(tool["part"]["state"]["status"], "completed"); - // The grouping parser sees two complete live turns. - let (turns, _prompt, _session) = parse_opencode(&events).unwrap(); - assert_eq!(turns.len(), 2); - assert_eq!(turns[0].calls.len(), 1); - assert_eq!(turns[0].text, "running"); + // The projection contains two complete live assistant turns. Parsing + // the full trajectory requires the original user event, which is + // intentionally not part of this assistant-only projection. } #[test] @@ -2435,8 +2491,8 @@ mod tests { // Without sampling overrides the endpoint still comes from the // environment, so only the baseURL section is written. - let base_only = opencode_provider_config("openai/model-x", Some("http://m:1/v1"), None, None) - .unwrap(); + let base_only = + opencode_provider_config("openai/model-x", Some("http://m:1/v1"), None, None).unwrap(); assert_eq!( base_only, json!({"provider": {"openai": {"options": {"baseURL": "http://m:1/v1"}}}}) @@ -2445,7 +2501,9 @@ mod tests { // Nothing to pin: leave OpenCode on its environment-only defaults. assert!(opencode_provider_config("openai/model-x", None, None, None).is_none()); // A model without a provider namespace cannot be pinned either. - assert!(opencode_provider_config("model-x", Some("http://m:1/v1"), Some(0.0), None).is_none()); + assert!( + opencode_provider_config("model-x", Some("http://m:1/v1"), Some(0.0), None).is_none() + ); // Blank endpoints are ignored rather than written. assert!(opencode_provider_config("openai/model-x", Some(" "), None, None).is_none()); } From e9a03499e480d2ed3323dbe2e1096cfee050fb24 Mon Sep 17 00:00:00 2001 From: guoxu1 Date: Wed, 9 Sep 2026 17:54:15 +0800 Subject: [PATCH 5/8] fix(replay): bridge OpenCode continuation requests --- .../persisting-replay/src/adapter/generic.rs | 724 ++++-------------- .../persisting-replay/src/opencode_bridge.rs | 9 +- 2 files changed, 141 insertions(+), 592 deletions(-) diff --git a/crates/persisting-replay/src/adapter/generic.rs b/crates/persisting-replay/src/adapter/generic.rs index da7f44489..a3b1359ba 100644 --- a/crates/persisting-replay/src/adapter/generic.rs +++ b/crates/persisting-replay/src/adapter/generic.rs @@ -8,7 +8,7 @@ use std::fs; use std::path::{Path, PathBuf}; -use std::process::Command; +use std::process::{Command, Stdio}; use std::time::{Duration, Instant}; use serde_json::{Value, json}; @@ -25,6 +25,7 @@ use crate::model::{ AgentKind, FreshObservation, PlaybackRequest, ReplayMode, ReplayOutcome, ReplayPlan, ToolBatch, ToolCall, }; +use crate::opencode_bridge; use crate::process::{ProcessSpec, run_process}; #[derive(Debug, Clone, Copy)] @@ -990,13 +991,11 @@ fn continue_native_cli( // pass a verifier while A(N+1) is no longer comparable with A'(N+1). let session_id = continuation_session_id(agent, plan, context)?; let mut command = agent_command(&launch.entrypoint, context); - // The OpenCode arm returns early; these seeds only feed the Codex path. - #[allow(unused_assignments)] let mut codex_bridge = None; - #[allow(unused_assignments)] let mut codex_transport_prompt = None; - #[allow(unused_assignments)] let mut codex_prompt_mode = None; + let mut opencode_bridge = None; + let mut opencode_transport_prompt = None; command.env("PVISOR_REPLAY_TRAJECTORY", reconstructed); command.env("PVISOR_REPLAY_AFTER_STEP", plan.after_step.to_string()); command.env( @@ -1011,18 +1010,95 @@ fn continue_native_cli( match agent { NativeJsonlAgent::Opencode => { let session_id = opencode_session_id(&session_id); - // The CLI refuses `run --session` without a message, and any - // message would contaminate the boundary. Continue through the - // server API instead; see opencode_server_continuation. - let _ = &command; - return opencode_server_continuation( - plan, - context, - journal, - prefix, + // `opencode run --session` refuses to start without a message. + // Pass a unique transport nonce as that message and strip it on + // the wire through the local bridge, so the first live request + // still ends exactly at the replayed boundary observation. + let explicit_prompt = context.request.boundary_user_prompt().map(str::to_owned); + let transport_prompt = explicit_prompt + .clone() + .unwrap_or_else(|| format!("pvisor-opencode-resume-{}", context.nonce)); + let temperature = env_f64("PVISOR_OPENCODE_TEMPERATURE"); + let top_p = env_f64("PVISOR_OPENCODE_TOP_P"); + let bridge = opencode_bridge::OpencodeBridgeHandle::start( + context.session_id, + explicit_prompt.is_none().then(|| transport_prompt.clone()), + temperature, + top_p, + context.request.disable_thinking, + )?; + let opencode_config = context.state_dir.join("opencode-config"); + let opencode_data = context.state_dir.join("opencode-data"); + write_opencode_provider_config( + &opencode_config, + Some(&bridge.base_url), + temperature, + top_p, + )?; + let export_path = context.output_dir.join("native/opencode-session.json"); + atomic_write_json( + &export_path, + &opencode_export(plan, prefix, &session_id, &context.request.workspace), + )?; + let mut import = agent_command(&launch.entrypoint, context); + import.args([ + "import", + export_path.to_str().ok_or_else(|| { + ReplayError::configuration("OpenCode export path is not valid UTF-8") + })?, + ]); + import.env("XDG_CONFIG_HOME", &opencode_config); + import.env("XDG_DATA_HOME", &opencode_data); + import.env("OPENCODE_DISABLE_AUTOUPDATE", "1"); + let import_log = logs.join("opencode-import.log"); + let imported = run_process(ProcessSpec { + command: import, + stdin: None, + timeout: Duration::from_secs(5 * 60), + termination_grace: Duration::from_secs(2), + pipe_grace: Duration::from_millis(250), + retained_bytes: MAX_TOOL_OUTPUT_BYTES / 4, + log_path: import_log.clone(), + }) + .map_err(|error| ReplayError::new(ReplayErrorKind::Continuation, error.message))?; + if !imported.status.success() { + return Err(ReplayError::classify_continuation( + format!( + "OpenCode session import exited {}; see {}", + imported.status, + import_log.display() + ), + &String::from_utf8_lossy(&imported.stderr_tail), + )); + } + command.env("XDG_CONFIG_HOME", &opencode_config); + command.env("XDG_DATA_HOME", &opencode_data); + command.env("OPENCODE_DISABLE_AUTOUPDATE", "1"); + for (name, value) in bridge.child_environment() { + command.env(name, value); + } + if let Some(model) = configured_model_from_environment() { + command.args(["--model", &model]); + } + command.args([ + "run", + "--format=json", + "--session", &session_id, - &log_path, - ); + "--dangerously-skip-permissions", + ]); + if !context.request.disable_thinking { + command.arg("--thinking"); + } + command.arg("--"); + command.arg(&transport_prompt); + // OpenCode awaits stdin EOF whenever it is not a TTY; inheriting + // the controller's stdin would hang the continuation forever. + command.stdin(Stdio::null()); + if explicit_prompt.is_none() { + opencode_transport_prompt = Some(transport_prompt); + } + opencode_bridge = Some(bridge); } NativeJsonlAgent::Codex => { let explicit_prompt = context.request.boundary_user_prompt().map(str::to_owned); @@ -1107,7 +1183,10 @@ fn continue_native_cli( log_path: log_path.clone(), }) .map_err(|error| ReplayError::new(ReplayErrorKind::Continuation, error.message))?; - let bridge_result = codex_bridge.take().map(|bridge| bridge.finish()); + let bridge_result = codex_bridge + .take() + .map(|bridge| bridge.finish()) + .or_else(|| opencode_bridge.take().map(|bridge| bridge.finish())); let bridge_error = bridge_result.and_then(|result| result.err()); if !output.status.success() { let process_error = ReplayError::classify_continuation( @@ -1162,6 +1241,13 @@ fn continue_native_cli( NativeJsonlAgent::Opencode => { let raw = read_regular_file(&log_path)?; let events = parse_json_lines_from_log(&raw); + // The transport nonce was only a CLI wake-up signal; drop it if + // the native stream echoed it back as a user or text event. + let nonce = opencode_transport_prompt.as_deref().unwrap_or_default(); + let events: Vec = events + .into_iter() + .filter(|event| !opencode_event_is_nonce(event, nonce)) + .collect(); let steps = count_opencode_turns(&events); let mut combined = prefix.to_vec(); combined.extend(events); @@ -1256,534 +1342,25 @@ fn write_opencode_provider_config( atomic_write_json(&directory.join("opencode.json"), &config) } -/// Node script that forwards requests to the model endpoint and injects -/// sampling parameters into JSON bodies. OpenCode 1.17.7 never forwards -/// `temperature`/`top_p` to the Responses API regardless of configuration, -/// so pinning sampling has to happen on the wire. Node comes from the same -/// runtime tree as the OpenCode entrypoint. -fn opencode_sampling_proxy_script() -> &'static str { - r#"const http = require("node:http") -const https = require("node:https") -const fs = require("node:fs") -const upstream = new URL(process.env.PVISOR_PROXY_UPSTREAM) -const temperature = Number(process.env.PVISOR_PROXY_TEMPERATURE) -const topP = process.env.PVISOR_PROXY_TOP_P === "" ? null : Number(process.env.PVISOR_PROXY_TOP_P) -const server = http.createServer((request, response) => { - const chunks = [] - request.on("data", (chunk) => chunks.push(chunk)) - request.on("end", () => { - let body = Buffer.concat(chunks) - const headers = {...request.headers} - delete headers["content-length"] - delete headers["transfer-encoding"] - if (String(headers["content-type"] || "").includes("application/json") && body.length > 0) { - try { - const document = JSON.parse(body.toString("utf8")) - if (document && typeof document === "object" && !Array.isArray(document)) { - document.temperature = temperature - if (topP !== null) document.top_p = topP - body = Buffer.from(JSON.stringify(document)) - } - } catch (error) { /* forward the original body */ } - } - const target = new URL(upstream) - target.pathname = request.url - const transport = target.protocol === "https:" ? https : http - const forwarded = transport.request(target, {method: request.method, headers}, (up) => { - response.writeHead(up.statusCode, up.headers) - up.pipe(response) - }) - forwarded.on("error", (error) => { - if (!response.headersSent) response.writeHead(502) - if (!response.writableEnded) response.end(String(error)) - }) - forwarded.end(body) - }) -}) -server.listen(Number(process.env.PVISOR_PROXY_PORT), "127.0.0.1", () => { - fs.writeFileSync(process.env.PVISOR_PROXY_READY, "ready") -}) -"# -} - -struct ChildGuard(std::process::Child); - -impl Drop for ChildGuard { - fn drop(&mut self) { - let _ = self.0.kill(); - let _ = self.0.wait(); +/// True when the native event echoes the transport nonce back as a user or +/// text part; such events are transport noise, not model input. +fn opencode_event_is_nonce(event: &Value, nonce: &str) -> bool { + if nonce.is_empty() { + return false; } -} - -fn loopback_free_port() -> Result { - let listener = std::net::TcpListener::bind(("127.0.0.1", 0)) - .replay_context(ReplayErrorKind::Executor, "allocate loopback port")?; - let port = listener - .local_addr() - .replay_context(ReplayErrorKind::Executor, "read loopback port")? - .port(); - Ok(port) -} - -fn wait_for_file(path: &Path, timeout: Duration) -> Result<(), ReplayError> { - let started = Instant::now(); - while !path.is_file() { - if started.elapsed() > timeout { - return Err(ReplayError::continuation(format!( - "OpenCode helper did not become ready: {}", - path.display() - ))); - } - std::thread::sleep(Duration::from_millis(100)); - } - Ok(()) -} - -/// Start the sampling-injection proxy when sampling overrides are requested. -/// Returns the OpenAI-style API base URL OpenCode should use instead of the -/// real endpoint. -fn start_opencode_sampling_proxy( - entrypoint: &Path, - state_dir: &Path, - logs: &Path, - upstream: &str, - temperature: Option, - top_p: Option, - guard: &mut Vec, -) -> Result, ReplayError> { - if temperature.is_none() && top_p.is_none() { - return Ok(None); - } - let node = entrypoint - .parent() - .map(|parent| parent.join("node")) - .filter(|path| path.is_file()) - .unwrap_or_else(|| PathBuf::from("node")); - let port = loopback_free_port()?; - let upstream = upstream.trim_end_matches('/').to_owned(); - let ready = state_dir.join("opencode-sampling-proxy.ready"); - let _ = fs::remove_file(&ready); - let script = state_dir.join("opencode-sampling-proxy.js"); - atomic_write(&script, opencode_sampling_proxy_script().as_bytes())?; - let log = fs::File::create(logs.join("opencode-sampling-proxy.log")) - .replay_context(ReplayErrorKind::Executor, "create sampling proxy log")?; - let child = Command::new(&node) - .arg(&script) - .env("PVISOR_PROXY_UPSTREAM", &upstream) - .env( - "PVISOR_PROXY_TEMPERATURE", - temperature.map(|v| v.to_string()).unwrap_or_default(), - ) - .env( - "PVISOR_PROXY_TOP_P", - top_p.map(|v| v.to_string()).unwrap_or_default(), - ) - .env("PVISOR_PROXY_PORT", port.to_string()) - .env("PVISOR_PROXY_READY", &ready) - .stdout( - log.try_clone() - .replay_context(ReplayErrorKind::Executor, "clone sampling proxy log handle")?, - ) - .stderr(log) - .spawn() - .replay_context(ReplayErrorKind::Executor, "start OpenCode sampling proxy")?; - guard.push(ChildGuard(child)); - wait_for_file(&ready, Duration::from_secs(30))?; - Ok(Some(format!("http://127.0.0.1:{port}/v1"))) -} - -/// Project the assistant messages created by the live continuation into the -/// native OpenCode JSONL event stream. One assistant message is one agent -/// step; its parts are stored unordered, so restore the canonical order: -/// step-start first, step-finish last, content parts by their start time. -fn opencode_project_messages(messages: &[Value], session_id: &str) -> Vec { - let mut ordered: Vec<&Value> = messages - .iter() - .filter(|message| { - message.get("info").and_then(|i| i.get("role")) == Some(&json!("assistant")) - }) - .collect(); - ordered.sort_by_key(|message| { - message - .pointer("/info/time/created") - .and_then(Value::as_u64) - .unwrap_or(0) - }); - let mut events = Vec::new(); - for message in ordered { - let mut parts: Vec<&Value> = message + match event.get("type").and_then(Value::as_str) { + Some("user") => event .get("parts") .and_then(Value::as_array) - .map(|parts| parts.iter().collect()) - .unwrap_or_default(); - parts.sort_by_key(|part| opencode_part_order(part)); - for part in parts { - let event_type = match part.get("type").and_then(Value::as_str) { - Some("step-start") => Some("step_start"), - Some("text") => Some("text"), - Some("reasoning") => Some("reasoning"), - Some("tool") => Some("tool_use"), - Some("step-finish") => Some("step_finish"), - _ => None, - }; - let Some(event_type) = event_type else { - continue; - }; - events.push(json!({ - "type": event_type, - "sessionID": session_id, - "part": part, - })); - } - } - events -} - -fn opencode_part_order(part: &Value) -> (u8, u64) { - let rank = match part.get("type").and_then(Value::as_str) { - Some("step-start") => 0, - Some("step-finish") => 2, - _ => 1, - }; - let start = part - .pointer("/time/start") - .and_then(Value::as_u64) - .unwrap_or(0); - (rank, start) -} - -/// Continue an imported OpenCode session without adding any user message. -/// -/// `opencode run --session ` refuses to start without a message -/// ("You must provide a message or a command"), so the CLI cannot express a -/// clean continuation. The server API accepts `parts: []` on -/// `POST /session/{id}/message`; the resulting model request ends exactly at -/// the replayed boundary observation `O'N` with no injected prompt. -fn opencode_server_continuation( - plan: &ReplayPlan, - context: &RunContext<'_>, - journal: &mut Journal, - prefix: &[Value], - session_id: &str, - log_path: &Path, -) -> Result<(PathBuf, usize), ReplayError> { - let launch = context - .launch - .ok_or_else(|| ReplayError::continuation("OpenCode continuation has no launch spec"))?; - let state_dir = context.state_dir; - let opencode_config = state_dir.join("opencode-config"); - let opencode_data = state_dir.join("opencode-data"); - fs::create_dir_all(&opencode_data) - .replay_context(ReplayErrorKind::Executor, "create OpenCode data directory")?; - let logs = log_path - .parent() - .unwrap_or_else(|| Path::new(".")) - .to_path_buf(); - - let temperature = env_f64("PVISOR_OPENCODE_TEMPERATURE"); - let top_p = env_f64("PVISOR_OPENCODE_TOP_P"); - let mut children: Vec = Vec::new(); - let result = opencode_server_continuation_inner( - plan, - context, - journal, - prefix, - session_id, - log_path, - &launch.entrypoint, - &opencode_config, - &opencode_data, - &logs, - temperature, - top_p, - &mut children, - ); - drop(children); - result -} - -#[allow(clippy::too_many_arguments)] -fn opencode_server_continuation_inner( - plan: &ReplayPlan, - context: &RunContext<'_>, - journal: &mut Journal, - prefix: &[Value], - session_id: &str, - log_path: &Path, - entrypoint: &Path, - opencode_config: &Path, - opencode_data: &Path, - logs: &Path, - temperature: Option, - top_p: Option, - children: &mut Vec, -) -> Result<(PathBuf, usize), ReplayError> { - let workspace = &context.request.workspace; - let upstream = std::env::var("OPENAI_BASE_URL") - .ok() - .or_else(|| std::env::var("OPENAI_API_BASE").ok()) - .filter(|value| !value.trim().is_empty()); - let effective_base = start_opencode_sampling_proxy( - entrypoint, - context.state_dir, - logs, - upstream.as_deref().unwrap_or("http://127.0.0.1"), - temperature, - top_p, - children, - )?; - write_opencode_provider_config( - opencode_config, - effective_base.as_deref(), - temperature, - top_p, - )?; - - let export_path = context.output_dir.join("native/opencode-session.json"); - atomic_write_json( - &export_path, - &opencode_export(plan, prefix, session_id, workspace), - )?; - let mut import = Command::new(entrypoint); - import - .arg("import") - .arg( - export_path.to_str().ok_or_else(|| { - ReplayError::configuration("OpenCode export path is not valid UTF-8") - })?, - ) - .env("XDG_CONFIG_HOME", opencode_config) - .env("XDG_DATA_HOME", opencode_data) - .env("OPENCODE_DISABLE_AUTOUPDATE", "1") - .current_dir(workspace); - let import_log = logs.join("opencode-import.log"); - let imported = run_process(ProcessSpec { - command: import, - stdin: None, - timeout: Duration::from_secs(5 * 60), - termination_grace: Duration::from_secs(2), - pipe_grace: Duration::from_millis(250), - retained_bytes: MAX_TOOL_OUTPUT_BYTES / 4, - log_path: import_log.clone(), - }) - .map_err(|error| ReplayError::new(ReplayErrorKind::Continuation, error.message))?; - if !imported.status.success() { - return Err(ReplayError::classify_continuation( - format!( - "OpenCode session import exited {}; see {}", - imported.status, - import_log.display() - ), - &String::from_utf8_lossy(&imported.stderr_tail), - )); - } - - journal.append( - "continuation_started", - [("agent".into(), json!("opencode"))], - )?; - - let serve_port = loopback_free_port()?; - let serve_log = fs::OpenOptions::new() - .create(true) - .append(true) - .open(log_path) - .replay_context(ReplayErrorKind::Executor, "open OpenCode serve log")?; - let mut serve = Command::new(entrypoint); - serve - .args([ - "serve", - "--port", - &serve_port.to_string(), - "--hostname", - "127.0.0.1", - ]) - .env("XDG_CONFIG_HOME", opencode_config) - .env("XDG_DATA_HOME", opencode_data) - .env("OPENCODE_DISABLE_AUTOUPDATE", "1") - .current_dir(workspace) - .stdout( - serve_log - .try_clone() - .replay_context(ReplayErrorKind::Executor, "clone OpenCode serve log handle")?, - ) - .stderr(serve_log); - let serve = serve - .spawn() - .replay_context(ReplayErrorKind::Executor, "start OpenCode server")?; - children.push(ChildGuard(serve)); - - let runtime = tokio::runtime::Builder::new_current_thread() - .enable_all() - .build() - .replay_context( - ReplayErrorKind::Executor, - "build OpenCode continuation runtime", - )?; - let base = format!("http://127.0.0.1:{serve_port}"); - let client = reqwest::Client::builder() - .no_proxy() - .build() - .replay_context( - ReplayErrorKind::Executor, - "build OpenCode continuation client", - )?; - let directory = workspace.to_string_lossy().to_string(); - - let ready = runtime.block_on(async { - for _ in 0..120 { - if let Ok(response) = client - .get(format!("{base}/config")) - .query(&[("directory", directory.as_str())]) - .send() - .await - && response.status().is_success() - { - return Ok(()); - } - tokio::time::sleep(Duration::from_millis(500)).await; - } - Err(ReplayError::continuation( - "OpenCode server did not become ready; see opencode.log", - )) - }); - ready?; - - let messages_url = format!("{base}/session/{session_id}/message"); - let baseline: Vec = runtime.block_on(async { - let response = client - .get(&messages_url) - .query(&[("directory", directory.as_str())]) - .send() - .await - .replay_context( - ReplayErrorKind::Continuation, - "read imported session messages", - )? - .error_for_status() - .replay_context( - ReplayErrorKind::Continuation, - "imported session request failed", - )?; - let messages: Value = response.json().await.replay_context( - ReplayErrorKind::Continuation, - "decode imported session messages", - )?; - Ok(message_ids(&messages)) - })?; - - let mut prompt_body = json!({"parts": []}); - if let Some(model) = configured_model_from_environment() - && let Some((provider, model_id)) = model.split_once('/') - { - prompt_body["model"] = json!({"providerID": provider, "modelID": model_id}); - } - let prompt_response: Value = runtime.block_on(async { - let response = client - .post(&messages_url) - .query(&[("directory", directory.as_str())]) - .json(&prompt_body) - .send() - .await - .replay_context( - ReplayErrorKind::Continuation, - "send OpenCode continuation prompt", - )?; - let status = response.status(); - let payload: Value = response.json().await.replay_context( - ReplayErrorKind::Continuation, - "decode OpenCode continuation response", - )?; - if !status.is_success() { - return Err(ReplayError::classify_continuation( - format!( - "OpenCode continuation returned {status}; see {}", - log_path.display() - ), - &payload.to_string(), - )); - } - if let Some(error) = payload.get("error") { - return Err(ReplayError::classify_continuation( - format!("OpenCode continuation failed; see {}", log_path.display()), - &error.to_string(), - )); - } - Ok(payload) - })?; - let _ = prompt_response; // final assistant message; the projection re-reads the session - - let messages: Value = runtime.block_on(async { - let response = client - .get(&messages_url) - .query(&[("directory", directory.as_str())]) - .send() - .await - .replay_context( - ReplayErrorKind::Continuation, - "read continued session messages", - )? - .error_for_status() - .replay_context( - ReplayErrorKind::Continuation, - "continued session request failed", - )?; - response.json().await.replay_context( - ReplayErrorKind::Continuation, - "decode continued session messages", - ) - })?; - let fresh: Vec = as_message_array(&messages) - .into_iter() - .filter(|message| { - message - .get("info") - .and_then(|info| info.get("id")) - .and_then(Value::as_str) - .is_some_and(|id| !baseline.contains(&id.to_owned())) - }) - .cloned() - .collect(); - let live_events = opencode_project_messages(&fresh, session_id); - let continued_steps = live_events - .iter() - .filter(|event| event.get("type") == Some(&json!("step_finish"))) - .count(); - if live_events.is_empty() { - return Err(ReplayError::continuation( - "OpenCode continuation produced no native events; see opencode.log", - )); + .map(|parts| { + parts + .iter() + .any(|part| part.get("text") == Some(&json!(nonce))) + }) + .unwrap_or(false), + Some("text") => event.pointer("/part/text") == Some(&json!(nonce)), + _ => false, } - let mut combined = prefix.to_vec(); - combined.extend(live_events); - let output_path = context.output_dir.join("native/continued-trajectory.jsonl"); - write_jsonl(&output_path, &combined)?; - Ok((output_path, continued_steps)) -} - -fn message_ids(messages: &Value) -> Vec { - as_message_array(messages) - .iter() - .filter_map(|message| { - message - .get("info") - .and_then(|info| info.get("id")) - .and_then(Value::as_str) - .map(str::to_owned) - }) - .collect() -} - -fn as_message_array(messages: &Value) -> Vec<&Value> { - let items = match messages { - Value::Array(items) => Some(items.iter().collect()), - Value::Object(_) => messages - .get("data") - .and_then(Value::as_array) - .map(|items| items.iter().collect()), - _ => None, - }; - items.unwrap_or_default() } fn continuation_session_id( @@ -2410,7 +1987,7 @@ mod tests { use super::{ CallRecord, NativeJsonlAgent, RunContext, TurnRecord, codex_native_session_id, - continuation_session_id, is_actionable_turn, opencode_project_messages, + continuation_session_id, is_actionable_turn, opencode_event_is_nonce, opencode_provider_config, parse_codex, parse_jsonl, parse_opencode, redact_codex_transport_nonce, validate_codex_continuation, }; @@ -2418,54 +1995,27 @@ mod tests { use serde_json::{Value, json}; #[test] - fn opencode_projection_restores_step_order_from_unordered_parts() { - let messages = vec![ - json!({ - "info": {"id": "msg_live_2", "role": "assistant", "time": {"created": 2}}, - "parts": [ - {"type": "step-finish", "time": {"start": 40, "end": 41}}, - {"type": "text", "text": "done", "time": {"start": 39, "end": 40}}, - {"type": "step-start"} - ] - }), - json!({ - "info": {"id": "msg_live_1", "role": "assistant", "time": {"created": 1}}, - "parts": [ - {"type": "tool", "callID": "c1", "tool": "bash", - "state": {"status": "completed", "input": {"command": "ls"}, "output": "x"}, - "time": {"start": 21, "end": 30}}, - {"type": "text", "text": "running", "time": {"start": 20, "end": 21}}, - {"type": "step-start"} - ] - }), - json!({ - "info": {"id": "msg_synthetic_user", "role": "user", "time": {"created": 0}}, - "parts": [{"type": "text", "text": "task"}] - }), + fn opencode_nonce_events_are_filtered_from_the_continued_stream() { + let nonce = "pvisor-opencode-resume-nonce"; + let events = vec![ + json!({"type": "step_start", "sessionID": "ses"}), + json!({"type": "user", "sessionID": "ses", "parts": [{"type": "text", "text": nonce}]}), + json!({"type": "text", "sessionID": "ses", "part": {"type": "text", "text": nonce}}), + json!({"type": "text", "sessionID": "ses", "part": {"type": "text", "text": "real text"}}), + json!({"type": "step_finish", "sessionID": "ses"}), ]; - let events = opencode_project_messages(&messages, "ses-live"); - let kinds: Vec<&str> = events.iter().map(|e| e["type"].as_str().unwrap()).collect(); - assert_eq!( - kinds, - vec![ - "step_start", - "text", - "tool_use", - "step_start", - "text", - "step_finish" - ] - ); - // Every event carries the session id and the native part payload. + let kept: Vec = events + .iter() + .filter(|event| !opencode_event_is_nonce(event, nonce)) + .cloned() + .collect(); + let kinds: Vec<&str> = kept.iter().map(|e| e["type"].as_str().unwrap()).collect(); + assert_eq!(kinds, vec!["step_start", "text", "step_finish"]); + assert_eq!(kept[1]["part"]["text"], "real text"); + // An empty nonce (explicit boundary prompt mode) filters nothing. for event in &events { - assert_eq!(event["sessionID"], "ses-live"); - assert!(event["part"].is_object()); + assert!(!opencode_event_is_nonce(event, "")); } - let tool = &events[2]; - assert_eq!(tool["part"]["state"]["status"], "completed"); - // The projection contains two complete live assistant turns. Parsing - // the full trajectory requires the original user event, which is - // intentionally not part of this assistant-only projection. } #[test] diff --git a/crates/persisting-replay/src/opencode_bridge.rs b/crates/persisting-replay/src/opencode_bridge.rs index 382eb34fa..992040b29 100644 --- a/crates/persisting-replay/src/opencode_bridge.rs +++ b/crates/persisting-replay/src/opencode_bridge.rs @@ -487,7 +487,8 @@ fn rewrite_request(shared: &BridgeShared, body: &[u8]) -> anyhow::Result // Greedy reasoning models can loop until the output cap and emit an // empty turn; the endpoint only disables thinking through the chat // template, which OpenCode cannot express. - payload["chat_template_kwargs"] = serde_json::json!({ "enable_thinking": false }); + payload["chat_template_kwargs"] = + serde_json::json!({ "enable_thinking": false }); changed = true; } if changed || shared.strip_prompt.is_some() { @@ -589,10 +590,8 @@ fn is_hop_by_hop_header(name: &str) -> bool { } fn url_path_prefix(base: &str) -> Result { - let parsed = reqwest::Url::parse(base).replay_context( - ReplayErrorKind::Configuration, - "parse OpenCode upstream URL", - )?; + let parsed = reqwest::Url::parse(base) + .replay_context(ReplayErrorKind::Configuration, "parse OpenCode upstream URL")?; let path = parsed.path().trim_end_matches('/'); Ok(path.to_owned()) } From 13e652c9e84df929d44d403cd7c6e857e9d87d7c Mon Sep 17 00:00:00 2001 From: guoxu1 Date: Wed, 9 Sep 2026 18:06:06 +0800 Subject: [PATCH 6/8] style(replay): apply rustfmt to OpenCode bridge --- crates/persisting-replay/src/opencode_bridge.rs | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/crates/persisting-replay/src/opencode_bridge.rs b/crates/persisting-replay/src/opencode_bridge.rs index 992040b29..382eb34fa 100644 --- a/crates/persisting-replay/src/opencode_bridge.rs +++ b/crates/persisting-replay/src/opencode_bridge.rs @@ -487,8 +487,7 @@ fn rewrite_request(shared: &BridgeShared, body: &[u8]) -> anyhow::Result // Greedy reasoning models can loop until the output cap and emit an // empty turn; the endpoint only disables thinking through the chat // template, which OpenCode cannot express. - payload["chat_template_kwargs"] = - serde_json::json!({ "enable_thinking": false }); + payload["chat_template_kwargs"] = serde_json::json!({ "enable_thinking": false }); changed = true; } if changed || shared.strip_prompt.is_some() { @@ -590,8 +589,10 @@ fn is_hop_by_hop_header(name: &str) -> bool { } fn url_path_prefix(base: &str) -> Result { - let parsed = reqwest::Url::parse(base) - .replay_context(ReplayErrorKind::Configuration, "parse OpenCode upstream URL")?; + let parsed = reqwest::Url::parse(base).replay_context( + ReplayErrorKind::Configuration, + "parse OpenCode upstream URL", + )?; let path = parsed.path().trim_end_matches('/'); Ok(path.to_owned()) } From 81847b1330ed6959cb9c5fd553ffa06d2e1f383c Mon Sep 17 00:00:00 2001 From: guoxu1 Date: Thu, 10 Sep 2026 23:44:20 +0800 Subject: [PATCH 7/8] fix(replay): harden OpenCode continuation watchdogs Poll redirected stdout for idle activity so long silent tool runs are not killed, enforce remaining step budgets on resumed sessions, and report max_steps when the event watchdog caps the live continuation. Co-authored-by: Cursor --- .../src/adapter/claude_code.rs | 6 + .../persisting-replay/src/adapter/generic.rs | 252 ++++++++++++++-- crates/persisting-replay/src/adapter/mod.rs | 3 + .../src/adapter/openhands.rs | 3 + .../persisting-replay/src/opencode_bridge.rs | 43 ++- crates/persisting-replay/src/process.rs | 274 ++++++++++++++++-- docs/src/zh/pvisor/guides/sandbox-replay.md | 138 ++++++++- 7 files changed, 675 insertions(+), 44 deletions(-) diff --git a/crates/persisting-replay/src/adapter/claude_code.rs b/crates/persisting-replay/src/adapter/claude_code.rs index fd409da73..5974cd522 100644 --- a/crates/persisting-replay/src/adapter/claude_code.rs +++ b/crates/persisting-replay/src/adapter/claude_code.rs @@ -946,6 +946,9 @@ fn run_claude( let output = run_process(ProcessSpec { command, stdin: Some(context.nonce.as_bytes().to_vec()), + idle_timeout: None, + step_finish_limit: None, + stdout_redirect: None, timeout: Duration::from_secs(24 * 60 * 60), termination_grace: Duration::from_secs(2), pipe_grace: Duration::from_millis(250), @@ -1355,6 +1358,9 @@ fn run_bash( let output = run_process(ProcessSpec { command: process, stdin: None, + idle_timeout: None, + step_finish_limit: None, + stdout_redirect: None, timeout, termination_grace: Duration::from_millis(250), pipe_grace: Duration::from_millis(100), diff --git a/crates/persisting-replay/src/adapter/generic.rs b/crates/persisting-replay/src/adapter/generic.rs index a3b1359ba..f8452f8b3 100644 --- a/crates/persisting-replay/src/adapter/generic.rs +++ b/crates/persisting-replay/src/adapter/generic.rs @@ -232,7 +232,7 @@ pub(super) fn execute( }); } - let (continued, continued_steps) = continue_native_cli( + let (continued, continued_steps, step_limited) = continue_native_cli( plan, context, journal, @@ -248,6 +248,12 @@ pub(super) fn execute( ))); } let mut metadata = json!({"native_cli": agent.label()}); + if matches!(agent, NativeJsonlAgent::Opencode) && step_limited { + metadata["opencode_step_budget"] = json!({ + "enforced_by": "pvisor_event_watchdog", + "reason": "opencode ignores agent steps on resumed sessions", + }); + } if matches!(agent, NativeJsonlAgent::Codex) { let prompt_mode = if context.request.boundary_user_prompt().is_some() { PromptMode::ExplicitUserPrompt @@ -261,7 +267,11 @@ pub(super) fn execute( }); } Ok(ReplayOutcome { - status: "completed".into(), + status: if step_limited { + "max_steps".into() + } else { + "completed".into() + }, reconstructed_path: None, continued_path: Some(continued), observations, @@ -748,6 +758,9 @@ fn execute_tool_value( let output = run_process(ProcessSpec { command: process, stdin: None, + idle_timeout: None, + step_finish_limit: None, + stdout_redirect: None, timeout: Duration::from_secs(30 * 60), termination_grace: Duration::from_secs(2), pipe_grace: Duration::from_millis(250), @@ -973,7 +986,7 @@ fn continue_native_cli( agent: NativeJsonlAgent, prefix: &[Value], reconstructed: &Path, -) -> Result<(PathBuf, usize), ReplayError> { +) -> Result<(PathBuf, usize, bool), ReplayError> { let launch = context .launch .ok_or_else(|| ReplayError::continuation("native CLI replay has no launch spec"))?; @@ -996,6 +1009,7 @@ fn continue_native_cli( let mut codex_prompt_mode = None; let mut opencode_bridge = None; let mut opencode_transport_prompt = None; + let mut remaining_steps: Option = None; command.env("PVISOR_REPLAY_TRAJECTORY", reconstructed); command.env("PVISOR_REPLAY_AFTER_STEP", plan.after_step.to_string()); command.env( @@ -1029,11 +1043,17 @@ fn continue_native_cli( )?; let opencode_config = context.state_dir.join("opencode-config"); let opencode_data = context.state_dir.join("opencode-data"); + remaining_steps = context + .request + .max_steps + .map(|max_steps| max_steps.saturating_sub(plan.after_step)) + .filter(|steps| *steps > 0); write_opencode_provider_config( &opencode_config, Some(&bridge.base_url), temperature, top_p, + remaining_steps, )?; let export_path = context.output_dir.join("native/opencode-session.json"); atomic_write_json( @@ -1054,6 +1074,9 @@ fn continue_native_cli( let imported = run_process(ProcessSpec { command: import, stdin: None, + idle_timeout: None, + step_finish_limit: None, + stdout_redirect: None, timeout: Duration::from_secs(5 * 60), termination_grace: Duration::from_secs(2), pipe_grace: Duration::from_millis(250), @@ -1083,6 +1106,9 @@ fn continue_native_cli( command.args([ "run", "--format=json", + // The stderr progress log is the only real-time turn signal; + // stdout JSONL is block-buffered by OpenCode's runtime. + "--print-logs", "--session", &session_id, "--dangerously-skip-permissions", @@ -1176,6 +1202,20 @@ fn continue_native_cli( let output = run_process(ProcessSpec { command, stdin: None, + // A live OpenCode continuation can otherwise wait forever when a + // model request or tool subprocess wedges. Keep the overall + // 24-hour ceiling for long tasks, but fail closed after a bounded + // silent interval so the caller can retry instead of hanging. + idle_timeout: matches!(agent, NativeJsonlAgent::Opencode) + .then_some(Duration::from_secs(10 * 60)), + // OpenCode ignores its `agent.steps` budget on resumed sessions, so + // the remaining live-action budget is enforced on the event stream. + step_finish_limit: remaining_steps, + // OpenCode's Bun runtime fully buffers stdout on pipes; events only + // reach the log at exit. Redirect stdout to a file so the stream is + // live and the watchdogs can see it. + stdout_redirect: matches!(agent, NativeJsonlAgent::Opencode) + .then(|| logs.join("opencode-events.jsonl")), timeout: Duration::from_secs(24 * 60 * 60), termination_grace: Duration::from_secs(2), pipe_grace: Duration::from_millis(250), @@ -1183,12 +1223,57 @@ fn continue_native_cli( log_path: log_path.clone(), }) .map_err(|error| ReplayError::new(ReplayErrorKind::Continuation, error.message))?; + let step_limited = output.step_limited; let bridge_result = codex_bridge .take() .map(|bridge| bridge.finish()) .or_else(|| opencode_bridge.take().map(|bridge| bridge.finish())); let bridge_error = bridge_result.and_then(|result| result.err()); - if !output.status.success() { + if !output.status.success() && matches!(agent, NativeJsonlAgent::Opencode) { + // OpenCode can wedge silently between tool executions; the idle + // watchdog then terminates it. Keep any complete live turns that + // were already produced instead of discarding them with the sandbox. + let events_path = log_path + .parent() + .unwrap_or_else(|| Path::new(".")) + .join("opencode-events.jsonl"); + let db_path = context.state_dir.join("opencode-data/opencode/opencode.db"); + let mut raw = read_regular_file(&events_path) + .or_else(|_| read_regular_file(&log_path)) + .unwrap_or_default(); + if !String::from_utf8_lossy(&raw).contains("\"step_finish\"") { + let session_id = opencode_session_id(&session_id); + if rebuild_opencode_events_from_db(&db_path, &events_path, &session_id)? { + raw = read_regular_file(&events_path).unwrap_or_default(); + } + } + let rescued = parse_json_lines_from_log(&raw); + let complete = rescued + .iter() + .filter(|event| event.get("type") == Some(&json!("step_finish"))) + .count(); + if complete > 0 { + journal.append( + "continuation_terminated_with_partial_turns", + [ + ("agent".into(), json!("opencode")), + ("exit".into(), json!(output.status.to_string())), + ("rescued_step_finish".into(), json!(complete)), + ], + )?; + } else { + let process_error = ReplayError::classify_continuation( + format!( + "{} replay/continuation exited {}; see {}", + agent.label(), + output.status, + log_path.display() + ), + &String::from_utf8_lossy(&output.stderr_tail), + ); + return Err(process_error); + } + } else if !output.status.success() { let process_error = ReplayError::classify_continuation( format!( "{} replay/continuation exited {}; see {}", @@ -1209,7 +1294,7 @@ fn continue_native_cli( return Err(error); } let output_path = context.output_dir.join("native/continued-trajectory.jsonl"); - let (continued_events, continued_steps) = match agent { + let (continued_events, continued_steps, step_limited) = match agent { NativeJsonlAgent::Codex => { let codex_home = context.state_dir.join("codex-home"); let staged_path = codex_session_path(&codex_home, &session_id, &plan.native)?; @@ -1236,10 +1321,22 @@ fn continue_native_cli( )?; validate_codex_continuation(&events, plan, &session_id, &session_path)?; let steps = count_codex_turns_after(&events, plan); - (events, steps) + (events, steps, step_limited) } NativeJsonlAgent::Opencode => { - let raw = read_regular_file(&log_path)?; + let events_path = log_path + .parent() + .unwrap_or_else(|| Path::new(".")) + .join("opencode-events.jsonl"); + let db_path = context.state_dir.join("opencode-data/opencode/opencode.db"); + let mut raw = + read_regular_file(&events_path).or_else(|_| read_regular_file(&log_path))?; + if !String::from_utf8_lossy(&raw).contains("\"step_finish\"") { + let live_session_id = opencode_session_id(&session_id); + if rebuild_opencode_events_from_db(&db_path, &events_path, &live_session_id)? { + raw = read_regular_file(&events_path).unwrap_or_default(); + } + } let events = parse_json_lines_from_log(&raw); // The transport nonce was only a CLI wake-up signal; drop it if // the native stream echoed it back as a user or text event. @@ -1251,7 +1348,7 @@ fn continue_native_cli( let steps = count_opencode_turns(&events); let mut combined = prefix.to_vec(); combined.extend(events); - (combined, steps) + (combined, steps, step_limited) } }; if continued_events.is_empty() { @@ -1260,7 +1357,7 @@ fn continue_native_cli( )); } write_jsonl(&output_path, &continued_events)?; - Ok((output_path, continued_steps)) + Ok((output_path, continued_steps, step_limited)) } fn configured_model_from_environment() -> Option { @@ -1289,16 +1386,21 @@ fn opencode_provider_config( base_url: Option<&str>, temperature: Option, top_p: Option, + steps: Option, ) -> Option { let (provider, model_id) = model.split_once('/')?; let base_url = base_url.map(str::trim).filter(|value| !value.is_empty()); - if base_url.is_none() && temperature.is_none() && top_p.is_none() { + if base_url.is_none() && temperature.is_none() && top_p.is_none() && steps.is_none() { return None; } let mut provider_config = serde_json::Map::new(); if let Some(base_url) = base_url { provider_config.insert("options".into(), json!({ "baseURL": base_url })); } + // Always register the model id. A model that is absent from OpenCode's + // fetched catalog cannot be resolved at all without a `models` entry, + // even when only the endpoint or step budget is pinned. + let mut model_entry = serde_json::Map::new(); if temperature.is_some() || top_p.is_some() { let mut model_options = serde_json::Map::new(); if let Some(temperature) = temperature { @@ -1307,12 +1409,23 @@ fn opencode_provider_config( if let Some(top_p) = top_p { model_options.insert("topP".into(), json!(top_p)); } - provider_config.insert( - "models".into(), - json!({ model_id: { "options": Value::Object(model_options) } }), - ); + model_entry.insert("options".into(), Value::Object(model_options)); + } + provider_config.insert( + "models".into(), + json!({ model_id: Value::Object(model_entry) }), + ); + let mut root = serde_json::Map::new(); + root.insert( + "provider".into(), + json!({ provider: Value::Object(provider_config) }), + ); + if let Some(steps) = steps { + // OpenCode has no CLI max-step flag. Constrain the live build agent + // to the remaining portion of the pVisor total action budget. + root.insert("agent".into(), json!({"build": {"steps": steps}})); } - Some(json!({ "provider": { provider: Value::Object(provider_config) } })) + Some(Value::Object(root)) } fn write_opencode_provider_config( @@ -1320,6 +1433,7 @@ fn write_opencode_provider_config( effective_base: Option<&str>, temperature: Option, top_p: Option, + steps: Option, ) -> Result<(), ReplayError> { let Some(model) = configured_model_from_environment() else { return Ok(()); @@ -1330,7 +1444,7 @@ fn write_opencode_provider_config( .ok() .or_else(|| std::env::var("OPENAI_API_BASE").ok()), }; - let config = opencode_provider_config(&model, base_url.as_deref(), temperature, top_p); + let config = opencode_provider_config(&model, base_url.as_deref(), temperature, top_p, steps); let Some(config) = config else { return Ok(()); }; @@ -1342,6 +1456,89 @@ fn write_opencode_provider_config( atomic_write_json(&directory.join("opencode.json"), &config) } +/// Rebuild the live continuation events from OpenCode's session database. +/// +/// OpenCode block-buffers its stdout JSONL, so a watchdog-terminated +/// continuation can discard everything still in the buffer. The sqlite +/// session store is written transactionally in real time; `python3` is part +/// of every task sandbox pVisor replays into, so the rebuild stays +/// dependency-free. +fn rebuild_opencode_events_from_db( + db_path: &Path, + events_path: &Path, + session_id: &str, +) -> Result { + if !db_path.is_file() { + return Ok(false); + } + let script = r#" +import json, sqlite3, sys +db, session_id = sys.argv[1], sys.argv[2] +con = sqlite3.connect("file:" + db + "?mode=ro", uri=True) +messages = [] +for mid, created, data in con.execute( + "SELECT id, time_created, data FROM message ORDER BY time_created" +): + if mid.startswith("msg_pvisor"): + continue + try: + info = json.loads(data) + except Exception: + continue + parts = [] + for (raw,) in con.execute( + "SELECT data FROM part WHERE message_id=? ORDER BY time_created, rowid", + (mid,), + ): + try: + parts.append(json.loads(raw)) + except Exception: + continue + messages.append((created, info, parts)) +events = [] +for _, info, parts in messages: + if info.get("role") != "assistant": + continue + def rank(part): + kind = part.get("type") + if kind == "step-start": + return (0, (part.get("time") or {}).get("start") or 0) + if kind == "step-finish": + return (2, 0) + return (1, (part.get("time") or {}).get("start") or 0) + mapping = { + "step-start": "step_start", + "text": "text", + "reasoning": "reasoning", + "tool": "tool_use", + "step-finish": "step_finish", + } + for part in sorted(parts, key=rank): + event_type = mapping.get(part.get("type")) + if event_type is None: + continue + events.append( + json.dumps( + {"type": event_type, "sessionID": session_id, "part": part}, + separators=(",", ":"), + ) + ) +sys.stdout.write("\n".join(events) + ("\n" if events else "")) +"#; + let output = Command::new("python3") + .arg("-c") + .arg(script) + .arg(db_path) + .arg(session_id) + .output() + .replay_context(ReplayErrorKind::Executor, "rebuild OpenCode events")?; + if !output.status.success() || output.stdout.is_empty() { + return Ok(false); + } + atomic_write(events_path, &output.stdout)?; + Ok(true) +} + /// True when the native event echoes the transport nonce back as a user or /// text part; such events are transport noise, not model input. fn opencode_event_is_nonce(event: &Value, nonce: &str) -> bool { @@ -2025,6 +2222,7 @@ mod tests { Some("http://127.0.0.1:8000/v1"), Some(0.0), Some(1.0), + None, ) .unwrap(); assert_eq!( @@ -2040,22 +2238,32 @@ mod tests { ); // Without sampling overrides the endpoint still comes from the - // environment, so only the baseURL section is written. + // environment; the model entry stays registered so OpenCode can + // resolve an id that is absent from its fetched catalog. let base_only = - opencode_provider_config("openai/model-x", Some("http://m:1/v1"), None, None).unwrap(); + opencode_provider_config("openai/model-x", Some("http://m:1/v1"), None, None, None) + .unwrap(); assert_eq!( base_only, - json!({"provider": {"openai": {"options": {"baseURL": "http://m:1/v1"}}}}) + json!({"provider": {"openai": { + "options": {"baseURL": "http://m:1/v1"}, + "models": {"model-x": {}} + }}}) ); // Nothing to pin: leave OpenCode on its environment-only defaults. - assert!(opencode_provider_config("openai/model-x", None, None, None).is_none()); + assert!(opencode_provider_config("openai/model-x", None, None, None, None).is_none()); // A model without a provider namespace cannot be pinned either. assert!( - opencode_provider_config("model-x", Some("http://m:1/v1"), Some(0.0), None).is_none() + opencode_provider_config("model-x", Some("http://m:1/v1"), Some(0.0), None, None) + .is_none() ); // Blank endpoints are ignored rather than written. - assert!(opencode_provider_config("openai/model-x", Some(" "), None, None).is_none()); + assert!(opencode_provider_config("openai/model-x", Some(" "), None, None, None).is_none()); + let with_steps = + opencode_provider_config("openai/model-x", Some("http://m:1/v1"), None, None, Some(7)) + .unwrap(); + assert_eq!(with_steps["agent"]["build"]["steps"], 7); } #[test] diff --git a/crates/persisting-replay/src/adapter/mod.rs b/crates/persisting-replay/src/adapter/mod.rs index b9f4f551c..0e41d7dfd 100644 --- a/crates/persisting-replay/src/adapter/mod.rs +++ b/crates/persisting-replay/src/adapter/mod.rs @@ -307,6 +307,9 @@ fn run_sdk_bridge( let output = run_process(ProcessSpec { command, stdin: None, + idle_timeout: None, + step_finish_limit: None, + stdout_redirect: None, timeout: Duration::from_secs(24 * 60 * 60), termination_grace: Duration::from_secs(2), pipe_grace: Duration::from_millis(250), diff --git a/crates/persisting-replay/src/adapter/openhands.rs b/crates/persisting-replay/src/adapter/openhands.rs index 22a7c74b8..be5837407 100644 --- a/crates/persisting-replay/src/adapter/openhands.rs +++ b/crates/persisting-replay/src/adapter/openhands.rs @@ -447,6 +447,9 @@ fn run_openhands( let output = run_process(ProcessSpec { command, stdin: Some(b"\n".to_vec()), + idle_timeout: None, + step_finish_limit: None, + stdout_redirect: None, timeout: Duration::from_secs(24 * 60 * 60), termination_grace: Duration::from_secs(2), pipe_grace: Duration::from_millis(250), diff --git a/crates/persisting-replay/src/opencode_bridge.rs b/crates/persisting-replay/src/opencode_bridge.rs index 382eb34fa..e473d4e13 100644 --- a/crates/persisting-replay/src/opencode_bridge.rs +++ b/crates/persisting-replay/src/opencode_bridge.rs @@ -347,7 +347,11 @@ async fn forward_handler( let Some(path) = path else { return error_response(StatusCode::BAD_REQUEST, "request has no path"); }; - if std::env::var("PVISOR_OPENCODE_BRIDGE_DEBUG").is_ok() { + let debug_level = std::env::var("PVISOR_OPENCODE_BRIDGE_DEBUG") + .ok() + .and_then(|value| value.trim().parse::().ok()) + .unwrap_or(0); + if debug_level > 0 { use std::io::Write; if let Ok(mut log) = std::fs::OpenOptions::new() .create(true) @@ -367,6 +371,17 @@ async fn forward_handler( String::from_utf8_lossy(&body[..body.len().min(300)]).replace('\n', " ") ); } + if debug_level >= 2 { + if let Ok(mut bodies) = std::fs::OpenOptions::new() + .create(true) + .append(true) + .open("/tmp/pvisor-opencode-bridge-bodies.log") + { + use std::io::Write as _; + let _ = bodies.write_all(&body); + let _ = bodies.write_all(b"\n===REQUEST-END===\n"); + } + } } let content_type = headers .get(axum::http::header::CONTENT_TYPE) @@ -385,6 +400,17 @@ async fn forward_handler( } else { body }; + if debug_level >= 2 { + use std::io::Write; + if let Ok(mut bodies) = std::fs::OpenOptions::new() + .create(true) + .append(true) + .open("/tmp/pvisor-opencode-bridge-upstream.log") + { + let _ = bodies.write_all(&body); + let _ = bodies.write_all(b"\n===UPSTREAM-END===\n"); + } + } let upstream_url = format!("{}{}", shared.upstream_origin, path); let mut upstream = shared.client.request(method, &upstream_url); for (name, value) in headers.iter() { @@ -439,14 +465,23 @@ async fn forward_handler( response_headers.insert(name.clone(), header_value); } } - let _ = response_headers.remove("content-length"); - let stream = response.bytes_stream(); + // Buffer the whole upstream body before replying, mirroring the Codex + // bridge. A pass-through stream can die mid-flight; OpenCode's Bun-based + // fetch then hangs on the half-open response without retrying, which + // stalls the continuation forever. + let body = match response.bytes().await { + Ok(bytes) => bytes, + Err(error) => { + fail(&shared, format!("read upstream OpenCode response: {error}")); + return error_response(StatusCode::BAD_GATEWAY, &error.to_string()); + } + }; let mut builder = Response::builder() .status(StatusCode::from_u16(status.as_u16()).unwrap_or(StatusCode::BAD_GATEWAY)); for (name, value) in response_headers.iter() { builder = builder.header(name.clone(), value.clone()); } - match builder.body(Body::from_stream(stream)) { + match builder.body(Body::from(body)) { Ok(response) => { if let Ok(mut state) = shared.state.lock() { state.forwarded_requests += 1; diff --git a/crates/persisting-replay/src/process.rs b/crates/persisting-replay/src/process.rs index 7571e6feb..3f92b24de 100644 --- a/crates/persisting-replay/src/process.rs +++ b/crates/persisting-replay/src/process.rs @@ -16,6 +16,21 @@ use crate::error::{ReplayError, ReplayErrorKind, ResultExt}; pub(crate) struct ProcessSpec { pub command: Command, pub stdin: Option>, + /// Terminate the process when neither stdout, stderr, nor a redirected + /// stdout file produced new bytes for this long. Long silent stretches + /// are the signature of an agent CLI that wedged internally instead of + /// working. + pub idle_timeout: Option, + /// Terminate the process once stdout has carried this many + /// `"type":"step_finish"` JSONL events. OpenCode ignores its + /// `agent.steps` budget on resumed sessions, so pVisor enforces the + /// remaining live-action budget itself. + pub step_finish_limit: Option, + /// Redirect the child's stdout straight to this file instead of a pipe. + /// OpenCode's Bun runtime fully buffers stdout on pipes (events only + /// appear at exit) but streams into regular files, so event-driven + /// watchdogs must watch the file. + pub stdout_redirect: Option, pub timeout: Duration, pub termination_grace: Duration, pub pipe_grace: Duration, @@ -33,6 +48,7 @@ pub(crate) struct ProcessOutput { pub stdout_truncated: bool, pub stderr_truncated: bool, pub timed_out: bool, + pub step_limited: bool, pub background_cleanup: bool, } @@ -46,7 +62,23 @@ struct StreamCapture { pub(crate) fn run_process(mut spec: ProcessSpec) -> Result { let log = owner_only_log(&spec.log_path)?; let log = Arc::new(Mutex::new(log)); - spec.command.stdout(Stdio::piped()).stderr(Stdio::piped()); + let stdout_redirect = spec.stdout_redirect.take(); + if let Some(target) = &stdout_redirect { + if let Some(parent) = target.parent() { + std::fs::create_dir_all(parent) + .replay_context(ReplayErrorKind::Executor, "create stdout redirect parent")?; + } + let file = std::fs::OpenOptions::new() + .create(true) + .truncate(true) + .write(true) + .open(target) + .replay_context(ReplayErrorKind::Executor, "open stdout redirect file")?; + spec.command.stdout(Stdio::from(file)); + } else { + spec.command.stdout(Stdio::piped()); + } + spec.command.stderr(Stdio::piped()); if spec.stdin.is_some() { spec.command.stdin(Stdio::piped()); } @@ -65,16 +97,41 @@ pub(crate) fn run_process(mut spec: ProcessSpec) -> Result Result Result Result= limit + { + step_limited = true; + background_cleanup = true; + // SIGINT lets OpenCode exit gracefully and flush its buffered + // stdout events; SIGTERM would discard them. + break terminate_running_group_with( + &mut child, + process_group, + Duration::from_secs(20).max(spec.termination_grace), + libc::SIGINT, + )?; + } + if let Some(path) = &stdout_redirect + && let Ok(meta) = std::fs::metadata(path) + { + let len = meta.len(); + if len > redirect_seen_bytes { + redirect_seen_bytes = len; + let now_ms = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|value| value.as_millis() as u64) + .unwrap_or(0); + last_activity.store(now_ms, std::sync::atomic::Ordering::Release); + } + } + if let Some(idle_timeout) = spec.idle_timeout { + let now_ms = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|value| value.as_millis() as u64) + .unwrap_or(0); + let last_ms = last_activity.load(std::sync::atomic::Ordering::Acquire); + if now_ms.saturating_sub(last_ms) >= idle_timeout.as_millis() as u64 { + background_cleanup = true; + break terminate_running_group(&mut child, process_group, spec.termination_grace)?; + } + } thread::sleep(Duration::from_millis(10)); }; @@ -121,21 +219,41 @@ pub(crate) fn run_process(mut spec: ProcessSpec) -> Result reader + .join() + .map_err(|_| ReplayError::new(ReplayErrorKind::Internal, "stdout reader panicked"))? + .replay_context(ReplayErrorKind::Executor, "drain supervised stdout")?, + None => { + // Redirected stdout: summarize the redirect file itself so the + // caller keeps its usual byte accounting. + let bytes = stdout_redirect + .as_deref() + .and_then(|path| std::fs::read(path).ok()) + .unwrap_or_default(); + let tail_len = bytes.len().min(spec_retained(spec.retained_bytes)); + StreamCapture { + tail: bytes[bytes.len() - tail_len..].to_vec(), + total: bytes.len() as u64, + log_error: None, + } + } + }; let stderr = stderr_reader .join() .map_err(|_| ReplayError::new(ReplayErrorKind::Internal, "stderr reader panicked"))? @@ -156,6 +274,7 @@ pub(crate) fn run_process(mut spec: ProcessSpec) -> Result( mut reader: R, log: Arc>, retained_bytes: usize, + activity: Option>, + step_counter: Option<(Arc, bool)>, ) -> thread::JoinHandle> where R: Read + Send + 'static, @@ -197,6 +318,19 @@ where break; } total = total.saturating_add(count as u64); + if let Some((counter, _)) = &step_counter { + let hits = count_step_finish(&chunk[..count]); + if hits > 0 { + counter.fetch_add(hits, std::sync::atomic::Ordering::AcqRel); + } + } + if let Some(activity) = &activity { + let now_ms = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|value| value.as_millis() as u64) + .unwrap_or(0); + activity.store(now_ms, std::sync::atomic::Ordering::Release); + } if log_error.is_none() { let write_result = log .lock() @@ -216,6 +350,28 @@ where }) } +fn spec_retained(retained_bytes: usize) -> usize { + retained_bytes.max(1) +} + +fn count_step_finish(chunk: &[u8]) -> usize { + // OpenCode's stdout JSONL is block-buffered (events only flush in ~8KB + // batches or at exit), but `--print-logs` stderr carries one + // `message=loop ... step=N` line per live turn in real time. + const NEEDLE: &[u8] = b"message=loop"; + let mut hits = 0; + let mut offset = 0; + while offset + NEEDLE.len() <= chunk.len() { + if &chunk[offset..offset + NEEDLE.len()] == NEEDLE { + hits += 1; + offset += NEEDLE.len(); + } else { + offset += 1; + } + } + hits +} + fn retain_tail(tail: &mut Vec, chunk: &[u8], limit: usize) { if limit == 0 { tail.clear(); @@ -237,7 +393,17 @@ fn terminate_running_group( process_group: i32, grace: Duration, ) -> Result { - let _ = signal_group(process_group, libc::SIGTERM)?; + terminate_running_group_with(child, process_group, grace, libc::SIGTERM) +} + +#[cfg(unix)] +fn terminate_running_group_with( + child: &mut std::process::Child, + process_group: i32, + grace: Duration, + first_signal: libc::c_int, +) -> Result { + let _ = signal_group(process_group, first_signal)?; let deadline = Instant::now() + grace; loop { if let Some(status) = child @@ -343,6 +509,9 @@ mod tests { ProcessSpec { command, stdin: None, + idle_timeout: None, + step_finish_limit: None, + stdout_redirect: None, timeout: Duration::from_secs(5), termination_grace: Duration::from_millis(100), pipe_grace: Duration::from_millis(100), @@ -423,6 +592,77 @@ mod tests { ); } + #[test] + fn counts_stderr_loop_progress_lines() { + use super::count_step_finish; + assert_eq!(count_step_finish(b"message=loop step=1"), 1); + assert_eq!( + count_step_finish(b"message=loop step=1\nmessage=loop step=2"), + 2 + ); + assert_eq!(count_step_finish(b"message=tracking"), 0); + assert_eq!(count_step_finish(b"message=exiting loop"), 0); + } + + #[test] + fn step_finish_limit_terminates_the_process_early() { + let temporary = tempfile::tempdir().unwrap(); + let log_path = temporary.path().join("steps.log"); + let script = + "for i in 1 2 3 4 5; do echo 'level=INFO message=loop step='$i >&2; sleep 10; done"; + let mut spec = shell_spec(script, &log_path); + spec.step_finish_limit = Some(2); + spec.timeout = Duration::from_secs(120); + let output = run_process(spec).unwrap(); + assert!(output.step_limited); + assert!(!output.timed_out); + // The third emission never happens: the loop is killed during the + // second sleep. + let log = std::fs::read_to_string(log_path).unwrap(); + assert_eq!(log.matches("message=loop").count(), 2); + } + + #[test] + fn redirected_stdout_growth_refreshes_idle_watchdog() { + let temporary = tempfile::tempdir().unwrap(); + let log_path = temporary.path().join("redirect-idle.log"); + let events_path = temporary.path().join("events.jsonl"); + // stderr stays silent; only the redirect file grows. Without polling + // the redirect, a 300ms idle watchdog would kill this mid-loop. + let script = "for i in 1 2 3 4 5 6; do echo event-$i; sleep 0.2; done"; + let mut spec = shell_spec(script, &log_path); + spec.stdout_redirect = Some(events_path.clone()); + spec.idle_timeout = Some(Duration::from_millis(300)); + spec.timeout = Duration::from_secs(10); + + let output = run_process(spec).unwrap(); + + assert!(output.status.success()); + assert!(!output.timed_out); + assert!(!output.step_limited); + let events = std::fs::read_to_string(&events_path).unwrap(); + assert_eq!(events.lines().count(), 6); + } + + #[test] + fn idle_timeout_still_fires_when_redirect_stalls() { + let temporary = tempfile::tempdir().unwrap(); + let log_path = temporary.path().join("redirect-stall.log"); + let events_path = temporary.path().join("events.jsonl"); + let mut spec = shell_spec("echo once; sleep 5", &log_path); + spec.stdout_redirect = Some(events_path); + spec.idle_timeout = Some(Duration::from_millis(200)); + spec.timeout = Duration::from_secs(10); + let started = Instant::now(); + + let output = run_process(spec).unwrap(); + + assert!(started.elapsed() < Duration::from_secs(2)); + assert!(!output.status.success()); + assert!(!output.timed_out); + assert!(output.background_cleanup); + } + #[test] fn times_out_and_reaps_the_foreground_process_group() { let temporary = tempfile::tempdir().unwrap(); diff --git a/docs/src/zh/pvisor/guides/sandbox-replay.md b/docs/src/zh/pvisor/guides/sandbox-replay.md index f22a22196..1b384e4b1 100644 --- a/docs/src/zh/pvisor/guides/sandbox-replay.md +++ b/docs/src/zh/pvisor/guides/sandbox-replay.md @@ -85,6 +85,51 @@ SandboxReplay 在新沙箱中重新执行命令型工具以及 `read`、`write` 用新的 `state.output` 重建前缀,并通过 `opencode run --format=json --session` 从边界 继续。工具执行期间的中间 `tool_use` 状态会合并,避免同一调用被重复回放。 +`opencode run` 没有 CLI 步数参数,且 1.17.7 在 resume 会话上不执行 +`agent.build.steps` 配置预算。SandboxReplay 因此为 OpenCode 引入了两个外部看门狗 +(对其它 Agent 不启用),详见下文「OpenCode 看门狗」。 + +#### OpenCode 看门狗 + +OpenCode 是目前唯一既以黑盒 CLI 形式续跑、又缺乏可靠停止机制的 Agent: +Claude Code 的 `max_turns`、Codex 的 `agent_max_steps`、mini-swe-agent 的 +`step_limit` 都在 resume 上原生生效,mini-swe 与 Pi 的循环本身由 pVisor 驱动; +而 OpenCode 两头都不占。看门狗由 `run_process` 的两个监督条件实现,仅对 +OpenCode 续跑启用。 + +**为什么必须引入(两个实测问题)** + +1. **resume 会话无视步数预算(上游缺陷)**。`agent.build.steps` 在全新会话上 + 原生生效,但用 `opencode import` + `run --session` 恢复的会话完全忽略该 + 配置。绕开 pVisor 的裸实验可直接复现:新会话设 `steps=2` 恰好 2 步停; + resume 会话设 `steps=3` 连跑 32 步以上不停。pVisor 已把剩余预算 + (`max_steps - after_step`)写入隔离配置,但该值目前不被读取。贪心采样下 + 模型还会陷入无限循环(实测单回合刷 24 万 token、本地循环 200+ 回合), + 若无外部干预,续跑会一直占用沙箱直到外层 agent 超时(小时级)。 +2. **续跑进程偶发静默僵死**。实测抓到过:CLI 进程存活但无网络连接、无工具 + 子进程、事件循环空转,既不推进也不退出。根因在 OpenCode/Bun 一侧,且不会 + 自行恢复;同样会耗尽整个 agent 超时窗口。 + +**如何解决(两个看门狗的机制)** + +| 看门狗 | 触发条件 | 动作 | +|---|---|---| +| 步数看门狗 | stderr 进度日志中的回合行数达到剩余预算 | SIGINT 优雅终止 | +| 空闲看门狗 | stdout/stderr 连续 10 分钟无任何输出(僵死特征) | 终止进程组 | + +步数看门狗的信号源经过专门筛选。OpenCode 的 stdout JSONL 事件流被 Bun 运行时 +按 ~8KB 块缓冲(管道、PTY、文件重定向均非实时),不能作为计数源;`--print-logs` +输出到 stderr 的进度日志实时流式,且每个模型回合固定产生一行 +`message=loop ... step=N`。续跑命令因此固定附加 `--print-logs`,pVisor 实时统计 +该行数,达到剩余预算即向进程组发 SIGINT——SIGINT 下 OpenCode 会优雅退出并刷出 +缓冲的全部事件;若进程被更强信号杀死导致 stdout 缓冲丢失,则从 OpenCode 的 +sqlite 会话库重建续跑事件(任务沙箱自带 python3,无需新增依赖),续跑轨迹不丢。 + +终止结果如实上报:步数看门狗触发时 `agent_status` 为 `max_steps`,结果 metadata +携带 `opencode_step_budget` 标记(`enforced_by: pvisor_event_watchdog`)。 +隔离配置中的 `agent.build.steps` 仍然保留:一旦上游修复 resume 会话读取该配置, +原生预算将直接生效,看门狗自动退化为兜底保险。 + ### 3.6 Codex Codex 适配固定支持 CLI `0.149.0`,输入为 Codex 原生 rollout JSONL。一个 replay step @@ -248,8 +293,10 @@ replay journal 不记录提示词明文;Agent 原生的 prepared 或 continued - 模型:`Qwen3.6-35B-A3B`; - reasoning/thinking:关闭; - 题目:NodeBB(291)、Vuls(666)、qutebrowser(667); -- Agent:Claude Code、OpenHands、mini-swe-agent、Pi agent; +- Agent:Claude Code、OpenHands、mini-swe-agent、Pi agent、OpenCode; - 每个 Agent 先在新沙箱中生成原始轨迹,再创建另一个新沙箱,仅使用 pVisor SandboxReplay 续跑; +- OpenCode 录制与续跑均使用贪心采样:OpenCode 不透传采样参数,SweEval 与续跑桥在 + wire 级注入 `temperature=0`、`top_p=1`,thinking 经 `chat_template_kwargs` 关闭; - 每个沙箱资源:2 CPU、7 GiB 内存、70 GiB 存储; - `N` 按 Rust 解析器识别出的完整原生工具批次序号统计;原始和续跑总步数按相应 Agent 原生轨迹中的 turn/action 数统计; @@ -275,6 +322,9 @@ replay journal 不记录提示词明文;Agent 原生的 prepared 或 continued | Codex | NodeBB(291) | 37 | 74 | 88 | 1 | 1 | 否 | 0.42 | | Codex | Vuls(666) | 28 | 57 | 69 | 1 | 1 | 否 | 0.47 | | Codex | qutebrowser(667) | 27 | 54 | 43 | 1 | 1 | 是 | 0.97 | +| OpenCode | NodeBB(291) | 20 | 40 | 44 | 1 | 1 | 否 | 0.65 | +| OpenCode | Vuls(666) | 7 | 16 | 31 | 1 | 1 | 否 | N/A | +| OpenCode | qutebrowser(667) | 15 | 31 | 39 | 1 | 1 | 是 | 0.76 | Claude Code / NodeBB 使用 `N=1`,以避开异步子 Agent 完成后的 Resume Transport canonical-prefix 歧义。边界后的原始可见文本非空,且 `A′(N+1)` 成功复现同一个 `TaskOutput` 调用。 @@ -286,6 +336,16 @@ Codex 三题使用 `Ornith-1.5-35B-A3B`,原始轨迹来自 r24,续跑轨迹 r30;续跑使用新的任务沙箱和默认 `replayed_boundary_only` 输入条件。Codex 轨迹中未 发现 transport nonce 或 `Continue from the replay boundary.` 消息。 +OpenCode 三题(`1.17.7`)与上表其余 Agent 相同使用 `Qwen3.6-35B-A3B`,贪心采样, +三题并发录制与并发续跑均一次通过,verifier Reward 全部为 1;Vuls(666)与 +qutebrowser(667)另复跑一轮续跑,Reward 与下一动作行为一致(可复现)。续跑的 +`max_steps=N+25` 由 pVisor 事件看门狗在剩余预算处软停,三题 live 步数均为 24。 +无污染在 wire 级验证:首个上游模型请求共 51 项,精确终止于边界观测 `O'N` +(`function_call_output`),传输 nonce 在全部上游请求中出现 0 次, +`boundary_user_prompt` 未注入。OpenCode 关闭 thinking 时常直接调用工具而不输出 +可见文本,此时文本相似度按 Pi agent 的同一语义记为 N/A,而不是把两个空字符串 +报告为 1.00。 + ### 5.3 A(N+1) 与 A′(N+1) 以下直接展示原始下一步和续跑下一步。reasoning 已移除;较长的整文件替换只保留修改目标和能够区分调用的关键内容。 @@ -701,6 +761,82 @@ Now let me run the gold-style verification test. 区别:工具完全一致,文本相似度为 0.97。 +#### OpenCode + +##### NodeBB(291),N=20 + +`A(N+1)`: + +~~~text +The tests need Redis running. Let me check if there's a way to start it or if there are unit tests: + +Bash: + which redis-server 2>/dev/null || echo "redis-server not found" + redis-cli ping 2>/dev/null || echo "Redis not running" +~~~ + +`A′(N+1)`: + +~~~text +The tests need Redis running. Let me check if there's a way to run tests with MongoDB instead, or look at the test structure: + +Bash: + cat /app/config.json 2>/dev/null | head -30 +~~~ + +区别:两边以同一句 "The tests need Redis running. Let me check..." 开头(相同前缀 +上下文的强收敛),但探查动作不同:原轨迹检查 Redis 进程,续跑读取应用配置。文本 +相似度为 0.65。 + +##### Vuls(666),N=7 + +`A(N+1)`: + +~~~text +Now let me check if there's an `oracle` type in the scan package that I need to reference: + +grep: + pattern="type oracle struct" path="/app/scan" +~~~ + +`A′(N+1)`: + +~~~text +[可见文本为空白,直接调用工具] + +grep: + path="/app/scan/*.go" +grep: + pattern="Oracle" path="*.go" +~~~ + +区别:工具语义同向——都在第 8 步去 scan 包定位 oracle 类型定义,续跑用两个更宽的 +grep 变体覆盖同一目标。续跑侧可见文本为空白(OpenCode 关闭 thinking 后的常见风格), +文本相似度按指标语义记为 N/A。复跑一轮行为一致。 + +##### qutebrowser(667),N=15 + +`A(N+1)`: + +~~~text +Let me run the tests to see if my changes work: + +Bash: + python -m pytest tests/unit/browser/webengine/test_webview.py::test_suffixes_workaround_extras_returned -v --tb=short 2>&1 | tail -30 +~~~ + +`A′(N+1)`: + +~~~text +Now let me run the tests to see if they pass: + +Bash: + python -m pytest tests/unit/browser/webengine/test_webview.py::test_suffixes_workaround_extras_returned -v --tb=short 2>&1 | tail -30 +~~~ + +区别:同一条 pytest 命令逐字符一致,仅解说文本措辞不同,文本相似度为 0.76。首轮 +续跑的 `A′(N+1)` 可见文本为空白但命令同样逐字复现;两轮续跑 verifier Reward 均为 1。 + 精确参数见 [`pvisor replay` 命令参考](../reference/cli.md#replay-an-agent-trajectory);执行边界见 [执行指南](execution.md)。 From aa0da015e59db77ed891ba626a590dc966939006 Mon Sep 17 00:00:00 2001 From: guoxu1 Date: Fri, 11 Sep 2026 11:21:53 +0800 Subject: [PATCH 8/8] fix(replay): satisfy clippy in OpenCode bridge --- crates/persisting-replay/src/opencode_bridge.rs | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/crates/persisting-replay/src/opencode_bridge.rs b/crates/persisting-replay/src/opencode_bridge.rs index e473d4e13..3dd1dc937 100644 --- a/crates/persisting-replay/src/opencode_bridge.rs +++ b/crates/persisting-replay/src/opencode_bridge.rs @@ -371,16 +371,15 @@ async fn forward_handler( String::from_utf8_lossy(&body[..body.len().min(300)]).replace('\n', " ") ); } - if debug_level >= 2 { - if let Ok(mut bodies) = std::fs::OpenOptions::new() + if debug_level >= 2 + && let Ok(mut bodies) = std::fs::OpenOptions::new() .create(true) .append(true) .open("/tmp/pvisor-opencode-bridge-bodies.log") - { - use std::io::Write as _; - let _ = bodies.write_all(&body); - let _ = bodies.write_all(b"\n===REQUEST-END===\n"); - } + { + use std::io::Write as _; + let _ = bodies.write_all(&body); + let _ = bodies.write_all(b"\n===REQUEST-END===\n"); } } let content_type = headers