diff --git a/app/game-api/src/config.ts b/app/game-api/src/config.ts index dc08bdc4..4d8ba1e4 100644 --- a/app/game-api/src/config.ts +++ b/app/game-api/src/config.ts @@ -2,6 +2,7 @@ export interface GameApiConfig { host: string; port: number; trpcPath: string; + eventsPath: string; profile: string; scenario: string; profileName: string; @@ -37,6 +38,7 @@ export const resolveGameApiConfigFromEnv = (env: NodeJS.ProcessEnv = process.env host: env.GAME_API_HOST ?? '0.0.0.0', port: parseNumber(env.GAME_API_PORT, 14000, 'GAME_API_PORT'), trpcPath: env.TRPC_PATH ?? '/trpc', + eventsPath: env.GAME_API_EVENTS_PATH ?? '/events', profile, scenario, profileName, diff --git a/app/game-api/src/realtime/eventHub.ts b/app/game-api/src/realtime/eventHub.ts new file mode 100644 index 00000000..94efa694 --- /dev/null +++ b/app/game-api/src/realtime/eventHub.ts @@ -0,0 +1,64 @@ +import type { RedisConnector } from '@sammo-ts/infra'; +import type { RealtimeEvent } from '@sammo-ts/common'; + +export type RealtimeListener = (event: RealtimeEvent) => void; + +export const parseRealtimeEvent = (message: string): RealtimeEvent | null => { + if (!message) { + return null; + } + try { + const parsed = JSON.parse(message) as RealtimeEvent; + if (!parsed || typeof parsed !== 'object') { + return null; + } + if (typeof parsed.type !== 'string') { + return null; + } + return parsed; + } catch { + return null; + } +}; + +// Redis pub/sub 이벤트를 SSE 구독자에게 전달하는 중계 허브. +export class RedisRealtimeEventHub { + private readonly listeners = new Set(); + private subscribed = false; + + constructor( + private readonly redis: RedisConnector['client'], + private readonly channel: string + ) {} + + async start(): Promise { + if (this.subscribed) { + return; + } + await this.redis.subscribe(this.channel, (message) => { + const event = parseRealtimeEvent(message); + if (!event) { + return; + } + for (const listener of this.listeners) { + listener(event); + } + }); + this.subscribed = true; + } + + subscribe(listener: RealtimeListener): () => void { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + } + + async stop(): Promise { + if (this.subscribed) { + await this.redis.unsubscribe(this.channel); + this.subscribed = false; + } + await this.redis.quit(); + } +} diff --git a/app/game-api/src/realtime/publisher.ts b/app/game-api/src/realtime/publisher.ts new file mode 100644 index 00000000..3607d75f --- /dev/null +++ b/app/game-api/src/realtime/publisher.ts @@ -0,0 +1,12 @@ +import type { RedisConnector } from '@sammo-ts/infra'; +import { buildGameEventChannel, type RealtimeEvent } from '@sammo-ts/common'; + +// 게임 서버의 실시간 이벤트를 Redis pub/sub 채널로 송신한다. +export const publishRealtimeEvent = async ( + redis: RedisConnector['client'], + profileName: string, + event: RealtimeEvent +): Promise => { + const channel = buildGameEventChannel(profileName); + await redis.publish(channel, JSON.stringify(event)); +}; diff --git a/app/game-api/src/realtime/sse.ts b/app/game-api/src/realtime/sse.ts new file mode 100644 index 00000000..4b0089d9 --- /dev/null +++ b/app/game-api/src/realtime/sse.ts @@ -0,0 +1,30 @@ +export interface SseFrame { + event?: string; + data?: string; + id?: string; + retry?: number; +} + +const splitLines = (value: string): string[] => value.split(/\r?\n/); + +export const formatSseFrame = (frame: SseFrame): string => { + const lines: string[] = []; + + if (frame.event) { + lines.push(`event: ${frame.event}`); + } + if (frame.id) { + lines.push(`id: ${frame.id}`); + } + if (frame.retry !== undefined) { + lines.push(`retry: ${frame.retry}`); + } + + const dataLines = splitLines(frame.data ?? ''); + for (const line of dataLines) { + lines.push(`data: ${line}`); + } + + lines.push(''); + return lines.join('\n'); +}; diff --git a/app/game-api/src/router/messages/index.ts b/app/game-api/src/router/messages/index.ts index adee5958..b51b37f9 100644 --- a/app/game-api/src/router/messages/index.ts +++ b/app/game-api/src/router/messages/index.ts @@ -17,6 +17,7 @@ import { insertMessage, type MessageView, } from '../../messages/store.js'; +import { publishRealtimeEvent } from '../../realtime/publisher.js'; const zMessageType = z.enum(['private', 'public', 'national', 'diplomacy']); @@ -255,6 +256,19 @@ export const messagesRouter = router({ draft ); + try { + await publishRealtimeEvent(ctx.redis, ctx.profile.name, { + type: 'messageCreated', + at: now.toISOString(), + mailbox: input.mailbox, + msgType, + messageId: result.receiverId, + senderId: general.id, + }); + } catch { + // 실시간 알림 실패는 메시지 전송 실패로 취급하지 않는다. + } + return { msgType, msgId: result.receiverId }; }), }); diff --git a/app/game-api/src/server.ts b/app/game-api/src/server.ts index b759b390..af620d81 100644 --- a/app/game-api/src/server.ts +++ b/app/game-api/src/server.ts @@ -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 => { + 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(); }); diff --git a/app/game-api/test/realtimeSse.test.ts b/app/game-api/test/realtimeSse.test.ts new file mode 100644 index 00000000..2412a38c --- /dev/null +++ b/app/game-api/test/realtimeSse.test.ts @@ -0,0 +1,57 @@ +import { describe, expect, it } from 'vitest'; + +import { buildGameEventChannel, type RealtimeEvent } from '@sammo-ts/common'; +import { parseRealtimeEvent } from '../src/realtime/eventHub.js'; +import { formatSseFrame } from '../src/realtime/sse.js'; + +describe('formatSseFrame', () => { + it('renders basic SSE payloads', () => { + const output = formatSseFrame({ + event: 'ping', + id: '1', + retry: 1500, + data: 'ok', + }); + + expect(output).toBe(['event: ping', 'id: 1', 'retry: 1500', 'data: ok', ''].join('\n')); + }); + + it('splits multiline data', () => { + const output = formatSseFrame({ + event: 'notice', + data: 'first\nsecond', + }); + + expect(output).toBe(['event: notice', 'data: first', 'data: second', ''].join('\n')); + }); +}); + +describe('parseRealtimeEvent', () => { + it('accepts valid realtime events', () => { + const payload: RealtimeEvent = { + type: 'turnCompleted', + at: '2026-01-01T00:00:00.000Z', + result: { + lastTurnTime: '2026-01-01T00:00:00.000Z', + processedGenerals: 2, + processedTurns: 1, + durationMs: 1200, + partial: false, + }, + }; + const parsed = parseRealtimeEvent(JSON.stringify(payload)); + + expect(parsed?.type).toBe('turnCompleted'); + }); + + it('rejects invalid payloads', () => { + expect(parseRealtimeEvent('not-json')).toBeNull(); + expect(parseRealtimeEvent(JSON.stringify({}))).toBeNull(); + }); +}); + +describe('buildGameEventChannel', () => { + it('namespaces channels by profile', () => { + expect(buildGameEventChannel('che:default')).toBe('sammo:che:default:realtime:events'); + }); +}); diff --git a/app/game-engine/src/turn/turnDaemon.ts b/app/game-engine/src/turn/turnDaemon.ts index d84726dd..4d8ded42 100644 --- a/app/game-engine/src/turn/turnDaemon.ts +++ b/app/game-engine/src/turn/turnDaemon.ts @@ -1,4 +1,5 @@ import type { TurnCommandProfile, TurnSchedule } from '@sammo-ts/logic'; +import { buildGameEventChannel, type RealtimeEvent } from '@sammo-ts/common'; import { createRedisConnector, resolveRedisConfigFromEnv } from '@sammo-ts/infra'; import { SystemClock } from '../lifecycle/clock.js'; @@ -148,6 +149,7 @@ export const createTurnDaemonRuntime = async (options: TurnDaemonRuntimeOptions) const clock = options.clock ?? new SystemClock(); let hooks: TurnDaemonHooks | undefined; + let publishRealtimeEvent: ((event: RealtimeEvent) => Promise) | null = null; let close = async () => {}; let redisCommandStream: RedisTurnDaemonCommandStream | null = null; let redisConnector: ReturnType | null = null; @@ -214,10 +216,35 @@ export const createTurnDaemonRuntime = async (options: TurnDaemonRuntimeOptions) if (redisConfig) { redisConnector = createRedisConnector(redisConfig); await redisConnector.connect(); - redisCommandStream = new RedisTurnDaemonCommandStream(redisConnector.client, { + const redisClient = redisConnector.client; + redisCommandStream = new RedisTurnDaemonCommandStream(redisClient, { keys: buildTurnDaemonStreamKeys(options.profileName ?? options.profile), startId: options.commandStreamStartId, }); + const realtimeChannel = buildGameEventChannel(options.profileName ?? options.profile); + publishRealtimeEvent = async (event: RealtimeEvent) => { + await redisClient.publish(realtimeChannel, JSON.stringify(event)); + }; + } + + if (publishRealtimeEvent) { + const basePublishEvents = hooks?.publishEvents; + // 턴 처리 완료 이벤트를 실시간 채널로 전파한다. + hooks = { + ...hooks, + publishEvents: async (result) => { + try { + await publishRealtimeEvent({ + type: 'turnCompleted', + at: new Date().toISOString(), + result, + }); + } catch { + // 실시간 이벤트 전송 실패는 턴 처리 결과에 영향을 주지 않는다. + } + await basePublishEvents?.(result); + }, + }; } const baseClose = close; diff --git a/app/game-frontend/src/env.d.ts b/app/game-frontend/src/env.d.ts index 665c9d1f..1d1c18cd 100644 --- a/app/game-frontend/src/env.d.ts +++ b/app/game-frontend/src/env.d.ts @@ -9,6 +9,7 @@ declare module '*.vue' { interface ImportMetaEnv { readonly VITE_GATEWAY_API_URL?: string; readonly VITE_GAME_API_URL?: string; + readonly VITE_GAME_SSE_URL?: string; readonly VITE_GAME_ASSET_URL?: string; readonly VITE_GAME_PROFILE?: string; } diff --git a/app/game-frontend/src/stores/mainDashboard.ts b/app/game-frontend/src/stores/mainDashboard.ts index dcae8890..b3e196b6 100644 --- a/app/game-frontend/src/stores/mainDashboard.ts +++ b/app/game-frontend/src/stores/mainDashboard.ts @@ -1,8 +1,10 @@ -import { computed, ref } from 'vue'; +import { computed, ref, watch } from 'vue'; import { defineStore } from 'pinia'; import { MESSAGE_MAILBOX_NATIONAL_BASE, MESSAGE_MAILBOX_PUBLIC, type MessageType } from '@sammo-ts/logic'; +import type { RealtimeEvent } from '@sammo-ts/common'; import { trpc } from '../utils/trpc'; import { useMapViewerStore } from './mapViewer'; +import { useSessionStore } from './session'; const resolveErrorMessage = (value: unknown): string => { if (value instanceof Error) { @@ -25,7 +27,7 @@ export const useMainDashboardStore = defineStore('mainDashboard', () => { const loading = ref(false); const error = ref(null); const realtimeEnabled = ref(true); - const realtimeStatus = ref<'idle' | 'connected' | 'paused'>('connected'); + const realtimeStatus = ref<'idle' | 'connected' | 'paused'>('idle'); const generalContext = ref(null); const lobbyInfo = ref(null); @@ -43,6 +45,7 @@ export const useMainDashboardStore = defineStore('mainDashboard', () => { const generalId = computed(() => general.value?.id ?? null); const nationId = computed(() => nation.value?.id ?? null); const mapViewer = useMapViewerStore(); + const session = useSessionStore(); const selectedCity = computed(() => { const layout = mapLayout.value; @@ -108,7 +111,9 @@ export const useMainDashboardStore = defineStore('mainDashboard', () => { const setRealtimeEnabled = (enabled: boolean) => { realtimeEnabled.value = enabled; - realtimeStatus.value = enabled ? 'connected' : 'paused'; + if (!enabled) { + realtimeStatus.value = 'paused'; + } }; const loadMainData = async () => { @@ -216,6 +221,139 @@ export const useMainDashboardStore = defineStore('mainDashboard', () => { } }; + let realtimeSource: EventSource | null = null; + let realtimeToken: string | null = null; + + const isAccessToken = (token: string | null): boolean => Boolean(token?.startsWith('ga_')); + + const buildRealtimeUrl = (token: string): string => { + const base = import.meta.env.VITE_GAME_SSE_URL ?? '/events'; + const url = new URL(base, window.location.origin); + url.searchParams.set('token', token); + return url.toString(); + }; + + const parseRealtimePayload = (raw: MessageEvent): RealtimeEvent | null => { + if (!raw.data || typeof raw.data !== 'string') { + return null; + } + try { + const parsed = JSON.parse(raw.data) as RealtimeEvent; + if (!parsed || typeof parsed !== 'object') { + return null; + } + if (typeof parsed.type !== 'string') { + return null; + } + return parsed; + } catch { + return null; + } + }; + + const isMailboxRelevant = (mailbox: number): boolean => { + if (mailbox === MESSAGE_MAILBOX_PUBLIC) { + return true; + } + const currentGeneralId = generalId.value; + if (currentGeneralId && mailbox === currentGeneralId) { + return true; + } + const currentNationId = nationId.value; + if (currentNationId && mailbox === MESSAGE_MAILBOX_NATIONAL_BASE + currentNationId) { + return true; + } + return false; + }; + + const closeRealtimeSource = () => { + if (!realtimeSource) { + return; + } + realtimeSource.close(); + realtimeSource = null; + realtimeToken = null; + }; + + const ensureAccessToken = async (): Promise => { + if (!session.gameToken) { + return null; + } + if (isAccessToken(session.gameToken)) { + return session.gameToken; + } + const exchanged = await session.exchangeGatewayToken(); + if (!exchanged) { + return null; + } + return session.gameToken && isAccessToken(session.gameToken) ? session.gameToken : null; + }; + + const connectRealtime = async () => { + if (typeof window === 'undefined') { + return; + } + if (!realtimeEnabled.value || !session.isReady || !session.hasGeneral) { + return; + } + const token = await ensureAccessToken(); + if (!token) { + realtimeStatus.value = 'idle'; + return; + } + if (realtimeSource && realtimeToken === token) { + return; + } + closeRealtimeSource(); + realtimeToken = token; + realtimeStatus.value = 'idle'; + + const source = new EventSource(buildRealtimeUrl(token)); + realtimeSource = source; + + source.addEventListener('open', () => { + realtimeStatus.value = 'connected'; + }); + source.addEventListener('error', () => { + realtimeStatus.value = realtimeEnabled.value ? 'idle' : 'paused'; + }); + source.addEventListener('turnCompleted', () => { + void loadMainData(); + }); + source.addEventListener('messageCreated', (event) => { + const payload = parseRealtimePayload(event); + if (!payload || payload.type !== 'messageCreated') { + return; + } + if (isMailboxRelevant(payload.mailbox)) { + void refreshMessages(); + } + }); + source.addEventListener('ping', () => { + if (realtimeEnabled.value) { + realtimeStatus.value = 'connected'; + } + }); + }; + + watch( + () => [realtimeEnabled.value, session.isReady, session.hasGeneral, session.gameToken], + ([enabled, ready, hasGeneral]) => { + if (!enabled) { + closeRealtimeSource(); + realtimeStatus.value = 'paused'; + return; + } + if (!ready || !hasGeneral) { + closeRealtimeSource(); + realtimeStatus.value = 'idle'; + return; + } + void connectRealtime(); + }, + { immediate: true } + ); + return { loading, error, diff --git a/packages/common/src/index.ts b/packages/common/src/index.ts index 739565b7..7093d8f4 100644 --- a/packages/common/src/index.ts +++ b/packages/common/src/index.ts @@ -11,3 +11,5 @@ export * from './util/RandUtil.js'; export * from './util/TestRNG.js'; export * from './util/sha512.js'; export * from './turnDaemon/types.js'; +export * from './realtime/keys.js'; +export * from './realtime/types.js'; diff --git a/packages/common/src/realtime/keys.ts b/packages/common/src/realtime/keys.ts new file mode 100644 index 00000000..2b638cc7 --- /dev/null +++ b/packages/common/src/realtime/keys.ts @@ -0,0 +1,9 @@ +const normalizeProfileName = (profileName: string): string => { + const trimmed = profileName.trim(); + return trimmed.length > 0 ? trimmed : 'unknown'; +}; + +export const buildGameEventChannel = (profileName: string): string => { + const normalized = normalizeProfileName(profileName); + return `sammo:${normalized}:realtime:events`; +}; diff --git a/packages/common/src/realtime/types.ts b/packages/common/src/realtime/types.ts new file mode 100644 index 00000000..ad191d3a --- /dev/null +++ b/packages/common/src/realtime/types.ts @@ -0,0 +1,18 @@ +import type { TurnRunResult } from '../turnDaemon/types.js'; + +export type MessageTypeKey = 'public' | 'private' | 'national' | 'diplomacy'; + +export type RealtimeEvent = + | { + type: 'turnCompleted'; + at: string; + result: TurnRunResult; + } + | { + type: 'messageCreated'; + at: string; + mailbox: number; + msgType: MessageTypeKey; + messageId: number; + senderId: number; + };