fix(code-mode): release canceled timers and parked execution slots (#131186)

This commit is contained in:
Peter Steinberger
2026-08-27 14:29:30 -07:00
committed by GitHub
parent 0e071abb1e
commit 3f4cefc4d3
13 changed files with 463 additions and 125 deletions
+6
View File
@@ -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<void>;
```
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
+11 -12
View File
@@ -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,
@@ -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",
});
});
});
+31 -100
View File
@@ -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,
+2
View File
@@ -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;
+65 -6
View File
@@ -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" {
+2
View File
@@ -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,
+230 -2
View File
@@ -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<never>(() => {});
});
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" },
+26 -1
View File
@@ -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<string, unk
return result.details as Record<string, unknown>;
}
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;
+3 -2
View File
@@ -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(
+1
View File
@@ -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(
+3 -1
View File
@@ -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()) {
+2 -1
View File
@@ -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";