Files
buzz-sheet/lib/events/redis-stream.ts
T

59 lines
1.4 KiB
TypeScript

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<Response> {
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" },
},
);
}