From 5a03d998a5dcf12a105fd96f834e4a14d77e7a57 Mon Sep 17 00:00:00 2001 From: Efuntoye Victor Date: Mon, 28 Sep 2026 08:17:21 +0100 Subject: [PATCH] feat: add a scheduled worker to prune expired idempotency keys (#1495) Closes #1495 --- .../migration.sql | 14 ++ backend/prisma/schema.prisma | 12 ++ .../src/workers/idempotency-prune-worker.ts | 79 ++++++++++++ backend/src/workers/index.ts | 11 ++ .../tests/idempotency-prune-worker.test.ts | 122 ++++++++++++++++++ 5 files changed, 238 insertions(+) create mode 100644 backend/prisma/migrations/20260928000000_add_idempotency_key/migration.sql create mode 100644 backend/src/workers/idempotency-prune-worker.ts create mode 100644 backend/tests/idempotency-prune-worker.test.ts diff --git a/backend/prisma/migrations/20260928000000_add_idempotency_key/migration.sql b/backend/prisma/migrations/20260928000000_add_idempotency_key/migration.sql new file mode 100644 index 00000000..ca4072e3 --- /dev/null +++ b/backend/prisma/migrations/20260928000000_add_idempotency_key/migration.sql @@ -0,0 +1,14 @@ +-- CreateTable +CREATE TABLE "IdempotencyKey" ( + "id" TEXT NOT NULL, + "key" TEXT NOT NULL, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "IdempotencyKey_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE UNIQUE INDEX "IdempotencyKey_key_key" ON "IdempotencyKey"("key"); + +-- CreateIndex +CREATE INDEX "IdempotencyKey_createdAt_idx" ON "IdempotencyKey"("createdAt"); diff --git a/backend/prisma/schema.prisma b/backend/prisma/schema.prisma index b30f85c6..6f90e2f0 100644 --- a/backend/prisma/schema.prisma +++ b/backend/prisma/schema.prisma @@ -156,6 +156,18 @@ model AlertHistory { @@index([sentAt]) } +// IdempotencyKey model - records the idempotency keys of mutating requests so a +// client retry cannot be applied twice. A key is only consulted while the +// original request may still be retried, so the retention worker (Issue #1495) +// prunes rows older than seven days to keep the table and its index bounded. +model IdempotencyKey { + id String @id @default(uuid()) + key String @unique + createdAt DateTime @default(now()) + + @@index([createdAt]) +} + // IndexerDeadLetterEvent model - quarantines Soroban events whose processing // failed (unexpected payload shape, transient DB lock, RPC timeout mid-handler). // A quarantined event is never retried inline by the poll loop, so one malformed diff --git a/backend/src/workers/idempotency-prune-worker.ts b/backend/src/workers/idempotency-prune-worker.ts new file mode 100644 index 00000000..2e5f3500 --- /dev/null +++ b/backend/src/workers/idempotency-prune-worker.ts @@ -0,0 +1,79 @@ +/** + * Idempotency key retention worker (Issue #1495) + * + * Idempotency keys are only consulted while a retry of the original request is + * still possible, so rows older than the retention window are dead weight: they + * keep growing the table (and its indexes) without ever being read again. This + * worker deletes them on a fixed interval, mirroring the stream runway worker's + * structure — an immediate first run plus a repeating timer. + */ +import { prisma } from "../lib/prisma.js"; +import logger from "../logger.js"; + +/** Keys older than seven days can no longer be replayed and are pruned. */ +export const IDEMPOTENCY_KEY_RETENTION_MS = 7 * 24 * 60 * 60 * 1000; + +/** How often the retention sweep runs. */ +export const IDEMPOTENCY_PRUNE_INTERVAL_MS = 24 * 60 * 60 * 1000; // 24 hours + +/** + * The oldest `createdAt` still inside the retention window. Rows created before + * this instant are expired. + */ +export function idempotencyRetentionCutoff(nowMs: number = Date.now()): Date { + return new Date(nowMs - IDEMPOTENCY_KEY_RETENTION_MS); +} + +/** True when a key created at `createdAt` has fallen out of the retention window. */ +export function isExpiredIdempotencyKey( + createdAt: Date, + nowMs: number = Date.now(), +): boolean { + return createdAt.getTime() < idempotencyRetentionCutoff(nowMs).getTime(); +} + +/** + * Delete every idempotency key older than the retention window and return the + * number of pruned rows so callers (and tests) can act on the result. + * + * Errors propagate: a failed sweep must not be reported as "0 pruned". + */ +export async function pruneExpiredIdempotencyKeys( + nowMs: number = Date.now(), +): Promise { + const cutoff = idempotencyRetentionCutoff(nowMs); + + const { count } = await prisma.idempotencyKey.deleteMany({ + where: { createdAt: { lt: cutoff } }, + }); + + logger.info( + `[IdempotencyPruneWorker] Pruned ${count} expired idempotency key(s) older than ${cutoff.toISOString()}`, + ); + + return count; +} + +/** + * Start the retention worker: prune once immediately, then on every interval. + * Returns the timer so `stopWorkers` can clear it on shutdown. + */ +export function startIdempotencyPruneWorker( + intervalMs: number = IDEMPOTENCY_PRUNE_INTERVAL_MS, +): NodeJS.Timeout { + const run = (phase: "Initial" | "Scheduled") => { + pruneExpiredIdempotencyKeys().catch((error) => { + logger.error(`[IdempotencyPruneWorker] ${phase} run failed:`, error); + }); + }; + + run("Initial"); + + const timer = setInterval(() => run("Scheduled"), intervalMs); + + logger.info( + `[IdempotencyPruneWorker] Worker started - pruning keys older than ${IDEMPOTENCY_KEY_RETENTION_MS}ms every ${intervalMs}ms`, + ); + + return timer; +} diff --git a/backend/src/workers/index.ts b/backend/src/workers/index.ts index 855657fb..570abcb5 100644 --- a/backend/src/workers/index.ts +++ b/backend/src/workers/index.ts @@ -7,9 +7,11 @@ import { sorobanEventWorker } from "./soroban-event-worker.js"; import { startStreamRunwayWorker } from "./stream-runway-worker.js"; +import { startIdempotencyPruneWorker } from "./idempotency-prune-worker.js"; import logger from "../logger.js"; let runwayWorkerTimer: NodeJS.Timeout | null = null; +let idempotencyPruneWorkerTimer: NodeJS.Timeout | null = null; export async function startWorkers(): Promise { logger.info("[Workers] Starting background workers..."); @@ -17,6 +19,9 @@ export async function startWorkers(): Promise { // Start stream runway alert worker (Issue #1190) runwayWorkerTimer = startStreamRunwayWorker(); + + // Start idempotency key retention worker (Issue #1495) + idempotencyPruneWorkerTimer = startIdempotencyPruneWorker(); } export function stopWorkers(): void { @@ -27,4 +32,10 @@ export function stopWorkers(): void { runwayWorkerTimer = null; logger.info("[Workers] Stream runway worker stopped"); } + + if (idempotencyPruneWorkerTimer) { + clearInterval(idempotencyPruneWorkerTimer); + idempotencyPruneWorkerTimer = null; + logger.info("[Workers] Idempotency prune worker stopped"); + } } diff --git a/backend/tests/idempotency-prune-worker.test.ts b/backend/tests/idempotency-prune-worker.test.ts new file mode 100644 index 00000000..11307f96 --- /dev/null +++ b/backend/tests/idempotency-prune-worker.test.ts @@ -0,0 +1,122 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import logger from '../src/logger.js'; +import { + IDEMPOTENCY_KEY_RETENTION_MS, + idempotencyRetentionCutoff, + isExpiredIdempotencyKey, + pruneExpiredIdempotencyKeys, + startIdempotencyPruneWorker, +} from '../src/workers/idempotency-prune-worker.js'; + +const DAY_MS = 24 * 60 * 60 * 1000; + +// In-memory stand-in for the IdempotencyKey table. `deleteMany` honours the +// `createdAt < cutoff` filter exactly like Postgres so the tests can prove +// which rows survive a sweep. +const { store, deleteMany } = vi.hoisted(() => { + const rows: { id: string; createdAt: Date }[] = []; + const deleteManyMock = vi.fn( + async (args: { where: { createdAt: { lt: Date } } }) => { + const cutoff = args.where.createdAt.lt.getTime(); + let count = 0; + for (let i = rows.length - 1; i >= 0; i--) { + const row = rows[i]!; + if (row.createdAt.getTime() < cutoff) { + rows.splice(i, 1); + count++; + } + } + return { count }; + }, + ); + return { store: rows, deleteMany: deleteManyMock }; +}); + +vi.mock('../src/lib/prisma.js', () => ({ + prisma: { + idempotencyKey: { deleteMany }, + }, +})); + +vi.mock('../src/logger.js', () => ({ + default: { + info: vi.fn(), + error: vi.fn(), + warn: vi.fn(), + }, +})); + +describe('Idempotency prune worker (#1495)', () => { + beforeEach(() => { + vi.clearAllMocks(); + store.length = 0; + }); + + it('keeps the retention window at seven days', () => { + expect(IDEMPOTENCY_KEY_RETENTION_MS).toBe(7 * DAY_MS); + }); + + it('deletes only records older than the retention window', async () => { + const now = Date.now(); + store.push( + { id: 'stale', createdAt: new Date(now - 8 * DAY_MS) }, + { id: 'just-expired', createdAt: new Date(now - 7 * DAY_MS - 1000) }, + { id: 'fresh', createdAt: new Date(now - 1 * DAY_MS) }, + { id: 'just-inside', createdAt: new Date(now - 7 * DAY_MS + 60_000) }, + ); + + const pruned = await pruneExpiredIdempotencyKeys(now); + + expect(pruned).toBe(2); + expect(store.map((row) => row.id).sort()).toEqual(['fresh', 'just-inside']); + expect(deleteMany).toHaveBeenCalledWith({ + where: { createdAt: { lt: idempotencyRetentionCutoff(now) } }, + }); + }); + + it('preserves records inside the retention window', async () => { + const now = Date.now(); + store.push( + { id: 'within-window', createdAt: new Date(now - 6 * DAY_MS) }, + { id: 'brand-new', createdAt: new Date(now) }, + ); + + const pruned = await pruneExpiredIdempotencyKeys(now); + + expect(pruned).toBe(0); + expect(store.map((row) => row.id).sort()).toEqual(['brand-new', 'within-window']); + }); + + it('logs the number of pruned records', async () => { + const now = Date.now(); + store.push({ id: 'ancient', createdAt: new Date(now - 30 * DAY_MS) }); + + await pruneExpiredIdempotencyKeys(now); + + expect(logger.info).toHaveBeenCalledWith(expect.stringContaining('Pruned 1')); + }); + + it('matches the delete boundary in isExpiredIdempotencyKey', () => { + const now = Date.now(); + expect(isExpiredIdempotencyKey(new Date(now - 8 * DAY_MS), now)).toBe(true); + expect(isExpiredIdempotencyKey(new Date(now - 6 * DAY_MS), now)).toBe(false); + }); + + it('runs an immediate sweep and schedules a repeating timer', async () => { + const now = Date.now(); + store.push({ id: 'old', createdAt: new Date(now - 10 * DAY_MS) }); + + const setIntervalSpy = vi.spyOn(global, 'setInterval').mockReturnValue(999 as unknown as NodeJS.Timeout); + + const timer = startIdempotencyPruneWorker(60_000); + await Promise.resolve(); + + expect(timer).toBe(999); + expect(setIntervalSpy).toHaveBeenCalledWith(expect.any(Function), 60_000); + expect(deleteMany).toHaveBeenCalledTimes(1); + expect(store).toHaveLength(0); + + clearInterval(timer); + setIntervalSpy.mockRestore(); + }); +});