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.
This commit is contained in:
@@ -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<Promise<void>>();
|
||||
#shutdown: Promise<void> | undefined;
|
||||
readonly #isShutdownRequested: () => boolean;
|
||||
readonly #performShutdown: () => Promise<void>;
|
||||
|
||||
constructor(options: { isShutdownRequested: () => boolean; performShutdown: () => Promise<void> }) {
|
||||
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>): 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<void> {
|
||||
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<void> {
|
||||
if (!this.#shutdown) {
|
||||
if (!this.#isShutdownRequested()) return Promise.resolve();
|
||||
this.#shutdown = this.drain().then(() => this.#performShutdown());
|
||||
}
|
||||
return this.#shutdown;
|
||||
}
|
||||
}
|
||||
|
||||
export type RpcSubagentResetRegistry = Pick<RpcSubagentRegistry, "clear">;
|
||||
|
||||
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<void> {
|
||||
if (!shutdownState.requested) return;
|
||||
|
||||
if (session.extensionRunner?.hasHandlers("session_shutdown")) {
|
||||
await session.extensionRunner.emit({ type: "session_shutdown" });
|
||||
}
|
||||
|
||||
process.exit(0);
|
||||
}
|
||||
|
||||
const backgroundInputTasks = new Set<Promise<void>>();
|
||||
const trackBackgroundTask = (task: Promise<void>) => {
|
||||
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");
|
||||
|
||||
+5
-7
@@ -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<string, unknown>;
|
||||
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);
|
||||
|
||||
@@ -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<RpcResponse>();
|
||||
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<void>();
|
||||
const gateB = Promise.withResolvers<void>();
|
||||
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<void>();
|
||||
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<void>();
|
||||
const gateB = Promise.withResolvers<void>();
|
||||
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);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user