fix(runner): restore verified native provider startup

This commit is contained in:
Dotta 2026-09-02 17:23:10 -05:00
parent 910de2f2f1
commit 03bab72ee8
6 changed files with 305 additions and 43 deletions

View File

@ -23,7 +23,9 @@ use crate::durable::{
AcpxLaunchProfile, Command, CommandExecution, CommandExecutor, DurableRunnerConfig,
DurableRunnerError, EventPriority, PolledEvent,
};
use crate::process_supervisor::{VerifiedProcessArgument, VerifiedProcessLaunch};
use crate::process_supervisor::{
VerifiedProcessArgument, VerifiedProcessLaunch, VERIFIED_NODE_ESM_DEFAULT_TYPE_ARG,
};
use crate::provider_bridge::{
authorized_tool_catalog_digest, AuthorizedToolSet, ToolResult, TOOL_SET_SCHEMA,
};
@ -210,7 +212,7 @@ impl AcpxProviderDescriptor {
"ACPX runner launch profile does not authenticate its command",
)
})?;
let verified_args = launch_profile
let mut verified_args = launch_profile
.args
.iter()
.map(|argument| {
@ -229,6 +231,15 @@ impl AcpxProviderDescriptor {
})
})
.collect::<Result<Vec<_>, _>>()?;
// Verified scripts are exposed to Node through an immutable descriptor
// path (for example /proc/self/fd/12), which intentionally has no .js
// suffix. Pin ESM interpretation so Node does not misclassify the
// bundled sidecar as CommonJS merely because its verified path is
// extensionless.
verified_args.insert(
0,
VerifiedProcessArgument::Literal(VERIFIED_NODE_ESM_DEFAULT_TYPE_ARG.to_owned()),
);
Ok(AcpxSidecarTransportConfig {
command: launch_profile.command.clone(),
args: launch_profile.args.clone(),
@ -1416,7 +1427,16 @@ mod tests {
let transport = descriptor.verified_transport(Some(&profile)).unwrap();
assert_eq!(transport.command, profile.command);
assert_eq!(transport.args[0], sidecar.to_string_lossy());
assert!(transport.verified_launch.is_some());
let verified_launch = transport.verified_launch.as_ref().unwrap();
assert!(matches!(
verified_launch.arguments().first(),
Some(VerifiedProcessArgument::Literal(argument))
if argument == VERIFIED_NODE_ESM_DEFAULT_TYPE_ARG
));
assert!(matches!(
verified_launch.arguments().get(1),
Some(VerifiedProcessArgument::Artifact(_))
));
let mut drifted_path = descriptor.clone();
drifted_path.sidecar_command = directory.join("other-node");

View File

@ -825,6 +825,15 @@ fn validate_reserved_terminal_value(
}
fn reserved_terminal_tool_bridge() -> Result<ProviderToolBridge, LocalRunnerError> {
let tool_set = reserved_terminal_tool_set()?;
let mut bridge = ProviderToolBridge::default();
bridge.prepare(tool_set).map_err(|error| {
LocalRunnerError::invalid(format!("ACPX reserved terminal tools are invalid: {error}"))
})?;
Ok(bridge)
}
fn reserved_terminal_tool_set() -> Result<AuthorizedToolSet, LocalRunnerError> {
let result_schema: Value = serde_json::from_str(include_str!(
"../../../../protocol/schemas/result.schema.json"
))
@ -848,18 +857,20 @@ fn reserved_terminal_tool_bridge() -> Result<ProviderToolBridge, LocalRunnerErro
let catalog_digest = authorized_tool_catalog_digest(&operations).map_err(|error| {
LocalRunnerError::invalid(format!("ACPX reserved terminal tools are invalid: {error}"))
})?;
let mut bridge = ProviderToolBridge::default();
bridge
.prepare(AuthorizedToolSet {
schema: TOOL_SET_SCHEMA.to_owned(),
schema_version: 1,
catalog_digest,
operations,
})
.map_err(|error| {
LocalRunnerError::invalid(format!("ACPX reserved terminal tools are invalid: {error}"))
})?;
Ok(bridge)
Ok(AuthorizedToolSet {
schema: TOOL_SET_SCHEMA.to_owned(),
schema_version: 1,
catalog_digest,
operations,
})
}
fn sidecar_tool_operations(
run_tool_set: &AuthorizedToolSet,
) -> Result<Vec<AuthorizedTool>, LocalRunnerError> {
let mut operations = run_tool_set.operations.clone();
operations.extend(reserved_terminal_tool_set()?.operations);
Ok(operations)
}
fn validate_prp_run_result(value: &Value) -> Result<(), LocalRunnerError> {
@ -891,6 +902,7 @@ fn bootstrap(
transport: &mut AcpxSidecarTransport,
config: &AcpxProviderSessionConfig,
) -> Result<(AcpxProviderSessionIdentity, AcpxProviderState), LocalRunnerError> {
let sidecar_tools = sidecar_tool_operations(&config.tool_set)?;
let initialized = transport.request(
GeneratedAcpxSidecarCommand::Initialize,
json!({"agent": config.agent, "model": config.model}),
@ -909,7 +921,7 @@ fn bootstrap(
"permissionModePinned": config.permission_mode_pinned,
"systemInstructions": config.system_instructions,
"runtimeContext": Value::Null,
"tools": config.tool_set.operations,
"tools": &sidecar_tools,
"expectedIdentity": config.expected_identity,
}),
)?;
@ -920,7 +932,7 @@ fn bootstrap(
json!({
"runId": config.run_id,
"catalogRevision": config.catalog_revision,
"tools": config.tool_set.operations,
"tools": &sidecar_tools,
}),
)?;
if attached.get("runId").and_then(Value::as_str) != Some(config.run_id.as_str())
@ -1119,3 +1131,41 @@ fn with_cleanup_error(
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn sidecar_catalog_combines_run_tools_with_trusted_terminal_tools() {
let operations = vec![AuthorizedTool {
operation_id: "get_task_context".to_owned(),
version: 1,
description: "Read the task context.".to_owned(),
input_schema: json!({"type":"object"}),
response_schema: json!({"type":"object"}),
}];
let run_tool_set = AuthorizedToolSet {
schema: TOOL_SET_SCHEMA.to_owned(),
schema_version: 1,
catalog_digest: authorized_tool_catalog_digest(&operations).unwrap(),
operations,
};
let sidecar_tools = sidecar_tool_operations(&run_tool_set).unwrap();
assert_eq!(
sidecar_tools
.iter()
.map(|tool| tool.operation_id.as_str())
.collect::<Vec<_>>(),
vec![
"get_task_context",
PRP_COMPLETION_TOOL_NAME,
PRP_BLOCK_TOOL_NAME,
]
);
assert_eq!(run_tool_set.operations.len(), 1);
assert!(sidecar_tools[1].input_schema.is_object());
assert!(sidecar_tools[2].input_schema.is_object());
}
}

View File

@ -10,10 +10,13 @@ use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use sha2::{Digest, Sha256};
#[cfg(test)]
use crate::durable::QualifiedLaunchArtifact;
use crate::durable::{redact_text, OpenCodeLaunchProfile};
use crate::local_runner::LocalRunnerError;
use crate::process_supervisor::{
SupervisedProcess, VerifiedProcessArgument, VerifiedProcessLaunch,
BoundedLogBuffer, ProcessOutput, SupervisedProcess, VerifiedProcessArgument,
VerifiedProcessLaunch, VERIFIED_NODE_ESM_DEFAULT_TYPE_ARG,
};
use crate::provider_bridge::{AuthorizedTool, DurableReplayFilter, ToolResult};
use crate::provider_events::normalized_codex_terminal_event_type;
@ -38,6 +41,8 @@ const OPENCODE_PROVIDER_ENVIRONMENT_KEYS: &[&str] = &[
"PAPERCLIP_NATIVE_RUNTIME_CONTEXT_PATH",
];
const TRUSTED_OPENCODE_EXECUTABLE_ARG: &str = "--paperclip-trusted-opencode-executable";
const MAX_PROVIDER_STDERR_LINES: usize = 32;
const MAX_PROVIDER_STDERR_BYTES: usize = 8 * 1024;
const MAX_INSTRUCTIONS_BYTES: usize = 1024 * 1024;
const MAX_PENDING_TOOL_REQUESTS: usize = 4_096;
const MAX_PENDING_TOOL_REQUEST_BYTES: usize = 16 * 1024 * 1024;
@ -489,6 +494,7 @@ enum AmbiguousTurnMessage {
pub struct CodexProvider {
process: SupervisedProcess,
stderr_tail: BoundedLogBuffer,
config: CodexProviderConfig,
authorized_tools: Vec<AuthorizedTool>,
next_request_id: u64,
@ -591,21 +597,7 @@ impl CodexProvider {
"OpenCode launch does not match the runner-owned qualified profile",
));
}
let command = verify_launch_artifact(&profile.command, "OpenCode proxy command")
.map_err(|error| LocalRunnerError::invalid(error.to_string()))?;
let proxy = verify_launch_artifact(&profile.proxy_script, "OpenCode proxy script")
.map_err(|error| LocalRunnerError::invalid(error.to_string()))?;
let executable =
verify_launch_artifact(&profile.executable, "OpenCode provider executable")
.map_err(|error| LocalRunnerError::invalid(error.to_string()))?;
let launch = VerifiedProcessLaunch::new(
command,
vec![
VerifiedProcessArgument::Artifact(proxy),
VerifiedProcessArgument::Literal(TRUSTED_OPENCODE_EXECUTABLE_ARG.to_owned()),
VerifiedProcessArgument::ExecutableArtifact(executable),
],
);
let launch = verified_opencode_launch(profile)?;
SupervisedProcess::spawn_verified_with_environment_keys(
&launch,
Duration::from_secs(2),
@ -623,6 +615,10 @@ impl CodexProvider {
};
let mut provider = Self {
process,
stderr_tail: BoundedLogBuffer::new(
MAX_PROVIDER_STDERR_LINES,
MAX_PROVIDER_STDERR_BYTES,
),
config: config.clone(),
authorized_tools,
next_request_id: 1,
@ -1795,14 +1791,29 @@ impl CodexProvider {
.map_err(ProviderRequestError::Ambiguous)?;
loop {
let line = self
.process
.receive_stdout_line(Duration::from_secs(30))
.map_err(ProviderRequestError::Ambiguous)?
.ok_or_else(|| {
ProviderRequestError::Ambiguous(LocalRunnerError::invalid(format!(
"Codex {method} response timed out"
)))
})?;
.receive_provider_stdout_line(Duration::from_secs(30))
.map_err(ProviderRequestError::Ambiguous)?;
let Some(line) = line else {
let exit = self
.process
.try_wait()
.map_err(ProviderRequestError::Ambiguous)?;
if exit.is_some() {
self.drain_provider_diagnostics(Duration::from_millis(50));
}
let diagnostic_suffix = self.provider_diagnostic_suffix();
let message = if let Some(exit) = exit {
format!(
"Codex {method} process exited before responding (exitCode={:?}, signal={:?}){diagnostic_suffix}",
exit.exit_code, exit.signal
)
} else {
format!("Codex {method} response timed out{diagnostic_suffix}")
};
return Err(ProviderRequestError::Ambiguous(LocalRunnerError::invalid(
message,
)));
};
let trace_frame_id = self.trace_inbound(&line);
let message = parse_provider_message(&line).map_err(|error| {
if let (Some(trace), Some(frame_id)) = (self.trace.as_mut(), trace_frame_id) {
@ -1870,6 +1881,85 @@ impl CodexProvider {
self.pending_message_bytes = next_retained_bytes;
}
}
fn receive_provider_stdout_line(
&mut self,
timeout: Duration,
) -> Result<Option<String>, LocalRunnerError> {
let deadline = std::time::Instant::now() + timeout;
loop {
let remaining = deadline.saturating_duration_since(std::time::Instant::now());
if remaining.is_zero() {
return Ok(None);
}
match self.process.recv_timeout(remaining) {
Ok(ProcessOutput::Stdout(line)) => return Ok(Some(line)),
Ok(ProcessOutput::Stderr(line)) => {
self.stderr_tail.push(redact_text(&line));
}
Ok(ProcessOutput::StdoutError(message)) => {
return Err(LocalRunnerError::invalid(message));
}
Ok(ProcessOutput::StdoutClosed) => return Ok(None),
Ok(ProcessOutput::StderrClosed) => {}
Err(mpsc::RecvTimeoutError::Timeout) => return Ok(None),
Err(mpsc::RecvTimeoutError::Disconnected) => return Ok(None),
}
}
}
fn drain_provider_diagnostics(&mut self, max_wait: Duration) {
let deadline = std::time::Instant::now() + max_wait;
loop {
let remaining = deadline.saturating_duration_since(std::time::Instant::now());
if remaining.is_zero() {
break;
}
match self.process.recv_timeout(remaining) {
Ok(ProcessOutput::Stderr(line)) => {
self.stderr_tail.push(redact_text(&line));
}
Ok(ProcessOutput::StderrClosed)
| Err(mpsc::RecvTimeoutError::Timeout)
| Err(mpsc::RecvTimeoutError::Disconnected) => break,
Ok(ProcessOutput::Stdout(_))
| Ok(ProcessOutput::StdoutError(_))
| Ok(ProcessOutput::StdoutClosed) => {}
}
}
}
fn provider_diagnostic_suffix(&self) -> String {
let diagnostics = self.stderr_tail.snapshot().lines.join("\n");
if diagnostics.is_empty() {
String::new()
} else {
format!(" stderrTail={diagnostics:?}")
}
}
}
fn verified_opencode_launch(
profile: &OpenCodeLaunchProfile,
) -> Result<VerifiedProcessLaunch, LocalRunnerError> {
let command = verify_launch_artifact(&profile.command, "OpenCode proxy command")
.map_err(|error| LocalRunnerError::invalid(error.to_string()))?;
let proxy = verify_launch_artifact(&profile.proxy_script, "OpenCode proxy script")
.map_err(|error| LocalRunnerError::invalid(error.to_string()))?;
let executable = verify_launch_artifact(&profile.executable, "OpenCode provider executable")
.map_err(|error| LocalRunnerError::invalid(error.to_string()))?;
Ok(VerifiedProcessLaunch::new(
command,
vec![
// The immutable verified proxy path is extensionless, so
// explicitly retain the ESM semantics of the bundled .js
// artifact instead of letting Node infer CommonJS.
VerifiedProcessArgument::Literal(VERIFIED_NODE_ESM_DEFAULT_TYPE_ARG.to_owned()),
VerifiedProcessArgument::Artifact(proxy),
VerifiedProcessArgument::Literal(TRUSTED_OPENCODE_EXECUTABLE_ARG.to_owned()),
VerifiedProcessArgument::ExecutableArtifact(executable),
],
))
}
fn json_size(value: &Value, label: &str) -> Result<usize, LocalRunnerError> {
@ -2583,6 +2673,58 @@ fn codex_question_response(
mod tests {
use super::*;
fn qualified_artifact(path: &Path) -> QualifiedLaunchArtifact {
QualifiedLaunchArtifact {
path: path.to_owned(),
sha256: format!("sha256:{:x}", Sha256::digest(fs::read(path).unwrap())),
}
}
#[test]
fn verified_opencode_proxy_retains_esm_semantics_for_extensionless_snapshot() {
let nonce = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos();
let directory = std::env::temp_dir().join(format!(
"paperclip-opencode-launch-{}-{nonce}",
std::process::id()
));
fs::create_dir_all(&directory).unwrap();
let command = directory.join("node");
let proxy = directory.join("proxy.js");
let executable = directory.join("opencode");
fs::write(&command, b"qualified node").unwrap();
fs::write(&proxy, b"export {};\n").unwrap();
fs::write(&executable, b"qualified opencode").unwrap();
let profile = OpenCodeLaunchProfile {
command: qualified_artifact(&command),
proxy_script: qualified_artifact(&proxy),
executable: qualified_artifact(&executable),
};
let launch = verified_opencode_launch(&profile).unwrap();
assert!(matches!(
launch.arguments().first(),
Some(VerifiedProcessArgument::Literal(argument))
if argument == VERIFIED_NODE_ESM_DEFAULT_TYPE_ARG
));
assert!(matches!(
launch.arguments().get(1),
Some(VerifiedProcessArgument::Artifact(_))
));
assert!(matches!(
launch.arguments().get(2),
Some(VerifiedProcessArgument::Literal(argument))
if argument == TRUSTED_OPENCODE_EXECUTABLE_ARG
));
assert!(matches!(
launch.arguments().get(3),
Some(VerifiedProcessArgument::ExecutableArtifact(_))
));
fs::remove_dir_all(directory).unwrap();
}
#[test]
fn admits_only_exact_local_facade_provider_driver_pairs() {
let mut config = CodexProviderConfig {

View File

@ -26,6 +26,7 @@ use sha2::{Digest, Sha256};
use crate::local_runner::LocalRunnerError;
const PROCESS_OUTPUT_QUEUE_CAPACITY: usize = 256;
pub(crate) const VERIFIED_NODE_ESM_DEFAULT_TYPE_ARG: &str = "--experimental-default-type=module";
#[derive(Clone, Debug)]
pub struct VerifiedProcessArtifact {
@ -199,6 +200,11 @@ impl VerifiedProcessLaunch {
Self { program, args }
}
#[cfg(test)]
pub(crate) fn arguments(&self) -> &[VerifiedProcessArgument] {
&self.args
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
fn inherited_command(&self) -> Result<InheritedCommand, LocalRunnerError> {
let mut inherited = Vec::with_capacity(self.args.len() + 1);

View File

@ -33,6 +33,7 @@ import {
import { releaseMaterializedNativeRuntimeSkills } from "../drivers/runtime-context-materializer.js";
import {
authorizedToolSetForProvider,
createCapabilityRunnerdCodexTransport,
createCapabilityRunnerdProviderEnvironment,
defaultCapabilityRunnerdBinary,
@ -54,6 +55,28 @@ import {
withCodexCollaborationRuntimeInstructions,
} from "./runnerd-codex-transport.js";
it("keeps ACPX terminal tools under the reserved runner-owned catalog", () => {
const tools = [
{
name: "get_task_context",
description: "Read the task context.",
inputSchema: { type: "object" },
},
...codexSemanticToolSpecs(),
];
expect(authorizedToolSetForProvider("acpx", tools)).toMatchObject({
operations: [{ operationId: "get_task_context" }],
});
expect(authorizedToolSetForProvider("codex", tools)).toMatchObject({
operations: [
{ operationId: "get_task_context" },
{ operationId: "paperclip_block" },
{ operationId: "paperclip_finish" },
],
});
});
it("defaults runnerd ACPX permissions to approve reads", () => {
expect(resolveRunnerdAcpxPermissionMode(undefined)).toBe("approve-reads");
expect(resolveRunnerdAcpxPermissionMode("deny-all")).toBe("deny-all");

View File

@ -1266,6 +1266,24 @@ function authorizedToolSet(
};
}
const ACPX_RESERVED_TERMINAL_TOOLS = new Set([
"paperclip_finish",
"paperclip_block",
]);
export function authorizedToolSetForProvider(
provider: CapabilityRunnerdCodexTransportOptions["provider"],
tools: readonly Readonly<Record<string, unknown>>[],
): Record<string, unknown> {
return authorizedToolSet(
provider === "acpx"
? tools.filter(
(tool) => !ACPX_RESERVED_TERMINAL_TOOLS.has(String(tool.name ?? "")),
)
: tools,
);
}
/**
* Raw provider tracing is consumed by runnerd itself. The provider child still
* receives the narrower allowlist enforced by Rust's `SupervisedProcess`, so
@ -1525,7 +1543,7 @@ class DurablePrpCodexTransport implements CodexAppServerTransport {
options.stateDirectory ??
mkdtempSync(resolve(tmpdir(), "paperclip-runner-lab-prp-"));
if (options.resumeDynamicTools !== undefined) {
this.#authorizedTools = authorizedToolSet([
this.#authorizedTools = authorizedToolSetForProvider(options.provider, [
...options.resumeDynamicTools,
...codexSemanticToolSpecs(),
]);
@ -2152,7 +2170,10 @@ class DurablePrpCodexTransport implements CodexAppServerTransport {
);
}
}
this.#authorizedTools = authorizedToolSet(dynamicTools);
this.#authorizedTools = authorizedToolSetForProvider(
provider,
dynamicTools,
);
const acpxAgent =
provider === "acpx" ? (this.options.acpxAgent ?? "codex") : null;
const requestedModel = typeof params.model === "string" ? params.model : "";