refactor(api): queue resume update notifications with events.on

This commit is contained in:
Amruth Pillai
2026-09-29 10:42:33 +02:00
parent 22dbbdc547
commit 29647ec634
2 changed files with 51 additions and 88 deletions
+16 -26
View File
@@ -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<string, (...args: string[]) => 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<Listener>();
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", () => {
+35 -62
View File
@@ -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<void>((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<typeof client>, "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<never>(() => {});
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<void>((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}`);