Skip to content

Commit e6b3f96

Browse files
panvanodejs-github-bot
authored andcommitted
worker: discard queued messages on termination
Close the outside port and remove its message forwarding listeners synchronously. Closing the port alone is asynchronous and can leave queued messages dispatching when terminate() runs inside a listener. Allow the current event to finish while discarding subsequent messages, including those drained when the backing thread exits. Signed-off-by: Filip Skokan <panva.ip@gmail.com> Assisted-by: Codex PR-URL: #66354 Reviewed-By: Anna Henningsen <anna@addaleax.net> Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Aviv Keller <me@aviv.sh>
1 parent 74564d9 commit e6b3f96

2 files changed

Lines changed: 64 additions & 5 deletions

File tree

‎lib/internal/webworker.js‎

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -109,6 +109,7 @@ const kName = Symbol('kName');
109109
const kNavigator = Symbol('kNavigator');
110110
const kNavigatorBrand = Symbol('kNavigatorBrand');
111111
const kOutsidePort = Symbol('kOutsidePort');
112+
const kRemoveMessageListeners = Symbol('kRemoveMessageListeners');
112113
const kType = Symbol('kType');
113114
const kURL = Symbol('kURL');
114115
const kWorker = Symbol('kWorker');
@@ -133,15 +134,21 @@ function createErrorEvent(init) {
133134
// "If this throws an exception, catch it, fire an event named messageerror
134135
// at messageEventTarget, using MessageEvent, and then return."
135136
function forwardMessageEvents(port, target) {
136-
port.on('message', (data) =>
137+
const onMessage = (data) =>
137138
target.dispatchEvent(lazyMessageEvent('message', {
138139
data,
139140
ports: port[kCurrentlyReceivingPorts],
140-
})));
141-
port.on('messageerror', (data) =>
141+
}));
142+
const onMessageError = (data) =>
142143
target.dispatchEvent(lazyMessageEvent('messageerror', {
143144
data,
144-
})));
145+
}));
146+
port.on('message', onMessage);
147+
port.on('messageerror', onMessageError);
148+
return () => {
149+
port.removeListener('message', onMessage);
150+
port.removeListener('messageerror', onMessageError);
151+
};
145152
}
146153

147154
// Web IDL [Replaceable]: the setter replaces the accessor with an own
@@ -846,7 +853,7 @@ class Worker extends EventTarget {
846853
});
847854

848855
this[kOutsidePort] = this[kWorker][kPublicPort];
849-
forwardMessageEvents(this[kOutsidePort], this);
856+
this[kRemoveMessageListeners] = forwardMessageEvents(this[kOutsidePort], this);
850857

851858
// "Set notHandled to the result of firing an event named error at
852859
// workerObject, using ErrorEvent, with the cancelable attribute
@@ -865,6 +872,10 @@ class Worker extends EventTarget {
865872
validateThisInternalField(this, kWorker, 'Worker');
866873
// "The terminate() method steps are to terminate a worker given this's
867874
// worker."
875+
// Closing a port is asynchronous. Remove its forwarding listeners too,
876+
// since terminate() can run while the backing thread drains queued messages.
877+
this[kRemoveMessageListeners]?.();
878+
this[kOutsidePort].close();
868879
this[kWorker]?.terminate();
869880
}
870881

‎test/parallel/test-webworker-postmessage-lifecycle.js‎

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,3 +28,51 @@ function checkSerialization(worker) {
2828
thread.once('exit', common.mustCall(() => checkSerialization(worker)));
2929
}));
3030
}
31+
32+
// Queue messages before terminating so the test does not depend on thread speed.
33+
{
34+
const source = `
35+
onmessage = ({ data }) => {
36+
for (let i = 0; i < 3; i++) postMessage(i);
37+
const state = new Int32Array(data);
38+
Atomics.store(state, 0, 1);
39+
Atomics.notify(state, 0);
40+
};
41+
`;
42+
const worker = new Worker(`data:text/javascript,${encodeURIComponent(source)}`);
43+
worker.onerror = common.mustNotCall('worker failed');
44+
worker.onmessage = common.mustNotCall('message delivered after terminate()');
45+
const state = new Int32Array(new SharedArrayBuffer(4));
46+
worker.postMessage(state.buffer);
47+
assert.notStrictEqual(Atomics.wait(state, 0, 0, common.platformTimeout(10000)), 'timed-out');
48+
worker.terminate();
49+
worker.terminate();
50+
checkSerialization(worker);
51+
}
52+
53+
// terminate() can run during message dispatch, including while the backing
54+
// thread is draining messages on exit. Finish this event but discard later ones.
55+
for (const closeAfterPosting of [false, true]) {
56+
const source = `
57+
onmessage = ({ data }) => {
58+
for (let i = 0; i < 3; i++) postMessage(i);
59+
const state = new Int32Array(data);
60+
Atomics.store(state, 0, 1);
61+
Atomics.notify(state, 0);
62+
if (${closeAfterPosting}) close();
63+
};
64+
`;
65+
const worker = new Worker(`data:text/javascript,${encodeURIComponent(source)}`);
66+
worker.onerror = common.mustNotCall('worker failed');
67+
worker.addEventListener('message', common.mustCall(({ data }) => {
68+
assert.strictEqual(data, 0);
69+
worker.terminate();
70+
checkSerialization(worker);
71+
}));
72+
worker.addEventListener('message', common.mustCall(({ data }) => {
73+
assert.strictEqual(data, 0);
74+
}));
75+
const state = new Int32Array(new SharedArrayBuffer(4));
76+
worker.postMessage(state.buffer);
77+
assert.notStrictEqual(Atomics.wait(state, 0, 0, common.platformTimeout(10000)), 'timed-out');
78+
}

0 commit comments

Comments
 (0)