59 lines
1.4 KiB
TypeScript
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" },
|
|
},
|
|
);
|
|
}
|