feat: implement real-time event handling and SSE support; add event publishing and subscription mechanisms

This commit is contained in:
2026-01-17 01:39:06 +00:00
parent f3c0f492be
commit 56b1c40642
13 changed files with 473 additions and 13 deletions
+95 -9
View File
@@ -1,6 +1,7 @@
import fastify, { type FastifyRequest } from 'fastify';
import cors from '@fastify/cors';
import { fastifyTRPCPlugin } from '@trpc/server/adapters/fastify';
import { buildGameEventChannel, type GameSessionTokenPayload } from '@sammo-ts/common';
import {
createGamePostgresConnector,
createRedisConnector,
@@ -12,11 +13,13 @@ import { resolveGameApiConfigFromEnv } from './config.js';
import { createGameApiContext, type DatabaseClient as _DatabaseClient } from './context.js';
import { buildTurnDaemonStreamKeys } from './daemon/streamKeys.js';
import { RedisTurnDaemonTransport } from './daemon/redisTransport.js';
import { InMemoryFlushStore, RedisGatewayFlushSubscriber } from './auth/flushStore.js';
import { InMemoryFlushStore, RedisGatewayFlushSubscriber, type FlushStore } from './auth/flushStore.js';
import { RedisAccessTokenStore } from './auth/accessTokenStore.js';
import { appRouter } from './router.js';
import { buildBattleSimQueueKeys } from './battleSim/keys.js';
import { RedisBattleSimTransport } from './battleSim/redisTransport.js';
import { RedisRealtimeEventHub } from './realtime/eventHub.js';
import { formatSseFrame } from './realtime/sse.js';
const extractBearerToken = (value: string | string[] | undefined): string | null => {
if (!value) {
@@ -33,6 +36,25 @@ const extractBearerToken = (value: string | string[] | undefined): string | null
return header.trim();
};
const resolveAuthFromToken = async (
token: string | null,
accessTokenStore: RedisAccessTokenStore,
flushStore: FlushStore
): Promise<GameSessionTokenPayload | null> => {
if (!token) {
return null;
}
const stored = await accessTokenStore.get(token);
if (!stored) {
return null;
}
const flushedAt = flushStore.getFlushedAt(stored.user.id);
if (flushedAt && new Date(stored.issuedAt) <= flushedAt) {
return null;
}
return stored;
};
export const createGameApiServer = async () => {
const config = resolveGameApiConfigFromEnv();
const postgres = createGamePostgresConnector(resolvePostgresConfigFromEnv({ schema: config.profile }));
@@ -56,6 +78,13 @@ export const createGameApiServer = async () => {
const flushSubscriber = new RedisGatewayFlushSubscriber(flushSubscriberClient, config.flushChannel, flushStore);
await flushSubscriber.start();
const accessTokenStore = new RedisAccessTokenStore(redis.client, config.profileName);
const realtimeSubscriberClient = redis.client.duplicate();
await realtimeSubscriberClient.connect();
const realtimeHub = new RedisRealtimeEventHub(
realtimeSubscriberClient,
buildGameEventChannel(config.profileName)
);
await realtimeHub.start();
const app = fastify({
logger: true,
@@ -72,14 +101,7 @@ export const createGameApiServer = async () => {
router: appRouter,
createContext: async ({ req }: { req: FastifyRequest }) => {
const token = extractBearerToken(req.headers.authorization);
let auth = null;
if (token) {
const stored = await accessTokenStore.get(token);
if (stored) {
const flushedAt = flushStore.getFlushedAt(stored.user.id);
auth = flushedAt && new Date(stored.issuedAt) <= flushedAt ? null : stored;
}
}
const auth = await resolveAuthFromToken(token, accessTokenStore, flushStore);
return createGameApiContext({
db: postgres.prisma,
redis: redis.client,
@@ -99,6 +121,69 @@ export const createGameApiServer = async () => {
},
});
app.get(config.eventsPath, async (request, reply) => {
const query = request.query as { token?: string };
const tokenFromHeader = extractBearerToken(request.headers.authorization);
const tokenFromQuery = typeof query.token === 'string' ? query.token : null;
const auth = await resolveAuthFromToken(tokenFromHeader ?? tokenFromQuery, accessTokenStore, flushStore);
if (!auth) {
await reply.status(401).send({ ok: false, error: 'unauthorized' });
return;
}
reply.hijack();
reply.raw.setHeader('Content-Type', 'text/event-stream');
reply.raw.setHeader('Cache-Control', 'no-cache');
reply.raw.setHeader('Connection', 'keep-alive');
reply.raw.setHeader('X-Accel-Buffering', 'no');
request.raw.setTimeout(0);
reply.raw.setTimeout?.(0);
reply.raw.flushHeaders?.();
const sendFrame = (payload: string) => {
try {
reply.raw.write(payload);
} catch {
return;
}
};
sendFrame(
formatSseFrame({
event: 'ready',
data: JSON.stringify({ at: new Date().toISOString() }),
})
);
const unsubscribe = realtimeHub.subscribe((event) => {
sendFrame(
formatSseFrame({
event: event.type,
data: JSON.stringify(event),
id: event.at,
})
);
});
const heartbeat = setInterval(() => {
sendFrame(
formatSseFrame({
event: 'ping',
data: JSON.stringify({ at: new Date().toISOString() }),
})
);
}, 15000);
const close = () => {
clearInterval(heartbeat);
unsubscribe();
};
request.raw.on('close', close);
request.raw.on('aborted', close);
});
app.get('/healthz', async () => ({
ok: true,
profile: config.profileName,
@@ -107,6 +192,7 @@ export const createGameApiServer = async () => {
app.addHook('onClose', async () => {
await flushSubscriber.stop();
await flushSubscriberClient.quit();
await realtimeHub.stop();
await redis.disconnect();
await postgres.disconnect();
});