feat : 6 astra improve it
This commit is contained in:
+74
-50
@@ -1,58 +1,82 @@
|
||||
import {
|
||||
expectedTopic,
|
||||
parseInvalidationEvent,
|
||||
} from "@/lib/events/invalidation";
|
||||
import {
|
||||
createEventStream,
|
||||
eventStreamHeaders,
|
||||
type EventStreamChannel,
|
||||
} from "@/lib/events/sse";
|
||||
import {
|
||||
connectRedisClient,
|
||||
createRedisClient,
|
||||
redisEventChannel,
|
||||
} from "@/lib/redis/client";
|
||||
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;
|
||||
}
|
||||
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;
|
||||
}
|
||||
|
||||
const onMessage = (receivedChannel: string, value: string) => {
|
||||
if (receivedChannel !== channel) return;
|
||||
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 && expectedTopic(event) === topic) eventStream.send(event);
|
||||
};
|
||||
const eventStream: EventStreamChannel = createEventStream({
|
||||
signal,
|
||||
onClose: () => {
|
||||
redis.off("message", onMessage);
|
||||
redis.disconnect();
|
||||
},
|
||||
if (!event || redisEventChannel(expectedTopic(event)) !== channel) return;
|
||||
for (const stream of listeners.get(channel) ?? []) stream.send(event);
|
||||
});
|
||||
redis.on("message", onMessage);
|
||||
|
||||
return new Response(eventStream.stream, {
|
||||
headers: eventStreamHeaders(),
|
||||
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 function eventsUnavailableResponse(): Response {
|
||||
return Response.json(
|
||||
{ error: "events-unavailable" },
|
||||
{
|
||||
status: 503,
|
||||
headers: { "Cache-Control": "no-store" },
|
||||
},
|
||||
);
|
||||
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();
|
||||
}
|
||||
|
||||
+14
-2
@@ -1,9 +1,13 @@
|
||||
import type { InvalidationEvent } from "./invalidation";
|
||||
import { HttpError } from "@/lib/security/http";
|
||||
|
||||
export const SSE_HEARTBEAT_MS = 90_000;
|
||||
export const SSE_MAX_LIFETIME_MS = 30 * 60 * 1_000;
|
||||
|
||||
const encoder = new TextEncoder();
|
||||
let openStreams = 0;
|
||||
const MAX_OPEN_STREAMS = 400;
|
||||
const MAX_QUEUED_BYTES = 64 * 1024;
|
||||
|
||||
function namedEvent(eventName: string, data: string): string {
|
||||
if (!/^[a-z0-9_-]+$/iu.test(eventName)) {
|
||||
@@ -45,6 +49,8 @@ export function eventStreamHeaders(): Headers {
|
||||
export function createEventStream(
|
||||
options: EventStreamOptions = {},
|
||||
): EventStreamChannel {
|
||||
if (openStreams >= MAX_OPEN_STREAMS) throw new HttpError(429, "streams-busy", 10);
|
||||
openStreams += 1;
|
||||
const heartbeatMs = options.heartbeatMs ?? SSE_HEARTBEAT_MS;
|
||||
const maximumLifetimeMs =
|
||||
options.maximumLifetimeMs ?? SSE_MAX_LIFETIME_MS;
|
||||
@@ -56,7 +62,12 @@ export function createEventStream(
|
||||
const enqueue = (value: string): boolean => {
|
||||
if (closed || !controller) return false;
|
||||
try {
|
||||
controller.enqueue(encoder.encode(value));
|
||||
const bytes = encoder.encode(value);
|
||||
if ((controller.desiredSize ?? 0) < bytes.byteLength) {
|
||||
close(true);
|
||||
return false;
|
||||
}
|
||||
controller.enqueue(bytes);
|
||||
return true;
|
||||
} catch {
|
||||
close(false);
|
||||
@@ -68,6 +79,7 @@ export function createEventStream(
|
||||
const close = (closeController: boolean) => {
|
||||
if (closed) return;
|
||||
closed = true;
|
||||
openStreams -= 1;
|
||||
if (heartbeat) clearInterval(heartbeat);
|
||||
if (lifetime) clearTimeout(lifetime);
|
||||
options.signal?.removeEventListener("abort", abort);
|
||||
@@ -95,7 +107,7 @@ export function createEventStream(
|
||||
cancel() {
|
||||
close(false);
|
||||
},
|
||||
});
|
||||
}, { highWaterMark: MAX_QUEUED_BYTES, size: (chunk) => chunk?.byteLength ?? 0 });
|
||||
|
||||
return {
|
||||
stream,
|
||||
|
||||
Reference in New Issue
Block a user