diff --git a/src/jsc/bindings/webcore/MessagePort.cpp b/src/jsc/bindings/webcore/MessagePort.cpp index 1e1681c17d85..43c58b6def8e 100644 --- a/src/jsc/bindings/webcore/MessagePort.cpp +++ b/src/jsc/bindings/webcore/MessagePort.cpp @@ -50,6 +50,8 @@ Ref MessagePort::create(ScriptExecutionContext& context, RefsuspendIfNeeded(); + // So the peer's close() reaches this port (peerClosed()) even if it never gets a listener. + messagePort->m_pipe->registerCloseContext(side, context.identifier(), ThreadSafeWeakPtr { messagePort.get() }); return messagePort; } @@ -105,14 +107,11 @@ ExceptionOr MessagePort::postMessage(JSC::JSGlobalObject& state, JSC::JSVa } RETURN_IF_EXCEPTION(warnScope, {}); - if (!isEntangled()) - return {}; - Vector transferredPorts; + // Posting a port's own entangled peer targets the message at itself. + // (The source port itself was rejected before serialization above.) + bool targetsEntangledPeer = false; if (!ports.isEmpty()) { - // Posting a port's own entangled peer targets the message at itself. - // (The source port itself was rejected before serialization above.) - bool targetsEntangledPeer = false; for (auto& port : ports) { if (port->pipe() == m_pipe.ptr()) { targetsEntangledPeer = true; @@ -125,20 +124,25 @@ ExceptionOr MessagePort::postMessage(JSC::JSGlobalObject& state, JSC::JSVa if (disentangled.hasException()) return disentangled.releaseException(); transferredPorts = disentangled.releaseReturnValue(); + } - if (targetsEntangledPeer) { - // Posting the port's own entangled peer: node warns and loses the channel - // rather than throwing. Transferables were already detached above; drop the - // message and close so the dead channel stops reffing the loop. - Bun__Process__emitWarning(defaultGlobalObject(&state), - JSC::JSValue::encode(JSC::jsString(vm, String("The target port was posted to itself, and the communication channel was lost"_s))), - JSC::JSValue::encode(JSC::jsString(vm, String("Warning"_s))), - JSC::JSValue::encode(JSC::jsUndefined()), - JSC::JSValue::encode(JSC::jsUndefined())); - CLEAR_IF_EXCEPTION(warnScope); - close(); - return {}; - } + // A closed port drops the message; the ports taken out of the transfer list go down with it + // (~TransferredMessagePort closes their sides, so their peers see a close), as in node. + if (!isEntangled()) + return {}; + + if (targetsEntangledPeer) { + // Posting the port's own entangled peer: node warns and loses the channel + // rather than throwing. Transferables were already detached above; drop the + // message and close so the dead channel stops reffing the loop. + Bun__Process__emitWarning(defaultGlobalObject(&state), + JSC::JSValue::encode(JSC::jsString(vm, String("The target port was posted to itself, and the communication channel was lost"_s))), + JSC::JSValue::encode(JSC::jsString(vm, String("Warning"_s))), + JSC::JSValue::encode(JSC::jsUndefined()), + JSC::JSValue::encode(JSC::jsUndefined())); + CLEAR_IF_EXCEPTION(warnScope); + close(); + return {}; } m_pipe->send(m_side, MessageWithMessagePorts { messageData.releaseReturnValue(), WTF::move(transferredPorts) }); @@ -154,20 +158,25 @@ void MessagePort::entrySettled() void MessagePort::start() { - if (m_started || !isEntangled()) + if (!isEntangled()) return; - // A node worker's parentPort delivers nothing until the entry module has evaluated; keep the - // request and let entrySettled() perform it. Messages stay buffered in the pipe meanwhile. - if (auto* context = scriptExecutionContext()) { - if (auto* jsGlobal = context->globalObject()) { - auto* globalObject = defaultGlobalObject(jsGlobal); - if (globalObject->nodeParentPort() == this && !globalObject->nodeWorkerEntrySettled()) { - m_startDeferredUntilEntrySettled = true; - return; + if (!m_started) { + // A node worker's parentPort delivers nothing until the entry module has evaluated; keep the + // request and let entrySettled() perform it. Messages stay buffered in the pipe meanwhile. + if (auto* context = scriptExecutionContext()) { + if (auto* jsGlobal = context->globalObject()) { + auto* globalObject = defaultGlobalObject(jsGlobal); + if (globalObject->nodeParentPort() == this && !globalObject->nodeWorkerEntrySettled()) { + m_startDeferredUntilEntrySettled = true; + return; + } } } + m_started = true; } - m_started = true; + // Every call turns delivery (back) on, as node's Start() does for a port whose last 'message' + // listener was removed. attach() re-schedules the drain and re-reports a peer that has closed. + m_receiving = true; auto* context = scriptExecutionContext(); ASSERT(context); @@ -191,10 +200,11 @@ void MessagePort::flushQueuedMessagesBeforeClose() // re-injecting into this closing port (via its entangled peer) can't starve the loop. size_t limit = std::max(MessagePortPipe::queuedCount(m_pipe->state(m_side)), 1000); for (size_t i = 0; i < limit; ++i) { - // A handler (or a microtask it queued) may have transferred this port; the - // remaining inbox now belongs to the new owner. drainAndDispatch()'s - // per-iteration ctxId/port re-check guards the same case. - if (m_isDetached) + // A handler (or a microtask it queued) may have transferred this port (the remaining + // inbox now belongs to the new owner) or removed the last listener, e.g. once(): like + // drainAndDispatch(), stop and leave the rest queued. peerClosed() then keeps the port + // open for a later listener; close() drops it. + if (m_isDetached || !hasMessageEventListener()) break; auto message = m_pipe->takeOne(m_side); if (!message) @@ -271,24 +281,35 @@ void MessagePort::dispatchCloseEvent() void MessagePort::peerClosed() { - if (m_isDetached) + if (m_isDetached || m_isClosing) return; auto* context = scriptExecutionContext(); if (!context || !context->globalObject()) return; Ref protectedThis { *this }; - // Deliver whatever the peer sent before it closed, then fire 'close'. Node orders - // them that way, and registerCloseContext()'s retroactive notify can land before any - // drain is scheduled -- e.g. on('close') registered before on('message'). + // Node delivers what the peer sent before closing first; this task can run before the drain. if (m_started && hasMessageEventListener()) flushQueuedMessagesBeforeClose(); - // Fire 'close' (guarded against a double dispatch) and release this side's loop refs - // so the loop can idle, matching node. - dispatchCloseEvent(); - // jsUnref() clears both the listener loop-ref (m_isRefd) and the onmessage/ref() - // keepalive (m_hasRef), so a listening transferred port stops pinning the loop. - auto* globalObject = defaultGlobalObject(context->globalObject()); - jsUnref(globalObject); + // A handler in that flush closed or transferred this port. + if (m_isDetached || m_isClosing) + return; + if (!m_pipe->isOtherSideClosedByRequest(m_side)) { + // Peer collected, not closed: node never closes a channel over that (see jsRef()), so + // only release the loop refs of a port that is started or listening for 'close'. An + // idle port keeps its one 'close' event for its own close(cb); it is notified again + // if it ever gets a 'close' listener (addEventListener) or starts (attach()). + if (m_started || m_hasCloseEventListener.load(std::memory_order_relaxed)) { + dispatchCloseEvent(); + jsUnref(defaultGlobalObject(context->globalObject())); + } + return; + } + // Not receiving with messages queued: node's close notification waits behind them. The + // next 'message' listener (attach()) or the receiveMessageOnPort() that empties the queue + // (tryTakeMessage()) completes the close. + if (!m_receiving && MessagePortPipe::queuedCount(m_pipe->state(m_side)) > 0) + return; + close(); } TransferredMessagePort MessagePort::disentangle() @@ -381,8 +402,14 @@ JSValue MessagePort::tryTakeMessage(JSGlobalObject* lexicalGlobalObject, bool& h return jsUndefined(); auto message = m_pipe->takeOne(m_side); - if (!message) + if (!message) { + // This read reached the peer's close: node closes the port here (see peerClosed()). + // Peer state before our queue, as in virtualHasPendingActivity(): a message the peer + // sent between takeOne() and its close() is then visible and must not be dropped. + if (m_pipe->isOtherSideClosedByRequest(m_side) && !MessagePortPipe::queuedCount(m_pipe->state(m_side))) + close(); return jsUndefined(); + } hadMessage = true; auto ports = MessagePort::entanglePorts(*context, WTF::move(message->transferredPorts)); @@ -515,19 +542,13 @@ void MessagePort::onDidChangeListenerImpl(EventTarget& self, const AtomString& e bool MessagePort::addEventListener(const AtomString& eventType, Ref&& listener, const AddEventListenerOptions& options) { if (eventType == eventNames().messageEvent) { + // Also resumes a port paused by removing its listeners (see start()). start(); m_hasMessageEventListener = true; - // start() no-ops after the first call; re-attach so a listener re-added after a - // pause re-schedules the drain for messages buffered meanwhile. - if (m_started && isEntangled()) { - if (auto* context = scriptExecutionContext()) - m_pipe->attach(m_side, context->identifier(), ThreadSafeWeakPtr { *this }); - } } else if (eventType == eventNames().closeEvent) { m_hasCloseEventListener.store(true, std::memory_order_release); + // Reports a peer that went away while this port had no listener (see peerClosed()). if (isEntangled()) { - // Record our context with the pipe so the peer's close() can deliver a - // 'close' event even if we never started (no 'message' listener). if (auto* context = scriptExecutionContext()) m_pipe->registerCloseContext(m_side, context->identifier(), ThreadSafeWeakPtr { *this }); } @@ -538,8 +559,12 @@ bool MessagePort::addEventListener(const AtomString& eventType, Ref { // already queued (e.g. after transfer). Passing a null port is allowed and // means "just buffer, don't dispatch" (used before start()). void attach(uint8_t side, ScriptExecutionContextIdentifier, ThreadSafeWeakPtr); - // Like attach() but only records ctxId/port so the peer's close() can deliver - // a 'close' event to a port that never started (no 'message' listener). Does - // NOT enable drains. No-op if already attached/registered/closed. + // Records ctxId/port (unless already attached/registered) so the peer's close() reaches a port + // that never started, without enabling drains, and reports a peer that is already closed. Called + // when a MessagePort is created for this side and when it gains a 'close' listener. void registerCloseContext(uint8_t side, ScriptExecutionContextIdentifier, ThreadSafeWeakPtr); void detach(uint8_t side); // Explicit == a real, permanent close: close(), context teardown, or an orphaned diff --git a/test/js/node/worker_threads/worker_threads.test.ts b/test/js/node/worker_threads/worker_threads.test.ts index 023ea833ffa1..c62468c66693 100644 --- a/test/js/node/worker_threads/worker_threads.test.ts +++ b/test/js/node/worker_threads/worker_threads.test.ts @@ -1785,6 +1785,447 @@ test("MessagePort: peer closing while a port is in transit still delivers 'close }); }); +// When one port of a channel is close()d, node closes the other one as well (once the +// close has propagated): it fires 'close', drops its loop refs, and from then on it is +// rejected as a transferable like any closed port. Bun fired 'close' on the peer but +// left it open, so a dead endpoint could be transferred; a peer with no listeners was +// not told at all. +describe("a port whose peer closed is closed too", () => { + const detached = expect.objectContaining({ + name: "DataCloneError", + message: "MessagePort in transfer list is already detached", + }); + const tick = () => new Promise(r => setImmediate(r)); + + // Probe without transferring anything: the trailing plain object is itself invalid, so the + // post always throws, and the message says whether `port` (checked first, in list order) + // was rejected as detached. + function isRejectedAsDetached(port: MessagePort): boolean { + const carrier = new MessageChannel(); + let message: string | undefined; + try { + carrier.port1.postMessage(null, [port, {} as any]); + } catch (e: any) { + message = e.message; + } finally { + carrier.port1.close(); + carrier.port2.close(); + } + if (message === undefined) throw new Error("the probe post must throw"); + return message === "MessagePort in transfer list is already detached"; + } + + test("with a 'close' listener: rejected by postMessage and structuredClone, transfer stays atomic", async () => { + const { port1, port2 } = new MessageChannel(); + const carrier = new MessageChannel(); + const closed = once(port1, "close"); + port2.close(); + await closed; + + const ab = new ArrayBuffer(8); + expect(() => carrier.port1.postMessage(null, [ab, port1])).toThrow(detached); + // Same as node: a bad port later in the list means nothing earlier in it is transferred. + expect(ab.byteLength).toBe(8); + expect(() => structuredClone(port1, { transfer: [port1] })).toThrow(detached); + expect(port1.hasRef()).toBe(false); + carrier.port1.close(); + carrier.port2.close(); + }); + + test("with a 'message' listener: what the peer sent first is delivered, then the port closes", async () => { + const { port1, port2 } = new MessageChannel(); + const events: string[] = []; + port1.on("message", m => events.push(m)); + const closed = once(port1, "close"); + port2.postMessage("m1"); + port2.close(); + await closed; + events.push("close"); + + expect(events).toEqual(["m1", "close"]); + expect(isRejectedAsDetached(port1)).toBe(true); + expect(port1.hasRef()).toBe(false); + }); + + test("the port is already closed when its own 'close' listener runs", async () => { + const { port1, port2 } = new MessageChannel(); + const { promise, resolve } = Promise.withResolvers(); + port1.on("close", () => resolve(isRejectedAsDetached(port1))); + port2.close(); + expect(await promise).toBe(true); + }); + + test("with no listeners at all", async () => { + const { port1, port2 } = new MessageChannel(); + port2.close(); + // The close reaches port1 as a task; poll for it rather than assume how many turns it takes. + let rejected = false; + for (let i = 0; i < 100 && !rejected; i++) { + await tick(); + rejected = isRejectedAsDetached(port1); + } + expect(rejected).toBe(true); + }); + + test("a 'close' listener added once the port has closed is not called (the event has passed)", async () => { + const { port1, port2 } = new MessageChannel(); + port2.close(); + let rejected = false; + for (let i = 0; i < 100 && !rejected; i++) { + await tick(); + rejected = isRejectedAsDetached(port1); + } + expect(rejected).toBe(true); + + let closes = 0; + port1.on("close", () => closes++); + for (let i = 0; i < 4; i++) await tick(); + expect(closes).toBe(0); + }); + + // postMessage() on a closed port still serializes (node does, so transfers keep their side + // effects): the ArrayBuffer is detached and a port in the list is consumed, and since the + // message goes nowhere that port's own peer sees it close. + test("postMessage() on the closed port is dropped but still consumes its transfer list", async () => { + const { port1, port2 } = new MessageChannel(); + const closed = once(port1, "close"); + port2.close(); + await closed; + + const third = new MessageChannel(); + let thirdPeerCloses = 0; + third.port2.on("close", () => thirdPeerCloses++); + const ab = new ArrayBuffer(8); + expect(() => port1.postMessage({ ab, port: third.port1 }, [ab, third.port1])).not.toThrow(); + for (let i = 0; i < 20 && thirdPeerCloses === 0; i++) await tick(); + expect({ byteLength: ab.byteLength, consumed: isRejectedAsDetached(third.port1), thirdPeerCloses }).toEqual({ + byteLength: 0, + consumed: true, + thirdPeerCloses: 1, + }); + }); + + test("Worker.postMessage() and the Worker constructor's transferList reject it as well", async () => { + const { port1, port2 } = new MessageChannel(); + const closed = once(port1, "close"); + port2.close(); + await closed; + + const worker = new Worker("", { eval: true }); + try { + expect(() => worker.postMessage(port1, [port1])).toThrow(detached); + expect(() => new Worker("", { eval: true, workerData: port1, transferList: [port1] })).toThrow(detached); + } finally { + await worker.terminate(); + } + }); + + // node's close notification sits behind whatever the peer posted before closing: a port + // that is not receiving keeps those messages and stays open until they are consumed. + test("a port that is not receiving stays open until a 'message' listener drains the queue", async () => { + const { port1, port2 } = new MessageChannel(); + const events: string[] = []; + const closed = once(port1, "close"); + port2.postMessage("m1"); + port2.postMessage("m2"); + port2.close(); + for (let i = 0; i < 4; i++) await tick(); + expect({ events, rejected: isRejectedAsDetached(port1) }).toEqual({ events: [], rejected: false }); + + port1.on("message", m => events.push(m)); + await closed; + events.push("close"); + expect(events).toEqual(["m1", "m2", "close"]); + expect(isRejectedAsDetached(port1)).toBe(true); + }); + + test("a port that is not receiving stays open until receiveMessageOnPort() drains the queue", async () => { + const { port1, port2 } = new MessageChannel(); + let closes = 0; + const closed = once(port1, "close"); + port1.on("close", () => closes++); + port2.postMessage("m1"); + port2.close(); + for (let i = 0; i < 4; i++) await tick(); + expect({ closes, rejected: isRejectedAsDetached(port1) }).toEqual({ closes: 0, rejected: false }); + + expect(receiveMessageOnPort(port1)).toEqual({ message: "m1" }); + // Taking the last message does not close it yet (node closes on reaching the notification). + expect({ closes, rejected: isRejectedAsDetached(port1) }).toEqual({ closes: 0, rejected: false }); + + // This call reaches the peer's close: no message, and the port is closed on the spot. + expect(receiveMessageOnPort(port1)).toBeUndefined(); + expect({ closes, rejected: isRejectedAsDetached(port1) }).toEqual({ closes: 0, rejected: true }); + await closed; + expect(closes).toBe(1); + }); + + // Removing the last 'message' listener stops a port in node (setupPortReferencing calls + // stopMessagePort), so a port paused that way is "not receiving" as well: what the peer sent + // before closing waits for the next listener, and only then does the port close. + test.each([ + [ + "on() then off()", + async (port: MessagePort) => { + const handler = () => {}; + port.on("message", handler); + port.off("message", handler); + }, + ], + [ + "a once() listener that has fired", + async (port: MessagePort, peer: MessagePort) => { + const fired = once(port, "message"); + peer.postMessage("consumed by once()"); + await fired; + }, + ], + [ + "onmessage set, then cleared", + async (port: MessagePort) => { + port.onmessage = () => {}; + port.onmessage = null; + }, + ], + ])("a port paused by %s stays open until the next 'message' listener drains the queue", async (_name, pause) => { + const { port1, port2 } = new MessageChannel(); + await pause(port1, port2); + let closes = 0; + port1.on("close", () => closes++); + port2.postMessage("while paused"); + port2.close(); + for (let i = 0; i < 4; i++) await tick(); + expect({ closes, rejected: isRejectedAsDetached(port1) }).toEqual({ closes: 0, rejected: false }); + + const [message] = await once(port1, "message"); + for (let i = 0; i < 20 && closes === 0; i++) await tick(); + expect({ message, closes, rejected: isRejectedAsDetached(port1) }).toEqual({ + message: "while paused", + closes: 1, + rejected: true, + }); + }); + + // Whereas a port that was start()ed and never given a listener is receiving: node emits its + // messages into the void and reaches the close. Bun keeps them buffered instead, so nothing + // would ever drain it; the close must not wait for that, or a ref()'d one pins the loop forever. + test("a start()ed port without a 'message' listener is closed even with a message queued", async () => { + const { port1, port2 } = new MessageChannel(); + port1.ref(); + port1.start(); + let closes = 0; + port1.on("close", () => closes++); + port2.postMessage("dropped"); + port2.close(); + for (let i = 0; i < 20 && closes === 0; i++) await tick(); + expect({ closes, rejected: isRejectedAsDetached(port1), hasRef: port1.hasRef() }).toEqual({ + closes: 1, + rejected: true, + hasRef: false, + }); + }); + + // start() turns a paused port back on as well (node's Start() always sets receiving), whether + // it is called before the peer closes or only after the close was already put on hold. + test.each([ + ["before the peer closes", true], + ["after the peer closed", false], + ])("start() on a paused port resumes it (%s), so the port is closed", async (_name, startFirst) => { + const { port1, port2 } = new MessageChannel(); + port1.ref(); + const handler = () => {}; + port1.on("message", handler); + port1.off("message", handler); + let closes = 0; + port1.on("close", () => closes++); + if (startFirst) port1.start(); + port2.postMessage("dropped"); + port2.close(); + if (!startFirst) { + for (let i = 0; i < 4; i++) await tick(); + expect({ closes, rejected: isRejectedAsDetached(port1) }).toEqual({ closes: 0, rejected: false }); + port1.start(); + } + for (let i = 0; i < 20 && closes === 0; i++) await tick(); + expect({ closes, rejected: isRejectedAsDetached(port1), hasRef: port1.hasRef() }).toEqual({ + closes: 1, + rejected: true, + hasRef: false, + }); + }); + + test("a port left open that way is still a valid transferable; the receiver gets the queue and the close", async () => { + const { port1, port2 } = new MessageChannel(); + const carrier = new MessageChannel(); + let senderCloses = 0; + port1.on("close", () => senderCloses++); + port2.postMessage("m1"); + port2.close(); + for (let i = 0; i < 4; i++) await tick(); + expect(senderCloses).toBe(0); + + carrier.port1.postMessage(port1, [port1]); + const [received] = (await once(carrier.port2, "message")) as [MessagePort]; + const events: string[] = []; + const closed = once(received, "close"); + received.on("message", m => events.push(m)); + await closed; + events.push("close"); + expect(events).toEqual(["m1", "close"]); + carrier.port1.close(); + carrier.port2.close(); + }); + + // For a received port whose peer is already gone, the close notification is processed before + // the first drain, so peerClosed() itself delivers the queue. A once() handler stops the port + // after one message; the delivery must stop there too and leave the rest for the next listener. + test("a once() handler on a received port stops the delivery; the next listener gets the rest", async () => { + const { port1, port2 } = new MessageChannel(); + const carrier = new MessageChannel(); + port2.postMessage("m1"); + port2.postMessage("m2"); + port2.close(); + carrier.port1.postMessage(port1, [port1]); + + const { promise, resolve } = Promise.withResolvers<{ received: MessagePort; first: unknown }>(); + carrier.port2.on("message", (received: MessagePort) => { + // Registered synchronously on receipt, like a handler that sets up the port it was handed. + received.once("message", first => resolve({ received, first })); + }); + const { received, first } = await promise; + let closes = 0; + received.on("close", () => closes++); + for (let i = 0; i < 4; i++) await tick(); + expect({ first, closes, rejected: isRejectedAsDetached(received) }).toEqual({ + first: "m1", + closes: 0, + rejected: false, + }); + + const [second] = await once(received, "message"); + for (let i = 0; i < 20 && closes === 0; i++) await tick(); + expect({ second, closes, rejected: isRejectedAsDetached(received) }).toEqual({ + second: "m2", + closes: 1, + rejected: true, + }); + carrier.port1.close(); + carrier.port2.close(); + }); + + // Makes a channel, drops one port and collects it (its WeakRef clearing is the proof; the + // collection runs ~MessagePort, which posts the peer notification), then gives the surviving + // port a turn to process that notification and returns it. + async function survivorOfCollectedPeer(prepare: (survivor: MessagePort, doomed: MessagePort) => void) { + const { survivor, doomed } = (() => { + const { port1, port2 } = new MessageChannel(); + prepare(port1, port2); + return { survivor: port1, doomed: new WeakRef(port2) }; + })(); + let collected = false; + for (let i = 0; i < 20 && !collected; i++) { + // new WeakRef() / deref() keep the target alive for the rest of the current turn. + await tick(); + Bun.gc(true); + collected = doomed.deref() === undefined; + } + expect(collected).toBe(true); + await tick(); + return survivor; + } + + // Only an explicit close counts. Bun collects an entangled port whose wrapper is unreachable + // (node never does) and notifies the peer so it releases its loop refs; closing the peer on + // top of that would make its transfer fail or not depending on when GC ran. + test("a peer that was garbage collected rather than closed leaves the port transferable", async () => { + const carrier = new MessageChannel(); + const port1 = await survivorOfCollectedPeer(survivor => survivor.on("close", () => {})); + expect(isRejectedAsDetached(port1)).toBe(false); + carrier.port1.postMessage(port1, [port1]); + carrier.port1.close(); + carrier.port2.close(); + }); + + // The wait-for-the-queue-to-drain rule is for an explicit close only: nothing would ever + // re-notify a 'close'-listening port that never starts, and until 'close' has been dispatched + // such a port is pinned in the heap. This is what main did before the change, kept as is. + test("a 'close'-listening port with an unread message is still told when its peer is garbage collected", async () => { + let closes = 0; + const port1 = await survivorOfCollectedPeer((survivor, doomed) => { + survivor.on("close", () => closes++); + doomed.postMessage("unread"); + }); + for (let i = 0; i < 20 && closes === 0; i++) await tick(); + expect({ closes, unread: receiveMessageOnPort(port1), rejected: isRejectedAsDetached(port1) }).toEqual({ + closes: 1, + unread: { message: "unread" }, + rejected: false, + }); + }); + + // ...and an idle port is not even told: its one 'close' event has to stay available for + // the close() that eventually happens (node's test-worker-message-port.js relies on this). + test("a port whose peer was garbage collected still delivers close(cb) when it is closed later", async () => { + const idle = await survivorOfCollectedPeer(() => {}); + const unread = await survivorOfCollectedPeer((_survivor, doomed) => doomed.postMessage("queued")); + + // 'close' is dispatched from a task queued by close(); give it a few turns, then give up. + function closeCallbackFires(port: MessagePort): Promise { + return Promise.race([ + new Promise(resolve => port.close(() => resolve("called"))), + (async () => { + for (let i = 0; i < 20; i++) await tick(); + return "not called"; + })(), + ]); + } + + expect(await closeCallbackFires(idle)).toBe("called"); + // close() from inside the 'message' handler, as node's test does. + const unreadResult = await new Promise(resolve => { + unread.on("message", m => closeCallbackFires(unread).then(cb => resolve(`${m}, close(cb) ${cb}`))); + }); + expect(unreadResult).toBe("queued, close(cb) called"); + }); + + // ...but once it does start listening for 'close' it is told after all (as on main): unlike the + // explicitly closed case above, the port is still open, and until that event has been delivered + // a 'close'-listening port is kept alive in the heap. + test("a 'close' listener added after the peer was garbage collected is still told", async () => { + const port1 = await survivorOfCollectedPeer(() => {}); + let closes = 0; + port1.on("close", () => closes++); + for (let i = 0; i < 20 && closes === 0; i++) await tick(); + expect({ closes, rejected: isRejectedAsDetached(port1) }).toEqual({ closes: 1, rejected: false }); + }); + + // The same missing notification left a ref()'d port with no listeners pinning the loop + // forever after its peer closed; node's port closes and the process exits. + test("a ref()'d port with no listeners lets the process exit once its peer closes", async () => { + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + `const { MessageChannel } = require("worker_threads"); + const { port1, port2 } = new MessageChannel(); + port1.ref(); + port2.close(); + // unref'd, so it only ever fires if something else (the port) is still holding the loop open. + setTimeout(() => { console.log("still running, hasRef=" + port1.hasRef()); process.exit(1); }, 3000).unref(); + process.on("exit", code => console.log("exit " + code + ", hasRef=" + port1.hasRef()));`, + ], + env: bunEnv, + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stdout: stdout.trim(), stderr, exitCode }).toEqual({ + stdout: "exit 0, hasRef=false", + stderr: "", + exitCode: 0, + }); + }); +}); + test("workerData is not unwrapped for a non-node globalThis.Worker", async () => { await using proc = Bun.spawn({ cmd: [