diff --git a/src/jsc/bindings/webcore/WorkerMessagingProxy.cpp b/src/jsc/bindings/webcore/WorkerMessagingProxy.cpp index fee6c0d4f0b6..b2f088e04d3e 100644 --- a/src/jsc/bindings/webcore/WorkerMessagingProxy.cpp +++ b/src/jsc/bindings/webcore/WorkerMessagingProxy.cpp @@ -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(); @@ -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, {}); + Deque 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()); @@ -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 @@ -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; } diff --git a/src/jsc/bindings/webcore/WorkerMessagingProxy.h b/src/jsc/bindings/webcore/WorkerMessagingProxy.h index 273a0f129f95..b990ab03413e 100644 --- a/src/jsc/bindings/webcore/WorkerMessagingProxy.h +++ b/src/jsc/bindings/webcore/WorkerMessagingProxy.h @@ -124,6 +124,7 @@ class WorkerMessagingProxy final : public ThreadSafeRefCounted`), 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"); + }); + } +}); diff --git a/test/js/web/workers/message-port-pipe.test.ts b/test/js/web/workers/message-port-pipe.test.ts index 9ff401768d06..a0cbf7fa89f6 100644 --- a/test/js/web/workers/message-port-pipe.test.ts +++ b/test/js/web/workers/message-port-pipe.test.ts @@ -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: @@ -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"]); + }); + }); });