From 1244ccae5e8975dd32560e83c914d0b6cca6dad4 Mon Sep 17 00:00:00 2001 From: nasalehj Date: Thu, 20 Aug 2026 23:31:33 +0000 Subject: [PATCH] fix: add processing lease to idempotency store Processing records previously held a key forever: if the process died between acquire and complete, every retry returned 409 until a manual cleanup. Add a lease_expires_at column refreshed by a middleware heartbeat while the handler runs, atomically re-acquire abandoned leases, and persist the result before delivering the response so a replay never observes processing. --- app/backend/src/idempotency/middleware.ts | 43 ++++- app/backend/src/idempotency/store.ts | 95 ++++++++++- app/backend/test/coverage-baseline.json | 4 +- app/backend/test/idempotency.spec.ts | 185 +++++++++++++++++++++- 4 files changed, 307 insertions(+), 20 deletions(-) diff --git a/app/backend/src/idempotency/middleware.ts b/app/backend/src/idempotency/middleware.ts index ac36faed..4ea5f814 100644 --- a/app/backend/src/idempotency/middleware.ts +++ b/app/backend/src/idempotency/middleware.ts @@ -21,27 +21,60 @@ export function idempotencyMiddleware(store: IdempotencyStore) { const existingRecord = await store.tryAcquire(key, fingerprint); if (!existingRecord) { - // First time: Intercept the response to cache it + // First time (or an abandoned lease was re-acquired): intercept the + // response to cache it, and keep the processing lease alive while the + // handler runs so a slow request is never mistaken for an abandoned + // one. const originalSend = res.send.bind(res); + const heartbeatMs = Math.max(1, Math.floor(store.leaseDurationMs / 2)); + const heartbeat = setInterval(() => { + store + .heartbeat(key) + .catch(err => + console.error( + `Failed to refresh idempotency lease for key ${key.asString()}:`, + err, + ), + ); + }, heartbeatMs); + heartbeat.unref?.(); + + const stopHeartbeat = () => clearInterval(heartbeat); + res.once('close', stopHeartbeat); + + // Persist the result BEFORE the response is delivered, so a retry with + // the same key can never observe `processing` after the first request + // has responded. The write is deferred until the record is committed; + // failures are logged, never thrown at the client. res.send = (body: any) => { + stopHeartbeat(); const status = res.statusCode; const recordStatus = status >= 200 && status < 300 ? 'succeeded' : 'failed'; const bodyString = typeof body === 'string' ? body : JSON.stringify(body); - // Fire and forget cache save (log on failure) - store + void store .complete(key, recordStatus, status, bodyString) .catch(err => console.error( `Failed to save idempotency record for key ${key.asString()}:`, err, ), - ); + ) + .finally(() => { + try { + originalSend(body); + } catch (err) { + console.error( + `Failed to deliver response for key ${key.asString()}:`, + err, + ); + } + }); - return originalSend(body); + return res; }; return next(); diff --git a/app/backend/src/idempotency/store.ts b/app/backend/src/idempotency/store.ts index 0e8c36ce..b41cb62a 100644 --- a/app/backend/src/idempotency/store.ts +++ b/app/backend/src/idempotency/store.ts @@ -10,34 +10,90 @@ export interface IdempotencyRecord { status: RecordStatus; responseBody: Buffer | null; responseStatus: number | null; + leaseExpiresAt: Date | null; } +export interface IdempotencyStoreOptions { + /** + * How long a `processing` record may hold the key before it is considered + * abandoned and becomes re-acquirable. The middleware refreshes this lease + * with `heartbeat()` while the handler runs. + */ + leaseDurationMs?: number; +} + +const DEFAULT_LEASE_DURATION_MS = 30_000; + +/** + * SQL fragment that computes a lease expiry from `now()`. The lease duration + * (milliseconds) is passed as the query parameter at `paramIndex` — callers + * must pass it as the last parameter of their statement. + */ +const leaseExpirySql = (paramIndex: number) => + `now() + make_interval(secs => $${paramIndex}::float8 / 1000.0)`; + export class IdempotencyStore { private pool: Pool; - constructor(pool: Pool) { + /** Lease duration for `processing` records, in milliseconds. */ + public readonly leaseDurationMs: number; + + constructor(pool: Pool, options: IdempotencyStoreOptions = {}) { this.pool = pool; + this.leaseDurationMs = options.leaseDurationMs ?? DEFAULT_LEASE_DURATION_MS; } + /** + * Attempts to acquire the idempotency key. + * + * - Returns `undefined` when the caller may proceed (fresh key, or an + * abandoned `processing` lease was atomically re-acquired). + * - Returns the existing record otherwise; the caller must decide between + * replaying a cached response and returning 409 for a live `processing` + * record. + */ public async tryAcquire( key: IdempotencyKey, fingerprint: RequestFingerprint, ): Promise { const insertResult = await this.pool.query( - `INSERT INTO idempotency_records (idempotency_key, request_fingerprint, status) - VALUES ($1, $2, 'processing') + `INSERT INTO idempotency_records (idempotency_key, request_fingerprint, status, lease_expires_at) + VALUES ($1, $2, 'processing', ${leaseExpirySql(3)}) ON CONFLICT (idempotency_key) DO NOTHING RETURNING idempotency_key`, - [key.asString(), fingerprint.asString()], + [key.asString(), fingerprint.asString(), this.leaseDurationMs], ); if (insertResult.rows.length > 0) { return undefined; // Fresh key — proceed! } - // Key exists — fetch it + // Key exists. Atomically claim it if the previous lease has expired: the + // row-level lock guarantees only one concurrent request can re-acquire an + // abandoned `processing` record, so an abandoned lease can never cause the + // handler to run twice concurrently. + const claimResult = await this.pool.query( + `UPDATE idempotency_records + SET request_fingerprint = $2, + status = 'processing', + response_body = NULL, + response_status = NULL, + lease_expires_at = ${leaseExpirySql(3)}, + updated_at = now() + WHERE idempotency_key = $1 + AND status = 'processing' + AND (lease_expires_at IS NULL OR lease_expires_at <= now()) + RETURNING idempotency_key`, + [key.asString(), fingerprint.asString(), this.leaseDurationMs], + ); + + if (claimResult.rows.length > 0) { + return undefined; // Abandoned lease re-acquired — proceed! + } + + // Key exists with a live lease or a terminal status — fetch it const { rows } = await this.pool.query( - `SELECT idempotency_key, request_fingerprint, status, response_body, response_status + `SELECT idempotency_key, request_fingerprint, status, response_body, response_status, lease_expires_at FROM idempotency_records WHERE idempotency_key = $1`, [key.asString()], ); @@ -49,9 +105,27 @@ export class IdempotencyStore { status: row.status, responseBody: row.response_body, responseStatus: row.response_status, + leaseExpiresAt: row.lease_expires_at, }; } + /** + * Refreshes the lease on a `processing` record. Called periodically by the + * middleware while the handler runs so a slow request is never mistaken for + * an abandoned one. A lease that has already expired is left untouched so a + * crashed request's record can still be re-acquired. + */ + public async heartbeat(key: IdempotencyKey): Promise { + await this.pool.query( + `UPDATE idempotency_records + SET lease_expires_at = ${leaseExpirySql(2)}, updated_at = now() + WHERE idempotency_key = $1 + AND status = 'processing' + AND (lease_expires_at IS NULL OR lease_expires_at > now())`, + [key.asString(), this.leaseDurationMs], + ); + } + public async complete( key: IdempotencyKey, status: RecordStatus, @@ -60,8 +134,13 @@ export class IdempotencyStore { ): Promise { await this.pool.query( `UPDATE idempotency_records - SET status = $2, response_status = $3, response_body = $4, updated_at = now() - WHERE idempotency_key = $1`, + SET status = $2, + response_status = $3, + response_body = $4, + lease_expires_at = NULL, + updated_at = now() + WHERE idempotency_key = $1 + AND status = 'processing'`, [key.asString(), status, responseStatus, Buffer.from(responseBody)], ); } diff --git a/app/backend/test/coverage-baseline.json b/app/backend/test/coverage-baseline.json index 5bb6b226..3101de99 100644 --- a/app/backend/test/coverage-baseline.json +++ b/app/backend/test/coverage-baseline.json @@ -60,8 +60,8 @@ ["src/idempotency/error.ts",-12,100,-5,-12], ["src/idempotency/fingerprint.ts",-16,-6,-6,-17], ["src/idempotency/key.ts",-17,-8,-3,-17], - ["src/idempotency/middleware.ts",-32,-18,-4,-32], - ["src/idempotency/store.ts",-11,-4,-4,-11], + ["src/idempotency/middleware.ts",-43,-18,-8,-44], + ["src/idempotency/store.ts",-19,-9,-6,-19], ["src/interceptors/idempotency.interceptor.ts",-19,-8,-4,-21], ["src/interceptors/logging.interceptor.ts",100,-1,100,100], ["src/jobs/dlq.service.ts",-6,-7,-1,-6], diff --git a/app/backend/test/idempotency.spec.ts b/app/backend/test/idempotency.spec.ts index 4dba36f5..15e2934e 100644 --- a/app/backend/test/idempotency.spec.ts +++ b/app/backend/test/idempotency.spec.ts @@ -1,10 +1,11 @@ import 'dotenv/config'; import { describe, it, expect, beforeAll, afterAll } from '@jest/globals'; import request from 'supertest'; -import express from 'express'; +import express, { Request } from 'express'; import { Pool } from 'pg'; import { IdempotencyStore } from '../src/idempotency/store'; +import { IdempotencyKey } from '../src/idempotency/key'; import { idempotencyMiddleware } from '../src/idempotency/middleware'; import { submitTransaction } from '../src/handlers/transaction'; import { RequestFingerprint } from '../src/idempotency/fingerprint'; @@ -34,7 +35,8 @@ const validBody = { transactionXdr: 'AAAAAAABLC0=' }; status TEXT NOT NULL DEFAULT 'processing', response_body BYTEA, response_status SMALLINT, - created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + lease_expires_at TIMESTAMPTZ DEFAULT now() + interval '30 seconds', + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); `); @@ -114,7 +116,7 @@ const validBody = { transactionXdr: 'AAAAAAABLC0=' }; expect(res.body.error).toContain('fingerprint'); }); - it('Processing record returns 409', async () => { + it('Live processing record returns 409', async () => { const validFingerprint = RequestFingerprint.fromBody(validBody).asString(); @@ -123,9 +125,10 @@ const validBody = { transactionXdr: 'AAAAAAABLC0=' }; INSERT INTO idempotency_records ( idempotency_key, request_fingerprint, - status + status, + lease_expires_at ) - VALUES ($1, $2, 'processing') + VALUES ($1, $2, 'processing', now() + interval '1 minute') `, ['key-4', validFingerprint], ); @@ -139,6 +142,178 @@ const validBody = { transactionXdr: 'AAAAAAABLC0=' }; expect(res.body.error).toContain('processed'); }); + it('Expired processing lease is re-acquired instead of stuck at 409', async () => { + const validFingerprint = + RequestFingerprint.fromBody(validBody).asString(); + + // Simulate a crash: the record is `processing` with an expired lease, + // exactly what is left behind when the process dies between acquire and + // complete. + await pool.query( + ` + INSERT INTO idempotency_records ( + idempotency_key, + request_fingerprint, + status, + lease_expires_at + ) + VALUES ($1, $2, 'processing', now() - interval '1 minute') + `, + ['key-crash', validFingerprint], + ); + + const res = await request(app) + .post('/v1/transactions/submit') + .set('Idempotency-Key', 'key-crash') + .send(validBody); + + expect(res.status).toBe(200); + expect(res.headers['x-idempotent-replayed']).toBeUndefined(); + + const { rows } = await pool.query( + `SELECT status, lease_expires_at + FROM idempotency_records + WHERE idempotency_key = $1`, + ['key-crash'], + ); + expect(rows[0].status).toBe('succeeded'); + expect(rows[0].lease_expires_at).toBeNull(); + }); + + it('heartbeat refreshes the lease on a live processing record', async () => { + const key = IdempotencyKey.fromHeaders({ + headers: { 'idempotency-key': 'key-heartbeat' }, + } as unknown as Request); + const fingerprint = RequestFingerprint.fromBody(validBody); + + // A live record with a lease that is still valid but about to expire. + await pool.query( + ` + INSERT INTO idempotency_records ( + idempotency_key, + request_fingerprint, + status, + lease_expires_at + ) + VALUES ($1, $2, 'processing', now() + interval '5 seconds') + `, + [key.asString(), fingerprint.asString()], + ); + + const before = await pool.query( + `SELECT lease_expires_at + FROM idempotency_records + WHERE idempotency_key = $1`, + [key.asString()], + ); + const beforeMs = (before.rows[0].lease_expires_at as Date).getTime(); + + await store.heartbeat(key); + + const after = await pool.query( + `SELECT lease_expires_at + FROM idempotency_records + WHERE idempotency_key = $1`, + [key.asString()], + ); + const afterMs = (after.rows[0].lease_expires_at as Date).getTime(); + + // The lease must be pushed out to roughly now + leaseDurationMs (30s), + // not left to expire while the handler is still running. + expect(afterMs).toBeGreaterThan(beforeMs + 5_000); + expect(afterMs).toBeGreaterThan(Date.now() + 10_000); + }); + + it('Concurrent request while the first is processing returns 409, then replays', async () => { + let started!: () => void; + let release!: () => void; + const startedPromise = new Promise(resolve => { + started = resolve; + }); + const gate = new Promise(resolve => { + release = resolve; + }); + + app.post( + '/v1/gated', + idempotencyMiddleware(store), + async (_req: Request, res: express.Response) => { + started(); + await gate; + res.status(200).json({ gated: true }); + }, + ); + + // Dispatch the first request without awaiting it, so the handler runs + // (and holds a live lease) while we probe concurrency with `second`. + const first = request(app) + .post('/v1/gated') + .set('Idempotency-Key', 'key-gated') + .send(validBody) + .then(res => res); + + // Wait until the first handler is actually in flight (lease live). + await startedPromise; + + const second = await request(app) + .post('/v1/gated') + .set('Idempotency-Key', 'key-gated') + .send(validBody); + + expect(second.status).toBe(409); + expect(second.body.error).toContain('processed'); + + release(); + const firstRes = await first; + expect(firstRes.status).toBe(200); + + const third = await request(app) + .post('/v1/gated') + .set('Idempotency-Key', 'key-gated') + .send(validBody); + + expect(third.status).toBe(200); + expect(third.headers['x-idempotent-replayed']).toBe('true'); + expect(third.body.gated).toBe(true); + }); + + it('Concurrent tryAcquire on an expired lease: exactly one request wins', async () => { + const key = IdempotencyKey.fromHeaders({ + headers: { 'idempotency-key': 'key-race' }, + } as unknown as Request); + const fingerprint = RequestFingerprint.fromBody(validBody); + + await pool.query( + ` + INSERT INTO idempotency_records ( + idempotency_key, + request_fingerprint, + status, + lease_expires_at + ) + VALUES ($1, $2, 'processing', now() - interval '1 minute') + `, + ['key-race', fingerprint.asString()], + ); + + const results = await Promise.all([ + store.tryAcquire(key, fingerprint), + store.tryAcquire(key, fingerprint), + ]); + + // An abandoned lease must never let two handlers run concurrently. + const winners = results.filter(r => r === undefined).length; + expect(winners).toBe(1); + + const record = results.find(r => r !== undefined); + expect(record).toBeDefined(); + expect(record!.status).toBe('processing'); + expect(record!.leaseExpiresAt).not.toBeNull(); + expect((record!.leaseExpiresAt as Date).getTime()).toBeGreaterThan( + Date.now(), + ); + }); + it('GET /v1/transactions/:hash returns 404', async () => { const res = await request(app).get('/v1/transactions/some-hash');