import { expectedTopic, parseInvalidationEvent, } from "@/lib/events/invalidation"; import { createEventStream, eventStreamHeaders, type EventStreamChannel, } from "@/lib/events/sse"; import { connectRedisClient, createRedisClient, redisEventChannel, } from "@/lib/redis/client"; export async function createRedisEventResponse( topic: string, signal: AbortSignal, ): Promise { const redis = createRedisClient(`buzz-sheet-sse:${topic}`); const channel = redisEventChannel(topic); try { await connectRedisClient(redis); if (signal.aborted) throw new Error("Request closed before subscription."); await redis.subscribe(channel); } catch (cause) { redis.disconnect(); throw cause; } const onMessage = (receivedChannel: string, value: string) => { if (receivedChannel !== channel) return; const event = parseInvalidationEvent(value); if (event && expectedTopic(event) === topic) eventStream.send(event); }; const eventStream: EventStreamChannel = createEventStream({ signal, onClose: () => { redis.off("message", onMessage); redis.disconnect(); }, }); redis.on("message", onMessage); return new Response(eventStream.stream, { headers: eventStreamHeaders(), }); } export function eventsUnavailableResponse(): Response { return Response.json( { error: "events-unavailable" }, { status: 503, headers: { "Cache-Control": "no-store" }, }, ); }