Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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");
12 changes: 12 additions & 0 deletions backend/prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,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
Expand Down
79 changes: 79 additions & 0 deletions backend/src/workers/idempotency-prune-worker.ts
Original file line number Diff line number Diff line change
@@ -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<number> {
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;
}
11 changes: 11 additions & 0 deletions backend/src/workers/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,16 +7,21 @@

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<void> {
logger.info("[Workers] Starting background workers...");
await sorobanEventWorker.start();

// Start stream runway alert worker (Issue #1190)
runwayWorkerTimer = startStreamRunwayWorker();

// Start idempotency key retention worker (Issue #1495)
idempotencyPruneWorkerTimer = startIdempotencyPruneWorker();
}

export function stopWorkers(): void {
Expand All @@ -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");
}
}
122 changes: 122 additions & 0 deletions backend/tests/idempotency-prune-worker.test.ts
Original file line number Diff line number Diff line change
@@ -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();
});
});
Loading