diff --git a/NativeScript/runtime/js/node-worker-threads.js b/NativeScript/runtime/js/node-worker-threads.js index e4b8ed1d..e54a8f04 100644 --- a/NativeScript/runtime/js/node-worker-threads.js +++ b/NativeScript/runtime/js/node-worker-threads.js @@ -45,6 +45,7 @@ const { BroadcastChannel } = require("internal/broadcast-channel"); const { EventTarget, defineEventHandler, + dispatchEventRethrowing, globalEventTarget, } = require("internal/events"); @@ -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, @@ -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; } @@ -124,15 +124,23 @@ 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); @@ -140,6 +148,7 @@ class WorkerEmitter { } FunctionPrototypeCall(entry.listener, this, arg); } + return true; } } @@ -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 @@ -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) ); diff --git a/TestRunner/app/tests/MessagingTests.js b/TestRunner/app/tests/MessagingTests.js index f6695984..16e0faef 100644 --- a/TestRunner/app/tests/MessagingTests.js +++ b/TestRunner/app/tests/MessagingTests.js @@ -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 = []; diff --git a/TestRunner/app/tests/messaging/parentPortThrowingWorker.js b/TestRunner/app/tests/messaging/parentPortThrowingWorker.js new file mode 100644 index 00000000..7e2f08c3 --- /dev/null +++ b/TestRunner/app/tests/messaging/parentPortThrowingWorker.js @@ -0,0 +1,4 @@ +var parentPort = require("node:worker_threads").parentPort; +parentPort.on("message", function () { + throw new Error("thrown by a parentPort listener"); +});