feat : notify new update income

This commit is contained in:
2026-09-26 23:34:42 +07:00 Unverified
parent e99a5b76b1
commit 4ef1968b04
14 changed files with 499 additions and 33 deletions
+65 -11
View File
@@ -3,8 +3,12 @@ import { createEventStream, eventStreamHeaders, type EventStreamChannel } from "
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<EventStreamChannel>>();
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> {
@@ -19,15 +23,15 @@ function getSubscriber(): Subscriber {
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);
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 stream of [...streams]) stream.close();
for (const streams of listeners.values()) {
for (const listener of [...streams]) listener.stream.close();
}
listeners.clear();
client.disconnect();
});
@@ -39,7 +43,18 @@ function getSubscriber(): Subscriber {
return state;
}
export async function createRedisEventResponse(topic: string, signal: AbortSignal): Promise<Response> {
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;
@@ -50,9 +65,9 @@ export async function createRedisEventResponse(topic: string, signal: AbortSigna
await state.client.subscribe(channel);
listeners.set(channel, new Set());
}
let eventStream: EventStreamChannel | undefined;
let listener: StreamListener | undefined;
const release = () => {
if (eventStream) listeners.get(channel)?.delete(eventStream);
if (listener) listeners.get(channel)?.delete(listener);
void ordered(async () => {
if (subscriber === state && listeners.get(channel)?.size === 0) {
listeners.delete(channel);
@@ -62,10 +77,21 @@ export async function createRedisEventResponse(topic: string, signal: AbortSigna
};
try {
signal.throwIfAborted();
eventStream = createEventStream({ signal, onClose: release });
listeners.get(channel)?.add(eventStream);
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;
}
@@ -73,10 +99,38 @@ export async function createRedisEventResponse(topic: string, signal: AbortSigna
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 stream of [...streams]) stream.close();
for (const streams of listeners.values()) {
for (const listener of [...streams]) listener.stream.close();
}
listeners.clear();
state?.client.disconnect();
}