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
27 changes: 20 additions & 7 deletions NativeScript/runtime/js/node-worker-threads.js
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ const { BroadcastChannel } = require("internal/broadcast-channel");
const {
EventTarget,
defineEventHandler,
dispatchEventRethrowing,
globalEventTarget,
} = require("internal/events");

Expand All @@ -64,7 +65,6 @@ const globalPostMessage = g.postMessage;

const addEventListener = EventTarget.prototype.addEventListener;
const removeEventListener = EventTarget.prototype.removeEventListener;
const dispatchEvent = EventTarget.prototype.dispatchEvent;

// Runs `fn` after the caller returns. Node reports 'online' and 'exit' from
// the thread's own lifecycle; the runtime's Worker has no equivalent signal,
Expand Down Expand Up @@ -101,7 +101,7 @@ class WorkerEmitter {
}
const key = `${type}`;
const list = this.#listeners[key] || (this.#listeners[key] = []);
ArrayPrototypePush(list, { listener, once: true });
ArrayPrototypePush(list, { listener, once: true, fired: false });
return this;
}

Expand All @@ -124,22 +124,31 @@ class WorkerEmitter {
return this.removeListener(type, listener);
}

// Whether a listener was registered, as Node's EventEmitter reports it.
emit(type, arg) {
const list = this.#listeners[type];
if (list === undefined) {
return;
if (list === undefined || list.length === 0) {
return false;
}
const snapshot = ArrayPrototypeSlice(list);
for (let i = 0; i < snapshot.length; i++) {
const entry = snapshot[i];
if (entry.once) {
// Node's once wrapper: a registration fires at most once, even when
// an earlier listener emits the same event again and the nested emit
// fires it first.
if (entry.fired) {
continue;
}
entry.fired = true;
const index = ArrayPrototypeIndexOf(list, entry);
if (index !== -1) {
ArrayPrototypeSplice(list, index, 1);
}
}
FunctionPrototypeCall(entry.listener, this, arg);
}
return true;
}
}

Expand Down Expand Up @@ -182,8 +191,10 @@ class Worker extends WorkerEmitter {
worker.onmessageerror = function (event) {
self.emit("messageerror", event.data);
};
// A truthy return cancels the error, so one an 'error' listener took is
// not reported to the parent's global scope as well.
worker.onerror = function (error) {
self.emit("error", error);
return self.emit("error", error);
};
// The runtime's end-of-worker event: the one place 'exit' comes from, for
// a worker's own close() and for terminate() alike, so nothing the worker
Expand Down Expand Up @@ -344,9 +355,11 @@ ObjectDefineProperty(ParentPort.prototype, SymbolToStringTag, {
let parentPort = null;
if (!isMainThread) {
parentPort = new ParentPort();
// Rethrowing, so a listener that throws reaches the worker's error chain
// (the scope's onerror, then the parent's Worker) the way a throwing
// onmessage on the global scope does.
const relay = function (event) {
FunctionPrototypeCall(
dispatchEvent,
dispatchEventRethrowing(
parentPort,
getCreateMessageEvent()(event.type, event.data, event.ports)
);
Expand Down
59 changes: 59 additions & 0 deletions TestRunner/app/tests/MessagingTests.js
Original file line number Diff line number Diff line change
Expand Up @@ -330,6 +330,65 @@ describe("Messaging runtime edges", function () {
worker = new Worker("./messaging/throwingWorker.js");
});

it("lets a node:worker_threads 'error' listener consume the error", function (done) {
var wt = require("node:worker_threads");
var globalErrors = [];
var listener = function (event) {
globalErrors.push(event.message);
event.preventDefault();
};
addEventListener("error", listener);
var worker = new wt.Worker("~/tests/messaging/throwingWorker.js");
worker.on("error", function (error) {
setTimeout(function () {
removeEventListener("error", listener);
expect(error.message).toContain("boom from worker");
expect(globalErrors).toEqual([]);
worker.terminate();
done();
}, SETTLE);
});
});

it("calls a node:worker_threads once listener once when an earlier listener emits again", function () {
var wt = require("node:worker_threads");
var worker = new wt.Worker("~/tests/eventLoopEchoWorker.js");
var calls = 0;
var nested = false;
worker.on("probe", function () {
if (!nested) {
nested = true;
worker.emit("probe");
}
});
worker.once("probe", function () { calls++; });
worker.emit("probe");
worker.terminate();
expect(calls).toBe(1);
});

it("routes a throw from a parentPort listener to the parent's 'error' listeners", function (done) {
var wt = require("node:worker_threads");
var worker = new wt.Worker("~/tests/messaging/parentPortThrowingWorker.js");
var messages = [];
var finish = function () {
expect(messages.length).toBe(1);
expect(messages[0]).toContain("thrown by a parentPort listener");
worker.terminate();
done();
};
// Nothing else settles the spec when the error never arrives.
var guard = setTimeout(finish, 10000);
worker.on("error", function (error) {
messages.push(error.message);
if (messages.length === 1) {
clearTimeout(guard);
setTimeout(finish, SETTLE);
}
});
worker.postMessage("go");
});

it("forwards the error a throwing scope onerror raised for a rejection, once", function (done) {
var worker = new Worker("./messaging/rejectingWorker.js");
var messages = [];
Expand Down
4 changes: 4 additions & 0 deletions TestRunner/app/tests/messaging/parentPortThrowingWorker.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
var parentPort = require("node:worker_threads").parentPort;
parentPort.on("message", function () {
throw new Error("thrown by a parentPort listener");
});
Loading