From ee15990840fa9db51e45827653bd4c675e04b330 Mon Sep 17 00:00:00 2001 From: Seth Date: Tue, 11 Aug 2026 11:28:57 -0700 Subject: [PATCH 1/3] fix(security): cancel admitted RLM work on host teardown Propagate kernel host abort signals through rlm.run and cancel admitted child runs when their host is disposed, while removing listeners after settlement. Extracted from the independently authored security stack commits: - a3ba5dba6e8c280caa20cdca88f057810c35bb1a - 8d47a2e9191b6d2b516c979b19c7a8fb9b3b2999 --- .../coding-agent/src/core/agent-session.ts | 21 ++++- .../coding-agent/src/core/kernel/index.ts | 28 ++++-- packages/coding-agent/src/core/rlm-runtime.ts | 17 ++-- .../test/agent-session-recursion.test.ts | 89 +++++++++++++------ 4 files changed, 114 insertions(+), 41 deletions(-) diff --git a/packages/coding-agent/src/core/agent-session.ts b/packages/coding-agent/src/core/agent-session.ts index 46d2bd4f8..d5d8803e2 100644 --- a/packages/coding-agent/src/core/agent-session.ts +++ b/packages/coding-agent/src/core/agent-session.ts @@ -8760,8 +8760,8 @@ export class AgentSession { /** Typed handlers for host requests arriving from the IPython kernel comm bridge. */ private _createKernelHostHandlers(): HostRequestHandlers { const handlers: HostRequestHandlers = { - "rlm.run": createRlmRunHostHandler(async ({ prompt, kwargs, cellSourceCode }) => ({ - ...(await this.runRlmChild(prompt, kwargs, cellSourceCode)), + "rlm.run": createRlmRunHostHandler(async ({ prompt, kwargs, cellSourceCode }, signal) => ({ + ...(await this.runRlmChild(prompt, kwargs, cellSourceCode, signal)), })), "rlm.find_models": createRlmFindModelsHostHandler((query, limit) => this.findRlmModels(query, limit)), "rlm.list_subagents": createRlmListSubagentsHostHandler(() => this.listRlmSubagents()), @@ -9685,7 +9685,9 @@ export class AgentSession { prompt: string, kwargs: Record = {}, spawnCode?: string, + signal?: AbortSignal, ): Promise { + signal?.throwIfAborted(); const { name: rawName, model: rawModel, ...unsupported } = kwargs; const unsupportedKwargs = Object.keys(unsupported); if (unsupportedKwargs.length > 0) { @@ -9712,12 +9714,14 @@ export class AgentSession { } finally { if (requestedSessionName) this._pendingRlmSubagentSessionNames.delete(requestedSessionName); } + signal?.throwIfAborted(); if (this._disposed || this._disposing) throw new Error("Cannot spawn a subagent after its parent was disposed"); const childSessionDir = this._createChildRlmSessionDir(); const childNodeId = basename(childSessionDir); const sessionName = requestedSessionName ?? createDefaultRlmSubagentSessionName(prompt, childNodeId); if (!requestedSessionName) await this._assertRlmSubagentSessionNameAvailable(sessionName); + signal?.throwIfAborted(); const startedAt = Date.now(); const parentAssistantForUsage = this._findLastAssistantMessage(); const label = rlmChildLabel(prompt); @@ -9741,6 +9745,15 @@ export class AgentSession { if (run.status === "cancelled") throw new Error(run.error ?? "RLM child cancelled"); }; this._activeRlmChildRuns.set(run.id, run); + const abortFromHost = () => { + const reason = signal?.reason; + this._cancelRlmChildRun(run, reason instanceof Error ? reason.message : "IPython kernel host request aborted"); + }; + if (signal?.aborted) { + abortFromHost(); + } else { + signal?.addEventListener("abort", abortFromHost, { once: true }); + } const emitChildUpdate = () => { const childModel = childSession?.model ?? modelSelection.model; this._emit({ @@ -9995,6 +10008,7 @@ export class AgentSession { } } } finally { + signal?.removeEventListener("abort", abortFromHost); if (run.detachedDeletion && childRuntime) { try { await this._deleteRlmSubagentSession(run.id, childRuntime.session); @@ -10036,8 +10050,9 @@ export class AgentSession { prompt: string, kwargs: Record = {}, spawnCode?: string, + signal?: AbortSignal, ): Promise { - return this._startRlmChildRun(prompt, kwargs, spawnCode); + return this._startRlmChildRun(prompt, kwargs, spawnCode, signal); } // ========================================================================= diff --git a/packages/coding-agent/src/core/kernel/index.ts b/packages/coding-agent/src/core/kernel/index.ts index b760a2e1e..c1a8e8d7d 100644 --- a/packages/coding-agent/src/core/kernel/index.ts +++ b/packages/coding-agent/src/core/kernel/index.ts @@ -59,7 +59,10 @@ export const HOST_COMM_TARGET = "host.request"; * Handles one typed request from Python code running in the kernel. * The returned record is sent back verbatim as the comm reply payload. */ -export type HostRequestHandler = (payload: Record) => Promise>; +export type HostRequestHandler = ( + payload: Record, + signal?: AbortSignal, +) => Promise>; /** Host request handlers keyed by request type (e.g. "rlm.run", "goal.complete"). */ export type HostRequestHandlers = Record; @@ -540,6 +543,7 @@ export class KernelManager { // attribute their spawning program. private lastCellCode?: string; private readonly inFlightHostRequests = new Set>(); + private hostRequestController = new AbortController(); private state: "idle" | "starting" | "running" | "shutdown" = "idle"; /** Memoized so concurrent callers all await the same in-flight startup. */ private startPromise?: Promise; @@ -582,6 +586,9 @@ export class KernelManager { private async doStart(startOptions: KernelStartOptions): Promise { if (this.state !== "idle") return; + if (this.hostRequestController.signal.aborted) { + this.hostRequestController = new AbortController(); + } this.state = "starting"; installSignalHandlersOnce(); // Tracked from the moment startup begins so session cleanup and signal @@ -1224,9 +1231,10 @@ export class KernelManager { } this.handledHostRequestCommIds.add(commId); + const signal = this.hostRequestController.signal; const task = (async () => { try { - const result = await this.handleHostRequest(data); + const result = await this.handleHostRequest(data, signal); try { await this.sendCommMessage(commId, { status: "ok", ...result }); } catch (replyError) { @@ -1251,7 +1259,7 @@ export class KernelManager { }); } - private async handleHostRequest(data: unknown): Promise> { + private async handleHostRequest(data: unknown, signal: AbortSignal): Promise> { if (!isRecord(data)) { throw new Error("host request payload must be an object"); } @@ -1267,7 +1275,7 @@ export class KernelManager { // the in-flight execution; detached spawns (asyncio.create_task) fire after // the scheduling cell goes idle, so fall back to that last cell's source. const cellSourceCode = this.activeExecution?.code ?? this.lastCellCode; - return handler({ ...data, cellSourceCode }); + return handler({ ...data, cellSourceCode }, signal); } private async sendCommMessage(commId: string, data: Record): Promise { @@ -1286,6 +1294,7 @@ export class KernelManager { } private cleanupResources(killSignal: NodeJS.Signals = "SIGTERM"): void { + this.abortHostRequests("IPython kernel stopped"); this.clearSnapshotTimer(); this.lateSentAgentMessageHandlers.clear(); if (this.forkedLivenessTimer) { @@ -1325,6 +1334,12 @@ export class KernelManager { this.startPromise = undefined; } + private abortHostRequests(message: string): void { + if (!this.hostRequestController.signal.aborted) { + this.hostRequestController.abort(new Error(message)); + } + } + private async waitForHostRequestsToSettle(tasks: Promise[], timeoutMs: number): Promise { let timeout: ReturnType | undefined; const timeoutPromise = new Promise<"timeout">((resolve) => { @@ -1356,6 +1371,7 @@ export class KernelManager { if (opts.snapshot) { await this.flushSnapshotForDispose(); } + this.abortHostRequests("IPython kernel shut down"); this.state = "shutdown"; liveKernels.delete(this); @@ -1393,6 +1409,7 @@ export class KernelManager { } async kill(): Promise { + this.abortHostRequests("IPython kernel killed"); this.state = "shutdown"; liveKernels.delete(this); this.cleanupResources("SIGKILL"); @@ -1501,10 +1518,10 @@ export class KernelManager { return (async () => { // Final namespace flush while the kernel is still live (session end / reload). await this.flushSnapshotForDispose(); + this.abortHostRequests("IPython kernel disposed"); this.state = "shutdown"; liveKernels.delete(this); const inFlightHostRequests = [...this.inFlightHostRequests]; - // TODO: plumb AbortSignal through AgentSession.prompt so disposal can cancel long-running child loops. try { if (inFlightHostRequests.length > 0) { await this.waitForHostRequestsToSettle(inFlightHostRequests, HOST_REQUEST_DISPOSE_TIMEOUT_MS); @@ -1517,6 +1534,7 @@ export class KernelManager { /** Synchronous best-effort cleanup. Safe to call from `process.on('exit')`. */ disposeSync(): void { + this.abortHostRequests("IPython kernel disposed"); this.state = "shutdown"; liveKernels.delete(this); // TODO: replace this best-effort hard-exit path if Node exposes an awaitable process-exit cleanup hook. diff --git a/packages/coding-agent/src/core/rlm-runtime.ts b/packages/coding-agent/src/core/rlm-runtime.ts index e89472fce..9a5c66652 100644 --- a/packages/coding-agent/src/core/rlm-runtime.ts +++ b/packages/coding-agent/src/core/rlm-runtime.ts @@ -49,7 +49,7 @@ export interface RlmFindModelsResult { models: RlmModelMatch[]; } -export type RlmRunHandler = (request: RlmRunRequest) => Promise>; +export type RlmRunHandler = (request: RlmRunRequest, signal?: AbortSignal) => Promise>; export type RlmListSubagentsHandler = () => RlmListSubagentsResult | Promise; export type RlmDeleteSubagentHandler = (target: string) => Promise; export type RlmFindModelsHandler = (query: string, limit: number) => RlmFindModelsResult | Promise; @@ -150,17 +150,20 @@ export function findRlmModelMatches(query: string, models: Model[], limit: /** Adapt an RlmRunHandler into the typed "rlm.run" handler for the kernel host bridge. */ export function createRlmRunHostHandler(handler: RlmRunHandler): HostRequestHandler { - return async (payload) => { + return async (payload, signal) => { if (typeof payload.prompt !== "string") { throw new Error("rlm.run prompt must be a string"); } const kwargs = isRecord(payload.kwargs) ? payload.kwargs : {}; const cellSourceCode = typeof payload.cellSourceCode === "string" ? payload.cellSourceCode : undefined; - const result = await handler({ - prompt: payload.prompt, - kwargs, - cellSourceCode, - }); + const result = await handler( + { + prompt: payload.prompt, + kwargs, + cellSourceCode, + }, + signal, + ); return result as unknown as Record; }; } diff --git a/packages/coding-agent/test/agent-session-recursion.test.ts b/packages/coding-agent/test/agent-session-recursion.test.ts index f066f02b1..bd1867930 100644 --- a/packages/coding-agent/test/agent-session-recursion.test.ts +++ b/packages/coding-agent/test/agent-session-recursion.test.ts @@ -2893,6 +2893,57 @@ describe("AgentSession rlm recursion", () => { await waitFor(() => rootRun.status === "done"); }); + it("cancels an admitted child in promptAndWait when its kernel host is disposed", async () => { + let releaseChild: () => void = () => {}; + const release = new Promise((resolve) => { + releaseChild = resolve; + }); + let childStarted = false; + const root = createSession({ + streamFn: (_model, context) => { + const text = userText(context); + const stream = createAssistantMessageEventStream(); + if (text === "kernel-owned shard") { + childStarted = true; + void release.then(() => { + stream.push({ type: "done", reason: "stop", message: assistantMessage(`child answer: ${text}`) }); + }); + } + return stream; + }, + }); + const manager = new KernelManager({ + python: process.execPath, + hostHandlers: (root as unknown as InspectableRlmSession)._createKernelHostHandlers(), + }); + const replies: CapturedCommReply[] = []; + const kernel = manager as unknown as KernelCommTestApi; + kernel.sendCommMessage = async (commId, data) => { + replies.push({ commId, data }); + }; + + try { + kernel.handleCommMessage(rlmCommOpen("comm-real-child", "kernel-owned shard")); + await waitFor(() => replies.some((reply) => reply.commId === "comm-real-child")); + await waitFor(() => childStarted); + const runs = (root as unknown as InspectableRlmSession)._activeRlmChildRuns; + const run = [...runs.values()][0]; + if (!run?.session) throw new Error("Missing admitted child session"); + const childAbort = vi.spyOn(run.session, "abort"); + + await manager.dispose(); + + expect(run.status).toBe("cancelled"); + expect(run.error).toBe("IPython kernel disposed"); + expect(childAbort).toHaveBeenCalledTimes(1); + releaseChild(); + await waitFor(() => !runs.has(run.id)); + } finally { + releaseChild(); + await manager.dispose(); + } + }); + it("runs parallel rlm comm requests independently", async () => { let active = 0; let maxActive = 0; @@ -3122,26 +3173,23 @@ print(_result.name) } }); - it("waits for in-flight rlm comm work during dispose and buffers failures", async () => { + it("aborts in-flight rlm comm work during dispose and buffers failures", async () => { let started = false; let handlerSettled = false; - let released = false; - let releaseChild: () => void = () => {}; - const release = new Promise((resolve) => { - releaseChild = () => { - if (released) return; - released = true; - resolve(); - }; - }); + let receivedSignal: AbortSignal | undefined; const manager = new KernelManager({ python: process.execPath, hostHandlers: { - "rlm.run": createRlmRunHostHandler(async () => { + "rlm.run": createRlmRunHostHandler(async (_request, signal) => { started = true; + receivedSignal = signal; try { - await release; - throw new Error("child failed after dispose"); + await new Promise((_resolve, reject) => { + const onAbort = () => reject(signal?.reason ?? new Error("aborted")); + if (signal?.aborted) onAbort(); + else signal?.addEventListener("abort", onAbort, { once: true }); + }); + return {}; } finally { handlerSettled = true; } @@ -3152,21 +3200,11 @@ print(_result.name) try { const kernel = manager as unknown as KernelCommTestApi; - kernel.handleCommMessage(rlmCommOpen("comm-dispose", "slow child")); await waitFor(() => started); - const disposePromise = manager.dispose(); - let disposeSettled = false; - const trackedDispose = disposePromise.then(() => { - disposeSettled = true; - }); - - await sleep(25); - expect(disposeSettled).toBe(false); - - releaseChild(); - await expectSettlesWithin(trackedDispose, 1000); + await expectSettlesWithin(manager.dispose(), 1000); + expect(receivedSignal?.aborted).toBe(true); expect(handlerSettled).toBe(true); const kernelStderr = (manager as unknown as { kernelStderr: string }).kernelStderr; @@ -3174,7 +3212,6 @@ print(_result.name) expect(kernelStderr).toContain("[kernel] failed to send host request error reply for comm comm-dispose"); expect(stderrSpy).not.toHaveBeenCalled(); } finally { - releaseChild(); await manager.dispose(); stderrSpy.mockRestore(); } From 0c0c75c11be214623474a9f7b406cb29302ccc7a Mon Sep 17 00:00:00 2001 From: Seth Date: Thu, 13 Aug 2026 16:30:08 -0700 Subject: [PATCH 2/3] fix(coding-agent): abort host work before snapshots --- .../coding-agent/src/core/kernel/index.ts | 4 +- .../test/agent-session-recursion.test.ts | 37 +++++++++++++++++++ 2 files changed, 39 insertions(+), 2 deletions(-) diff --git a/packages/coding-agent/src/core/kernel/index.ts b/packages/coding-agent/src/core/kernel/index.ts index 1f22e5699..757eeea65 100644 --- a/packages/coding-agent/src/core/kernel/index.ts +++ b/packages/coding-agent/src/core/kernel/index.ts @@ -1444,10 +1444,10 @@ export class KernelManager { } // Best-effort final flush (bounded) before teardown — used by signal handlers // so a SIGINT/SIGTERM exit doesn't lose work the debounced snapshot hasn't saved. + this.abortHostRequests("IPython kernel shut down"); if (opts.snapshot) { await this.flushSnapshotForDispose(); } - this.abortHostRequests("IPython kernel shut down"); this.state = "shutdown"; liveKernels.delete(this); @@ -1592,9 +1592,9 @@ export class KernelManager { /** Graceful cleanup. Waits briefly for in-flight host request handlers before closing sockets. */ dispose(): Promise { return (async () => { + this.abortHostRequests("IPython kernel disposed"); // Final namespace flush while the kernel is still live (session end / reload). await this.flushSnapshotForDispose(); - this.abortHostRequests("IPython kernel disposed"); this.state = "shutdown"; liveKernels.delete(this); const inFlightHostRequests = [...this.inFlightHostRequests]; diff --git a/packages/coding-agent/test/agent-session-recursion.test.ts b/packages/coding-agent/test/agent-session-recursion.test.ts index bd1867930..c43cdccab 100644 --- a/packages/coding-agent/test/agent-session-recursion.test.ts +++ b/packages/coding-agent/test/agent-session-recursion.test.ts @@ -144,6 +144,7 @@ interface KernelExecuteTestApi { start: () => Promise; state: "idle" | "starting" | "running" | "shutdown"; activeExecution?: unknown; + snapshotState?: () => Promise; shell?: { send(frames: Buffer[]): Promise; close(): void; @@ -3173,6 +3174,42 @@ print(_result.name) } }); + it.each(["dispose", "shutdown"] as const)( + "aborts host work before the final snapshot during %s", + async (teardown) => { + let hostSignal: AbortSignal | undefined; + const manager = new KernelManager({ + python: process.execPath, + snapshot: { path: "/tmp/unused.dill", manifestPath: "/tmp/unused.json" }, + hostHandlers: { + "rlm.run": createRlmRunHostHandler(async (_request, signal) => { + hostSignal = signal; + await new Promise((resolve) => { + if (signal?.aborted) resolve(); + else signal?.addEventListener("abort", () => resolve(), { once: true }); + }); + return {}; + }), + }, + }); + const kernel = manager as unknown as KernelCommTestApi & KernelExecuteTestApi; + kernel.state = "running"; + kernel.sendCommMessage = async () => {}; + kernel.snapshotState = async () => { + expect(hostSignal?.aborted).toBe(true); + return null; + }; + + kernel.handleCommMessage(rlmCommOpen(`comm-${teardown}`, "slow child")); + await waitFor(() => hostSignal !== undefined); + + if (teardown === "dispose") await manager.dispose(); + else await manager.shutdown({ snapshot: true }); + + expect(hostSignal?.aborted).toBe(true); + }, + ); + it("aborts in-flight rlm comm work during dispose and buffers failures", async () => { let started = false; let handlerSettled = false; From 6a6ed7d36b76b3c2d4e3fc69d5bfbd883a286cbf Mon Sep 17 00:00:00 2001 From: Seth Date: Thu, 13 Aug 2026 16:38:15 -0700 Subject: [PATCH 3/3] fix(coding-agent): clean up aborted child admission --- .../coding-agent/src/core/agent-session.ts | 11 ++++-- .../test/agent-session-recursion.test.ts | 37 +++++++++++++++++++ 2 files changed, 45 insertions(+), 3 deletions(-) diff --git a/packages/coding-agent/src/core/agent-session.ts b/packages/coding-agent/src/core/agent-session.ts index d5d8803e2..3e4f1fe81 100644 --- a/packages/coding-agent/src/core/agent-session.ts +++ b/packages/coding-agent/src/core/agent-session.ts @@ -15,7 +15,7 @@ import { AsyncLocalStorage } from "node:async_hooks"; import { randomUUID } from "node:crypto"; -import { existsSync, mkdirSync, mkdtempSync, readFileSync, writeFileSync } from "node:fs"; +import { existsSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { basename, dirname, join, resolve } from "node:path"; import { @@ -9720,8 +9720,13 @@ export class AgentSession { const childSessionDir = this._createChildRlmSessionDir(); const childNodeId = basename(childSessionDir); const sessionName = requestedSessionName ?? createDefaultRlmSubagentSessionName(prompt, childNodeId); - if (!requestedSessionName) await this._assertRlmSubagentSessionNameAvailable(sessionName); - signal?.throwIfAborted(); + try { + if (!requestedSessionName) await this._assertRlmSubagentSessionNameAvailable(sessionName); + signal?.throwIfAborted(); + } catch (error) { + rmSync(childSessionDir, { recursive: true, force: true }); + throw error; + } const startedAt = Date.now(); const parentAssistantForUsage = this._findLastAssistantMessage(); const label = rlmChildLabel(prompt); diff --git a/packages/coding-agent/test/agent-session-recursion.test.ts b/packages/coding-agent/test/agent-session-recursion.test.ts index c43cdccab..48e66e002 100644 --- a/packages/coding-agent/test/agent-session-recursion.test.ts +++ b/packages/coding-agent/test/agent-session-recursion.test.ts @@ -2169,6 +2169,43 @@ describe("AgentSession rlm recursion", () => { ); }); + it("removes an unadmitted child directory when host cancellation wins name validation", async () => { + let releaseNameCheck: () => void = () => {}; + const nameCheckGate = new Promise((resolve) => { + releaseNameCheck = resolve; + }); + let nameCheckStarted = false; + const controller = new AbortController(); + const root = createSession({ + agentMessageController: { + assertSessionNameAvailable: async () => { + nameCheckStarted = true; + await nameCheckGate; + }, + listAgents: async () => ({ + current: { activeSessionId: "parent-active", sessionId: "parent-session" }, + agents: [], + }), + sendAgentMessage: async () => { + throw new Error("unexpected send"); + }, + }, + }); + + const run = root.runRlmChild("cancel during default name validation", {}, undefined, controller.signal); + await waitFor(() => nameCheckStarted); + const childDirs = readdirSync(tempDir, { recursive: true, encoding: "utf8" }).filter((path) => + basename(path).startsWith("sub-"), + ); + expect(childDirs).toHaveLength(1); + + controller.abort(new Error("host disposed")); + releaseNameCheck(); + + await expect(run).rejects.toThrow("host disposed"); + expect(existsSync(join(tempDir, childDirs[0] ?? ""))).toBe(false); + }); + it("cancels active rlm children when the parent session is disposed", async () => { let releaseChild: () => void = () => {}; const release = new Promise((resolve) => {