feat: add Redis-backed SSE updates

This commit is contained in:
2026-07-24 21:31:47 +07:00 Unverified
parent 5cc7edaa13
commit 29df784aeb
22 changed files with 912 additions and 96 deletions
+7 -2
View File
@@ -58,8 +58,13 @@ export async function getFollowerCounts(): Promise<FetchResult> {
const { getRedisClient } = await import("@/lib/redis");
const redis = await getRedisClient();
await Promise.all([
redis.set(CACHE_KEY, JSON.stringify(result.counts), { ex: CACHE_TTL }),
redis.set(CACHE_KEY_TIKTOK, JSON.stringify(result.tiktokProfile), { ex: CACHE_TTL }),
redis.set(CACHE_KEY, JSON.stringify(result.counts), "EX", CACHE_TTL),
redis.set(
CACHE_KEY_TIKTOK,
JSON.stringify(result.tiktokProfile),
"EX",
CACHE_TTL,
),
]);
} catch {
// Redis unavailable, skip caching
+56 -8
View File
@@ -1,15 +1,63 @@
// @ts-expect-error Bun exposes RedisClient at runtime.
import { RedisClient } from "bun";
import type { RedisClient } from "bun";
const globalForRedis = global as unknown as { redis: RedisClient };
type RedisConnections = {
client: RedisClient;
publisher?: Promise<RedisClient>;
subscriber?: Promise<RedisClient>;
};
export const getRedisClient = async () => {
if (!globalForRedis.redis) {
globalForRedis.redis = new RedisClient(process.env.REDIS_URL || "");
const globalForRedis = globalThis as typeof globalThis & {
erikaRedis?: RedisConnections;
};
function getRedisConnections() {
if (typeof Bun === "undefined") {
throw new Error("Redis is only available on the Bun server runtime.");
}
const client = globalForRedis.redis;
if (!globalForRedis.erikaRedis) {
globalForRedis.erikaRedis = { client: Bun.redis };
}
return globalForRedis.erikaRedis;
}
async function connect(client: RedisClient) {
if (!client.connected) {
await client.connect();
}
return client;
};
}
function duplicateConnection(kind: "publisher" | "subscriber") {
const connections = getRedisConnections();
const existing = connections[kind];
if (existing) {
return existing;
}
const pending = connect(connections.client)
.then((client) => client.duplicate())
.then(connect);
connections[kind] = pending;
void pending.catch(() => {
if (connections[kind] === pending) {
connections[kind] = undefined;
}
});
return pending;
}
export async function getRedisClient() {
return connect(getRedisConnections().client);
}
export function getRedisPublisher() {
return duplicateConnection("publisher");
}
export function getRedisSubscriber() {
return duplicateConnection("subscriber");
}
+313
View File
@@ -0,0 +1,313 @@
import ReconnectingEventSource from "reconnecting-eventsource";
import { z } from "zod";
import { getRedisPublisher, getRedisSubscriber } from "@/lib/redis";
const CHANNEL_PREFIX = "erika:sse:";
const DEFAULT_RETRY_MS = 3_000;
const DEFAULT_PUBLISH_TIMEOUT_MS = 2_000;
const DEFAULT_HEARTBEAT_MS = 90_000;
const DEFAULT_CONNECTION_MAX_MS = 30 * 60_000;
type EventMap = Record<string, z.ZodType>;
type EventHandlers<T extends EventMap> = Partial<{
[K in keyof T]: (data: z.infer<T[K]>) => void;
}>;
type SubscriptionOptions = {
endpoint?: string | URL;
onerror?: (event: Event) => void;
onopen?: (event: Event) => void;
};
type StreamOptions = {
signal?: AbortSignal;
};
const redisEnvelopeSchema = z.object({
event: z.string(),
data: z.unknown(),
});
function durationFromEnvironment(name: string, fallback: number) {
const value = Number(process.env[name]);
return Number.isFinite(value) && value > 0 ? value : fallback;
}
function channelName(topic: string) {
return `${CHANNEL_PREFIX}${topic}`;
}
function encodeEvent(event: string, data: unknown) {
return `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`;
}
async function publish(topic: string, event: string, data: unknown) {
let timeout: ReturnType<typeof setTimeout> | undefined;
const operation = getRedisPublisher().then((publisher) =>
publisher.publish(
channelName(topic),
JSON.stringify({ event, data }),
),
);
const deadline = new Promise<never>((_, reject) => {
timeout = setTimeout(
() => reject(new Error("Redis SSE publish timed out")),
DEFAULT_PUBLISH_TIMEOUT_MS,
);
});
try {
await Promise.race([operation, deadline]);
} finally {
if (timeout) clearTimeout(timeout);
}
}
export class SseEndpoint<T extends EventMap> {
constructor(
readonly topic: string,
private readonly events: T,
readonly adminOnly = false,
) {}
pub<K extends Extract<keyof T, string>>(event: K, data: z.input<T[K]>) {
const payload = this.events[event].parse(data);
return publish(this.topic, event, payload).catch((error: unknown) => {
console.error(`Failed to publish SSE event ${this.topic}.${event}`, error);
});
}
subMany(
handlers: EventHandlers<T>,
{
endpoint = `/sse/${encodeURIComponent(this.topic)}`,
onerror,
onopen,
}: SubscriptionOptions = {},
) {
const eventSource = new ReconnectingEventSource(endpoint);
const listeners: Array<[string, (event: MessageEvent<string>) => void]> = [];
for (const event of Object.keys(handlers) as Array<Extract<keyof T, string>>) {
const callback = handlers[event];
if (!callback) continue;
const listener = (message: MessageEvent<string>) => {
try {
const payload = this.events[event].safeParse(JSON.parse(message.data));
if (payload.success) {
callback(payload.data);
} else {
console.error(
`Ignored invalid SSE event ${this.topic}.${event}`,
payload.error,
);
}
} catch (error) {
console.error(`Ignored malformed SSE event ${this.topic}.${event}`, error);
}
};
listeners.push([event, listener]);
eventSource.addEventListener(event, listener);
}
if (onerror) eventSource.onerror = onerror;
if (onopen) eventSource.onopen = onopen;
const clean = () => {
for (const [event, listener] of listeners) {
eventSource.removeEventListener(event, listener);
}
eventSource.close();
};
return { clean, eventSource };
}
sub<K extends Extract<keyof T, string>>(
event: K,
handler: (data: z.infer<T[K]>) => void,
options?: SubscriptionOptions,
) {
return this.subMany(
{ [event]: handler } as unknown as EventHandlers<T>,
options,
);
}
async stream({ signal }: StreamOptions = {}) {
const subscriber = await getRedisSubscriber();
const topic = this.topic;
const channel = channelName(this.topic);
const pendingMessages: string[] = [];
let receive: ((message: string) => void) | undefined;
const listener = (message: string) => {
if (receive) {
receive(message);
} else {
pendingMessages.push(message);
}
};
await subscriber.subscribe(channel, listener);
let heartbeat: ReturnType<typeof setInterval> | undefined;
let renewal: ReturnType<typeof setTimeout> | undefined;
let controller: ReadableStreamDefaultController<Uint8Array> | undefined;
let closed = false;
const encoder = new TextEncoder();
const unsubscribe = () => {
void subscriber.unsubscribe(channel, listener).catch((error: unknown) => {
console.error(`Failed to unsubscribe SSE topic ${this.topic}`, error);
});
};
const clean = () => {
if (closed) return;
closed = true;
if (heartbeat) clearInterval(heartbeat);
if (renewal) clearTimeout(renewal);
signal?.removeEventListener("abort", close);
unsubscribe();
};
const close = () => {
clean();
try {
controller?.close();
} catch {
// The request or reader may have already closed the stream.
}
};
const write = (value: string) => {
if (closed) return;
try {
controller?.enqueue(encoder.encode(value));
} catch {
close();
}
};
const handleMessage = (message: string) => {
try {
const envelope = redisEnvelopeSchema.safeParse(JSON.parse(message));
if (!envelope.success || !Object.hasOwn(this.events, envelope.data.event)) {
console.error(`Ignored invalid Redis SSE message for ${this.topic}`);
return;
}
const event = envelope.data.event as Extract<keyof T, string>;
const payload = this.events[event].safeParse(envelope.data.data);
if (!payload.success) {
console.error(
`Ignored invalid Redis SSE event ${this.topic}.${event}`,
payload.error,
);
return;
}
write(encodeEvent(event, payload.data));
} catch (error) {
console.error(`Ignored malformed Redis SSE message for ${this.topic}`, error);
}
};
const stream = new ReadableStream<Uint8Array>({
start(streamController) {
controller = streamController;
write(`retry: ${DEFAULT_RETRY_MS}\n\n: connected ${topic}\n\n`);
receive = handleMessage;
for (const message of pendingMessages.splice(0)) {
handleMessage(message);
}
heartbeat = setInterval(
() => write(": heartbeat\n\n"),
durationFromEnvironment("SSE_HEARTBEAT_MS", DEFAULT_HEARTBEAT_MS),
);
renewal = setTimeout(
close,
durationFromEnvironment(
"SSE_CONNECTION_MAX_MS",
DEFAULT_CONNECTION_MAX_MS,
),
);
signal?.addEventListener("abort", close, { once: true });
if (signal?.aborted) close();
},
cancel() {
clean();
},
});
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream; charset=utf-8",
"Cache-Control": "no-cache, no-transform",
"Connection": "keep-alive",
"Content-Encoding": "identity",
"X-Accel-Buffering": "no",
},
});
}
}
type EndpointDefinition<T extends EventMap> = {
events: T;
adminOnly?: boolean;
};
export function createSseEndpoints<
T extends Record<string, EndpointDefinition<EventMap>>,
>(definitions: T) {
const endpoints = {} as {
[K in keyof T]: SseEndpoint<T[K]["events"]>;
};
for (const topic of Object.keys(definitions) as Array<keyof T>) {
const definition = definitions[topic];
endpoints[topic] = new SseEndpoint(
String(topic),
definition.events,
definition.adminOnly,
);
}
return endpoints;
}
const updateActionSchema = z.enum(["created", "updated", "deleted"]);
export const sse = createSseEndpoints({
forms: {
events: {
update: z.object({
formId: z.string().min(1),
entity: z.enum(["form", "question"]),
action: updateActionSchema,
}).strict(),
},
},
submissions: {
adminOnly: true,
events: {
update: z.object({
formId: z.string().min(1),
submissionId: z.string().min(1).optional(),
action: updateActionSchema,
}).strict(),
},
},
});
export type SseTopic = keyof typeof sse;
export function isSseTopic(topic: string): topic is SseTopic {
return Object.hasOwn(sse, topic);
}