fix(coding-agent/eval): made cells own every bridge call they start
- Split the Python bridge signal: the raw abort reaches tools (so subagents die with the turn) while a shielded signal governs how long the host waits, so a cancel can no longer settle a cell on top of a still-running isolation merge. - Mirrored the contract in the JS runtime: abort in-flight tool calls immediately, then drain any deferExternalAbort phase before killing the worker, and refuse new bridge calls once cancelled. - Held a finished worker result until the run's tool calls drain. A floated agent() previously settled the cell at once, dropping the run's abort listener and leaving the subagent running with nothing able to cancel it. - Added regression coverage for all three, each verified to fail without its fix.
This commit is contained in:
@@ -437,18 +437,21 @@ export async function executeWithKernelBase<
|
||||
const emitStatus: (event: JsStatusEvent) => void =
|
||||
options?.emitStatus ?? (event => collectDisplay({ type: "status", event }));
|
||||
const runId = `${runIdPrefix}-${crypto.randomUUID()}`;
|
||||
// The shield is a *kernel* protection: it keeps the runtime alive across a
|
||||
// critical bridge phase (isolation worktree setup, merge/cherry-pick) so the
|
||||
// cell can't be torn down mid-git-operation. Host-side work reached through
|
||||
// the bridge — above all the subagents `agent()` spawns — must stay directly
|
||||
// cancellable, so it gets the caller's real signal. Shielding it here made
|
||||
// Python/Ruby/Julia `agent()` fan-outs survive a turn cancel indefinitely,
|
||||
// while JS cells (which never route through this shield) cancelled fine.
|
||||
// Two aborts cross the bridge, and conflating them is what let a cancelled
|
||||
// turn keep working. Delegated work (above all the subagents `agent()`
|
||||
// spawns) gets the caller's real signal so it dies with the turn — shielding
|
||||
// it here made Python/Ruby/Julia fan-outs outlive a cancel indefinitely,
|
||||
// while JS cells, which never route through this shield, stopped fine. The
|
||||
// shielded signal only governs how long the host waits on a call, holding
|
||||
// the cell open across a critical phase (isolation worktree setup,
|
||||
// merge/cherry-pick) so a cancel can't settle it on top of a half-applied
|
||||
// git operation.
|
||||
const unregisterBridge =
|
||||
options?.toolSession && options?.bridgeSessionId
|
||||
? registerPyToolBridge(options.bridgeSessionId, runId, {
|
||||
toolSession: options.toolSession,
|
||||
signal: options.signal,
|
||||
shieldedSignal: abortShield.signal,
|
||||
emitStatus,
|
||||
abortRequested: () => {
|
||||
return abortShield.abortRequested;
|
||||
|
||||
@@ -8,6 +8,7 @@ import {
|
||||
import type { ToolSession } from "../../tools";
|
||||
import { ToolAbortError, ToolError } from "../../tools/tool-errors";
|
||||
import { safeSend as safeSendIpc } from "../../utils/ipc";
|
||||
import { EVAL_TIMEOUT_PAUSE_OP, EVAL_TIMEOUT_RESUME_OP } from "../bridge-timeout";
|
||||
import { attachSessionOwner, resolveOwnerScopedSessionKey, type SessionOwners } from "../executor-base";
|
||||
import { shouldDetachKernel } from "../py/spawn-options";
|
||||
import { callSessionTool, type JsStatusEvent } from "./tool-bridge";
|
||||
@@ -48,6 +49,26 @@ interface PendingRun {
|
||||
resolve(value: { value: unknown }): void;
|
||||
reject(error: Error): void;
|
||||
toolCalls: Map<string, AbortController>;
|
||||
/**
|
||||
* Host calls currently inside a `deferExternalAbort` phase — `agent()`
|
||||
* isolation worktree setup and merge/cherry-pick, which ignore their abort
|
||||
* once started. Settling the run while one is live would return the cell on
|
||||
* top of a git operation still rewriting the repo, so the abort path drains
|
||||
* this first. Mirrors the Python bridge's shielded-signal contract.
|
||||
*/
|
||||
deferDepth: number;
|
||||
/** Resolves once {@link deferDepth} falls back to zero. */
|
||||
deferDrained?: PromiseWithResolvers<void>;
|
||||
/** Set once the turn was cancelled; blocks new bridge calls during the drain. */
|
||||
aborted: boolean;
|
||||
/**
|
||||
* A worker `result` withheld because the cell still has bridge calls in
|
||||
* flight. `#runOne` reports a finished run without awaiting its pending
|
||||
* tools, so a floated or caught `agent()` would otherwise settle the run —
|
||||
* tearing down the abort listener — while the subagent kept going with
|
||||
* nothing left able to cancel it. Delivered once the last call drains.
|
||||
*/
|
||||
heldResult?: Extract<WorkerOutbound, { type: "result" }>;
|
||||
settled: boolean;
|
||||
}
|
||||
|
||||
@@ -267,6 +288,8 @@ async function runOnce(
|
||||
resolve,
|
||||
reject,
|
||||
toolCalls: new Map(),
|
||||
deferDepth: 0,
|
||||
aborted: false,
|
||||
settled: false,
|
||||
};
|
||||
session.pending.set(runId, pending);
|
||||
@@ -274,9 +297,21 @@ async function runOnce(
|
||||
const onAbort = (): void => {
|
||||
const reason = options.runState.signal?.reason;
|
||||
const abortError = reasonToError(reason, "Execution aborted");
|
||||
// Cancel any in-flight tool calls first.
|
||||
// Stop delegated work at once — this is what kills spawned subagents —
|
||||
// and refuse further bridge calls so the drain below stays bounded to
|
||||
// phases that had already started.
|
||||
pending.aborted = true;
|
||||
for (const ctrl of pending.toolCalls.values()) ctrl.abort(abortError);
|
||||
// Hard-kill the worker — only way to interrupt synchronous user code.
|
||||
// A critical host phase ignores its abort once started (isolation
|
||||
// worktree setup, merge/cherry-pick). Killing the worker now would
|
||||
// settle the cell on top of a git operation still in progress, so wait
|
||||
// for it. Hard-kill is still the only way to interrupt synchronous user
|
||||
// code, hence it stays the terminal step either way.
|
||||
const drained = pending.deferDepth > 0 ? pending.deferDrained?.promise : undefined;
|
||||
if (drained) {
|
||||
void drained.then(() => killSessionFor(session, abortError, { force: true }));
|
||||
return;
|
||||
}
|
||||
void killSessionFor(session, abortError, { force: true });
|
||||
};
|
||||
|
||||
@@ -458,6 +493,25 @@ function handleSessionMessage(session: JsSession, msg: WorkerOutbound): void {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Maintain {@link PendingRun.deferDepth} from the bridge's pause/resume status
|
||||
* events so an abort can wait out a critical `agent()` phase instead of
|
||||
* settling the cell over a half-applied merge.
|
||||
*/
|
||||
function trackDeferPhase(pending: PendingRun, event: JsStatusEvent): void {
|
||||
if (event.deferExternalAbort !== true) return;
|
||||
if (event.op === EVAL_TIMEOUT_PAUSE_OP) {
|
||||
pending.deferDepth++;
|
||||
pending.deferDrained ??= Promise.withResolvers<void>();
|
||||
return;
|
||||
}
|
||||
if (event.op !== EVAL_TIMEOUT_RESUME_OP || pending.deferDepth === 0) return;
|
||||
pending.deferDepth--;
|
||||
if (pending.deferDepth > 0) return;
|
||||
pending.deferDrained?.resolve();
|
||||
pending.deferDrained = undefined;
|
||||
}
|
||||
|
||||
async function handleToolCall(session: JsSession, msg: Extract<WorkerOutbound, { type: "tool-call" }>): Promise<void> {
|
||||
const pending = session.pending.get(msg.runId);
|
||||
if (!pending) {
|
||||
@@ -468,26 +522,42 @@ async function handleToolCall(session: JsSession, msg: Extract<WorkerOutbound, {
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (pending.aborted) {
|
||||
safeSend(session, {
|
||||
type: "tool-reply",
|
||||
id: msg.id,
|
||||
reply: { ok: false, error: { message: "Run was interrupted" } },
|
||||
});
|
||||
return;
|
||||
}
|
||||
const ctrl = new AbortController();
|
||||
pending.toolCalls.set(msg.id, ctrl);
|
||||
try {
|
||||
const value = await callSessionTool(msg.name, msg.args, {
|
||||
session: pending.toolSession,
|
||||
signal: ctrl.signal,
|
||||
emitStatus: (event: JsStatusEvent) => pending.runState.onDisplay?.({ type: "status", event }),
|
||||
emitStatus: (event: JsStatusEvent) => {
|
||||
trackDeferPhase(pending, event);
|
||||
pending.runState.onDisplay?.({ type: "status", event });
|
||||
},
|
||||
});
|
||||
safeSend(session, { type: "tool-reply", id: msg.id, reply: { ok: true, value } });
|
||||
} catch (error) {
|
||||
safeSend(session, { type: "tool-reply", id: msg.id, reply: { ok: false, error: toErrorPayload(error) } });
|
||||
} finally {
|
||||
pending.toolCalls.delete(msg.id);
|
||||
// Last call of a run whose worker result was withheld: settle it now.
|
||||
const held = pending.heldResult;
|
||||
if (held && !pending.settled && !pending.aborted && pending.toolCalls.size === 0) {
|
||||
finishPending(pending, held);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function settlePending(session: JsSession, msg: Extract<WorkerOutbound, { type: "result" }>): void {
|
||||
const pending = session.pending.get(msg.runId);
|
||||
if (!pending || pending.settled) return;
|
||||
/** Deliver a worker `result` to the waiting {@link runOnce}. */
|
||||
function finishPending(pending: PendingRun, msg: Extract<WorkerOutbound, { type: "result" }>): void {
|
||||
pending.settled = true;
|
||||
pending.heldResult = undefined;
|
||||
if (msg.ok) {
|
||||
pending.resolve({ value: undefined });
|
||||
return;
|
||||
@@ -495,6 +565,24 @@ function settlePending(session: JsSession, msg: Extract<WorkerOutbound, { type:
|
||||
pending.reject(errorFromPayload(msg.error));
|
||||
}
|
||||
|
||||
function settlePending(session: JsSession, msg: Extract<WorkerOutbound, { type: "result" }>): void {
|
||||
const pending = session.pending.get(msg.runId);
|
||||
if (!pending || pending.settled) return;
|
||||
// Once the turn is cancelled the scheduled kill is the sole settler, so a
|
||||
// late worker result can't cut the abort drain short.
|
||||
if (pending.aborted) return;
|
||||
// A cell owns every bridge call it starts. The worker finishes a run without
|
||||
// awaiting its outstanding tool calls, so `agent(...)` that is floated or
|
||||
// caught would settle the run here — `runOnce` then drops the abort listener
|
||||
// and the pending entry, leaving the subagent running with nothing able to
|
||||
// cancel it. Hold the result until the last call drains.
|
||||
if (pending.toolCalls.size > 0) {
|
||||
pending.heldResult = msg;
|
||||
return;
|
||||
}
|
||||
finishPending(pending, msg);
|
||||
}
|
||||
|
||||
async function killSessionFor(session: JsSession, error: Error, options: { force: boolean }): Promise<void> {
|
||||
if (sessions.get(session.sessionKey) === session) {
|
||||
sessions.delete(session.sessionKey);
|
||||
|
||||
@@ -13,7 +13,19 @@ import { callSessionTool, type JsStatusEvent } from "../js/tool-bridge";
|
||||
|
||||
export interface PyToolBridgeEntry {
|
||||
toolSession: ToolSession;
|
||||
/**
|
||||
* Turn-cancel handed to the tool implementation. Raw and never deferred, so
|
||||
* delegated work — above all the subagents `agent()` spawns — stops at once.
|
||||
*/
|
||||
signal?: AbortSignal;
|
||||
/**
|
||||
* Kernel-side abort, held back while a critical `agent()` phase (isolation
|
||||
* worktree setup, merge/cherry-pick) is in flight. Decides only when the host
|
||||
* may stop waiting on a call and let the kernel unwind; it is never given to
|
||||
* a tool. Keeping these separate is what stops a cancel from settling the
|
||||
* cell on top of a still-running, abort-insensitive merge.
|
||||
*/
|
||||
shieldedSignal?: AbortSignal;
|
||||
emitStatus?: (event: JsStatusEvent) => void;
|
||||
abortRequested?: () => boolean;
|
||||
}
|
||||
@@ -36,12 +48,20 @@ let serverPromise: Promise<BridgeServer> | null = null;
|
||||
* has been interrupted.
|
||||
*
|
||||
* Python invokes this bridge with blocking `urllib` requests from worker threads
|
||||
* (each `agent()` / `tool.*` call). The registered signal is the caller's real
|
||||
* abort signal, not the executor's kernel shield, so a turn cancel tears down
|
||||
* delegated work — subagents included — instead of leaving it running past the
|
||||
* cell. Calls that arrive after an abort are rejected before starting, and an
|
||||
* abort mid-call resolves the HTTP request promptly so the kernel can unwind
|
||||
* without being hard-killed.
|
||||
* (each `agent()` / `tool.*` call). Two different aborts meet here:
|
||||
*
|
||||
* - {@link PyToolBridgeEntry.signal} goes to the tool, so a turn cancel tears
|
||||
* down delegated work — subagents included — instead of leaving it running
|
||||
* past the cell.
|
||||
* - {@link PyToolBridgeEntry.shieldedSignal} decides when we may stop waiting.
|
||||
* It is deferred across a critical `agent()` phase, so a cancel landing
|
||||
* mid-merge cannot return early and let the cell settle while an
|
||||
* abort-insensitive cherry-pick is still rewriting the repo.
|
||||
*
|
||||
* Calls arriving after an abort are rejected before starting. Otherwise the
|
||||
* usual path is that the tool observes its own abort and rejects; the race only
|
||||
* matters for tools that ignore the signal, keeping the kernel unwinding
|
||||
* promptly instead of being hard-killed.
|
||||
*/
|
||||
async function callSessionToolPromptOnAbort(name: string, args: unknown, entry: PyToolBridgeEntry): Promise<unknown> {
|
||||
if (entry.abortRequested?.()) {
|
||||
@@ -52,7 +72,7 @@ async function callSessionToolPromptOnAbort(name: string, args: unknown, entry:
|
||||
signal: entry.signal,
|
||||
emitStatus: entry.emitStatus,
|
||||
});
|
||||
const signal = entry.signal;
|
||||
const signal = entry.shieldedSignal ?? entry.signal;
|
||||
if (!signal) return await call;
|
||||
if (signal.aborted) {
|
||||
void call.catch(() => {});
|
||||
|
||||
@@ -672,7 +672,7 @@ describe("agent() through eval runtimes", () => {
|
||||
expect(barrier.maxInFlight()).toBeLessThanOrEqual(2);
|
||||
});
|
||||
|
||||
it("interrupting a Python parallel() fan-out settles the kernel cleanly and preserves session state", async () => {
|
||||
it("interrupting a Python parallel() fan-out aborts in-flight subagents and preserves session state", async () => {
|
||||
using tempDir = TempDir.createSync("@omp-eval-agent-py-interrupt-");
|
||||
const settings = Settings.isolated({
|
||||
"async.enabled": false,
|
||||
@@ -683,20 +683,34 @@ describe("agent() through eval runtimes", () => {
|
||||
const { session, sessionFile, sessionId } = makeEvalSession(tempDir, "py-agent-interrupt", settings);
|
||||
mockAgents();
|
||||
// Each kernel worker thread blocks in a synchronous `urllib` bridge call,
|
||||
// joined by `parallel()`'s ThreadPoolExecutor exit. The host must keep
|
||||
// those already-started calls attached until they settle, then interrupt
|
||||
// the kernel before `parallel()` launches another wave.
|
||||
// joined by `parallel()`'s ThreadPoolExecutor exit. A turn cancel must
|
||||
// reach the subagents those calls started — the bridge is handed the real
|
||||
// signal, not the executor's kernel shield — while the kernel itself is
|
||||
// still interrupted cleanly before `parallel()` launches another wave.
|
||||
let inFlight = 0;
|
||||
let completed = 0;
|
||||
let abortedSubagents = 0;
|
||||
let markSaturated: (() => void) | undefined;
|
||||
const saturated = new Promise<void>(resolve => {
|
||||
markSaturated = resolve;
|
||||
});
|
||||
const releaseAgents = Promise.withResolvers<void>();
|
||||
// Mirrors the real executor: park until the run finishes *or* the caller's
|
||||
// signal aborts. Nothing releases these agents, so the only way the cell
|
||||
// can settle is the abort actually reaching them.
|
||||
const runSpy = vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options => {
|
||||
// task.maxConcurrency=6 → six bridge calls block at once; signal then.
|
||||
if (++inFlight >= 6) markSaturated?.();
|
||||
await releaseAgents.promise;
|
||||
const aborted = Promise.withResolvers<never>();
|
||||
const onAbort = () => aborted.reject(new Error("subagent aborted"));
|
||||
options.signal?.addEventListener("abort", onAbort, { once: true });
|
||||
try {
|
||||
await aborted.promise;
|
||||
} catch {
|
||||
abortedSubagents++;
|
||||
return singleResult(options, { output: "", aborted: true, abortReason: "aborted by user" });
|
||||
} finally {
|
||||
options.signal?.removeEventListener("abort", onAbort);
|
||||
}
|
||||
completed++;
|
||||
return singleResult(options, { output: options.assignment ?? "" });
|
||||
});
|
||||
@@ -735,14 +749,15 @@ describe("agent() through eval runtimes", () => {
|
||||
await saturated;
|
||||
await Promise.resolve();
|
||||
expect(completed).toBe(0);
|
||||
releaseAgents.resolve();
|
||||
const result = await resultPromise;
|
||||
|
||||
// Cancelled, but cleanly: no hard-kill, no orphaned bridge calls, and no
|
||||
// second fan-out wave started after the deferred abort was delivered.
|
||||
// The interrupt reached every in-flight subagent: nothing here released
|
||||
// them, so the cell could only settle because the abort propagated.
|
||||
expect(abortedSubagents).toBe(6);
|
||||
expect(completed).toBe(0);
|
||||
// Cancelled, but cleanly: no hard-kill, and no second fan-out wave started.
|
||||
expect(result.cancelled).toBe(true);
|
||||
expect(result.output).not.toContain("Python kernel shutdown");
|
||||
expect(completed).toBe(6);
|
||||
expect(runSpy).toHaveBeenCalledTimes(6);
|
||||
|
||||
// The persistent kernel survived the interrupt: prior state is intact.
|
||||
@@ -993,6 +1008,65 @@ describe("agent() through eval runtimes", () => {
|
||||
expect(ops.at(-1)).toBe(EVAL_TIMEOUT_RESUME_OP);
|
||||
expect(idle.signal.aborted).toBe(false);
|
||||
});
|
||||
|
||||
it("interrupting a JavaScript agent() aborts it at once but waits out its critical phase", async () => {
|
||||
// Regression: `onAbort` used to hard-kill the worker straight away, which
|
||||
// rejected the run while the untracked `handleToolCall` promise carried on
|
||||
// — so an isolation merge could keep cherry-picking after the cell had
|
||||
// already returned. Mirrors the Python bridge's shielded-signal contract.
|
||||
//
|
||||
// Asserted as an ordering, not a duration: the agent call must finish
|
||||
// before the cell settles. Killing early inverts the two.
|
||||
using tempDir = TempDir.createSync("@omp-eval-agent-js-interrupt-");
|
||||
const { session, sessionFile } = makeEvalSession(tempDir, "js-agent-interrupt");
|
||||
mockAgents();
|
||||
|
||||
const order: string[] = [];
|
||||
const inFlight = Promise.withResolvers<void>();
|
||||
const sawAbort = Promise.withResolvers<void>();
|
||||
const release = Promise.withResolvers<void>();
|
||||
// Stands in for the isolation merge: notices the abort, then keeps going.
|
||||
vi.spyOn(taskExecutor, "runSubprocess").mockImplementation(async options => {
|
||||
options.signal?.addEventListener(
|
||||
"abort",
|
||||
() => {
|
||||
order.push("agent-saw-abort");
|
||||
sawAbort.resolve();
|
||||
},
|
||||
{ once: true },
|
||||
);
|
||||
inFlight.resolve();
|
||||
await release.promise;
|
||||
order.push("agent-returned");
|
||||
return singleResult(options, { output: "merged" });
|
||||
});
|
||||
|
||||
const ac = new AbortController();
|
||||
const cell = executeJs('return await agent("merge");', {
|
||||
cwd: tempDir.path(),
|
||||
sessionId: "agent-bridge-js-interrupt",
|
||||
session,
|
||||
sessionFile,
|
||||
signal: ac.signal,
|
||||
}).finally(() => {
|
||||
order.push("cell-settled");
|
||||
});
|
||||
|
||||
await inFlight.promise;
|
||||
ac.abort(new Error("external interrupt"));
|
||||
// Delegated work is notified immediately, before anything is released.
|
||||
await sawAbort.promise;
|
||||
// Drain the microtask queue. An abort that settled the run outright would
|
||||
// have resolved the cell by now; no wall clock is involved.
|
||||
for (let i = 0; i < 200; i++) await Promise.resolve();
|
||||
expect(order).toEqual(["agent-saw-abort"]);
|
||||
|
||||
release.resolve();
|
||||
const result = await cell;
|
||||
|
||||
expect(order).toEqual(["agent-saw-abort", "agent-returned", "cell-settled"]);
|
||||
expect(result.cancelled).toBe(true);
|
||||
}, 30_000);
|
||||
});
|
||||
|
||||
describe("runEvalAgent isolation", () => {
|
||||
|
||||
@@ -1,4 +1,7 @@
|
||||
import { afterEach, describe, expect, it, vi } from "bun:test";
|
||||
import { afterAll, afterEach, describe, expect, it, vi } from "bun:test";
|
||||
import { type } from "@oh-my-pi/omptype";
|
||||
import type { AgentTool } from "@oh-my-pi/pi-agent-core";
|
||||
import { Settings } from "../../src/config/settings";
|
||||
import {
|
||||
EVAL_TIMEOUT_PAUSE_OP,
|
||||
EVAL_TIMEOUT_RESUME_OP,
|
||||
@@ -11,10 +14,33 @@ import type { KernelDisplayOutput } from "../../src/eval/py/display";
|
||||
import * as pyToolBridge from "../../src/eval/py/tool-bridge";
|
||||
import type { ToolSession } from "../../src/tools";
|
||||
|
||||
/** Minimal `ToolSession` exposing `tools` to the eval tool bridge. */
|
||||
function makeToolSession(...tools: AgentTool[]): ToolSession {
|
||||
return {
|
||||
cwd: process.cwd(),
|
||||
hasUI: false,
|
||||
settings: Settings.isolated({ "async.enabled": false }),
|
||||
taskDepth: 0,
|
||||
enableLsp: false,
|
||||
getSessionFile: () => null,
|
||||
getSessionSpawns: () => "*",
|
||||
getActiveModelString: () => "p/active",
|
||||
getModelString: () => "p/fallback",
|
||||
getArtifactsDir: () => null,
|
||||
getSessionId: () => "bridge-timeout-session",
|
||||
getEvalSessionId: () => "bridge-timeout-eval-session",
|
||||
getToolByName: name => tools.find(tool => tool.name === name),
|
||||
};
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
await pyToolBridge.disposePyToolBridge();
|
||||
});
|
||||
|
||||
describe("withBridgeTimeoutPause", () => {
|
||||
it("emits one pause before the operation and one resume after it settles", async () => {
|
||||
const events: JsStatusEvent[] = [];
|
||||
@@ -214,7 +240,7 @@ it("hands the tool bridge the unshielded signal so a deferred phase still cancel
|
||||
code: "agent('slow')",
|
||||
options: {
|
||||
signal: abortController.signal,
|
||||
toolSession: {} as ToolSession,
|
||||
toolSession: makeToolSession(),
|
||||
bridgeSessionId: "bridge-session",
|
||||
},
|
||||
runIdPrefix: "test",
|
||||
@@ -237,3 +263,83 @@ it("hands the tool bridge the unshielded signal so a deferred phase still cancel
|
||||
const result = await resultPromise;
|
||||
expect(result.cancelled).toBe(true);
|
||||
});
|
||||
|
||||
it("holds the cell open through a deferred phase while still aborting the tool at once", async () => {
|
||||
// Regression: the tool call and the wait-for-result race must use different
|
||||
// aborts. Sharing the raw signal answered the HTTP call the moment the turn
|
||||
// was cancelled, letting the cell settle and the bridge unregister on top of
|
||||
// an abort-insensitive isolation merge still rewriting the repo.
|
||||
//
|
||||
// The discriminator is the reply itself, not its timing: losing the race
|
||||
// makes the bridge answer `ok: false, aborted`, while a correctly deferred
|
||||
// wait answers with the tool's real value once the phase releases.
|
||||
const bridge = await pyToolBridge.ensurePyToolBridge();
|
||||
const abortController = new AbortController();
|
||||
const toolStarted = Promise.withResolvers<void>();
|
||||
const toolSawAbort = Promise.withResolvers<void>();
|
||||
const releaseTool = Promise.withResolvers<void>();
|
||||
|
||||
// Stands in for `runStructuredSubagent`: observes its abort at once (the
|
||||
// subagent dies) but keeps working afterwards, exactly like a cherry-pick
|
||||
// that never looks at a signal.
|
||||
const parkedTool: AgentTool = {
|
||||
name: "merge",
|
||||
label: "Merge",
|
||||
description: "Parks until released, ignoring its abort",
|
||||
parameters: type({}),
|
||||
execute: async (_id, _args, signal) => {
|
||||
signal?.addEventListener("abort", () => toolSawAbort.resolve(), { once: true });
|
||||
toolStarted.resolve();
|
||||
await releaseTool.promise;
|
||||
return { content: [{ type: "text", text: "merged" }] };
|
||||
},
|
||||
};
|
||||
const toolSession = makeToolSession(parkedTool);
|
||||
|
||||
const bridgeSessionId = `bridge-${crypto.randomUUID()}`;
|
||||
let reply: { ok: boolean; value?: unknown; error?: string } | undefined;
|
||||
const kernel: GenericKernel<Record<string, string | null>> = {
|
||||
async execute(_code, options) {
|
||||
options.onDisplay({
|
||||
type: "status",
|
||||
event: { op: EVAL_TIMEOUT_PAUSE_OP, deferExternalAbort: true },
|
||||
} satisfies KernelDisplayOutput);
|
||||
// Mirrors the Python prelude's blocking loopback call.
|
||||
const pending = fetch(`${bridge.url}/v1/tool`, {
|
||||
method: "POST",
|
||||
headers: { authorization: `Bearer ${bridge.token}`, "content-type": "application/json" },
|
||||
body: JSON.stringify({ session: bridgeSessionId, run: options.id, name: "merge", args: {} }),
|
||||
}).then(res => res.json() as Promise<{ ok: boolean; value?: unknown; error?: string }>);
|
||||
|
||||
await toolStarted.promise;
|
||||
abortController.abort(new Error("external interrupt"));
|
||||
// The raw signal reached the tool: delegated work is already dying.
|
||||
await toolSawAbort.promise;
|
||||
|
||||
releaseTool.resolve();
|
||||
reply = await pending;
|
||||
options.onDisplay({
|
||||
type: "status",
|
||||
event: { op: EVAL_TIMEOUT_RESUME_OP, deferExternalAbort: true },
|
||||
} satisfies KernelDisplayOutput);
|
||||
return { status: "ok", cancelled: false, timedOut: false };
|
||||
},
|
||||
};
|
||||
|
||||
const result = await executeWithKernelBase({
|
||||
kernel,
|
||||
code: "agent('isolated')",
|
||||
options: { signal: abortController.signal, toolSession, bridgeSessionId },
|
||||
runIdPrefix: "test",
|
||||
errorLogLabel: "test",
|
||||
cancelledErrorClass: TestCancelledError,
|
||||
buildKernelEnvPatch: () => ({}),
|
||||
formatKernelTimeoutAnnotation: () => "kernel timed out",
|
||||
formatTimeoutAnnotation: () => "timed out",
|
||||
});
|
||||
|
||||
// The host waited out the critical phase instead of bailing on the cancel.
|
||||
expect(reply).toEqual({ ok: true, value: "merged" });
|
||||
// ...and the turn is still reported as cancelled.
|
||||
expect(result.cancelled).toBe(true);
|
||||
});
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
|
||||
import * as fs from "node:fs";
|
||||
import * as path from "node:path";
|
||||
import { type } from "@oh-my-pi/omptype";
|
||||
import type { AgentTool } from "@oh-my-pi/pi-agent-core";
|
||||
import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
import { Settings } from "../../src/config/settings";
|
||||
import {
|
||||
@@ -22,9 +24,14 @@ interface FakeWorkerBehavior {
|
||||
exitOnClose: boolean;
|
||||
settleRuns: boolean;
|
||||
errorOnStart?: boolean;
|
||||
/**
|
||||
* Reproduces `WorkerCore#runOne` for a floated bridge call: start a tool call
|
||||
* and report the run finished in the same turn, without awaiting the call.
|
||||
*/
|
||||
floatingToolCall?: string;
|
||||
}
|
||||
|
||||
function makeSession(cwd: string): ToolSession {
|
||||
function makeSession(cwd: string, ...tools: AgentTool[]): ToolSession {
|
||||
return {
|
||||
cwd,
|
||||
hasUI: false,
|
||||
@@ -42,6 +49,7 @@ function makeSession(cwd: string): ToolSession {
|
||||
getArtifactsDir: () => null,
|
||||
getSessionId: () => "test-session",
|
||||
getEvalSessionId: () => "test-eval-session",
|
||||
getToolByName: name => tools.find(tool => tool.name === name),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -108,6 +116,20 @@ function installFakeWorker(stats: FakeWorkerStats, behavior: FakeWorkerBehavior)
|
||||
postMessage(message: unknown): void {
|
||||
if (!message || typeof message !== "object") return;
|
||||
const typed = message as { type?: string; runId?: string };
|
||||
if (typed.type === "run" && typed.runId && behavior.floatingToolCall) {
|
||||
const runId = typed.runId;
|
||||
queueMicrotask(() => {
|
||||
this.#emitMessage({
|
||||
type: "tool-call",
|
||||
runId,
|
||||
id: `tc-${runId}`,
|
||||
name: behavior.floatingToolCall,
|
||||
args: {},
|
||||
});
|
||||
this.#emitMessage({ type: "result", runId, ok: true });
|
||||
});
|
||||
return;
|
||||
}
|
||||
if (typed.type === "run" && typed.runId && behavior.settleRuns) {
|
||||
queueMicrotask(() => this.#emitMessage({ type: "result", runId: typed.runId, ok: true }));
|
||||
return;
|
||||
@@ -366,6 +388,57 @@ describe("JavaScript eval worker lifecycle", () => {
|
||||
// The errored primary worker is torn down before the inline retry takes over.
|
||||
expect(stats.terminateCalls).toBe(1);
|
||||
});
|
||||
|
||||
it("holds a finished run until its floated tool call drains", async () => {
|
||||
// Regression, and the sharpest form of the runaway `agent()` fan-out — no
|
||||
// cancellation required. `WorkerCore#runOne` reports a finished run without
|
||||
// awaiting outstanding tool calls, so `agent(...)` with no `await` used to
|
||||
// settle the cell immediately; `runOnce` then dropped the run's abort
|
||||
// listener and pending entry, leaving the subagent running with nothing
|
||||
// able to cancel it. A cell must own every bridge call it starts.
|
||||
using tempDir = TempDir.createSync("@omp-js-worker-float-");
|
||||
const stats: FakeWorkerStats = { closeRequests: 0, terminateCalls: 0 };
|
||||
installFakeWorker(stats, { exitOnClose: true, settleRuns: false, floatingToolCall: "park" });
|
||||
|
||||
const started = Promise.withResolvers<void>();
|
||||
const release = Promise.withResolvers<void>();
|
||||
let toolReturned = false;
|
||||
const park: AgentTool = {
|
||||
name: "park",
|
||||
label: "Park",
|
||||
description: "Parks until released",
|
||||
parameters: type({}),
|
||||
execute: async () => {
|
||||
started.resolve();
|
||||
await release.promise;
|
||||
toolReturned = true;
|
||||
return { content: [{ type: "text", text: "parked" }] };
|
||||
},
|
||||
};
|
||||
const session = makeSession(tempDir.path(), park);
|
||||
|
||||
let settled = false;
|
||||
const cell = executeJs('tool.park({});\nreturn "floated";', {
|
||||
cwd: tempDir.path(),
|
||||
sessionId: `js-float:${crypto.randomUUID()}`,
|
||||
session,
|
||||
}).finally(() => {
|
||||
settled = true;
|
||||
});
|
||||
|
||||
await started.promise;
|
||||
// The fake worker already reported the run finished, in the same microtask
|
||||
// that started the call. Draining the microtask queue is exact here: the
|
||||
// harness is queueMicrotask-driven, with no IPC or timers in the path.
|
||||
for (let i = 0; i < 50; i++) await Promise.resolve();
|
||||
expect(settled).toBe(false);
|
||||
expect(toolReturned).toBe(false);
|
||||
|
||||
release.resolve();
|
||||
const result = await cell;
|
||||
expect(toolReturned).toBe(true);
|
||||
expect(result.exitCode).toBe(0);
|
||||
});
|
||||
});
|
||||
|
||||
describe.skipIf(process.platform === "win32")("JavaScript eval process isolation", () => {
|
||||
|
||||
Reference in New Issue
Block a user