feat: add realtime Discord leaderboards
This commit is contained in:
+31
-21
@@ -26,7 +26,7 @@ type StreamOptions = {
|
||||
const redisEnvelopeSchema = z.object({
|
||||
event: z.string(),
|
||||
data: z.unknown(),
|
||||
});
|
||||
}).strict();
|
||||
|
||||
function durationFromEnvironment(name: string, fallback: number) {
|
||||
const value = Number(process.env[name]);
|
||||
@@ -137,6 +137,23 @@ export class SseEndpoint<T extends EventMap> {
|
||||
);
|
||||
}
|
||||
|
||||
parseRedisMessage(message: string) {
|
||||
try {
|
||||
const envelope = redisEnvelopeSchema.safeParse(JSON.parse(message));
|
||||
if (!envelope.success || !Object.hasOwn(this.events, envelope.data.event)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const event = envelope.data.event as Extract<keyof T, string>;
|
||||
const payload = this.events[event].safeParse(envelope.data.data);
|
||||
if (!payload.success) return null;
|
||||
|
||||
return { event, data: payload.data };
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
async stream({ signal }: StreamOptions = {}) {
|
||||
const subscriber = await getRedisSubscriber();
|
||||
const topic = this.topic;
|
||||
@@ -194,27 +211,12 @@ export class SseEndpoint<T extends EventMap> {
|
||||
};
|
||||
|
||||
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 parsed = this.parseRedisMessage(message);
|
||||
if (!parsed) {
|
||||
console.error(`Ignored invalid Redis SSE message for ${this.topic}`);
|
||||
return;
|
||||
}
|
||||
write(encodeEvent(parsed.event, parsed.data));
|
||||
};
|
||||
|
||||
const stream = new ReadableStream<Uint8Array>({
|
||||
@@ -304,6 +306,14 @@ export const sse = createSseEndpoints({
|
||||
}).strict(),
|
||||
},
|
||||
},
|
||||
leaderboards: {
|
||||
events: {
|
||||
xp: z.object({}).strict(),
|
||||
vc: z.object({
|
||||
channelId: z.string().regex(/^\d+$/),
|
||||
}).strict(),
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
export type SseTopic = keyof typeof sse;
|
||||
|
||||
Reference in New Issue
Block a user