From d84a54e579fa1269377a3e44ed8a2bcc92436594 Mon Sep 17 00:00:00 2001 From: can1357 Date: Thu, 2 Jul 2026 02:40:07 +0200 Subject: [PATCH] fix(coding-agent): introduced a shutdown coordinator to manage background tasks - Introduced `RpcShutdownCoordinator` to track background tasks and manage deferred shutdowns safely. - Guaranteed all background bash task response frames are fully written before the process exits. - Re-checked shutdown requests automatically as each tracked background task settles. - Latched the shutdown sequence to prevent concurrent execution from duplicate triggers. - Updated `mock-rpc-agent` to consume stdin via an async iterator to match standard behavior. - Added comprehensive unit tests in `rpc-input-frame.test.ts` covering background task coordination. --- .../coding-agent/src/modes/rpc/rpc-mode.ts | 111 +++++++++--- .../test/fixtures/mock-rpc-agent.ts | 12 +- .../coding-agent/test/rpc-input-frame.test.ts | 160 ++++++++++++++++++ 3 files changed, 248 insertions(+), 35 deletions(-) diff --git a/packages/coding-agent/src/modes/rpc/rpc-mode.ts b/packages/coding-agent/src/modes/rpc/rpc-mode.ts index f56f3a408..5c4cabd89 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-mode.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-mode.ts @@ -237,7 +237,7 @@ export interface RpcInputFrameDeps { * `type === "extension_ui_response"` and a string `id`. Payload variants (value, * confirmed, cancelled) are validated at the read site. */ -export function isRpcExtensionUIResponse(value: unknown): value is RpcExtensionUIResponse { +function isRpcExtensionUIResponse(value: unknown): value is RpcExtensionUIResponse { if (!isRecord(value)) return false; return value.type === "extension_ui_response" && typeof value.id === "string"; } @@ -301,7 +301,6 @@ export function dispatchRpcInputFrame(parsed: unknown, deps: RpcInputFrameDeps): } })(); deps.trackBackgroundTask?.(task); - void task; return undefined; } @@ -310,6 +309,65 @@ export function dispatchRpcInputFrame(parsed: unknown, deps: RpcInputFrameDeps): })(); } +/** + * Coordinates deferred shutdown with in-flight background input tasks. + * + * `pi.shutdown()` from an extension only *requests* shutdown; the process must + * not exit while a background-dispatched command (`bash`, see + * {@link dispatchRpcInputFrame}) still owes the client a response frame. The + * coordinator tracks those tasks, re-checks the shutdown request whenever one + * settles (covering a shutdown requested mid-bash with no follow-up client + * frame), and drains every tracked task before invoking `performShutdown`. + * The shutdown sequence is latched so concurrent triggers (input loop and + * settling tasks) run it exactly once. + */ +export class RpcShutdownCoordinator { + #tasks = new Set>(); + #shutdown: Promise | undefined; + readonly #isShutdownRequested: () => boolean; + readonly #performShutdown: () => Promise; + + constructor(options: { isShutdownRequested: () => boolean; performShutdown: () => Promise }) { + this.#isShutdownRequested = options.isShutdownRequested; + this.#performShutdown = options.performShutdown; + } + + /** + * Track a background input task. When it settles it is untracked and the + * shutdown request is re-checked, so a deferred shutdown fires even when + * no further client frames arrive. + */ + track(task: Promise): void { + this.#tasks.add(task); + void task.finally(() => { + this.#tasks.delete(task); + // Fire-and-forget: performShutdown ends the process. Rejections are + // not expected — hook errors are caught inside extensionRunner.emit, + // and background tasks catch their own dispatch errors. + void this.checkShutdownRequested(); + }); + } + + /** Await every tracked task, including tasks tracked while draining. */ + async drain(): Promise { + while (this.#tasks.size > 0) { + await Promise.allSettled(Array.from(this.#tasks)); + } + } + + /** + * If shutdown was requested, drain background tasks (so every owed + * response frame is written) before running the shutdown sequence. + */ + checkShutdownRequested(): Promise { + if (!this.#shutdown) { + if (!this.#isShutdownRequested()) return Promise.resolve(); + this.#shutdown = this.drain().then(() => this.#performShutdown()); + } + return this.#shutdown; + } +} + export type RpcSubagentResetRegistry = Pick; export async function handleRpcSessionChange( @@ -1193,31 +1251,25 @@ export async function runRpcMode( } }; - /** - * Check if shutdown was requested and perform shutdown if so. - * Called after handling each command when waiting for the next command. - */ - async function checkShutdownRequested(): Promise { - if (!shutdownState.requested) return; - - if (session.extensionRunner?.hasHandlers("session_shutdown")) { - await session.extensionRunner.emit({ type: "session_shutdown" }); - } - - process.exit(0); - } - - const backgroundInputTasks = new Set>(); - const trackBackgroundTask = (task: Promise) => { - backgroundInputTasks.add(task); - void task.finally(() => backgroundInputTasks.delete(task)); - }; + // Deferred shutdown (pi.shutdown() from an extension) must not kill the + // process while a background-dispatched bash still owes the client its + // response frame. The coordinator drains tracked tasks before exiting and + // re-checks the request as each task settles. + const shutdownCoordinator = new RpcShutdownCoordinator({ + isShutdownRequested: () => shutdownState.requested, + performShutdown: async () => { + if (session.extensionRunner?.hasHandlers("session_shutdown")) { + await session.extensionRunner.emit({ type: "session_shutdown" }); + } + process.exit(0); + }, + }); const dispatchFrameDeps: RpcInputFrameDeps = { handleCommand, output, errorResponse: error, - trackBackgroundTask, + trackBackgroundTask: task => shutdownCoordinator.track(task), pendingExtensionRequests, onHostToolResult: frame => hostToolBridge.handleResult(frame), onHostToolUpdate: frame => hostToolBridge.handleUpdate(frame), @@ -1233,9 +1285,12 @@ export async function runRpcMode( if (awaited) { await awaited; // Check for deferred shutdown request (idle between commands). - // Skipped when a bash command was just dispatched in the - // background — shutdown checks run after subsequent frames. - await checkShutdownRequested(); + // Background-dispatched bash frames skip this check so a later + // abort_bash can still be read; the coordinator re-checks when + // each tracked task settles, so a shutdown requested mid-bash + // fires once the response frame is written even if no further + // client frames arrive. + await shutdownCoordinator.checkShutdownRequested(); } } catch (e: unknown) { const message = e instanceof Error ? e.message : String(e); @@ -1243,9 +1298,9 @@ export async function runRpcMode( } } - if (backgroundInputTasks.size > 0) { - await Promise.allSettled(Array.from(backgroundInputTasks)); - } + // Background bash tasks may still owe response frames; drain them before + // tearing down (stdin EOF ends the frame stream, not in-flight work). + await shutdownCoordinator.drain(); // stdin closed — RPC client is gone, exit cleanly hostToolBridge.rejectAllPending("RPC client disconnected before host tool execution completed"); diff --git a/packages/coding-agent/test/fixtures/mock-rpc-agent.ts b/packages/coding-agent/test/fixtures/mock-rpc-agent.ts index 5e19139fa..a1c3f30b3 100755 --- a/packages/coding-agent/test/fixtures/mock-rpc-agent.ts +++ b/packages/coding-agent/test/fixtures/mock-rpc-agent.ts @@ -7,13 +7,11 @@ * Used by rpc-client lifecycle tests that need to exercise start/stop/start * without booting the full agent runtime (which requires provider credentials). */ -import * as readline from "node:readline"; - process.stdout.write(`${JSON.stringify({ type: "ready" })}\n`); -const rl = readline.createInterface({ input: process.stdin }); -rl.on("line", raw => { - if (!raw) return; +// Bun's `console` is an AsyncIterable over stdin lines. +for await (const raw of console) { + if (!raw) continue; try { const frame = JSON.parse(raw) as Record; if (frame && typeof frame === "object" && typeof frame.type === "string") { @@ -31,5 +29,5 @@ rl.on("line", raw => { } catch { // ignore parse errors — the test harness sends well-formed frames. } -}); -rl.on("close", () => process.exit(0)); +} +process.exit(0); diff --git a/packages/coding-agent/test/rpc-input-frame.test.ts b/packages/coding-agent/test/rpc-input-frame.test.ts index 18b98e4e1..0b582fcb8 100644 --- a/packages/coding-agent/test/rpc-input-frame.test.ts +++ b/packages/coding-agent/test/rpc-input-frame.test.ts @@ -3,6 +3,7 @@ import { dispatchRpcInputFrame, type PendingExtensionRequest, type RpcInputFrameDeps, + RpcShutdownCoordinator, } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-mode"; import type { RpcCommand, RpcResponse } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-types"; @@ -200,3 +201,162 @@ describe("dispatchRpcInputFrame", () => { expect(outputs).toEqual([bashResponse]); }); }); + +describe("RpcShutdownCoordinator", () => { + /** performShutdown spy that records call count and outputs.length at the moment it ran. */ + const makeShutdownRecorder = (outputs: OutputFrame[]) => { + const state = { calls: 0, outputsAtShutdown: -1 }; + const performShutdown = async () => { + state.calls++; + state.outputsAtShutdown = outputs.length; + }; + return { state, performShutdown }; + }; + + /** + * Full production-shaped harness: a background-dispatched bash frame whose + * handler blocks on a gate, tracked by the coordinator exactly as + * `runRpcMode` wires it (`trackBackgroundTask: task => coordinator.track(task)`). + */ + const makeBashHarness = () => { + const gate = Promise.withResolvers(); + const { deps, outputs } = makeDeps(async command => { + if (command.type === "bash") return await gate.promise; + throw new Error(`unexpected: ${command.type}`); + }); + const shutdown = { requested: false }; + const recorder = makeShutdownRecorder(outputs); + const coordinator = new RpcShutdownCoordinator({ + isShutdownRequested: () => shutdown.requested, + performShutdown: recorder.performShutdown, + }); + deps.trackBackgroundTask = task => coordinator.track(task); + return { gate, deps, outputs, shutdown, recorder, coordinator }; + }; + + test("deferred shutdown drains an in-flight background bash before performShutdown", async () => { + const { gate, deps, outputs, shutdown, recorder, coordinator } = makeBashHarness(); + + const awaited = dispatchRpcInputFrame({ id: "s1", type: "bash", command: "sleep 9999" }, deps); + expect(awaited).toBeUndefined(); + + // Extension calls pi.shutdown() while bash is in flight; the input loop + // re-checks after its next serially-awaited frame. + shutdown.requested = true; + const check = coordinator.checkShutdownRequested(); + + // The check must stay pending while the background bash still owes its + // response frame. Race it against a flushed sentinel: if the check could + // resolve, its microtask would win before the setImmediate tick. + const winner = await Promise.race([check.then(() => "shutdown"), flushMicrotasks().then(() => "pending")]); + expect(winner).toBe("pending"); + expect(recorder.state.calls).toBe(0); + expect(outputs).toHaveLength(0); + + gate.resolve(cancelledBashResponse("s1")); + await check; + + expect(outputs).toEqual([cancelledBashResponse("s1")]); + expect(recorder.state.calls).toBe(1); + // The bash response frame was already written when performShutdown ran. + expect(recorder.state.outputsAtShutdown).toBe(1); + }); + + test("settle hook fires the deferred shutdown when no further client frames arrive", async () => { + const { gate, deps, outputs, shutdown, recorder } = makeBashHarness(); + + const awaited = dispatchRpcInputFrame({ id: "s2", type: "bash", command: "sleep 9999" }, deps); + expect(awaited).toBeUndefined(); + + // Shutdown requested mid-bash; the stdin loop is parked with no frames, + // so the test never calls checkShutdownRequested() — only track()'s + // settle hook can trigger it. + shutdown.requested = true; + await flushMicrotasks(); + expect(recorder.state.calls).toBe(0); + + gate.resolve(cancelledBashResponse("s2")); + await flushMicrotasks(); + await flushMicrotasks(); + + expect(recorder.state.calls).toBe(1); + expect(outputs).toEqual([cancelledBashResponse("s2")]); + expect(recorder.state.outputsAtShutdown).toBe(1); + }); + + test("concurrent triggers are latched: performShutdown runs exactly once", async () => { + const outputs: OutputFrame[] = []; + const recorder = makeShutdownRecorder(outputs); + const coordinator = new RpcShutdownCoordinator({ + isShutdownRequested: () => true, + performShutdown: recorder.performShutdown, + }); + + const gateA = Promise.withResolvers(); + const gateB = Promise.withResolvers(); + coordinator.track(gateA.promise); + coordinator.track(gateB.promise); + + // Explicit trigger (input loop) races the settle hooks of both tasks. + const check = coordinator.checkShutdownRequested(); + gateA.resolve(); + gateB.resolve(); + await check; + await flushMicrotasks(); + await flushMicrotasks(); + + expect(recorder.state.calls).toBe(1); + // A later re-check reuses the latched sequence instead of re-running it. + await coordinator.checkShutdownRequested(); + expect(recorder.state.calls).toBe(1); + }); + + test("no-op when shutdown was not requested", async () => { + const outputs: OutputFrame[] = []; + const recorder = makeShutdownRecorder(outputs); + const coordinator = new RpcShutdownCoordinator({ + isShutdownRequested: () => false, + performShutdown: recorder.performShutdown, + }); + + await coordinator.checkShutdownRequested(); + expect(recorder.state.calls).toBe(0); + + // A tracked task settling with the flag false never triggers shutdown. + const gate = Promise.withResolvers(); + coordinator.track(gate.promise); + gate.resolve(); + await flushMicrotasks(); + await flushMicrotasks(); + expect(recorder.state.calls).toBe(0); + }); + + test("drain() waits for tasks tracked while draining", async () => { + const coordinator = new RpcShutdownCoordinator({ + isShutdownRequested: () => false, + performShutdown: async () => {}, + }); + + const gateA = Promise.withResolvers(); + const gateB = Promise.withResolvers(); + coordinator.track(gateA.promise); + // When A settles, a new task B enters the set mid-drain. + void gateA.promise.then(() => { + coordinator.track(gateB.promise); + }); + + let drained = false; + const drain = coordinator.drain().then(() => { + drained = true; + }); + + gateA.resolve(); + await flushMicrotasks(); + // A settled and B was tracked mid-drain; drain must keep waiting on B. + expect(drained).toBe(false); + + gateB.resolve(); + await drain; + expect(drained).toBe(true); + }); +});