fix(realtime): invalidate committed global activity logs

This commit is contained in:
2026-08-12 16:45:48 +00:00
parent 719c0c0e35
commit d1efa425eb
6 changed files with 286 additions and 13 deletions
@@ -0,0 +1,63 @@
import { randomUUID } from 'node:crypto';
import { buildGameEventChannel, createEmptyRealtimeReadModelChanges, type RealtimeEvent } from '@sammo-ts/common';
import { createRedisConnector, resolveRedisConfigFromEnv } from '@sammo-ts/infra';
import { describe, expect, it } from 'vitest';
import { RedisRealtimeEventHub } from '../src/realtime/eventHub.js';
import { publishRealtimeEvent } from '../src/realtime/publisher.js';
const liveDescribe = process.env.REDIS_URL ? describe : describe.skip;
liveDescribe('realtime event hub with live Redis', () => {
it('forwards committed global record, history, and month-boundary flags intact', async () => {
const publisher = createRedisConnector(resolveRedisConfigFromEnv());
const subscriber = createRedisConnector(resolveRedisConfigFromEnv());
await Promise.all([publisher.connect(), subscriber.connect()]);
const runId = process.env.CONDITIONAL_INTEGRATION_RUN_ID ?? randomUUID();
const profileName = `hwe:global-events-${runId}-${randomUUID()}`;
const hub = new RedisRealtimeEventHub(subscriber.client, buildGameEventChannel(profileName));
let unsubscribe = () => {};
try {
const received = new Promise<RealtimeEvent>((resolve) => {
unsubscribe = hub.subscribe(resolve);
});
await hub.start();
await publishRealtimeEvent(publisher.client, profileName, {
type: 'turnCompleted',
at: '2026-08-12T00:00:00.000Z',
lastTurnTime: '0185-02-01T00:00:00.000Z',
changes: {
...createEmptyRealtimeReadModelChanges(),
worldChanged: true,
globalRecordsChanged: true,
worldHistoryChanged: true,
},
revision: 12,
});
await expect(
Promise.race([
received,
new Promise<never>((_, reject) =>
setTimeout(() => reject(new Error('Redis realtime event timeout')), 2_000)
),
])
).resolves.toMatchObject({
type: 'turnCompleted',
changes: {
worldChanged: true,
globalRecordsChanged: true,
worldHistoryChanged: true,
},
revision: 12,
});
} finally {
unsubscribe();
await hub.stop();
await publisher.disconnect();
}
});
});
@@ -1,10 +1,12 @@
import { afterAll, beforeAll, describe, expect, it } from 'vitest'; import { afterAll, beforeAll, describe, expect, it } from 'vitest';
import { RANK_DATA_TYPES } from '@sammo-ts/common'; import { buildGameEventChannel, RANK_DATA_TYPES, type RealtimeEvent } from '@sammo-ts/common';
import type { GameSessionTokenPayload } from '@sammo-ts/common/auth/gameToken'; import type { GameSessionTokenPayload } from '@sammo-ts/common/auth/gameToken';
import { createTurnDaemonRuntime, seedScenarioToDatabase, type TurnDaemonRuntime } from '@sammo-ts/game-engine'; import { createTurnDaemonRuntime, seedScenarioToDatabase, type TurnDaemonRuntime } from '@sammo-ts/game-engine';
import { import {
createGamePostgresConnector, createGamePostgresConnector,
createRedisConnector,
resolveRedisConfigFromEnv,
type GamePrisma, type GamePrisma,
type GamePrismaClient, type GamePrismaClient,
type RedisConnector, type RedisConnector,
@@ -16,6 +18,7 @@ import { InMemoryBattleSimTransport } from '../src/battleSim/inMemoryTransport.j
import type { GameApiContext } from '../src/context.js'; import type { GameApiContext } from '../src/context.js';
import { DatabaseTurnDaemonTransport } from '../src/daemon/databaseTransport.js'; import { DatabaseTurnDaemonTransport } from '../src/daemon/databaseTransport.js';
import type { TurnDaemonTransport } from '../src/daemon/transport.js'; import type { TurnDaemonTransport } from '../src/daemon/transport.js';
import { RedisRealtimeEventHub } from '../src/realtime/eventHub.js';
import { appRouter } from '../src/router.js'; import { appRouter } from '../src/router.js';
const databaseUrl = process.env.SELECT_POOL_DATABASE_URL; const databaseUrl = process.env.SELECT_POOL_DATABASE_URL;
@@ -92,6 +95,22 @@ integration('scenario 903 select pool through the durable turn daemon', () => {
let daemonLoop: Promise<void> | undefined; let daemonLoop: Promise<void> | undefined;
let turnDaemon: TurnDaemonTransport; let turnDaemon: TurnDaemonTransport;
let worldStateId: number; let worldStateId: number;
let realtimeHub: RedisRealtimeEventHub | undefined;
let unsubscribeRealtime = () => {};
const realtimeEvents: RealtimeEvent[] = [];
const waitForRealtimeEvent = async (
predicate: (event: RealtimeEvent) => boolean,
timeoutMs = 2_000
): Promise<RealtimeEvent> => {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
const event = realtimeEvents.find(predicate);
if (event) return event;
await new Promise((resolve) => setTimeout(resolve, 20));
}
throw new Error('Timed out waiting for select-pool realtime event.');
};
const buildContext = (requestId: string, actorAuth: GameSessionTokenPayload = auth): GameApiContext => { const buildContext = (requestId: string, actorAuth: GameSessionTokenPayload = auth): GameApiContext => {
const redisClient = { const redisClient = {
@@ -148,6 +167,14 @@ integration('scenario 903 select pool through the durable turn daemon', () => {
await db.logEntry.deleteMany(); await db.logEntry.deleteMany();
worldStateId = (await db.worldState.findFirstOrThrow()).id; worldStateId = (await db.worldState.findFirstOrThrow()).id;
if (process.env.REDIS_URL) {
const realtimeSubscriber = createRedisConnector(resolveRedisConfigFromEnv());
await realtimeSubscriber.connect();
realtimeHub = new RedisRealtimeEventHub(realtimeSubscriber.client, buildGameEventChannel(profile));
unsubscribeRealtime = realtimeHub.subscribe((event) => realtimeEvents.push(event));
await realtimeHub.start();
}
runtime = await createTurnDaemonRuntime({ runtime = await createTurnDaemonRuntime({
profile, profile,
databaseUrl: databaseUrl!, databaseUrl: databaseUrl!,
@@ -168,6 +195,8 @@ integration('scenario 903 select pool through the durable turn daemon', () => {
await daemonLoop; await daemonLoop;
await runtime.close(); await runtime.close();
} }
unsubscribeRealtime();
await realtimeHub?.stop();
await closeDb?.(); await closeDb?.();
}, 30_000); }, 30_000);
@@ -267,6 +296,18 @@ integration('scenario 903 select pool through the durable turn daemon', () => {
where: { meta: { path: ['ownerUserId'], equals: userId } }, where: { meta: { path: ['ownerUserId'], equals: userId } },
}) })
).toBe(2); ).toBe(2);
if (realtimeHub) {
const creationEvent = await waitForRealtimeEvent(
(event) => event.type === 'readModelChanged' && event.changes.globalRecordsChanged
);
expect(creationEvent).toMatchObject({
type: 'readModelChanged',
changes: {
generalIds: [initial.id],
globalRecordsChanged: true,
},
});
}
await expect( await expect(
appRouter.createCaller(buildContext('select-pool-cooldown')).join.getSelectionPool() appRouter.createCaller(buildContext('select-pool-cooldown')).join.getSelectionPool()
+71 -6
View File
@@ -66,6 +66,13 @@ export interface RealtimeReadModelBaseline {
nations: Map<number, ReadModelSignatures>; nations: Map<number, ReadModelSignatures>;
} }
export type PersistedVisibleLogRow = {
id: number;
scope: LogScope;
category: LogCategory;
generalId: number | null;
};
const canonicalizeReadModelValue = (value: unknown): unknown => { const canonicalizeReadModelValue = (value: unknown): unknown => {
if (value instanceof Date) { if (value instanceof Date) {
return value.toISOString(); return value.toISOString();
@@ -321,6 +328,31 @@ export const summarizeRealtimeReadModelChanges = (
}; };
}; };
export const mergePersistedVisibleLogChanges = (
changes: RealtimeReadModelChanges,
rows: readonly PersistedVisibleLogRow[]
): RealtimeReadModelChanges => ({
...changes,
recordGeneralIds: uniqueSortedIds([
...changes.recordGeneralIds,
...rows.flatMap((entry) =>
entry.scope === LogScope.GENERAL && entry.category === LogCategory.ACTION && entry.generalId
? [entry.generalId]
: []
),
]),
globalRecordsChanged:
changes.globalRecordsChanged ||
rows.some(
(entry) =>
entry.scope === LogScope.SYSTEM &&
(entry.category === LogCategory.SUMMARY || entry.category === LogCategory.ACTION)
),
worldHistoryChanged:
changes.worldHistoryChanged ||
rows.some((entry) => entry.scope === LogScope.SYSTEM && entry.category === LogCategory.HISTORY),
});
export const excludeDeletedReservedTurnQueues = ( export const excludeDeletedReservedTurnQueues = (
changes: ReservedTurnChanges, changes: ReservedTurnChanges,
deletedGeneralIds: readonly number[], deletedGeneralIds: readonly number[],
@@ -871,10 +903,13 @@ export const createDatabaseTurnHooks = async (
const persistChanges = async ( const persistChanges = async (
transaction?: GamePrisma.TransactionClient, transaction?: GamePrisma.TransactionClient,
commandCompletion?: { requestId: string; result: TurnDaemonCommandResult } commandCompletion?: { requestId: string; result: TurnDaemonCommandResult },
directLogFloor?: number
): Promise<{ acknowledge: () => void; readModelChanges: RealtimeReadModelChanges }> => { ): Promise<{ acknowledge: () => void; readModelChanges: RealtimeReadModelChanges }> => {
const state = world.getState(); const state = world.getState();
const changes = world.peekDirtyState(); const changes = world.peekDirtyState();
let persistedVisibleLogs: PersistedVisibleLogRow[] = [];
let visibleLogFloor = directLogFloor;
const { const {
generals, generals,
cities, cities,
@@ -918,6 +953,13 @@ export const createDatabaseTurnHooks = async (
meta: asJson(state.meta), meta: asJson(state.meta),
}; };
const persist = async (prisma: GamePrisma.TransactionClient): Promise<void> => { const persist = async (prisma: GamePrisma.TransactionClient): Promise<void> => {
visibleLogFloor ??=
(
await prisma.logEntry.findFirst({
orderBy: { id: 'desc' },
select: { id: true },
})
)?.id ?? 0;
// Lock and validate the fencing row in the same transaction as every // Lock and validate the fencing row in the same transaction as every
// world mutation. A stale daemon can finish calculating, but it can // world mutation. A stale daemon can finish calculating, but it can
// never commit after another owner has advanced the epoch. // never commit after another owner has advanced the epoch.
@@ -1329,6 +1371,23 @@ export const createDatabaseTurnHooks = async (
}, },
}); });
} }
persistedVisibleLogs = await prisma.logEntry.findMany({
where: {
id: { gt: visibleLogFloor },
OR: [
{
scope: LogScope.GENERAL,
category: LogCategory.ACTION,
},
{
scope: LogScope.SYSTEM,
category: { in: [LogCategory.SUMMARY, LogCategory.ACTION, LogCategory.HISTORY] },
},
],
},
orderBy: { id: 'asc' },
select: { id: true, scope: true, category: true, generalId: true },
});
}; };
if (transaction) { if (transaction) {
await persist(transaction); await persist(transaction);
@@ -1339,10 +1398,9 @@ export const createDatabaseTurnHooks = async (
); );
} }
const readModelChanges = summarizeRealtimeReadModelChanges( const readModelChanges = mergePersistedVisibleLogChanges(
changes, summarizeRealtimeReadModelChanges(changes, persistedReservedTurnChanges, readModelBaseline),
persistedReservedTurnChanges, persistedVisibleLogs
readModelBaseline
); );
return { return {
acknowledge: () => { acknowledge: () => {
@@ -1370,8 +1428,15 @@ export const createDatabaseTurnHooks = async (
executeCommand: async (requestId, execute) => { executeCommand: async (requestId, execute) => {
const committed = await prisma.$transaction( const committed = await prisma.$transaction(
async (transaction) => { async (transaction) => {
const directLogFloor =
(
await transaction.logEntry.findFirst({
orderBy: { id: 'desc' },
select: { id: true },
})
)?.id ?? 0;
const result = await execute({ db: transaction }); const result = await execute({ db: transaction });
const persisted = await persistChanges(transaction, { requestId, result }); const persisted = await persistChanges(transaction, { requestId, result }, directLogFloor);
return { result, persisted }; return { result, persisted };
}, },
options?.transactionTimeoutMs ? { timeout: options.transactionTimeoutMs } : undefined options?.transactionTimeoutMs ? { timeout: options.transactionTimeoutMs } : undefined
@@ -271,6 +271,7 @@ integration('monthly nation betting persistence', () => {
durationMs: 0, durationMs: 0,
partial: false, partial: false,
}); });
expect(hooks.takeCommittedReadModelChanges()?.worldHistoryChanged).toBe(true);
expect(await db.nationBetting.findUniqueOrThrow({ where: { id: bettingId } })).toMatchObject({ expect(await db.nationBetting.findUniqueOrThrow({ where: { id: bettingId } })).toMatchObject({
name: '천통국 예상', name: '천통국 예상',
@@ -320,6 +321,7 @@ integration('monthly nation betting persistence', () => {
}) })
).toBe(true); ).toBe(true);
await world.advanceMonth(new Date('0200-02-01T00:00:00.000Z')); await world.advanceMonth(new Date('0200-02-01T00:00:00.000Z'));
expect(world.peekDirtyState().logs).toEqual([]);
await hooks.hooks.flushChanges?.({ await hooks.hooks.flushChanges?.({
lastTurnTime: state.lastTurnTime.toISOString(), lastTurnTime: state.lastTurnTime.toISOString(),
processedGenerals: 0, processedGenerals: 0,
@@ -327,6 +329,7 @@ integration('monthly nation betting persistence', () => {
durationMs: 0, durationMs: 0,
partial: false, partial: false,
}); });
expect(hooks.takeCommittedReadModelChanges()?.worldHistoryChanged).toBe(true);
expect(await db.nationBetting.findUniqueOrThrow({ where: { id: bettingId } })).toMatchObject({ expect(await db.nationBetting.findUniqueOrThrow({ where: { id: bettingId } })).toMatchObject({
finished: true, finished: true,
@@ -1,9 +1,11 @@
import { describe, expect, it } from 'vitest'; import { describe, expect, it } from 'vitest';
import { createEmptyRealtimeReadModelChanges } from '@sammo-ts/common';
import { LogCategory, LogFormat, LogScope } from '@sammo-ts/logic'; import { LogCategory, LogFormat, LogScope } from '@sammo-ts/logic';
import { import {
applyRealtimeReadModelBaseline, applyRealtimeReadModelBaseline,
createRealtimeReadModelBaseline, createRealtimeReadModelBaseline,
mergePersistedVisibleLogChanges,
summarizeRealtimeReadModelChanges, summarizeRealtimeReadModelChanges,
} from '../src/turn/databaseHooks.js'; } from '../src/turn/databaseHooks.js';
import type { InMemoryTurnWorld } from '../src/turn/inMemoryWorld.js'; import type { InMemoryTurnWorld } from '../src/turn/inMemoryWorld.js';
@@ -11,6 +13,35 @@ import type { TurnWorldChanges } from '../src/turn/inMemoryWorld.js';
import type { ReservedTurnChanges } from '../src/turn/reservedTurnStore.js'; import type { ReservedTurnChanges } from '../src/turn/reservedTurnStore.js';
describe('summarizeRealtimeReadModelChanges', () => { describe('summarizeRealtimeReadModelChanges', () => {
it('classifies the committed log rows even when they bypass in-memory log drafts', () => {
expect(
mergePersistedVisibleLogChanges(createEmptyRealtimeReadModelChanges(), [
{
id: 10,
scope: LogScope.GENERAL,
category: LogCategory.ACTION,
generalId: 7,
},
{
id: 11,
scope: LogScope.SYSTEM,
category: LogCategory.SUMMARY,
generalId: null,
},
{
id: 12,
scope: LogScope.SYSTEM,
category: LogCategory.HISTORY,
generalId: null,
},
])
).toMatchObject({
recordGeneralIds: [7],
globalRecordsChanged: true,
worldHistoryChanged: true,
});
});
it('emits deterministic entity and record invalidations from committed changes', () => { it('emits deterministic entity and record invalidations from committed changes', () => {
const worldChanges = { const worldChanges = {
generals: [{ id: 9 }, { id: 7 }], generals: [{ id: 9 }, { id: 7 }],
+76 -6
View File
@@ -31,6 +31,11 @@ type NavigationFixture = {
forceSnapshotCalls?: number; forceSnapshotCalls?: number;
refreshDelayMs?: number; refreshDelayMs?: number;
largeCommandTable?: boolean; largeCommandTable?: boolean;
currentYear?: number;
currentMonth?: number;
globalRecords?: Array<{ id: number; text: string }>;
generalRecords?: Array<{ id: number; text: string }>;
worldHistory?: Array<{ id: number; text: string }>;
reservedTurns?: Array<{ index: number; action: string; args: Record<string, unknown> }>; reservedTurns?: Array<{ index: number; action: string; args: Record<string, unknown> }>;
messages?: unknown; messages?: unknown;
messageContacts?: unknown; messageContacts?: unknown;
@@ -318,7 +323,12 @@ const installFixture = async (page: Page, state: NavigationFixture) => {
const results = operations.map((operation, index) => { const results = operations.map((operation, index) => {
if (operation === 'auth.status') return response({ ok: true }); if (operation === 'auth.status') return response({ ok: true });
if (operation === 'lobby.info') { if (operation === 'lobby.info') {
return response({ myGeneral: { id: 7, name: '메뉴검증장수' }, year: 185, month: 1, turnTerm: 10 }); return response({
myGeneral: { id: 7, name: '메뉴검증장수' },
year: state.currentYear ?? 185,
month: state.currentMonth ?? 1,
turnTerm: 10,
});
} }
if (operation === 'dashboard.getContextBundleDelta') { if (operation === 'dashboard.getContextBundleDelta') {
state.generalMeCalls += 1; state.generalMeCalls += 1;
@@ -394,8 +404,8 @@ const installFixture = async (page: Page, state: NavigationFixture) => {
result: true, result: true,
version: 0, version: 0,
startYear: 180, startYear: 180,
year: 185, year: state.currentYear ?? 185,
month: 1, month: state.currentMonth ?? 1,
cityList: [[1, 8, state.cityState ?? 0, 1, 1, 1]], cityList: [[1, 8, state.cityState ?? 0, 1, 1, 1]],
nationList: [[1, '위', '#008000', 1]], nationList: [[1, '위', '#008000', 1]],
spyList: {}, spyList: {},
@@ -418,9 +428,9 @@ const installFixture = async (page: Page, state: NavigationFixture) => {
if (operation === 'messages.getContacts') return response(state.messageContacts ?? { nation: [] }); if (operation === 'messages.getContacts') return response(state.messageContacts ?? { nation: [] });
if (operation === 'general.getRecentRecords') { if (operation === 'general.getRecentRecords') {
return response({ return response({
global: [{ id: 3, text: '장수 동향 기록' }], global: state.globalRecords ?? [{ id: 3, text: '장수 동향 기록' }],
general: [{ id: 2, text: '개인 기록' }], general: state.generalRecords ?? [{ id: 2, text: '개인 기록' }],
history: [{ id: 1, text: '중원 정세 기록' }], history: state.worldHistory ?? [{ id: 1, text: '중원 정세 기록' }],
}); });
} }
if (operation === 'general.getFrontStatus') { if (operation === 'general.getFrontStatus') {
@@ -1562,6 +1572,66 @@ test('realtime read-model events skip clock-only work, merge bursts, patch in pl
expect(state.generalMeCalls).toBe(callsAfterLeavingMain); expect(state.generalMeCalls).toBe(callsAfterLeavingMain);
}); });
test('global activity, world history, and a month boundary refresh their visible main slices', async ({ page }) => {
const state: NavigationFixture = {
officerLevel: 5,
permission: 2,
nationLevel: 3,
stage: 0,
npcMode: 1,
generalMeCalls: 0,
operations: [],
};
await installRealtimeHarness(page);
await installFixture(page, state);
await page.setViewportSize({ width: 1200, height: 900 });
await waitForMain(page);
await expect
.poll(() =>
page.evaluate(() => (window as unknown as { __hasMainRealtime: () => boolean }).__hasMainRealtime())
)
.toBe(true);
state.globalRecords = [
{ id: 4, text: '자동 갱신된 장수 동향' },
{ id: 3, text: '장수 동향 기록' },
];
const operationsBeforeGlobal = state.operations.length;
await emitReadModelChanges(page, readModelChanges({ globalRecordsChanged: true }));
await expect(page.locator('[data-main-target="global-records"]')).toContainText('자동 갱신된 장수 동향');
expect(state.operations.slice(operationsBeforeGlobal)).toEqual(['general.getRecentRecords']);
state.worldHistory = [
{ id: 5, text: '자동 갱신된 중원 정세' },
{ id: 1, text: '중원 정세 기록' },
];
const operationsBeforeHistory = state.operations.length;
await emitReadModelChanges(page, readModelChanges({ worldHistoryChanged: true }));
await expect(page.locator('[data-main-target="world-history"]')).toContainText('자동 갱신된 중원 정세');
expect(state.operations.slice(operationsBeforeHistory)).toEqual(['general.getRecentRecords']);
state.currentMonth = 2;
const operationsBeforeMonth = state.operations.length;
await page.evaluate(
(changes) => {
(window as unknown as { __emitMainRealtime: (type: string, payload: unknown) => void }).__emitMainRealtime(
'turnCompleted',
{
at: new Date().toISOString(),
lastTurnTime: '0185-02-01T00:00:00.000Z',
changes,
}
);
},
readModelChanges({ worldChanged: true })
);
await expect(page.getByText('현재: 185년 2월')).toBeVisible();
await expect(page.locator('.map-viewer')).toContainText('185年 2月');
expect(state.operations.slice(operationsBeforeMonth).sort()).toEqual(
['dashboard.getContextBundleDelta', 'lobby.info', 'world.getMap'].sort()
);
});
test('same-account main tabs share one realtime diff and exclude a tab while sync is off', async ({ test('same-account main tabs share one realtime diff and exclude a tab while sync is off', async ({
context, context,
page, page,