Files
buzz-sheet/lib/events/redis-stream.ts
T
gunshiz 87e6bcd96f
CI / Verify and audit (push) Successful in 2m33s
CI / Build, scan and deploy immutable images (push) Failing after 1m29s
feat : 6 astra improve it
2026-09-22 18:28:18 +07:00

83 lines
3.2 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";
type Subscriber = { client: ReturnType<typeof createRedisClient>; ready: Promise<void> };
let subscriber: Subscriber | undefined;
const listeners = new Map<string, Set<EventStreamChannel>>();
// Serialize subscribe/unsubscribe so a last-reader disconnect cannot race a join.
let operations = Promise.resolve();
function ordered<T>(task: () => Promise<T>): Promise<T> {
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<Response> {
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();
}