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
8 changes: 4 additions & 4 deletions src/jsc/bindings/webcore/MessagePort.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -427,13 +427,13 @@ bool MessagePort::hasPendingActivity() const

// Keep alive while a drain task is pending or mid-dispatch. drainAndDispatch
// pops each message (queued -> 0) before invoking listeners, so the in-hand
// message is invisible to the queued count; without this bit a concurrent GC
// message is invisible to the queued count; without these bits a concurrent GC
// running inside that window (queue empty, peer already closed) severs the
// wrapper weak and the dispatch hits a dead JSEventListener wrapper (debug
// ASSERT m_wrapper). DrainScheduled is set from schedule until the inbox is
// observed empty, covering every dispatch.
// ASSERT m_wrapper). DrainScheduled covers schedule -> drain start;
// Dispatching covers the dispatch loop itself.
Comment thread
robobun marked this conversation as resolved.
uint64_t s = m_pipe->state(m_side);
if (s & MessagePortPipe::DrainScheduled)
if (s & (MessagePortPipe::DrainScheduled | MessagePortPipe::Dispatching))
return true;

// Keep alive if the peer is still open and could send more, or messages are
Expand Down
70 changes: 48 additions & 22 deletions src/jsc/bindings/webcore/MessagePortPipe.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ void MessagePortPipe::drainAndDispatch(uint8_t side, ScriptExecutionContextIdent

RefPtr<MessagePort> port;
size_t limit;
bool ownsDispatching = false;
{
Locker locker { s.lock };
// This task was posted to `expectedCtx` (and is running there). If
Expand All @@ -115,20 +116,36 @@ void MessagePortPipe::drainAndDispatch(uint8_t side, ScriptExecutionContextIdent
return;
}
limit = std::max<size_t>(s.inbox.size(), 1000);
// Trade DrainScheduled for Dispatching before user JS runs: a handler
// can park this loop in a nested event-loop wait, and a send() arriving
// then must post a fresh drain task (#37189). Dispatching keeps
// hasPendingActivity() true across the dispatch window DrainScheduled
// used to cover. A nested drain (posted by such a send) finds the bit
// already set and leaves clearing it to this outer invocation.
Comment thread
robobun marked this conversation as resolved.
ownsDispatching = !(st & Dispatching);
s.state.store((st & ~DrainScheduled) | Dispatching, std::memory_order_release);
}

// Clear Dispatching only if this invocation set it and still owns the
// side; after a detach the bit belongs to the next owner's drain.
Comment thread
robobun marked this conversation as resolved.
auto finish = [&] {
if (!ownsDispatching)
return;
Locker locker { s.lock };
if (s.ctxId == expectedCtx && s.port.get() == port)
s.state.fetch_and(~uint64_t(Dispatching), std::memory_order_acq_rel);
};

// All 'message' listeners removed: the port is paused. Leave the inbox buffered
// and stop draining; a later addEventListener re-schedules this drain.
if (!port->hasMessageEventListener()) {
Locker locker { s.lock };
s.state.fetch_and(~uint64_t(DrainScheduled), std::memory_order_acq_rel);
finish();
return;
}

auto* context = port->scriptExecutionContext();
if (!context || !context->globalObject()) {
Locker locker { s.lock };
s.state.fetch_and(~uint64_t(DrainScheduled), std::memory_order_acq_rel);
finish();
return;
}
auto* globalObject = defaultGlobalObject(context->globalObject());
Expand All @@ -144,18 +161,26 @@ void MessagePortPipe::drainAndDispatch(uint8_t side, ScriptExecutionContextIdent
// MessagePort, so compare port identity too — dispatching to
// the stale (now m_isDetached) `port` would silently drop.
// The new owner's attach() scheduled its own drain; leave the
// inbox for that.
// inbox (and the flags, which detach/close reset) for that.
if (s.ctxId != expectedCtx || s.port.get() != port)
break;
return;
uint64_t st = s.state.load(std::memory_order_relaxed);
if (!(st & Attached) || s.inbox.isEmpty()) {
s.state.store(st & ~DrainScheduled, std::memory_order_release);
break;
if (ownsDispatching)
s.state.store(st & ~uint64_t(Dispatching), std::memory_order_release);
return;
}
if (limit-- == 0) {
// Yield to the rest of the event loop; DrainScheduled stays
// set so concurrent sends don't double-schedule.
rescheduleCtx = s.ctxId;
// Budget spent; yield. If a racing send already posted a
// wakeup let that task drain the rest, else claim the flag
// and reschedule.
Comment thread
robobun marked this conversation as resolved.
if (!(st & DrainScheduled)) {
st |= DrainScheduled;
rescheduleCtx = s.ctxId;
}
if (ownsDispatching)
st &= ~uint64_t(Dispatching);
s.state.store(st, std::memory_order_release);
break;
}
message = s.inbox.takeFirst();
Expand All @@ -167,15 +192,16 @@ void MessagePortPipe::drainAndDispatch(uint8_t side, ScriptExecutionContextIdent
// Node's MakeCallback wraps each emit in an InternalCallbackScope,
// which drains nextTick + microtasks on exit; match that so
// queueMicrotask(cb) inside onmessage runs before the next message.
if (globalObject->drainMicrotasks())
break; // termination pending
if (globalObject->drainMicrotasks()) {
finish();
return; // termination pending
}

// Listeners may have been removed mid-drain (port.off()); pause like the
// pre-loop check instead of dispatching the rest to zero listeners.
if (!port->hasMessageEventListener()) {
Locker locker { s.lock };
s.state.fetch_and(~uint64_t(DrainScheduled), std::memory_order_acq_rel);
break;
finish();
return;
}
}

Expand Down Expand Up @@ -247,12 +273,12 @@ void MessagePortPipe::detach(uint8_t side)
Locker locker { s.lock };
s.ctxId = 0;
s.port = nullptr;
// Drop Attached and DrainScheduled. A drain task already in flight on
// the old context can't be recalled, but it captured the old ctxId and
// drainAndDispatch()'s s.ctxId != expectedCtx check makes it a no-op —
// even if a new owner attach()es to a different context before it runs.
// Messages remain queued for the next owner.
s.state.fetch_and(~uint64_t(Attached | ContextKnown | DrainScheduled), std::memory_order_acq_rel);
// Drop Attached, DrainScheduled and Dispatching. A drain task already in
// flight on the old context can't be recalled, but it captured the old
// ctxId and drainAndDispatch()'s s.ctxId != expectedCtx check makes it a
// no-op — even if a new owner attach()es to a different context before
// it runs. Messages remain queued for the next owner.
Comment thread
robobun marked this conversation as resolved.
s.state.fetch_and(~uint64_t(Attached | ContextKnown | DrainScheduled | Dispatching), std::memory_order_acq_rel);
}

void MessagePortPipe::close(uint8_t side, CloseKind kind)
Expand Down
3 changes: 2 additions & 1 deletion src/jsc/bindings/webcore/MessagePortPipe.h
Original file line number Diff line number Diff line change
Expand Up @@ -42,10 +42,11 @@ class MessagePortPipe final : public ThreadSafeRefCounted<MessagePortPipe> {
// lives in the upper bits so it can be bumped with fetch_add(QueuedOne).
enum State : uint64_t {
Closed = 1ull << 0, // close() was called on this side; drops further deliveries.
DrainScheduled = 1ull << 1, // a drain task for this side is in flight.
DrainScheduled = 1ull << 1, // a posted drain task for this side has not yet started draining.
Attached = 1ull << 2, // ctxId/port are valid; ok to schedule drains.
ContextKnown = 1ull << 3, // ctxId/port are valid for close-notification only (no drains).
ClosedByRequest = 1ull << 4, // the Closed above came from close(), not from the port being collected.
Dispatching = 1ull << 5, // drainAndDispatch is mid-dispatch on this side (GC liveness; see hasPendingActivity).

QueuedShift = 8,
QueuedOne = 1ull << QueuedShift,
Expand Down
72 changes: 36 additions & 36 deletions src/jsc/bindings/webcore/Worker.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -273,60 +273,60 @@ void Worker::enqueueToParent(MessageWithMessagePorts&& message)
// queueMicrotask/Promise callbacks observe messages one at a time, then
// yields and reschedules if more remain.
//
// Unlike MessagePortPipe, Worker sides never transfer, so we don't need to
// re-check port identity each iteration — which lets us swap the whole inbox
// into a local deque under the lock and dispatch without contending with the
// sender. A sustained producer (e.g. a tight postMessage loop) would otherwise
// make every per-message pop a contended acquire.
// drainScheduled is only set while a posted drain task has not yet started
// draining. It must be clear while user JS runs: a handler can park this loop
// in a nested event-loop wait, and a send arriving then has to post a fresh
// wakeup or that wait never wakes (#37189). Per-message pops under the lock
// (rather than swapping the queue out) keep delivery FIFO when such a nested
// drain runs.
Comment thread
robobun marked this conversation as resolved.
template<typename Dispatch>
static inline bool drainInbox(Worker::MessageInbox& inbox, Zig::GlobalObject* globalObject, ScriptExecutionContext& context, Dispatch&& dispatch)
{
size_t limit;
Deque<MessageWithMessagePorts> batch;
{
Locker locker { inbox.lock };
if (inbox.queue.isEmpty()) {
inbox.drainScheduled.store(false, std::memory_order_relaxed);
inbox.drainScheduled.store(false, std::memory_order_relaxed);
if (inbox.queue.isEmpty())
return false;
}
limit = std::max<size_t>(inbox.queue.size(), 1000);
batch = std::exchange(inbox.queue, {});
}

while (true) {
while (!batch.isEmpty()) {
std::optional<MessageWithMessagePorts> message;
{
Locker locker { inbox.lock };
if (inbox.queue.isEmpty())
return false;
if (limit-- == 0) {
// Yield to the rest of the event loop. Return the undrained
// tail to the front of the inbox so it stays ahead of
// anything enqueued concurrently; caller reschedules.
Locker locker { inbox.lock };
while (!batch.isEmpty())
inbox.queue.prepend(batch.takeLast());
// Budget spent; yield. If a racing send already posted a
// wakeup let that task drain the rest, else claim the flag
// and have the caller reschedule.
Comment thread
robobun marked this conversation as resolved.
if (inbox.drainScheduled.load(std::memory_order_relaxed))
return false;
inbox.drainScheduled.store(true, std::memory_order_relaxed);
return true;
}
auto message = batch.takeFirst();

auto ports = MessagePort::entanglePorts(context, WTF::move(message.transferredPorts));
auto event = MessageEvent::create(*context.jsGlobalObject(), message.message.releaseNonNull(), nullptr, WTF::move(ports));
dispatch(event.event);

if (globalObject->drainMicrotasks()) {
// Termination pending. Drop the rest — dispatch is a no-op
// once m_terminateRequested is set (drainToParent), and the
// worker thread is tearing down (drainToWorker).
return false;
}
message = inbox.queue.takeFirst();
}

// Batch exhausted — see if more arrived while we were dispatching.
Locker locker { inbox.lock };
if (inbox.queue.isEmpty()) {
inbox.drainScheduled.store(false, std::memory_order_relaxed);
auto ports = MessagePort::entanglePorts(context, WTF::move(message->transferredPorts));
auto event = MessageEvent::create(*context.jsGlobalObject(), message->message.releaseNonNull(), nullptr, WTF::move(ports));
dispatch(event.event);

if (globalObject->drainMicrotasks()) {
// Termination pending. Drop everything still queued — dispatch
// is a no-op once m_terminateRequested is set (drainToParent),
// and the worker thread is tearing down (drainToWorker).
// Destructing the messages now (outside the lock) closes any
// transferred ports so their peers see 'close' promptly instead
// of waiting for ~Worker.
Comment thread
robobun marked this conversation as resolved.
Deque<MessageWithMessagePorts> dropped;
{
Locker locker { inbox.lock };
dropped = std::exchange(inbox.queue, {});
}
return false;
}
if (limit == 0)
return true; // budget spent; caller reschedules
batch = std::exchange(inbox.queue, {});
}
}

Expand Down
33 changes: 33 additions & 0 deletions test/js/web/workers/message-port-pipe.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -363,6 +363,39 @@ describe("MessagePort pipe", () => {
expect(stdout.trim()).toBe("OK");
expect(exitCode).toBe(0);
});

// Each onmessage enqueues the next message mid-drain, so the drain's budget
// (1000 when the inbox starts with one message) runs out with the inbox
// non-empty while the in-handler send has already posted the continuation
// drain task. Exercises the budget-yield handoff; out-of-order or missing
// delivery fails.
test("self-feeding chain outlives the drain budget and stays in order", async () => {
await using proc = Bun.spawn({
cmd: [
bunExe(),
"-e",
`
const { port1, port2 } = new MessageChannel();
const N = 2500;
let next = 0;
port1.onmessage = e => {
if (e.data !== next) { console.error("out of order", e.data, next); process.exit(1); }
next++;
if (next === N) { console.log("OK"); port1.close(); port2.close(); return; }
port2.postMessage(next);
};
port2.postMessage(0);
`,
],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect(stderr).toBe("");
expect(stdout.trim()).toBe("OK");
expect(exitCode).toBe(0);
});
});

// worker.postMessage / parentPort.postMessage go through the same coalesced
Expand Down
100 changes: 100 additions & 0 deletions test/regression/issue/37189.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
import { expect, test } from "bun:test";
import { bunEnv, bunExe, tempDir } from "harness";

// https://github.com/oven-sh/bun/issues/37189
// Regression in 1.3.14: a Worker/MessagePort message arriving while the
// receiver was parked in a nested event-loop wait (expect().rejects) reached
// from a previous message's continuation was enqueued without posting a
// wakeup, deadlocking the process. The first awaited call below makes the
// second call's assertion run as a continuation of the message dispatch,
// with the drain loop still on the native stack.

const channelMain = `import { expect } from "bun:test";
const { port1, port2 } = new MessageChannel();
port2.onmessage = e => {
port2.postMessage({ type: "reply", id: e.data.id, error: "boom" });
};
const pending = new Map();
port1.onmessage = e => {
const msg = e.data;
if (msg.type !== "reply") return;
const reject = pending.get(msg.id);
pending.delete(msg.id);
reject?.(new Error(msg.error));
};
let seq = 0;
const call = () => new Promise((_resolve, reject) => {
const id = ++seq;
pending.set(id, reject);
port1.postMessage({ id });
});

await call().catch(() => {});
await expect(call()).rejects.toThrow("boom");
port1.close();
port2.close();
console.log("OK");`;

const workerMain = `import { expect } from "bun:test";
const worker = new Worker(new URL("./worker.js", import.meta.url).href);
const pending = new Map();
worker.onmessage = e => {
const msg = e.data;
if (msg.type !== "reply") return;
const reject = pending.get(msg.id);
pending.delete(msg.id);
reject?.(new Error(msg.error));
};
let seq = 0;
const call = () => new Promise((_resolve, reject) => {
const id = ++seq;
pending.set(id, reject);
worker.postMessage({ id });
});

await call().catch(() => {});
await expect(call()).rejects.toThrow("boom");
worker.terminate();
console.log("OK");`;

async function expectExitsCleanly(proc: Bun.Subprocess<"ignore", "pipe", "pipe">) {
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect({ stdout, stderr, exitCode }).toEqual({ stdout: "OK\n", stderr: "", exitCode: 0 });
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

test.concurrent(
"expect().rejects settles when the rejection arrives from a Worker message during a nested wait",
async () => {
using dir = tempDir("issue-37189-worker", {
"worker.js": `self.onmessage = e => {
self.postMessage({ type: "reply", id: e.data.id, error: "boom" });
};`,
"main.js": workerMain,
});
await using proc = Bun.spawn({
cmd: [bunExe(), "main.js"],
env: bunEnv,
cwd: String(dir),
stdout: "pipe",
stderr: "pipe",
timeout: 15_000,
killSignal: "SIGKILL",
});
await expectExitsCleanly(proc);
},
);

test.concurrent(
"expect().rejects settles when the rejection arrives from a MessageChannel message during a nested wait",
async () => {
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", channelMain],
env: bunEnv,
stdout: "pipe",
stderr: "pipe",
timeout: 15_000,
killSignal: "SIGKILL",
});
await expectExitsCleanly(proc);
},
);