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
19 changes: 17 additions & 2 deletions src/jsc/bindings/webcore/WorkerMessagingProxy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -321,8 +321,7 @@ static bool drainInbox(WorkerMessagingProxy::MessageInbox& inbox, Zig::GlobalObj
remaining -= batch.size();

while (!batch.isEmpty()) {
// The receiving VM is being stopped: nothing more is delivered (the
// rest is dropped with the proxy).
// The receiving VM is being stopped: nothing more is delivered, and what is left is dropped unread.
if (context.isJSExecutionForbidden())
return false;
auto message = batch.takeFirst();
Expand Down Expand Up @@ -543,6 +542,19 @@ void WorkerMessagingProxy::releaseWorkerThread()
deref();
}

void WorkerMessagingProxy::dropUndeliveredWorkerMessages()
{
// Ports a worker that never started did not take. Closing its parentPort end also closes the ports queued on it.
auto droppedDataPorts = std::exchange(m_options.dataMessagePorts, {});
Comment thread
robobun marked this conversation as resolved.
Deque<MessageWithMessagePorts> droppedMessages;
{
Locker locker { m_toWorker.lock };
droppedMessages = std::exchange(m_toWorker.queue, {});
m_toWorker.drainScheduled = false;
}
// Destroyed here, outside the lock: ~TransferredMessagePort closes its pipe side and notifies the peer.
}

void WorkerMessagingProxy::workerGlobalScopeDestroyedInternal(int32_t exitCode, bool stoppedByParent)
{
ASSERT(m_scriptExecutionContext && m_scriptExecutionContext->isContextThread());
Expand All @@ -561,6 +573,7 @@ void WorkerMessagingProxy::workerGlobalScopeDestroyedInternal(int32_t exitCode,
m_state.store(State::Closing);
m_pendingTasks.clear();
}
dropUndeliveredWorkerMessages();
rejectAllCrossVMRequests();

// Everything the worker posted before it exited is delivered before 'close' (Node: before
Expand Down Expand Up @@ -595,6 +608,8 @@ void WorkerMessagingProxy::parentContextWillDestroy()
m_pendingCrossVMRequests.clear();
}
releaseWorkerThread();
// After the join: a live worker thread can still be taking its ports (createNodeWorkerThreadsBinding).
dropUndeliveredWorkerMessages();
m_scriptExecutionContext = nullptr;
}

Expand Down
1 change: 1 addition & 0 deletions src/jsc/bindings/webcore/WorkerMessagingProxy.h
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,7 @@ class WorkerMessagingProxy final : public ThreadSafeRefCounted<WorkerMessagingPr

void workerGlobalScopeDestroyedInternal(int32_t exitCode, bool stoppedByParent);
void releaseWorkerThread();
void dropUndeliveredWorkerMessages();
void drainMessagesToWorkerObject(ScriptExecutionContext&, DrainBudget);
void rejectAllCrossVMRequests();
void postMessageErrorToWorkerObject(String&& message);
Expand Down
50 changes: 50 additions & 0 deletions test/js/node/worker_threads/worker-transferred-port-close.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
// Also runs in Node.js (`node --test <file>`), so it uses node:test and imports only Node modules.
import assert from "node:assert";
import { once } from "node:events";
import { join } from "node:path";
import { describe, test } from "node:test";
import { MessageChannel, Worker } from "node:worker_threads";

describe("a MessagePort transferred to a worker that never runs is closed", () => {
// Stays referenced, as a Worker in a pool does: a collected Worker drops its ports too.
let worker: Worker;

for (const route of ["workerData", "postMessage"]) {
// Node reports 'exit' and the port's 'close' in either order, so that order is not asserted.
test(`the entry does not resolve, port transferred through ${route}`, async () => {
const { port1, port2 } = new MessageChannel();
worker = new Worker(
join(import.meta.dirname, "worker-transferred-port-close-missing-entry.cjs"),
route === "workerData" ? { workerData: { port: port2 }, transferList: [port2] } : undefined,
);
const events: string[] = [];
worker.on("online", () => events.push("online"));
worker.on("error", (error: NodeJS.ErrnoException) => events.push(`error:${error.code}`));
const portClosed = once(port1, "close").then(() => events.push("port-close"));
const exited = new Promise(resolve => worker.on("exit", code => resolve(events.push(`exit:${code}`))));
if (route === "postMessage") worker.postMessage({ port: port2 }, [port2]);

await Promise.all([exited, portClosed]);
assert.deepStrictEqual(
{ first: events.slice(0, 2), rest: events.slice(2).sort() },
{ first: ["online", "error:MODULE_NOT_FOUND"], rest: ["exit:1", "port-close"] },
);
});

// terminate() in the same tick as the constructor stops the thread before it takes its ports.
// The thread takes them before any user code runs, so nothing can hold it back, and a thread
// that wins that race closes them as it exits. The exit code depends on that timing too.
test(`terminate() stops the worker before it starts, port transferred through ${route}`, async () => {
const { port1, port2 } = new MessageChannel();
worker = new Worker(
"setInterval(() => {}, 1000)",
route === "workerData" ? { eval: true, workerData: { port: port2 }, transferList: [port2] } : { eval: true },
);
const portClosed = once(port1, "close").then(() => "port-close");
if (route === "postMessage") worker.postMessage({ port: port2 }, [port2]);

await worker.terminate();
assert.strictEqual(await portClosed, "port-close");
});
}
});
40 changes: 39 additions & 1 deletion test/js/web/workers/message-port-pipe.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import { describe, expect, test } from "bun:test";
import { bunEnv, bunExe, isASAN, isDebug } from "harness";
import { bunEnv, bunExe, isASAN, isDebug, tempDir } from "harness";
import { once } from "node:events";
import { join } from "node:path";
import { receiveMessageOnPort } from "node:worker_threads";

// Exercises the MessagePortPipe layer that backs MessagePort/MessageChannel:
Expand Down Expand Up @@ -478,4 +480,40 @@ describe("Worker postMessage inbox", () => {
expect(stdout.trim()).toBe("OK");
expect(exitCode).toBe(0);
});

// A message the worker never took from its inbox is dropped when the worker is gone. A port in
// that message is closed with it, so the port's peer hears 'close'.
describe("a MessagePort in a message the worker never reads is closed", () => {
// Stays referenced, as a Worker in a pool does: a collected Worker drops its inbox too.
let worker: Worker;

test("the entry point does not resolve", async () => {
using dir = tempDir("web-worker-missing-entry-port", {});
const { port1, port2 } = new MessageChannel();
worker = new Worker(join(String(dir), "missing.js"));
const events: string[] = [];
worker.addEventListener("error", () => events.push("error"));
worker.addEventListener("close", e => events.push(`close:${e.code}`));
const portClosed = once(port1, "close").then(() => events.push("port-close"));
worker.postMessage({ port: port2 }, [port2]);

await portClosed;
expect(events).toEqual(["error", "close:1", "port-close"]);
});

// The entry never returns, so the worker never reads its inbox, whether terminate() lands
// before the thread starts or while the entry runs.
test("terminate() stops a worker whose entry is still running", async () => {
const { port1, port2 } = new MessageChannel();
worker = new Worker("data:text/javascript,for(;;){}");
const events: string[] = [];
worker.addEventListener("close", () => events.push("close"));
const portClosed = once(port1, "close").then(() => events.push("port-close"));
worker.postMessage({ port: port2 }, [port2]);
worker.terminate();

await portClosed;
expect(events).toEqual(["close", "port-close"]);
});
});
});
Loading