Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 8 additions & 4 deletions src/jsc/web_worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,7 @@ struct WorkerVmInit {

enum EntryOutcome {
Continue,
/// The entry module rejected and no handler took it: the worker exits.
/// The entry module rejected and no handler took it, or the handler that took it stopped the worker.
Stop,
}

Expand Down Expand Up @@ -884,6 +884,11 @@ impl WebWorker {
// so buffered postMessageToThread deliveries drain and the sender's
// Atomics.waitAsync settles. WebWorker__entrySettled re-calls it as a no-op.
WebWorker__entrySettled(vm.global());
// The hook ran script: a stop requested there ends the start sequence before its GC and first tick.
if self.has_requested_terminate() {
self.flush_logs(vm);
return self.shutdown();
}

// The entry's evaluation outcome is checked once now and then after every
// loop turn: a rejection (immediate, or a top-level await rejecting
Expand All @@ -908,16 +913,15 @@ impl WebWorker {
(*promise).result(vm.jsc_vm()),
is_rejection,
);
if handled {
if handled && !self.has_requested_terminate() {
EntryOutcome::Continue
} else {
EntryOutcome::Stop
}
}
};
if let EntryOutcome::Stop = observe_entry(vm) {
// exit_code is already 1 from uncaught_exception; re-setting it here
// would clobber a process.on('exit') change to process.exitCode.
// exit_code is already set (1, or the handler's process.exit() code), and an 'exit' listener may have changed it.
return self.shutdown();
}
// A still-pending entry promise is an unsettled top-level await: as in
Expand Down
92 changes: 92 additions & 0 deletions test/js/node/worker_threads/worker_threads.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -717,6 +717,98 @@ describe("error event", () => {
});
});

// Runs a worker in a child process (under `bun test` a worker's uncaught error goes to the test runner) and
// returns the child's stdout lines, sorted.
async function linesFromWorkerInChild(source: string, parentCode: string) {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`const { Worker, postMessageToThread } = require("node:worker_threads");
const worker = new Worker(${JSON.stringify(source)}, { eval: true });
worker.on("error", err => console.log("error", err.message));
worker.on("exit", code => console.log("exit", code));
${parentCode}`,
],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
return { lines: stdout.trim().split("\n").sort(), stderr, exitCode, signalCode: proc.signalCode };
}

// The start sequence calls these handlers with no script frame beneath them, so a stop that is requested inside one
// has nothing to unwind: the start sequence has to stand down. If it goes on to its GC, debug and ASAN builds abort
// at `vm.hasTerminationRequest()`. A release build prints the same lines either way, except in the last row.
const spinUntilTerminated = `require("node:worker_threads").parentPort.postMessage("in the handler"); for (;;) {}`;
const terminateOnMessage = `worker.on("message", () => worker.terminate());`;
test.concurrent.each([
[
"process.exit() in an uncaughtException capture callback, CommonJS entry point throws",
`require("node:fs"); process.setUncaughtExceptionCaptureCallback(() => process.exit(42)); throw new Error("boom");`,
"",
["exit 42"],
],
[
"process.exit() in an uncaughtException capture callback, ES module entry point rejects",
`process.setUncaughtExceptionCaptureCallback(() => process.exit(42)); await 0; throw new Error("boom");`,
"",
["exit 42"],
],
[
"terminate() landing in an uncaughtException capture callback",
`process.setUncaughtExceptionCaptureCallback(() => { ${spinUntilTerminated} }); throw new Error("boom");`,
terminateOnMessage,
["exit 1"],
],
// The abort needs an Error whose stack nothing has read yet: the GC formats it.
[
"process.exit() in a 'workerMessage' listener, message buffered while the entry point loaded",
`globalThis.unreadStack = new Error("kept"); process.on("workerMessage", () => process.exit(7)); setInterval(() => {}, 1000);`,
`postMessageToThread(worker.threadId, "hello").catch(() => {});`,
["exit 7"],
],
[
"terminate() landing in a 'workerMessage' listener, message buffered while the entry point loaded",
`globalThis.unreadStack = new Error("kept"); process.on("workerMessage", () => { ${spinUntilTerminated} }); setInterval(() => {}, 1000);`,
`postMessageToThread(worker.threadId, "hello").catch(() => {}); ${terminateOnMessage}`,
Comment on lines +764 to +775

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 nit (optional): on a release build, deleting the new check at src/jsc/web_worker.rs:888 breaks no test, so a release-only CI lane cannot catch its regression. The two 'workerMessage' rows here print "exit 7" / "exit 1" with or without that check; only debug and ASAN builds abort, and the single row that differs on release (getHeapSnapshot) exercises the observe_entry clause, not this one. Fix: make each load-bearing clause fail a test on every build, e.g. add a 'workerMessage' row whose parent also holds a pending worker.getHeapSnapshot() and expects "snapshot rejected ERR_WORKER_NOT_RUNNING", since without the check the worker still reaches Running and the snapshot resolves.

Why this was flagged

Without the check at src/jsc/web_worker.rs:888 on a release build, the 'workerMessage' process.exit(7) row runs: WebWorker__entrySettled runs the listener, process.exit(7) sets exit_code and requests termination, observe_entry at web_worker.rs:923 sees a fulfilled entry promise and returns Continue, WebWorker__workerGlobalScopeStarted at web_worker.rs:939 moves the proxy to Running, run_gc at web_worker.rs:949 does not assert on release, the loop breaks at web_worker.rs:959 and shutdown reports exit 7. The row's expected lines are identical, so the test passes both ways on release; the getHeapSnapshot row only covers the observe_entry change at web_worker.rs:916. REVIEW.md asks that deleting each load-bearing clause of the fix break at least one test. A pending getHeapSnapshot() in the 'workerMessage' rows would distinguish: with the check the proxy never reaches Running and rejectAllCrossVMRequests (WorkerMessagingProxy.cpp:564) rejects it; without it the pending task runs after workerGlobalScopeStarted and resolves.

Verification: Without the new check at src/jsc/web_worker.rs:888-891, the two 'workerMessage' rows take observe_entry line 902 -> Continue, run_gc line 949 (the ASSERT in VMTraps.cpp is compiled out in release), then the loop at 957-961 breaks and shutdown() runs with exit_code already 7. Stdout is "exit 7"/"exit 1", exactly what the test at test/js/node/worker_threads/worker_threads.test.ts:775 expects.

["exit 1"],
],
[
"a getHeapSnapshot() that waits for the worker to run rejects, as in Node",
`process.setUncaughtExceptionCaptureCallback(() => process.exit(42)); throw new Error("boom");`,
`worker.getHeapSnapshot().then(() => console.log("snapshot resolved"), err => console.log("snapshot rejected", err.code));`,
["exit 42", "snapshot rejected ERR_WORKER_NOT_RUNNING"],
],
])("a worker stopped inside a handler that its start sequence calls: %s", async (_name, source, parentCode, lines) => {
expect(await linesFromWorkerInChild(source, parentCode)).toEqual({
lines,
stderr: "",
exitCode: 0,
signalCode: null,
});
});

// The error reporter leaves the termination pending, and that is what stops the nextTick drain above it.
test.concurrent(
"process.exit() in an uncaughtException capture callback stops the ticks queued behind the one that threw",
async () => {
const source = `const { parentPort } = require("node:worker_threads");
process.setUncaughtExceptionCaptureCallback(() => process.exit(42));
setTimeout(() => {
Comment thread
coderabbitai[bot] marked this conversation as resolved.
process.nextTick(() => { throw new Error("boom"); });
process.nextTick(() => parentPort.postMessage("a tick ran after process.exit()"));
}, 1);`;
expect(await linesFromWorkerInChild(source, `worker.on("message", message => console.log(message));`)).toEqual({
lines: ["exit 42"],
stderr: "",
exitCode: 0,
signalCode: null,
});
},
);

describe("getHeapSnapshot", () => {
test("throws if the wrong options are passed", () => {
const worker = new Worker("", { eval: true });
Expand Down
Loading