diff --git a/app/game-api/src/turns/reservedTurns.ts b/app/game-api/src/turns/reservedTurns.ts index d6abe7a..ffb4439 100644 --- a/app/game-api/src/turns/reservedTurns.ts +++ b/app/game-api/src/turns/reservedTurns.ts @@ -229,16 +229,25 @@ const claimGeneralRevision = async ( expectedRevision: number ): Promise => { const nextRevision = expectedRevision + 1; - const claimed = - expectedRevision === 0 - ? await db.generalTurnRevision.createMany({ - data: [{ generalId, revision: nextRevision }], - skipDuplicates: true, - }) - : await db.generalTurnRevision.updateMany({ - where: { generalId, revision: expectedRevision }, - data: { revision: nextRevision }, - }); + const now = new Date(); + let claimed = await db.generalTurnRevision.updateMany({ + where: { + generalId, + revision: expectedRevision, + OR: [{ leaseOwner: null }, { leaseExpiresAt: { lte: now } }], + }, + data: { + revision: nextRevision, + leaseOwner: null, + leaseExpiresAt: null, + }, + }); + if (claimed.count === 0 && expectedRevision === 0) { + claimed = await db.generalTurnRevision.createMany({ + data: [{ generalId, revision: nextRevision }], + skipDuplicates: true, + }); + } if (claimed.count === 1) { return nextRevision; } @@ -253,16 +262,26 @@ const claimNationRevision = async ( expectedRevision: number ): Promise => { const nextRevision = expectedRevision + 1; - const claimed = - expectedRevision === 0 - ? await db.nationTurnRevision.createMany({ - data: [{ nationId, officerLevel, revision: nextRevision }], - skipDuplicates: true, - }) - : await db.nationTurnRevision.updateMany({ - where: { nationId, officerLevel, revision: expectedRevision }, - data: { revision: nextRevision }, - }); + const now = new Date(); + let claimed = await db.nationTurnRevision.updateMany({ + where: { + nationId, + officerLevel, + revision: expectedRevision, + OR: [{ leaseOwner: null }, { leaseExpiresAt: { lte: now } }], + }, + data: { + revision: nextRevision, + leaseOwner: null, + leaseExpiresAt: null, + }, + }); + if (claimed.count === 0 && expectedRevision === 0) { + claimed = await db.nationTurnRevision.createMany({ + data: [{ nationId, officerLevel, revision: nextRevision }], + skipDuplicates: true, + }); + } if (claimed.count === 1) { return nextRevision; } diff --git a/app/game-api/test/reservedTurnRevision.integration.test.ts b/app/game-api/test/reservedTurnRevision.integration.test.ts index a7a7336..06398a2 100644 --- a/app/game-api/test/reservedTurnRevision.integration.test.ts +++ b/app/game-api/test/reservedTurnRevision.integration.test.ts @@ -11,6 +11,7 @@ import { const databaseUrl = process.env.RESERVED_TURN_DATABASE_URL; const describeIntegration = databaseUrl ? describe : describe.skip; const GENERAL_ID = 2_147_400_001; +const REVISION_ZERO_GENERAL_ID = 2_147_400_002; describeIntegration('reserved turn queue revision integration', () => { const connector = databaseUrl ? createGamePostgresConnector({ url: databaseUrl }) : null; @@ -20,16 +21,18 @@ describeIntegration('reserved turn queue revision integration', () => { return; } await connector.connect(); - await connector.prisma.generalTurn.deleteMany({ where: { generalId: GENERAL_ID } }); - await connector.prisma.generalTurnRevision.deleteMany({ where: { generalId: GENERAL_ID } }); + const generalIds = [GENERAL_ID, REVISION_ZERO_GENERAL_ID]; + await connector.prisma.generalTurn.deleteMany({ where: { generalId: { in: generalIds } } }); + await connector.prisma.generalTurnRevision.deleteMany({ where: { generalId: { in: generalIds } } }); }); afterAll(async () => { if (!connector) { return; } - await connector.prisma.generalTurn.deleteMany({ where: { generalId: GENERAL_ID } }); - await connector.prisma.generalTurnRevision.deleteMany({ where: { generalId: GENERAL_ID } }); + const generalIds = [GENERAL_ID, REVISION_ZERO_GENERAL_ID]; + await connector.prisma.generalTurn.deleteMany({ where: { generalId: { in: generalIds } } }); + await connector.prisma.generalTurnRevision.deleteMany({ where: { generalId: { in: generalIds } } }); await connector.disconnect(); }); @@ -60,4 +63,46 @@ describeIntegration('reserved turn queue revision integration', () => { expect(['che_훈련', 'che_사기진작']).toContain(snapshot.turns[0]?.action); expect(await connector.prisma.generalTurn.count({ where: { generalId: GENERAL_ID } })).toBe(30); }); + + it('uses an unlocked revision-zero lease row as the first API revision', async () => { + if (!connector) { + throw new Error('integration connector is unavailable'); + } + + await connector.prisma.generalTurnRevision.create({ + data: { + generalId: REVISION_ZERO_GENERAL_ID, + revision: 0, + leaseOwner: null, + leaseExpiresAt: null, + }, + }); + + const result = await connector.prisma.$transaction((transaction) => + setGeneralTurn( + transaction as unknown as DatabaseClient, + REVISION_ZERO_GENERAL_ID, + 0, + 'che_훈련', + {}, + 0 + ) + ); + + expect(result.revision).toBe(1); + expect( + await connector.prisma.generalTurn.count({ + where: { generalId: REVISION_ZERO_GENERAL_ID }, + }) + ).toBe(30); + expect( + await connector.prisma.generalTurnRevision.findUnique({ + where: { generalId: REVISION_ZERO_GENERAL_ID }, + }) + ).toMatchObject({ + revision: 1, + leaseOwner: null, + leaseExpiresAt: null, + }); + }); }); diff --git a/app/game-api/test/reservedTurns.test.ts b/app/game-api/test/reservedTurns.test.ts index 869b1c6..cce2678 100644 --- a/app/game-api/test/reservedTurns.test.ts +++ b/app/game-api/test/reservedTurns.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it } from 'vitest'; +import { describe, expect, it, vi } from 'vitest'; import type { DatabaseClient, GeneralTurnRow, NationTurnRow } from '../src/context.js'; import { @@ -338,4 +338,33 @@ describe('reservedTurns', () => { const noOpPush = await shiftNationTurns(db, 5, 12, 12, repeated.revision); expect(noOpPush).toEqual(repeated); }); + + it('rejects an API writer while the daemon holds the queue lease without touching turns', async () => { + const deleteMany = vi.fn(async () => ({})); + const createMany = vi.fn(async () => ({})); + const db = { + generalTurnRevision: { + updateMany: vi.fn(async () => ({ count: 0 })), + createMany: vi.fn(async () => ({ count: 0 })), + findUnique: vi.fn(async () => ({ + generalId: 9, + revision: 0, + leaseOwner: 'daemon-1', + leaseExpiresAt: new Date(Date.now() + 60_000), + updatedAt: new Date(), + })), + }, + generalTurn: { + findMany: vi.fn(async () => []), + deleteMany, + createMany, + }, + } as unknown as DatabaseClient; + + await expect(setGeneralTurn(db, 9, 0, 'che_훈련', {}, 0)).rejects.toBeInstanceOf( + ReservedTurnRevisionConflictError + ); + expect(deleteMany).not.toHaveBeenCalled(); + expect(createMany).not.toHaveBeenCalled(); + }); }); diff --git a/app/game-engine/src/turn/reservedTurnStore.ts b/app/game-engine/src/turn/reservedTurnStore.ts index 9390b1b..9efcb94 100644 --- a/app/game-engine/src/turn/reservedTurnStore.ts +++ b/app/game-engine/src/turn/reservedTurnStore.ts @@ -1,5 +1,6 @@ import { createGamePostgresConnector, type InputJsonValue, type TurnEngineDatabaseClient } from '@sammo-ts/infra'; import { isRecord } from '@sammo-ts/common'; +import { randomUUID } from 'node:crypto'; export interface ReservedTurnEntry { action: string; @@ -10,6 +11,8 @@ export interface ReservedTurnStoreOptions { databaseUrl: string; maxGeneralTurns?: number; maxNationTurns?: number; + leaseOwner?: string; + leaseDurationMs?: number; } export interface ReservedTurnStoreHandle { @@ -20,6 +23,7 @@ export interface ReservedTurnStoreHandle { const DEFAULT_TURN_ACTION = '휴식'; const DEFAULT_GENERAL_TURNS = 30; const DEFAULT_NATION_TURNS = 12; +const DEFAULT_LEASE_DURATION_MS = 5 * 60_000; const asJson = (value: unknown): InputJsonValue => value as InputJsonValue; @@ -71,15 +75,30 @@ const buildTurnListFromRows = ( const buildNationKey = (nationId: number, officerLevel: number): string => `${nationId}:${officerLevel}`; type ReservedTurnDatabaseClient = Pick & { - generalTurnRevision?: Pick, 'upsert'>; - nationTurnRevision?: Pick, 'upsert'>; + generalTurnRevision?: Pick< + NonNullable, + 'findUnique' | 'createMany' | 'updateMany' + >; + nationTurnRevision?: Pick< + NonNullable, + 'findUnique' | 'createMany' | 'updateMany' + >; }; export interface ReservedTurnChanges { generalIds: number[]; generalInitializationIds: number[]; + generalLeaseIds: number[]; nationKeys: string[]; nationInitializationKeys: string[]; + nationLeaseKeys: string[]; +} + +export class ReservedTurnLeaseConflictError extends Error { + constructor(readonly queueKey: string) { + super(`Reserved turn queue lease conflict: ${queueKey}.`); + this.name = 'ReservedTurnLeaseConflictError'; + } } export class InMemoryReservedTurnStore { @@ -89,15 +108,26 @@ export class InMemoryReservedTurnStore { private readonly dirtyNationKeys = new Set(); private readonly pendingGeneralInitializationIds = new Set(); private readonly pendingNationInitializationKeys = new Set(); + private readonly leasedGeneralIds = new Set(); + private readonly leasedNationKeys = new Set(); private readonly maxGeneralTurns: number; private readonly maxNationTurns: number; + private readonly leaseOwner: string; + private readonly leaseDurationMs: number; constructor( private readonly prisma: ReservedTurnDatabaseClient, - options: { maxGeneralTurns: number; maxNationTurns: number } + options: { + maxGeneralTurns: number; + maxNationTurns: number; + leaseOwner?: string; + leaseDurationMs?: number; + } ) { this.maxGeneralTurns = options.maxGeneralTurns; this.maxNationTurns = options.maxNationTurns; + this.leaseOwner = options.leaseOwner ?? randomUUID(); + this.leaseDurationMs = options.leaseDurationMs ?? DEFAULT_LEASE_DURATION_MS; } async loadAll(): Promise { @@ -134,6 +164,153 @@ export class InMemoryReservedTurnStore { } } + private getLeaseExpiresAt(): Date { + return new Date(Date.now() + this.leaseDurationMs); + } + + private async acquireGeneralLease(generalId: number): Promise { + const revisionStore = this.prisma.generalTurnRevision; + if (!revisionStore) { + return; + } + const now = new Date(); + const leaseExpiresAt = this.getLeaseExpiresAt(); + let claimed = await revisionStore.updateMany({ + where: { + generalId, + OR: [ + { leaseOwner: this.leaseOwner }, + { leaseOwner: null }, + { leaseExpiresAt: { lte: now } }, + ], + }, + data: { + leaseOwner: this.leaseOwner, + leaseExpiresAt, + }, + }); + if (claimed.count === 0) { + claimed = await revisionStore.createMany({ + data: [{ generalId, revision: 0, leaseOwner: this.leaseOwner, leaseExpiresAt }], + skipDuplicates: true, + }); + if (claimed.count === 0) { + claimed = await revisionStore.updateMany({ + where: { + generalId, + OR: [ + { leaseOwner: this.leaseOwner }, + { leaseOwner: null }, + { leaseExpiresAt: { lte: now } }, + ], + }, + data: { + leaseOwner: this.leaseOwner, + leaseExpiresAt, + }, + }); + } + } + if (claimed.count !== 1) { + throw new ReservedTurnLeaseConflictError(`general:${generalId}`); + } + this.leasedGeneralIds.add(generalId); + } + + private async acquireNationLease(nationId: number, officerLevel: number): Promise { + const revisionStore = this.prisma.nationTurnRevision; + if (!revisionStore) { + return; + } + const now = new Date(); + const leaseExpiresAt = this.getLeaseExpiresAt(); + let claimed = await revisionStore.updateMany({ + where: { + nationId, + officerLevel, + OR: [ + { leaseOwner: this.leaseOwner }, + { leaseOwner: null }, + { leaseExpiresAt: { lte: now } }, + ], + }, + data: { + leaseOwner: this.leaseOwner, + leaseExpiresAt, + }, + }); + if (claimed.count === 0) { + claimed = await revisionStore.createMany({ + data: [{ nationId, officerLevel, revision: 0, leaseOwner: this.leaseOwner, leaseExpiresAt }], + skipDuplicates: true, + }); + if (claimed.count === 0) { + claimed = await revisionStore.updateMany({ + where: { + nationId, + officerLevel, + OR: [ + { leaseOwner: this.leaseOwner }, + { leaseOwner: null }, + { leaseExpiresAt: { lte: now } }, + ], + }, + data: { + leaseOwner: this.leaseOwner, + leaseExpiresAt, + }, + }); + } + } + if (claimed.count !== 1) { + throw new ReservedTurnLeaseConflictError(`nation:${nationId}:${officerLevel}`); + } + this.leasedNationKeys.add(buildNationKey(nationId, officerLevel)); + } + + private async releaseGeneralLease(generalId: number): Promise { + await this.prisma.generalTurnRevision?.updateMany({ + where: { generalId, leaseOwner: this.leaseOwner }, + data: { leaseOwner: null, leaseExpiresAt: null }, + }); + this.leasedGeneralIds.delete(generalId); + } + + private async releaseNationLease(nationId: number, officerLevel: number): Promise { + await this.prisma.nationTurnRevision?.updateMany({ + where: { nationId, officerLevel, leaseOwner: this.leaseOwner }, + data: { leaseOwner: null, leaseExpiresAt: null }, + }); + this.leasedNationKeys.delete(buildNationKey(nationId, officerLevel)); + } + + async prepareTurnsForExecution( + generalId: number, + nation?: { nationId: number; officerLevel: number } + ): Promise { + const hadGeneralLease = this.leasedGeneralIds.has(generalId); + const nationKey = nation ? buildNationKey(nation.nationId, nation.officerLevel) : null; + const hadNationLease = nationKey ? this.leasedNationKeys.has(nationKey) : false; + try { + await this.acquireGeneralLease(generalId); + if (nation) { + await this.acquireNationLease(nation.nationId, nation.officerLevel); + } + await Promise.all([ + this.refreshGeneralTurns(generalId), + nation ? this.refreshNationTurns(nation.nationId, nation.officerLevel) : Promise.resolve(), + ]); + } catch (error) { + if (nation && nationKey !== null && !hadNationLease && this.leasedNationKeys.has(nationKey)) { + await this.releaseNationLease(nation.nationId, nation.officerLevel); + } + if (!hadGeneralLease && this.leasedGeneralIds.has(generalId)) { + await this.releaseGeneralLease(generalId); + } + throw error; + } + } + async refreshGeneralTurns(generalId: number): Promise { if (this.dirtyGeneralIds.has(generalId)) { return; @@ -252,8 +429,10 @@ export class InMemoryReservedTurnStore { return { generalIds: Array.from(this.dirtyGeneralIds), generalInitializationIds: Array.from(this.pendingGeneralInitializationIds), + generalLeaseIds: Array.from(this.leasedGeneralIds), nationKeys: Array.from(this.dirtyNationKeys), nationInitializationKeys: Array.from(this.pendingNationInitializationKeys), + nationLeaseKeys: Array.from(this.leasedNationKeys), }; } @@ -264,17 +443,172 @@ export class InMemoryReservedTurnStore { for (const generalId of changes.generalInitializationIds) { this.pendingGeneralInitializationIds.delete(generalId); } + for (const generalId of changes.generalLeaseIds) { + this.leasedGeneralIds.delete(generalId); + } for (const key of changes.nationKeys) { this.dirtyNationKeys.delete(key); } for (const key of changes.nationInitializationKeys) { this.pendingNationInitializationKeys.delete(key); } + for (const key of changes.nationLeaseKeys) { + this.leasedNationKeys.delete(key); + } + } + + private async claimGeneralFlushLease( + prisma: ReservedTurnDatabaseClient, + generalId: number + ): Promise { + const revisionStore = prisma.generalTurnRevision; + if (!revisionStore) { + return false; + } + const leaseExpiresAt = this.getLeaseExpiresAt(); + const where = this.leasedGeneralIds.has(generalId) + ? { generalId, leaseOwner: this.leaseOwner } + : { + generalId, + OR: [ + { leaseOwner: this.leaseOwner }, + { leaseOwner: null }, + { leaseExpiresAt: { lte: new Date() } }, + ], + }; + let claimed = await revisionStore.updateMany({ + where, + data: { + leaseOwner: this.leaseOwner, + leaseExpiresAt, + }, + }); + if (claimed.count === 0 && !this.leasedGeneralIds.has(generalId)) { + claimed = await revisionStore.createMany({ + data: [{ generalId, revision: 0, leaseOwner: this.leaseOwner, leaseExpiresAt }], + skipDuplicates: true, + }); + if (claimed.count === 0) { + claimed = await revisionStore.updateMany({ + where, + data: { + leaseOwner: this.leaseOwner, + leaseExpiresAt, + }, + }); + } + } + if (claimed.count !== 1) { + throw new ReservedTurnLeaseConflictError(`general:${generalId}`); + } + return true; + } + + private async finalizeGeneralFlush( + prisma: ReservedTurnDatabaseClient, + generalId: number, + claimedLease: boolean + ): Promise { + const revisionStore = prisma.generalTurnRevision; + if (!revisionStore) { + return; + } + if (!claimedLease) { + throw new ReservedTurnLeaseConflictError(`general:${generalId}`); + } + const finalized = await revisionStore.updateMany({ + where: { generalId, leaseOwner: this.leaseOwner }, + data: { + revision: { increment: 1 }, + leaseOwner: null, + leaseExpiresAt: null, + }, + }); + if (finalized.count !== 1) { + throw new ReservedTurnLeaseConflictError(`general:${generalId}`); + } + } + + private async claimNationFlushLease( + prisma: ReservedTurnDatabaseClient, + nationId: number, + officerLevel: number + ): Promise { + const revisionStore = prisma.nationTurnRevision; + const key = buildNationKey(nationId, officerLevel); + if (!revisionStore) { + return false; + } + const leaseExpiresAt = this.getLeaseExpiresAt(); + const where = this.leasedNationKeys.has(key) + ? { nationId, officerLevel, leaseOwner: this.leaseOwner } + : { + nationId, + officerLevel, + OR: [ + { leaseOwner: this.leaseOwner }, + { leaseOwner: null }, + { leaseExpiresAt: { lte: new Date() } }, + ], + }; + let claimed = await revisionStore.updateMany({ + where, + data: { + leaseOwner: this.leaseOwner, + leaseExpiresAt, + }, + }); + if (claimed.count === 0 && !this.leasedNationKeys.has(key)) { + claimed = await revisionStore.createMany({ + data: [{ nationId, officerLevel, revision: 0, leaseOwner: this.leaseOwner, leaseExpiresAt }], + skipDuplicates: true, + }); + if (claimed.count === 0) { + claimed = await revisionStore.updateMany({ + where, + data: { + leaseOwner: this.leaseOwner, + leaseExpiresAt, + }, + }); + } + } + if (claimed.count !== 1) { + throw new ReservedTurnLeaseConflictError(`nation:${nationId}:${officerLevel}`); + } + return true; + } + + private async finalizeNationFlush( + prisma: ReservedTurnDatabaseClient, + nationId: number, + officerLevel: number, + claimedLease: boolean + ): Promise { + const revisionStore = prisma.nationTurnRevision; + if (!revisionStore) { + return; + } + if (!claimedLease) { + throw new ReservedTurnLeaseConflictError(`nation:${nationId}:${officerLevel}`); + } + const finalized = await revisionStore.updateMany({ + where: { nationId, officerLevel, leaseOwner: this.leaseOwner }, + data: { + revision: { increment: 1 }, + leaseOwner: null, + leaseExpiresAt: null, + }, + }); + if (finalized.count !== 1) { + throw new ReservedTurnLeaseConflictError(`nation:${nationId}:${officerLevel}`); + } } async persistChanges(prisma: ReservedTurnDatabaseClient, changes: ReservedTurnChanges): Promise { for (const generalId of changes.generalIds) { const turns = this.getGeneralTurns(generalId); + const claimedLease = await this.claimGeneralFlushLease(prisma, generalId); await prisma.generalTurn.deleteMany({ where: { generalId } }); await prisma.generalTurn.createMany({ data: turns.map((entry, turnIdx) => ({ @@ -284,13 +618,7 @@ export class InMemoryReservedTurnStore { arg: asJson(normalizeArgs(entry.args)), })), }); - if (prisma.generalTurnRevision) { - await prisma.generalTurnRevision.upsert({ - where: { generalId }, - create: { generalId, revision: 1 }, - update: { revision: { increment: 1 } }, - }); - } + await this.finalizeGeneralFlush(prisma, generalId, claimedLease); } for (const generalId of changes.generalInitializationIds) { @@ -314,6 +642,7 @@ export class InMemoryReservedTurnStore { const nationId = Number(nationIdRaw); const officerLevel = Number(officerLevelRaw); const turns = this.getNationTurns(nationId, officerLevel); + const claimedLease = await this.claimNationFlushLease(prisma, nationId, officerLevel); await prisma.nationTurn.deleteMany({ where: { nationId, officerLevel }, }); @@ -326,18 +655,7 @@ export class InMemoryReservedTurnStore { arg: asJson(normalizeArgs(entry.args)), })), }); - if (prisma.nationTurnRevision) { - await prisma.nationTurnRevision.upsert({ - where: { - nationId_officerLevel: { - nationId, - officerLevel, - }, - }, - create: { nationId, officerLevel, revision: 1 }, - update: { revision: { increment: 1 } }, - }); - } + await this.finalizeNationFlush(prisma, nationId, officerLevel, claimedLease); } for (const key of changes.nationInitializationKeys) { @@ -359,6 +677,31 @@ export class InMemoryReservedTurnStore { skipDuplicates: true, }); } + + for (const generalId of changes.generalLeaseIds) { + if (changes.generalIds.includes(generalId)) { + continue; + } + await prisma.generalTurnRevision?.updateMany({ + where: { generalId, leaseOwner: this.leaseOwner }, + data: { leaseOwner: null, leaseExpiresAt: null }, + }); + } + + for (const key of changes.nationLeaseKeys) { + if (changes.nationKeys.includes(key)) { + continue; + } + const [nationIdRaw, officerLevelRaw] = key.split(':'); + await prisma.nationTurnRevision?.updateMany({ + where: { + nationId: Number(nationIdRaw), + officerLevel: Number(officerLevelRaw), + leaseOwner: this.leaseOwner, + }, + data: { leaseOwner: null, leaseExpiresAt: null }, + }); + } } async flushChanges(): Promise { @@ -374,6 +717,8 @@ export const createReservedTurnStore = async (options: ReservedTurnStoreOptions) const store = new InMemoryReservedTurnStore(connector.prisma, { maxGeneralTurns: options.maxGeneralTurns ?? DEFAULT_GENERAL_TURNS, maxNationTurns: options.maxNationTurns ?? DEFAULT_NATION_TURNS, + leaseOwner: options.leaseOwner, + leaseDurationMs: options.leaseDurationMs, }); await store.loadAll(); return { diff --git a/app/game-engine/src/turn/turnDaemon.ts b/app/game-engine/src/turn/turnDaemon.ts index 2362c52..7dc898b 100644 --- a/app/game-engine/src/turn/turnDaemon.ts +++ b/app/game-engine/src/turn/turnDaemon.ts @@ -503,13 +503,18 @@ const createTurnDaemonRuntimeWithLease = async ( beforeExecuteGeneral: reservedTurnStoreHandle ? async (general) => { const promises: Promise[] = []; - promises.push(reservedTurnStoreHandle.store.refreshGeneralTurns(general.id)); + promises.push( + reservedTurnStoreHandle.store.prepareTurnsForExecution( + general.id, + general.nationId > 0 && general.officerLevel >= 5 + ? { + nationId: general.nationId, + officerLevel: general.officerLevel, + } + : undefined + ) + ); promises.push(refreshOccupiedAuctionUniqueItemKeys()); - if (general.nationId > 0 && general.officerLevel >= 5) { - promises.push( - reservedTurnStoreHandle.store.refreshNationTurns(general.nationId, general.officerLevel) - ); - } if (general.nationId > 0 && general.officerLevel >= 5 && shouldUseAi(general, world.getState())) { const key = `${general.nationId}:${world.getState().currentYear}:${world.getState().currentMonth}`; if (!prefetchedNationTurns.has(key)) { diff --git a/app/game-engine/test/inputEventAtomicity.test.ts b/app/game-engine/test/inputEventAtomicity.test.ts index 2dcca08..acd6345 100644 --- a/app/game-engine/test/inputEventAtomicity.test.ts +++ b/app/game-engine/test/inputEventAtomicity.test.ts @@ -1,4 +1,5 @@ import { describe, expect, it, vi } from 'vitest'; +import { asRecord } from '@sammo-ts/common'; import { InMemoryControlQueue, @@ -27,7 +28,34 @@ const processor: TurnProcessor = { describe('input event atomicity', () => { it('keeps reserved-turn dirty state when persistence fails', async () => { let failCreate = true; - const revisionUpsert = vi.fn(async () => ({ revision: 1 })); + let revision: { revision: number; leaseOwner: string | null; leaseExpiresAt: Date | null } | null = null; + const revisionUpdate = vi.fn(async (rawArgs: unknown) => { + if (!revision) { + return { count: 0 }; + } + const data = asRecord(asRecord(rawArgs).data); + const revisionChange = asRecord(data.revision); + if (typeof revisionChange.increment === 'number') { + revision.revision += revisionChange.increment; + } + if ('leaseOwner' in data) { + revision.leaseOwner = typeof data.leaseOwner === 'string' ? data.leaseOwner : null; + } + if ('leaseExpiresAt' in data) { + revision.leaseExpiresAt = data.leaseExpiresAt instanceof Date ? data.leaseExpiresAt : null; + } + return { count: 1 }; + }); + const revisionCreate = vi.fn(async (rawArgs: unknown) => { + const rawData = asRecord(rawArgs).data; + const row = Array.isArray(rawData) ? asRecord(rawData[0]) : asRecord(rawData); + revision = { + revision: typeof row.revision === 'number' ? row.revision : 0, + leaseOwner: typeof row.leaseOwner === 'string' ? row.leaseOwner : null, + leaseExpiresAt: row.leaseExpiresAt instanceof Date ? row.leaseExpiresAt : null, + }; + return { count: 1 }; + }); const prisma = { generalTurn: { findMany: vi.fn(async () => []), @@ -40,7 +68,9 @@ describe('input event atomicity', () => { }), }, generalTurnRevision: { - upsert: revisionUpsert, + findUnique: vi.fn(async () => null), + createMany: revisionCreate, + updateMany: revisionUpdate, }, nationTurn: { findMany: vi.fn(async () => []), @@ -51,6 +81,7 @@ describe('input event atomicity', () => { const store = new InMemoryReservedTurnStore(prisma, { maxGeneralTurns: 1, maxNationTurns: 1, + leaseOwner: 'test-daemon', }); store.shiftGeneralTurns(7, -1); store.ensureGeneralTurns(8); @@ -59,23 +90,38 @@ describe('input event atomicity', () => { expect(store.peekDirtyState()).toEqual({ generalIds: [7], generalInitializationIds: [8], + generalLeaseIds: [], nationKeys: [], nationInitializationKeys: [], + nationLeaseKeys: [], }); failCreate = false; await store.flushChanges(); - expect(revisionUpsert).toHaveBeenCalledOnce(); - expect(revisionUpsert).toHaveBeenCalledWith({ - where: { generalId: 7 }, - create: { generalId: 7, revision: 1 }, - update: { revision: { increment: 1 } }, + expect(revisionCreate).toHaveBeenCalledOnce(); + expect(revisionCreate).toHaveBeenCalledWith({ + data: [ + { + generalId: 7, + revision: 0, + leaseOwner: 'test-daemon', + leaseExpiresAt: expect.any(Date), + }, + ], + skipDuplicates: true, + }); + expect(revision).toMatchObject({ + revision: 1, + leaseOwner: null, + leaseExpiresAt: null, }); expect(store.peekDirtyState()).toEqual({ generalIds: [], generalInitializationIds: [], + generalLeaseIds: [], nationKeys: [], nationInitializationKeys: [], + nationLeaseKeys: [], }); }); diff --git a/app/game-engine/test/monthlyCreateManyNpcAction.test.ts b/app/game-engine/test/monthlyCreateManyNpcAction.test.ts index 73b6336..12f92f4 100644 --- a/app/game-engine/test/monthlyCreateManyNpcAction.test.ts +++ b/app/game-engine/test/monthlyCreateManyNpcAction.test.ts @@ -197,8 +197,10 @@ describe('CreateManyNPC monthly action', () => { expect(reservedTurns.peekDirtyState()).toEqual({ generalIds: [], generalInitializationIds: [created.id], + generalLeaseIds: [], nationKeys: [], nationInitializationKeys: [], + nationLeaseKeys: [], }); expect(world.peekDirtyState().logs).toMatchInlineSnapshot(` [ diff --git a/app/game-engine/test/reservedTurnLease.integration.test.ts b/app/game-engine/test/reservedTurnLease.integration.test.ts new file mode 100644 index 0000000..5b96fad --- /dev/null +++ b/app/game-engine/test/reservedTurnLease.integration.test.ts @@ -0,0 +1,135 @@ +import { afterAll, beforeAll, describe, expect, it } from 'vitest'; +import { createGamePostgresConnector } from '@sammo-ts/infra'; + +import { + InMemoryReservedTurnStore, + ReservedTurnLeaseConflictError, +} from '../src/turn/reservedTurnStore.js'; + +const databaseUrl = process.env.RESERVED_TURN_DATABASE_URL; +const describeIntegration = databaseUrl ? describe : describe.skip; +const GENERAL_ID = 2_147_400_002; + +describeIntegration('reserved turn daemon/API lease integration', () => { + const connector = databaseUrl ? createGamePostgresConnector({ url: databaseUrl }) : null; + + beforeAll(async () => { + if (!connector) { + return; + } + await connector.connect(); + await connector.prisma.generalTurn.deleteMany({ where: { generalId: GENERAL_ID } }); + await connector.prisma.generalTurnRevision.deleteMany({ where: { generalId: GENERAL_ID } }); + await connector.prisma.generalTurn.createMany({ + data: [ + { generalId: GENERAL_ID, turnIdx: 0, actionCode: 'che_훈련', arg: {} }, + { generalId: GENERAL_ID, turnIdx: 1, actionCode: 'che_사기진작', arg: {} }, + ], + }); + }); + + afterAll(async () => { + if (!connector) { + return; + } + await connector.prisma.generalTurn.deleteMany({ where: { generalId: GENERAL_ID } }); + await connector.prisma.generalTurnRevision.deleteMany({ where: { generalId: GENERAL_ID } }); + await connector.disconnect(); + }); + + it('blocks the API CAS during execution and releases the lease atomically with the shifted queue', async () => { + if (!connector) { + throw new Error('integration connector is unavailable'); + } + const store = new InMemoryReservedTurnStore(connector.prisma, { + maxGeneralTurns: 2, + maxNationTurns: 1, + leaseOwner: 'integration-daemon', + leaseDurationMs: 60_000, + }); + + await store.prepareTurnsForExecution(GENERAL_ID); + const leased = await connector.prisma.generalTurnRevision.findUniqueOrThrow({ + where: { generalId: GENERAL_ID }, + }); + expect(leased).toMatchObject({ revision: 0, leaseOwner: 'integration-daemon' }); + + const blockedApiClaim = await connector.prisma.generalTurnRevision.updateMany({ + where: { + generalId: GENERAL_ID, + revision: 0, + OR: [{ leaseOwner: null }, { leaseExpiresAt: { lte: new Date() } }], + }, + data: { + revision: 1, + leaseOwner: null, + leaseExpiresAt: null, + }, + }); + expect(blockedApiClaim.count).toBe(0); + + store.shiftGeneralTurns(GENERAL_ID, -1); + const changes = store.peekDirtyState(); + await connector.prisma.$transaction((transaction) => store.persistChanges(transaction, changes)); + store.acknowledgeDirtyState(changes); + + const committed = await connector.prisma.generalTurnRevision.findUniqueOrThrow({ + where: { generalId: GENERAL_ID }, + }); + expect(committed).toMatchObject({ + revision: 1, + leaseOwner: null, + leaseExpiresAt: null, + }); + const turns = await connector.prisma.generalTurn.findMany({ + where: { generalId: GENERAL_ID }, + orderBy: { turnIdx: 'asc' }, + }); + expect(turns.map((turn) => turn.actionCode)).toEqual(['che_사기진작', '휴식']); + + const nextApiClaim = await connector.prisma.generalTurnRevision.updateMany({ + where: { + generalId: GENERAL_ID, + revision: 1, + OR: [{ leaseOwner: null }, { leaseExpiresAt: { lte: new Date() } }], + }, + data: { revision: 2 }, + }); + expect(nextApiClaim.count).toBe(1); + }); + + it('blocks a daemon-side replacement before deleting turns while an API lease is active', async () => { + if (!connector) { + throw new Error('integration connector is unavailable'); + } + await connector.prisma.generalTurnRevision.update({ + where: { generalId: GENERAL_ID }, + data: { + leaseOwner: 'integration-api', + leaseExpiresAt: new Date(Date.now() + 60_000), + }, + }); + const before = await connector.prisma.generalTurn.findMany({ + where: { generalId: GENERAL_ID }, + orderBy: { turnIdx: 'asc' }, + }); + const store = new InMemoryReservedTurnStore(connector.prisma, { + maxGeneralTurns: 2, + maxNationTurns: 1, + leaseOwner: 'integration-daemon-replacement', + leaseDurationMs: 60_000, + }); + store.replaceGeneralTurns(GENERAL_ID, { action: 'che_훈련', args: {} }); + const changes = store.peekDirtyState(); + + await expect( + connector.prisma.$transaction((transaction) => store.persistChanges(transaction, changes)) + ).rejects.toBeInstanceOf(ReservedTurnLeaseConflictError); + + const after = await connector.prisma.generalTurn.findMany({ + where: { generalId: GENERAL_ID }, + orderBy: { turnIdx: 'asc' }, + }); + expect(after.map((turn) => turn.actionCode)).toEqual(before.map((turn) => turn.actionCode)); + }); +}); diff --git a/app/game-engine/test/reservedTurnLease.test.ts b/app/game-engine/test/reservedTurnLease.test.ts new file mode 100644 index 0000000..ed215d2 --- /dev/null +++ b/app/game-engine/test/reservedTurnLease.test.ts @@ -0,0 +1,288 @@ +import { describe, expect, it, vi } from 'vitest'; +import { asRecord } from '@sammo-ts/common'; + +import { + InMemoryReservedTurnStore, + ReservedTurnLeaseConflictError, +} from '../src/turn/reservedTurnStore.js'; + +interface RevisionRow { + revision: number; + leaseOwner: string | null; + leaseExpiresAt: Date | null; +} + +const buildHarness = (initialRevision: RevisionRow | null = null) => { + let revision = initialRevision; + let nationRevision: RevisionRow | null = null; + const writtenTurns: unknown[] = []; + const generalFindMany = vi.fn(async () => [ + { + id: 1, + generalId: 7, + turnIdx: 0, + actionCode: 'che_훈련', + arg: {}, + createdAt: new Date(), + }, + ]); + const generalDeleteMany = vi.fn(async () => ({ count: 1 })); + const generalCreateMany = vi.fn(async (args: unknown) => { + writtenTurns.push(args); + return { count: 2 }; + }); + const updateMany = vi.fn(async (rawArgs: unknown) => { + const args = asRecord(rawArgs); + const where = asRecord(args.where); + const data = asRecord(args.data); + if (!revision) { + return { count: 0 }; + } + if (typeof where.generalId === 'number' && where.generalId !== 7) { + return { count: 0 }; + } + if (typeof where.leaseOwner === 'string' && where.leaseOwner !== revision.leaseOwner) { + return { count: 0 }; + } + if (Array.isArray(where.OR)) { + const now = new Date(); + const eligible = + revision.leaseOwner === null || + revision.leaseOwner === 'daemon-1' || + (revision.leaseExpiresAt !== null && revision.leaseExpiresAt.getTime() <= now.getTime()); + if (!eligible) { + return { count: 0 }; + } + } + const revisionChange = asRecord(data.revision); + if (typeof revisionChange.increment === 'number') { + revision.revision += revisionChange.increment; + } + if ('leaseOwner' in data) { + revision.leaseOwner = typeof data.leaseOwner === 'string' ? data.leaseOwner : null; + } + if ('leaseExpiresAt' in data) { + revision.leaseExpiresAt = data.leaseExpiresAt instanceof Date ? data.leaseExpiresAt : null; + } + return { count: 1 }; + }); + const createMany = vi.fn(async (rawArgs: unknown) => { + if (revision) { + return { count: 0 }; + } + const args = asRecord(rawArgs); + const rows = Array.isArray(args.data) ? args.data : []; + const data = asRecord(rows[0]); + revision = { + revision: typeof data.revision === 'number' ? data.revision : 0, + leaseOwner: typeof data.leaseOwner === 'string' ? data.leaseOwner : null, + leaseExpiresAt: data.leaseExpiresAt instanceof Date ? data.leaseExpiresAt : null, + }; + return { count: 1 }; + }); + const nationRevisionUpdateMany = vi.fn(async (rawArgs: unknown) => { + const args = asRecord(rawArgs); + const where = asRecord(args.where); + const data = asRecord(args.data); + if (!nationRevision) { + return { count: 0 }; + } + if (typeof where.leaseOwner === 'string' && where.leaseOwner !== nationRevision.leaseOwner) { + return { count: 0 }; + } + const revisionChange = asRecord(data.revision); + if (typeof revisionChange.increment === 'number') { + nationRevision.revision += revisionChange.increment; + } + if ('leaseOwner' in data) { + nationRevision.leaseOwner = typeof data.leaseOwner === 'string' ? data.leaseOwner : null; + } + if ('leaseExpiresAt' in data) { + nationRevision.leaseExpiresAt = data.leaseExpiresAt instanceof Date ? data.leaseExpiresAt : null; + } + return { count: 1 }; + }); + const nationRevisionCreateMany = vi.fn(async (rawArgs: unknown) => { + if (nationRevision) { + return { count: 0 }; + } + const args = asRecord(rawArgs); + const rows = Array.isArray(args.data) ? args.data : []; + const data = asRecord(rows[0]); + nationRevision = { + revision: typeof data.revision === 'number' ? data.revision : 0, + leaseOwner: typeof data.leaseOwner === 'string' ? data.leaseOwner : null, + leaseExpiresAt: data.leaseExpiresAt instanceof Date ? data.leaseExpiresAt : null, + }; + return { count: 1 }; + }); + const prisma = { + generalTurn: { + findMany: generalFindMany, + deleteMany: generalDeleteMany, + createMany: generalCreateMany, + }, + generalTurnRevision: { + findUnique: vi.fn(async () => revision), + createMany, + updateMany, + }, + nationTurn: { + findMany: vi.fn(async () => [ + { + id: 1, + nationId: 3, + officerLevel: 12, + turnIdx: 0, + actionCode: 'che_포상', + arg: {}, + createdAt: new Date(), + }, + ]), + deleteMany: vi.fn(async () => ({ count: 0 })), + createMany: vi.fn(async () => ({ count: 0 })), + }, + nationTurnRevision: { + findUnique: vi.fn(async () => nationRevision), + createMany: nationRevisionCreateMany, + updateMany: nationRevisionUpdateMany, + }, + }; + const store = new InMemoryReservedTurnStore(prisma, { + maxGeneralTurns: 2, + maxNationTurns: 1, + leaseOwner: 'daemon-1', + leaseDurationMs: 60_000, + }); + return { + store, + prisma, + generalFindMany, + generalDeleteMany, + generalCreateMany, + writtenTurns, + getRevision: () => revision, + getNationRevision: () => nationRevision, + stealLease: () => { + if (revision) { + revision.leaseOwner = 'other-writer'; + } + }, + }; +}; + +describe('reserved turn daemon lease', () => { + it('holds the queue lease from refresh through shift and releases it with the revision increment', async () => { + const harness = buildHarness(); + + await harness.store.prepareTurnsForExecution(7); + expect(harness.getRevision()).toMatchObject({ revision: 0, leaseOwner: 'daemon-1' }); + expect(harness.store.getGeneralTurn(7, 0).action).toBe('che_훈련'); + + harness.store.shiftGeneralTurns(7, -1); + const changes = harness.store.peekDirtyState(); + expect(changes).toMatchObject({ generalIds: [7], generalLeaseIds: [7] }); + await harness.store.persistChanges(harness.prisma, changes); + harness.store.acknowledgeDirtyState(changes); + + expect(harness.getRevision()).toMatchObject({ + revision: 1, + leaseOwner: null, + leaseExpiresAt: null, + }); + expect(harness.generalDeleteMany).toHaveBeenCalledOnce(); + expect(harness.generalCreateMany).toHaveBeenCalledOnce(); + expect(harness.store.peekDirtyState()).toEqual({ + generalIds: [], + generalInitializationIds: [], + generalLeaseIds: [], + nationKeys: [], + nationInitializationKeys: [], + nationLeaseKeys: [], + }); + }); + + it('rejects an active foreign lease before reading the queue', async () => { + const harness = buildHarness({ + revision: 4, + leaseOwner: 'api-writer', + leaseExpiresAt: new Date(Date.now() + 60_000), + }); + + await expect(harness.store.prepareTurnsForExecution(7)).rejects.toBeInstanceOf( + ReservedTurnLeaseConflictError + ); + expect(harness.generalFindMany).not.toHaveBeenCalled(); + }); + + it('takes over an expired foreign lease before reading the queue', async () => { + const harness = buildHarness({ + revision: 4, + leaseOwner: 'stopped-daemon', + leaseExpiresAt: new Date(Date.now() - 1_000), + }); + + await harness.store.prepareTurnsForExecution(7); + + expect(harness.getRevision()).toMatchObject({ + revision: 4, + leaseOwner: 'daemon-1', + }); + expect(harness.generalFindMany).toHaveBeenCalledOnce(); + }); + + it('leases and releases the matching nation officer queue with the general queue', async () => { + const harness = buildHarness(); + + await harness.store.prepareTurnsForExecution(7, { nationId: 3, officerLevel: 12 }); + expect(harness.getRevision()).toMatchObject({ revision: 0, leaseOwner: 'daemon-1' }); + expect(harness.getNationRevision()).toMatchObject({ revision: 0, leaseOwner: 'daemon-1' }); + + harness.store.shiftNationTurns(3, 12, -1); + const changes = harness.store.peekDirtyState(); + expect(changes).toMatchObject({ + generalLeaseIds: [7], + nationKeys: ['3:12'], + nationLeaseKeys: ['3:12'], + }); + await harness.store.persistChanges(harness.prisma, changes); + harness.store.acknowledgeDirtyState(changes); + + expect(harness.getRevision()).toMatchObject({ revision: 0, leaseOwner: null }); + expect(harness.getNationRevision()).toMatchObject({ revision: 1, leaseOwner: null }); + }); + + it('detects a lost lease before deleting or replacing turns', async () => { + const harness = buildHarness(); + await harness.store.prepareTurnsForExecution(7); + harness.store.shiftGeneralTurns(7, -1); + harness.stealLease(); + const changes = harness.store.peekDirtyState(); + + await expect(harness.store.persistChanges(harness.prisma, changes)).rejects.toBeInstanceOf( + ReservedTurnLeaseConflictError + ); + expect(harness.generalDeleteMany).not.toHaveBeenCalled(); + expect(harness.generalCreateMany).not.toHaveBeenCalled(); + expect(harness.store.peekDirtyState()).toMatchObject({ + generalIds: [7], + generalLeaseIds: [7], + }); + }); + + it('blocks a daemon-side queue replacement behind an active foreign lease', async () => { + const harness = buildHarness({ + revision: 6, + leaseOwner: 'api-writer', + leaseExpiresAt: new Date(Date.now() + 60_000), + }); + harness.store.replaceGeneralTurns(7, { action: 'che_훈련', args: {} }); + const changes = harness.store.peekDirtyState(); + + await expect(harness.store.persistChanges(harness.prisma, changes)).rejects.toBeInstanceOf( + ReservedTurnLeaseConflictError + ); + expect(harness.generalDeleteMany).not.toHaveBeenCalled(); + expect(harness.generalCreateMany).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/infra/prisma/game.prisma b/packages/infra/prisma/game.prisma index 2a3155f..3622302 100644 --- a/packages/infra/prisma/game.prisma +++ b/packages/infra/prisma/game.prisma @@ -78,11 +78,11 @@ model InputEvent { } model TurnDaemonLease { - profile String @id - ownerId String @map("owner_id") - leaseUntil DateTime @map("lease_until") - fencingEpoch BigInt @default(1) @map("fencing_epoch") - heartbeatAt DateTime @default(now()) @map("heartbeat_at") + profile String @id + ownerId String @map("owner_id") + leaseUntil DateTime @map("lease_until") + fencingEpoch BigInt @default(1) @map("fencing_epoch") + heartbeatAt DateTime @default(now()) @map("heartbeat_at") @@map("turn_daemon_lease") } @@ -369,9 +369,11 @@ model GeneralTurn { } model GeneralTurnRevision { - generalId Int @id @map("general_id") - revision Int @default(0) - updatedAt DateTime @updatedAt @map("updated_at") + generalId Int @id @map("general_id") + revision Int @default(0) + leaseOwner String? @map("lease_owner") + leaseExpiresAt DateTime? @map("lease_expires_at") + updatedAt DateTime @updatedAt @map("updated_at") @@map("general_turn_revision") } @@ -390,10 +392,12 @@ model NationTurn { } model NationTurnRevision { - nationId Int @map("nation_id") - officerLevel Int @map("officer_level") - revision Int @default(0) - updatedAt DateTime @updatedAt @map("updated_at") + nationId Int @map("nation_id") + officerLevel Int @map("officer_level") + revision Int @default(0) + leaseOwner String? @map("lease_owner") + leaseExpiresAt DateTime? @map("lease_expires_at") + updatedAt DateTime @updatedAt @map("updated_at") @@id([nationId, officerLevel]) @@map("nation_turn_revision") diff --git a/packages/infra/prisma/migrations/20260726190000_add_reserved_turn_lease/migration.sql b/packages/infra/prisma/migrations/20260726190000_add_reserved_turn_lease/migration.sql new file mode 100644 index 0000000..bf03156 --- /dev/null +++ b/packages/infra/prisma/migrations/20260726190000_add_reserved_turn_lease/migration.sql @@ -0,0 +1,7 @@ +ALTER TABLE "general_turn_revision" + ADD COLUMN "lease_owner" TEXT, + ADD COLUMN "lease_expires_at" TIMESTAMP(3); + +ALTER TABLE "nation_turn_revision" + ADD COLUMN "lease_owner" TEXT, + ADD COLUMN "lease_expires_at" TIMESTAMP(3); diff --git a/packages/infra/src/turnEngineDb.ts b/packages/infra/src/turnEngineDb.ts index b653b4f..a6af1cd 100644 --- a/packages/infra/src/turnEngineDb.ts +++ b/packages/infra/src/turnEngineDb.ts @@ -437,7 +437,9 @@ export interface TurnEngineDatabaseClient { createMany(args?: unknown): Promise; }; generalTurnRevision?: { - upsert(args: unknown): Promise; + findUnique(args: unknown): Promise; + createMany(args: unknown): Promise<{ count: number }>; + updateMany(args: unknown): Promise<{ count: number }>; deleteMany(args?: unknown): Promise; }; nationTurn: { @@ -446,7 +448,9 @@ export interface TurnEngineDatabaseClient { createMany(args?: unknown): Promise; }; nationTurnRevision?: { - upsert(args: unknown): Promise; + findUnique(args: unknown): Promise; + createMany(args: unknown): Promise<{ count: number }>; + updateMany(args: unknown): Promise<{ count: number }>; deleteMany(args?: unknown): Promise; }; }