diff --git a/packages/utils/src/postmortem.ts b/packages/utils/src/postmortem.ts index 33dabe3d3..c2784c1e4 100644 --- a/packages/utils/src/postmortem.ts +++ b/packages/utils/src/postmortem.ts @@ -27,6 +27,7 @@ const callbackList: ((reason: Reason) => Promise | void)[] = []; // Tracks cleanup run state (to prevent recursion/reentry issues) let cleanupStage: "idle" | "running" | "complete" = "idle"; const CLEANUP_DEADLINE_MS = 10_000; +let cleanupPromise: Promise | undefined; let stdioDisconnectRegistrations = 0; /** @@ -41,7 +42,7 @@ function runCleanup(reason: Reason): Promise { cleanupStage = "running"; break; case "running": - return Promise.resolve(); + return cleanupPromise ?? Promise.resolve(); case "complete": return Promise.resolve(); } @@ -68,9 +69,10 @@ function runCleanup(reason: Reason): Promise { deadline.resolve(); }, CLEANUP_DEADLINE_MS); deadlineTimer.unref(); - return Promise.race([cleanupSettled, deadline.promise]).finally(() => { + cleanupPromise = Promise.race([cleanupSettled, deadline.promise]).finally(() => { clearTimeout(deadlineTimer); }); + return cleanupPromise; } // Register signal and error event handlers to trigger cleanup before exit. @@ -92,6 +94,11 @@ export function classifyBrokenPipe(err: Error): BrokenPipeSource | undefined { return undefined; } +/** Whether an EPIPE came from an IPC `send()` to an optional worker. */ +export function isIpcSendEpipe(err: Error): boolean { + return classifyBrokenPipe(err) === "ipc-send"; +} + /** * Treat unhandled stdout EPIPE rejections as a graceful peer disconnect. * diff --git a/packages/utils/test/postmortem-epipe.test.ts b/packages/utils/test/postmortem-epipe.test.ts index 25fdcd8a1..c6df49a9f 100644 --- a/packages/utils/test/postmortem-epipe.test.ts +++ b/packages/utils/test/postmortem-epipe.test.ts @@ -2,6 +2,7 @@ import { describe, expect, it } from "bun:test"; import { postmortem } from "@oh-my-pi/pi-utils"; const childFlag = "--stdio-epipe-child"; +const raceChildFlag = "--stdio-epipe-race-child"; const childFlagIndex = process.argv.indexOf(childFlag); if (childFlagIndex >= 0) { const marker = process.argv[childFlagIndex + 1]; @@ -17,6 +18,33 @@ if (childFlagIndex >= 0) { const keepAlive = Promise.withResolvers(); await keepAlive.promise; } +else if (process.argv.includes(raceChildFlag)) { + const marker = process.argv[process.argv.indexOf(raceChildFlag) + 1]; + if (!marker) throw new Error("Missing cleanup marker path"); + let cleanupComplete = false; + let exitAttempted = false; + const exit = process.exit; + process.exit = ((code?: number) => { + if (!exitAttempted) { + exitAttempted = true; + void Bun.write(marker, cleanupComplete ? "after cleanup" : "before cleanup").then(() => exit(code)); + } + return undefined as never; + }) as typeof process.exit; + postmortem.registerStdioDisconnectHandling(); + postmortem.register("stdio-epipe-race-test", async () => { + process.stderr.write("cleanup started\n"); + void Promise.reject(Object.assign(new Error("broken pipe"), { code: "EPIPE", syscall: "write" })); + await new Response(Bun.stdin.stream()).text(); + cleanupComplete = true; + }); + let rejectionCount = 0; + process.on("unhandledRejection", () => { + if (++rejectionCount === 2) process.stderr.write("second rejection observed\n"); + }); + void Promise.reject(Object.assign(new Error("broken pipe"), { code: "EPIPE", syscall: "write" })); + await new Promise(() => {}); +} describe("postmortem broken-pipe handling", () => { function makeErr(props: { code?: string; syscall?: string; message?: string }): Error { @@ -28,6 +56,8 @@ describe("postmortem broken-pipe handling", () => { it("classifies worker IPC and stdio EPIPE errors", () => { expect(postmortem.classifyBrokenPipe(makeErr({ code: "EPIPE", syscall: "send" }))).toBe("ipc-send"); expect(postmortem.classifyBrokenPipe(makeErr({ code: "EPIPE", syscall: "write" }))).toBe("stdio-write"); + expect(postmortem.isIpcSendEpipe(makeErr({ code: "EPIPE", syscall: "send" }))).toBe(true); + expect(postmortem.isIpcSendEpipe(makeErr({ code: "EPIPE", syscall: "write" }))).toBe(false); }); it("does not classify unrelated errors as recoverable broken pipes", () => { @@ -66,4 +96,40 @@ describe("postmortem broken-pipe handling", () => { .catch(() => {}); } }); + + it("keeps waiting for active cleanup when another stdio EPIPE arrives", async () => { + const marker = `/tmp/omp-postmortem-stdio-race-${process.pid}-${Date.now()}`; + const child = Bun.spawn([process.execPath, import.meta.path, raceChildFlag, marker], { + stdin: "pipe", + stdout: "pipe", + stderr: "pipe", + }); + try { + const decoder = new TextDecoder(); + const stderrReader = child.stderr.getReader(); + let stderr = ""; + while (!stderr.includes("second rejection observed\n")) { + const chunk = await stderrReader.read(); + if (chunk.done) throw new Error("Child exited before observing the second rejection"); + stderr += decoder.decode(chunk.value); + } + stderrReader.releaseLock(); + expect(stderr).toContain("cleanup started\n"); + child.stdin.end(); + const [exitCode, stdout] = await Promise.all([child.exited, new Response(child.stdout).text()]); + expect(stdout).toBe(""); + expect(exitCode).toBe(0); + expect(await Bun.file(marker).text()).toBe("after cleanup"); + } finally { + try { + child.stdin.end(); + } catch { + // Already closed after completing teardown. + } + await child.exited; + await Bun.file(marker) + .delete() + .catch(() => {}); + } + }); });