From 3f4cefc4d3acfbce7409ba14bd46abb675fa2f37 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 27 Aug 2026 14:29:30 -0700 Subject: [PATCH] fix(code-mode): release canceled timers and parked execution slots (#131186) --- docs/tools/code-mode.md | 6 + src/agents/code-mode-execution.ts | 23 +- .../code-mode-headless.cancellation.test.ts | 81 ++++++ src/agents/code-mode-headless.test.ts | 131 +++------- src/agents/code-mode-headless.ts | 2 + src/agents/code-mode-state.ts | 71 +++++- src/agents/code-mode-swarm.test.ts | 2 + src/agents/code-mode.bridge.lifecycle.test.ts | 232 +++++++++++++++++- src/agents/code-mode.test-support.ts | 27 +- src/agents/code-mode.ts | 5 +- src/agents/code-mode.wait.test.ts | 1 + src/agents/tool-search-catalog.ts | 4 +- src/agents/tool-search-types.ts | 3 +- 13 files changed, 463 insertions(+), 125 deletions(-) create mode 100644 src/agents/code-mode-headless.cancellation.test.ts diff --git a/docs/tools/code-mode.md b/docs/tools/code-mode.md index 6ac76495e87d..c398f7290c95 100644 --- a/docs/tools/code-mode.md +++ b/docs/tools/code-mode.md @@ -562,6 +562,10 @@ QuickJS-WASI snapshot/restore is the resume mechanism: Snapshots are runtime state, not user artifacts: they live only in an in-process map (no database or disk write), are size-limited, expire, and are scoped to the run and session that created them. +Canceling the owning run or tool call, or closing its tool catalog at attempt +teardown, immediately releases parked snapshots and cancels their pending host +work, even if no `wait` call follows. Catalog description refreshes and client +tool additions do not close the owner. `wait` fails (as a `failed` result) when: @@ -590,6 +594,8 @@ declare function yield_control(reason?: string): Promise; ``` Guest timers are bridged through the host, so they survive QuickJS snapshot/resume and remain bounded by the Code Mode execution and snapshot limits. +`clearTimeout` also cancels a timer created before an earlier suspension; this +applies to interactive Code Mode and headless automation scripts. Every effective non-MCP tool is also installed as an async global function. The model-visible `exec` description includes a bounded, deterministic subset diff --git a/src/agents/code-mode-execution.ts b/src/agents/code-mode-execution.ts index 28e9bfc2506f..cec4774a4e07 100644 --- a/src/agents/code-mode-execution.ts +++ b/src/agents/code-mode-execution.ts @@ -35,6 +35,7 @@ import { activeRuns, cancelPendingBridgeStates, cancelPendingBridgeStatesById, + codeModeAbortedResult, codeModeWaitingReason, createCodeModeBridgeDispatchState, createPendingBridgeStates, @@ -295,16 +296,7 @@ async function settleCodeModeResult(params: { // rounds; maxPendingToolCalls stays a per-batch concurrency cap enforced in // the worker. const settleDeadline = () => params.deadlineMs + params.approvalWait.pausedMs; - const abortedResult = () => ({ - status: "failed" as const, - error: "code mode execution aborted", - code: "aborted" as const, - failurePhase: params.bridgeDispatch.started ? ("bridge" as const) : ("host" as const), - bridgeDispatchStarted: params.bridgeDispatch.started, - output: output.slice(deliveredOutputCount), - replaySafe: params.replaySafe, - telemetry: telemetry(params.runtime), - }); + const abortedResult = () => codeModeAbortedResult({ ...params, output, deliveredOutputCount }); // Bridge tool calls (search/describe/call/namespace) run through the same // policy-checked executor whether the model awaits them one at a time or in a // batch, so resolve them inline within the exec deadline and resume the VM @@ -408,6 +400,7 @@ async function settleCodeModeResult(params: { output, deliveredOutputCount, bridgeDispatch: params.bridgeDispatch, + signal: params.signal, }); } // Deliver the settled frontier only. Unresolved sibling promises remain @@ -524,6 +517,7 @@ async function settleCodeModeResult(params: { output, deliveredOutputCount, bridgeDispatch: params.bridgeDispatch, + signal: params.signal, }); } catch (error) { cancelPendingBridgeStates(pending); @@ -612,6 +606,11 @@ export async function runWait(params: { // pending calls, the resume worker, and the inline settle phase. const deadlineMs = Date.now() + state.config.timeoutMs; const approvalWait = observeAgentRunApprovalWait(state.ctx); + // Snapshot closure wakes an observing wait even if its guest budget just expired. + // Transfer releases this signal; worker execution keeps its normal per-call signal. + const parkedSignal = params.signal + ? AbortSignal.any([params.signal, state.ownerSignal]) + : state.ownerSignal; let releaseActiveRunSlot: (() => void) | undefined; try { const ready = await waitForPending( @@ -619,7 +618,7 @@ export async function runWait(params: { state.settlementMode, Math.max(1, deadlineMs - Date.now()), approvalWait, - params.signal, + parkedSignal, ); const resumeBudgetMs = ready ? usableResumeBudgetMs(deadlineMs + approvalWait.pausedMs, state.config) @@ -627,7 +626,7 @@ export async function runWait(params: { if (!ready || resumeBudgetMs === undefined) { // An aborted wait drops the suspended run: nothing will resume it, and // parking it would pin a process-global active-run slot until TTL expiry. - if (params.signal?.aborted) { + if (parkedSignal.aborted) { disposeCodeModeRun(state.runId); return { status: "failed" as const, diff --git a/src/agents/code-mode-headless.cancellation.test.ts b/src/agents/code-mode-headless.cancellation.test.ts new file mode 100644 index 000000000000..85bcbfae3411 --- /dev/null +++ b/src/agents/code-mode-headless.cancellation.test.ts @@ -0,0 +1,81 @@ +import { afterEach, describe, expect, it } from "vitest"; +import { runCodeModeScriptHeadless } from "./code-mode.js"; +import { + createHeadlessCodeModeHarness, + resetCodeModeTestState, + testing, +} from "./code-mode.test-support.js"; + +describe("headless Code Mode cancellation", () => { + afterEach(() => { + try { + expect(testing.activeRuns.size).toBe(0); + } finally { + resetCodeModeTestState(); + } + }); + + it("completes after canceling a guest timer across two resumes", async () => { + const result = await runCodeModeScriptHeadless({ + ctx: createHeadlessCodeModeHarness(), + code: ` + const timer = setTimeout(() => {}, 60_000); + await new Promise((resolve) => setTimeout(resolve, 1)); + clearTimeout(timer); + await new Promise((resolve) => setTimeout(resolve, 1)); + return "done"; + `, + wallClockMs: 5_000, + }); + + expect(result).toEqual({ + status: "completed", + value: "done", + output: [], + toolCallCount: 0, + }); + }); + + it("terminates an in-flight worker leg when aborted", async () => { + const ctx = createHeadlessCodeModeHarness(); + const config = testing.resolveCodeModeHeadlessConfig(ctx); + const controller = new AbortController(); + const resultPromise = testing.runCodeModeWorker( + { + kind: "exec", + source: "while (true) {}", + config, + catalog: [], + apiFiles: [], + namespaces: [], + }, + 5000, + undefined, + controller.signal, + ); + setTimeout(() => controller.abort(), 100); + + await expect(resultPromise).resolves.toMatchObject({ + status: "failed", + code: "aborted", + error: "code mode execution aborted", + }); + }); + + it("classifies caller aborts before the worker leg as aborted", async () => { + const controller = new AbortController(); + controller.abort(); + + const result = await runCodeModeScriptHeadless({ + ctx: createHeadlessCodeModeHarness(), + code: "return true;", + signal: controller.signal, + }); + + expect(result).toMatchObject({ + status: "failed", + code: "aborted", + error: "code mode execution aborted", + }); + }); +}); diff --git a/src/agents/code-mode-headless.test.ts b/src/agents/code-mode-headless.test.ts index 3ae54c13e85c..7b1917865506 100644 --- a/src/agents/code-mode-headless.test.ts +++ b/src/agents/code-mode-headless.test.ts @@ -11,12 +11,7 @@ import { createDeferred } from "../../test/helpers/promise.js"; import type { CodeModeNamespaceDescriptor } from "./code-mode-namespaces.js"; import { prepareSource } from "./code-mode-runtime.js"; import { runCodeModeScriptHeadless, type CodeModeHeadlessResult } from "./code-mode.js"; -import { testing } from "./code-mode.test-support.js"; -import { - createToolSearchCatalogRef, - registerHeadlessToolSearchCatalog, - type ToolSearchToolContext, -} from "./tool-search.js"; +import { createHeadlessCodeModeHarness, testing } from "./code-mode.test-support.js"; import { jsonResult, type AnyAgentTool } from "./tools/common.js"; function fakeTool(name: string, execute: AnyAgentTool["execute"]): AnyAgentTool { @@ -29,26 +24,6 @@ function fakeTool(name: string, execute: AnyAgentTool["execute"]): AnyAgentTool }; } -function createHeadlessHarness( - tools: AnyAgentTool[] = [], - options: { swarmEnabled?: boolean } = {}, -): ToolSearchToolContext { - const config = { - tools: { - codeMode: { enabled: false, timeoutMs: 60_000 }, - ...(options.swarmEnabled ? { swarm: true } : {}), - }, - } as never; - const catalogRef = createToolSearchCatalogRef(); - registerHeadlessToolSearchCatalog({ catalogRef, tools }); - return { - config, - runtimeConfig: config, - agentId: "main", - catalogRef, - }; -} - function expectCompleted(result: CodeModeHeadlessResult) { expect(result.status).toBe("completed"); if (result.status !== "completed") { @@ -87,7 +62,7 @@ describe("headless Code Mode", () => { expect(testing.activeRuns.size).toBe(0); return jsonResult({ input }); }); - const ctx = createHeadlessHarness([first, second]); + const ctx = createHeadlessCodeModeHarness([first, second]); const result = expectCompleted( await runCodeModeScriptHeadless({ @@ -143,7 +118,7 @@ describe("headless Code Mode", () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([first, second, release]), + ctx: createHeadlessCodeModeHarness([first, second, release]), code: `const value = await Promise.race([ headless_first_race({}), headless_second_race({}), @@ -197,7 +172,7 @@ describe("headless Code Mode", () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([never, fast, release]), + ctx: createHeadlessCodeModeHarness([never, fast, release]), code: `const value = await Promise.race([ Promise.all([headless_nested_race_never({})]), headless_nested_race_fast({}), @@ -280,7 +255,7 @@ describe("headless Code Mode", () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([audit, fast, release]), + ctx: createHeadlessCodeModeHarness([audit, fast, release]), code: `${auditCode} const value = await headless_awaited_fast({}); void headless_early_audit_release({}); @@ -341,7 +316,7 @@ describe("headless Code Mode", () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([winner, loser, audit, release]), + ctx: createHeadlessCodeModeHarness([winner, loser, audit, release]), code: `const value = await Promise.race([ headless_race_winner({}), headless_race_loser({}), @@ -375,7 +350,7 @@ describe("headless Code Mode", () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([first, second]), + ctx: createHeadlessCodeModeHarness([first, second]), code: `void headless_detached_first({}); void headless_detached_second({}); return "done";`, @@ -427,7 +402,7 @@ describe("headless Code Mode", () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([fast, slow, release]), + ctx: createHeadlessCodeModeHarness([fast, slow, release]), code: `const value = await Promise.${combinator}([ headless_slow({}), headless_fast({}), @@ -486,7 +461,7 @@ describe("headless Code Mode", () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([failed, slow, release]), + ctx: createHeadlessCodeModeHarness([failed, slow, release]), code: `try { await Promise.all([ headless_failed({}), @@ -520,7 +495,7 @@ describe("headless Code Mode", () => { it("does not expose collector globals without resumable snapshot state", async () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([], { swarmEnabled: true }), + ctx: createHeadlessCodeModeHarness([], { swarmEnabled: true }), code: "return [typeof agents, typeof phase, typeof log];", }), ); @@ -571,14 +546,14 @@ describe("headless Code Mode", () => { "preserves harmless $name in headless source validation", async ({ code, value, realHeadless }) => { if (!realHeadless) { - const ctx = createHeadlessHarness(); + const ctx = createHeadlessCodeModeHarness(); const config = testing.resolveCodeModeHeadlessConfig(ctx); await expect(prepareSource({ code, config })).resolves.toBe(code); return; } const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), code, }), ); @@ -591,7 +566,7 @@ describe("headless Code Mode", () => { it("executes module-shaped regular expressions in a TypeScript headless guest", async () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), language: "typescript", code: 'const value: number = 1; return /import.meta/.test("import.meta");', }), @@ -627,7 +602,7 @@ describe("headless Code Mode", () => { ])("rejects executable module access in a headless guest: %s", async (code) => { const result = expectFailed( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), code, }), ); @@ -642,7 +617,7 @@ describe("headless Code Mode", () => { async (moduleAccess) => { const result = expectFailed( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), language: "typescript", code: `const padding: string = "${"😀".repeat(96)}"; return ${moduleAccess};`, }), @@ -657,7 +632,7 @@ describe("headless Code Mode", () => { it("injects deeply frozen trigger state and emits replacement state through json", async () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), code: ` json({ fire: true, @@ -700,7 +675,7 @@ describe("headless Code Mode", () => { it("keeps an injected namespace while calling a colliding tool by its advertised global", async () => { const tool = fakeTool("trigger", async () => jsonResult({ owner: "tool" })); - const ctx = createHeadlessHarness([tool]); + const ctx = createHeadlessCodeModeHarness([tool]); const extraNamespaces: CodeModeNamespaceDescriptor[] = [ { id: "cron:trigger", @@ -744,7 +719,7 @@ describe("headless Code Mode", () => { it("rejects colliding injected namespace globals", async () => { const result = expectFailed( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), code: "return true;", extraNamespaces: [ { @@ -769,7 +744,7 @@ describe("headless Code Mode", () => { const tool = fakeTool("budgeted", async () => jsonResult({ ok: true })); const result = expectFailed( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([tool]), + ctx: createHeadlessCodeModeHarness([tool]), code: ` await budgeted({}); await budgeted({}); @@ -790,7 +765,7 @@ describe("headless Code Mode", () => { const result = expectFailed( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([nodesTool]), + ctx: createHeadlessCodeModeHarness([nodesTool]), code: ` await nodes.list(); await nodes.list(); @@ -809,7 +784,7 @@ describe("headless Code Mode", () => { it("fails an awaiting promise without bridge work before resuming a worker", async () => { const result = expectFailed( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), code: "await new Promise(() => {}); return true;", wallClockMs: 5_000, }), @@ -825,7 +800,7 @@ describe("headless Code Mode", () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([tool]), + ctx: createHeadlessCodeModeHarness([tool]), code: ` text("x".repeat(700)); await output_boundary({}); @@ -847,7 +822,7 @@ describe("headless Code Mode", () => { const tool = fakeTool("budgeted", async () => jsonResult({ ok: true })); const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([tool]), + ctx: createHeadlessCodeModeHarness([tool]), code: ` const calls = Array.from({ length: 129 }, () => () => budgeted({}), @@ -887,7 +862,7 @@ describe("headless Code Mode", () => { return jsonResult({ ok: true }); }); const resultPromise = runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([slow]), + ctx: createHeadlessCodeModeHarness([slow]), code: ` await slow_leg({}); return true; @@ -915,7 +890,7 @@ describe("headless Code Mode", () => { return jsonResult({ ok: true }); }); const resultPromise = runCodeModeScriptHeadless({ - ctx: createHeadlessHarness([slow]), + ctx: createHeadlessCodeModeHarness([slow]), code: ` await slow_leg({}); return true; @@ -940,7 +915,7 @@ describe("headless Code Mode", () => { it("settles yield_control inline and resumes to completion", async () => { const result = expectCompleted( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), code: ` const yielded = await yield_control("pause"); return { yielded, resumed: true }; @@ -955,56 +930,12 @@ describe("headless Code Mode", () => { expect(result.toolCallCount).toBe(0); }); - it("terminates an in-flight worker leg when aborted", async () => { - const ctx = createHeadlessHarness(); - const config = testing.resolveCodeModeHeadlessConfig(ctx); - const controller = new AbortController(); - const resultPromise = testing.runCodeModeWorker( - { - kind: "exec", - source: "while (true) {}", - config, - catalog: [], - apiFiles: [], - namespaces: [], - }, - 5000, - undefined, - controller.signal, - ); - setTimeout(() => controller.abort(), 100); - - await expect(resultPromise).resolves.toMatchObject({ - status: "failed", - code: "aborted", - error: "code mode execution aborted", - }); - }); - - it("classifies caller aborts before the worker leg as aborted", async () => { - const controller = new AbortController(); - controller.abort(); - - const result = expectFailed( - await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), - code: "return true;", - signal: controller.signal, - }), - ); - - expect(result).toMatchObject({ - code: "aborted", - error: "code mode execution aborted", - }); - }); - it("times out an unfinished headless TypeScript runtime load", async () => { loadCodeModeTypeScriptRuntime.mockReturnValue(new Promise(() => {})); const result = expectFailed( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), language: "typescript", code: "return 42;", wallClockMs: 25, @@ -1023,7 +954,7 @@ describe("headless Code Mode", () => { loadCodeModeTypeScriptRuntime.mockReturnValue(new Promise(() => {})); const controller = new AbortController(); const resultPromise = runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), language: "typescript", code: "return 42;", signal: controller.signal, @@ -1040,7 +971,7 @@ describe("headless Code Mode", () => { }); it("keeps worker-leg wall-clock expiry classified as timeout", async () => { - const ctx = createHeadlessHarness(); + const ctx = createHeadlessCodeModeHarness(); expectCompleted(await runCodeModeScriptHeadless({ ctx, code: "return true;" })); const result = expectFailed( @@ -1058,7 +989,7 @@ describe("headless Code Mode", () => { it("classifies syntax errors", async () => { const result = expectFailed( await runCodeModeScriptHeadless({ - ctx: createHeadlessHarness(), + ctx: createHeadlessCodeModeHarness(), code: "return (;", }), ); @@ -1067,7 +998,7 @@ describe("headless Code Mode", () => { }); it("clamps headless limit overrides to worker-safe bounds", () => { - const config = testing.resolveCodeModeHeadlessConfig(createHeadlessHarness(), { + const config = testing.resolveCodeModeHeadlessConfig(createHeadlessCodeModeHarness(), { timeoutMs: 1, memoryLimitBytes: 1, maxOutputBytes: 1, diff --git a/src/agents/code-mode-headless.ts b/src/agents/code-mode-headless.ts index c391d6575012..fa1b43deab5d 100644 --- a/src/agents/code-mode-headless.ts +++ b/src/agents/code-mode-headless.ts @@ -30,6 +30,7 @@ import { } from "./code-mode-runtime.js"; import { cancelPendingBridgeStates, + cancelPendingBridgeStatesById, createCodeModeBridgeDispatchState, createPendingBridgeStates, pendingBridgeStatesForSettlement, @@ -283,6 +284,7 @@ export async function runCodeModeScriptHeadless(params: { } enforceSnapshotPayloadLimits({ snapshotBytes: result.snapshotBytes, config }); + cancelPendingBridgeStatesById(pending, result.canceledRequestIds); const pendingIds = new Set(pending.map((entry) => entry.id)); const newRequests = result.pendingRequests.filter((request) => !pendingIds.has(request.id)); // Node discovery invokes the generic nodes tool for live status too; diff --git a/src/agents/code-mode-state.ts b/src/agents/code-mode-state.ts index 464afe79f7bf..4ed851a4af54 100644 --- a/src/agents/code-mode-state.ts +++ b/src/agents/code-mode-state.ts @@ -52,6 +52,8 @@ type CodeModeRunState = { catalogProjection: CodeModeCatalogProjection; namespaceRuntime: CodeModeNamespaceRuntime; bridgeDispatch: CodeModeBridgeDispatchState; + ownerSignal: AbortSignal; + releaseOwner: () => void; }; const MAX_ACTIVE_CODE_MODE_RUNS = 64; @@ -120,15 +122,19 @@ export function removeExpiredRuns(now = Date.now()): void { export function disposeCodeModeRun(runId: string): void { const state = activeRuns.get(runId); - cancelPendingBridgeStates(state?.pending ?? []); activeRuns.delete(runId); + state?.releaseOwner(); + cancelPendingBridgeStates(state?.pending ?? []); resumingRunIds.delete(runId); scheduleActiveRunExpiry(); } /** Cancel suspended bridge work before its Gateway-owned runtimes disappear. */ export function disposeAllCodeModeRuns(): void { - activeRuns.forEach((state) => cancelPendingBridgeStates(state.pending)); + activeRuns.forEach((state) => { + state.releaseOwner(); + cancelPendingBridgeStates(state.pending); + }); activeRuns.clear(); resumingRunIds.clear(); scheduleActiveRunExpiry(); @@ -221,8 +227,13 @@ function enforceActiveRunLimit(): void { export function reserveActiveRunSlot(ownedRunId?: string): () => void { if (ownedRunId === undefined) { enforceActiveRunLimit(); - } else if (!activeRuns.delete(ownedRunId)) { - throw new ToolInputError("code mode run is unavailable or expired."); + } else { + const state = activeRuns.get(ownedRunId); + if (!state) { + throw new ToolInputError("code mode run is unavailable or expired."); + } + activeRuns.delete(ownedRunId); + state.releaseOwner(); } // Resume transfers an existing slot without exposing a free capacity window // to concurrent exec calls or rejecting its own run at the global limit. @@ -443,7 +454,19 @@ export function storeSnapshotState(params: { output: unknown[]; deliveredOutputCount?: number; bridgeDispatch: CodeModeBridgeDispatchState; + signal?: AbortSignal; }) { + const catalogRef = params.ctx.catalogRef; + const closed = new AbortController(); + const ownerSignal = AbortSignal.any( + [params.signal, params.ctx.abortSignal, closed.signal].filter( + (signal): signal is AbortSignal => signal !== undefined, + ), + ); + if (!catalogRef?.current || ownerSignal.aborted) { + cancelPendingBridgeStates(params.pending); + return codeModeAbortedResult(params); + } const now = Date.now(); const expiresAt = resolveCodeModeSnapshotExpiresAt(now, params.config.snapshotTtlSeconds); if (expiresAt === undefined) { @@ -458,7 +481,15 @@ export function storeSnapshotState(params: { params.config.snapshotTtlSeconds * MAX_AGENT_WAIT_SNAPSHOT_TTL_WINDOWS, ) : undefined; - activeRuns.set(params.runId, { + const disposers = (catalogRef.onDispose ??= new Set()); + const onClose = () => { + // A transferred snapshot may reuse this cell id; stale observers own only + // the exact parked state they subscribed for, never its replacement. + if (activeRuns.get(params.runId) === state) { + disposeCodeModeRun(params.runId); + } + }; + const state: CodeModeRunState = { runId: params.runId, replayId: params.replayId, parentToolCallId: params.parentToolCallId, @@ -476,7 +507,16 @@ export function storeSnapshotState(params: { catalogProjection: params.catalogProjection, namespaceRuntime: params.namespaceRuntime, bridgeDispatch: params.bridgeDispatch, - }); + ownerSignal, + releaseOwner: () => { + disposers.delete(onClose); + ownerSignal.removeEventListener("abort", onClose); + closed.abort(); + }, + }; + activeRuns.set(params.runId, state); + disposers.add(onClose); + ownerSignal.addEventListener("abort", onClose, { once: true }); scheduleActiveRunExpiry(); return { status: "waiting" as const, @@ -489,6 +529,25 @@ export function storeSnapshotState(params: { }; } +export function codeModeAbortedResult(params: { + bridgeDispatch: CodeModeBridgeDispatchState; + output: unknown[]; + deliveredOutputCount?: number; + replaySafe: boolean; + runtime: ToolSearchRuntime; +}) { + return { + status: "failed" as const, + error: "code mode execution aborted", + code: "aborted" as const, + failurePhase: params.bridgeDispatch.started ? ("bridge" as const) : ("host" as const), + bridgeDispatchStarted: params.bridgeDispatch.started, + output: params.output.slice(params.deliveredOutputCount ?? 0), + replaySafe: params.replaySafe, + telemetry: telemetry(params.runtime), + }; +} + export function codeModeWaitingReason( pending: readonly PendingBridgeState[], ): "pending_tools" | "yield" { diff --git a/src/agents/code-mode-swarm.test.ts b/src/agents/code-mode-swarm.test.ts index a308ddf266cc..441de1cd960d 100644 --- a/src/agents/code-mode-swarm.test.ts +++ b/src/agents/code-mode-swarm.test.ts @@ -494,6 +494,7 @@ describe("Code Mode swarm host bridge", () => { it("renews expired snapshots while agentWait remains pending", () => { const now = 10_000; testing.activeRuns.set("cm-pending-agent", { + releaseOwner: () => undefined, config: { ...config, snapshotTtlSeconds: 60 }, expiresAt: now - 1, agentWaitRetainUntil: now + 120_000, @@ -516,6 +517,7 @@ describe("Code Mode swarm host bridge", () => { const now = 10_000; const cancel = vi.fn(); testing.activeRuns.set("cm-expired-agent", { + releaseOwner: () => undefined, config: { ...config, snapshotTtlSeconds: 60 }, expiresAt: now - 1, agentWaitRetainUntil: now - 1, diff --git a/src/agents/code-mode.bridge.lifecycle.test.ts b/src/agents/code-mode.bridge.lifecycle.test.ts index c2dc07102e46..6f4ab760884f 100644 --- a/src/agents/code-mode.bridge.lifecycle.test.ts +++ b/src/agents/code-mode.bridge.lifecycle.test.ts @@ -1,12 +1,18 @@ /** Subscribed embedded tool lifecycles, including real QuickJS bridge coverage. */ +import { getEventListeners } from "node:events"; import { expectDefined } from "@openclaw/normalization-core"; import { afterEach, describe, expect, it, vi } from "vitest"; import { createDeferred } from "../../test/helpers/promise.js"; import { createDiagnosticEmbeddedRunOwner } from "../logging/diagnostic-run-activity.js"; import { buildExecApprovalPendingToolResult } from "./bash-tools.exec-host-shared.js"; import { disposeAllCodeModeRuns } from "./code-mode-state.js"; -import { applyCodeModeCatalog, createCodeModeTools } from "./code-mode.js"; import { + addClientToolsToCodeModeCatalog, + applyCodeModeCatalog, + createCodeModeTools, +} from "./code-mode.js"; +import { + fakeTool, pluginToolWithExecute, resetCodeModeTestState, resultDetails, @@ -21,7 +27,7 @@ import { emitAssistantTextDeltaAndEnd, } from "./embedded-agent-subscribe.e2e-harness.js"; import { countActiveToolExecutions } from "./embedded-agent-subscribe.handlers.tools.js"; -import { createToolSearchCatalogRef } from "./tool-search.js"; +import { clearToolSearchCatalog, createToolSearchCatalogRef } from "./tool-search.js"; import { jsonResult } from "./tools/common.js"; function createSubscribedCodeModeHarness(params: { @@ -333,6 +339,228 @@ describe("Code Mode subscribed bridge lifecycle", () => { } }); + it.each(["context", "tool", "catalog"] as const)( + "releases only the parked owner after %s abort without wait", + async (signalSource) => { + const owner = createSubscribedCodeModeHarness({ name: `parked-${signalSource}` }); + const survivor = createSubscribedCodeModeHarness({ name: `survivor-${signalSource}` }); + const toolAbortController = new AbortController(); + const controller = + signalSource === "context" ? owner.runAbortController : toolAbortController; + const code = 'setTimeout(() => {}, 60_000); await yield_control("pause"); return "done";'; + applyCodeModeCatalog(owner); + applyCodeModeCatalog(survivor); + + try { + const parked = resultDetails( + await expectDefined(owner.tools[0], "owner exec").execute( + "code-call-parked", + { code }, + signalSource === "tool" ? toolAbortController.signal : undefined, + ), + ); + const other = resultDetails( + await expectDefined(survivor.tools[0], "survivor exec").execute("code-call-survivor", { + code, + }), + ); + for (const result of [parked, other]) { + expect(result).toMatchObject({ status: "waiting", runId: expect.any(String) }); + } + const ownerId = parked.runId as string; + const survivorId = other.runId as string; + const ownerState = expectDefined(testing.activeRuns.get(ownerId), "parked owner snapshot"); + const survivorState = expectDefined( + testing.activeRuns.get(survivorId), + "survivor snapshot", + ); + const pending = expectDefined( + ownerState.pending.find((entry) => entry.method === "sleep"), + "owner timer", + ); + const otherPending = expectDefined( + survivorState.pending.find((entry) => entry.method === "sleep"), + "survivor timer", + ); + expect(pending.settled).toBeUndefined(); + expect(otherPending.settled).toBeUndefined(); + expect(ownerState.snapshotBytes.byteLength).toBeGreaterThan(0); + expect(testing.resumingRunIds.size).toBe(0); + + // Both exec calls have returned; no wait is in flight to perform owner cleanup. + if (signalSource === "catalog") { + clearToolSearchCatalog(owner); + } else { + controller.abort(new Error("parked owner closed")); + } + expect([...testing.activeRuns.keys()]).toEqual([survivorId]); + await expect(pending.promise).resolves.toMatchObject({ id: pending.id, ok: false }); + expect(testing.activeRuns.get(survivorId)).toBe(survivorState); + expect(otherPending.settled).toBeUndefined(); + } finally { + owner.dispose(); + survivor.dispose(); + } + }, + ); + + it.each(["complete", "context", "tool", "catalog"] as const)( + "transfers parked ownership across refresh and repeated resumes before %s", + async (close) => { + const owner = createSubscribedCodeModeHarness({ name: `transfer-${close}` }); + applyCodeModeCatalog(owner); + const exec = expectDefined(owner.tools[0], "owner exec"); + const wait = expectDefined(owner.tools[1], "owner wait"); + let controller = new AbortController(); + try { + let result = resultDetails( + await exec.execute( + "transfer-exec", + { + code: `const timer = setTimeout(() => {}, 60_000); + await yield_control("first"); + await yield_control("second"); + await yield_control("third"); + clearTimeout(timer); + return "done";`, + }, + controller.signal, + ), + ); + expect(result.status).toBe("waiting"); + const runId = result.runId as string; + const initial = expectDefined(testing.activeRuns.get(runId), "initial snapshot"); + expect(applyCodeModeCatalog(owner).catalogReused).toBe(true); + addClientToolsToCodeModeCatalog({ + ...owner, + tools: [fakeTool("client_probe", "Client probe")], + }); + expect(testing.activeRuns.get(runId)).toBe(initial); + expect(exec.description).toContain("client_probe"); + + for (let index = 0; index < 2; index += 1) { + const previous = expectDefined(testing.activeRuns.get(runId), "previous snapshot"); + const staleClose = expectDefined( + owner.catalogRef.onDispose?.values().next().value, + "parked owner subscription", + ); + controller = new AbortController(); + result = resultDetails( + await wait.execute(`transfer-wait-${index}`, { runId }, controller.signal), + ); + expect(result).toMatchObject({ status: "waiting", runId }); + const replacement = expectDefined(testing.activeRuns.get(runId), "replacement snapshot"); + expect(replacement).not.toBe(previous); + staleClose(); + expect(testing.activeRuns.get(runId)).toBe(replacement); + expect(getEventListeners(previous.ownerSignal, "abort")).toHaveLength(0); + expect(previous.ownerSignal.aborted).toBe(true); + expect(getEventListeners(replacement.ownerSignal, "abort")).toHaveLength(1); + expect(owner.catalogRef.onDispose?.size).toBe(1); + } + + const finalState = expectDefined(testing.activeRuns.get(runId), "final snapshot"); + const pending = finalState.pending; + if (close === "complete") { + expect(resultDetails(await wait.execute("transfer-complete", { runId }))).toMatchObject({ + status: "completed", + value: "done", + }); + } else if (close === "catalog") { + clearToolSearchCatalog(owner); + } else { + (close === "context" ? owner.runAbortController : controller).abort(); + } + expect(testing.activeRuns.size).toBe(0); + expect(testing.resumingRunIds.size).toBe(0); + await Promise.all(pending.map((entry) => entry.promise)); + expect(getEventListeners(finalState.ownerSignal, "abort")).toHaveLength(0); + expect(finalState.ownerSignal.aborted).toBe(true); + expect(owner.catalogRef.onDispose?.size ?? 0).toBe(0); + } finally { + clearToolSearchCatalog(owner); + owner.dispose(); + } + }, + ); + + it.each(["exec", "wait"] as const)( + "does not publish a snapshot after its catalog closes during %s", + async (phase) => { + const owner = createSubscribedCodeModeHarness({ name: `close-during-${phase}` }); + const closeOwner = pluginToolWithExecute("close_owner", "Close the run catalog", async () => { + clearToolSearchCatalog(owner); + return jsonResult({ closed: true }); + }); + applyCodeModeCatalog({ ...owner, tools: [...owner.tools, closeOwner] }); + try { + const execute = () => + expectDefined(owner.tools[0], "owner exec").execute("close-during-exec", { + code: `${phase === "wait" ? 'await yield_control("initial");' : ""} + await close_owner({}); + await yield_control("closed"); + return "unreachable";`, + }); + let completion; + if (phase === "wait") { + const parked = resultDetails(await execute()); + expect(parked.status).toBe("waiting"); + completion = expectDefined(owner.tools[1], "owner wait").execute("close-during-wait", { + runId: parked.runId, + }); + } else { + completion = execute(); + } + await expect(completion).rejects.toThrow( + "Tool Search catalog is unavailable for this run.", + ); + expect(closeOwner.execute).toHaveBeenCalledOnce(); + expect(testing.activeRuns.size).toBe(0); + expect(testing.resumingRunIds.size).toBe(0); + expect(countActiveToolExecutions(owner.runId)).toBe(0); + } finally { + clearToolSearchCatalog(owner); + owner.dispose(); + } + }, + ); + + it("does not return a closed snapshot when owner abort races the wait deadline", async () => { + vi.useFakeTimers({ toFake: ["Date", "setTimeout", "clearTimeout"] }); + const owner = createSubscribedCodeModeHarness({ + name: "abort-wait-deadline", + timeoutMs: 1_500, + }); + const started = createDeferred(); + const stalled = pluginToolWithExecute("stalled", "Await cancellation", async () => { + started.resolve(); + return await new Promise(() => {}); + }); + applyCodeModeCatalog({ ...owner, tools: [...owner.tools, stalled] }); + try { + const execution = expectDefined(owner.tools[0], "owner exec").execute("deadline-exec", { + code: "return await stalled({});", + }); + await started.promise; + await vi.advanceTimersByTimeAsync(1_500); + const parked = resultDetails(await execution); + expect(parked.status).toBe("waiting"); + const waiting = expectDefined(owner.tools[1], "owner wait").execute("deadline-wait", { + runId: parked.runId, + }); + vi.advanceTimersByTime(1_499); + owner.runAbortController.abort(); + expect(resultDetails(await waiting)).toMatchObject({ status: "failed", code: "aborted" }); + expect(testing.activeRuns.size).toBe(0); + expect(testing.resumingRunIds.size).toBe(0); + expect(countActiveToolExecutions(owner.runId)).toBe(0); + } finally { + clearToolSearchCatalog(owner); + owner.dispose(); + vi.useRealTimers(); + } + }); + it.each([ { kind: "explicit cancellation", close: "cancel" }, { kind: "run-owner loss", close: "abort" }, diff --git a/src/agents/code-mode.test-support.ts b/src/agents/code-mode.test-support.ts index e15d9f199b79..0bb2b4442586 100644 --- a/src/agents/code-mode.test-support.ts +++ b/src/agents/code-mode.test-support.ts @@ -11,7 +11,12 @@ import { } from "./code-mode-state.js"; import { normalizeCodeModeWorkerResult, runCodeModeWorker } from "./code-mode-worker.js"; import { createCodeModeTools } from "./code-mode.js"; -import { createToolSearchCatalogRef, type ToolSearchCatalogRef } from "./tool-search.js"; +import { + createToolSearchCatalogRef, + registerHeadlessToolSearchCatalog, + type ToolSearchCatalogRef, + type ToolSearchToolContext, +} from "./tool-search.js"; import { jsonResult, type AnyAgentTool } from "./tools/common.js"; export const testing = { @@ -115,6 +120,26 @@ export function resultDetails(result: { details?: unknown }): Record; } +export function createHeadlessCodeModeHarness( + tools: AnyAgentTool[] = [], + options: { swarmEnabled?: boolean } = {}, +): ToolSearchToolContext { + const config = { + tools: { + codeMode: { enabled: false, timeoutMs: 60_000 }, + ...(options.swarmEnabled ? { swarm: true } : {}), + }, + } as never; + const catalogRef = createToolSearchCatalogRef(); + registerHeadlessToolSearchCatalog({ catalogRef, tools }); + return { + config, + runtimeConfig: config, + agentId: "main", + catalogRef, + }; +} + export function createCodeModeHarness( params: { agentId?: string; diff --git a/src/agents/code-mode.ts b/src/agents/code-mode.ts index c34c40a1dd13..72e337043128 100644 --- a/src/agents/code-mode.ts +++ b/src/agents/code-mode.ts @@ -324,9 +324,10 @@ export function applyCodeModeCatalog(params: { const catalogRef = params.catalogRef; const execTool = compacted.tools.find((tool) => tool.name === CODE_MODE_EXEC_TOOL_NAME); if (catalogRef?.current && execTool) { - catalogRef.onDispose?.(); + // Refreshing descriptions replaces their observer, not the catalog's parked consumers. + catalogRef.disposeObserver?.(); const descriptionUpdater = createCodeModeExecDescriptionUpdater(execTool); - catalogRef.onDispose = descriptionUpdater.dispose; + catalogRef.disposeObserver = descriptionUpdater.dispose; catalogRef.onChange = () => { descriptionUpdater.update( createCodeModeExecDescription( diff --git a/src/agents/code-mode.wait.test.ts b/src/agents/code-mode.wait.test.ts index 7a5834e15ae8..a6e71eb11372 100644 --- a/src/agents/code-mode.wait.test.ts +++ b/src/agents/code-mode.wait.test.ts @@ -538,6 +538,7 @@ describe("Code Mode wait, scope, and suspended runs", () => { const { tools: codeModeTools } = createCodeModeHarness(); testing.activeRuns.set("invalid-expiry-run", { expiresAt: 8_640_000_000_000_001, + releaseOwner: () => undefined, } as never); await expect( diff --git a/src/agents/tool-search-catalog.ts b/src/agents/tool-search-catalog.ts index 289f40040952..db621ec42c56 100644 --- a/src/agents/tool-search-catalog.ts +++ b/src/agents/tool-search-catalog.ts @@ -334,9 +334,11 @@ export function clearToolSearchCatalog(params: { catalogRef?: ToolSearchCatalogRef; }): void { if (params.catalogRef) { - params.catalogRef.onDispose?.(); params.catalogRef.current = undefined; + params.catalogRef.disposeObserver?.(); + params.catalogRef.onDispose?.forEach((dispose) => dispose()); delete params.catalogRef.onChange; + delete params.catalogRef.disposeObserver; delete params.catalogRef.onDispose; } if (!params.runId?.trim()) { diff --git a/src/agents/tool-search-types.ts b/src/agents/tool-search-types.ts index 47dcbed8fa85..04ed89dc6c89 100644 --- a/src/agents/tool-search-types.ts +++ b/src/agents/tool-search-types.ts @@ -132,7 +132,8 @@ export type ToolSearchCatalogSession = { export type ToolSearchCatalogRef = { current?: ToolSearchCatalogSession; onChange?: () => void; - onDispose?: () => void; + disposeObserver?: () => void; + onDispose?: Set<() => void>; }; export type CodeModeBridgeMethod = "search" | "describe" | "call";