From edcc361dd4b213fd9d8761db163f9edea037a877 Mon Sep 17 00:00:00 2001 From: Taras Mankovski Date: Thu, 3 Sep 2026 06:49:20 -0400 Subject: [PATCH 1/3] =?UTF-8?q?=F0=9F=90=9B=20Bind=20`once()`=20listener?= =?UTF-8?q?=20lifetime=20to=20its=20interpreting=20scope?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `once()` registered its listener when the operation was constructed and removed it only when the event fired. A halted scope — a losing `race()` branch is the common case — left the listener attached to the emitter, where it kept firing into a dead scope. Interpret the operation lazily instead: the listener is registered when the operation runs and removed in a synchronous `finally`, so it goes away on the event, on failure, and on halt alike. Bumps @effectionx/node to 0.2.5. --- node/events.test.ts | 162 ++++++++++++++++++++++++++++++++++++++++++++ node/events.ts | 49 ++++++++------ node/package.json | 2 +- 3 files changed, 193 insertions(+), 20 deletions(-) create mode 100644 node/events.test.ts diff --git a/node/events.test.ts b/node/events.test.ts new file mode 100644 index 00000000..947b63b5 --- /dev/null +++ b/node/events.test.ts @@ -0,0 +1,162 @@ +import { EventEmitter } from "node:events"; +import { describe, it } from "@effectionx/vitest"; +import { type Operation, race, spawn } from "effection"; +import { expect } from "expect"; + +import { type EventTargetLike, once } from "./events.ts"; + +/** + * A spawned task does not start until its parent suspends, so anything that + * has to happen while another task is registering its listener must run from + * inside a task of its own. + */ +function* observe(read: () => T): Operation { + const task = yield* spawn(function* () { + return read(); + }); + return yield* task; +} + +/** + * `EventTarget` has no listener count, so count registrations ourselves. + */ +function createEventTarget(): EventTargetLike & { + listenerCount(): number; + dispatch(eventName: string): void; +} { + const target = new EventTarget(); + const listeners = new Set<(event: unknown) => void>(); + + return { + addEventListener(eventName, listener) { + listeners.add(listener); + target.addEventListener(eventName, listener as EventListener); + }, + removeEventListener(eventName, listener) { + listeners.delete(listener); + target.removeEventListener(eventName, listener as EventListener); + }, + listenerCount: () => listeners.size, + dispatch: (eventName) => void target.dispatchEvent(new Event(eventName)), + }; +} + +describe("once", () => { + describe("with an EventEmitter", () => { + it("registers nothing until it is interpreted", function* () { + const emitter = new EventEmitter(); + + once(emitter, "done"); + + expect(emitter.listenerCount("done")).toBe(0); + }); + + it("registers one listener and yields the event arguments", function* () { + const emitter = new EventEmitter(); + + const task = yield* spawn(function* () { + return yield* once<[number, string]>(emitter, "exit"); + }); + + expect(yield* observe(() => emitter.listenerCount("exit"))).toBe(1); + + emitter.emit("exit", 42, "SIGTERM"); + + expect(yield* task).toEqual([42, "SIGTERM"]); + }); + + it("removes the listener before its owner continues", function* () { + const emitter = new EventEmitter(); + + const task = yield* spawn(function* () { + yield* once(emitter, "done"); + return emitter.listenerCount("done"); + }); + + yield* observe(() => emitter.emit("done")); + + expect(yield* task).toBe(0); + }); + + it("removes the listener when the interpreting task is halted", function* () { + const emitter = new EventEmitter(); + + const task = yield* spawn(function* () { + yield* once(emitter, "done"); + }); + + expect(yield* observe(() => emitter.listenerCount("done"))).toBe(1); + + yield* task.halt(); + + expect(emitter.listenerCount("done")).toBe(0); + + emitter.emit("done"); + + expect(emitter.listenerCount("done")).toBe(0); + }); + + it("deregisters the losing branch of a race", function* () { + const emitter = new EventEmitter(); + let racing = 0; + + const value = yield* race([ + once<[string]>(emitter, "loser"), + { + *[Symbol.iterator]() { + racing = emitter.listenerCount("loser"); + return "won"; + }, + }, + ]); + + expect(value).toBe("won"); + expect(racing).toBe(1); + expect(emitter.listenerCount("loser")).toBe(0); + }); + }); + + describe("with an EventTarget", () => { + it("registers nothing until it is interpreted", function* () { + const target = createEventTarget(); + + once(target, "done"); + + expect(target.listenerCount()).toBe(0); + }); + + it("registers one listener and yields the event", function* () { + const target = createEventTarget(); + + const task = yield* spawn(function* () { + return yield* once<[Event]>(target, "done"); + }); + + expect(yield* observe(() => target.listenerCount())).toBe(1); + + target.dispatch("done"); + + const [event] = yield* task; + expect(event.type).toBe("done"); + expect(target.listenerCount()).toBe(0); + }); + + it("removes the listener when the interpreting task is halted", function* () { + const target = createEventTarget(); + + const task = yield* spawn(function* () { + yield* once(target, "done"); + }); + + expect(yield* observe(() => target.listenerCount())).toBe(1); + + yield* task.halt(); + + expect(target.listenerCount()).toBe(0); + + target.dispatch("done"); + + expect(target.listenerCount()).toBe(0); + }); + }); +}); diff --git a/node/events.ts b/node/events.ts index 2e702cc4..7bd67bb2 100644 --- a/node/events.ts +++ b/node/events.ts @@ -109,6 +109,11 @@ export function on( * For EventEmitters, returns an array of arguments. * For EventTargets, returns a single-element array containing the event object. * + * The listener's lifetime is bound to the scope that interprets this operation, + * never to the event firing: constructing the operation registers nothing, and + * the listener is removed when the event arrives, the operation fails, or the + * scope is halted first. + * * @example * ```ts * import { once } from "@effectionx/node/events"; @@ -126,25 +131,31 @@ export function once( target: EventSourceLike | null, eventName: string, ): Operation { - const result = withResolvers(); + return { + *[Symbol.iterator]() { + const result = withResolvers(); - if (target) { - if (isEventTarget(target)) { - // EventTarget style (DOM, web-worker self) - const listener = (event: unknown) => { - result.resolve([event] as TArgs); - target.removeEventListener(eventName, listener); - }; - target.addEventListener(eventName, listener); - } else { - // EventEmitter style (Node.js) - const listener = (...args: unknown[]) => { - result.resolve(args as TArgs); - target.off(eventName, listener); - }; - target.on(eventName, listener); - } - } + if (!target) { + return yield* result.operation; + } - return result.operation; + if (isEventTarget(target)) { + const listener = (event: unknown) => result.resolve([event] as TArgs); + target.addEventListener(eventName, listener); + try { + return yield* result.operation; + } finally { + target.removeEventListener(eventName, listener); + } + } + + const listener = (...args: unknown[]) => result.resolve(args as TArgs); + target.on(eventName, listener); + try { + return yield* result.operation; + } finally { + target.off(eventName, listener); + } + }, + }; } diff --git a/node/package.json b/node/package.json index 7106e7fe..e084c548 100644 --- a/node/package.json +++ b/node/package.json @@ -1,7 +1,7 @@ { "name": "@effectionx/node", "description": "Node.js stream and event emitter adapters for Effection", - "version": "0.2.4", + "version": "0.2.5", "keywords": ["io", "streams"], "type": "module", "main": "./dist/mod.js", From 3668e4095e254f15573efe07957e26b4ccf4fb72 Mon Sep 17 00:00:00 2001 From: Taras Mankovski Date: Thu, 3 Sep 2026 06:55:56 -0400 Subject: [PATCH 2/3] =?UTF-8?q?=F0=9F=93=A6=20Bump=20dependents=20of=20@ef?= =?UTF-8?q?fectionx/node?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `workspace:*` is replaced with an exact version at publish time, so @effectionx/process and @effectionx/watch need releases of their own to ship the scope-bound `once()` fix to consumers. --- process/package.json | 2 +- watch/package.json | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/process/package.json b/process/package.json index 6215178a..9a8e72cf 100644 --- a/process/package.json +++ b/process/package.json @@ -1,7 +1,7 @@ { "name": "@effectionx/process", "description": "Spawn and manage child processes with structured concurrency", - "version": "0.8.2", + "version": "0.8.3", "keywords": ["process"], "type": "module", "main": "./dist/mod.js", diff --git a/watch/package.json b/watch/package.json index 63003039..d570c046 100644 --- a/watch/package.json +++ b/watch/package.json @@ -1,7 +1,7 @@ { "name": "@effectionx/watch", "description": "Run commands and restart them gracefully when source files change", - "version": "0.4.6", + "version": "0.4.7", "keywords": ["process"], "type": "module", "main": "./dist/main.js", From aae99a4c30cd89092142b9708183378ef49becee Mon Sep 17 00:00:00 2001 From: Taras Mankovski Date: Thu, 3 Sep 2026 07:00:51 -0400 Subject: [PATCH 3/3] =?UTF-8?q?=E2=9C=85=20Cover=20owner=20failure=20and?= =?UTF-8?q?=20redelivery=20in=20`once()`=20tests?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Scope failure unwinds `once()` through the same `finally` as a halt, but it reaches it by a different path, so give it a test of its own. Also pin down the window the fix opens: the listener now stays attached between `resolve()` and the owner resuming, so assert that an event redelivered inside that window cannot change what the operation returns. --- node/events.test.ts | 49 ++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 48 insertions(+), 1 deletion(-) diff --git a/node/events.test.ts b/node/events.test.ts index 947b63b5..dac24c5f 100644 --- a/node/events.test.ts +++ b/node/events.test.ts @@ -1,6 +1,6 @@ import { EventEmitter } from "node:events"; import { describe, it } from "@effectionx/vitest"; -import { type Operation, race, spawn } from "effection"; +import { type Operation, race, scoped, spawn, withResolvers } from "effection"; import { expect } from "expect"; import { type EventTargetLike, once } from "./events.ts"; @@ -96,6 +96,53 @@ describe("once", () => { expect(emitter.listenerCount("done")).toBe(0); }); + it("removes the listener when its owning scope fails", function* () { + const emitter = new EventEmitter(); + const failure = withResolvers(); + let registered = 0; + let caught: unknown; + + try { + yield* scoped(function* () { + yield* spawn(function* () { + yield* once(emitter, "done"); + }); + yield* spawn(function* () { + registered = emitter.listenerCount("done"); + failure.reject(new Error("boom")); + }); + yield* failure.operation; + }); + } catch (error) { + caught = error; + } + + expect((caught as Error).message).toBe("boom"); + expect(registered).toBe(1); + expect(emitter.listenerCount("done")).toBe(0); + }); + + it("ignores an event redelivered before its owner resumes", function* () { + const emitter = new EventEmitter(); + + const task = yield* spawn(function* () { + return yield* once<[string]>(emitter, "done"); + }); + + yield* observe(() => { + let reentered = false; + emitter.on("done", () => { + if (!reentered) { + reentered = true; + emitter.emit("done", "second"); + } + }); + emitter.emit("done", "first"); + }); + + expect(yield* task).toEqual(["first"]); + }); + it("deregisters the losing branch of a race", function* () { const emitter = new EventEmitter(); let racing = 0;