diff --git a/docs/plugins/codex-computer-use.md b/docs/plugins/codex-computer-use.md index 7edd534a464c..3a6d9896c795 100644 --- a/docs/plugins/codex-computer-use.md +++ b/docs/plugins/codex-computer-use.md @@ -278,7 +278,7 @@ reconciliation so OpenClaw does not override that selection. ## Remote marketplaces Remote marketplace support was introduced in Codex 0.146.1 and remains -available in OpenClaw's pinned Codex 0.153.4. OpenClaw passes the opaque remote +available in OpenClaw's pinned Codex 0.154.0. OpenClaw passes the opaque remote plugin ID returned by Codex to `plugin/read` and `plugin/install`; a human-readable plugin name is not a valid substitute. diff --git a/docs/plugins/codex-harness-reference/app-server-transport.md b/docs/plugins/codex-harness-reference/app-server-transport.md index c9dc3c1c2491..1de4384a115c 100644 --- a/docs/plugins/codex-harness-reference/app-server-transport.md +++ b/docs/plugins/codex-harness-reference/app-server-transport.md @@ -13,7 +13,7 @@ How OpenClaw starts and reaches the Codex app-server, and every `appServer` fiel ## App-server transport For ordinary harness turns, OpenClaw starts the managed Codex binary shipped -with the official plugin (currently `@openai/codex` `0.153.4`): +with the official plugin (currently `@openai/codex` `0.154.0`): ```bash codex app-server --listen stdio:// @@ -182,7 +182,7 @@ If the normal app-server runtime would be `danger-full-access`, enabling permission profile instead. Codex-managed network enforcement is sandboxed networking, so a full-access profile would not protect outbound traffic. -The plugin manages stable Codex app-server `0.153.4`. Explicit custom +The plugin manages stable Codex app-server `0.154.0`. Explicit custom executables, remote app-servers, and macOS desktop binaries must report a parseable semantic version of `0.149.0` or newer. Older, malformed, and unversioned handshakes are rejected. Newer versions log a compatibility warning diff --git a/docs/plugins/codex-harness-reference/approval-and-sandbox.md b/docs/plugins/codex-harness-reference/approval-and-sandbox.md index a37b590c9ba7..585a205c2a86 100644 --- a/docs/plugins/codex-harness-reference/approval-and-sandbox.md +++ b/docs/plugins/codex-harness-reference/approval-and-sandbox.md @@ -82,7 +82,7 @@ The stable default is fail-closed: active OpenClaw sandboxing disables native Codex execution surfaces that would otherwise run from the Codex app-server host. Use `appServer.experimental.sandboxExecServer: true` only when you want to try Codex's remote environment support with OpenClaw's sandbox backend. -This preview path uses the pinned Codex `0.153.4` app-server. +This preview path uses the pinned Codex `0.154.0` app-server. ```json5 { diff --git a/docs/plugins/codex-harness-reference/model-discovery.md b/docs/plugins/codex-harness-reference/model-discovery.md index b7408b787431..3c28e113bd1b 100644 --- a/docs/plugins/codex-harness-reference/model-discovery.md +++ b/docs/plugins/codex-harness-reference/model-discovery.md @@ -77,7 +77,7 @@ response remains authoritative even if it contains no visible models; HTTP `401` and `403` return an empty catalog rather than exposing fallback models. -The current bundled harness is `@openai/codex` `0.153.4`. A live `model/list` +The current bundled harness is `@openai/codex` `0.154.0`. A live `model/list` probe against the official `0.153.4` app-server, using an isolated, unauthenticated Codex home and `includeHidden: true`, returned this public subset of catalog metadata: diff --git a/docs/plugins/codex-harness.md b/docs/plugins/codex-harness.md index 8cb63cbcdf1c..330d628b4482 100644 --- a/docs/plugins/codex-harness.md +++ b/docs/plugins/codex-harness.md @@ -132,8 +132,8 @@ Proxy launch arguments are rejected to avoid changing a shared daemon's login. - The official `@openclaw/codex` plugin installed. Include `codex` in `plugins.allow` if your config uses an allowlist. -- Managed Codex app-server `0.153.4`. The plugin ships and manages - `@openai/codex` `0.153.4` by default, so a `codex` command on `PATH` does not +- Managed Codex app-server `0.154.0`. The plugin ships and manages + `@openai/codex` `0.154.0` by default, so a `codex` command on `PATH` does not affect normal startup. Explicit custom, remote, and macOS desktop-owned app-servers must report a parseable semantic version of `0.149.0` or newer. Newer versions continue with a compatibility warning and normal runtime diff --git a/docs/plugins/codex-native-plugins.md b/docs/plugins/codex-native-plugins.md index 849a6bfa823e..af8e210f2be3 100644 --- a/docs/plugins/codex-native-plugins.md +++ b/docs/plugins/codex-native-plugins.md @@ -23,7 +23,7 @@ working. - `plugins.entries.codex.enabled` is `true`. - `plugins.entries.codex.config.codexPlugins.enabled` is `true`. - Codex app-server reports `0.149.0` or newer. The official plugin ships - `@openai/codex` `0.153.4`; newer custom, remote, and macOS desktop-owned + `@openai/codex` `0.154.0`; newer custom, remote, and macOS desktop-owned binaries continue with a compatibility warning and normal runtime validation. - The target Codex app-server can see the expected marketplace, plugin, and app inventory. diff --git a/extensions/codex/package.json b/extensions/codex/package.json index 7a10cef8f972..07d862749274 100644 --- a/extensions/codex/package.json +++ b/extensions/codex/package.json @@ -8,7 +8,7 @@ }, "type": "module", "dependencies": { - "@openai/codex": "0.153.4", + "@openai/codex": "0.154.0", "semver": "7.8.5", "smol-toml": "1.8.0", "typebox": "1.3.27", diff --git a/extensions/codex/src/app-server/native-hook-relay.ts b/extensions/codex/src/app-server/native-hook-relay.ts index 9cefd3c2d5fa..4139e868c6ad 100644 --- a/extensions/codex/src/app-server/native-hook-relay.ts +++ b/extensions/codex/src/app-server/native-hook-relay.ts @@ -10,6 +10,7 @@ import type { registerNativeHookRelay, } from "openclaw/plugin-sdk/agent-harness-runtime"; import { emitTrustedDiagnosticEvent } from "openclaw/plugin-sdk/diagnostic-runtime"; +import { toErrorObject } from "openclaw/plugin-sdk/error-runtime"; import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; import { registerNativeHookRelayForBundledRuntime } from "openclaw/plugin-sdk/native-hook-relay-runtime"; import type { NativeHookRelayCommandPlan } from "openclaw/plugin-sdk/native-hook-relay-runtime"; @@ -223,6 +224,7 @@ export function createCodexNativeHookRelay(params: { promise: Promise; resolve: (claim: symbol) => void; reject: (reason: Error) => void; + waiters: number; } >(); let foregroundClosed = false; @@ -266,6 +268,7 @@ export function createCodexNativeHookRelay(params: { }), signal: params.signal, runBeforeToolCall: params.hostCapabilities.runBeforeToolCall, + approvalHost: params.hostCapabilities, assertActive: () => { params.hostCapabilities.assertActive(); params.assertCurrent?.(); @@ -277,7 +280,7 @@ export function createCodexNativeHookRelay(params: { shouldRetainAfterForegroundClose: () => successfulYieldRetentionAuthorized && directChildClaims.size > 0, allowPreToolUse: (childThreadId) => directChildClaims.has(childThreadId), - awaitForegroundAdmission: (childThreadId) => { + awaitForegroundAdmission: (childThreadId, signal) => { if (foregroundClosed) { return Promise.reject(new Error("native hook relay foreground admission unavailable")); } @@ -285,22 +288,44 @@ export function createCodexNativeHookRelay(params: { if (existingClaim) { return Promise.resolve(assertClaim(childThreadId, existingClaim)); } - const existingPending = pendingDirectChildAdmissions.get(childThreadId); - if (existingPending) { - return existingPending.promise.then((claim) => assertClaim(childThreadId, claim)); + let pending = pendingDirectChildAdmissions.get(childThreadId); + if (!pending) { + if (pendingDirectChildAdmissions.size >= MAX_PENDING_DIRECT_CHILD_ADMISSIONS) { + return Promise.reject( + new Error("native hook relay foreground admission capacity reached"), + ); + } + pending = { ...createDeferred(), waiters: 0 }; + pendingDirectChildAdmissions.set(childThreadId, pending); } - if (pendingDirectChildAdmissions.size >= MAX_PENDING_DIRECT_CHILD_ADMISSIONS) { - return Promise.reject( - new Error("native hook relay foreground admission capacity reached"), - ); - } - const { promise, resolve, reject } = createDeferred(); - pendingDirectChildAdmissions.set(childThreadId, { - promise, - resolve, - reject, + const admission = pending; + admission.waiters++; + let onAbort: (() => void) | undefined; + const wait = new Promise((resolve, reject) => { + void admission.promise.then(resolve, reject); + onAbort = () => + reject(toErrorObject(signal?.reason, "native hook relay admission aborted")); + signal?.addEventListener("abort", onAbort, { once: true }); + if (signal?.aborted) { + onAbort(); + } }); - return promise.then((claim) => assertClaim(childThreadId, claim)); + return wait + .then((claim) => assertClaim(childThreadId, claim)) + .finally(() => { + if (onAbort) { + signal?.removeEventListener("abort", onAbort); + } + // Duplicate callbacks share admission, but each owns its wait. A + // disconnected last waiter releases capacity without revoking a child. + admission.waiters--; + if ( + admission.waiters === 0 && + pendingDirectChildAdmissions.get(childThreadId) === admission + ) { + pendingDirectChildAdmissions.delete(childThreadId); + } + }); }, onDispose: () => { foregroundClosed = true; diff --git a/extensions/codex/src/app-server/run-attempt.native-hook-relay-retention.test.ts b/extensions/codex/src/app-server/run-attempt.native-hook-relay-retention.test.ts index fa6228db4292..767486b93a33 100644 --- a/extensions/codex/src/app-server/run-attempt.native-hook-relay-retention.test.ts +++ b/extensions/codex/src/app-server/run-attempt.native-hook-relay-retention.test.ts @@ -6,6 +6,7 @@ import { } from "openclaw/plugin-sdk/agent-harness-runtime"; import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; import { initializeGlobalHookRunner } from "openclaw/plugin-sdk/hook-runtime"; +import * as relayRuntime from "openclaw/plugin-sdk/native-hook-relay-runtime"; import { createAdmittedHostCapabilityTestFixture, createMockPluginRegistry, @@ -13,6 +14,7 @@ import { import { beforeEach, describe, expect, it, vi } from "vitest"; import { readAttemptTerminal } from "./attempt-terminal.test-helper.js"; import { nativeHookRelayUnregisterQueue } from "./native-hook-relay-state.js"; +import { createCodexNativeHookRelay } from "./native-hook-relay.js"; import type { CodexServerNotification } from "./protocol.js"; import { createParams, @@ -35,6 +37,125 @@ describe("runCodexAppServerAttempt native hook relay retention", () => { vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] }); }); + it.each([undefined, "callback cancelled"])( + "releases abandoned admission capacity while preserving duplicate waiters and retained children (reason: %s)", + async (abortReason) => { + const host = await createAdmittedHostCapabilityTestFixture({ + runId: "admission-cancellation", + }); + const source = new AbortController(); + const admissionWaits = new Map[]>(); + const register = relayRuntime.registerNativeHookRelayForBundledRuntime; + vi.spyOn(relayRuntime, "registerNativeHookRelayForBundledRuntime").mockImplementation( + (params) => { + const retention = params.retention; + if (!retention?.awaitForegroundAdmission) { + throw new Error("fixture admission missing"); + } + const admit = retention.awaitForegroundAdmission; + return register({ + ...params, + retention: { + ...retention, + awaitForegroundAdmission: (child, signal) => { + const waiting = admit(child, signal); + void waiting.catch(() => undefined); + admissionWaits.set(child, [...(admissionWaits.get(child) ?? []), waiting]); + return waiting; + }, + }, + }); + }, + ); + const relay = createCodexNativeHookRelay({ + options: { enabled: true }, + events: ["pre_tool_use"], + agentId: undefined, + sessionId: "admission-cancellation", + sessionKey: undefined, + config: {}, + runId: "admission-cancellation", + attemptTimeoutMs: 30_000, + startupTimeoutMs: 1_000, + turnStartTimeoutMs: 1_000, + loopDetectionPreToolUseRelay: false, + signal: source.signal, + hostCapabilities: host.hostCapabilities, + onPreToolUseFailure: () => {}, + }); + if (!relay) { + throw new Error("fixture relay missing"); + } + const pending: Promise[] = []; + const invoke = (child: string, signal?: AbortSignal) => { + const invocation = invokeNativeHookRelay( + { + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "pre_tool_use", + rawPayload: { + agent_id: child, + tool_name: "Bash", + tool_input: { command: "echo fixture" }, + }, + }, + signal, + ); + void invocation.catch(() => undefined); + pending.push(invocation); + return invocation; + }; + let releaseChild: (() => void) | undefined; + try { + await relay.ready; + const firstAbort = new AbortController(); + const duplicateAbort = new AbortController(); + const first = invoke("pending-0", firstAbort.signal); + const duplicate = invoke("pending-0", duplicateAbort.signal); + for (let i = 1; i < 32; i++) { + void invoke(`pending-${i}`); + } + // The rejected 33rd distinct child proves all earlier waits reached admission. + await expect(invoke("over-capacity")).rejects.toThrow("capacity reached"); + expect(admissionWaits.get("pending-0")).toHaveLength(2); + firstAbort.abort(abortReason); + await expect(first).rejects.toThrow(/abort/i); + await expect(admissionWaits.get("pending-0")![0]).rejects.toBeInstanceOf(Error); + await expect(invoke("still-over-capacity")).rejects.toThrow("capacity reached"); + duplicateAbort.abort(); + await expect(duplicate).rejects.toThrow(/abort/i); + await Promise.allSettled(admissionWaits.get("pending-0")!); + const replacement = invoke("replacement-child"); + await expect(invoke("capacity-control")).rejects.toThrow("capacity reached"); + releaseChild = relay.claimDirectChild("replacement-child"); + await expect(replacement).resolves.toEqual({ stdout: "", stderr: "", exitCode: 0 }); + relay.authorizeRetentionAfterSuccessfulYield(); + relay.unregister(); + await expect(invoke("replacement-child")).resolves.toEqual({ + stdout: "", + stderr: "", + exitCode: 0, + }); + releaseChild(); + releaseChild = undefined; + await expect(invoke("replacement-child")).rejects.toThrow(/not found|inactive/); + expect( + nativeHookRelayTesting.getNativeHookRelayRegistrationForTests(relay.relayId), + ).toBeUndefined(); + } finally { + releaseChild?.(); + relay.unregister(); + source.abort(); + await Promise.allSettled(pending); + await relay.drain(); + await nativeHookRelayUnregisterQueue.flush(); + host.closeHost(); + host.closeAdmission(); + } + }, + ); + it.each([ { name: "Codex multi-agent V1", diff --git a/extensions/codex/src/app-server/version.ts b/extensions/codex/src/app-server/version.ts index 6716ca3c52bc..3b9bcd60f7bd 100644 --- a/extensions/codex/src/app-server/version.ts +++ b/extensions/codex/src/app-server/version.ts @@ -2,7 +2,7 @@ * Version and package pins for the managed Codex app-server runtime. */ /** Exact Codex app-server version shipped by the OpenClaw Codex bridge. */ -export const CODEX_APP_SERVER_VERSION = "0.153.4"; +export const CODEX_APP_SERVER_VERSION = "0.154.0"; /** Inclusive runtime compatibility floor for external app-server binaries. */ export const MIN_SUPPORTED_CODEX_APP_SERVER_VERSION = "0.149.0"; /** npm package name for the managed Codex app-server binary. */ diff --git a/extensions/openai/openai-provider.ts b/extensions/openai/openai-provider.ts index f935769c4e54..9e9f61c39cb8 100644 --- a/extensions/openai/openai-provider.ts +++ b/extensions/openai/openai-provider.ts @@ -92,7 +92,7 @@ function classifyOpenAiFailoverCode(code: string | undefined) { const OPENAI_MODELS_ENDPOINT = "https://api.openai.com/v1/models"; // Keep synchronized with extensions/codex's exact @openai/codex dependency; // the provider contract test fails when that managed-runtime pin changes. -const OPENAI_CODEX_CLIENT_VERSION = "0.153.4"; +const OPENAI_CODEX_CLIENT_VERSION = "0.154.0"; const OPENAI_CODEX_MODELS_ENDPOINT = `${OPENAI_CODEX_RESPONSES_BASE_URL}/models?client_version=${OPENAI_CODEX_CLIENT_VERSION}`; const OPENAI_MODELS_CACHE_TTL_MS = 60_000; const OPENAI_CODEX_MODELS_CACHE_TTL_MS = 60_000; diff --git a/extensions/qa-lab/src/codex-native-hook-pressure.e2e.test.ts b/extensions/qa-lab/src/codex-native-hook-pressure.e2e.test.ts new file mode 100644 index 000000000000..463a61cc6364 --- /dev/null +++ b/extensions/qa-lab/src/codex-native-hook-pressure.e2e.test.ts @@ -0,0 +1,628 @@ +import { execFile } from "node:child_process"; +import { createHash, randomUUID } from "node:crypto"; +import { once } from "node:events"; +import fs from "node:fs/promises"; +import { createServer } from "node:http"; +import os from "node:os"; +import path from "node:path"; +import { setTimeout as sleep } from "node:timers/promises"; +import { promisify } from "node:util"; +import { extractErrorCode } from "openclaw/plugin-sdk/error-runtime"; +import { useAutoCleanupTempDirTracker } from "openclaw/plugin-sdk/test-env"; +import { afterEach, describe, expect, it } from "vitest"; +import { createQaGatewayChild } from "./gateway-child.js"; +import { attachQaMockResponsesWebSocketServer } from "./providers/mock-openai/mock-openai-responses-websocket.js"; +import { MockResponseStream } from "./providers/mock-openai/mock-openai-stream.js"; +import { listMockCodexModelInfos } from "./providers/shared/mock-model-config.js"; + +const execFileAsync = promisify(execFile); +const REPO_ROOT = path.resolve(import.meta.dirname, "../../.."); +const PLUGIN_ID = "qa-native-hook-pressure"; +const tempDirs = useAutoCleanupTempDirTracker(afterEach); + +type Tool = { + type?: string; + name?: string; + tools?: Tool[]; + parameters?: { properties?: Record }; +}; +type Scenario = { + id: string; + count: number; + mode: "allow" | "deny"; + issued: boolean; + outputs: unknown[]; +}; +type ProcessIdentity = { pid: number; startTimeTicks: number }; +type ProcessRow = ProcessIdentity & { + state: string; + comm: string; + argv: string[]; + rss: number; + ticks: number; +}; + +const processKey = (row: ProcessIdentity) => `${row.pid}:${row.startTimeTicks}`; +const isRelay = (row: ProcessRow) => + row.comm === "openclaw-hooks" || + row.argv.some((arg) => arg.endsWith("/native-hook-relay/entry.js")); + +async function readProcess(pid: number): Promise { + try { + const [stat, status, cmdline] = await Promise.all([ + fs.readFile(`/proc/${pid}/stat`, "utf8"), + fs.readFile(`/proc/${pid}/status`, "utf8"), + fs.readFile(`/proc/${pid}/cmdline`, "utf8"), + ]); + const end = stat.lastIndexOf(")"); + const fields = stat.slice(end + 2).split(" "); + return { + pid, + startTimeTicks: Number(fields[19]), + state: fields[0]!, + comm: stat.slice(stat.indexOf("(") + 1, end), + argv: cmdline.split("\0"), + rss: Number(/VmRSS:\s+(\d+)/.exec(status)?.[1] ?? 0) * 1024, + ticks: Number(fields[11]) + Number(fields[12]), + }; + } catch (error) { + const code = extractErrorCode(error); + if (code !== "ENOENT" && code !== "ESRCH") { + throw error; + } + return undefined; + } +} + +async function processTree(root: number): Promise { + const rows: ProcessRow[] = []; + const pending = [root]; + const seen = new Set(); + while (pending.length) { + const pid = pending.pop()!; + if (seen.has(pid)) { + continue; + } + seen.add(pid); + const row = await readProcess(pid); + if (!row) { + continue; + } + rows.push(row); + try { + // Rust can spawn from any executor thread, not only the process leader. + for (const tid of await fs.readdir(`/proc/${pid}/task`)) { + const children = await fs + .readFile(`/proc/${pid}/task/${tid}/children`, "utf8") + .catch((error: unknown) => { + const code = extractErrorCode(error); + if (code === "ENOENT" || code === "ESRCH") { + return ""; + } + throw error; + }); + pending.push(...children.trim().split(/\s+/).filter(Boolean).map(Number)); + } + } catch (error) { + const code = extractErrorCode(error); + if (code !== "ENOENT" && code !== "ESRCH") { + throw error; + } + } + } + return rows; +} + +async function inspectOwnedRelays(root: number, observed: Map) { + // Keep observed identities after reparenting; also discover any remaining descendants. + for (const row of await processTree(root)) { + if (isRelay(row)) { + observed.set(processKey(row), { pid: row.pid, startTimeTicks: row.startTimeTicks }); + } + } + const live: ProcessIdentity[] = []; + const zombies: ProcessIdentity[] = []; + for (const identity of observed.values()) { + const row = await readProcess(identity.pid); + if (!row || row.startTimeTicks !== identity.startTimeTicks) { + continue; + } + (row.state === "Z" ? zombies : live).push(identity); + } + return { live, zombies }; +} + +function toolsIn(body: Record): Array<{ tool: Tool; namespace?: string }> { + const entries: Array<{ tool: Tool; namespace?: string }> = []; + for (const tool of (Array.isArray(body.tools) ? body.tools : []) as Tool[]) { + if (tool.type === "namespace") { + for (const child of tool.tools ?? []) { + entries.push({ tool: child, namespace: tool.name }); + } + } else { + entries.push({ tool }); + } + } + return entries; +} + +describe.skipIf(process.platform !== "linux")( + "Codex native-hook pressure real Gateway proof", + () => { + it.each(["matched", "unmatched", "none"] as const)( + "records %s policy work through the real Gateway", + async (selection) => { + expect(process.platform).toBe("linux"); + const root = tempDirs.make("openclaw-native-hook-pressure-"); + const pluginDir = path.join(root, "plugin"); + await fs.mkdir(pluginDir); + await fs.writeFile( + path.join(pluginDir, "openclaw.plugin.json"), + JSON.stringify({ + id: PLUGIN_ID, + activation: { onStartup: true }, + configSchema: { type: "object", additionalProperties: false, properties: {} }, + }), + ); + await fs.writeFile( + path.join(pluginDir, "index.js"), + ` +import { hasBeforeToolCallPolicy, nativeHookRelayTesting } from "openclaw/plugin-sdk/agent-harness-runtime"; +// Plugin generation modules have distinct instances; the fixture observer is process-owned. +const key = Symbol.for("openclaw.test.native-hook-pressure.${randomUUID()}"); +const calls = globalThis[key] ??= []; +export default { + id: "${PLUGIN_ID}", + register(api) { + if ("${selection}" !== "none") api.on("before_tool_call", async (event) => { + const serialized = JSON.stringify(event.params); + calls.push({ toolName: event.toolName, denied: serialized.includes("PRESSURE_DENIED") }); + await new Promise((resolve) => setTimeout(resolve, 75)); + return serialized.includes("PRESSURE_DENIED") + ? { block: true, blockReason: "PRESSURE_POLICY_DENIED" } + : undefined; + }, { matcher: ["${selection === "unmatched" ? "web_fetch" : "exec"}"] }); + api.registerHttpRoute({ + path: "/qa/native-hook-pressure", auth: "gateway", match: "exact", + gatewayRuntimeScopeSurface: "trusted-operator", + async handler(_req, res) { res.setHeader("Content-Type", "application/json"); res.end(JSON.stringify({ calls, pid: process.pid, nodeVersion: process.version, hasPolicy: hasBeforeToolCallPolicy(), invocations: nativeHookRelayTesting.getNativeHookRelayInvocationsForTests() })); return true; } + }); + } +};`, + ); + let scenario: Scenario | undefined; + const wireTools = new Map(); + const dispatch = async ({ body }: { body: Record }) => { + if (!scenario) { + throw new Error("provider request outside scenario"); + } + const current = scenario; + const stream = new MockResponseStream(`resp_${current.id}_${randomUUID()}`); + const tools = toolsIn(body); + for (const { tool } of tools) { + wireTools.set(tool.name ?? "unknown", tool); + } + const input = Array.isArray(body.input) ? body.input : []; + current.outputs.push( + ...input.filter( + (item) => item && typeof item === "object" && String(item.type).endsWith("_output"), + ), + ); + if (!current.issued) { + current.issued = true; + const native = tools.find(({ tool }) => + ["exec_command", "shell_command", "shell"].includes(tool.name ?? ""), + ); + if (!native) { + throw new Error( + `native shell tool unavailable: ${JSON.stringify(tools.map(({ tool }) => ({ name: tool.name, type: tool.type })))}`, + ); + } + const properties = native.tool.parameters?.properties ?? {}; + expect(Object.hasOwn(properties, "login")).toBe(true); + for (let i = 0; i < current.count; i++) { + const command = + current.mode === "deny" + ? "printf PRESSURE_DENIED > pressure-denied.txt" + : `printf PRESSURE_ALLOW_${current.id}_${i}_END; printf PRESSURE_ALLOW_${current.id}_${i}_END > pressure-${current.id}-${i}.txt`; + // Host login profiles can leave the granted workspace before the command runs. + // Use the advertised non-login option equally for every measured scenario. + const args = { + ...(Object.hasOwn(properties, "cmd") + ? { cmd: command } + : { command: native.tool.name === "shell" ? ["sh", "-c", command] : command }), + login: false, + }; + stream.tool({ + type: "function_call", + id: `fc_${current.id}_${i}`, + call_id: `call_${current.id}_${i}`, + name: native.tool.name!, + ...(native.namespace ? { namespace: native.namespace } : {}), + arguments: JSON.stringify(args), + }); + } + } else { + stream.message({ id: `msg_${current.id}`, text: `PRESSURE_DONE_${current.id}` }); + } + return { + events: stream.complete(16), + model: typeof body.model === "string" ? body.model : "", + }; + }; + const server = createServer((req, res) => { + const handle = async () => { + if (req.method === "GET" && req.url?.startsWith("/v1/models")) { + res.setHeader("Content-Type", "application/json"); + res.end( + JSON.stringify({ + data: [{ id: "gpt-5.6-luna", object: "model" }], + models: listMockCodexModelInfos(), + }), + ); + return; + } + if (req.method !== "POST" || req.url !== "/v1/responses") { + res.writeHead(404).end(); + return; + } + const chunks: Buffer[] = []; + for await (const chunk of req) { + chunks.push(Buffer.from(chunk)); + } + const body = JSON.parse(Buffer.concat(chunks).toString()) as Record; + const { events } = await dispatch({ body }); + res.writeHead(200, { "Content-Type": "text/event-stream" }); + res.end( + events.map((event) => `data: ${JSON.stringify(event)}\n\n`).join("") + + "data: [DONE]\n\n", + ); + }; + void handle().catch((error: unknown) => { + res + .writeHead(500, { "Content-Type": "application/json" }) + .end(JSON.stringify({ error: String(error) })); + }); + }); + const sockets = attachQaMockResponsesWebSocketServer({ server, dispatch }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const address = server.address(); + if (!address || typeof address === "string") { + throw new Error("mock provider failed to bind"); + } + const owner = createQaGatewayChild(); + const reports: unknown[] = []; + const observedRelays = new Map(); + let gatewayPid: number | undefined; + try { + const gateway = await owner.start({ + repoRoot: REPO_ROOT, + command: { + executablePath: process.execPath, + argsPrefix: [path.join(REPO_ROOT, "dist/index.js")], + cwd: REPO_ROOT, + usePackagedPlugins: true, + }, + providerBaseUrl: `http://127.0.0.1:${address.port}/v1`, + transportBaseUrl: "", + providerMode: "mock-openai", + primaryModel: "mock-openai/gpt-5.6-luna", + alternateModel: "mock-openai/gpt-5.6-luna-alt", + forcedRuntime: "codex", + controlUiEnabled: false, + mutateConfig: (config) => ({ + ...config, + plugins: { + ...config.plugins, + allow: [...(config.plugins?.allow ?? []), PLUGIN_ID], + load: { paths: [pluginDir] }, + entries: { + ...config.plugins?.entries, + codex: { + enabled: true, + config: { + appServer: { + // Keep the QA sandbox; isolate policy work from optional loop detection. + sandbox: "workspace-write", + loopDetectionPreToolUseRelay: false, + }, + }, + }, + [PLUGIN_ID]: { enabled: true }, + }, + }, + }), + }); + if (!gateway.pid) { + throw new Error("Gateway has no PID"); + } + gatewayPid = gateway.pid; + const readCalls = async () => + ( + (await ( + await fetch(`${gateway.baseUrl}/qa/native-hook-pressure`, { + headers: { Authorization: `Bearer ${gateway.token}` }, + }) + ).json()) as { calls: Array<{ denied: boolean }> } + ).calls; + const registry = (await ( + await fetch(`${gateway.baseUrl}/qa/native-hook-pressure`, { + headers: { Authorization: `Bearer ${gateway.token}` }, + }) + ).json()) as { hasPolicy: boolean; pid: number; nodeVersion: string }; + expect(registry.hasPolicy).toBe(selection !== "none"); + expect(registry.pid).toBe(gateway.pid); + const cpuTickRate = Number((await execFileAsync("getconf", ["CLK_TCK"])).stdout.trim()); + expect(cpuTickRate).toBeGreaterThan(0); + const status = await fs.readFile(`/proc/${gateway.pid}/status`, "utf8"); + console.log( + "NATIVE_HOOK_HARDWARE " + + JSON.stringify({ + node: registry.nodeVersion, + kernel: os.release(), + cpuModel: os.cpus()[0]?.model, + logicalCpus: os.cpus().length, + availableParallelism: os.availableParallelism(), + affinity: /Cpus_allowed_list:\s*(.+)/.exec(status)?.[1], + cpuTickRate, + compileCacheConfigured: Boolean(process.env.NODE_COMPILE_CACHE), + compileCacheDisabled: process.env.NODE_DISABLE_COMPILE_CACHE === "1", + }), + ); + for (const count of selection === "matched" ? [1, 5, 20, 1] : [5]) { + const mode = selection === "matched" && reports.length === 3 ? "deny" : "allow"; + scenario = { id: randomUUID(), count, mode, issued: false, outputs: [] }; + const callsBefore = (await readCalls()).length; + const monitoring = new AbortController(); + let peakRelays = 0; + let peakRelayRss = 0; + let peakTreeRss = 0; + const relayIdentities = new Set(); + const relaySamples = new Map< + string, + { + firstAtMs: number; + lastAtMs: number; + peakRss: number; + firstTicks: number; + lastTicks: number; + } + >(); + let gatewayFirstTicks: number | undefined; + let gatewayLastTicks: number | undefined; + const healthMs: number[] = []; + const healthErrors: string[] = []; + const processSampler = (async () => { + while (!monitoring.signal.aborted) { + const tree = await processTree(gateway.pid!); + const relays = tree.filter(isRelay); + const sampledAtMs = performance.now(); + const gatewayRow = tree.find((row) => row.pid === gateway.pid); + gatewayFirstTicks ??= gatewayRow?.ticks; + gatewayLastTicks = gatewayRow?.ticks ?? gatewayLastTicks; + for (const relay of relays) { + const key = processKey(relay); + observedRelays.set(key, { pid: relay.pid, startTimeTicks: relay.startTimeTicks }); + relayIdentities.add(key); + const previous = relaySamples.get(key); + relaySamples.set(key, { + firstAtMs: previous?.firstAtMs ?? sampledAtMs, + lastAtMs: sampledAtMs, + peakRss: Math.max(previous?.peakRss ?? 0, relay.rss), + firstTicks: previous?.firstTicks ?? relay.ticks, + lastTicks: relay.ticks, + }); + } + peakRelays = Math.max(peakRelays, relays.filter((row) => row.state !== "Z").length); + peakRelayRss = Math.max( + peakRelayRss, + relays.filter((row) => row.state !== "Z").reduce((sum, row) => sum + row.rss, 0), + ); + peakTreeRss = Math.max( + peakTreeRss, + tree.filter((row) => row.state !== "Z").reduce((sum, row) => sum + row.rss, 0), + ); + await sleep(20); + } + })(); + const healthSampler = (async () => { + while (!monitoring.signal.aborted) { + const started = performance.now(); + try { + await gateway.call("health", {}, { timeoutMs: 5_000 }); + healthMs.push(performance.now() - started); + } catch (error) { + healthErrors.push(String(error)); + } + await sleep(25); + } + })(); + // Observe both failures immediately and join both samplers before retiring the Gateway. + const observers = Promise.allSettled([processSampler, healthSampler]); + const failures: unknown[] = []; + const started = performance.now(); + let turnElapsedMs = 0; + try { + const turn = (await gateway.call("chat.send", { + sessionKey: `agent:qa:pressure-${scenario.id}`, + message: "Run the bounded native shell pressure fixture.", + deliver: false, + idempotencyKey: randomUUID(), + })) as { runId: string; status: string }; + expect(turn.status).toBe("started"); + const terminal = (await gateway.call( + "agent.wait", + { runId: turn.runId, timeoutMs: 90_000 }, + { timeoutMs: 95_000 }, + )) as { status: string }; + turnElapsedMs = Math.round(performance.now() - started); + expect(terminal.status, gateway.logs()).toBe("ok"); + } catch (error) { + failures.push(error); + } finally { + monitoring.abort(); + for (const result of await observers) { + if (result.status === "rejected") { + failures.push(result.reason); + } + } + } + if (failures.length > 0) { + throw new AggregateError(failures, "native hook pressure scenario failed"); + } + const calls = (await readCalls()).slice(callsBefore); + const expectedCalls = selection === "matched" ? count : 0; + if (calls.length !== expectedCalls) { + console.log( + "NATIVE_HOOK_PRESSURE_FAILURE " + + JSON.stringify({ + scenario, + tools: [...wireTools.values()].map(({ name, type }) => ({ name, type })), + peakRelays, + registry: await ( + await fetch(`${gateway.baseUrl}/qa/native-hook-pressure`, { + headers: { Authorization: `Bearer ${gateway.token}` }, + }) + ).json(), + logs: gateway.logs(), + }), + ); + } + expect(calls).toHaveLength(expectedCalls); + expect(healthErrors).toEqual([]); + if (selection !== "matched") { + expect(peakRelays).toBe(0); + } + expect(calls.every((call) => call.denied === (mode === "deny"))).toBe(true); + if (mode === "deny") { + await expect( + fs.access(path.join(gateway.workspaceDir, "pressure-denied.txt")), + ).rejects.toMatchObject({ code: "ENOENT" }); + expect(JSON.stringify(scenario.outputs)).toContain("PRESSURE_POLICY_DENIED"); + } else { + for (let i = 0; i < count; i++) { + const output = scenario.outputs.find( + (value) => + value !== null && + typeof value === "object" && + "call_id" in value && + value.call_id === `call_${scenario!.id}_${i}`, + ); + expect(output).toBeDefined(); + const marker = `PRESSURE_ALLOW_${scenario.id}_${i}_END`; + console.log( + "NATIVE_HOOK_TOOL_RESULT " + + JSON.stringify({ output, workspaceDir: gateway.workspaceDir }), + ); + try { + expect(output).toMatchObject({ + output: expect.stringContaining("\nProcess exited with code 0\n"), + }); + expect(JSON.stringify(output)).toContain(marker); + const content = await fs.readFile( + path.join(gateway.workspaceDir, `pressure-${scenario.id}-${i}.txt`), + "utf8", + ); + expect(content).toBe(marker); + } catch (error) { + console.log( + "NATIVE_HOOK_PRESSURE_FAILURE " + + JSON.stringify({ + output, + workspaceDir: gateway.workspaceDir, + gatewayLogs: gateway.logs(), + }), + ); + throw error; + } + } + } + const binaryIdentities: Array = + []; + for (const row of await processTree(gateway.pid)) { + if (row.state === "Z" || !row.argv.includes("app-server")) { + continue; + } + const executable = `/proc/${row.pid}/exe`; + if (path.basename(await fs.readlink(executable)) !== "codex") { + continue; + } + const sha256 = createHash("sha256") + .update(await fs.readFile(executable)) + .digest("hex"); + const version = ( + await execFileAsync(executable, ["--version"], { timeout: 10_000 }) + ).stdout.trim(); + const current = await readProcess(row.pid); + expect(current?.startTimeTicks).toBe(row.startTimeTicks); + expect(version).toBe("codex-cli 0.154.0"); + binaryIdentities.push({ + pid: row.pid, + startTimeTicks: row.startTimeTicks, + sha256, + version, + }); + } + expect(binaryIdentities.length).toBeGreaterThan(0); + await expect + .poll(() => inspectOwnedRelays(gateway.pid!, observedRelays), { + timeout: 3_000, + interval: 25, + }) + .toEqual({ live: [], zombies: [] }); + healthMs.sort((a, b) => a - b); + reports.push({ + selection, + count, + mode, + binaries: binaryIdentities, + syntheticPolicyDelayMs: 75, + sampling: + "Non-atomic /proc snapshots; RSS can count shared pages; CPU omits work outside samples; elapsed windows are observation intervals, not process lifetimes.", + // /proc sampling misses short-lived children and shares pages between RSS values. + sampledRelayWindows: [...relaySamples.values()].map((sample) => ({ + observedWindowMs: sample.lastAtMs - sample.firstAtMs, + peakRss: sample.peakRss, + observedCpuMs: ((sample.lastTicks - sample.firstTicks) * 1000) / cpuTickRate, + })), + observedGatewayCpuMs: + gatewayFirstTicks === undefined || gatewayLastTicks === undefined + ? undefined + : ((gatewayLastTicks - gatewayFirstTicks) * 1000) / cpuTickRate, + turnElapsedMs, + policyCalls: calls.length, + uniqueSampledRelayProcesses: relayIdentities.size, + peakRelays, + peakRelayRss, + peakTreeRss, + healthSamples: healthMs.length, + healthErrors, + healthMedianMs: healthMs[Math.floor(healthMs.length / 2)], + healthP95Ms: healthMs[Math.floor(healthMs.length * 0.95)], + healthMaxMs: healthMs.at(-1), + remainingObservedOrDescendantRelays: { live: 0, zombies: 0 }, + }); + console.log("NATIVE_HOOK_PRESSURE " + JSON.stringify(reports.at(-1))); + } + } finally { + const stopped = await owner.stop(); + await sockets.close(); + await new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }); + if (gatewayPid) { + await expect + .poll(() => inspectOwnedRelays(gatewayPid!, observedRelays), { + timeout: 3_000, + interval: 25, + }) + .toEqual({ live: [], zombies: [] }); + console.log("NATIVE_HOOK_CLEANUP " + JSON.stringify({ live: [], zombies: [] })); + } + expect(stopped.errors).toEqual([]); + } + }, + 480_000, + ); + }, +); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 4a60d3990a86..87267c0c7b63 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -108,7 +108,7 @@ settings: overrides: '@lancedb/lancedb>@huggingface/transformers': '-' baileys>sharp: '-' - '@agentclientprotocol/codex-acp@1.10.0>@openai/codex': 0.153.4 + '@agentclientprotocol/codex-acp@1.10.0>@openai/codex': 0.154.0 '@codemirror/commands@6.11.0>@codemirror/view': 6.43.11 '@anthropic-ai/sdk': 0.124.0 '@anthropic-ai/vertex-sdk@0.19.7>google-auth-library': 11.0.2 @@ -823,8 +823,8 @@ importers: extensions/codex: dependencies: '@openai/codex': - specifier: 0.153.4 - version: 0.153.4 + specifier: 0.154.0 + version: 0.154.0 semver: specifier: 7.8.5 version: 7.8.5 @@ -4121,43 +4121,43 @@ packages: resolution: {integrity: sha512-3zcN5Q3yEmeyxXBzqB6fXPQFzYa2ROsGFSr69W0ArXIAGJqxl/aFECOVPD2kbkYPm0U/EHxFKgclK3UA9WQg5A==} engines: {node: ^22.22.2 || ^24.15.0 || >=26.0.0} - '@openai/codex@0.153.4': - resolution: {integrity: sha512-wbHDmit7S/YvBGVX1DQmk13xtWblZ2cApeJ/pB7xDZ10Cna+DZc5ij7f0F4OxdsXN4FW1oLT48OpogUI1+8Y2w==} + '@openai/codex@0.154.0': + resolution: {integrity: sha512-FV/x1OHXYv/ifjf3mXj9ThTTAWcUZN6cGIRQRhRxkKNOPuImu1WW0c8ev1vUkE9XGH90dEnYG1tBjIkxRikg0w==} engines: {node: '>=16'} hasBin: true - '@openai/codex@0.153.4-darwin-arm64': - resolution: {integrity: sha512-B1qhN3fa1ay0R0wGziXqgwSkB5icpYChNKHhtBHff/0UtSTC7z+l8aTtvMlGjH3E8HEvY3+njIJelM9CAAoVWg==} + '@openai/codex@0.154.0-darwin-arm64': + resolution: {integrity: sha512-HP/vJCH/t2hB9Kg6hotN9UglClJ6/z584fal5lEP14C9gNAgAQS4/kTQC7l5V+BA3TqwDPwINSjul28cX8AYXg==} engines: {node: '>=16'} cpu: [arm64] os: [darwin] - '@openai/codex@0.153.4-darwin-x64': - resolution: {integrity: sha512-vnSbbPzfoDZmmyzsxswsDDXQ06IVFBzkQU7/hroB3ji93Ok2utcsq8Psfk2tjF5r9mEx8RWFJhzuTGHG26/NDA==} + '@openai/codex@0.154.0-darwin-x64': + resolution: {integrity: sha512-2aqz+72Hop8PF2RYglQ4JnGjm3OlRIrTykJIT0hyLeUgM6NCFy09RgTmqRCoWliKQZjEn9jjZqUEp7QujAj77g==} engines: {node: '>=16'} cpu: [x64] os: [darwin] - '@openai/codex@0.153.4-linux-arm64': - resolution: {integrity: sha512-QKdjYLYV4hXIuUQDP3P6F4NXuWFoKo9WUoV4nAREIx55kiUyi8UsYdsVobkeXir5n/maEQgYMCKLHVma4rNPiw==} + '@openai/codex@0.154.0-linux-arm64': + resolution: {integrity: sha512-KmTCB6ST484zeYlPpKP/K5P/gRaYmt6TihVD+zotoe6O9q0JSBP+FYvCz4A/zZXR7xDOHURTSjHp0sD8wWS0YQ==} engines: {node: '>=16'} cpu: [arm64] os: [linux] - '@openai/codex@0.153.4-linux-x64': - resolution: {integrity: sha512-x1EcwBlY3AObM1VTUHNM2AzAJQsyreGdagpF+qFiYi/Oa30VBktvvG0C6tLtCzqW6hjZNWkGZQWmeVk7MuJKWg==} + '@openai/codex@0.154.0-linux-x64': + resolution: {integrity: sha512-a4FI3A8sGtwGrOqltrPbrS2hajrHQG591EwmRfiRoLMb10VxdBtUGW4gu6IJVYENiYGA7k3P4jlRHEoCZU/s9Q==} engines: {node: '>=16'} cpu: [x64] os: [linux] - '@openai/codex@0.153.4-win32-arm64': - resolution: {integrity: sha512-/FBh42976ltF1kxDoPQBg1Q6+hwChRU5/sm5dfeC8kFVQMvOCGoGeY5d8rRZGVJE8XojlXo74VQb0sHowcfgBw==} + '@openai/codex@0.154.0-win32-arm64': + resolution: {integrity: sha512-CRUmZnE0Y/a8aLMrrA681EytOGaPaF659wJAiI4I3hsbQjaeYBSPV7PkCjy4Qn5LR/fmwIUORVH+6JaBNQL+tw==} engines: {node: '>=16'} cpu: [arm64] os: [win32] - '@openai/codex@0.153.4-win32-x64': - resolution: {integrity: sha512-lMkB43kJZH0VFr+hoXc11qqR7QtQIbkr07ALgj4urKL1osNyUyuy1iXd3Vzz2iCYvBUCSw7I0l/W1cEPGx9euQ==} + '@openai/codex@0.154.0-win32-x64': + resolution: {integrity: sha512-Stg2KEJPIKVqPPR1wCverGOR4ey3RR3cvakR07w7FNKQUMzmHaOZomRsP2bR1qOT/67yHsks9rB+MCMfIWXcRA==} engines: {node: '>=16'} cpu: [x64] os: [win32] @@ -9983,25 +9983,21 @@ snapshots: '@agentclientprotocol/claude-agent-acp@0.75.1(@anthropic-ai/sdk@0.124.0(zod@4.5.4))(@modelcontextprotocol/sdk@1.30.0(supports-color@10.2.2)(zod@4.5.4))': dependencies: - '@agentclientprotocol/sdk': 1.4.0(zod@4.4.3) - '@anthropic-ai/claude-agent-sdk': 0.3.257(@anthropic-ai/sdk@0.124.0(zod@4.5.4))(@modelcontextprotocol/sdk@1.30.0(supports-color@10.2.2)(zod@4.5.4))(zod@4.4.3) - zod: 4.4.3 + '@agentclientprotocol/sdk': 1.4.0(zod@4.5.4) + '@anthropic-ai/claude-agent-sdk': 0.3.257(@anthropic-ai/sdk@0.124.0(zod@4.5.4))(@modelcontextprotocol/sdk@1.30.0(supports-color@10.2.2)(zod@4.5.4))(zod@4.5.4) + zod: 4.5.4 transitivePeerDependencies: - '@anthropic-ai/sdk' - '@modelcontextprotocol/sdk' '@agentclientprotocol/codex-acp@1.10.0': dependencies: - '@agentclientprotocol/sdk': 1.4.0(zod@4.4.3) - '@openai/codex': 0.153.4 + '@agentclientprotocol/sdk': 1.4.0(zod@4.5.4) + '@openai/codex': 0.154.0 diff: 9.0.0 open: 11.0.1 vscode-jsonrpc: 9.0.1 - zod: 4.4.3 - - '@agentclientprotocol/sdk@1.4.0(zod@4.4.3)': - dependencies: - zod: 4.4.3 + zod: 4.5.4 '@agentclientprotocol/sdk@1.4.0(zod@4.5.4)': dependencies: @@ -10036,11 +10032,11 @@ snapshots: '@anthropic-ai/claude-agent-sdk-win32-x64@0.3.257': optional: true - '@anthropic-ai/claude-agent-sdk@0.3.257(@anthropic-ai/sdk@0.124.0(zod@4.5.4))(@modelcontextprotocol/sdk@1.30.0(supports-color@10.2.2)(zod@4.5.4))(zod@4.4.3)': + '@anthropic-ai/claude-agent-sdk@0.3.257(@anthropic-ai/sdk@0.124.0(zod@4.5.4))(@modelcontextprotocol/sdk@1.30.0(supports-color@10.2.2)(zod@4.5.4))(zod@4.5.4)': dependencies: '@anthropic-ai/sdk': 0.124.0(zod@4.5.4) '@modelcontextprotocol/sdk': 1.30.0(supports-color@10.2.2)(zod@4.5.4) - zod: 4.4.3 + zod: 4.5.4 optionalDependencies: '@anthropic-ai/claude-agent-sdk-darwin-arm64': 0.3.257 '@anthropic-ai/claude-agent-sdk-darwin-x64': 0.3.257 @@ -11695,31 +11691,31 @@ snapshots: '@npmcli/redact@5.0.0': {} - '@openai/codex@0.153.4': + '@openai/codex@0.154.0': optionalDependencies: - '@openai/codex-darwin-arm64': '@openai/codex@0.153.4-darwin-arm64' - '@openai/codex-darwin-x64': '@openai/codex@0.153.4-darwin-x64' - '@openai/codex-linux-arm64': '@openai/codex@0.153.4-linux-arm64' - '@openai/codex-linux-x64': '@openai/codex@0.153.4-linux-x64' - '@openai/codex-win32-arm64': '@openai/codex@0.153.4-win32-arm64' - '@openai/codex-win32-x64': '@openai/codex@0.153.4-win32-x64' + '@openai/codex-darwin-arm64': '@openai/codex@0.154.0-darwin-arm64' + '@openai/codex-darwin-x64': '@openai/codex@0.154.0-darwin-x64' + '@openai/codex-linux-arm64': '@openai/codex@0.154.0-linux-arm64' + '@openai/codex-linux-x64': '@openai/codex@0.154.0-linux-x64' + '@openai/codex-win32-arm64': '@openai/codex@0.154.0-win32-arm64' + '@openai/codex-win32-x64': '@openai/codex@0.154.0-win32-x64' - '@openai/codex@0.153.4-darwin-arm64': + '@openai/codex@0.154.0-darwin-arm64': optional: true - '@openai/codex@0.153.4-darwin-x64': + '@openai/codex@0.154.0-darwin-x64': optional: true - '@openai/codex@0.153.4-linux-arm64': + '@openai/codex@0.154.0-linux-arm64': optional: true - '@openai/codex@0.153.4-linux-x64': + '@openai/codex@0.154.0-linux-x64': optional: true - '@openai/codex@0.153.4-win32-arm64': + '@openai/codex@0.154.0-win32-arm64': optional: true - '@openai/codex@0.153.4-win32-x64': + '@openai/codex@0.154.0-win32-x64': optional: true '@openclaw/crabline@0.1.21': diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 33a90f6a5079..cc7c7225ab57 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -40,7 +40,7 @@ overrides: # Memory LanceDB supplies vectors; WhatsApp prepares thumbnails with Rastermill. "@lancedb/lancedb>@huggingface/transformers": "-" "baileys>sharp": "-" - "@agentclientprotocol/codex-acp@1.10.0>@openai/codex": 0.153.4 + "@agentclientprotocol/codex-acp@1.10.0>@openai/codex": 0.154.0 "@codemirror/commands@6.11.0>@codemirror/view": 6.43.11 "@anthropic-ai/sdk": 0.124.0 "@anthropic-ai/vertex-sdk@0.19.7>google-auth-library": 11.0.2 diff --git a/scripts/e2e/lib/codex-media-path/fake-codex-app-server.mjs b/scripts/e2e/lib/codex-media-path/fake-codex-app-server.mjs index 218d0f2be4ac..500f5fd6ddc0 100644 --- a/scripts/e2e/lib/codex-media-path/fake-codex-app-server.mjs +++ b/scripts/e2e/lib/codex-media-path/fake-codex-app-server.mjs @@ -5,7 +5,7 @@ import { runFakeCodexAppServer, } from "../codex-app-server-fixture.mjs"; -const version = "0.153.4"; +const version = "0.154.0"; const requestLog = process.env.OPENCLAW_CODEX_MEDIA_PATH_APP_SERVER_LOG ?? "/tmp/openclaw-codex-media-path-app-server.jsonl"; diff --git a/src/agents/harness/native-hook-relay-approval-presentation.ts b/src/agents/harness/native-hook-relay-approval-presentation.ts new file mode 100644 index 000000000000..b4f7ed1b4272 --- /dev/null +++ b/src/agents/harness/native-hook-relay-approval-presentation.ts @@ -0,0 +1,75 @@ +import { stripAnsi } from "../../../packages/terminal-core/src/ansi.js"; +import type { + NativeHookRelayPermissionApprovalRequest, + NativeHookRelayProvider, +} from "./native-hook-relay-types.js"; +import { readOptionalNonEmptyString, truncateRelayText } from "./native-hook-relay-utils.js"; + +const MAX_APPROVAL_TITLE_LENGTH = 80; +const MAX_APPROVAL_DESCRIPTION_LENGTH = 700; + +export function formatNativeHookRelayApprovalPresentation( + request: NativeHookRelayPermissionApprovalRequest, +): { title: string; description: string } { + return { + title: truncateRelayText( + `${nativeHookRelayProviderDisplayName(request.provider)} permission request`, + MAX_APPROVAL_TITLE_LENGTH, + ), + description: truncateRelayText( + formatPermissionApprovalDescription(request), + MAX_APPROVAL_DESCRIPTION_LENGTH, + ), + }; +} + +export function formatPermissionApprovalDescription( + request: NativeHookRelayPermissionApprovalRequest, +): string { + const lines = [ + `Tool: ${sanitizeApprovalText(request.toolName)}`, + request.cwd ? `Cwd: ${sanitizeApprovalText(request.cwd)}` : undefined, + request.model ? `Model: ${sanitizeApprovalText(request.model)}` : undefined, + formatToolInputPreview(request.toolInput), + ].filter((line): line is string => Boolean(line)); + return lines.join("\n"); +} + +function formatToolInputPreview(toolInput: Record): string | undefined { + const command = readOptionalNonEmptyString(toolInput.command); + if (command) { + return `Command: ${truncateRelayText(sanitizeApprovalText(command), 240)}`; + } + const keys = Object.keys(toolInput).map(sanitizeApprovalText).filter(Boolean).toSorted(); + if (!keys.length) { + return undefined; + } + const shownKeys = keys.slice(0, 12).join(", "); + const omitted = keys.length > 12 ? ` (${keys.length - 12} omitted)` : ""; + return `Input keys: ${shownKeys}${omitted}`; +} + +function sanitizeApprovalText(value: string): string { + let sanitized = ""; + for (const char of stripAnsi(value)) { + const codePoint = char.codePointAt(0); + sanitized += codePoint != null && isUnsafeApprovalCodePoint(codePoint) ? " " : char; + } + return sanitized.replace(/\s+/g, " ").trim(); +} + +function isUnsafeApprovalCodePoint(codePoint: number): boolean { + return ( + (codePoint >= 0 && codePoint <= 8) || + codePoint === 11 || + codePoint === 12 || + (codePoint >= 14 && codePoint <= 31) || + (codePoint >= 127 && codePoint <= 159) || + (codePoint >= 0x202a && codePoint <= 0x202e) || + (codePoint >= 0x2066 && codePoint <= 0x2069) + ); +} + +function nativeHookRelayProviderDisplayName(provider: NativeHookRelayProvider): string { + return provider === "codex" ? "Codex" : provider; +} diff --git a/src/agents/harness/native-hook-relay-bridge.ts b/src/agents/harness/native-hook-relay-bridge.ts index f072ceabbfcd..3a6dff294c17 100644 --- a/src/agents/harness/native-hook-relay-bridge.ts +++ b/src/agents/harness/native-hook-relay-bridge.ts @@ -1,6 +1,7 @@ import { randomUUID } from "node:crypto"; import { createServer, type IncomingMessage, type ServerResponse } from "node:http"; import { hasErrnoCode } from "../../infra/errno.js"; +import { createHttpRequestAbortSignal } from "../../infra/http-request-lifecycle.js"; import { createSubsystemLogger } from "../../logging/subsystem.js"; import { createDeferredCore } from "../../shared/deferred.js"; import { isPidDefinitelyDead } from "../../shared/pid-alive.js"; @@ -44,6 +45,7 @@ const { relays, relayBridges, pendingOperations } = nativeHookRelayState; type InvokeNativeHookRelay = ( params: InvokeNativeHookRelayParams, + signal?: AbortSignal, ) => Promise; type NativeHookRelayBridgeRenewalResult = "renewed" | "unavailable" | "ownership-changed"; @@ -279,6 +281,7 @@ async function handleNativeHookRelayBridgeRequest( res: ServerResponse, auth: NativeHookRelayBridgeRequestAuth, ): Promise { + const requestAbort = createHttpRequestAbortSignal(req, res); try { if (req.method !== "POST" || req.url !== "/invoke") { writeNativeHookRelayBridgeJson(res, 404, { ok: false, error: "not found" }); @@ -311,14 +314,22 @@ async function handleNativeHookRelayBridgeRequest( }); return; } - const result = await auth.invokeRelay({ ...payload, requireGeneration: true }); + const result = await auth.invokeRelay( + { ...payload, requireGeneration: true }, + requestAbort.signal, + ); writeNativeHookRelayBridgeJson(res, 200, { ok: true, result }); } catch (error) { + if (requestAbort.signal.aborted) { + return; + } writeNativeHookRelayBridgeJson( res, isNativeHookRelayBridgeStaleRegistrationError(error) ? 410 : 500, { ok: false, error: error instanceof Error ? error.message : String(error) }, ); + } finally { + requestAbort.cleanup(); } } diff --git a/src/agents/harness/native-hook-relay-events.ts b/src/agents/harness/native-hook-relay-events.ts index 364f1f4a1ac7..ecddc4848207 100644 --- a/src/agents/harness/native-hook-relay-events.ts +++ b/src/agents/harness/native-hook-relay-events.ts @@ -156,6 +156,16 @@ async function runNativeHookRelayPreToolUse(params: { : {}), }, }); + try { + params.registration.signal?.throwIfAborted(); + params.registration.assertActive?.(); + } catch (error) { + // A disconnected request cannot leave an approval for a later tool to consume. + if (!outcome.blocked && outcome.deferredApproval) { + cancelDeferredPluginToolApproval(outcome.deferredApproval); + } + throw error; + } if (outcome.blocked) { return params.adapter.renderPreToolUseBlockResponse( outcome.reason, diff --git a/src/agents/harness/native-hook-relay-permissions.ts b/src/agents/harness/native-hook-relay-permissions.ts index 3316699c2f24..77bfbc6fd8aa 100644 --- a/src/agents/harness/native-hook-relay-permissions.ts +++ b/src/agents/harness/native-hook-relay-permissions.ts @@ -1,8 +1,7 @@ import { createHash } from "node:crypto"; import { resolveExpiresAtMsFromDurationMs } from "@openclaw/normalization-core/number-coercion"; -import { stripAnsi } from "../../../packages/terminal-core/src/ansi.js"; +import { racePromiseWithAbortSignal } from "../../infra/abort-signal.js"; import { isApprovalNotFoundError } from "../../infra/approval-errors.js"; -import { toErrorObject } from "../../infra/errors.js"; import { pruneMapToMaxSize } from "../../infra/map-size.js"; import { prepareSystemRunMutableFileBinding, @@ -18,6 +17,7 @@ import { } from "../agent-tools.before-tool-call.js"; import { formatMcpCodexApprovalRemedy } from "../mcp-codex-tool-approval.js"; import { callGatewayTool } from "../tools/gateway.js"; +import { formatNativeHookRelayApprovalPresentation } from "./native-hook-relay-approval-presentation.js"; import { nativeHookRelayParamsWereRewritten, normalizeNativeHookToolName, @@ -33,9 +33,9 @@ import type { NativeHookRelayPermissionApprovalRequest, NativeHookRelayPermissionApprovalRequester, NativeHookRelayPermissionApprovalResult, + NativeHookRelayPendingPermissionApproval, NativeHookRelayPreToolUseApproval, NativeHookRelayProcessResponse, - NativeHookRelayProvider, NativeHookRelayProviderAdapter, NativeHookRelayRegistration, } from "./native-hook-relay-types.js"; @@ -48,8 +48,6 @@ const PERMISSION_ALLOW_ALWAYS_TTL_MS = 30 * 60 * 1000; const MAX_PERMISSION_FALLBACK_KEYS = 200; const MAX_PERMISSION_FALLBACK_KEY_CHARS = 240; const MAX_PERMISSION_FINGERPRINT_SORT_KEYS = 200; -const MAX_APPROVAL_TITLE_LENGTH = 80; -const MAX_APPROVAL_DESCRIPTION_LENGTH = 700; const MAX_PERMISSION_APPROVALS_PER_WINDOW = 12; const PERMISSION_APPROVAL_WINDOW_MS = 60_000; const MAX_PERMISSION_ALLOW_ALWAYS_ENTRIES = 512; @@ -79,7 +77,7 @@ function nativeHookRelayPreToolUseApprovalKey(params: { toolUseId?: string; }): string | undefined { const toolUseId = params.toolUseId?.trim(); - return toolUseId ? `${params.relayId}:${toolUseId}` : undefined; + return toolUseId ? JSON.stringify([params.relayId, toolUseId]) : undefined; } export function setNativeHookRelayPreToolUseApproval(params: { @@ -93,34 +91,49 @@ export function setNativeHookRelayPreToolUseApproval(params: { return false; } const previousApproval = pendingPreToolUseApprovals.get(key); - if (previousApproval) { - cancelDeferredPluginToolApproval(previousApproval.deferredApproval); - } pendingPreToolUseApprovals.set(key, { + relayId: params.relayId, deferredApproval: params.deferredApproval, originalParamsFingerprint: params.originalParamsFingerprint, }); + let evictedApproval: NativeHookRelayPreToolUseApproval | undefined; if (pendingPreToolUseApprovals.size > MAX_NATIVE_HOOK_RELAY_INVOCATIONS) { const oldestKey = pendingPreToolUseApprovals.keys().next().value; if (oldestKey) { - const oldestApproval = pendingPreToolUseApprovals.get(oldestKey); - if (oldestApproval) { - cancelDeferredPluginToolApproval(oldestApproval.deferredApproval); - } + evictedApproval = pendingPreToolUseApprovals.get(oldestKey); pendingPreToolUseApprovals.delete(oldestKey); } } + // Publish/detach before notifying: cancellation callbacks may replace or + // retire this relay synchronously, and their successor must remain authoritative. + if (previousApproval) { + cancelDeferredPluginToolApproval(previousApproval.deferredApproval); + } + if (evictedApproval) { + cancelDeferredPluginToolApproval(evictedApproval.deferredApproval); + } return true; } -export function removeNativeHookRelayPreToolUseApprovals(relayId: string): void { - const prefix = `${relayId}:`; - for (const [key, pendingApproval] of pendingPreToolUseApprovals) { - if (key.startsWith(prefix)) { - cancelDeferredPluginToolApproval(pendingApproval.deferredApproval); +export function detachNativeHookRelayApprovalState(relayId: string): () => void { + const preToolUseApprovals: NativeHookRelayPreToolUseApproval[] = []; + for (const [key, approval] of pendingPreToolUseApprovals) { + if (approval.relayId === relayId) { pendingPreToolUseApprovals.delete(key); + preToolUseApprovals.push(approval); } } + const permissionApprovals = detachNativeHookRelayPermissionState(relayId); + // Detach every old entry before any callback can register a same-id successor. + // Completion finalizers still compare object identity before deleting entries. + return () => { + for (const approval of preToolUseApprovals) { + cancelDeferredPluginToolApproval(approval.deferredApproval); + } + for (const approval of permissionApprovals) { + approval.controller.abort(); + } + }; } export async function resolveNativeHookRelayDeferredToolApproval(params: { @@ -206,6 +219,8 @@ export async function runNativeHookRelayPermissionRequest(params: { }; const mcpServerName = /^mcp__(.+?)__/.exec(request.toolName)?.[1]; const mutableFileBinding = await prepareNativeHookMutableFileBinding(request); + // File preparation yields; a disconnected callback must not create a new approval. + params.registration.assertActive?.(); if (!mutableFileBinding.ok) { return params.adapter.renderPermissionDecisionResponse("deny", mutableFileBinding.message); } @@ -233,14 +248,12 @@ export async function runNativeHookRelayPermissionRequest(params: { } return params.adapter.renderPermissionDecisionResponse("allow"); } - const pendingApproval = pendingPermissionApprovals.get(approvalKey); try { - const decision = await (pendingApproval ?? - startNativeHookRelayPermissionApprovalWithBudget({ - registration: params.registration, - approvalKey, - request, - })); + const decision = await waitForNativeHookRelayPermissionApproval({ + registration: params.registration, + approvalKey, + request, + }); params.registration.assertActive?.(); if ((decision === "allow" || decision === "allow-always") && mutableFileBinding.binding) { // PermissionRequest is OpenClaw's last boundary before the native runtime @@ -273,6 +286,7 @@ export async function runNativeHookRelayPermissionRequest(params: { ); } } catch (error) { + params.registration.assertActive?.(); log.warn( `native hook permission approval failed; deferring to provider approval path: ${String(error)}`, ); @@ -282,25 +296,63 @@ export async function runNativeHookRelayPermissionRequest(params: { return params.adapter.renderNoopResponse(params.invocation.event); } -async function startNativeHookRelayPermissionApprovalWithBudget(params: { +async function waitForNativeHookRelayPermissionApproval(params: { registration: NativeHookRelayRegistration; approvalKey: string; request: NativeHookRelayPermissionApprovalRequest; }): Promise { - if (!consumeNativeHookRelayPermissionBudget(params.registration.relayId)) { - log.warn( - `native hook permission approval rate limit exceeded; deferring to provider approval path: relay=${params.registration.relayId} run=${params.registration.runId}`, - ); - return "defer"; + let approval = pendingPermissionApprovals.get(params.approvalKey); + if (!approval) { + if (!consumeNativeHookRelayPermissionBudget(params.registration.relayId)) { + log.warn( + `native hook permission approval rate limit exceeded; deferring to provider approval path: relay=${params.registration.relayId} run=${params.registration.runId}`, + ); + return "defer"; + } + const controller = new AbortController(); + const pending: NativeHookRelayPendingPermissionApproval = { + relayId: params.registration.relayId, + controller, + waiters: 0, + cancelWhenUnobserved: params.registration.approvalHost !== undefined, + promise: racePromiseWithAbortSignal( + Promise.resolve().then(() => { + controller.signal.throwIfAborted(); + const request = { + ...params.request, + signal: controller.signal, + }; + return params.registration.approvalHost + ? requestNativeHookRelayPermissionApproval(request, params.registration.approvalHost) + : nativeHookRelayPermissionApprovalRequester(request); + }), + controller.signal, + ).finally(() => { + if (pendingPermissionApprovals.get(params.approvalKey) === pending) { + pendingPermissionApprovals.delete(params.approvalKey); + } + }), + }; + pendingPermissionApprovals.set(params.approvalKey, pending); + approval = pending; + } + approval.waiters++; + try { + return await racePromiseWithAbortSignal(approval.promise, params.registration.signal); + } finally { + // A duplicate callback owns its wait, not the shared approval request. + // Only the admitted host can retract an accepted Gateway approval. Public + // callers keep dedup until decision/expiry so a retry rejoins that prompt. + approval.waiters--; + if ( + approval.cancelWhenUnobserved && + approval.waiters === 0 && + pendingPermissionApprovals.get(params.approvalKey) === approval + ) { + pendingPermissionApprovals.delete(params.approvalKey); + approval.controller.abort(); + } } - const approval: Promise = - nativeHookRelayPermissionApprovalRequester(params.request).finally(() => { - if (pendingPermissionApprovals.get(params.approvalKey) === approval) { - pendingPermissionApprovals.delete(params.approvalKey); - } - }); - pendingPermissionApprovals.set(params.approvalKey, approval); - return approval; } function nativeHookRelayPermissionApprovalKey(params: { @@ -308,15 +360,15 @@ function nativeHookRelayPermissionApprovalKey(params: { request: NativeHookRelayPermissionApprovalRequest; binding?: SystemRunMutableFileBinding; }): string { - return [ + return JSON.stringify([ params.registration.relayId, params.registration.runId, params.request.toolCallId - ? `call:${params.request.toolCallId}` - : permissionRequestFallbackKey(params.request), + ? ["call", params.request.toolCallId] + : ["fallback", permissionRequestFallbackKey(params.request)], permissionRequestContentFingerprint(params.request), params.binding ? permissionRequestBindingFingerprint(params.binding) : "no-file-binding", - ].join(":"); + ]); } async function prepareNativeHookMutableFileBinding( @@ -549,52 +601,67 @@ export function pruneNativeHookRelayPermissionAllowAlways(now = Date.now()): voi } } -export function removeNativeHookRelayPermissionState(relayId: string): void { +function detachNativeHookRelayPermissionState( + relayId: string, +): NativeHookRelayPendingPermissionApproval[] { + const approvals: NativeHookRelayPendingPermissionApproval[] = []; permissionApprovalWindows.delete(relayId); for (const [key, entry] of permissionAllowAlwaysApprovals) { if (entry.relayId === relayId) { permissionAllowAlwaysApprovals.delete(key); } } - for (const key of pendingPermissionApprovals.keys()) { - if (key.startsWith(`${relayId}:`)) { + for (const [key, approval] of pendingPermissionApprovals) { + if (approval.relayId === relayId) { pendingPermissionApprovals.delete(key); + approvals.push(approval); } } + return approvals; +} + +export function removeNativeHookRelayPermissionState(relayId: string): void { + for (const approval of detachNativeHookRelayPermissionState(relayId)) { + approval.controller.abort(); + } } async function requestNativeHookRelayPermissionApproval( request: NativeHookRelayPermissionApprovalRequest, + approvalHost?: NativeHookRelayRegistration["approvalHost"], ): Promise { const timeoutMs = DEFAULT_PERMISSION_TIMEOUT_MS; - const requestResult: { id?: string; decision?: string | null } = await callGatewayTool( - "plugin.approval.request", - { timeoutMs: timeoutMs + 10_000 }, - { - pluginId: `openclaw-native-hook-relay-${request.provider}`, - title: truncateRelayText( - `${nativeHookRelayProviderDisplayName(request.provider)} permission request`, - MAX_APPROVAL_TITLE_LENGTH, - ), - description: truncateRelayText( - formatPermissionApprovalDescription(request), - MAX_APPROVAL_DESCRIPTION_LENGTH, - ), - severity: "warning", - toolName: request.toolName, - toolCallId: request.toolCallId, - allowedDecisions: [ - PluginApprovalResolutions.ALLOW_ONCE, - PluginApprovalResolutions.ALLOW_ALWAYS, - PluginApprovalResolutions.DENY, - ], - agentId: request.agentId, - sessionKey: request.sessionKey, - timeoutMs, - twoPhase: true, - }, - { expectFinal: false }, - ); + const approvalRequest = { + ...formatNativeHookRelayApprovalPresentation(request), + severity: "warning" as const, + toolName: request.toolName, + toolCallId: request.toolCallId, + allowedDecisions: [ + PluginApprovalResolutions.ALLOW_ONCE, + PluginApprovalResolutions.ALLOW_ALWAYS, + PluginApprovalResolutions.DENY, + ], + timeoutMs, + }; + const requestResult = approvalHost + ? await approvalHost.requestApproval({ + ...approvalRequest, + transportTimeoutMs: timeoutMs + 10_000, + signal: request.signal, + }) + : await callGatewayTool<{ id?: string; decision?: string | null }>( + "plugin.approval.request", + { timeoutMs: timeoutMs + 10_000 }, + { + ...approvalRequest, + pluginId: `openclaw-native-hook-relay-${request.provider}`, + agentId: request.agentId, + sessionKey: request.sessionKey, + twoPhase: true, + }, + { expectFinal: false, signal: request.signal }, + ); + request.signal?.throwIfAborted(); const approvalId = requestResult?.id; if (!approvalId) { return "defer"; @@ -607,6 +674,7 @@ async function requestNativeHookRelayPermissionApproval( approvalId, signal: request.signal, timeoutMs, + approvalHost, }); // Bind the verdict to the request that parked this call. A stale or // misrouted reply must never release a different tool gate. @@ -631,94 +699,31 @@ async function waitForNativeHookRelayApprovalDecision(params: { approvalId: string; signal?: AbortSignal; timeoutMs: number; + approvalHost?: NativeHookRelayRegistration["approvalHost"]; }): Promise<{ id?: string; decision?: string | null } | undefined> { - const waitPromise: Promise<{ id?: string; decision?: string | null } | undefined> = - callGatewayTool( - "plugin.approval.waitDecision", - { timeoutMs: params.timeoutMs + 10_000 }, - { id: params.approvalId }, - ).catch((error: unknown) => { - if (isApprovalNotFoundError(error)) { - return undefined; - } - throw error; - }); - if (!params.signal) { - return waitPromise; - } - let onAbort: (() => void) | undefined; - const abortPromise = new Promise((_, reject) => { - if (params.signal!.aborted) { - reject(toErrorObject(params.signal!.reason, "Non-Error rejection")); - return; + const pending = params.approvalHost + ? params.approvalHost + .waitForApproval({ + approvalId: params.approvalId, + timeoutMs: params.timeoutMs, + transportTimeoutMs: params.timeoutMs + 10_000, + signal: params.signal, + }) + .then((result) => + result ? { id: params.approvalId, decision: result.decision } : undefined, + ) + : callGatewayTool<{ id?: string; decision?: string | null }>( + "plugin.approval.waitDecision", + { timeoutMs: params.timeoutMs + 10_000 }, + { id: params.approvalId }, + { signal: params.signal }, + ); + return pending.catch((error: unknown) => { + if (isApprovalNotFoundError(error)) { + return undefined; } - onAbort = () => reject(toErrorObject(params.signal!.reason, "Non-Error rejection")); - params.signal!.addEventListener("abort", onAbort, { once: true }); + throw error; }); - try { - return await Promise.race([waitPromise, abortPromise]); - } finally { - if (onAbort) { - params.signal.removeEventListener("abort", onAbort); - } - } -} - -export function formatPermissionApprovalDescriptionForTests( - request: NativeHookRelayPermissionApprovalRequest, -): string { - return formatPermissionApprovalDescription(request); -} - -function formatPermissionApprovalDescription( - request: NativeHookRelayPermissionApprovalRequest, -): string { - const lines = [ - `Tool: ${sanitizeApprovalText(request.toolName)}`, - request.cwd ? `Cwd: ${sanitizeApprovalText(request.cwd)}` : undefined, - request.model ? `Model: ${sanitizeApprovalText(request.model)}` : undefined, - formatToolInputPreview(request.toolInput), - ].filter((line): line is string => Boolean(line)); - return lines.join("\n"); -} - -function formatToolInputPreview(toolInput: Record): string | undefined { - const command = readOptionalNonEmptyString(toolInput.command); - if (command) { - return `Command: ${truncateRelayText(sanitizeApprovalText(command), 240)}`; - } - const keys = Object.keys(toolInput).map(sanitizeApprovalText).filter(Boolean).toSorted(); - if (!keys.length) { - return undefined; - } - const shownKeys = keys.slice(0, 12).join(", "); - const omitted = keys.length > 12 ? ` (${keys.length - 12} omitted)` : ""; - return `Input keys: ${shownKeys}${omitted}`; -} - -function sanitizeApprovalText(value: string): string { - let sanitized = ""; - for (const char of stripAnsi(value)) { - const codePoint = char.codePointAt(0); - sanitized += codePoint != null && isUnsafeApprovalCodePoint(codePoint) ? " " : char; - } - return sanitized.replace(/\s+/g, " ").trim(); -} - -function isUnsafeApprovalCodePoint(codePoint: number): boolean { - return ( - (codePoint >= 0 && codePoint <= 8) || - codePoint === 11 || - codePoint === 12 || - (codePoint >= 14 && codePoint <= 31) || - (codePoint >= 127 && codePoint <= 159) || - (codePoint >= 0x202a && codePoint <= 0x202e) || - (codePoint >= 0x2066 && codePoint <= 0x2069) - ); -} - -function nativeHookRelayProviderDisplayName(provider: NativeHookRelayProvider): string { - return provider === "codex" ? "Codex" : provider; } export function setNativeHookRelayPermissionApprovalRequesterForTests( @@ -734,6 +739,9 @@ export function setNativeHookRelayDeferredToolApprovalRequesterForTests( } export function clearNativeHookRelayPermissionsForTests(): void { + for (const approval of pendingPermissionApprovals.values()) { + approval.controller.abort(); + } pendingPermissionApprovals.clear(); for (const pendingApproval of pendingPreToolUseApprovals.values()) { cancelDeferredPluginToolApproval(pendingApproval.deferredApproval); diff --git a/src/agents/harness/native-hook-relay-types.ts b/src/agents/harness/native-hook-relay-types.ts index 8c8bab66573b..1eaaf52a2ffc 100644 --- a/src/agents/harness/native-hook-relay-types.ts +++ b/src/agents/harness/native-hook-relay-types.ts @@ -8,6 +8,7 @@ import type { } from "../agent-tools.before-tool-call.js"; import type { CodexMcpServersConfig } from "../codex-mcp-config.types.js"; import type { AgentHarnessHostCapabilities } from "./host-capability-types.js"; +import type { retainBeforeToolCallForNativeHookRelay } from "./host-private-capabilities.js"; type NativeHookRelayApprovalContext = Pick< HookContext, @@ -87,6 +88,8 @@ export type NativeHookRelayRegistration = { signal?: AbortSignal; /** Exact host policy capability for authority-bearing native callbacks. */ runBeforeToolCall?: AgentHarnessHostCapabilities["runBeforeToolCall"]; + /** Foreground-only approval authority supplied by the admitted bundled host. */ + approvalHost?: Pick; /** Revalidates the exact admitted owner after authority-bearing awaits. */ assertActive?: AgentHarnessHostCapabilities["assertActive"]; onPreToolUseFailure?: (failure: { @@ -239,7 +242,16 @@ export type NativeHookRelayPermissionApprovalRequester = ( request: NativeHookRelayPermissionApprovalRequest, ) => Promise; +export type NativeHookRelayPendingPermissionApproval = { + relayId: string; + promise: Promise; + controller: AbortController; + waiters: number; + cancelWhenUnobserved: boolean; +}; + export type NativeHookRelayPreToolUseApproval = { + relayId: string; deferredApproval: DeferredPluginToolApproval; originalParamsFingerprint: string; resolutionPromise?: Promise; @@ -273,8 +285,30 @@ export type NativeHookRelaySharedState = { relayBridges: Map; pendingOperations: Set>; invocations: NativeHookRelayInvocation[]; - pendingPermissionApprovals: Map>; + pendingPermissionApprovals: Map; pendingPreToolUseApprovals: Map; permissionApprovalWindows: Map; permissionAllowAlwaysApprovals: Map; }; + +/** Private bundled-runtime callbacks for retained direct-child hook policy. */ +export type NativeHookRelayRetention = Readonly<{ + readClaim: (rawPayload: unknown) => string | undefined; + shouldRetainAfterForegroundClose: () => boolean; + allowPreToolUse: (claim: string) => boolean; + awaitForegroundAdmission?: ( + claim: string, + signal?: AbortSignal, + ) => Promise<(() => boolean) | undefined>; + onDispose: () => void; +}>; + +export type RelayLifetime = { + foregroundOpen: boolean; + foregroundToken: symbol; + policyReady: Promise; + retained?: ReturnType; + retention?: NativeHookRelayRetention; + removeAbortListener?: () => void; + expiryTimer?: ReturnType; +}; diff --git a/src/agents/harness/native-hook-relay.approval-wait.test.ts b/src/agents/harness/native-hook-relay.approval-wait.test.ts index 9db14a22a87c..aefb16ee4fbd 100644 --- a/src/agents/harness/native-hook-relay.approval-wait.test.ts +++ b/src/agents/harness/native-hook-relay.approval-wait.test.ts @@ -1,5 +1,8 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { createDeferredCore } from "../../shared/deferred.js"; import { callGatewayTool } from "../tools/gateway.js"; +import { createAdmittedHostCapabilityTestFixture } from "./host-capability.test-support.js"; +import { nativeHookRelayState } from "./native-hook-relay-state.js"; import { invokeNativeHookRelay, registerNativeHookRelay, @@ -30,6 +33,375 @@ afterEach(async () => { }); describe("native hook relay approval wait handling", () => { + it.each([ + { + name: "relay/run tuple", + relayIds: ["a", "a:b"], + runIds: ["b:c", "c"], + callIds: ["call", "call"], + }, + { + name: "relay prefix", + relayIds: ["a", "a:b"], + runIds: ["run", "run"], + callIds: ["one", "two"], + }, + { + name: "explicit call versus fallback", + relayIds: ["a", "a"], + runIds: ["run", "run"], + callIds: [undefined, "keys:none"], + }, + ])("isolates permission ownership across $name keys", async ({ relayIds, runIds, callIds }) => { + const host = await createAdmittedHostCapabilityTestFixture({ runId: "tuple-owner" }); + const held = [ + createDeferredCore<{ id: string; decision: string }>(), + createDeferredCore<{ id: string; decision: string }>(), + ]; + const signals: AbortSignal[] = []; + mockCallGatewayTool.mockImplementation(async (method, _opts, _params, extra) => { + if (method !== "plugin.approval.request" || !extra?.signal) { + throw new Error("unexpected approval request"); + } + signals.push(extra.signal); + return await held[signals.length - 1]!.promise; + }); + const firstRelay = registerOwnedNativeHookRelay({ + provider: "codex", + relayId: relayIds[0], + sessionId: "tuple-session", + runId: runIds[0]!, + approvalHost: host.hostCapabilities, + }); + const secondRelay = + relayIds[0] === relayIds[1] + ? firstRelay + : registerOwnedNativeHookRelay({ + provider: "codex", + relayId: relayIds[1], + sessionId: "tuple-session", + runId: runIds[1]!, + approvalHost: host.hostCapabilities, + }); + await Promise.all([firstRelay.ready, secondRelay.ready]); + const invoke = ( + relay: typeof firstRelay, + toolUseId: string | undefined, + signal?: AbortSignal, + ) => + invokeNativeHookRelay( + { + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "permission_request", + rawPayload: { + tool_name: "call", + ...(toolUseId ? { tool_use_id: toolUseId } : {}), + tool_input: {}, + }, + }, + signal, + ); + const controller = new AbortController(); + const first = invoke(firstRelay, callIds[0], controller.signal); + void first.catch(() => {}); + let second: ReturnType | undefined; + let duplicate: ReturnType | undefined; + try { + await vi.waitFor(() => expect(signals).toHaveLength(1)); + second = invoke(secondRelay, callIds[1]); + void second.catch(() => {}); + await vi.waitFor(() => expect(signals).toHaveLength(2)); + expect(nativeHookRelayState.pendingPermissionApprovals.size).toBe(2); + duplicate = invoke(secondRelay, callIds[1]); + await vi.waitFor(() => + expect( + [...nativeHookRelayState.pendingPermissionApprovals.values()] + .map((entry) => entry.waiters) + .toSorted((a, b) => a - b), + ).toEqual([1, 2]), + ); + if (firstRelay === secondRelay) { + controller.abort(); + } else { + firstRelay.unregister(); + } + await expect(first).rejects.toThrow(firstRelay === secondRelay ? /abort/i : /inactive/); + await vi.waitFor(() => expect(signals[0]?.aborted).toBe(true)); + expect(signals[1]?.aborted).toBe(false); + host.hostCapabilities.assertActive(); + held[1]!.resolve({ id: "second", decision: "allow-always" }); + for (const response of await Promise.all([second, duplicate])) { + expect(JSON.parse(response.stdout).hookSpecificOutput.decision.behavior).toBe("allow"); + } + held[0]!.resolve({ id: "first", decision: "allow-always" }); + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect([...nativeHookRelayState.permissionAllowAlwaysApprovals.values()]).toEqual([ + expect.objectContaining({ relayId: secondRelay.relayId }), + ]); + const cached = await invoke(secondRelay, "next-call"); + expect(JSON.parse(cached.stdout).hookSpecificOutput.decision.behavior).toBe("allow"); + expect(signals).toHaveLength(2); + } finally { + controller.abort(); + firstRelay.unregister(); + secondRelay.unregister(); + for (const decision of held) { + decision.resolve({ id: "cleanup", decision: "deny" }); + } + await Promise.allSettled([first, second, ...(duplicate ? [duplicate] : [])]); + await Promise.all([firstRelay.drain(), secondRelay.drain()]); + host.closeHost(); + host.closeAdmission(); + } + }); + + it.each(["request", "waitDecision"])( + "rejoins a public approval after every %s waiter disconnects", + async (phase) => { + const held = createDeferredCore<{ id: string; decision: string }>(); + let requestSignal: AbortSignal | undefined; + mockCallGatewayTool.mockImplementation(async (method, _opts, _params, extra) => { + if (method === `plugin.approval.${phase}`) { + requestSignal = extra?.signal; + return await held.promise; + } + if (method === "plugin.approval.request") { + return { id: "public-approval", status: "accepted" }; + } + throw new Error(`unexpected gateway method: ${method}`); + }); + const relay = registerNativeHookRelay({ + provider: "codex", + sessionId: "public-approval", + runId: "public-approval", + }); + const invoke = (signal?: AbortSignal) => + invokeNativeHookRelay( + { + provider: "codex", + relayId: relay.relayId, + event: "permission_request", + rawPayload: { tool_name: "fixture", tool_use_id: "public-call", tool_input: {} }, + }, + signal, + ); + const controller = new AbortController(); + const first = invoke(controller.signal); + void first.catch(() => {}); + try { + await vi.waitFor(() => expect(requestSignal).toBeDefined()); + controller.abort(); + await expect(first).rejects.toThrow(/abort/i); + await vi.waitFor(() => + expect([...nativeHookRelayState.pendingPermissionApprovals.values()][0]?.waiters).toBe(0), + ); + expect(requestSignal?.aborted).toBe(false); + const retry = invoke(); + await vi.waitFor(() => + expect([...nativeHookRelayState.pendingPermissionApprovals.values()][0]?.waiters).toBe(1), + ); + expect( + mockCallGatewayTool.mock.calls.filter( + ([method]) => method === `plugin.approval.${phase}`, + ), + ).toHaveLength(1); + held.resolve({ id: "public-approval", decision: "allow-once" }); + expect(JSON.parse((await retry).stdout).hookSpecificOutput.decision.behavior).toBe("allow"); + expect(nativeHookRelayState.pendingPermissionApprovals.size).toBe(0); + } finally { + held.resolve({ id: "public-approval", decision: "deny" }); + controller.abort(); + relay.unregister(); + await Promise.allSettled([first]); + } + }, + ); + + it.each([false, true])( + "fences foreground permission results without retiring a retained child (retained: %s)", + async (retainChild) => { + const host = await createAdmittedHostCapabilityTestFixture({ runId: "permission-owner" }); + const entered = createDeferredCore(); + const decision = createDeferredCore<{ id: string; decision: string }>(); + mockCallGatewayTool.mockImplementation(async (method, _opts, _params, extra) => { + if (method === "plugin.approval.request") { + return { id: "approval-1", status: "accepted" }; + } + if (method !== "plugin.approval.waitDecision" || !extra?.signal) { + throw new Error("fixture wait missing"); + } + entered.resolve(extra.signal); + return await decision.promise; + }); + let retained = retainChild; + const relay = registerOwnedNativeHookRelay({ + provider: "codex", + sessionId: "permission-owner", + runId: "permission-owner", + runBeforeToolCall: host.hostCapabilities.runBeforeToolCall, + approvalHost: host.hostCapabilities, + assertActive: host.hostCapabilities.assertActive, + retention: { + readClaim: () => "child", + shouldRetainAfterForegroundClose: () => retained, + allowPreToolUse: () => true, + onDispose: () => {}, + }, + }); + await relay.ready; + const pending = invokeNativeHookRelay({ + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "permission_request", + rawPayload: { agent_id: "child", tool_name: "fixture", tool_input: {} }, + }); + void pending.catch(() => undefined); + try { + const signal = await entered.promise; + relay.unregister(); + expect(signal.aborted).toBe(true); + if (retainChild) { + decision.resolve({ id: "approval-1", decision: "allow-once" }); + await expect(pending).rejects.toThrow("foreground invocation not allowed"); + await expect( + invokeNativeHookRelay({ + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "pre_tool_use", + rawPayload: { agent_id: "child", tool_name: "fixture", tool_input: {} }, + }), + ).resolves.toMatchObject({ exitCode: 0 }); + } else { + await expect(pending).rejects.toThrow(/inactive/); + } + } finally { + retained = false; + relay.unregister(); + decision.resolve({ id: "approval-1", decision: "deny" }); + await Promise.allSettled([pending]); + await relay.drain(); + host.closeHost(); + host.closeAdmission(); + } + }, + ); + + it.each([ + { phase: "request", cancelAll: false }, + { phase: "waitDecision", cancelAll: false }, + { phase: "request", cancelAll: true }, + { phase: "waitDecision", cancelAll: true }, + ])( + "owns shared approval $phase work independently of duplicate callers (all cancelled: $cancelAll)", + async ({ phase, cancelAll }) => { + const host = await createAdmittedHostCapabilityTestFixture({ runId: "approval-cancel" }); + const held = [ + createDeferredCore<{ id: string; decision: string }>(), + createDeferredCore<{ id: string; decision: string }>(), + ]; + const signals: AbortSignal[] = []; + mockCallGatewayTool.mockImplementation(async (method, _opts, _params, extra) => { + if (method === `plugin.approval.${phase}`) { + if (!extra?.signal) { + throw new Error("fixture shared approval signal missing"); + } + signals.push(extra.signal); + // Ignore cancellation deliberately: late transport completion must be harmless. + return await held[signals.length - 1]!.promise; + } + if (method === "plugin.approval.request") { + return { id: "approval-1", status: "accepted" }; + } + throw new Error(`unexpected gateway method: ${method}`); + }); + const relay = registerOwnedNativeHookRelay({ + provider: "codex", + sessionId: "approval-cancel", + runId: "approval-cancel", + approvalHost: host.hostCapabilities, + }); + await relay.ready; + const invoke = (signal?: AbortSignal) => + invokeNativeHookRelay( + { + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "permission_request", + rawPayload: { tool_name: "fixture", tool_use_id: "duplicate-call", tool_input: {} }, + }, + signal, + ); + const firstAbort = new AbortController(); + const secondAbort = new AbortController(); + const first = invoke(firstAbort.signal); + const second = invoke(secondAbort.signal); + void first.catch(() => undefined); + void second.catch(() => undefined); + let successor: ReturnType | undefined; + try { + await vi.waitFor(() => + expect([...nativeHookRelayState.pendingPermissionApprovals.values()][0]?.waiters).toBe(2), + ); + expect(signals).toHaveLength(1); + expect(signals[0]).not.toBe(firstAbort.signal); + firstAbort.abort(); + await expect(first).rejects.toThrow(/abort/i); + expect(signals[0]?.aborted).toBe(false); + expect( + mockCallGatewayTool.mock.calls.filter( + ([method]) => method === `plugin.approval.${phase}`, + ), + ).toHaveLength(1); + if (cancelAll) { + secondAbort.abort(); + await expect(second).rejects.toThrow(/abort/i); + await vi.waitFor(() => expect(signals[0]?.aborted).toBe(true)); + expect(nativeHookRelayState.pendingPermissionApprovals.size).toBe(0); + successor = invoke(); + await vi.waitFor(() => expect(signals).toHaveLength(2)); + const successorEntry = [...nativeHookRelayState.pendingPermissionApprovals.values()][0]; + held[0]!.resolve({ id: "approval-1", decision: "allow-always" }); + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect([...nativeHookRelayState.pendingPermissionApprovals.values()][0]).toBe( + successorEntry, + ); + expect(nativeHookRelayState.permissionAllowAlwaysApprovals.size).toBe(0); + expect(signals[1]?.aborted).toBe(false); + held[1]!.resolve({ id: "approval-1", decision: "deny" }); + expect(JSON.parse((await successor).stdout).hookSpecificOutput.decision.behavior).toBe( + "deny", + ); + } else { + held[0]!.resolve({ id: "approval-1", decision: "allow-once" }); + expect(JSON.parse((await second).stdout).hookSpecificOutput.decision.behavior).toBe( + "allow", + ); + } + expect(nativeHookRelayState.pendingPermissionApprovals.size).toBe(0); + } finally { + for (const pending of held) { + pending.resolve({ id: "approval-1", decision: "deny" }); + } + firstAbort.abort(); + secondAbort.abort(); + relay.unregister(); + await Promise.allSettled([first, second, successor]); + await relay.drain(); + host.closeHost(); + host.closeAdmission(); + } + }, + ); + it("defers all native MCP names to Codex when the exact agent has a prepared durable grant", async () => { const grant = { server: "raw-server_", tool: "_raw.tool", source: "allow-always", addedAt: 1 }; approvalMocks.loadExecApprovalsReadOnly.mockReturnValue({ diff --git a/src/agents/harness/native-hook-relay.lifecycle.test.ts b/src/agents/harness/native-hook-relay.lifecycle.test.ts index c298ca931311..e24bd315a2ba 100644 --- a/src/agents/harness/native-hook-relay.lifecycle.test.ts +++ b/src/agents/harness/native-hook-relay.lifecycle.test.ts @@ -1,5 +1,6 @@ -import { Agent, Server } from "node:http"; +import { Agent, Server, request } from "node:http"; import { afterEach, expect, it, vi } from "vitest"; +import * as mutableFileBinding from "../../infra/system-run-approval-binding.js"; import { initializeGlobalHookRunner, resetGlobalHookRunner, @@ -11,11 +12,14 @@ import { createAdmittedHostCapabilityTestFixture } from "./host-capability.test- import * as relayBridge from "./native-hook-relay-bridge.js"; import * as clientStore from "./native-hook-relay-client-store.js"; import { invokeNativeHookRelayBridge } from "./native-hook-relay-client.js"; +import { setNativeHookRelayPreToolUseApproval } from "./native-hook-relay-permissions.js"; +import { nativeHookRelayState } from "./native-hook-relay-state.js"; import * as store from "./native-hook-relay-store.js"; import { invokeNativeHookRelay, registerNativeHookRelay, registerOwnedNativeHookRelay, + resolveNativeHookRelayDeferredToolApproval, testing, } from "./native-hook-relay.js"; @@ -25,6 +29,334 @@ afterEach(async () => { vi.restoreAllMocks(); }); +it.each(["deferred outcome", "rejection"] as const)( + "observes policy %s after synchronous cancellation", + async (outcome) => { + const controller = new AbortController(); + const reason = new Error("policy cancelled its run"); + const onResolution = vi.fn(); + const policy = vi.fn(async () => { + controller.abort(reason); + if (outcome === "rejection") { + throw new Error("policy failed after cancellation"); + } + return { + blocked: false as const, + params: {}, + deferredApproval: { + approval: { title: "fixture", description: "fixture", onResolution }, + toolName: "fixture", + baseParams: {}, + }, + }; + }); + const relay = registerNativeHookRelay({ + provider: "codex", + sessionId: "synchronous-abort", + runId: "synchronous-abort", + signal: controller.signal, + runBeforeToolCall: policy, + }); + await expect( + invokeNativeHookRelay({ + provider: "codex", + relayId: relay.relayId, + event: "pre_tool_use", + rawPayload: { tool_name: "fixture", tool_use_id: "call", tool_input: {} }, + }), + ).rejects.toMatchObject({ name: "AbortError", cause: reason }); + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(policy).toHaveBeenCalledOnce(); + expect(onResolution.mock.calls).toEqual(outcome === "deferred outcome" ? [["cancelled"]] : []); + expect(nativeHookRelayState.pendingPreToolUseApprovals.size).toBe(0); + expect(testing.getNativeHookRelayRegistrationForTests(relay.relayId)).toBeUndefined(); + }, +); + +it.each([ + { name: "tuple collision", toolIds: ["b:c", "c"] }, + { name: "relay prefix", toolIds: ["one", "two"] }, +])("keeps deferred approvals with their exact relay across $name", async ({ toolIds }) => { + const callbacks = [vi.fn(), vi.fn()]; + const relays = ["a", "a:b"].map((relayId, index) => + registerNativeHookRelay({ + provider: "codex", + relayId, + sessionId: "tuple-session", + runId: "tuple-run", + runBeforeToolCall: async () => ({ + blocked: false, + params: {}, + deferredApproval: { + approval: { title: "fixture", description: "fixture", onResolution: callbacks[index] }, + toolName: "fixture", + baseParams: {}, + }, + }), + }), + ); + for (const [index, relay] of relays.entries()) { + await invokeNativeHookRelay({ + provider: "codex", + relayId: relay.relayId, + event: "pre_tool_use", + rawPayload: { tool_name: "fixture", tool_use_id: toolIds[index], tool_input: {} }, + }); + } + expect(nativeHookRelayState.pendingPreToolUseApprovals.size).toBe(2); + expect(callbacks[0]).not.toHaveBeenCalled(); + expect(callbacks[1]).not.toHaveBeenCalled(); + relays[0]!.unregister(); + expect(callbacks[0]).toHaveBeenCalledExactlyOnceWith("cancelled"); + expect(callbacks[1]).not.toHaveBeenCalled(); + testing.setNativeHookRelayDeferredToolApprovalRequesterForTests(async () => ({ + blocked: false, + params: {}, + approvalResolution: "allow-once", + })); + await expect( + resolveNativeHookRelayDeferredToolApproval({ + relayId: relays[1]!.relayId, + toolUseId: toolIds[1], + }), + ).resolves.toEqual({ handled: true, outcome: "approved-once" }); + expect(nativeHookRelayState.pendingPreToolUseApprovals.size).toBe(0); +}); + +it("detaches both approval maps before a cancellation callback installs a successor", async () => { + const relay = registerNativeHookRelay({ + provider: "codex", + relayId: "reentrant", + sessionId: "old", + runId: "old", + }); + const key = JSON.stringify([relay.relayId, "call"]); + const held = createDeferredCore<{ + blocked: false; + params: unknown; + approvalResolution: "allow-once"; + }>(); + testing.setNativeHookRelayDeferredToolApprovalRequesterForTests(() => held.promise); + const successorCancelled = vi.fn(); + const successorController = new AbortController(); + const oldController = new AbortController(); + const oldPermission = { + relayId: relay.relayId, + controller: oldController, + waiters: 1, + cancelWhenUnobserved: true, + promise: Promise.resolve("deny" as const), + }; + const permissionKey = "fixture-permission-entry"; + nativeHookRelayState.pendingPermissionApprovals.set(permissionKey, oldPermission); + let successor: ReturnType | undefined; + const onResolution = vi.fn(() => { + successor = registerNativeHookRelay({ + provider: "codex", + relayId: relay.relayId, + sessionId: "new", + runId: "new", + }); + setNativeHookRelayPreToolUseApproval({ + relayId: relay.relayId, + toolUseId: "call", + originalParamsFingerprint: "fixture", + deferredApproval: { + approval: { title: "new", description: "new", onResolution: successorCancelled }, + toolName: "fixture", + baseParams: {}, + }, + }); + nativeHookRelayState.pendingPermissionApprovals.set(permissionKey, { + ...oldPermission, + controller: successorController, + }); + nativeHookRelayState.permissionAllowAlwaysApprovals.set("new-grant", { + relayId: relay.relayId, + }); + nativeHookRelayState.permissionApprovalWindows.set(relay.relayId, [1]); + }); + setNativeHookRelayPreToolUseApproval({ + relayId: relay.relayId, + toolUseId: "call", + originalParamsFingerprint: "fixture", + deferredApproval: { + approval: { title: "old", description: "old", onResolution }, + toolName: "fixture", + baseParams: {}, + }, + }); + const pending = resolveNativeHookRelayDeferredToolApproval({ + relayId: relay.relayId, + toolUseId: "call", + }); + relay.unregister(); + expect(onResolution).toHaveBeenCalledExactlyOnceWith("cancelled"); + expect(oldController.signal.aborted).toBe(true); + expect(successorController.signal.aborted).toBe(false); + expect(successorCancelled).not.toHaveBeenCalled(); + expect(successor).toBeDefined(); + expect(nativeHookRelayState.relays.get(relay.relayId)?.generation).toBe(successor?.generation); + const successorApproval = nativeHookRelayState.pendingPreToolUseApprovals.get(key); + expect(successorApproval?.deferredApproval.approval.title).toBe("new"); + held.resolve({ blocked: false, params: {}, approvalResolution: "allow-once" }); + await pending; + expect(nativeHookRelayState.pendingPreToolUseApprovals.get(key)).toBe(successorApproval); + expect(nativeHookRelayState.pendingPermissionApprovals.get(permissionKey)?.controller).toBe( + successorController, + ); + expect(nativeHookRelayState.permissionAllowAlwaysApprovals.has("new-grant")).toBe(true); + expect(nativeHookRelayState.permissionApprovalWindows.get(relay.relayId)).toEqual([1]); + successor?.unregister(); +}); + +it("cancels disconnected HTTP policy work without retiring the relay or storing a late approval", async () => { + await withOpenClawTestState({ label: "relay-http-disconnect" }, async () => { + const entered = createDeferredCore(); + const release = createDeferredCore(); + const cancelled = createDeferredCore(); + const onResolution = vi.fn(() => cancelled.resolve()); + const relay = registerOwnedNativeHookRelay({ + provider: "codex", + sessionId: "http-disconnect", + runId: "http-disconnect", + runBeforeToolCall: async ({ signal }) => { + entered.resolve(signal); + await release.promise; + return { + blocked: false, + params: {}, + deferredApproval: { + approval: { title: "fixture", description: "fixture", onResolution }, + toolName: "exec", + baseParams: {}, + }, + }; + }, + }); + await relay.ready; + const record = await store.readNativeHookRelayBridgeRecord({ relayId: relay.relayId }); + if (!record) { + throw new Error("fixture bridge missing"); + } + const outgoing = request({ + host: record.hostname, + port: record.port, + method: "POST", + path: "/invoke", + headers: { authorization: `Bearer ${record.token}`, "content-type": "application/json" }, + }); + outgoing.on("error", () => undefined); + outgoing.end( + JSON.stringify({ + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "pre_tool_use", + rawPayload: { tool_name: "Bash", tool_input: {}, tool_use_id: "disconnected-call" }, + }), + ); + try { + const signal = await entered.promise; + if (!signal) { + throw new Error("fixture invocation signal missing"); + } + expect(signal.aborted).toBe(false); + const aborted = new Promise((resolve) => { + signal.addEventListener("abort", () => resolve(), { once: true }); + }); + outgoing.destroy(); + await aborted; + expect(testing.getNativeHookRelayRegistrationForTests(relay.relayId)).toBeDefined(); + release.resolve(); + await cancelled.promise; + expect(onResolution).toHaveBeenCalledExactlyOnceWith("cancelled"); + expect( + nativeHookRelayState.pendingPreToolUseApprovals.has( + JSON.stringify([relay.relayId, "disconnected-call"]), + ), + ).toBe(false); + await expect( + invokeNativeHookRelayBridge({ + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "post_tool_use", + rawPayload: { tool_name: "Bash", tool_response: {} }, + }), + ).resolves.toEqual({ stdout: "", stderr: "", exitCode: 0 }); + } finally { + outgoing.destroy(); + release.resolve(); + relay.unregister(); + await relay.drain(); + } + }); +}); + +it("does not start a permission approval after cancellation during file preparation", async () => { + await withOpenClawTestState({ label: "relay-permission-disconnect" }, async () => { + const entered = createDeferredCore(); + const release = createDeferredCore(); + const completed = createDeferredCore(); + const prepare = mutableFileBinding.prepareSystemRunMutableFileBinding; + vi.spyOn(mutableFileBinding, "prepareSystemRunMutableFileBinding").mockImplementation( + async (...args) => { + entered.resolve(); + await release.promise; + try { + return await prepare(...args); + } finally { + completed.resolve(); + } + }, + ); + const requester = vi.fn(async () => "allow" as const); + testing.setNativeHookRelayPermissionApprovalRequesterForTests(requester); + const relay = registerOwnedNativeHookRelay({ + provider: "codex", + sessionId: "permission-disconnect", + runId: "permission-disconnect", + }); + await relay.ready; + const abort = new AbortController(); + const pending = invokeNativeHookRelay( + { + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "permission_request", + rawPayload: { + tool_name: "Bash", + tool_input: { command: "echo fixture" }, + tool_use_id: "cancelled-permission", + }, + }, + abort.signal, + ); + void pending.catch(() => undefined); + try { + await entered.promise; + abort.abort(); + await expect(pending).rejects.toThrow(/abort/i); + release.resolve(); + await completed.promise; + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(requester).not.toHaveBeenCalled(); + expect(testing.getNativeHookRelayRegistrationForTests(relay.relayId)).toBeDefined(); + } finally { + release.resolve(); + await Promise.allSettled([pending, completed.promise]); + relay.unregister(); + await relay.drain(); + } + }); +}); + it("keeps command preparation synchronous while readiness waits for locator publication", async () => { await withOpenClawTestState({ label: "relay-ready-publication" }, async () => { const entered = createDeferredCore(); diff --git a/src/agents/harness/native-hook-relay.test.ts b/src/agents/harness/native-hook-relay.test.ts index f682f9ac9ea7..cd7670b8f161 100644 --- a/src/agents/harness/native-hook-relay.test.ts +++ b/src/agents/harness/native-hook-relay.test.ts @@ -4371,6 +4371,7 @@ describe("native hook relay registry", () => { expect(getNativeHookRelaySharedStateForTests().pendingPermissionApprovals.size).toBe(1); firstRelay.unregister(); + await expect(firstApproval).rejects.toThrow("registration is inactive"); registerNativeHookRelay({ provider: "codex", relayId, @@ -4387,7 +4388,9 @@ describe("native hook relay registry", () => { expect(getNativeHookRelaySharedStateForTests().pendingPermissionApprovals.size).toBe(1); resolvers[0]?.("allow"); - await expect(firstApproval).rejects.toThrow("registration is inactive"); + await new Promise((resolve) => { + setImmediate(resolve); + }); expect(getNativeHookRelaySharedStateForTests().pendingPermissionApprovals.size).toBe(1); const duplicateSecondApproval = invokeNativeHookRelay({ diff --git a/src/agents/harness/native-hook-relay.ts b/src/agents/harness/native-hook-relay.ts index 4e4d9b9773e3..4e52b50731a3 100644 --- a/src/agents/harness/native-hook-relay.ts +++ b/src/agents/harness/native-hook-relay.ts @@ -4,9 +4,11 @@ import { MAX_TIMER_TIMEOUT_MS, resolveExpiresAtMsFromDurationMs, } from "@openclaw/normalization-core/number-coercion"; +import { racePromiseWithAbortSignal } from "../../infra/abort-signal.js"; import { createSubsystemLogger } from "../../logging/subsystem.js"; import { resolveOpenClawStateSqlitePath } from "../../state/openclaw-state-db.paths.js"; import { retainBeforeToolCallForNativeHookRelay } from "./host-private-capabilities.js"; +import { formatPermissionApprovalDescription as formatPermissionApprovalDescriptionForTestsImpl } from "./native-hook-relay-approval-presentation.js"; import { clearNativeHookRelayBridgesForTests, NATIVE_HOOK_BRIDGE_REPLACEMENT_RECORD_GRACE_MS, @@ -27,12 +29,11 @@ import { import { processNativeHookRelayInvocation } from "./native-hook-relay-events.js"; import { clearNativeHookRelayPermissionsForTests, - formatPermissionApprovalDescriptionForTests as formatPermissionApprovalDescriptionForTestsImpl, permissionRequestContentFingerprintForTests as permissionRequestContentFingerprintForTestsImpl, permissionRequestToolInputKeyFingerprintForTests as permissionRequestToolInputKeyFingerprintForTestsImpl, pruneNativeHookRelayPermissionAllowAlways, removeNativeHookRelayPermissionState, - removeNativeHookRelayPreToolUseApprovals, + detachNativeHookRelayApprovalState, setNativeHookRelayDeferredToolApprovalRequesterForTests as setNativeHookRelayDeferredToolApprovalRequesterForTestsImpl, setNativeHookRelayPermissionApprovalRequesterForTests as setNativeHookRelayPermissionApprovalRequesterForTestsImpl, } from "./native-hook-relay-permissions.js"; @@ -48,12 +49,13 @@ import type { InvokeNativeHookRelayParams, NativeHookRelayEvent, NativeHookRelayInvocation, - NativeHookRelayPermissionApprovalRequest, NativeHookRelayPermissionApprovalRequester, NativeHookRelayProcessResponse, NativeHookRelayRegistration, + NativeHookRelayRetention, OwnedNativeHookRelayRegistrationHandle, RegisterNativeHookRelayParams, + RelayLifetime, } from "./native-hook-relay-types.js"; import { NATIVE_HOOK_RELAY_EVENTS } from "./native-hook-relay-types.js"; import { @@ -82,38 +84,22 @@ const DEFAULT_RELAY_TTL_MS = 30 * 60 * 1000; const log = createSubsystemLogger("agents/harness/native-hook-relay"); const { relays, relayBridges, invocations } = nativeHookRelayState; -type RelayLifetime = { - foregroundOpen: boolean; - foregroundToken: symbol; - policyReady: Promise; - retained?: ReturnType; - retention?: NativeHookRelayRetention; - removeAbortListener?: () => void; - expiryTimer?: ReturnType; -}; - const RELAY_LIFETIME = "__openclawNativeHookRelayLifetimeV1"; -/** Private bundled-runtime callbacks for retained direct-child hook policy. */ -export type NativeHookRelayRetention = Readonly<{ - readClaim: (rawPayload: unknown) => string | undefined; - shouldRetainAfterForegroundClose: () => boolean; - allowPreToolUse: (claim: string) => boolean; - awaitForegroundAdmission?: (claim: string) => Promise<(() => boolean) | undefined>; - onDispose: () => void; -}>; - type OwnedNativeHookRelayParams = RegisterNativeHookRelayParams & { retention?: NativeHookRelayRetention; + approvalHost?: NativeHookRelayRegistration["approvalHost"]; +}; + +type RelayLifetimeRegistration = ActiveNativeHookRelayRegistration & { + [RELAY_LIFETIME]?: RelayLifetime; }; function readRelayLifetime( registration: ActiveNativeHookRelayRegistration, ): RelayLifetime | undefined { - // SAFETY: this private symbol-keyed expando is installed only by setRelayLifetime below. - return (registration as ActiveNativeHookRelayRegistration & { [RELAY_LIFETIME]?: RelayLifetime })[ - RELAY_LIFETIME - ]; + // SAFETY: this private expando is installed only by setRelayLifetime below. + return (registration as RelayLifetimeRegistration)[RELAY_LIFETIME]; } function setRelayLifetime( @@ -166,13 +152,14 @@ export function registerNativeHookRelay( export function registerOwnedNativeHookRelay( params: OwnedNativeHookRelayParams, ): OwnedNativeHookRelayRegistrationHandle { - const { retention, ...registrationParams } = params; - return registerNativeHookRelayInternal(registrationParams, retention); + const { retention, approvalHost, ...registrationParams } = params; + return registerNativeHookRelayInternal(registrationParams, retention, approvalHost); } function registerNativeHookRelayInternal( params: RegisterNativeHookRelayParams, retention: NativeHookRelayRetention | undefined, + approvalHost?: NativeHookRelayRegistration["approvalHost"], ): OwnedNativeHookRelayRegistrationHandle { pruneExpiredNativeHookRelays(); pruneNativeHookRelayPermissionAllowAlways(); @@ -226,6 +213,7 @@ function registerNativeHookRelayInternal( preToolUseFailureProjections: new Map(), ...(params.signal ? { signal: params.signal } : {}), ...(params.runBeforeToolCall ? { runBeforeToolCall: params.runBeforeToolCall } : {}), + ...(approvalHost ? { approvalHost } : {}), ...(params.assertActive ? { assertActive: params.assertActive } : {}), ...(params.onPreToolUseFailure ? { onPreToolUseFailure: params.onPreToolUseFailure } : {}), // SAFETY: the literal supplies the complete mutable internal registration contract. @@ -363,16 +351,14 @@ function unregisterNativeHookRelay( lifetime?.removeAbortListener?.(); lifetime?.retained?.release(); // SAFETY: this deletes the same private expando installed by setRelayLifetime. - delete (registration as ActiveNativeHookRelayRegistration & { [RELAY_LIFETIME]?: RelayLifetime })[ - RELAY_LIFETIME - ]; + delete (registration as RelayLifetimeRegistration)[RELAY_LIFETIME]; void unregisterNativeHookRelayBridge(relayId, { ...options, ...(bridge ? { expectedBridge: bridge } : {}), }); removeNativeHookRelayInvocations(relayId); - removeNativeHookRelayPreToolUseApprovals(relayId); - removeNativeHookRelayPermissionState(relayId); + const cancelApprovals = detachNativeHookRelayApprovalState(relayId); + cancelApprovals(); const deliverOnUnregister = () => { try { lifetime?.retention?.onDispose(); @@ -417,6 +403,8 @@ function deactivateNativeHookRelayForeground( } } if (shouldRetain) { + // Retention covers child PreToolUse only; foreground approval authority ends now. + removeNativeHookRelayPermissionState(relayId); return; } unregisterNativeHookRelay(relayId, registration); @@ -426,13 +414,15 @@ async function resolveNativeHookRelayInvocationBinding( registration: ActiveNativeHookRelayRegistration, event: NativeHookRelayEvent, rawPayload: unknown, + signal?: AbortSignal, ): Promise { const lifetime = readRelayLifetime(registration); if (!lifetime) { throw new Error("native hook relay registration is inactive"); } // Gateway fallback shares policy readiness without depending on HTTP locator publication. - await lifetime.policyReady; + await racePromiseWithAbortSignal(lifetime.policyReady, signal); + signal?.throwIfAborted(); if (relays.get(registration.relayId) !== registration || Date.now() > registration.expiresAtMs) { throw new Error("native hook relay registration is inactive"); } @@ -442,6 +432,7 @@ async function resolveNativeHookRelayInvocationBinding( const retention = lifetime.retention; let assertAdmission: (() => boolean) | undefined; const assertRetainedAuthority = () => { + signal?.throwIfAborted(); if ( relays.get(registration.relayId) !== registration || Date.now() > registration.expiresAtMs @@ -458,7 +449,10 @@ async function resolveNativeHookRelayInvocationBinding( } }; if (lifetime.foregroundOpen && retention.awaitForegroundAdmission) { - assertAdmission = await retention.awaitForegroundAdmission(claim); + assertAdmission = await racePromiseWithAbortSignal( + retention.awaitForegroundAdmission(claim, signal), + signal, + ); if (!assertAdmission) { throw new Error("native hook relay retained invocation not allowed"); } @@ -470,15 +464,18 @@ async function resolveNativeHookRelayInvocationBinding( ...registration, assertActive: assertRetainedAuthority, runBeforeToolCall: retained.runBeforeToolCall, + signal, }; } if (!lifetime.foregroundOpen) { throw new Error("native hook relay foreground invocation not allowed"); } const foregroundToken = lifetime.foregroundToken; - const assertActive = () => + const assertActive = () => { + signal?.throwIfAborted(); assertNativeHookRelayForegroundCurrent(registration, lifetime, foregroundToken); - return { ...registration, assertActive }; + }; + return { ...registration, assertActive, signal }; } function normalizeRelayKey( @@ -497,6 +494,7 @@ function normalizeRelayKey( export async function invokeNativeHookRelay( params: InvokeNativeHookRelayParams, + invocationSignal?: AbortSignal, ): Promise { const provider = readNativeHookRelayProvider(params.provider); const relayId = readNonEmptyString(params.relayId, "relayId"); @@ -506,6 +504,11 @@ export async function invokeNativeHookRelay( pruneExpiredNativeHookRelays(); throw new Error("native hook relay not found"); } + const signal = + invocationSignal && registration.signal + ? AbortSignal.any([invocationSignal, registration.signal]) + : (invocationSignal ?? registration.signal); + signal?.throwIfAborted(); if (Date.now() > registration.expiresAtMs) { unregisterNativeHookRelay(relayId, registration); throw new Error("native hook relay expired"); @@ -542,17 +545,21 @@ export async function invokeNativeHookRelay( registration, event, params.rawPayload, + signal, ); if (event === "pre_tool_use" || event === "permission_request") { effectiveRegistration.assertActive?.(); } recordNativeHookRelayInvocation(normalized); const startedAt = Date.now(); - const response = await processNativeHookRelayInvocation({ - registration: effectiveRegistration, - invocation: normalized, - adapter: getNativeHookRelayProviderAdapter(provider), - }); + const response = await racePromiseWithAbortSignal( + processNativeHookRelayInvocation({ + registration: effectiveRegistration, + invocation: normalized, + adapter: getNativeHookRelayProviderAdapter(provider), + }), + signal, + ); // Policy and approval callbacks may yield while their admitted run closes. // Never let a late allow cross back into the native runtime. if (event === "pre_tool_use" || event === "permission_request") { @@ -715,16 +722,8 @@ export const testing = { isNativeHookRelayBridgeLookupRetryableForTests(error: unknown, elapsedMs = 0): boolean { return isRetryableNativeHookRelayBridgeLookupError({ error, elapsedMs }); }, - formatPermissionApprovalDescriptionForTests( - request: NativeHookRelayPermissionApprovalRequest, - ): string { - return formatPermissionApprovalDescriptionForTestsImpl(request); - }, - permissionRequestContentFingerprintForTests( - request: NativeHookRelayPermissionApprovalRequest, - ): string { - return permissionRequestContentFingerprintForTestsImpl(request); - }, + formatPermissionApprovalDescriptionForTests: formatPermissionApprovalDescriptionForTestsImpl, + permissionRequestContentFingerprintForTests: permissionRequestContentFingerprintForTestsImpl, permissionRequestToolInputKeyFingerprintForTests: permissionRequestToolInputKeyFingerprintForTestsImpl, setNativeHookRelayPermissionApprovalRequesterForTests( diff --git a/src/gateway/mcp-http.ts b/src/gateway/mcp-http.ts index b5f864b6935c..516a06e661c4 100644 --- a/src/gateway/mcp-http.ts +++ b/src/gateway/mcp-http.ts @@ -1,11 +1,7 @@ // MCP loopback HTTP server. // Exposes Gateway-scoped tools to local MCP clients over bearer-auth loopback. import crypto from "node:crypto"; -import { - createServer as createHttpServer, - type IncomingMessage, - type ServerResponse, -} from "node:http"; +import { createServer as createHttpServer, type ServerResponse } from "node:http"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { withAgentQuestionAnswerAuthority } from "../agents/harness/host-private-capabilities.js"; import { acknowledgeInternalToolResult } from "../agents/runtime/internal-hooks.js"; @@ -20,7 +16,10 @@ import { resolveSessionEntryAccessTarget } from "../config/sessions/session-acce import { isTruthyEnvValue } from "../infra/env.js"; import { formatErrorMessage } from "../infra/errors.js"; import { isRequestBodyLimitError, readRequestBodyWithLimit } from "../infra/http-body.js"; -import { sendHttpRequestRejection } from "../infra/http-request-lifecycle.js"; +import { + createHttpRequestAbortSignal, + sendHttpRequestRejection, +} from "../infra/http-request-lifecycle.js"; import { logDebug, logWarn } from "../logger.js"; import { AGENT_HARNESS_SESSION_KEY_RESERVED_MESSAGE, @@ -130,39 +129,6 @@ function logMcpLoopbackTraffic(step: string, details: Record): console.error(`[mcp-loopback] ${step} ${JSON.stringify(details)}`); } -// Abort tool calls when the request disconnects before completion, but keep -// completed responses alive through normal response close notifications. -function createRequestAbortSignal(req: IncomingMessage, res: ServerResponse) { - const controller = new AbortController(); - const abort = () => { - if (!controller.signal.aborted) { - controller.abort(); - } - }; - const abortIfRequestIncomplete = () => { - if (!req.complete) { - abort(); - } - }; - const abortIfResponseStillOpen = () => { - if (!res.writableEnded) { - abort(); - } - }; - req.once("close", abortIfRequestIncomplete); - res.once("close", abortIfResponseStillOpen); - if (req.destroyed && !req.complete) { - abort(); - } - return { - signal: controller.signal, - cleanup: () => { - req.off("close", abortIfRequestIncomplete); - res.off("close", abortIfResponseStillOpen); - }, - }; -} - /** Starts a new MCP loopback HTTP server and registers its bearer tokens. */ async function startMcpLoopbackServer(port = 0): Promise<() => Promise> { const ownerToken = crypto.randomBytes(32).toString("hex"); @@ -209,7 +175,7 @@ async function startMcpLoopbackServer(port = 0): Promise<() => Promise> { // an accepted request is still uploading, and retries must not outrun it. const cliCaptureKey = resolveMcpCliCaptureKey(req, auth); const cliRequestCaptureHandle = markMcpLoopbackRequestStarted(cliCaptureKey); - const requestAbort = createRequestAbortSignal(req, res); + const requestAbort = createHttpRequestAbortSignal(req, res); void (async () => { let parsed: unknown; let cliCaptureHandles: Array> = []; diff --git a/src/gateway/server-aux-handlers.authority.test.ts b/src/gateway/server-aux-handlers.authority.test.ts index 41898dabf47c..5f961cf79d62 100644 --- a/src/gateway/server-aux-handlers.authority.test.ts +++ b/src/gateway/server-aux-handlers.authority.test.ts @@ -152,6 +152,79 @@ describe("gateway auxiliary authority lifecycle", () => { ); }); + it("retires one request approval while its sibling and admitted run remain live", async () => { + const onAgentRunAuthorityClosed = + vi.fn< + ( + authority: ReturnType, + approvalReason?: string, + ) => void + >(); + const gatewayAux = createAuthorityHarness({ + onAgentRunAuthorityClosed, + validateAgentRuntimeDelegatedAuthority: validateAgentRunDelegatedAuthority, + }); + const publishResolved = vi.fn(); + gatewayAux.bindApprovalPublicationContext({ + broadcast: vi.fn(), + broadcastToConnIds: vi.fn(), + approvalEvents: { publishResolved }, + logGateway: { error: vi.fn() }, + } as never); + const authority = claimAgentRunDelegatedAuthority({ + instanceId: "native-permission-owner", + runId: "native-permission-run", + }); + const host = new AbortController(); + const first = new AbortController(); + const second = new AbortController(); + const records = [first, second].map((request, index) => { + const scoped = claimAgentRunApprovalAuthority(authority, [host.signal, request.signal]); + const record = gatewayAux.pluginApprovalManager.create( + { title: `Request ${index}`, description: "Independent native approval" }, + 60_000, + ); + record.agentRuntimeDelegatedAuthority = { ...scoped, kind: "local" }; + const decision = gatewayAux.pluginApprovalManager.register(record, 60_000); + return { record, decision }; + }); + try { + first.abort(); + await expect(records[0]!.decision).resolves.toBeNull(); + expect(getOperatorApprovalDetailed({ id: records[0]!.record.id })).toMatchObject({ + outcome: "found", + record: { status: "cancelled", terminalReason: "run-aborted" }, + }); + await vi.waitFor(() => expect(publishResolved).toHaveBeenCalledTimes(1)); + expect(publishResolved).toHaveBeenCalledWith( + "plugin", + expect.objectContaining({ id: records[0]!.record.id }), + ); + expect(getOperatorApprovalDetailed({ id: records[1]!.record.id })).toMatchObject({ + outcome: "found", + record: { status: "pending" }, + }); + // Scoped approval notifications do not retire whole-run capabilities. + expect( + onAgentRunAuthorityClosed.mock.calls.filter( + ([, approvalReason]) => approvalReason === undefined, + ), + ).toEqual([]); + expect(validateAgentRunDelegatedAuthority(authority)).toBe(true); + expect( + gatewayAux.pluginApprovalManager.resolve( + records[1]!.record.id, + "allow-once", + "fixture reviewer", + ), + ).toBe(true); + await expect(records[1]!.decision).resolves.toBe("allow-once"); + } finally { + host.abort(); + releaseAgentRunDelegatedAuthority(authority); + } + }); + it.each(["release", "replacement", "generation"] as const)( "settles credential questions on authority %s without waiting for a read", async (closure) => { diff --git a/src/gateway/server-methods/native-hook-relay.test.ts b/src/gateway/server-methods/native-hook-relay.test.ts index 1821ebc256f6..166bacd3a289 100644 --- a/src/gateway/server-methods/native-hook-relay.test.ts +++ b/src/gateway/server-methods/native-hook-relay.test.ts @@ -8,6 +8,7 @@ import path from "node:path"; import { setImmediate } from "node:timers/promises"; import { expectDefined } from "@openclaw/normalization-core"; import { afterEach, describe, expect, it, vi } from "vitest"; +import { nativeHookRelayState } from "../../agents/harness/native-hook-relay-state.js"; import { testing, registerNativeHookRelay, @@ -34,6 +35,73 @@ afterEach(async () => { }); describe("native hook relay gateway method", () => { + it("cancels a closed one-shot Gateway connection and discards its late approval", async () => { + await withOpenClawTestState({ label: "relay-gateway-disconnect" }, async () => { + const entered = createDeferredCore(); + const release = createDeferredCore(); + const cancelled = createDeferredCore(); + const onResolution = vi.fn(() => cancelled.resolve()); + const relay = registerOwnedNativeHookRelay({ + provider: "codex", + sessionId: "gateway-disconnect", + runId: "gateway-disconnect", + runBeforeToolCall: async ({ signal }) => { + entered.resolve(signal); + await release.promise; + return { + blocked: false, + params: {}, + deferredApproval: { + approval: { title: "fixture", description: "fixture", onResolution }, + toolName: "exec", + baseParams: {}, + }, + }; + }, + }); + await relay.ready; + const connection = new AbortController(); + const pending = invokeNativeHook( + { + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "pre_tool_use", + rawPayload: { tool_name: "Bash", tool_input: {}, tool_use_id: "disconnected-call" }, + }, + connection.signal, + ); + try { + const signal = await entered.promise; + connection.abort(); + expectInvalidRequest(await pending, "aborted"); + expect(signal?.aborted).toBe(true); + expect(testing.getNativeHookRelayRegistrationForTests(relay.relayId)).toBeDefined(); + release.resolve(); + await cancelled.promise; + expect(onResolution).toHaveBeenCalledExactlyOnceWith("cancelled"); + expect( + nativeHookRelayState.pendingPreToolUseApprovals.has( + JSON.stringify([relay.relayId, "disconnected-call"]), + ), + ).toBe(false); + const next = await invokeNativeHook({ + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "post_tool_use", + rawPayload: POST_TOOL_USE_PAYLOAD, + }); + expect(next).toHaveBeenCalledWith(true, { stdout: "", stderr: "", exitCode: 0 }); + } finally { + release.resolve(); + await pending; + relay.unregister(); + await relay.drain(); + } + }); + }); + it("returns its synchronous handle before reading stored MCP policy", async () => { await withOpenClawTestState({ label: "relay-policy-registration" }, async () => { const read = vi.spyOn(mcpGrants, "loadMcpToolGrants"); @@ -329,7 +397,8 @@ describe("native hook relay gateway method", () => { closure === "owner-close" ? "fixture native owner closed" : "registration is inactive"; await expect(preparation).rejects.toThrow(expectedError); const respond = await invocation; - expectInvalidRequest(respond, expectedError); + // Invocation cancellation stops its wait before preparation rechecks the owner. + expectInvalidRequest(respond, closure === "abort" ? "aborted" : expectedError); expect(requester.mock.calls.length).toBe(0); } finally { paused.resume.resolve(); @@ -449,7 +518,7 @@ describe("native hook relay gateway method", () => { }); }); -async function invokeNativeHook(params: Record) { +async function invokeNativeHook(params: Record, connectionSignal?: AbortSignal) { const respond = viRespond(); await expectDefined( nativeHookRelayHandlers["nativeHook.invoke"], @@ -457,7 +526,16 @@ async function invokeNativeHook(params: Record) { )({ req: { type: "req", id: "1", method: "nativeHook.invoke" }, params, - client: null, + client: connectionSignal + ? { + connectionSignal, + connect: { + minProtocol: 1, + maxProtocol: 1, + client: { id: "cli", version: "test", platform: "test", mode: "cli" }, + }, + } + : null, isWebchatConnect: () => false, respond, context: {} as never, diff --git a/src/gateway/server-methods/native-hook-relay.ts b/src/gateway/server-methods/native-hook-relay.ts index bad8705db7a9..bbf6b5761962 100644 --- a/src/gateway/server-methods/native-hook-relay.ts +++ b/src/gateway/server-methods/native-hook-relay.ts @@ -8,19 +8,22 @@ import type { GatewayRequestHandlers } from "./types.js"; /** Gateway request handlers for invoking registered native hook relays. */ export const nativeHookRelayHandlers: GatewayRequestHandlers = { - "nativeHook.invoke": async ({ params, respond }) => { + "nativeHook.invoke": async ({ params, respond, client }) => { try { // Relay invocations are one-shot bridges into a live native harness. // Require the current generation so stale clients cannot post into a // newly registered relay with the same id. - const result: NativeHookRelayProcessResponse = await invokeNativeHookRelay({ - provider: params.provider, - relayId: params.relayId, - generation: params.generation, - event: params.event, - rawPayload: params.rawPayload, - requireGeneration: true, - }); + const result: NativeHookRelayProcessResponse = await invokeNativeHookRelay( + { + provider: params.provider, + relayId: params.relayId, + generation: params.generation, + event: params.event, + rawPayload: params.rawPayload, + requireGeneration: true, + }, + client?.connectionSignal, + ); respond(true, result); } catch (error) { respond( diff --git a/src/gateway/server.native-hook-request-lifetime.test.ts b/src/gateway/server.native-hook-request-lifetime.test.ts new file mode 100644 index 000000000000..6fc6b534b291 --- /dev/null +++ b/src/gateway/server.native-hook-request-lifetime.test.ts @@ -0,0 +1,340 @@ +// Real Gateway proof: run only with isolated SQLite coordination. +import { once } from "node:events"; +import { rawDataToString } from "@openclaw/gateway-client/websocket-data"; +import { describe, expect, it, vi } from "vitest"; +import type { WebSocket } from "ws"; +import { GATEWAY_CLIENT_CAPS } from "../../packages/gateway-protocol/src/client-info.js"; +import { createAdmittedHostCapabilityTestFixture } from "../agents/harness/host-capability.test-support.js"; +import { nativeHookRelayState } from "../agents/harness/native-hook-relay-state.js"; +import { + invokeNativeHookRelay, + registerOwnedNativeHookRelay, + testing, +} from "../agents/harness/native-hook-relay.js"; +import { racePromiseWithAbortSignal } from "../infra/abort-signal.js"; +import { createDeferredCore } from "../shared/deferred.js"; +import { getOperatorApprovalDetailed } from "./operator-approval-store.js"; +import * as approvalShared from "./server-methods/approval-shared.js"; +import { + connectOk, + createGatewaySuiteHarness, + installGatewayTestHooks, + rpcReq, +} from "./test-helpers.server.js"; + +installGatewayTestHooks({ scope: "suite" }); + +type GatewayHarness = Awaited>; + +describe("native hook relay WebSocket request lifetime", () => { + it.for([false, true])( + "keeps accepted approval ownership across a lost callback (admitted host: %s)", + async (owned, { signal }) => { + const accepted = createDeferredCore(); + const records = new Map>(); + const responses: Array<() => void> = []; + const handlers: Promise[] = []; + const pending: Promise[] = []; + const resolved: string[] = []; + const originalRequest = approvalShared.handlePendingApprovalRequest; + const observation = vi + .spyOn(approvalShared, "handlePendingApprovalRequest") + .mockImplementation((params) => { + const toolCallId = + "toolCallId" in params.record.request ? params.record.request.toolCallId : undefined; + if (params.approvalKind !== "plugin" || typeof toolCallId !== "string") { + return originalRequest(params); + } + const calls = records.get(toolCallId) ?? []; + calls.push({ + id: params.record.id, + claimId: params.record.agentRuntimeDelegatedAuthority?.claimId, + }); + records.set(toolCallId, calls); + const holdAck = toolCallId === "call-a" && calls.length === 1; + const run = originalRequest({ + ...params, + respond: (...args) => { + const payload = args[1]; + if ( + holdAck && + payload && + typeof payload === "object" && + "status" in payload && + payload.status === "accepted" + ) { + // Delay only delivery of the real accepted acknowledgement, after + // the owner has persisted and published its approval request. + responses.push(() => params.respond(...args)); + accepted.resolve(params.record.id); + } else { + params.respond(...args); + } + }, + }); + handlers.push(run.catch(() => {})); + return run; + }); + let gateway: GatewayHarness | undefined; + let host: Awaited> | undefined; + let relay: ReturnType | undefined; + let reviewer: WebSocket | undefined; + const firstAbort = new AbortController(); + const releaseResponses = () => { + for (const respond of responses.splice(0)) { + respond(); + } + }; + signal.addEventListener("abort", releaseResponses, { once: true }); + try { + gateway = await createGatewaySuiteHarness({ serverOptions: { bind: "loopback" } }); + await gateway.server.startupSettled; + reviewer = await gateway.openWs(); + await connectOk(reviewer, { + scopes: ["operator.admin"], + caps: [GATEWAY_CLIENT_CAPS.APPROVALS], + }); + reviewer.on("message", (data) => { + const message = JSON.parse(rawDataToString(data)); + if (message.type === "event" && message.event === "plugin.approval.resolved") { + resolved.push(message.payload.id); + } + }); + if (owned) { + host = await createAdmittedHostCapabilityTestFixture({ + agentId: "main", + sessionId: "native-approval-owner", + sessionKey: "agent:main:native-approval-owner", + runId: "native-approval-owner", + }); + } + relay = registerOwnedNativeHookRelay({ + provider: "codex", + sessionId: "native-approval-owner", + runId: "native-approval-owner", + ...(host + ? { + approvalHost: host.hostCapabilities, + assertActive: host.hostCapabilities.assertActive, + } + : {}), + }); + await relay.ready; + const invoke = (toolCallId: string, invocationSignal?: AbortSignal) => { + const request = invokeNativeHookRelay( + { + provider: "codex", + relayId: relay!.relayId, + event: "permission_request", + rawPayload: { tool_name: "fixture", tool_use_id: toolCallId, tool_input: {} }, + }, + invocationSignal, + ); + pending.push(request.catch(() => {})); + return request; + }; + const first = invoke("call-a", firstAbort.signal); + const id = await racePromiseWithAbortSignal( + Promise.race([ + accepted.promise, + first.then(() => { + throw new Error("relay completed before the accepted acknowledgement"); + }), + ]), + signal, + ); + expect(getOperatorApprovalDetailed({ id })).toMatchObject({ + outcome: "found", + record: { status: "pending" }, + }); + const other = owned ? invoke("call-b") : undefined; + if (owned) { + await vi.waitFor(() => expect(records.get("call-b")).toHaveLength(1)); + expect(records.get("call-a")?.[0]?.claimId).toBeDefined(); + expect(records.get("call-b")?.[0]?.claimId).toBeDefined(); + expect(records.get("call-b")?.[0]?.claimId).not.toBe(records.get("call-a")?.[0]?.claimId); + } + firstAbort.abort(); + await expect(first).rejects.toThrow(/abort/i); + if (owned) { + await vi.waitFor(() => + expect(getOperatorApprovalDetailed({ id })).toMatchObject({ + outcome: "found", + record: { status: "cancelled", terminalReason: "run-aborted" }, + }), + ); + await vi.waitFor(() => expect(resolved).toContain(id)); + expect(getOperatorApprovalDetailed({ id: records.get("call-b")![0]!.id })).toMatchObject({ + outcome: "found", + record: { status: "pending" }, + }); + expect(() => host!.hostCapabilities.assertActive()).not.toThrow(); + } else { + expect(getOperatorApprovalDetailed({ id })).toMatchObject({ + outcome: "found", + record: { status: "pending" }, + }); + expect(resolved).not.toContain(id); + } + const retry = invoke("call-a"); + await vi.waitFor(() => + expect( + [...nativeHookRelayState.pendingPermissionApprovals.values()].some( + (entry) => entry.waiters === 1, + ), + ).toBe(true), + ); + if (owned) { + await vi.waitFor(() => expect(records.get("call-a")).toHaveLength(2)); + } else { + expect(records.get("call-a")).toHaveLength(1); + } + releaseResponses(); + const retryId = records.get("call-a")!.at(-1)!.id; + expect( + ( + await rpcReq(reviewer, "plugin.approval.resolve", { + id: retryId, + decision: "allow-once", + }) + ).ok, + ).toBe(true); + expect(JSON.parse((await retry).stdout).hookSpecificOutput.decision.behavior).toBe("allow"); + if (other) { + const otherId = records.get("call-b")![0]!.id; + expect( + ( + await rpcReq(reviewer, "plugin.approval.resolve", { + id: otherId, + decision: "allow-once", + }) + ).ok, + ).toBe(true); + expect(JSON.parse((await other).stdout).hookSpecificOutput.decision.behavior).toBe( + "allow", + ); + expect(() => host!.hostCapabilities.assertActive()).not.toThrow(); + } + } finally { + firstAbort.abort(); + releaseResponses(); + relay?.unregister(); + host?.closeHost(); + host?.closeAdmission(); + // Public requests intentionally survive disconnection; explicitly settle + // only this fixture's still-pending approvals before closing the server. + if (reviewer?.readyState === 1) { + for (const values of records.values()) { + for (const { id } of values) { + const stored = getOperatorApprovalDetailed({ id }); + if (stored.outcome === "found" && stored.record.status === "pending") { + await rpcReq(reviewer, "plugin.approval.resolve", { id, decision: "deny" }).catch( + () => {}, + ); + } + } + } + } + await Promise.all(pending); + await relay?.drain(); + reviewer?.terminate(); + await gateway?.server.close(); + await Promise.all(handlers); + observation.mockRestore(); + signal.removeEventListener("abort", releaseResponses); + } + }, + ); + + it("cancels a disconnected callback without revoking the relay or keeping its late approval", async ({ + signal, + }) => { + const entered = createDeferredCore(); + const release = createDeferredCore(); + const cancelled = createDeferredCore(); + const onResolution = vi.fn(() => cancelled.resolve()); + let gateway: GatewayHarness | undefined; + let relay: ReturnType | undefined; + const clients: WebSocket[] = []; + const unblock = () => release.resolve(); + signal.addEventListener("abort", unblock, { once: true }); + try { + gateway = await createGatewaySuiteHarness({ + serverOptions: { bind: "loopback" }, + }); + await gateway.server.startupSettled; + relay = registerOwnedNativeHookRelay({ + provider: "codex", + sessionId: "native-ws-disconnect", + runId: "native-ws-disconnect", + runBeforeToolCall: async ({ signal: invocationSignal }) => { + entered.resolve(invocationSignal); + await release.promise; + return { + blocked: false, + params: {}, + deferredApproval: { + approval: { title: "fixture", description: "fixture", onResolution }, + toolName: "exec", + baseParams: {}, + }, + }; + }, + }); + await relay.ready; + const requester = await gateway.openWs(); + clients.push(requester); + await connectOk(requester, { scopes: ["operator.admin"] }); + const params = { + provider: "codex", + relayId: relay.relayId, + generation: relay.generation, + event: "pre_tool_use", + rawPayload: { tool_name: "Bash", tool_input: {}, tool_use_id: "disconnected-call" }, + }; + requester.send( + JSON.stringify({ + type: "req", + id: "native-callback", + method: "nativeHook.invoke", + params, + }), + ); + const invocationSignal = await racePromiseWithAbortSignal(entered.promise, signal); + expect(invocationSignal).toBeDefined(); + const disconnected = once(requester, "close"); + requester.close(); + await disconnected; + await vi.waitFor(() => expect(invocationSignal?.aborted).toBe(true)); + expect(testing.getNativeHookRelayRegistrationForTests(relay.relayId)).toBeDefined(); + release.resolve(); + await racePromiseWithAbortSignal(cancelled.promise, signal); + expect(onResolution).toHaveBeenCalledExactlyOnceWith("cancelled"); + expect( + nativeHookRelayState.pendingPreToolUseApprovals.has( + JSON.stringify([relay.relayId, "disconnected-call"]), + ), + ).toBe(false); + + const next = await gateway.openWs(); + clients.push(next); + await connectOk(next, { scopes: ["operator.admin"] }); + const result = await rpcReq(next, "nativeHook.invoke", { + ...params, + event: "post_tool_use", + rawPayload: { tool_name: "Bash", tool_input: {}, tool_response: { output: "ok" } }, + }); + expect(result).toMatchObject({ ok: true, payload: { stdout: "", stderr: "", exitCode: 0 } }); + expect((await rpcReq(next, "health", {})).ok).toBe(true); + } finally { + unblock(); + relay?.unregister(); + await relay?.drain(); + for (const ws of clients) { + ws.terminate(); + } + await gateway?.server.close(); + signal.removeEventListener("abort", unblock); + } + }); +}); diff --git a/src/infra/abort-signal.test.ts b/src/infra/abort-signal.test.ts index a0e2fb00d479..cdd6d561db69 100644 --- a/src/infra/abort-signal.test.ts +++ b/src/infra/abort-signal.test.ts @@ -64,6 +64,57 @@ describe("waitForAbortSignal", () => { }); describe("racePromiseWithAbortSignal", () => { + it.each(["rejected", "pending", "fulfilled"] as const)( + "observes a %s source when an existing abort wins", + async (state) => { + const controller = new AbortController(); + const reason = new Error("already stopped"); + controller.abort(reason); + const deferred = createDeferred(); + const sourceError = new Error("source failed"); + const source = + state === "rejected" + ? Promise.reject(sourceError) + : state === "fulfilled" + ? Promise.resolve("done") + : deferred.promise; + + await expect(racePromiseWithAbortSignal(source, controller.signal)).rejects.toMatchObject({ + name: "AbortError", + cause: reason, + }); + if (state === "pending") { + deferred.reject(sourceError); + } + // Let Node report an unobserved rejection before this regression finishes. + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(getEventListeners(controller.signal, "abort")).toHaveLength(0); + }, + ); + + it("observes a later source rejection after an active signal aborts", async () => { + const controller = new AbortController(); + const source = createDeferred(); + const raced = racePromiseWithAbortSignal(source.promise, controller.signal); + controller.abort(); + await expect(raced).rejects.toMatchObject({ name: "AbortError" }); + source.reject(new Error("late source failure")); + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(getEventListeners(controller.signal, "abort")).toHaveLength(0); + }); + + it("returns the source unchanged when no signal is supplied", async () => { + const source = Promise.resolve("done"); + expect(racePromiseWithAbortSignal(source)).toBe(source); + await expect(source).resolves.toBe("done"); + const failure = new Error("source failed"); + await expect(racePromiseWithAbortSignal(Promise.reject(failure))).rejects.toBe(failure); + }); + it("preserves source settlement and removes the listener", async () => { const signal = new AbortController().signal; const sourceError = new Error("source failed"); diff --git a/src/infra/abort-signal.ts b/src/infra/abort-signal.ts index 5851c706eec7..a22e5212eca7 100644 --- a/src/infra/abort-signal.ts +++ b/src/infra/abort-signal.ts @@ -25,7 +25,9 @@ export function racePromiseWithAbortSignal( } const abortError = () => createAbortError("Operation aborted", { cause: signal.reason }); if (signal.aborted) { - return Promise.reject(abortError()); + // The source may already be running. Observe its rejection while preserving + // the existing abort's precedence, even over an already-settled source. + return Promise.race([Promise.reject(abortError()), promise]); } let onAbort!: () => void; const aborted = new Promise((_, reject) => { diff --git a/src/infra/http-request-lifecycle.ts b/src/infra/http-request-lifecycle.ts index 65d44a2cc704..a6d31c164818 100644 --- a/src/infra/http-request-lifecycle.ts +++ b/src/infra/http-request-lifecycle.ts @@ -24,6 +24,33 @@ type HttpConnection = { const connections = new WeakMap(); +/** Abort disconnected work without treating normal request/response completion as cancellation. */ +export function createHttpRequestAbortSignal(req: IncomingMessage, res: ServerResponse) { + const controller = new AbortController(); + const abortIfRequestIncomplete = () => { + if (!req.complete) { + controller.abort(); + } + }; + const abortIfResponseStillOpen = () => { + if (!res.writableEnded) { + controller.abort(); + } + }; + req.once("close", abortIfRequestIncomplete); + res.once("close", abortIfResponseStillOpen); + if ((req.destroyed && !req.complete) || (res.destroyed && !res.writableEnded)) { + controller.abort(); + } + return { + signal: controller.signal, + cleanup: () => { + req.off("close", abortIfRequestIncomplete); + res.off("close", abortIfResponseStillOpen); + }, + }; +} + function connectionFor(socket: Duplex): HttpConnection { const existing = connections.get(socket); if (existing) { diff --git a/src/plugin-sdk/native-hook-relay-runtime.ts b/src/plugin-sdk/native-hook-relay-runtime.ts index 939ae3cef56b..3ec449bcda9f 100644 --- a/src/plugin-sdk/native-hook-relay-runtime.ts +++ b/src/plugin-sdk/native-hook-relay-runtime.ts @@ -1,18 +1,12 @@ -import type { RegisterNativeHookRelayParams } from "../agents/harness/native-hook-relay-types.js"; // Private retained native-hook relay capability for bundled runtime owners. -import { - registerOwnedNativeHookRelay, - type NativeHookRelayRetention, -} from "../agents/harness/native-hook-relay.js"; +import { registerOwnedNativeHookRelay } from "../agents/harness/native-hook-relay.js"; export { buildNativeHookRelayCommandPlan, type NativeHookRelayCommandPlan, } from "../agents/harness/native-hook-relay-plan.js"; -export type OwnedNativeHookRelayParams = RegisterNativeHookRelayParams & { - retention?: NativeHookRelayRetention; -}; +export type OwnedNativeHookRelayParams = Parameters[0]; /** Bundled owners join publication and cleanup while preserving optional direct-child retention. */ export function registerNativeHookRelayForBundledRuntime(params: OwnedNativeHookRelayParams) { diff --git a/test/scripts/tsdown-runtime-config.test.ts b/test/scripts/tsdown-runtime-config.test.ts index 7858cc3ce7ad..7d0420e27704 100644 --- a/test/scripts/tsdown-runtime-config.test.ts +++ b/test/scripts/tsdown-runtime-config.test.ts @@ -21,7 +21,7 @@ type TsdownConfigEntry = { minify?: unknown; dts?: boolean | { emitDtsOnly?: boolean }; define?: Record; - outputOptions?: { codeSplitting?: boolean }; + outputOptions?: { codeSplitting?: boolean; chunkFileNames?: string }; outExtensions?: () => { js: string }; outDir?: string; plugins?: Array<{ name?: string }>; @@ -83,6 +83,14 @@ function requireSqliteReadOnlyChildGraph(): TsdownConfigEntry { return expectDefined(graphs[0], "read-only snapshot child graph"); } +function requireNativeHookRelayGraph(): TsdownConfigEntry { + const graphs = asConfigArray(tsdownConfig).filter((config) => + entryKeys(config).includes("native-hook-relay/entry"), + ); + expect(graphs).toHaveLength(1); + return expectDefined(graphs[0], "native hook relay graph"); +} + function bundledEntry(pluginId: string): string { return `${bundledPluginRoot(pluginId)}/index`; } @@ -203,13 +211,28 @@ describe("tsdown config", () => { expect(handoffGraph?.plugins).toContainEqual( expect.objectContaining({ name: STATE_SCHEMA_INLINE_PLUGIN_NAME }), ); - expect(entrySources(unifiedGraph)["native-hook-relay/entry"]).toBe( - "src/cli/native-hook-relay-entry.ts", + expect(requireNativeHookRelayGraph().plugins).toContainEqual( + expect.objectContaining({ name: STATE_SCHEMA_INLINE_PLUGIN_NAME }), ); expect(requireSqliteReadOnlyChildGraph().plugins).toContainEqual( expect.objectContaining({ name: STATE_SCHEMA_INLINE_PLUGIN_NAME }), ); - expect(inlinePlugins).toHaveLength(4); + expect(inlinePlugins).toHaveLength(5); + }); + + it("isolates relay startup from shared runtime chunks while retaining lazy fallback", () => { + const relay = requireNativeHookRelayGraph(); + expect(entrySources(relay)).toEqual({ + "native-hook-relay/entry": "src/cli/native-hook-relay-entry.ts", + }); + expect(relay).not.toBe(requireUnifiedDistGraph()); + expect(relay.dts).toBe(false); + expect(relay.outputOptions?.codeSplitting).not.toBe(false); + expect(relay.outputOptions?.chunkFileNames).toBe("native-hook-relay/[name]-[hash].mjs"); + // Only the shared graph may publish the global plugin ownership manifest. + expect(relay.plugins).not.toContainEqual( + expect.objectContaining({ name: "openclaw:runtime-dependency-ownership" }), + ); }); it("keeps core, plugin runtime, plugin-sdk, bundled root plugins, and bundled hooks in one dist graph", () => { diff --git a/tsdown.config.ts b/tsdown.config.ts index 342924701a91..8c5a7b536e26 100644 --- a/tsdown.config.ts +++ b/tsdown.config.ts @@ -856,7 +856,6 @@ const configs: UserConfig[] = [ ([name]) => !bundledInventoryEntryNames.has(name), ), ), - "native-hook-relay/entry": "src/cli/native-hook-relay-entry.ts", }, deps: unifiedDeps, // Explicit ESM chunks avoid repeated package-format parsing in Node; @@ -870,6 +869,18 @@ const configs: UserConfig[] = [ }, false, ), + nodeBuildConfig( + { + name: TSDOWN_UNIFIED_CONFIG_GROUP, + // One-shot relays must not load shared Gateway/SDK chunks just to read a locator. + // Keep splitting enabled so the existing Gateway fallback stays lazy. + entry: { "native-hook-relay/entry": "src/cli/native-hook-relay-entry.ts" }, + deps: unifiedDeps, + outputOptions: { chunkFileNames: "native-hook-relay/[name]-[hash].mjs" }, + plugins: [createStateSchemaInlinePlugin()], + }, + false, + ), ...bundledInventoryEntries.map((plugin) => { const entry = listBundledPluginEntrySources([plugin]); const privateChunks = `${bundledPluginRoot(plugin.id)}/.setup/[name]-[hash].mjs`;