import { expectedTopic, parseInvalidationEvent } from "@/lib/events/invalidation"; import { createEventStream, eventStreamHeaders, type EventStreamChannel } from "@/lib/events/sse"; import { connectRedisClient, createRedisClient, redisEventChannel } from "@/lib/redis/client"; type Subscriber = { client: ReturnType; ready: Promise }; type StreamListener = { stream: EventStreamChannel; receive(value: string): void; }; let subscriber: Subscriber | undefined; const listeners = new Map>(); // Serialize subscribe/unsubscribe so a last-reader disconnect cannot race a join. let operations = Promise.resolve(); function ordered(task: () => Promise): Promise { const result = operations.then(task); operations = result.then(() => undefined, () => undefined); return result; } function getSubscriber(): Subscriber { if (subscriber) return subscriber; const client = createRedisClient("buzz-sheet-sse"); const state: Subscriber = { client, ready: Promise.resolve() }; subscriber = state; client.on("message", (channel, value) => { for (const listener of listeners.get(channel) ?? []) listener.receive(value); }); client.on("close", () => { if (subscriber !== state) return; subscriber = undefined; // Force EventSource reconnect/refetch after a gap in the Pub/Sub delivery. for (const streams of listeners.values()) { for (const listener of [...streams]) listener.stream.close(); } listeners.clear(); client.disconnect(); }); state.ready = connectRedisClient(client).then(() => undefined).catch((cause) => { if (subscriber === state) subscriber = undefined; client.disconnect(); throw cause; }); return state; } interface RedisStreamOptions { initialEvents?: | ReadonlyArray<{ eventName: string; data: string }> | (() => Promise>); receive(stream: EventStreamChannel, value: string): void; } async function createRedisStreamResponse( topic: string, signal: AbortSignal, options: RedisStreamOptions, ): Promise { const channel = redisEventChannel(topic); const state = getSubscriber(); await state.ready; const stream = await ordered(async () => { signal.throwIfAborted(); if (subscriber !== state || state.client.status !== "ready") throw new Error("Subscription lost"); if (!listeners.has(channel)) { await state.client.subscribe(channel); listeners.set(channel, new Set()); } let listener: StreamListener | undefined; const release = () => { if (listener) listeners.get(channel)?.delete(listener); void ordered(async () => { if (subscriber === state && listeners.get(channel)?.size === 0) { listeners.delete(channel); await state.client.unsubscribe(channel); } }).catch(() => undefined); }; try { signal.throwIfAborted(); const eventStream = createEventStream({ signal, onClose: release }); listener = { stream: eventStream, receive: (value) => options.receive(eventStream, value), }; listeners.get(channel)?.add(listener); const initialEvents = typeof options.initialEvents === "function" ? await options.initialEvents() : options.initialEvents ?? []; for (const event of initialEvents) { eventStream.sendNamed(event.eventName, event.data); } return eventStream; } catch (cause) { listener?.stream.close(); release(); throw cause; } }); return new Response(stream.stream, { headers: eventStreamHeaders() }); } export async function createRedisEventResponse( topic: string, signal: AbortSignal, ): Promise { return createRedisStreamResponse(topic, signal, { receive(stream, value) { const event = parseInvalidationEvent(value); if (event && expectedTopic(event) === topic) stream.send(event); }, }); } export async function createRedisNamedEventResponse( topic: string, eventName: string, signal: AbortSignal, initialEvents: RedisStreamOptions["initialEvents"] = [], ): Promise { return createRedisStreamResponse(topic, signal, { initialEvents, receive(stream, value) { stream.sendNamed(eventName, value); }, }); } export function closeRedisEventStreams() { const state = subscriber; subscriber = undefined; for (const streams of listeners.values()) { for (const listener of [...streams]) listener.stream.close(); } listeners.clear(); state?.client.disconnect(); }