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 }; 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) => { const event = parseInvalidationEvent(value); if (!event || redisEventChannel(expectedTopic(event)) !== channel) return; for (const stream of listeners.get(channel) ?? []) stream.send(event); }); 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 stream of [...streams]) 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; } export async function createRedisEventResponse(topic: string, signal: AbortSignal): 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 eventStream: EventStreamChannel | undefined; const release = () => { if (eventStream) listeners.get(channel)?.delete(eventStream); void ordered(async () => { if (subscriber === state && listeners.get(channel)?.size === 0) { listeners.delete(channel); await state.client.unsubscribe(channel); } }).catch(() => undefined); }; try { signal.throwIfAborted(); eventStream = createEventStream({ signal, onClose: release }); listeners.get(channel)?.add(eventStream); return eventStream; } catch (cause) { release(); throw cause; } }); return new Response(stream.stream, { headers: eventStreamHeaders() }); } export function closeRedisEventStreams() { const state = subscriber; subscriber = undefined; for (const streams of listeners.values()) for (const stream of [...streams]) stream.close(); listeners.clear(); state?.client.disconnect(); }