feat: type-safe SSE rewrite

This commit is contained in:
2026-04-09 04:52:14 +07:00 Unverified
parent a59784e9d2
commit 94af52648c
16 changed files with 263 additions and 271 deletions
+10 -10
View File
@@ -9,7 +9,6 @@ import z from "zod";
import { adminCheck } from "./auth";
import { uidRegex } from "./const";
import { db } from "./db";
import { ps } from "./db/redis";
import { cdnReferences } from "./db/references";
import {
artifactSettings,
@@ -22,6 +21,7 @@ import {
tierlistStates,
tierlistVersions,
} from "./db/schema";
import { sse, tlSse } from "./db/sse-endpoints";
import { b2s } from "./utils";
export async function getCharacters(chars: string[]) {
@@ -120,7 +120,7 @@ export async function submitArtifact(
// clear card cache
await db.delete(cards).where(eq(cards.submission, queue.id));
revalidatePath("/artifact/admin");
ps.publish({ type: "submit" }, { topic: "artifact", event: "update" });
sse.artifact.pub("update", { type: "submit" });
return queue;
}
const [queue] = await db
@@ -131,7 +131,7 @@ export async function submitArtifact(
})
.returning({ queue: submissions.queue, id: submissions.id });
revalidatePath("/artifact/admin");
ps.publish({ type: "submit" }, { topic: "artifact", event: "update" });
sse.artifact.pub("update", { type: "submit" });
return queue;
}
@@ -158,7 +158,7 @@ export async function toggleCheck(submissionId: string) {
revalidatePath("/artifact/admin");
revalidatePath("/artifact");
ps.publish({ type: "toggleCheck" }, { topic: "artifact", event: "update" });
sse.artifact.pub("update", { type: "toggleCheck" });
await actionLog(`Toggled an artifact submission check mark`);
}
@@ -175,7 +175,7 @@ export async function toggleLock() {
revalidatePath("/artifact/admin");
revalidatePath("/artifact");
ps.publish({ type: "toggleLock" }, { topic: "artifact", event: "update" });
sse.artifact.pub("update", { type: "toggleLock" });
await actionLog(
`${(existing.length ? existing[0].locked : true) ? "Locked" : "Unlocked"} artifact submission`,
);
@@ -198,7 +198,7 @@ export async function setLimit(limit: number) {
revalidatePath("/artifact/admin");
revalidatePath("/artifact");
ps.publish({ type: "setLimit" }, { topic: "artifact", event: "update" });
sse.artifact.pub("update", { type: "setLimit" });
await actionLog(
`Set artifact submit limit to ${limit < 0 ? "unlimited" : limit}`,
);
@@ -212,7 +212,7 @@ export async function wipe() {
);
revalidatePath("/artifact");
ps.publish({ type: "wipe" }, { topic: "artifact", event: "update" });
sse.artifact.pub("update", { type: "wipe" });
await actionLog(`Deleted artifact submissions`);
redirect("/artifact/admin");
}
@@ -263,7 +263,7 @@ export async function tlState(
revalidatePath(`/api/tl/${list}/states`);
await actionLog(`Updated a state in tierlist ${list}`, data);
ps.publish(states, { topic: `tl.${list}`, event: "update_states" });
tlSse(list).pub("update_states", states);
}
export async function tlPlacements(
@@ -284,7 +284,7 @@ export async function tlPlacements(
revalidatePath(`/api/tl/${list}`);
await actionLog(`Updated a placement in tierlist ${list}`);
ps.publish(placements, { topic: `tl.${list}`, event: "update_placements" });
tlSse(list).pub("update_placements", placements);
}
export async function cdnDelete(ids: string[], force = false) {
@@ -376,5 +376,5 @@ export async function actionLog(text: string, details?: unknown) {
revalidatePath("/admin/log");
} catch {}
if (res) ps.publish(res, { topic: "log" });
if (res) sse.log.pub("update", res);
}
+105 -35
View File
@@ -1,37 +1,22 @@
export const redis = Bun.redis;
import ReconnectingEventSource from "reconnecting-eventsource";
import type z from "zod";
type PubMOTD = { event?: string; data: string };
export const redis = typeof Bun !== "undefined" ? Bun.redis : null;
export function optionalStringify(data: unknown | string) {
try {
if (typeof data !== "string") throw "";
JSON.parse(data);
return data;
} catch {
return JSON.stringify(data);
}
}
export function optionalParse(data: unknown | string) {
if (typeof data === "string") {
try {
return { parsed: true, data: JSON.parse(data) };
} catch {}
}
return { parsed: false, data };
}
type PubPayload = { event?: string; data: unknown };
export class PubSubManager {
private pub;
private sub;
private readonly prefix = "sse:";
constructor() {
if (!redis) throw new Error("Redis is not available in this environment.");
this.pub = redis.duplicate();
this.sub = redis.duplicate();
}
private constructMessage(data: string, event?: string) {
return `${event ? `event: ${event}\n` : ""}data: ${optionalStringify(data)}\n\n`;
private constructMessage(data: unknown, event?: string) {
return `${event ? `event: ${event}\n` : ""}data: ${JSON.stringify(data)}\n\n`;
}
async publish(
@@ -43,7 +28,7 @@ export class PubSubManager {
const pub = await this.pub;
return pub.publish(
`${this.prefix}${topic}`,
JSON.stringify({ data, event }),
JSON.stringify({ data, event } satisfies PubPayload),
);
}
@@ -54,7 +39,7 @@ export class PubSubManager {
motd,
}: {
signal?: AbortSignal;
motd?: PubMOTD | PubMOTD[];
motd?: PubPayload | PubPayload[];
} = {},
) {
console.log(` SUB ${topic}`);
@@ -71,15 +56,9 @@ export class PubSubManager {
ping();
}, 90000); // cloudflare timeout = 100s
}
const write: (payload: {
data: string;
event?: string;
}) => void | ((data: string, event?: string) => void) = (
p: { data: string; event?: string } | string,
event?: string,
) => {
if (typeof p === "string") writer.write(this.constructMessage(p, event));
else writer.write(this.constructMessage(p.data, p.event));
const write = (payload: PubPayload) => {
writer.write(this.constructMessage(payload.data, payload.event));
ping();
};
@@ -91,7 +70,7 @@ export class PubSubManager {
timeout = setTimeout(() => close(), 1.8e6); // 30 minutes
};
const handler = (payload: string) => {
const p = JSON.parse(payload);
const p: PubPayload = JSON.parse(payload);
write(p);
resetTimer();
};
@@ -121,4 +100,95 @@ export class PubSubManager {
}
}
export const ps = new PubSubManager();
export type EventSourceEventMap = Record<string, z.ZodTypeAny>;
export type SubOption = {
endpoint?: URL | string;
onerror?: () => void;
onopen?: () => void;
};
export class EventSourceEndpoint<T extends EventSourceEventMap> {
private manager = ps;
private readonly defaultEndpointUrl = new URL(
`/sse/${this.endpoint}`,
typeof location === "undefined" ? process.env.BASE_URL : location.href,
);
constructor(
private endpoint: string,
private eventMap: T,
) {}
pub<K extends keyof T>(event: K, data: z.infer<T[K]>) {
if (!this.manager)
throw new Error(
"EventSourceEndpoint.pub(...) can only be called on the server.",
);
return this.manager.publish(this.eventMap[event].parse(data), {
topic: this.endpoint,
event: String(event),
});
}
subMany(
events: Partial<{ [K in keyof T]: (data: z.infer<T[K]>) => void }>,
{
endpoint = this.defaultEndpointUrl,
onerror = () => {},
onopen = () => {},
}: SubOption = {},
) {
const es = new ReconnectingEventSource(endpoint);
for (const [event, callback] of Object.entries(events)) {
const listener = (e: MessageEvent<string>) =>
callback?.(JSON.parse(e.data));
es.addEventListener(event, listener);
}
es.onerror = onerror;
es.onopen = onopen;
return { clean: () => es.close(), es };
}
sub<K extends keyof T>(
event: K,
callback: (data: z.infer<T[K]>) => void,
{
endpoint = this.defaultEndpointUrl,
onerror = () => {},
onopen = () => {},
}: SubOption = {},
) {
const listener = (e: MessageEvent<string>) => callback(JSON.parse(e.data));
const es = new ReconnectingEventSource(endpoint);
es.addEventListener(String(event), listener);
es.onerror = onerror;
es.onopen = onopen;
return { clean: () => es.close(), es };
}
stream(opt?: {
signal?: AbortSignal | undefined;
motd?: PubPayload | PubPayload[] | undefined;
}) {
return this.manager?.new(this.endpoint, opt);
}
}
export function sseEndpoint<M extends EventSourceEventMap>(
endpoint: string,
eventMap: M,
): EventSourceEndpoint<M> {
return new EventSourceEndpoint(endpoint, eventMap);
}
export function sseEndpointMap<M extends Record<string, EventSourceEventMap>>(
map: M,
): { [K in keyof M]: EventSourceEndpoint<M[K]> } {
const endpoints = {} as { [K in keyof M]: EventSourceEndpoint<M[K]> };
for (const key in map) {
endpoints[key] = new EventSourceEndpoint(key, map[key]);
}
return endpoints;
}
export const ps = redis === null ? null : new PubSubManager();
+47
View File
@@ -0,0 +1,47 @@
import z from "zod/v4";
import { sseEndpoint, sseEndpointMap } from "./redis";
import type { auditLog, tierlistStates } from "./schema";
export const sse = sseEndpointMap({
// Artifact Admin Listener
artifact: {
update: z.object({
type: z.enum(["setLimit", "wipe", "submit", "toggleCheck", "toggleLock"]),
}),
},
// Rubgram Admin Listener
rubgram: {
update: z
.object({
type: z.enum([
"setLimit",
"wipe",
"toggleCheck",
"toggleLock",
"setFree",
"cancel",
"uploadSlip",
]),
})
.or(
z.object({
type: z.enum(["submit", "paid"]),
sub: z.string(),
}),
),
},
// Passive Update Checker
active: {
version: z.string(),
},
// Admin Live Log
log: {
update: z.custom<typeof auditLog.$inferSelect>(),
},
});
export function tlSse<T extends string>(list: T) {
return sseEndpoint(`tl.${list}`, {
update_states: z.custom<(typeof tierlistStates.$inferSelect)[]>(),
update_placements: z.record(z.string(), z.string().array()),
});
}
-124
View File
@@ -5,130 +5,6 @@ export function cn(...inputs: ClassValue[]) {
return twMerge(clsx(inputs));
}
type ServerEventSource = {
send: (data: unknown, event?: string) => void;
write: (msg: string) => void;
close: () => void;
topic: string;
id: number;
};
type EventSourceMOTD = { event?: string; data: string };
export class EventSourceManager {
private list: ServerEventSource[] = [];
private id = 0;
new(
topic = "_global",
{
onDisconnect,
signal,
motd,
}: {
onDisconnect?: () => void;
signal?: AbortSignal;
motd?: EventSourceMOTD | EventSourceMOTD[];
} = {},
): Response {
let timeout: NodeJS.Timeout | undefined;
this.id++;
if (this.id > 100000) this.id = 0;
const id = this.id;
console.log(` SUB ${topic}#${id} (C${this.list.length})`);
const push = (s: ServerEventSource) => {
this.list.push(s);
// Send MOTD
if (Array.isArray(motd)) motd.map((m) => s.send(m.data, m.event));
else if (motd) s.send(motd.data, motd.event);
};
const remove = (id: number) => {
if (this.list.findIndex((e) => e.id === id) < 0) return;
console.log(` DSC ${topic}#${id}`);
this.list = this.list.toSpliced(
this.list.findIndex((e) => e.id === id),
1,
);
onDisconnect?.();
};
signal?.addEventListener("abort", () => {
clearInterval(timeout);
remove(id);
});
return new Response(
new ReadableStream({
start(controller) {
const encoder = new TextEncoder();
// Register this client
const res = {
send: (data: unknown, event?: string) =>
res.write(
`${event ? `event: ${event}\n` : ""}data: ${JSON.stringify(data)}\n\n`,
),
write: (msg: string) => controller.enqueue(encoder.encode(msg)),
close: () => {
try {
controller.close();
} catch {}
clearTimeout(timeout);
remove(res.id);
},
topic,
id,
};
timeout = setInterval(() => {
if (controller.desiredSize === null) return res.close();
try {
res.write(`:ping\n\n`);
} catch {
res.close();
}
}, 5000);
push(res);
},
cancel() {
clearTimeout(timeout);
remove(id);
},
}),
{
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
Connection: "keep-alive",
"X-Accel-Buffering": "no",
},
},
);
}
publish(
data: unknown,
{ event, topic = "_global" }: { event?: string; topic?: string },
) {
console.log(` PUB #${topic}`);
queueMicrotask(() => {
for (const s of this.list) {
if (s.topic !== topic) continue;
try {
s.send(data, event);
} catch (error) {
console.error(` ERR #${topic}#${s.id}, ${error}`);
}
}
});
}
count(topic = "_global") {
return this.list.reduce((c, s) => c + (s.topic === topic ? 1 : 0), 0);
}
}
export const sse = new EventSourceManager();
export const b2s = (t: number) => {
let e = (Math.log2(t) / 10) | 0;
return `${(t / 1024 ** (e = e <= 0 ? 0 : e)).toFixed(1)}${" KMGP"[e]}B`;