diff --git a/packages/api/src/features/resume/events.test.ts b/packages/api/src/features/resume/events.test.ts index ace6f0705..4edc377bb 100644 --- a/packages/api/src/features/resume/events.test.ts +++ b/packages/api/src/features/resume/events.test.ts @@ -1,3 +1,4 @@ +import { EventEmitter } from "node:events"; import { beforeEach, describe, expect, it, vi } from "vitest"; const pool = vi.hoisted(() => ({ @@ -50,16 +51,14 @@ describe("publishResumeUpdated", () => { }); function makeSubscriber() { - const handlers = new Map void>(); - return { + const emitter = new EventEmitter(); + return Object.assign(emitter, { subscribe: vi.fn().mockResolvedValue(1), disconnect: vi.fn(), - on: vi.fn((event: string, handler: (...args: string[]) => void) => handlers.set(event, handler)), - off: vi.fn((event: string) => handlers.delete(event)), - emit(channel: string, payload: string) { - handlers.get("message")?.(channel, payload); + send(channel: string, payload: string) { + emitter.emit("message", channel, payload); }, - }; + }); } describe("Redis resume subscriptions", () => { @@ -71,10 +70,10 @@ describe("Redis resume subscriptions", () => { const iterator = subscribeResumeUpdated({ resumeId: "r1", userId: "u1", signal: controller.signal }); const first = iterator.next(); const channel = "reactive-resume:preview:resume_updated"; - subscriber.emit("other", JSON.stringify(exampleEvent)); - subscriber.emit(channel, "not json"); - subscriber.emit(channel, JSON.stringify({ ...exampleEvent, userId: "other" })); - subscriber.emit(channel, JSON.stringify(exampleEvent)); + subscriber.send("other", JSON.stringify(exampleEvent)); + subscriber.send(channel, "not json"); + subscriber.send(channel, JSON.stringify({ ...exampleEvent, userId: "other" })); + subscriber.send(channel, JSON.stringify(exampleEvent)); expect(await first).toEqual({ done: false, value: exampleEvent }); const next = iterator.next(); controller.abort(); @@ -82,8 +81,8 @@ describe("Redis resume subscriptions", () => { expect(redis.duplicate).toHaveBeenCalledWith({ commandTimeout: 5_000 }); expect(subscriber.subscribe).toHaveBeenCalledWith(channel); expect(subscriber.disconnect).toHaveBeenCalledTimes(1); - expect(subscriber.off).toHaveBeenCalledWith("message", expect.any(Function)); - expect(subscriber.off).toHaveBeenCalledWith("error", expect.any(Function)); + expect(subscriber.listenerCount("message")).toBe(0); + expect(subscriber.listenerCount("error")).toBe(0); expect(pool.connect).not.toHaveBeenCalled(); }); @@ -112,23 +111,14 @@ describe("Redis resume subscriptions", () => { }); const makeFakeClient = () => { - type Listener = (notification: { channel?: string; payload?: string }) => void; - const listeners = new Set(); - - const client = { + const emitter = new EventEmitter(); + return Object.assign(emitter, { query: vi.fn().mockResolvedValue(undefined), - on: vi.fn((event: string, fn: Listener) => { - if (event === "notification") listeners.add(fn); - }), - off: vi.fn((event: string, fn: Listener) => { - if (event === "notification") listeners.delete(fn); - }), release: vi.fn(), __notify(channel: string, payload: string) { - for (const fn of listeners) fn({ channel, payload }); + emitter.emit("notification", { channel, payload }); }, - }; - return client; + }); }; describe("subscribeResumeUpdated", () => { diff --git a/packages/api/src/features/resume/events.ts b/packages/api/src/features/resume/events.ts index a84d16322..06dad7832 100644 --- a/packages/api/src/features/resume/events.ts +++ b/packages/api/src/features/resume/events.ts @@ -1,3 +1,4 @@ +import { on, once } from "node:events"; import { getPool } from "@reactive-resume/db/client"; import { getRedis, redisKey } from "@reactive-resume/db/redis"; @@ -46,81 +47,53 @@ export async function publishResumeUpdated(event: ResumeUpdatedEvent) { await getPool().query("SELECT pg_notify($1, $2)", [RESUME_UPDATED_CHANNEL, JSON.stringify(event)]); } +/** The event a notification carries, when it's one for this subscription. */ +function readEvent(payload: string | undefined, resumeId: string, userId: string) { + if (!payload) return undefined; + try { + const event = JSON.parse(payload) as unknown; + if (isResumeUpdatedEvent(event) && event.resumeId === resumeId && event.userId === userId) return event; + } catch { + // Ignore malformed notifications; the refetch path is invalidation-only. + } + return undefined; +} + +const isAbort = (error: unknown) => error instanceof Error && error.name === "AbortError"; + export async function* subscribeResumeUpdated({ resumeId, userId, signal }: SubscribeResumeUpdatedInput) { if (signal?.aborted) return; const subscriber = getRedis()?.duplicate({ commandTimeout: 5_000 }); const client = subscriber ? undefined : await getPool().connect(); const channel = subscriber ? redisKey(RESUME_UPDATED_CHANNEL) : RESUME_UPDATED_CHANNEL; - const queue: ResumeUpdatedEvent[] = []; - let done = signal?.aborted ?? false; - let wake: (() => void) | undefined; - let failure: Error | undefined; - let stop: () => void = () => {}; - const stopped = new Promise((resolve) => { - stop = resolve; - }); - - const resolveWake = () => { - wake?.(); - wake = undefined; - }; - - const onAbort = () => { - done = true; - stop(); - resolveWake(); - }; - const onError = (error: Error) => { - failure = error; - onAbort(); - }; - - const onNotification = (notification: PgNotification) => { - if (notification.channel !== channel || !notification.payload) return; - - try { - const event = JSON.parse(notification.payload) as unknown; - if (!isResumeUpdatedEvent(event)) return; - if (event.resumeId !== resumeId || event.userId !== userId) return; - - queue.push(event); - resolveWake(); - } catch { - // Ignore malformed notifications; the refetch path is invalidation-only. - } - }; - const onMessage = (messageChannel: string, payload: string) => onNotification({ channel: messageChannel, payload }); - - signal?.addEventListener("abort", onAbort, { once: true }); - client?.on("notification", onNotification); - subscriber?.on("message", onMessage); - subscriber?.on("error", onError); + // `on` queues notifications from now on, so none that arrive while subscribing are lost. It rethrows the + // connection's "error" event and ends with an AbortError when the signal aborts. + const notifications = subscriber + ? on(subscriber, "message", { signal }) + : on(client as NonNullable, "notification", { signal }); try { - if (subscriber) await Promise.race([subscriber.subscribe(channel), stopped]); - else await client?.query(`LISTEN ${RESUME_UPDATED_CHANNEL}`); + if (subscriber) { + // A subscription that never connects still ends when the caller goes away. + const aborted = signal ? once(signal, "abort") : new Promise(() => {}); + await Promise.race([subscriber.subscribe(channel), aborted]); + if (signal?.aborted) return; + } else await client?.query(`LISTEN ${RESUME_UPDATED_CHANNEL}`); - while (!done) { - const event = queue.shift(); - if (event) { - yield event; - continue; - } - - await new Promise((resolve) => { - wake = resolve; - }); + for await (const args of notifications) { + const [notificationChannel, payload] = subscriber + ? (args as [string, string]) + : [(args[0] as PgNotification).channel, (args[0] as PgNotification).payload]; + if (notificationChannel !== channel) continue; + const event = readEvent(payload, resumeId, userId); + if (event) yield event; } - if (failure) throw failure; + } catch (error) { + if (!isAbort(error) || !signal?.aborted) throw error; } finally { - signal?.removeEventListener("abort", onAbort); - client?.off("notification", onNotification); - subscriber?.off("message", onMessage); - if (subscriber) { // The duplicate connection exists only for this subscription; closing it unsubscribes. subscriber.disconnect(); - subscriber.off("error", onError); } else if (client) { try { await client.query(`UNLISTEN ${RESUME_UPDATED_CHANNEL}`);