mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-27 21:07:01 -06:00
fix(process): preserve output across long event-loop stalls (#125381)
Refine the existing repair in #125381 with one cancellable release owner and continued-writer coverage. Fixes #125380 and #130504. Reported by Vyctor H. Brzezowski (@vyctorbrzezowski). Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
@@ -18,33 +18,35 @@ describe.skipIf(process.platform === "win32")("releaseChildProcessOutputAfterExi
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("drains active descendant output after the parent exits", async () => {
|
||||
const command = 'printf "HEAD\\n"; ( sleep 0.05; printf "TAIL\\n" ) &';
|
||||
child = execa("/bin/sh", ["-c", command], {
|
||||
buffer: false,
|
||||
detached: true,
|
||||
reject: false,
|
||||
stdio: ["ignore", "pipe", "pipe"],
|
||||
});
|
||||
const releaseOutput = releaseChildProcessOutputAfterExit(child.nodeChildProcess);
|
||||
let output = "";
|
||||
child.stdout?.on("data", (chunk: Buffer) => {
|
||||
output += chunk.toString();
|
||||
});
|
||||
|
||||
// Simulate a contended worker after the direct child exits. The descendant
|
||||
// writes while JS is parked, so its pipe data and the idle timer are both
|
||||
// ready when the event loop resumes.
|
||||
await new Promise<void>((resolve) => {
|
||||
child?.nodeChildProcess.once("exit", () => {
|
||||
Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 250);
|
||||
resolve();
|
||||
it.each([250, 1_250])(
|
||||
"drains descendant output across a %ims event-loop stall",
|
||||
async (stallMs) => {
|
||||
const command =
|
||||
'printf "HEAD\\n"; printf "HEAD\\n" >&2; ( sleep 0.05; printf "TAIL\\n"; printf "TAIL\\n" >&2 ) &';
|
||||
child = execa("/bin/sh", ["-c", command], {
|
||||
detached: true,
|
||||
reject: false,
|
||||
stdio: ["ignore", "pipe", "pipe"],
|
||||
stripFinalNewline: false,
|
||||
});
|
||||
});
|
||||
await child.finally(releaseOutput);
|
||||
expect(output).toContain("HEAD");
|
||||
expect(output).toContain("TAIL");
|
||||
});
|
||||
const releaseOutput = releaseChildProcessOutputAfterExit(child.nodeChildProcess);
|
||||
|
||||
// Simulate a contended worker after the direct child exits. The descendant
|
||||
// writes while JS is parked past the idle or hard deadline, so buffered
|
||||
// pipe data and release timers are ready when the event loop resumes.
|
||||
await new Promise<void>((resolve) => {
|
||||
child?.nodeChildProcess.once("exit", () => {
|
||||
Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, stallMs);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
expect(await child.finally(releaseOutput)).toMatchObject({
|
||||
exitCode: 0,
|
||||
stdout: "HEAD\nTAIL\n",
|
||||
stderr: "HEAD\nTAIL\n",
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
it("releases a quiet inherited pipe after the idle grace", async () => {
|
||||
child = execa("/bin/sh", ["-c", 'printf "DONE\\n"; ( sleep 30 ) &'], {
|
||||
@@ -77,13 +79,38 @@ describe.skipIf(process.platform === "win32")("releaseChildProcessOutputAfterExi
|
||||
fakeChild.emit("exit", 0);
|
||||
const writer = setInterval(() => stdout.write("TICK\n"), 30);
|
||||
|
||||
await vi.advanceTimersByTimeAsync(999);
|
||||
expect(stdout.destroyed).toBe(false);
|
||||
await vi.advanceTimersByTimeAsync(1);
|
||||
expect(stdout.destroyed).toBe(true);
|
||||
expect(stderr.destroyed).toBe(true);
|
||||
try {
|
||||
await vi.advanceTimersByTimeAsync(999);
|
||||
expect(stdout.destroyed).toBe(false);
|
||||
await vi.advanceTimersByTimeAsync(1);
|
||||
expect(stdout.destroyed).toBe(false);
|
||||
stdout.write("AFTER DEADLINE\n");
|
||||
await vi.advanceTimersByTimeAsync(1);
|
||||
expect(stdout.destroyed).toBe(true);
|
||||
expect(stderr.destroyed).toBe(true);
|
||||
} finally {
|
||||
clearInterval(writer);
|
||||
cleanup();
|
||||
}
|
||||
});
|
||||
|
||||
clearInterval(writer);
|
||||
it.each(["idle", "hard"])("cancels pending %s release when cleaned up", async (deadline) => {
|
||||
vi.useFakeTimers();
|
||||
const stdout = new PassThrough();
|
||||
const stderr = new PassThrough();
|
||||
const fakeChild = Object.assign(new EventEmitter(), {
|
||||
stdout,
|
||||
stderr,
|
||||
}) as unknown as ChildProcess;
|
||||
const cleanup = releaseChildProcessOutputAfterExit(fakeChild);
|
||||
fakeChild.emit("exit", 0);
|
||||
const writer = deadline === "hard" ? setInterval(() => stdout.write("TICK\n"), 30) : undefined;
|
||||
|
||||
await vi.advanceTimersByTimeAsync(deadline === "hard" ? 1_000 : 100);
|
||||
cleanup();
|
||||
clearInterval(writer);
|
||||
await vi.runAllTimersAsync();
|
||||
expect(stdout.destroyed).toBe(false);
|
||||
expect(stderr.destroyed).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -12,27 +12,14 @@ const EXIT_STDIO_MAX_DRAIN_MS = 1_000;
|
||||
* short output tails. The returned cleanup must run after awaiting the child.
|
||||
*/
|
||||
export function releaseChildProcessOutputAfterExit(child: ChildProcess): () => void {
|
||||
let exited = false;
|
||||
let idleTimer: NodeJS.Timeout | undefined;
|
||||
let idleReleaseImmediate: NodeJS.Immediate | undefined;
|
||||
let releaseImmediate: NodeJS.Immediate | undefined;
|
||||
let deadlineTimer: NodeJS.Timeout | undefined;
|
||||
|
||||
const clearTimers = () => {
|
||||
if (idleTimer) {
|
||||
clearTimeout(idleTimer);
|
||||
idleTimer = undefined;
|
||||
}
|
||||
if (idleReleaseImmediate) {
|
||||
clearImmediate(idleReleaseImmediate);
|
||||
idleReleaseImmediate = undefined;
|
||||
}
|
||||
if (deadlineTimer) {
|
||||
clearTimeout(deadlineTimer);
|
||||
deadlineTimer = undefined;
|
||||
}
|
||||
};
|
||||
const cleanup = () => {
|
||||
clearTimers();
|
||||
clearTimeout(idleTimer);
|
||||
clearImmediate(releaseImmediate);
|
||||
clearTimeout(deadlineTimer);
|
||||
child.removeListener("exit", onExit);
|
||||
child.stdout?.removeListener("data", onData);
|
||||
child.stderr?.removeListener("data", onData);
|
||||
@@ -42,36 +29,32 @@ export function releaseChildProcessOutputAfterExit(child: ChildProcess): () => v
|
||||
child.stdout?.destroy();
|
||||
child.stderr?.destroy();
|
||||
};
|
||||
const scheduleRelease = () => {
|
||||
// Either timer may run before already-buffered pipe data on a loaded loop.
|
||||
// Share one cancellable release after poll has had a turn to drain it.
|
||||
releaseImmediate ??= setImmediate(release);
|
||||
releaseImmediate.unref();
|
||||
};
|
||||
const armIdleTimer = () => {
|
||||
if (idleTimer) {
|
||||
clearTimeout(idleTimer);
|
||||
}
|
||||
if (idleReleaseImmediate) {
|
||||
clearImmediate(idleReleaseImmediate);
|
||||
idleReleaseImmediate = undefined;
|
||||
}
|
||||
idleTimer = setTimeout(() => {
|
||||
idleTimer = undefined;
|
||||
// A loaded event loop can observe the idle timer before already-buffered
|
||||
// pipe data. Give the poll phase one turn so that data can rearm the grace.
|
||||
idleReleaseImmediate = setImmediate(() => {
|
||||
idleReleaseImmediate = undefined;
|
||||
release();
|
||||
});
|
||||
idleReleaseImmediate.unref();
|
||||
}, EXIT_STDIO_GRACE_MS);
|
||||
clearTimeout(idleTimer);
|
||||
clearImmediate(releaseImmediate);
|
||||
releaseImmediate = undefined;
|
||||
idleTimer = setTimeout(scheduleRelease, EXIT_STDIO_GRACE_MS);
|
||||
idleTimer.unref();
|
||||
};
|
||||
const onData = () => {
|
||||
if (exited) {
|
||||
if (deadlineTimer) {
|
||||
armIdleTimer();
|
||||
}
|
||||
};
|
||||
const onExit = () => {
|
||||
exited = true;
|
||||
armIdleTimer();
|
||||
deadlineTimer = setTimeout(release, EXIT_STDIO_MAX_DRAIN_MS);
|
||||
deadlineTimer = setTimeout(() => {
|
||||
// Post-deadline data must not cancel release and extend the hard bound.
|
||||
deadlineTimer = undefined;
|
||||
scheduleRelease();
|
||||
}, EXIT_STDIO_MAX_DRAIN_MS);
|
||||
deadlineTimer.unref();
|
||||
armIdleTimer();
|
||||
};
|
||||
|
||||
child.stdout?.on("data", onData);
|
||||
|
||||
Reference in New Issue
Block a user