feat: read model outbox 재시도 dispatcher를 추가

This commit is contained in:
2026-08-16 18:09:04 +00:00
parent 54126cd67f
commit 8278ef5221
9 changed files with 611 additions and 1 deletions
+1
View File
@@ -7,3 +7,4 @@ export * from './db.js';
export * from './redis.js';
export * from './turnEngineDb.js';
export * from './readModelChangeJournal.js';
export * from './readModelOutboxDispatcher.js';
@@ -0,0 +1,175 @@
import { parseReadModelOutboxPayload, type ReadModelOutboxPayloadV1 } from '@sammo-ts/common';
import { GamePrisma, type GamePrismaClient } from './gamePrisma.js';
export interface ClaimedReadModelOutbox {
id: bigint;
payload: unknown;
attempts: number;
}
type ClaimedRow = {
id: bigint;
payload: unknown;
attempts: number;
};
export interface ReadModelOutboxDispatchOptions {
owner: string;
limit?: number;
leaseMs?: number;
retryBaseMs?: number;
retryMaxMs?: number;
now?: () => Date;
}
export interface ReadModelOutboxDispatchResult {
claimed: number;
delivered: number;
failed: number;
}
const normalizeLimit = (value: number | undefined): number =>
Math.min(500, Math.max(1, Math.floor(value ?? 50)));
const normalizeDuration = (value: number | undefined, fallback: number): number =>
Math.max(1, Math.floor(value ?? fallback));
const formatDispatchError = (error: unknown): string => {
const text = error instanceof Error ? `${error.name}: ${error.message}` : String(error);
return text.replaceAll(/\s+/gu, ' ').slice(0, 1_000);
};
const retryDelayMs = (attempts: number, baseMs: number, maxMs: number): number => {
const exponent = Math.min(20, Math.max(0, attempts - 1));
return Math.min(maxMs, baseMs * 2 ** exponent);
};
export const claimReadModelOutboxBatch = async (
db: GamePrismaClient,
options: Pick<ReadModelOutboxDispatchOptions, 'owner' | 'limit' | 'leaseMs'> & { now?: Date }
): Promise<ClaimedReadModelOutbox[]> => {
if (!options.owner.trim()) {
throw new Error('Read-model outbox owner must not be empty.');
}
const limit = normalizeLimit(options.limit);
const leaseMs = normalizeDuration(options.leaseMs, 30_000);
const now = options.now ?? new Date();
const leaseExpiredBefore = new Date(now.getTime() - leaseMs);
const rows = await db.$queryRaw<ClaimedRow[]>(GamePrisma.sql`
WITH candidates AS (
SELECT "id"
FROM "read_model_outbox"
WHERE "delivered_at" IS NULL
AND "available_at" <= ${now}
AND ("locked_at" IS NULL OR "locked_at" < ${leaseExpiredBefore})
ORDER BY "id"
LIMIT ${limit}
FOR UPDATE SKIP LOCKED
)
UPDATE "read_model_outbox" AS outbox
SET
"attempts" = outbox."attempts" + 1,
"locked_at" = ${now},
"lock_owner" = ${options.owner},
"last_error" = NULL
FROM candidates
WHERE outbox."id" = candidates."id"
RETURNING outbox."id", outbox."payload", outbox."attempts"
`);
return rows.map((row) => ({ id: BigInt(row.id), payload: row.payload, attempts: row.attempts }));
};
export const markReadModelOutboxDelivered = async (
db: GamePrismaClient,
input: { id: bigint; owner: string; deliveredAt?: Date }
): Promise<boolean> => {
const result = await db.readModelOutbox.updateMany({
where: { id: input.id, lockOwner: input.owner, deliveredAt: null },
data: {
deliveredAt: input.deliveredAt ?? new Date(),
lockedAt: null,
lockOwner: null,
lastError: null,
},
});
return result.count === 1;
};
export const releaseReadModelOutbox = async (
db: GamePrismaClient,
input: { id: bigint; owner: string; error: unknown; availableAt: Date }
): Promise<boolean> => {
const result = await db.readModelOutbox.updateMany({
where: { id: input.id, lockOwner: input.owner, deliveredAt: null },
data: {
availableAt: input.availableAt,
lockedAt: null,
lockOwner: null,
lastError: formatDispatchError(input.error),
},
});
return result.count === 1;
};
export const dispatchReadModelOutboxBatch = async (
db: GamePrismaClient,
publish: (payload: ReadModelOutboxPayloadV1, outboxId: bigint) => Promise<void>,
options: ReadModelOutboxDispatchOptions
): Promise<ReadModelOutboxDispatchResult> => {
const now = options.now ?? (() => new Date());
const retryBaseMs = normalizeDuration(options.retryBaseMs, 1_000);
const retryMaxMs = Math.max(retryBaseMs, normalizeDuration(options.retryMaxMs, 60_000));
const claimed = await claimReadModelOutboxBatch(db, {
owner: options.owner,
limit: options.limit,
leaseMs: options.leaseMs,
now: now(),
});
let delivered = 0;
let failed = 0;
for (const item of claimed) {
try {
const payload = parseReadModelOutboxPayload(item.payload);
if (!payload) {
throw new Error(`Read-model outbox ${item.id.toString()} has an invalid payload.`);
}
await publish(payload, item.id);
if (!(await markReadModelOutboxDelivered(db, { id: item.id, owner: options.owner, deliveredAt: now() }))) {
throw new Error(`Read-model outbox ${item.id.toString()} lost its delivery lease.`);
}
delivered += 1;
} catch (error) {
failed += 1;
await releaseReadModelOutbox(db, {
id: item.id,
owner: options.owner,
error,
availableAt: new Date(now().getTime() + retryDelayMs(item.attempts, retryBaseMs, retryMaxMs)),
});
}
}
return { claimed: claimed.length, delivered, failed };
};
export const pruneDeliveredReadModelOutbox = async (
db: GamePrismaClient,
input: { deliveredBefore: Date; limit?: number }
): Promise<number> => {
const limit = normalizeLimit(input.limit);
const rows = await db.$queryRaw<Array<{ id: bigint }>>(GamePrisma.sql`
DELETE FROM "read_model_outbox"
WHERE "id" IN (
SELECT "id"
FROM "read_model_outbox"
WHERE "delivered_at" < ${input.deliveredBefore}
ORDER BY "id"
LIMIT ${limit}
)
RETURNING "id"
`);
return rows.length;
};
@@ -0,0 +1,107 @@
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest';
import { createGamePostgresConnector, type GamePrismaClient } from '../src/gamePrisma.js';
import { writeReadModelChangeJournal } from '../src/readModelChangeJournal.js';
import {
claimReadModelOutboxBatch,
dispatchReadModelOutboxBatch,
pruneDeliveredReadModelOutbox,
} from '../src/readModelOutboxDispatcher.js';
const databaseUrl = process.env.READ_MODEL_JOURNAL_DATABASE_URL;
const integration = describe.skipIf(!databaseUrl);
integration('read-model outbox PostgreSQL delivery boundary', () => {
let disconnect: (() => Promise<void>) | undefined;
let prisma: GamePrismaClient;
beforeAll(async () => {
if (!databaseUrl) throw new Error('READ_MODEL_JOURNAL_DATABASE_URL is required.');
const connector = createGamePostgresConnector({ url: databaseUrl });
prisma = connector.prisma;
disconnect = connector.disconnect;
await connector.connect();
});
afterAll(async () => disconnect?.());
beforeEach(async () => {
await prisma.$executeRaw`TRUNCATE TABLE "read_model_outbox", "read_model_revision" RESTART IDENTITY`;
});
const enqueue = async (entityId: number): Promise<void> => {
await prisma.$transaction((transaction) =>
writeReadModelChangeJournal(transaction, [{ domain: 'general.content', entityId }])
);
};
it('claims concurrent worker batches without overlap', async () => {
await Promise.all(Array.from({ length: 20 }, (_, index) => enqueue(index + 1)));
const now = new Date('2099-08-16T00:00:00.000Z');
const [left, right] = await Promise.all([
claimReadModelOutboxBatch(prisma, { owner: 'left', limit: 10, now }),
claimReadModelOutboxBatch(prisma, { owner: 'right', limit: 10, now }),
]);
const ids = [...left, ...right].map(({ id }) => id);
expect(ids).toHaveLength(20);
expect(new Set(ids).size).toBe(20);
});
it('releases a publish failure and later delivers the same row', async () => {
await enqueue(7);
const failedAt = new Date('2099-08-16T00:00:00.000Z');
await expect(
dispatchReadModelOutboxBatch(prisma, vi.fn().mockRejectedValue(new Error('redis unavailable')), {
owner: 'worker-a',
now: () => failedAt,
retryBaseMs: 1_000,
})
).resolves.toEqual({ claimed: 1, delivered: 0, failed: 1 });
const released = await prisma.readModelOutbox.findUniqueOrThrow({ where: { id: 1n } });
expect(released).toMatchObject({ attempts: 1, lockedAt: null, lockOwner: null, deliveredAt: null });
expect(released.lastError).toContain('redis unavailable');
const retryAt = new Date('2099-08-16T00:00:01.000Z');
const publish = vi.fn().mockResolvedValue(undefined);
await expect(
dispatchReadModelOutboxBatch(prisma, publish, { owner: 'worker-b', now: () => retryAt })
).resolves.toEqual({ claimed: 1, delivered: 1, failed: 0 });
expect(publish).toHaveBeenCalledTimes(1);
await expect(prisma.readModelOutbox.findUniqueOrThrow({ where: { id: 1n } })).resolves.toMatchObject({
attempts: 2,
lockOwner: null,
deliveredAt: retryAt,
});
});
it('allows a lease-expired row to be republished after publish-before-ack crash', async () => {
await enqueue(7);
const first = await claimReadModelOutboxBatch(prisma, {
owner: 'crashed-worker',
leaseMs: 30_000,
now: new Date('2099-08-16T00:00:00.000Z'),
});
expect(first).toHaveLength(1);
const second = await claimReadModelOutboxBatch(prisma, {
owner: 'recovery-worker',
leaseMs: 30_000,
now: new Date('2099-08-16T00:00:31.000Z'),
});
expect(second.map(({ id }) => id)).toEqual(first.map(({ id }) => id));
expect(second[0]?.attempts).toBe(2);
});
it('prunes delivered rows in bounded batches', async () => {
await Promise.all([enqueue(1), enqueue(2), enqueue(3)]);
const deliveredAt = new Date('2099-08-15T00:00:00.000Z');
await prisma.readModelOutbox.updateMany({ data: { deliveredAt } });
await expect(
pruneDeliveredReadModelOutbox(prisma, {
deliveredBefore: new Date('2099-08-16T00:00:00.000Z'),
limit: 2,
})
).resolves.toBe(2);
await expect(prisma.readModelOutbox.count()).resolves.toBe(1);
});
});
@@ -0,0 +1,97 @@
import { describe, expect, it, vi } from 'vitest';
import type { GamePrismaClient } from '../src/gamePrisma.js';
import {
claimReadModelOutboxBatch,
dispatchReadModelOutboxBatch,
pruneDeliveredReadModelOutbox,
} from '../src/readModelOutboxDispatcher.js';
const validPayload = {
version: 1,
changes: [['general.content', 7, '3']],
};
const createDb = (rows: readonly object[]) => {
const queryRaw = vi.fn().mockResolvedValue(rows);
const updateMany = vi.fn().mockResolvedValue({ count: 1 });
return {
db: { $queryRaw: queryRaw, readModelOutbox: { updateMany } } as unknown as GamePrismaClient,
queryRaw,
updateMany,
};
};
describe('read-model outbox dispatcher', () => {
it('claims a bounded lease with one statement and preserves delivery identity', async () => {
const fixture = createDb([{ id: 41n, payload: validPayload, attempts: 2 }]);
await expect(
claimReadModelOutboxBatch(fixture.db, {
owner: 'worker-a',
limit: 25,
leaseMs: 15_000,
now: new Date('2026-08-16T00:00:00.000Z'),
})
).resolves.toEqual([{ id: 41n, payload: validPayload, attempts: 2 }]);
expect(fixture.queryRaw).toHaveBeenCalledTimes(1);
});
it('publishes and acknowledges a valid payload', async () => {
const fixture = createDb([{ id: 1n, payload: validPayload, attempts: 1 }]);
const publish = vi.fn().mockResolvedValue(undefined);
const result = await dispatchReadModelOutboxBatch(fixture.db, publish, {
owner: 'worker-a',
now: () => new Date('2026-08-16T00:00:00.000Z'),
});
expect(result).toEqual({ claimed: 1, delivered: 1, failed: 0 });
expect(publish).toHaveBeenCalledWith(validPayload, 1n);
expect(fixture.updateMany).toHaveBeenCalledWith(
expect.objectContaining({
where: { id: 1n, lockOwner: 'worker-a', deliveredAt: null },
data: expect.objectContaining({ deliveredAt: new Date('2026-08-16T00:00:00.000Z') }),
})
);
});
it('releases failed and malformed rows with bounded retry state', async () => {
const fixture = createDb([
{ id: 1n, payload: validPayload, attempts: 3 },
{ id: 2n, payload: { version: 99 }, attempts: 1 },
]);
const publish = vi.fn().mockRejectedValue(new Error('redis unavailable'));
const result = await dispatchReadModelOutboxBatch(fixture.db, publish, {
owner: 'worker-a',
retryBaseMs: 1_000,
retryMaxMs: 10_000,
now: () => new Date('2026-08-16T00:00:00.000Z'),
});
expect(result).toEqual({ claimed: 2, delivered: 0, failed: 2 });
expect(publish).toHaveBeenCalledTimes(1);
expect(fixture.updateMany).toHaveBeenNthCalledWith(
1,
expect.objectContaining({
where: { id: 1n, lockOwner: 'worker-a', deliveredAt: null },
data: expect.objectContaining({ availableAt: new Date('2026-08-16T00:00:04.000Z') }),
})
);
expect(fixture.updateMany).toHaveBeenNthCalledWith(
2,
expect.objectContaining({
where: { id: 2n, lockOwner: 'worker-a', deliveredAt: null },
data: expect.objectContaining({ availableAt: new Date('2026-08-16T00:00:01.000Z') }),
})
);
});
it('prunes only a bounded delivered batch', async () => {
const fixture = createDb([{ id: 1n }, { id: 2n }]);
await expect(
pruneDeliveredReadModelOutbox(fixture.db, {
deliveredBefore: new Date('2026-08-15T00:00:00.000Z'),
limit: 100,
})
).resolves.toBe(2);
});
});
+4
View File
@@ -7,6 +7,10 @@ export default defineConfig({
test: {
environment: 'node',
globals: true,
// Integration files share the explicitly supplied disposable schema.
// Keep file-level TRUNCATE/setup boundaries from racing each other;
// individual tests still create concurrent writers deliberately.
fileParallelism: false,
include: ['test/**/*.test.ts'],
},
});