fix(utils): share active postmortem cleanup
This commit is contained in:
@@ -27,6 +27,7 @@ const callbackList: ((reason: Reason) => Promise<void> | 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<void> | undefined;
|
||||
let stdioDisconnectRegistrations = 0;
|
||||
|
||||
/**
|
||||
@@ -41,7 +42,7 @@ function runCleanup(reason: Reason): Promise<void> {
|
||||
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<void> {
|
||||
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.
|
||||
*
|
||||
|
||||
@@ -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<void>();
|
||||
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<void>(() => {});
|
||||
}
|
||||
|
||||
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(() => {});
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user