137 lines
4.5 KiB
TypeScript
137 lines
4.5 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> };
|
|
type StreamListener = {
|
|
stream: EventStreamChannel;
|
|
receive(value: string): void;
|
|
};
|
|
let subscriber: Subscriber | undefined;
|
|
const listeners = new Map<string, Set<StreamListener>>();
|
|
// 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) => {
|
|
for (const listener of listeners.get(channel) ?? []) listener.receive(value);
|
|
});
|
|
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 listener of [...streams]) listener.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;
|
|
}
|
|
|
|
interface RedisStreamOptions {
|
|
initialEvents?:
|
|
| ReadonlyArray<{ eventName: string; data: string }>
|
|
| (() => Promise<ReadonlyArray<{ eventName: string; data: string }>>);
|
|
receive(stream: EventStreamChannel, value: string): void;
|
|
}
|
|
|
|
async function createRedisStreamResponse(
|
|
topic: string,
|
|
signal: AbortSignal,
|
|
options: RedisStreamOptions,
|
|
): 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 listener: StreamListener | undefined;
|
|
const release = () => {
|
|
if (listener) listeners.get(channel)?.delete(listener);
|
|
void ordered(async () => {
|
|
if (subscriber === state && listeners.get(channel)?.size === 0) {
|
|
listeners.delete(channel);
|
|
await state.client.unsubscribe(channel);
|
|
}
|
|
}).catch(() => undefined);
|
|
};
|
|
try {
|
|
signal.throwIfAborted();
|
|
const eventStream = createEventStream({ signal, onClose: release });
|
|
listener = {
|
|
stream: eventStream,
|
|
receive: (value) => options.receive(eventStream, value),
|
|
};
|
|
listeners.get(channel)?.add(listener);
|
|
const initialEvents = typeof options.initialEvents === "function"
|
|
? await options.initialEvents()
|
|
: options.initialEvents ?? [];
|
|
for (const event of initialEvents) {
|
|
eventStream.sendNamed(event.eventName, event.data);
|
|
}
|
|
return eventStream;
|
|
} catch (cause) {
|
|
listener?.stream.close();
|
|
release();
|
|
throw cause;
|
|
}
|
|
});
|
|
return new Response(stream.stream, { headers: eventStreamHeaders() });
|
|
}
|
|
|
|
export async function createRedisEventResponse(
|
|
topic: string,
|
|
signal: AbortSignal,
|
|
): Promise<Response> {
|
|
return createRedisStreamResponse(topic, signal, {
|
|
receive(stream, value) {
|
|
const event = parseInvalidationEvent(value);
|
|
if (event && expectedTopic(event) === topic) stream.send(event);
|
|
},
|
|
});
|
|
}
|
|
|
|
export async function createRedisNamedEventResponse(
|
|
topic: string,
|
|
eventName: string,
|
|
signal: AbortSignal,
|
|
initialEvents: RedisStreamOptions["initialEvents"] = [],
|
|
): Promise<Response> {
|
|
return createRedisStreamResponse(topic, signal, {
|
|
initialEvents,
|
|
receive(stream, value) {
|
|
stream.sendNamed(eventName, value);
|
|
},
|
|
});
|
|
}
|
|
|
|
export function closeRedisEventStreams() {
|
|
const state = subscriber;
|
|
subscriber = undefined;
|
|
for (const streams of listeners.values()) {
|
|
for (const listener of [...streams]) listener.stream.close();
|
|
}
|
|
listeners.clear();
|
|
state?.client.disconnect();
|
|
}
|