diff --git a/__tests__/idempotency/middleware.test.ts b/__tests__/idempotency/middleware.test.ts new file mode 100644 index 0000000..3b9eeea --- /dev/null +++ b/__tests__/idempotency/middleware.test.ts @@ -0,0 +1,146 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest' +import { NextRequest, NextResponse } from 'next/server' + +import { withIdempotency } from '@/lib/idempotency/middleware' +import { + IdempotencyInProgressError, + IdempotencyKeyReusedError, +} from '@/lib/idempotency/errors' + +vi.mock('@/lib/idempotency/service', () => ({ + idempotencyService: { run: vi.fn() }, + IdempotencyService: class {}, +})) + +import { idempotencyService } from '@/lib/idempotency/service' + +const mockRun = vi.mocked(idempotencyService.run) +const KEY = 'test-key-1234567890abcd' +const AUTH = { walletAddress: 'GABC' } + +function makeRequest(options: { + body?: unknown + key?: string + rawBody?: string +}): NextRequest { + const headers: Record = { 'content-type': 'application/json' } + if (options.key) headers['Idempotency-Key'] = options.key + return new NextRequest('http://localhost/api/escrow/fund', { + method: 'POST', + headers, + body: options.rawBody ?? JSON.stringify(options.body ?? {}), + }) +} + +beforeEach(() => { + vi.clearAllMocks() +}) + +describe('withIdempotency', () => { + it('returns 400 KEY_REQUIRED when no key is supplied', async () => { + const handler = vi.fn() + const wrapped = withIdempotency('escrow_fund', handler) + + const res = await wrapped(makeRequest({ body: { contractId: 'c-1' } }), AUTH) + + expect(res.status).toBe(400) + expect((await res.json()).code).toBe('IDEMPOTENCY_KEY_REQUIRED') + expect(mockRun).not.toHaveBeenCalled() + expect(handler).not.toHaveBeenCalled() + }) + + it('returns 400 KEY_INVALID for a malformed key', async () => { + const wrapped = withIdempotency('escrow_fund', vi.fn()) + + const res = await wrapped( + makeRequest({ body: { contractId: 'c-1' }, key: 'short' }), + AUTH + ) + + expect(res.status).toBe(400) + expect((await res.json()).code).toBe('IDEMPOTENCY_KEY_INVALID') + expect(mockRun).not.toHaveBeenCalled() + }) + + it('returns 400 INVALID_JSON for an unparsable body', async () => { + const wrapped = withIdempotency('escrow_fund', vi.fn()) + + const res = await wrapped( + makeRequest({ rawBody: 'not-json', key: KEY }), + AUTH + ) + + expect(res.status).toBe(400) + expect((await res.json()).code).toBe('INVALID_JSON') + expect(mockRun).not.toHaveBeenCalled() + }) + + it('executes the handler and passes the parsed body', async () => { + mockRun.mockImplementation(async (params) => ({ + ...(await params.handler()), + replayed: false, + })) + const handler = vi.fn(async () => NextResponse.json({ contractId: 'c-1' }, { status: 200 })) + const wrapped = withIdempotency('escrow_fund', handler) + + const res = await wrapped( + makeRequest({ body: { contractId: 'c-1' }, key: KEY }), + AUTH + ) + + expect(res.status).toBe(200) + expect(await res.json()).toEqual({ contractId: 'c-1' }) + expect(res.headers.get('Idempotency-Replayed')).toBe('false') + expect(handler).toHaveBeenCalledTimes(1) + expect(handler.mock.calls[0][2]).toEqual({ contractId: 'c-1' }) + expect(mockRun).toHaveBeenCalledWith( + expect.objectContaining({ key: KEY, operationType: 'escrow_fund' }) + ) + }) + + it('accepts the key from the request body', async () => { + mockRun.mockImplementation(async (params) => ({ + ...(await params.handler()), + replayed: false, + })) + const wrapped = withIdempotency('escrow_fund', vi.fn(async () => NextResponse.json({ ok: true }))) + + await wrapped(makeRequest({ body: { contractId: 'c-1', idempotencyKey: KEY } }), AUTH) + + expect(mockRun).toHaveBeenCalledWith(expect.objectContaining({ key: KEY })) + }) + + it('replays a stored response and marks the header', async () => { + mockRun.mockResolvedValue({ status: 201, body: { contractId: 'c-1' }, replayed: true }) + const handler = vi.fn() + const wrapped = withIdempotency('escrow_fund', handler) + + const res = await wrapped(makeRequest({ body: { contractId: 'c-1' }, key: KEY }), AUTH) + + expect(res.status).toBe(201) + expect(await res.json()).toEqual({ contractId: 'c-1' }) + expect(res.headers.get('Idempotency-Replayed')).toBe('true') + expect(handler).not.toHaveBeenCalled() + }) + + it('returns 409 with Retry-After while a duplicate is in progress', async () => { + mockRun.mockRejectedValue(new IdempotencyInProgressError()) + const wrapped = withIdempotency('escrow_fund', vi.fn()) + + const res = await wrapped(makeRequest({ body: { contractId: 'c-1' }, key: KEY }), AUTH) + + expect(res.status).toBe(409) + expect((await res.json()).code).toBe('IDEMPOTENCY_IN_PROGRESS') + expect(res.headers.get('Retry-After')).toBe('1') + }) + + it('returns 409 when a key is reused with a different payload', async () => { + mockRun.mockRejectedValue(new IdempotencyKeyReusedError()) + const wrapped = withIdempotency('escrow_fund', vi.fn()) + + const res = await wrapped(makeRequest({ body: { contractId: 'c-2' }, key: KEY }), AUTH) + + expect(res.status).toBe(409) + expect((await res.json()).code).toBe('IDEMPOTENCY_KEY_REUSED') + }) +}) diff --git a/__tests__/idempotency/repository.test.ts b/__tests__/idempotency/repository.test.ts new file mode 100644 index 0000000..2cd83fd --- /dev/null +++ b/__tests__/idempotency/repository.test.ts @@ -0,0 +1,114 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest' + +vi.mock('@/lib/db', () => ({ + sql: vi.fn(), +})) + +import { sql } from '@/lib/db' +import { IdempotencyRepository } from '@/lib/idempotency/repository' + +const mockSql = sql as unknown as ReturnType + +function row(overrides: Record = {}) { + return { + id: 'r-1', + idempotency_key: 'key-1234567890abcdef', + operation_type: 'escrow_fund', + request_hash: 'hash-1', + request_payload: { contractId: 'c-1' }, + response_payload: null, + response_status: null, + status: 'in_progress', + created_at: '2026-01-01T00:00:00.000Z', + updated_at: '2026-01-01T00:00:00.000Z', + expires_at: '2026-01-02T00:00:00.000Z', + ...overrides, + } +} + +beforeEach(() => { + vi.clearAllMocks() +}) + +describe('IdempotencyRepository', () => { + const repo = new IdempotencyRepository() + + it('claims a new key when the insert returns a row', async () => { + mockSql.mockResolvedValueOnce([row()]) + + const result = await repo.claim({ + key: 'key-1234567890abcdef', + operationType: 'escrow_fund', + requestHash: 'hash-1', + requestPayload: { contractId: 'c-1' }, + ttlHours: 24, + }) + + expect(result.claimed).toBe(true) + expect(result.record?.status).toBe('in_progress') + expect(mockSql).toHaveBeenCalledTimes(1) + }) + + it('falls back to the existing record when the insert conflicts', async () => { + mockSql + .mockResolvedValueOnce([]) // ON CONFLICT … DO UPDATE WHERE → no row + .mockResolvedValueOnce([row({ status: 'completed', response_status: 200, response_payload: { ok: true } })]) + + const result = await repo.claim({ + key: 'key-1234567890abcdef', + operationType: 'escrow_fund', + requestHash: 'hash-1', + requestPayload: { contractId: 'c-1' }, + ttlHours: 24, + }) + + expect(result.claimed).toBe(false) + expect(result.record?.status).toBe('completed') + expect(result.record?.responseStatus).toBe(200) + expect(result.record?.responsePayload).toEqual({ ok: true }) + expect(mockSql).toHaveBeenCalledTimes(2) + }) + + it('parses string-serialised JSON payloads', async () => { + mockSql.mockResolvedValueOnce([ + row({ request_payload: '{"contractId":"c-9"}', response_payload: '{"ok":true}' }), + ]) + + const record = await repo.get('key-1234567890abcdef', 'escrow_fund') + + expect(record?.requestPayload).toEqual({ contractId: 'c-9' }) + expect(record?.responsePayload).toEqual({ ok: true }) + }) + + it('returns null when no record exists', async () => { + mockSql.mockResolvedValueOnce([]) + expect(await repo.get('missing-key-123456', 'escrow_fund')).toBeNull() + }) + + it('marks a completed record and returns it', async () => { + mockSql.mockResolvedValueOnce([row({ status: 'completed', response_status: 201, response_payload: { ok: true } })]) + + const record = await repo.complete( + 'key-1234567890abcdef', + 'escrow_fund', + { ok: true }, + 201, + 24 + ) + + expect(record?.status).toBe('completed') + expect(record?.responseStatus).toBe(201) + }) + + it('purges expired records and returns the count', async () => { + mockSql.mockResolvedValueOnce([{ id: 'r-1' }, { id: 'r-2' }]) + expect(await repo.purgeExpired()).toBe(2) + }) + + it('marks failed and removes records', async () => { + mockSql.mockResolvedValue([]) + await repo.markFailed('key-1234567890abcdef', 'escrow_fund') + await repo.remove('key-1234567890abcdef', 'escrow_fund') + expect(mockSql).toHaveBeenCalledTimes(2) + }) +}) diff --git a/__tests__/idempotency/service.test.ts b/__tests__/idempotency/service.test.ts new file mode 100644 index 0000000..4435625 --- /dev/null +++ b/__tests__/idempotency/service.test.ts @@ -0,0 +1,202 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest' + +import { IdempotencyService } from '@/lib/idempotency/service' +import { + IdempotencyInProgressError, + IdempotencyKeyReusedError, + IdempotencyStorageError, +} from '@/lib/idempotency/errors' +import { hashRequestPayload } from '@/lib/idempotency/validation' +import type { + IdempotencyOperationType, + IdempotencyRecord, + IIdempotencyRepository, +} from '@/lib/idempotency/types' + +const OPERATION: IdempotencyOperationType = 'escrow_fund' +const KEY = 'test-key-1234567890abcd' + +function makeRecord(overrides: Partial = {}): IdempotencyRecord { + return { + id: 'record-1', + idempotencyKey: KEY, + operationType: OPERATION, + requestHash: hashRequestPayload({ contractId: 'c-1' }), + requestPayload: { contractId: 'c-1' }, + responsePayload: null, + responseStatus: null, + status: 'in_progress', + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + expiresAt: new Date(Date.now() + 3_600_000).toISOString(), + ...overrides, + } +} + +function makeRepo(): IIdempotencyRepository { + return { + claim: vi.fn(), + get: vi.fn().mockResolvedValue(null), + complete: vi.fn().mockResolvedValue(null), + markFailed: vi.fn().mockResolvedValue(undefined), + remove: vi.fn().mockResolvedValue(undefined), + purgeExpired: vi.fn().mockResolvedValue(0), + } +} + +describe('IdempotencyService.run', () => { + let repo: IIdempotencyRepository + let service: IdempotencyService + + beforeEach(() => { + repo = makeRepo() + service = new IdempotencyService(repo, () => 24) + }) + + it('executes and stores the response for a new key', async () => { + vi.mocked(repo.claim).mockResolvedValue({ claimed: true, record: makeRecord() }) + const handler = vi.fn().mockResolvedValue({ status: 200, body: { contractId: 'c-1' } }) + + const outcome = await service.run({ + key: KEY, + operationType: OPERATION, + requestPayload: { contractId: 'c-1' }, + handler, + }) + + expect(handler).toHaveBeenCalledTimes(1) + expect(outcome).toEqual({ status: 200, body: { contractId: 'c-1' }, replayed: false }) + expect(repo.complete).toHaveBeenCalledWith( + KEY, + OPERATION, + { contractId: 'c-1' }, + 200, + 24 + ) + expect(repo.markFailed).not.toHaveBeenCalled() + }) + + it('replays the stored response without re-running the handler', async () => { + vi.mocked(repo.claim).mockResolvedValue({ + claimed: false, + record: makeRecord({ + status: 'completed', + responseStatus: 201, + responsePayload: { contractId: 'c-1', deployTxHash: 'tx-1' }, + }), + }) + const handler = vi.fn() + + const outcome = await service.run({ + key: KEY, + operationType: OPERATION, + requestPayload: { contractId: 'c-1' }, + handler, + }) + + expect(handler).not.toHaveBeenCalled() + expect(outcome).toEqual({ + status: 201, + body: { contractId: 'c-1', deployTxHash: 'tx-1' }, + replayed: true, + }) + }) + + it('rejects a key reused with a different payload', async () => { + vi.mocked(repo.claim).mockResolvedValue({ + claimed: false, + record: makeRecord({ status: 'completed', requestHash: hashRequestPayload({ contractId: 'OTHER' }) }), + }) + + await expect( + service.run({ + key: KEY, + operationType: OPERATION, + requestPayload: { contractId: 'c-1' }, + handler: vi.fn(), + }) + ).rejects.toBeInstanceOf(IdempotencyKeyReusedError) + }) + + it('rejects a duplicate request while the original is in progress', async () => { + vi.mocked(repo.claim).mockResolvedValue({ + claimed: false, + record: makeRecord({ status: 'in_progress' }), + }) + + await expect( + service.run({ + key: KEY, + operationType: OPERATION, + requestPayload: { contractId: 'c-1' }, + handler: vi.fn(), + }) + ).rejects.toBeInstanceOf(IdempotencyInProgressError) + }) + + it('does not cache a non-2xx response and releases the claim', async () => { + vi.mocked(repo.claim).mockResolvedValue({ claimed: true, record: makeRecord() }) + const handler = vi.fn().mockResolvedValue({ + status: 409, + body: { error: 'conflict', code: 'ESCROW_INVALID_STATE' }, + }) + + const outcome = await service.run({ + key: KEY, + operationType: OPERATION, + requestPayload: { contractId: 'c-1' }, + handler, + }) + + expect(outcome.status).toBe(409) + expect(outcome.replayed).toBe(false) + expect(repo.complete).not.toHaveBeenCalled() + expect(repo.markFailed).toHaveBeenCalledWith(KEY, OPERATION) + }) + + it('releases the claim when the handler throws', async () => { + vi.mocked(repo.claim).mockResolvedValue({ claimed: true, record: makeRecord() }) + const handler = vi.fn().mockRejectedValue(new Error('boom')) + + await expect( + service.run({ + key: KEY, + operationType: OPERATION, + requestPayload: { contractId: 'c-1' }, + handler, + }) + ).rejects.toThrow('boom') + + expect(repo.markFailed).toHaveBeenCalledWith(KEY, OPERATION) + }) + + it('wraps store failures in IdempotencyStorageError', async () => { + vi.mocked(repo.claim).mockRejectedValue(new Error('db down')) + + await expect( + service.run({ + key: KEY, + operationType: OPERATION, + requestPayload: { contractId: 'c-1' }, + handler: vi.fn(), + }) + ).rejects.toBeInstanceOf(IdempotencyStorageError) + }) + + it('uses the configured TTL when storing a response', async () => { + const customService = new IdempotencyService(repo, () => 48) + vi.mocked(repo.claim).mockResolvedValue({ claimed: true, record: makeRecord() }) + + await customService.run({ + key: KEY, + operationType: OPERATION, + requestPayload: { contractId: 'c-1' }, + handler: vi.fn().mockResolvedValue({ status: 200, body: { ok: true } }), + }) + + expect(repo.claim).toHaveBeenCalledWith( + expect.objectContaining({ ttlHours: 48 }) + ) + expect(repo.complete).toHaveBeenCalledWith(KEY, OPERATION, { ok: true }, 200, 48) + }) +}) diff --git a/__tests__/idempotency/validation.test.ts b/__tests__/idempotency/validation.test.ts new file mode 100644 index 0000000..763b893 --- /dev/null +++ b/__tests__/idempotency/validation.test.ts @@ -0,0 +1,105 @@ +import { describe, it, expect } from 'vitest' +import { NextRequest } from 'next/server' + +import { + extractIdempotencyKey, + validateIdempotencyKey, + hashRequestPayload, + stableStringify, +} from '@/lib/idempotency/validation' +import { + IdempotencyKeyInvalidError, + IdempotencyKeyRequiredError, +} from '@/lib/idempotency/errors' + +const VALID_KEY = 'a4f8c2e0-6b1d-4f3a-9c7e-1234567890ab' + +function makeRequest(headers: Record = {}): NextRequest { + return new NextRequest('http://localhost/api/escrow/fund', { + method: 'POST', + headers, + }) +} + +describe('extractIdempotencyKey', () => { + it('prefers the Idempotency-Key header', () => { + const request = makeRequest({ 'Idempotency-Key': 'header-key-1234567890' }) + expect(extractIdempotencyKey(request, { idempotencyKey: 'body-key-1234567890' })).toBe( + 'header-key-1234567890' + ) + }) + + it('falls back to the body field', () => { + const request = makeRequest() + expect(extractIdempotencyKey(request, { idempotencyKey: 'body-key-1234567890' })).toBe( + 'body-key-1234567890' + ) + }) + + it('returns null when no key is provided', () => { + expect(extractIdempotencyKey(makeRequest(), {})).toBeNull() + }) + + it('trims whitespace around a key', () => { + const request = makeRequest({ 'Idempotency-Key': ' spaced-key-1234567890 ' }) + expect(extractIdempotencyKey(request)).toBe('spaced-key-1234567890') + }) +}) + +describe('validateIdempotencyKey', () => { + it('accepts a UUID key', () => { + expect(validateIdempotencyKey(VALID_KEY)).toBe(VALID_KEY) + }) + + it('accepts a hex hash key', () => { + const hash = 'f'.repeat(64) + expect(validateIdempotencyKey(hash)).toBe(hash) + }) + + it('throws KEY_REQUIRED when missing', () => { + expect(() => validateIdempotencyKey(null)).toThrow(IdempotencyKeyRequiredError) + expect(() => validateIdempotencyKey('')).toThrow(IdempotencyKeyRequiredError) + try { + validateIdempotencyKey(' ') + } catch (err) { + expect((err as IdempotencyKeyRequiredError).code).toBe('IDEMPOTENCY_KEY_REQUIRED') + expect((err as IdempotencyKeyRequiredError).status).toBe(400) + } + }) + + it('throws KEY_INVALID when too short', () => { + expect(() => validateIdempotencyKey('short')).toThrow(IdempotencyKeyInvalidError) + }) + + it('throws KEY_INVALID on illegal characters', () => { + expect(() => validateIdempotencyKey('invalid key with spaces!!')).toThrow( + IdempotencyKeyInvalidError + ) + }) + + it('throws KEY_INVALID when too long', () => { + expect(() => validateIdempotencyKey('a'.repeat(256))).toThrow( + IdempotencyKeyInvalidError + ) + }) +}) + +describe('stableStringify / hashRequestPayload', () => { + it('produces the same hash regardless of key order', () => { + const a = { b: 1, a: { d: 2, c: 3 } } + const b = { a: { c: 3, d: 2 }, b: 1 } + expect(stableStringify(a)).toBe(stableStringify(b)) + expect(hashRequestPayload(a)).toBe(hashRequestPayload(b)) + }) + + it('produces different hashes for different payloads', () => { + expect(hashRequestPayload({ amount: '10' })).not.toBe( + hashRequestPayload({ amount: '20' }) + ) + }) + + it('handles arrays, null and primitives', () => { + expect(stableStringify([1, null, 'x'])).toBe('[1,null,"x"]') + expect(typeof hashRequestPayload(null)).toBe('string') + }) +}) diff --git a/app/api/escrow/create/route.ts b/app/api/escrow/create/route.ts index f6ff67a..eeb3d5b 100644 --- a/app/api/escrow/create/route.ts +++ b/app/api/escrow/create/route.ts @@ -4,6 +4,12 @@ * Deploy a new escrow contract for a project. * The authenticated user must be the client (project owner). * + * Idempotency: + * Requires an "Idempotency-Key" header (or `idempotencyKey` body field) — + * a UUID or cryptographically random hash. Repeating the request with the + * same key returns the original response instead of deploying a second + * contract. + * * Body: * projectId string (UUID) * freelancerId string (UUID) @@ -15,8 +21,9 @@ */ import { NextRequest, NextResponse } from 'next/server' -import { withAuth } from '@/lib/auth/middleware' +import { withAuth, AuthContext } from '@/lib/auth/middleware' import { sql } from '@/lib/db' +import { withIdempotency } from '@/lib/idempotency' import { escrowService, EscrowError, @@ -24,17 +31,11 @@ import { EscrowAlreadyExistsError, } from '@/lib/escrow' -export const POST = withAuth(async (request: NextRequest, auth) => { - let body: Record - try { - body = await request.json() - } catch { - return NextResponse.json( - { error: 'Request body must be valid JSON', code: 'INVALID_JSON' }, - { status: 400 } - ) - } - +export const POST = withAuth(withIdempotency('escrow_create', async ( + _request: NextRequest, + auth: AuthContext, + body: Record +) => { // --- Resolve authenticated wallet to a DB user --- const users = await sql` SELECT id FROM users WHERE wallet_address = ${auth.walletAddress} LIMIT 1 @@ -103,4 +104,4 @@ export const POST = withAuth(async (request: NextRequest, auth) => { { status: 500 } ) } -}) +})) diff --git a/app/api/escrow/fund/route.ts b/app/api/escrow/fund/route.ts index 6fc6bf0..28746e9 100644 --- a/app/api/escrow/fund/route.ts +++ b/app/api/escrow/fund/route.ts @@ -4,6 +4,9 @@ * Record that the client has funded the escrow contract on-chain. * Verifies the funding transaction before updating state. * + * Idempotency: + * Requires an "Idempotency-Key" header (or `idempotencyKey` body field). + * * Body: * contractId string (UUID) * fundingTxHash string — on-chain transaction hash @@ -12,21 +15,15 @@ import { NextRequest, NextResponse } from 'next/server' import { withRbac, RbacContext } from '@/lib/auth/rbacMiddleware' -import { sql } from '@/lib/db' +import { withIdempotency } from '@/lib/idempotency' import { escrowService, EscrowError, escrowErrorToHttpStatus } from '@/lib/escrow' import { dispatchNotification } from '@/lib/notifications' -export const POST = withRbac('escrow:fund', async (request: NextRequest, auth: RbacContext) => { - let body: Record - try { - body = await request.json() - } catch { - return NextResponse.json( - { error: 'Request body must be valid JSON', code: 'INVALID_JSON' }, - { status: 400 } - ) - } - +export const POST = withRbac('escrow:fund', withIdempotency('escrow_fund', async ( + _request: NextRequest, + auth: RbacContext, + body: Record +) => { try { const result = await escrowService.fundEscrow({ contractId: body.contractId as string, @@ -70,4 +67,4 @@ export const POST = withRbac('escrow:fund', async (request: NextRequest, auth: R { status: 500 } ) } -}) +})) diff --git a/app/api/escrow/refund/route.ts b/app/api/escrow/refund/route.ts index aa3b47f..36c880e 100644 --- a/app/api/escrow/refund/route.ts +++ b/app/api/escrow/refund/route.ts @@ -4,6 +4,9 @@ * Refund all remaining escrowed funds back to the client. * Can be triggered by the client (voluntary cancellation) or an admin. * + * Idempotency: + * Requires an "Idempotency-Key" header (or `idempotencyKey` body field). + * * Body: * contractId string (UUID) * reason string @@ -11,20 +14,17 @@ import { NextRequest, NextResponse } from 'next/server' import { withAnyRbac, RbacContext } from '@/lib/auth/rbacMiddleware' +import { withIdempotency } from '@/lib/idempotency' import { escrowService, EscrowError, escrowErrorToHttpStatus } from '@/lib/escrow' import { dispatchNotification } from '@/lib/notifications' -export const POST = withAnyRbac(['escrow:refund', 'admin:contracts_freeze'], async (request: NextRequest, auth: RbacContext) => { - let body: Record - try { - body = await request.json() - } catch { - return NextResponse.json( - { error: 'Request body must be valid JSON', code: 'INVALID_JSON' }, - { status: 400 } - ) - } - +export const POST = withAnyRbac( + ['escrow:refund', 'admin:contracts_freeze'], + withIdempotency('escrow_refund', async ( + _request: NextRequest, + auth: RbacContext, + body: Record + ) => { try { const result = await escrowService.refundEscrow({ contractId: body.contractId as string, @@ -68,4 +68,5 @@ export const POST = withAnyRbac(['escrow:refund', 'admin:contracts_freeze'], asy { status: 500 } ) } -}) + }) +) diff --git a/app/api/escrow/release/route.ts b/app/api/escrow/release/route.ts index 580cafb..8f554f4 100644 --- a/app/api/escrow/release/route.ts +++ b/app/api/escrow/release/route.ts @@ -4,6 +4,9 @@ * Release funds for an approved milestone to the freelancer. * Only the contract client can trigger a release. * + * Idempotency: + * Requires an "Idempotency-Key" header (or `idempotencyKey` body field). + * * Body: * contractId string (UUID) * milestoneId string (UUID) @@ -11,20 +14,15 @@ import { NextRequest, NextResponse } from 'next/server' import { withRbac, RbacContext } from '@/lib/auth/rbacMiddleware' +import { withIdempotency } from '@/lib/idempotency' import { escrowService, EscrowError, escrowErrorToHttpStatus } from '@/lib/escrow' import { dispatchNotification } from '@/lib/notifications' -export const POST = withRbac('escrow:release', async (request: NextRequest, auth: RbacContext) => { - let body: Record - try { - body = await request.json() - } catch { - return NextResponse.json( - { error: 'Request body must be valid JSON', code: 'INVALID_JSON' }, - { status: 400 } - ) - } - +export const POST = withRbac('escrow:release', withIdempotency('escrow_release', async ( + _request: NextRequest, + auth: RbacContext, + body: Record +) => { try { const result = await escrowService.releaseFunds({ contractId: body.contractId as string, @@ -73,4 +71,4 @@ export const POST = withRbac('escrow:release', async (request: NextRequest, auth { status: 500 } ) } -}) +})) diff --git a/env.example b/env.example index b4fa902..42f6646 100644 --- a/env.example +++ b/env.example @@ -31,6 +31,12 @@ FILE_ENCRYPTION_KEY=replace_with_a_64_char_hex_string # Defaults to /uploads if not set. # FILE_UPLOAD_DIR=/data/taskchain/uploads +# ─── Idempotency ────────────────────────────────────────────────────────────── +# Retention window (in hours) for stored idempotency records on payment routes. +# Defaults to 24. Records older than this are purged by: +# pnpm idempotency:cleanup +# IDEMPOTENCY_TTL_HOURS=24 + # ─── Optional ───────────────────────────────────────────────────────────────── # Override the freelancer dashboard API base URL (defaults to /api/freelancer/dashboard). # NEXT_PUBLIC_DASHBOARD_API_URL=https://your-domain.com diff --git a/lib/db/migrations/010_idempotency_records.sql b/lib/db/migrations/010_idempotency_records.sql new file mode 100644 index 0000000..4070d53 --- /dev/null +++ b/lib/db/migrations/010_idempotency_records.sql @@ -0,0 +1,58 @@ +-- 010_idempotency_records.sql +-- +-- Idempotency layer for payment operations (issue #218). +-- +-- Guarantees that blockchain / payment requests are processed exactly once: +-- * every request carries a client-generated idempotency key +-- * the key + operation type is unique at the database level +-- * concurrent duplicates are rejected by the unique constraint instead of +-- racing through the escrow lifecycle +-- * stored responses are replayed for repeated requests +-- * records expire after a configurable TTL and are purged by a cleanup job + +-- ── Status enum ───────────────────────────────────────────────────────────── +DO $$ +BEGIN + IF NOT EXISTS (SELECT 1 FROM pg_type WHERE typname = 'idempotency_status') THEN + CREATE TYPE idempotency_status AS ENUM ('in_progress', 'completed', 'failed'); + END IF; +END; +$$; + +-- ── Core table ────────────────────────────────────────────────────────────── +CREATE TABLE IF NOT EXISTS idempotency_records ( + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + idempotency_key TEXT NOT NULL, + operation_type TEXT NOT NULL, + request_hash TEXT NOT NULL, + request_payload JSONB NOT NULL DEFAULT '{}'::jsonb, + response_payload JSONB, + response_status INTEGER, + status idempotency_status NOT NULL DEFAULT 'in_progress', + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + expires_at TIMESTAMPTZ NOT NULL, + + -- Race-safety: only one record may exist for a (key, operation) pair. + -- A concurrent duplicate INSERT fails here rather than executing the + -- payment operation a second time. + CONSTRAINT uq_idempotency_key_operation UNIQUE (idempotency_key, operation_type) +); + +-- ── Indexes ───────────────────────────────────────────────────────────────── +-- Supports the TTL cleanup sweep. +CREATE INDEX IF NOT EXISTS idx_idempotency_records_expires_at + ON idempotency_records (expires_at); + +-- Supports operational queries / dashboards for a given operation type. +CREATE INDEX IF NOT EXISTS idx_idempotency_records_operation_created + ON idempotency_records (operation_type, created_at DESC); + +-- ── Database-level duplicate protection for payment records ───────────────── +-- A given on-chain transaction hash must never be recorded twice for a money +-- movement (deposit / release / refund). Partial index so dispute rows and +-- rows without a hash are unaffected. +CREATE UNIQUE INDEX IF NOT EXISTS uq_escrow_transaction_logs_payment_hash + ON escrow_transaction_logs (transaction_hash) + WHERE transaction_hash IS NOT NULL + AND transaction_type IN ('deposit', 'milestone_release', 'refund'); diff --git a/lib/idempotency/config.ts b/lib/idempotency/config.ts new file mode 100644 index 0000000..28954f8 --- /dev/null +++ b/lib/idempotency/config.ts @@ -0,0 +1,44 @@ +/** + * Idempotency Layer — Configuration + * + * All tunables are read from the environment with safe defaults so the layer + * works out of the box in development and CI. + */ + +/** Standard header clients use to send their idempotency key. */ +export const IDEMPOTENCY_KEY_HEADER = 'Idempotency-Key' + +/** Response header indicating whether a request was served from storage. */ +export const IDEMPOTENCY_REPLAYED_HEADER = 'Idempotency-Replayed' + +/** Default retention window for completed idempotency records. */ +export const DEFAULT_IDEMPOTENCY_TTL_HOURS = 24 + +/** Hard bounds for the configurable TTL (1 hour … 30 days). */ +export const MIN_IDEMPOTENCY_TTL_HOURS = 1 +export const MAX_IDEMPOTENCY_TTL_HOURS = 720 + +/** Minimum accepted key length — enough entropy for a UUID or hash. */ +export const MIN_IDEMPOTENCY_KEY_LENGTH = 16 + +/** Maximum accepted key length (matches a generous TEXT sanity bound). */ +export const MAX_IDEMPOTENCY_KEY_LENGTH = 255 + +function clamp(value: number): number { + if (!Number.isFinite(value)) return DEFAULT_IDEMPOTENCY_TTL_HOURS + return Math.min( + MAX_IDEMPOTENCY_TTL_HOURS, + Math.max(MIN_IDEMPOTENCY_TTL_HOURS, Math.floor(value)) + ) +} + +/** + * Reads the idempotency TTL (in hours) from `IDEMPOTENCY_TTL_HOURS`. + * Falls back to {@link DEFAULT_IDEMPOTENCY_TTL_HOURS} when unset/invalid and + * clamps the result to the supported range. + */ +export function getIdempotencyTtlHours(): number { + const raw = process.env.IDEMPOTENCY_TTL_HOURS + if (!raw) return DEFAULT_IDEMPOTENCY_TTL_HOURS + return clamp(Number(raw)) +} diff --git a/lib/idempotency/errors.ts b/lib/idempotency/errors.ts new file mode 100644 index 0000000..aa79441 --- /dev/null +++ b/lib/idempotency/errors.ts @@ -0,0 +1,88 @@ +/** + * Idempotency Layer — Domain Errors + * + * Typed errors with machine-readable codes so route handlers can return + * precise HTTP responses (see the acceptance criteria in issue #218). + */ + +export class IdempotencyError extends Error { + /** Machine-readable error code returned to the client. */ + readonly code: string + /** HTTP status that should accompany the error. */ + readonly status: number + /** Optional structured details merged into the response body. */ + readonly details?: Record + + constructor( + message: string, + code: string, + status: number, + details?: Record + ) { + super(message) + this.name = 'IdempotencyError' + this.code = code + this.status = status + this.details = details + } +} + +/** No idempotency key was supplied with a payment request. */ +export class IdempotencyKeyRequiredError extends IdempotencyError { + constructor() { + super( + 'An "Idempotency-Key" header or "idempotencyKey" body field is required for this operation', + 'IDEMPOTENCY_KEY_REQUIRED', + 400 + ) + this.name = 'IdempotencyKeyRequiredError' + } +} + +/** The supplied key is present but malformed. */ +export class IdempotencyKeyInvalidError extends IdempotencyError { + constructor(reason: string) { + super( + `Invalid idempotency key: ${reason}`, + 'IDEMPOTENCY_KEY_INVALID', + 400 + ) + this.name = 'IdempotencyKeyInvalidError' + } +} + +/** The same key was reused with a different request payload. */ +export class IdempotencyKeyReusedError extends IdempotencyError { + constructor() { + super( + 'This idempotency key was already used with a different request payload', + 'IDEMPOTENCY_KEY_REUSED', + 409 + ) + this.name = 'IdempotencyKeyReusedError' + } +} + +/** + * A request with the same key is currently being processed. The client should + * retry after a short delay, at which point the stored result is returned. + */ +export class IdempotencyInProgressError extends IdempotencyError { + constructor(expiresAt?: string) { + super( + 'A request with this idempotency key is already in progress', + 'IDEMPOTENCY_IN_PROGRESS', + 409, + expiresAt ? { expiresAt } : undefined + ) + this.name = 'IdempotencyInProgressError' + } +} + +/** The idempotency store could not be read or written. */ +export class IdempotencyStorageError extends IdempotencyError { + constructor(message = 'Idempotency store is unavailable') { + super(message, 'IDEMPOTENCY_STORAGE_ERROR', 503) + this.name = 'IdempotencyStorageError' + } +} diff --git a/lib/idempotency/index.ts b/lib/idempotency/index.ts new file mode 100644 index 0000000..2ca477b --- /dev/null +++ b/lib/idempotency/index.ts @@ -0,0 +1,59 @@ +/** + * Idempotency Layer — Public API + * + * Usage (payment route): + * import { withIdempotency } from '@/lib/idempotency' + * + * export const POST = withRbac( + * 'escrow:fund', + * withIdempotency('escrow_fund', async (request, auth, body) => { + * // `auth` is the RBAC context, `body` is the parsed request body + * return NextResponse.json({ ... }) + * }) + * ) + */ + +export { IdempotencyService, idempotencyService } from './service' +export { withIdempotency } from './middleware' +export type { IdempotentRouteHandler } from './middleware' +export { + IdempotencyRepository, + idempotencyRepository, +} from './repository' + +export { + IdempotencyError, + IdempotencyKeyRequiredError, + IdempotencyKeyInvalidError, + IdempotencyKeyReusedError, + IdempotencyInProgressError, + IdempotencyStorageError, +} from './errors' + +export { + IDEMPOTENCY_KEY_HEADER, + IDEMPOTENCY_REPLAYED_HEADER, + DEFAULT_IDEMPOTENCY_TTL_HOURS, + MIN_IDEMPOTENCY_TTL_HOURS, + MAX_IDEMPOTENCY_TTL_HOURS, + getIdempotencyTtlHours, +} from './config' + +export { + extractIdempotencyKey, + validateIdempotencyKey, + hashRequestPayload, + stableStringify, +} from './validation' + +export { IDEMPOTENCY_OPERATIONS } from './types' +export type { + IdempotencyStatus, + IdempotencyOperationType, + IdempotencyRecord, + ClaimIdempotencyInput, + ClaimIdempotencyResult, + IdempotentResponse, + IdempotentOutcome, + IIdempotencyRepository, +} from './types' diff --git a/lib/idempotency/middleware.ts b/lib/idempotency/middleware.ts new file mode 100644 index 0000000..f884aa1 --- /dev/null +++ b/lib/idempotency/middleware.ts @@ -0,0 +1,113 @@ +/** + * Idempotency Layer — Route Middleware + * + * `withIdempotency` wraps a payment route handler and enforces the idempotency + * contract: + * + * 1. Parse + validate the JSON body. + * 2. Require a valid idempotency key (header or body field). + * 3. Claim the `(key, operation)` pair atomically. + * 4. Execute the wrapped handler only if the claim succeeded. + * 5. Persist a successful response and replay it for repeat requests. + * + * Compose it *inside* the auth / RBAC wrappers so unauthenticated requests are + * rejected before a key is ever claimed: + * + * export const POST = withRbac('escrow:fund', withIdempotency('escrow_fund', handler)) + */ + +import { NextRequest, NextResponse } from 'next/server' + +import { IDEMPOTENCY_REPLAYED_HEADER } from './config' +import { IdempotencyError } from './errors' +import { idempotencyService } from './service' +import { extractIdempotencyKey, validateIdempotencyKey } from './validation' +import type { IdempotencyOperationType, IdempotentResponse } from './types' + +export type IdempotentRouteHandler = ( + request: NextRequest, + auth: C, + body: Record +) => Promise + +function errorResponse(err: unknown): NextResponse { + if (err instanceof IdempotencyError) { + const headers: Record = {} + if (err.code === 'IDEMPOTENCY_IN_PROGRESS') { + headers['Retry-After'] = '1' + } + return NextResponse.json( + { error: err.message, code: err.code, ...(err.details ?? {}) }, + { status: err.status, headers } + ) + } + + console.error('[idempotency] Unexpected error:', err) + return NextResponse.json( + { error: 'Internal server error', code: 'INTERNAL_ERROR' }, + { status: 500 } + ) +} + +export function withIdempotency( + operationType: IdempotencyOperationType, + handler: IdempotentRouteHandler +): (request: NextRequest, auth: C) => Promise { + return async (request: NextRequest, auth: C): Promise => { + // --- Parse body (clone so the wrapped handler can still read it) --- + let body: Record + try { + const parsed = await request.clone().json() + if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) { + return NextResponse.json( + { error: 'Request body must be a JSON object', code: 'INVALID_JSON' }, + { status: 400 } + ) + } + body = parsed as Record + } catch { + return NextResponse.json( + { error: 'Request body must be valid JSON', code: 'INVALID_JSON' }, + { status: 400 } + ) + } + + // --- Validate key presence / format before any work happens --- + let key: string + try { + key = validateIdempotencyKey(extractIdempotencyKey(request, body)) + } catch (err) { + return errorResponse(err) + } + + try { + const outcome = await idempotencyService.run({ + key, + operationType, + requestPayload: body, + handler: async (): Promise => { + const response = await handler(request, auth, body) + let responseBody: Record = {} + try { + const parsed = await response.clone().json() + if (parsed && typeof parsed === 'object' && !Array.isArray(parsed)) { + responseBody = parsed as Record + } + } catch { + responseBody = {} + } + return { status: response.status, body: responseBody } + }, + }) + + return NextResponse.json(outcome.body, { + status: outcome.status, + headers: { + [IDEMPOTENCY_REPLAYED_HEADER]: String(outcome.replayed), + }, + }) + } catch (err) { + return errorResponse(err) + } + } +} diff --git a/lib/idempotency/repository.ts b/lib/idempotency/repository.ts new file mode 100644 index 0000000..778f081 --- /dev/null +++ b/lib/idempotency/repository.ts @@ -0,0 +1,181 @@ +/** + * Idempotency Layer — Database Repository + * + * All SQL for idempotency records lives here. The service layer depends on the + * {@link IIdempotencyRepository} interface, so tests can substitute an + * in-memory implementation without touching the database client. + * + * Concurrency model + * ----------------- + * `claim()` is a single atomic `INSERT … ON CONFLICT` statement. PostgreSQL's + * unique constraint on `(idempotency_key, operation_type)` serialises + * simultaneous requests: exactly one caller gets the row back and is allowed + * to execute the operation, every other caller observes the existing record. + * This is what prevents duplicate escrow/payment records at the database + * level. + */ + +import { sql } from '@/lib/db' +import type { + ClaimIdempotencyInput, + ClaimIdempotencyResult, + IdempotencyOperationType, + IdempotencyRecord, + IIdempotencyRepository, +} from './types' + +function parsePayload(value: unknown): Record { + if (!value) return {} + if (typeof value === 'string') { + try { + const parsed = JSON.parse(value) + return parsed && typeof parsed === 'object' + ? (parsed as Record) + : {} + } catch { + return {} + } + } + if (typeof value === 'object') return value as Record + return {} +} + +function rowToRecord(row: Record): IdempotencyRecord { + return { + id: row.id as string, + idempotencyKey: row.idempotency_key as string, + operationType: row.operation_type as IdempotencyOperationType, + requestHash: row.request_hash as string, + requestPayload: parsePayload(row.request_payload), + responsePayload: + row.response_payload === null || row.response_payload === undefined + ? null + : parsePayload(row.response_payload), + responseStatus: + row.response_status === null || row.response_status === undefined + ? null + : Number(row.response_status), + status: row.status as IdempotencyRecord['status'], + createdAt: new Date(row.created_at as string).toISOString(), + updatedAt: new Date(row.updated_at as string).toISOString(), + expiresAt: new Date(row.expires_at as string).toISOString(), + } +} + +export class IdempotencyRepository implements IIdempotencyRepository { + async claim(input: ClaimIdempotencyInput): Promise { + const rows = (await sql` + INSERT INTO idempotency_records ( + idempotency_key, + operation_type, + request_hash, + request_payload, + status, + expires_at, + created_at, + updated_at + ) + VALUES ( + ${input.key}, + ${input.operationType}, + ${input.requestHash}, + ${JSON.stringify(input.requestPayload)}::jsonb, + 'in_progress', + NOW() + make_interval(hours => ${input.ttlHours}::int), + NOW(), + NOW() + ) + ON CONFLICT (idempotency_key, operation_type) DO UPDATE + SET request_hash = EXCLUDED.request_hash, + request_payload = EXCLUDED.request_payload, + response_payload = NULL, + response_status = NULL, + status = 'in_progress', + updated_at = NOW(), + expires_at = EXCLUDED.expires_at + WHERE idempotency_records.expires_at <= NOW() + OR idempotency_records.status = 'failed' + RETURNING * + `) as Record[] + + if (rows.length > 0) { + return { claimed: true, record: rowToRecord(rows[0]) } + } + + // A live record already exists — return it so the caller can decide + // whether to replay the stored response or reject the request. + return { claimed: false, record: await this.get(input.key, input.operationType) } + } + + async get( + key: string, + operationType: IdempotencyOperationType + ): Promise { + const rows = (await sql` + SELECT * FROM idempotency_records + WHERE idempotency_key = ${key} + AND operation_type = ${operationType} + LIMIT 1 + `) as Record[] + + return rows[0] ? rowToRecord(rows[0]) : null + } + + async complete( + key: string, + operationType: IdempotencyOperationType, + responsePayload: Record, + responseStatus: number, + ttlHours: number + ): Promise { + const rows = (await sql` + UPDATE idempotency_records + SET response_payload = ${JSON.stringify(responsePayload)}::jsonb, + response_status = ${responseStatus}, + status = 'completed', + updated_at = NOW(), + expires_at = NOW() + make_interval(hours => ${ttlHours}::int) + WHERE idempotency_key = ${key} + AND operation_type = ${operationType} + RETURNING * + `) as Record[] + + return rows[0] ? rowToRecord(rows[0]) : null + } + + async markFailed( + key: string, + operationType: IdempotencyOperationType + ): Promise { + await sql` + UPDATE idempotency_records + SET status = 'failed', + updated_at = NOW() + WHERE idempotency_key = ${key} + AND operation_type = ${operationType} + ` + } + + async remove( + key: string, + operationType: IdempotencyOperationType + ): Promise { + await sql` + DELETE FROM idempotency_records + WHERE idempotency_key = ${key} + AND operation_type = ${operationType} + ` + } + + async purgeExpired(): Promise { + const rows = (await sql` + DELETE FROM idempotency_records + WHERE expires_at <= NOW() + RETURNING id + `) as Record[] + + return rows.length + } +} + +export const idempotencyRepository = new IdempotencyRepository() diff --git a/lib/idempotency/service.ts b/lib/idempotency/service.ts new file mode 100644 index 0000000..76f50fa --- /dev/null +++ b/lib/idempotency/service.ts @@ -0,0 +1,133 @@ +/** + * Idempotency Layer — Service + * + * Orchestrates the claim → execute → store lifecycle. It is persistence- + * agnostic (depends on {@link IIdempotencyRepository}) and free of Next.js + * imports, which keeps it straightforward to unit-test. + * + * Exactly-once semantics + * ---------------------- + * The operation is executed only by the caller that wins the atomic + * `repository.claim()`. Duplicate / concurrent callers never reach the handler: + * - completed record → the stored response is replayed + * - in-progress record → a 409 tells the client to retry shortly + * - failed record → the key is reclaimed so a retry can proceed + */ + +import { getIdempotencyTtlHours } from './config' +import { + IdempotencyInProgressError, + IdempotencyKeyReusedError, + IdempotencyStorageError, +} from './errors' +import { idempotencyRepository } from './repository' +import { hashRequestPayload } from './validation' +import type { + IdempotencyOperationType, + IdempotentOutcome, + IdempotentResponse, + IIdempotencyRepository, +} from './types' + +export interface RunIdempotentParams { + key: string + operationType: IdempotencyOperationType + requestPayload: Record + handler: () => Promise +} + +export class IdempotencyService { + constructor( + private readonly repo: IIdempotencyRepository = idempotencyRepository, + private readonly ttlHours: () => number = getIdempotencyTtlHours + ) {} + + /** + * Runs `handler` exactly once for the given `(key, operationType)` pair. + * + * @throws {IdempotencyKeyReusedError} same key, different payload + * @throws {IdempotencyInProgressError} duplicate request while in flight + * @throws {IdempotencyStorageError} the store could not be reached + */ + async run(params: RunIdempotentParams): Promise { + const ttlHours = this.ttlHours() + const requestHash = hashRequestPayload(params.requestPayload) + + let claim + try { + claim = await this.repo.claim({ + key: params.key, + operationType: params.operationType, + requestHash, + requestPayload: params.requestPayload, + ttlHours, + }) + } catch (err) { + throw new IdempotencyStorageError( + err instanceof Error ? `Failed to claim idempotency key: ${err.message}` : undefined + ) + } + + if (!claim.claimed) { + const existing = claim.record + if (!existing) { + throw new IdempotencyStorageError('Idempotency record disappeared during claim') + } + + // Reusing a key with a different payload is a client bug / replay attack. + if (existing.requestHash !== requestHash) { + throw new IdempotencyKeyReusedError() + } + + if (existing.status === 'completed') { + return { + status: existing.responseStatus ?? 200, + body: existing.responsePayload ?? {}, + replayed: true, + } + } + + if (existing.status === 'in_progress') { + throw new IdempotencyInProgressError(existing.expiresAt) + } + + throw new IdempotencyStorageError( + `Idempotency record is in an unexpected state: ${existing.status}` + ) + } + + try { + const response = await params.handler() + + if (response.status >= 200 && response.status < 300) { + try { + await this.repo.complete( + params.key, + params.operationType, + response.body, + response.status, + ttlHours + ) + } catch (err) { + throw new IdempotencyStorageError( + err instanceof Error + ? `Failed to store idempotent response: ${err.message}` + : undefined + ) + } + return { ...response, replayed: false } + } + + // Only successful responses are cached; a failed attempt may be retried + // with the same key, so release the claim. + await this.repo.markFailed(params.key, params.operationType).catch(() => {}) + return { ...response, replayed: false } + } catch (err) { + // Best-effort release so the client can retry the same key. + await this.repo.markFailed(params.key, params.operationType).catch(() => {}) + throw err + } + } +} + +export const idempotencyService = new IdempotencyService() diff --git a/lib/idempotency/types.ts b/lib/idempotency/types.ts new file mode 100644 index 0000000..daa6553 --- /dev/null +++ b/lib/idempotency/types.ts @@ -0,0 +1,92 @@ +/** + * Idempotency Layer — Shared Types + * + * Types used by the idempotency repository, service and route wrapper. + * Keeping them in a dedicated module lets consumers import types without + * pulling in the service (and therefore the DB client) as a side effect. + */ + +/** Lifecycle of a stored idempotency record. */ +export type IdempotencyStatus = 'in_progress' | 'completed' | 'failed' + +/** + * Logical operation the idempotency key is scoped to. + * The pair (idempotencyKey, operationType) is unique, so the same key can be + * reused safely across different operations. + */ +export const IDEMPOTENCY_OPERATIONS = [ + 'escrow_create', + 'escrow_fund', + 'escrow_release', + 'escrow_refund', + 'escrow_dispute', + 'dispute_resolve', +] as const + +export type IdempotencyOperationType = (typeof IDEMPOTENCY_OPERATIONS)[number] + +/** Persisted representation of an idempotency record. */ +export interface IdempotencyRecord { + id: string + idempotencyKey: string + operationType: IdempotencyOperationType + requestHash: string + requestPayload: Record + responsePayload: Record | null + responseStatus: number | null + status: IdempotencyStatus + createdAt: string + updatedAt: string + expiresAt: string +} + +/** Input for atomically claiming an idempotency key. */ +export interface ClaimIdempotencyInput { + key: string + operationType: IdempotencyOperationType + requestHash: string + requestPayload: Record + ttlHours: number +} + +/** Result of a claim attempt. */ +export interface ClaimIdempotencyResult { + /** True when this caller created / reclaimed the record and may execute. */ + claimed: boolean + /** The current record, whether freshly claimed or pre-existing. */ + record: IdempotencyRecord | null +} + +/** Repository abstraction — swapped for an in-memory stub in tests. */ +export interface IIdempotencyRepository { + claim(input: ClaimIdempotencyInput): Promise + get( + key: string, + operationType: IdempotencyOperationType + ): Promise + complete( + key: string, + operationType: IdempotencyOperationType, + responsePayload: Record, + responseStatus: number, + ttlHours: number + ): Promise + markFailed( + key: string, + operationType: IdempotencyOperationType + ): Promise + remove(key: string, operationType: IdempotencyOperationType): Promise + purgeExpired(): Promise +} + +/** HTTP response captured from a route handler. */ +export interface IdempotentResponse { + status: number + body: Record +} + +/** Outcome returned by the idempotency service. */ +export interface IdempotentOutcome extends IdempotentResponse { + /** True when the response came from a previously stored record. */ + replayed: boolean +} diff --git a/lib/idempotency/validation.ts b/lib/idempotency/validation.ts new file mode 100644 index 0000000..9a26e93 --- /dev/null +++ b/lib/idempotency/validation.ts @@ -0,0 +1,101 @@ +/** + * Idempotency Layer — Validation & Hashing Helpers + * + * These are pure functions (no DB / no Next.js side effects) so they can be + * unit-tested in isolation. + */ + +import { createHash } from 'node:crypto' +import type { NextRequest } from 'next/server' + +import { + IDEMPOTENCY_KEY_HEADER, + MIN_IDEMPOTENCY_KEY_LENGTH, + MAX_IDEMPOTENCY_KEY_LENGTH, +} from './config' +import { + IdempotencyKeyInvalidError, + IdempotencyKeyRequiredError, +} from './errors' + +/** Characters permitted in an idempotency key (UUIDs, hex hashes, base64url…). */ +const IDEMPOTENCY_KEY_PATTERN = /^[A-Za-z0-9._:+/-]+$/ + +/** + * Reads the idempotency key from the `Idempotency-Key` header, falling back to + * an `idempotencyKey` field on the request body (if the body is an object). + */ +export function extractIdempotencyKey( + request: NextRequest, + body?: unknown +): string | null { + const header = request.headers.get(IDEMPOTENCY_KEY_HEADER) + if (header && header.trim().length > 0) return header.trim() + + if (body && typeof body === 'object' && !Array.isArray(body)) { + const candidate = (body as Record).idempotencyKey + if (typeof candidate === 'string' && candidate.trim().length > 0) { + return candidate.trim() + } + } + + return null +} + +/** + * Validates a raw idempotency key and returns the normalised value. + * + * @throws {IdempotencyKeyRequiredError} when the key is null/undefined/empty + * @throws {IdempotencyKeyInvalidError} when the key is malformed + */ +export function validateIdempotencyKey(raw: string | null | undefined): string { + if (raw === null || raw === undefined || raw.trim().length === 0) { + throw new IdempotencyKeyRequiredError() + } + + const key = raw.trim() + + if (key.length < MIN_IDEMPOTENCY_KEY_LENGTH) { + throw new IdempotencyKeyInvalidError( + `must be at least ${MIN_IDEMPOTENCY_KEY_LENGTH} characters` + ) + } + + if (key.length > MAX_IDEMPOTENCY_KEY_LENGTH) { + throw new IdempotencyKeyInvalidError( + `must be at most ${MAX_IDEMPOTENCY_KEY_LENGTH} characters` + ) + } + + if (!IDEMPOTENCY_KEY_PATTERN.test(key)) { + throw new IdempotencyKeyInvalidError( + 'may only contain letters, digits and . _ : + / -' + ) + } + + return key +} + +/** Deterministic JSON serialisation (objects are key-sorted at every level). */ +export function stableStringify(value: unknown): string { + if (value === null || value === undefined) return 'null' + if (typeof value !== 'object') return JSON.stringify(value) as string + if (Array.isArray(value)) { + return `[${value.map((entry) => stableStringify(entry)).join(',')}]` + } + const record = value as Record + const keys = Object.keys(record).sort() + return `{${keys + .map((key) => `${JSON.stringify(key)}:${stableStringify(record[key])}`) + .join(',')}}` +} + +/** + * Hashes a request payload so a replayed key can be matched to the original + * request. Key order does not affect the hash. + */ +export function hashRequestPayload(payload: unknown): string { + return createHash('sha256') + .update(stableStringify(payload)) + .digest('hex') +} diff --git a/package.json b/package.json index be669c5..fab9a4a 100644 --- a/package.json +++ b/package.json @@ -14,6 +14,7 @@ "sync-worker": "tsx scripts/sync-worker.ts", "deadline-monitor": "tsx scripts/deadline-monitor.ts", "migrate": "tsx scripts/migrate.ts", + "idempotency:cleanup": "tsx scripts/idempotency-cleanup.ts", "build:production": "npm run migrate && next build", "start:production": "next start" }, diff --git a/scripts/idempotency-cleanup.ts b/scripts/idempotency-cleanup.ts new file mode 100644 index 0000000..7aecaff --- /dev/null +++ b/scripts/idempotency-cleanup.ts @@ -0,0 +1,47 @@ +/** + * scripts/idempotency-cleanup.ts + * + * Purges expired idempotency records. Run on a schedule (e.g. hourly cron, + * Railway scheduled job or GitHub Actions) so the `idempotency_records` table + * does not grow unbounded. + * + * Retention is controlled by `IDEMPOTENCY_TTL_HOURS` (default 24). The + * `expires_at` column is set when a record is created/completed, so this script + * simply deletes anything past its expiry. + * + * Usage: + * npx tsx scripts/idempotency-cleanup.ts + * # or + * pnpm idempotency:cleanup + */ + +import { neon } from '@neondatabase/serverless' +import * as dotenv from 'dotenv' + +dotenv.config() + +if (!process.env.DATABASE_URL) { + console.error( + 'FATAL: DATABASE_URL is not set. Cannot purge idempotency records.' + ) + process.exit(1) +} + +const sql = neon(process.env.DATABASE_URL) + +async function purgeExpired(): Promise { + const deleted = (await sql` + DELETE FROM idempotency_records + WHERE expires_at <= NOW() + RETURNING id + `) as Array<{ id: string }> + + console.log( + `\n✓ Idempotency cleanup complete — purged ${deleted.length} expired record(s).\n` + ) +} + +purgeExpired().catch((err) => { + console.error('\n✗ Idempotency cleanup failed:', err) + process.exit(1) +})