diff --git a/backend/openapi.json b/backend/openapi.json index e84421cb3..df8636f67 100644 --- a/backend/openapi.json +++ b/backend/openapi.json @@ -29,7 +29,7 @@ "type": "http", "scheme": "bearer", "bearerFormat": "JWT", - "description": "JWT issued by POST /auth/login or POST /auth/refresh." + "description": "JWT issued by POST /auth/login or POST /auth/refresh. `exp` is enforced with **zero** clock tolerance — a token is rejected the moment it expires. Only `nbf`/`iat` get a 5s tolerance for clock skew. The revocation list is checked on every authenticated request, so a token presented after `POST /auth/logout` (401 `TOKEN_REVOKED`) or after `POST /auth/logout-all` is refused immediately rather than at expiry." }, "apiKeyAuth": { "type": "apiKey", @@ -54,6 +54,31 @@ "type": "string", "format": "uuid" } + }, + "pageSize": { + "name": "limit", + "in": "query", + "required": false, + "description": "Maximum number of items to return. Hard ceiling is 50 (default 20). A larger value is **rejected**, not clamped, with `400` and `code: \"LIMIT_EXCEEDED\"` — the response body repeats the ceiling so clients can self-correct. Use `page` (clamped to 1..1000) to walk the rest.", + "schema": { + "type": "integer", + "minimum": 1, + "maximum": 50, + "default": 20, + "example": 20 + } + }, + "pageNumber": { + "name": "page", + "in": "query", + "required": false, + "description": "1-based page number for offset pagination. Values outside 1..1000 are clamped rather than rejected, and the effective page is echoed back in `pagination.currentPage`.", + "schema": { + "type": "integer", + "minimum": 1, + "maximum": 1000, + "default": 1 + } } }, "schemas": { @@ -106,6 +131,10 @@ "count": { "type": "integer" }, + "limit": { + "type": "integer", + "maximum": 50 + }, "total": { "type": "integer" }, @@ -118,7 +147,9 @@ "nullable": true }, "currentPage": { - "type": "integer" + "type": "integer", + "minimum": 1, + "maximum": 1000 }, "totalPages": { "type": "integer" @@ -131,6 +162,62 @@ } } }, + "Vault": { + "type": "object", + "required": [ + "id", + "aum", + "tvlUsd", + "createdAt", + "updatedAt" + ], + "properties": { + "id": { + "type": "string", + "format": "uuid" + }, + "aum": { + "type": "number", + "example": 0 + }, + "tvlUsd": { + "type": "string", + "nullable": true, + "example": "1250000.00" + }, + "createdAt": { + "type": "string", + "format": "date-time" + }, + "updatedAt": { + "type": "string", + "format": "date-time" + } + } + }, + "VaultListResponse": { + "type": "object", + "required": [ + "data", + "pagination", + "timestamp" + ], + "properties": { + "data": { + "type": "array", + "items": { + "$ref": "#/components/schemas/Vault" + } + }, + "pagination": { + "$ref": "#/components/schemas/PaginationMeta" + }, + "timestamp": { + "type": "string", + "format": "date-time" + } + } + }, "VaultSummary": { "type": "object", "properties": { @@ -143,15 +230,8 @@ "example": 0 }, "apy": { - "type": ["number", "null"], - "example": 8.45, - "description": "Annualised APY as a decimal percentage. null when the vault has insufficient price history (e.g. zero shares or fewer than 2 snapshots)." - }, - "apyStatus": { - "type": "string", - "enum": ["ok", "insufficient_data"], - "example": "ok", - "description": "ok when apy is a valid number; insufficient_data when apy is null." + "type": "number", + "example": 0 }, "timestamp": { "type": "string", @@ -363,11 +443,21 @@ "timestamp": "2024-01-01T00:00:00.000Z", "uptime": 123.4, "environment": "production", + "lastIndexedLedger": 12345678, "checks": { "api": "up", "cache": "up", "stellarRpc": "up", + "databasePrimary": "up", + "databaseReplica": "up", + "prisma": "up", + "jobs": "up", "indexer": "up" + }, + "sorobanCircuitBreaker": { + "state": "closed", + "failures": 0, + "retryAfterMs": 0 } } } @@ -522,16 +612,26 @@ "application/json": { "schema": { "type": "object", - "required": ["apy", "apyStatus", "timestamp"], + "required": [ + "apy", + "apyStatus", + "timestamp" + ], "properties": { "apy": { - "type": ["number", "null"], + "type": [ + "number", + "null" + ], "example": 8.45, "description": "Annualised APY as a decimal percentage. null when insufficient data." }, "apyStatus": { "type": "string", - "enum": ["ok", "insufficient_data"], + "enum": [ + "ok", + "insufficient_data" + ], "example": "ok" }, "timestamp": { @@ -546,6 +646,95 @@ } } }, + "/api/v1/vaults": { + "get": { + "tags": [ + "Vault" + ], + "summary": "List vaults", + "description": "Paginated listing of active (non-deleted) vaults. `limit` is capped at 50 and `page` at 1000; exceeding `limit` fails fast with `400` / `LIMIT_EXCEEDED` instead of silently clamping, so a client that asks for an oversized page is never handed a response that looks complete. Rate limited: global 100 req / 15 min per IP; authenticated APIs 30 req / min per key. `429` responses follow the standard error envelope with a `retryAfter` detail.", + "parameters": [ + { + "$ref": "#/components/parameters/pageSize" + }, + { + "$ref": "#/components/parameters/pageNumber" + }, + { + "$ref": "#/components/parameters/correlationId" + } + ], + "responses": { + "200": { + "description": "One page of vaults", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/VaultListResponse" + }, + "example": { + "data": [ + { + "id": "0f1b3a2c-6d4e-4a1b-9f0e-2c5d7e8b9a01", + "aum": 1250000, + "tvlUsd": "1250000.00", + "createdAt": "2026-01-01T00:00:00.000Z", + "updatedAt": "2026-01-02T00:00:00.000Z" + } + ], + "pagination": { + "count": 1, + "limit": 20, + "total": 1, + "nextCursor": null, + "prevCursor": null, + "currentPage": 1, + "totalPages": 1, + "hasNextPage": false, + "hasPrevPage": false + }, + "timestamp": "2026-01-01T00:00:00.000Z" + } + } + } + }, + "400": { + "description": "`limit` above the published ceiling", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ErrorEnvelope" + }, + "example": { + "error": "Bad Request", + "status": 400, + "code": "LIMIT_EXCEEDED", + "message": "limit must not exceed 50 (received 100000).", + "retryable": false, + "details": { + "field": "limit", + "requested": 100000, + "maxLimit": 50, + "defaultLimit": 20, + "maxPage": 1000 + } + } + } + } + }, + "429": { + "description": "Rate limited", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ErrorEnvelope" + } + } + } + } + } + } + }, "/api/v1/vault/deposits": { "post": { "tags": [ diff --git a/backend/prisma/dev.db b/backend/prisma/dev.db index 4816f3079..1f9b8f6f1 100644 Binary files a/backend/prisma/dev.db and b/backend/prisma/dev.db differ diff --git a/backend/prisma/migrations/20260929120000_add_session_audit_and_tenant_scopes/migration.sql b/backend/prisma/migrations/20260929120000_add_session_audit_and_tenant_scopes/migration.sql new file mode 100644 index 000000000..e92395c13 --- /dev/null +++ b/backend/prisma/migrations/20260929120000_add_session_audit_and_tenant_scopes/migration.sql @@ -0,0 +1,38 @@ +-- AlterTable +ALTER TABLE "Transaction" ADD COLUMN "deletedAt" DATETIME; + +-- AlterTable +ALTER TABLE "WebhookEndpoint" ADD COLUMN "tenantId" TEXT; + +-- CreateTable +CREATE TABLE "SessionAuditLog" ( + "id" TEXT NOT NULL PRIMARY KEY, + "walletAddress" TEXT NOT NULL, + "eventType" TEXT NOT NULL, + "reason" TEXT NOT NULL, + "sessionId" TEXT NOT NULL, + "ipAddress" TEXT, + "userAgent" TEXT, + "metadata" TEXT, + "correlationId" TEXT, + "traceId" TEXT, + "timestamp" DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +-- CreateIndex +CREATE INDEX "SessionAuditLog_walletAddress_timestamp_idx" ON "SessionAuditLog"("walletAddress", "timestamp" DESC); + +-- CreateIndex +CREATE INDEX "SessionAuditLog_eventType_idx" ON "SessionAuditLog"("eventType"); + +-- CreateIndex +CREATE INDEX "SessionAuditLog_timestamp_idx" ON "SessionAuditLog"("timestamp" DESC); + +-- CreateIndex +CREATE INDEX "Transaction_tenantId_idx" ON "Transaction"("tenantId"); + +-- CreateIndex +CREATE INDEX "Transaction_deletedAt_idx" ON "Transaction"("deletedAt"); + +-- CreateIndex +CREATE INDEX "WebhookEndpoint_tenantId_idx" ON "WebhookEndpoint"("tenantId"); diff --git a/backend/prisma/schema.prisma b/backend/prisma/schema.prisma index 20e8c2e41..5b4e16431 100644 --- a/backend/prisma/schema.prisma +++ b/backend/prisma/schema.prisma @@ -50,6 +50,9 @@ model Transaction { tenantId String? latencyMs Int? failureReason String? + // Soft-delete marker used by the tenant boundary guard so revoked rows are + // never returned to a tenant-scoped read. + deletedAt DateTime? timestamp DateTime @default(now()) // Issue #895 – query optimisation indices @@ -61,6 +64,8 @@ model Transaction { @@index([user, type, timestamp(sort: Desc)]) @@index([user, status, timestamp(sort: Desc)]) @@index([vaultId, tenantId, timestamp(sort: Desc)]) + @@index([tenantId]) + @@index([deletedAt]) } model Referral { @@ -221,6 +226,8 @@ model WebhookEndpoint { challengeExpiresAt DateTime? verifiedAt DateTime? lastVerificationError String? + // Owning tenant, used by the tenant boundary guard for cross-account reads. + tenantId String? createdAt DateTime @default(now()) updatedAt DateTime @updatedAt deletedAt DateTime? @@ -231,6 +238,7 @@ model WebhookEndpoint { @@index([createdAt]) @@index([deletedAt]) @@index([verificationStatus]) + @@index([tenantId]) } model WebhookDelivery { @@ -675,6 +683,8 @@ model WalletTenantAssociation { @@index([walletAddress, tenantId]) @@index([tenantId]) +} + model IdempotencyKey { id String @id @default(uuid()) keyId String @@ -693,3 +703,27 @@ model IdempotencyKey { @@unique([keyId, tenantId]) @@index([expiresAt]) } + +// ─── Session Audit ──────────────────────────────────────────────────────────── + +/// Append-only trail of session lifecycle events (login, refresh, logout, +/// expiry) written by `sessionAudit.ts`. Powers incident investigation and +/// the anomaly heuristics in `detectSuspiciousActivity`. +model SessionAuditLog { + id String @id @default(uuid()) + walletAddress String + eventType String // created | refreshed | revoked | expired | failed + reason String // user_logout | token_rotation | token_expiration | ... + sessionId String + ipAddress String? + userAgent String? + /// JSON-serialised context attached to the event. + metadata String? + correlationId String? + traceId String? + timestamp DateTime @default(now()) + + @@index([walletAddress, timestamp(sort: Desc)]) + @@index([eventType]) + @@index([timestamp(sort: Desc)]) +} diff --git a/backend/schema-snapshots/get-_api_v1_vaults.json b/backend/schema-snapshots/get-_api_v1_vaults.json new file mode 100644 index 000000000..604802c13 --- /dev/null +++ b/backend/schema-snapshots/get-_api_v1_vaults.json @@ -0,0 +1,89 @@ +{ + "type": "object", + "properties": { + "data": { + "type": "array", + "items": { + "type": "object", + "properties": { + "id": { + "type": "string" + }, + "aum": { + "type": "number" + }, + "tvlUsd": { + "type": "string" + }, + "createdAt": { + "type": "string" + }, + "updatedAt": { + "type": "string" + } + }, + "required": [ + "id", + "aum", + "tvlUsd", + "createdAt", + "updatedAt" + ], + "additionalProperties": false + } + }, + "pagination": { + "type": "object", + "properties": { + "count": { + "type": "number" + }, + "limit": { + "type": "number" + }, + "total": { + "type": "number" + }, + "nextCursor": { + "type": "string" + }, + "prevCursor": { + "type": "string" + }, + "currentPage": { + "type": "number" + }, + "totalPages": { + "type": "number" + }, + "hasNextPage": { + "type": "boolean" + }, + "hasPrevPage": { + "type": "boolean" + } + }, + "required": [ + "count", + "limit", + "total", + "nextCursor", + "prevCursor", + "currentPage", + "totalPages", + "hasNextPage", + "hasPrevPage" + ], + "additionalProperties": false + }, + "timestamp": { + "type": "string" + } + }, + "required": [ + "data", + "pagination", + "timestamp" + ], + "additionalProperties": false +} diff --git a/backend/schema-snapshots/get-_health.json b/backend/schema-snapshots/get-_health.json index c62ee2340..05101c579 100644 --- a/backend/schema-snapshots/get-_health.json +++ b/backend/schema-snapshots/get-_health.json @@ -135,4 +135,4 @@ "sorobanCircuitBreaker" ], "additionalProperties": false -} \ No newline at end of file +} diff --git a/backend/schema-snapshots/get-_ready.json b/backend/schema-snapshots/get-_ready.json index ce8460efd..cca011448 100644 --- a/backend/schema-snapshots/get-_ready.json +++ b/backend/schema-snapshots/get-_ready.json @@ -42,4 +42,4 @@ "dependencies" ], "additionalProperties": false -} \ No newline at end of file +} diff --git a/backend/src/__tests__/idempotencyStore.test.ts b/backend/src/__tests__/idempotencyStore.test.ts new file mode 100644 index 000000000..e40c6a951 --- /dev/null +++ b/backend/src/__tests__/idempotencyStore.test.ts @@ -0,0 +1,417 @@ +/** + * Unit tests for the idempotency replay cache (`idempotencyStore.ts`). + * + * `transferOrchestrator.test.ts` already exercises the happy path end to end. + * What it cannot reach is the Redis branch: with no `REDIS_URL` the + * `redisClientManager` never reports ready, so `IdempotencyStore.redis` + * returns `null` and every Redis helper is skipped. These tests drive the + * store through a stub client so the Redis code path — and, more importantly, + * its *fallback* when Redis errors — is actually executed. + */ + +import { + IdempotencyStore, + IdempotencyConflictError, + buildIdempotencyFingerprint, + getIdempotencyHashThreshold, + idempotencyStore, +} from '../idempotencyStore'; + +type RedisStub = Record; + +function createRedisStub(overrides: Partial = {}): RedisStub { + const store = new Map(); + + const base: RedisStub = { + get: jest.fn(async (key: string) => store.get(key) ?? null), + set: jest.fn(async (key: string, value: string) => { + store.set(key, value); + return 'OK'; + }), + del: jest.fn(async (key: string) => (store.delete(key) ? 1 : 0)), + exists: jest.fn(async (key: string) => (store.has(key) ? 1 : 0)), + scan: jest.fn(async () => ['0', [] as string[]]), + ttl: jest.fn(async () => -1), + }; + + return { ...base, ...overrides }; +} + +/** Points the module's `redisClientManager` at `client` for the test's duration. */ +async function withRedisClient(client: unknown, fn: () => Promise): Promise { + // eslint-disable-next-line @typescript-eslint/no-var-requires + const rateLimiter = require('../rateLimiter') as typeof import('../rateLimiter'); + const manager = rateLimiter.redisClientManager as unknown as { + getClient: () => unknown; + isReady: () => boolean; + }; + const originalGetClient = manager.getClient; + const originalIsReady = manager.isReady; + + manager.getClient = () => client; + manager.isReady = () => true; + + try { + return await fn(); + } finally { + manager.getClient = originalGetClient; + manager.isReady = originalIsReady; + } +} + +describe('IdempotencyStore', () => { + describe('execute() with the in-process cache', () => { + it('runs the operation once and replays the stored result afterwards', async () => { + const store = new IdempotencyStore(60_000); + const operation = jest.fn(async () => ({ statusCode: 201, body: { hash: 'abc' } })); + + const first = await store.execute('key-1', 'fp-1', operation); + const second = await store.execute('key-1', 'fp-1', operation); + + expect(first).toEqual({ result: { statusCode: 201, body: { hash: 'abc' } }, replayed: false }); + expect(second.replayed).toBe(true); + expect(second.result).toEqual(first.result); + expect(operation).toHaveBeenCalledTimes(1); + }); + + it('throws IdempotencyConflictError when the body changed for the same key', async () => { + const store = new IdempotencyStore(60_000); + await store.execute('key-2', 'fp-1', async () => ({ statusCode: 200, body: 'a' })); + + await expect( + store.execute('key-2', 'fp-2', async () => ({ statusCode: 200, body: 'b' })), + ).rejects.toBeInstanceOf(IdempotencyConflictError); + }); + + it('coalesces concurrent calls into a single operation', async () => { + const store = new IdempotencyStore(60_000); + let resolveOperation: (() => void) | undefined; + const gate = new Promise((resolve) => { + resolveOperation = resolve; + }); + const operation = jest.fn(async () => { + await gate; + return { statusCode: 200, body: 'once' }; + }); + + const inflight = Promise.all([ + store.execute('key-3', 'fp', operation), + store.execute('key-3', 'fp', operation), + ]); + resolveOperation!(); + const [a, b] = await inflight; + + expect(operation).toHaveBeenCalledTimes(1); + expect([a.replayed, b.replayed].filter(Boolean)).toHaveLength(1); + }); + + it('rejects a concurrent call with a different fingerprint', async () => { + const store = new IdempotencyStore(60_000); + let resolveOperation: (() => void) | undefined; + const gate = new Promise((resolve) => { + resolveOperation = resolve; + }); + + const inflight = store.execute('key-4', 'fp-a', async () => { + await gate; + return { statusCode: 200, body: 'a' }; + }); + const conflicting = store.execute('key-4', 'fp-b', async () => ({ statusCode: 200, body: 'b' })); + + resolveOperation!(); + await inflight; + await expect(conflicting).rejects.toBeInstanceOf(IdempotencyConflictError); + }); + + it('does not leave a key pending after the operation rejects', async () => { + const store = new IdempotencyStore(60_000); + await expect( + store.execute('key-5', 'fp', async () => { + throw new Error('boom'); + }), + ).rejects.toThrow('boom'); + + // The pending slot must be free, so a retry runs the operation again. + const retry = await store.execute('key-5', 'fp', async () => ({ statusCode: 200, body: 'ok' })); + expect(retry.replayed).toBe(false); + }); + }); + + describe('observability and maintenance', () => { + it('reports hits, conflicts and active keys', async () => { + const store = new IdempotencyStore(60_000); + await store.execute('m-1', 'fp', async () => ({ statusCode: 200, body: 1 })); + await store.execute('m-1', 'fp', async () => ({ statusCode: 200, body: 1 })); + await expect( + store.execute('m-2', 'other', async () => ({ statusCode: 200, body: 2 })), + ).resolves.toBeDefined(); + + const metrics = store.getMetrics(); + expect(metrics.hits).toBe(1); + expect(metrics.activeKeys).toBe(2); + expect(metrics.pendingKeys).toBe(0); + }); + + it('counts a conflict', async () => { + const store = new IdempotencyStore(60_000); + await store.execute('m-3', 'fp', async () => ({ statusCode: 200, body: 1 })); + await expect( + store.execute('m-3', 'fp-different', async () => ({ statusCode: 200, body: 2 })), + ).rejects.toBeInstanceOf(IdempotencyConflictError); + + expect(store.getMetrics().conflicts).toBe(1); + }); + + it('lists stored keys, optionally filtered by prefix', async () => { + const store = new IdempotencyStore(60_000); + await store.execute('alpha-1', 'fp', async () => ({ statusCode: 200, body: 1 })); + await store.execute('beta-1', 'fp', async () => ({ statusCode: 200, body: 2 })); + + expect(store.inspectKeys().map((entry) => entry.key).sort()).toEqual(['alpha-1', 'beta-1']); + expect(store.inspectKeys('alpha').map((entry) => entry.key)).toEqual(['alpha-1']); + expect(store.inspectKeys('alpha')[0].metadata.status).toBe('completed'); + }); + + it('deletes a single key and reports whether anything was removed', async () => { + const store = new IdempotencyStore(60_000); + await store.execute('d-1', 'fp', async () => ({ statusCode: 200, body: 1 })); + + expect(await store.deleteKey('d-1')).toBe(true); + expect(await store.deleteKey('d-1')).toBe(false); + expect(store.inspectKeys()).toHaveLength(0); + }); + + it('clear() empties the store and counts the evictions', async () => { + const store = new IdempotencyStore(60_000); + await store.execute('c-1', 'fp', async () => ({ statusCode: 200, body: 1 })); + await store.execute('c-2', 'fp', async () => ({ statusCode: 200, body: 2 })); + + const before = store.getMetrics().evictions; + store.clear(); + + expect(store.inspectKeys()).toHaveLength(0); + expect(store.getMetrics().evictions).toBe(before + 2); + }); + + it('prunes entries older than the retention window and honours dryRun', async () => { + const store = new IdempotencyStore(60_000); + await store.execute('p-1', 'fp', async () => ({ statusCode: 200, body: 1 })); + await store.execute('p-2', 'fp', async () => ({ statusCode: 200, body: 2 })); + + // A retention of -1 puts the cutoff in the future, so everything + // currently stored counts as stale. + const dry = await store.pruneStaleKeys(-1, true); + expect(dry.localPruned).toBe(2); + expect(dry.redisPruned).toBe(0); + expect(store.inspectKeys()).toHaveLength(2); + + const applied = await store.pruneStaleKeys(-1, false); + expect(applied.localPruned).toBe(2); + expect(applied.pruned).toBe(2); + expect(store.inspectKeys()).toHaveLength(0); + }); + + it('keeps entries that are newer than the retention window', async () => { + const store = new IdempotencyStore(60_000); + await store.execute('p-3', 'fp', async () => ({ statusCode: 200, body: 1 })); + + const result = await store.pruneStaleKeys(60_000, false); + expect(result.localPruned).toBe(0); + expect(store.inspectKeys()).toHaveLength(1); + }); + }); + + describe('Redis path', () => { + it('reads and writes through Redis when it is available', async () => { + const redis = createRedisStub(); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + const first = await store.execute('r-1', 'fp', async () => ({ statusCode: 200, body: 'v' })); + expect(first.replayed).toBe(false); + expect(redis.set).toHaveBeenCalled(); + + // A fresh store shares no local cache, so only Redis can serve the replay. + const other = new IdempotencyStore(60_000); + const second = await other.execute('r-1', 'fp', async () => ({ statusCode: 200, body: 'v' })); + expect(second.replayed).toBe(true); + expect(redis.get).toHaveBeenCalled(); + }); + }); + + it('detects a conflicting body via Redis', async () => { + const redis = createRedisStub(); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + await store.execute('r-2', 'fp', async () => ({ statusCode: 200, body: 'v' })); + + await expect( + store.execute('r-2', 'different', async () => ({ statusCode: 200, body: 'v' })), + ).rejects.toBeInstanceOf(IdempotencyConflictError); + }); + }); + + it('deletes via Redis and reports a removal', async () => { + const redis = createRedisStub(); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + await store.execute('r-3', 'fp', async () => ({ statusCode: 200, body: 'v' })); + + expect(await store.deleteKey('r-3')).toBe(true); + expect(redis.del).toHaveBeenCalled(); + expect(await store.deleteKey('never-stored')).toBe(false); + }); + }); + + it('survives Redis write failures by falling back to the local cache', async () => { + const redis = createRedisStub({ set: jest.fn(async () => { throw new Error('write fail'); }) }); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + const first = await store.execute('r-4', 'fp', async () => ({ statusCode: 200, body: 'v' })); + const second = await store.execute('r-4', 'fp', async () => ({ statusCode: 200, body: 'v' })); + + expect(first.replayed).toBe(false); + expect(second.replayed).toBe(true); + }); + }); + + it('survives Redis read failures', async () => { + const redis = createRedisStub({ get: jest.fn(async () => { throw new Error('read fail'); }) }); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + const first = await store.execute('r-5', 'fp', async () => ({ statusCode: 200, body: 'v' })); + const second = await store.execute('r-5', 'fp', async () => ({ statusCode: 200, body: 'v' })); + + expect(first.replayed).toBe(false); + expect(second.replayed).toBe(true); + }); + }); + + it('survives Redis delete failures', async () => { + const redis = createRedisStub({ del: jest.fn(async () => { throw new Error('del fail'); }) }); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + await store.execute('r-6', 'fp', async () => ({ statusCode: 200, body: 'v' })); + + // The local entry is still removed, so the key is gone either way. + expect(await store.deleteKey('r-6')).toBe(true); + }); + }); + + it('prunes stale Redis keys, honouring dryRun', async () => { + const keys = new Map(); + const redis = createRedisStub({ + scan: jest.fn(async () => ['0', ['idempotency:r-7']]), + get: jest.fn(async (key: string) => { + if (key !== 'idempotency:r-7') return null; + return ( + keys.get(key) ?? + JSON.stringify({ + statusCode: 200, + body: 'v', + fingerprint: 'fp', + metadata: { + createdAt: new Date(Date.now() - 120_000).toISOString(), + lastAccessedAt: new Date().toISOString(), + replayCount: 0, + status: 'completed', + }, + }) + ); + }), + ttl: jest.fn(async () => -1), + }); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + + const dry = await store.pruneStaleKeys(60_000, true); + expect(dry.redisPruned).toBe(1); + expect(redis.del).not.toHaveBeenCalled(); + + const applied = await store.pruneStaleKeys(60_000, false); + expect(applied.redisPruned).toBe(1); + expect(redis.del).toHaveBeenCalledWith('idempotency:r-7'); + }); + }); + + it('prunes a Redis key whose TTL has run out', async () => { + const redis = createRedisStub({ + scan: jest.fn(async () => ['0', ['idempotency:r-8']]), + get: jest.fn(async () => + JSON.stringify({ + statusCode: 200, + body: 'v', + fingerprint: 'fp', + metadata: { + createdAt: new Date().toISOString(), + lastAccessedAt: new Date().toISOString(), + replayCount: 0, + status: 'completed', + }, + }), + ), + ttl: jest.fn(async () => 0), + }); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + const result = await store.pruneStaleKeys(60_000, false); + expect(result.redisPruned).toBe(1); + }); + }); + + it('prunes a Redis key that cannot be parsed', async () => { + const redis = createRedisStub({ + scan: jest.fn(async () => ['0', ['idempotency:r-9']]), + get: jest.fn(async () => 'not json'), + }); + + await withRedisClient(redis, async () => { + const store = new IdempotencyStore(60_000); + const result = await store.pruneStaleKeys(60_000, false); + expect(result.redisPruned).toBe(1); + expect(redis.del).toHaveBeenCalledWith('idempotency:r-9'); + }); + }); + }); + + describe('buildIdempotencyFingerprint()', () => { + it('is stable regardless of key order', () => { + expect(buildIdempotencyFingerprint({ a: 1, b: 2 })).toBe( + buildIdempotencyFingerprint({ b: 2, a: 1 }), + ); + }); + + it('handles nested objects, arrays, dates and null', () => { + const date = new Date('2026-01-01T00:00:00.000Z'); + const fingerprint = buildIdempotencyFingerprint({ z: [1, { y: 2 }], d: date, n: null }); + + expect(fingerprint).toContain('2026-01-01T00:00:00.000Z'); + expect(fingerprint).toContain('null'); + expect(buildIdempotencyFingerprint({ z: [1, { y: 2 }], d: date, n: null })).toBe(fingerprint); + }); + + it('hashes payloads larger than the configured threshold', () => { + const oversized = { blob: 'x'.repeat(getIdempotencyHashThreshold() + 10) }; + const fingerprint = buildIdempotencyFingerprint(oversized); + + expect(fingerprint.startsWith('hashv1:')).toBe(true); + expect(fingerprint).toHaveLength('hashv1:'.length + 64); + }); + + it('does not hash payloads within the threshold', () => { + expect(buildIdempotencyFingerprint({ blob: 'small' })).not.toMatch(/^hashv1:/); + }); + }); + + it('exposes a module-level singleton', () => { + expect(idempotencyStore).toBeInstanceOf(IdempotencyStore); + expect(typeof idempotencyStore.getMetrics().activeKeys).toBe('number'); + }); +}); diff --git a/backend/src/__tests__/issues711.test.ts b/backend/src/__tests__/issues711.test.ts index 7feef6c6b..f7c019fc1 100644 --- a/backend/src/__tests__/issues711.test.ts +++ b/backend/src/__tests__/issues711.test.ts @@ -47,11 +47,14 @@ describe('#711 API contract schema snapshots', () => { }); it('detects newly added required fields as breaking changes', () => { - const baseline = zodToJsonShape(HealthResponseSchema); - const current = JSON.parse(JSON.stringify(baseline)) as typeof baseline; - // Simulate an older snapshot that is missing the 'indexer' field - delete current.properties?.checks?.properties?.indexer; - current.properties!.checks!.required = (current.properties!.checks!.required ?? []).filter( + const current = zodToJsonShape(HealthResponseSchema); + const baseline = JSON.parse(JSON.stringify(current)) as typeof current; + // Simulate an older snapshot that is missing the 'indexer' field, so the + // live schema now has a required field the committed snapshot does not know + // about. `baseline` is the committed snapshot and `current` is the live + // schema, so the deletion has to happen on the baseline. + delete baseline.properties?.checks?.properties?.indexer; + baseline.properties!.checks!.required = (baseline.properties!.checks!.required ?? []).filter( (k: string) => k !== 'indexer', ); diff --git a/backend/src/__tests__/jwtExpiryAndRevocation.test.ts b/backend/src/__tests__/jwtExpiryAndRevocation.test.ts new file mode 100644 index 000000000..80efac8b2 --- /dev/null +++ b/backend/src/__tests__/jwtExpiryAndRevocation.test.ts @@ -0,0 +1,508 @@ +/** + * Regression tests for Issue #1431 — the auth middleware accepted an expired + * JWT because the verification window was configured with a 60s + * `clockTolerance` shared across every time claim, so a stolen bearer token + * stayed valid for a full minute after logout/expiry and could still replay + * `POST /vault/:id/withdraw`. + * + * The fix splits the tolerance: `nbf`/`iat` get + * `CLOCK_SKEW_TOLERANCE_SECONDS` (5s) of slack for clock skew, `exp` gets + * none, and the revocation list is consulted on every authenticated request. + */ + +import crypto from 'crypto'; +import request from 'supertest'; +import app from '../index'; +import { + CLOCK_SKEW_TOLERANCE_SECONDS, + TOKEN_REVOKED_CODE, + TokenExpiredError, + TokenNotYetValidError, + assertTimeClaims, + isAccessTokenRevoked, + issueTokenPair, + revokeAccessToken, + revokeAllAccessTokens, + verifyJwt, + type JwtPayload, +} from '../auth'; +import { + InMemoryRevocationStore, + RedisRevocationStore, + getRevocationStore, + setRevocationStore, + type RevocationStore, +} from '../tokenRevocation'; +import { VALID_TEST_WALLET, SECOND_TEST_WALLET } from './setup'; + +const TEST_WALLET = VALID_TEST_WALLET; + +// ─── Helpers ───────────────────────────────────────────────────────────────── + +function base64UrlEncode(input: string): string { + return Buffer.from(input, 'utf8') + .toString('base64') + .replace(/\+/g, '-') + .replace(/\//g, '_') + .replace(/=/g, ''); +} + +function signingSecret(): string { + return process.env.JWT_SECRET || 'change-me-in-production-must-be-at-least-32-characters'; +} + +/** Signs an arbitrary payload with the same HS256 scheme auth.ts uses. */ +function signToken(payload: Record): string { + const header = base64UrlEncode(JSON.stringify({ alg: 'HS256', typ: 'JWT' })); + const body = base64UrlEncode(JSON.stringify(payload)); + const sig = crypto + .createHmac('sha256', signingSecret()) + .update(`${header}.${body}`) + .digest('base64') + .replace(/\+/g, '-') + .replace(/\//g, '_') + .replace(/=/g, ''); + return `${header}.${body}.${sig}`; +} + +function payloadAt( + expOffsetSeconds: number, + extra: Record = {}, + nowSeconds = Math.floor(Date.now() / 1000), +): JwtPayload { + return { + sub: TEST_WALLET, + iat: nowSeconds, + exp: nowSeconds + expOffsetSeconds, + jti: crypto.randomUUID(), + ...extra, + } as JwtPayload; +} + +describe('Issue #1431 — strict exp, 5s nbf skew, revocation on every request', () => { + const originalStore = getRevocationStore(); + + beforeEach(async () => { + await originalStore.clear(); + }); + + afterAll(async () => { + await originalStore.clear(); + setRevocationStore(originalStore); + }); + + // ─── exp is enforced with zero tolerance ─────────────────────────────────── + + describe('assertTimeClaims: exp tolerance is exactly 0', () => { + it('rejects a token whose exp is the current second', () => { + const now = 1_000_000; + expect(() => assertTimeClaims({ ...payloadAt(0, {}, now), exp: now }, now)).toThrow(TokenExpiredError); + }); + + it('rejects a token that expired 10 seconds ago', () => { + const now = 1_000_000; + expect(() => assertTimeClaims({ ...payloadAt(0, {}, now), exp: now - 10 }, now)).toThrow(TokenExpiredError); + }); + + it('accepts a token that has not expired yet', () => { + const now = 1_000_000; + expect(() => assertTimeClaims({ ...payloadAt(0, {}, now), exp: now + 1 }, now)).not.toThrow(); + }); + + it('tolerates 0s of clock drift, not the old 60s', () => { + const now = 1_000_000; + // 59 seconds past exp used to be accepted with clockTolerance: 60. + expect(() => assertTimeClaims({ ...payloadAt(0, {}, now), exp: now - 59 }, now)).toThrow( + TokenExpiredError, + ); + }); + + it('rejects a token with no usable exp claim', () => { + const now = 1_000_000; + expect(() => assertTimeClaims({ sub: TEST_WALLET } as JwtPayload, now)).toThrow( + 'Malformed JWT payload', + ); + expect(() => + assertTimeClaims({ ...payloadAt(0), exp: Number.NaN }, now), + ).toThrow('Malformed JWT payload'); + }); + }); + + // ─── nbf/iat keep a small, bounded skew allowance ────────────────────────── + + describe('assertTimeClaims: nbf/iat skew', () => { + it('accepts a token that becomes valid within the tolerated skew', () => { + const now = 1_000_000; + expect(CLOCK_SKEW_TOLERANCE_SECONDS).toBe(5); + expect(() => + assertTimeClaims( + { ...payloadAt(600, {}, now), nbf: now + CLOCK_SKEW_TOLERANCE_SECONDS }, + now, + ), + ).not.toThrow(); + }); + + it('rejects a token whose nbf is beyond the tolerated skew', () => { + const now = 1_000_000; + expect(() => + assertTimeClaims( + { ...payloadAt(600, {}, now), nbf: now + CLOCK_SKEW_TOLERANCE_SECONDS + 1 }, + now, + ), + ).toThrow(TokenNotYetValidError); + }); + + it('rejects a token issued further in the future than the skew', () => { + const now = 1_000_000; + expect(() => + assertTimeClaims( + { ...payloadAt(600, {}, now), iat: now + CLOCK_SKEW_TOLERANCE_SECONDS + 1 }, + now, + ), + ).toThrow(TokenNotYetValidError); + }); + }); + + // ─── verifyJwt end to end, with the clock moved forward ──────────────────── + + describe('verifyJwt with a moving clock', () => { + const realNow = Date.now; + + afterEach(() => { + Date.now = realNow; + }); + + function moveClockTo(unixSeconds: number): void { + Date.now = () => unixSeconds * 1000; + } + + it('rejects a token whose exp is the current second (RFC 7519: invalid at exp)', () => { + const exp = Math.floor(Date.now() / 1000); + const token = signToken(payloadAt(0, { exp })); + expect(() => verifyJwt(token)).toThrow(TokenExpiredError); + }); + + it('accepts a token one second before it expires', () => { + const exp = Math.floor(Date.now() / 1000) + 1; + const token = signToken(payloadAt(1, { exp })); + expect(() => verifyJwt(token)).not.toThrow(); + }); + + it('rejects that same token 10 seconds later — the reported 60s window', () => { + const exp = Math.floor(Date.now() / 1000); + const token = signToken(payloadAt(0, { exp })); + + moveClockTo(exp + 10); + + expect(() => verifyJwt(token)).toThrow(TokenExpiredError); + }); + + it('rejects a token 59 seconds past exp (old clockTolerance)', () => { + const exp = Math.floor(Date.now() / 1000); + const token = signToken(payloadAt(0, { exp })); + + moveClockTo(exp + 59); + + expect(() => verifyJwt(token)).toThrow(TokenExpiredError); + }); + + it('still accepts a valid token with a 5s nbf skew after the clock moves', () => { + const issuedAt = Math.floor(Date.now() / 1000); + const token = signToken({ + sub: TEST_WALLET, + iat: issuedAt, + nbf: issuedAt + 4, + exp: issuedAt + 900, + jti: crypto.randomUUID(), + }); + + moveClockTo(issuedAt); + + expect(() => verifyJwt(token)).not.toThrow(); + }); + }); + + // ─── Revocation list is consulted on every authenticated request ─────────── + + describe('revocation list', () => { + it('rejects a revoked token id on the next request', async () => { + const { accessToken } = await issueTokenPair(TEST_WALLET); + const payload = verifyJwt(accessToken); + + const before = await request(app) + .get('/api/v1/webhooks') + .set('Authorization', `Bearer ${accessToken}`); + expect(before.status).toBe(200); + + await revokeAccessToken(payload, 'logout'); + + const after = await request(app) + .get('/api/v1/webhooks') + .set('Authorization', `Bearer ${accessToken}`); + expect(after.status).toBe(401); + expect(after.body.code).toBe(TOKEN_REVOKED_CODE); + }); + + it('reports a token expiring now as 401 (not 200) 10s later over HTTP', async () => { + const issuedAt = Math.floor(Date.now() / 1000); + const token = signToken({ + sub: TEST_WALLET, + iat: issuedAt, + exp: issuedAt, + jti: crypto.randomUUID(), + }); + + const realNow = Date.now; + Date.now = () => (issuedAt + 10) * 1000; + try { + const res = await request(app) + .get('/api/v1/webhooks') + .set('Authorization', `Bearer ${token}`); + + expect(res.status).toBe(401); + expect(res.body.code).not.toBe(200); + expect(res.body.code).toBe('AUTH_TOKEN_INVALID'); + } finally { + Date.now = realNow; + } + }); + + it('logout revokes the presented access token', async () => { + const login = await request(app) + .post('/api/v1/auth/login') + .send({ walletAddress: TEST_WALLET }); + const { accessToken, refreshToken } = login.body; + + const logout = await request(app) + .post('/api/v1/auth/logout') + .set('Authorization', `Bearer ${accessToken}`) + .send({ refreshToken }); + expect(logout.status).toBe(200); + expect(logout.body.revokedAccessToken).toBe(true); + expect(logout.body.refreshSessionRevoked).toBe(true); + + const replay = await request(app) + .get('/api/v1/webhooks') + .set('Authorization', `Bearer ${accessToken}`); + expect(replay.status).toBe(401); + expect(replay.body.code).toBe(TOKEN_REVOKED_CODE); + + // The refresh family is dead too, so the session cannot be resurrected. + const refresh = await request(app) + .post('/api/v1/auth/refresh') + .send({ refreshToken }); + expect(refresh.status).toBe(401); + }); + + it('logout-all revokes other sessions of the same wallet', async () => { + const first = await issueTokenPair(TEST_WALLET); + const second = await issueTokenPair(TEST_WALLET); + + expect( + (await request(app) + .get('/api/v1/webhooks') + .set('Authorization', `Bearer ${first.accessToken}`)).status, + ).toBe(200); + + const logoutAll = await request(app) + .post('/api/v1/auth/logout-all') + .set('Authorization', `Bearer ${first.accessToken}`); + expect(logoutAll.status).toBe(200); + + // A token minted *before* logout-all is dead… + expect( + (await request(app) + .get('/api/v1/webhooks') + .set('Authorization', `Bearer ${first.accessToken}`)).body.code, + ).toBe(TOKEN_REVOKED_CODE); + + // …and so is every other token of that wallet. + expect( + (await request(app) + .get('/api/v1/webhooks') + .set('Authorization', `Bearer ${second.accessToken}`)).body.code, + ).toBe(TOKEN_REVOKED_CODE); + }); + + it('does not affect a different wallet', async () => { + const other = await issueTokenPair(SECOND_TEST_WALLET); + const mine = await issueTokenPair(TEST_WALLET); + + await revokeAllAccessTokens(TEST_WALLET, 'logout'); + + expect(await isAccessTokenRevoked(verifyJwt(mine.accessToken))).toBe(true); + expect(await isAccessTokenRevoked(verifyJwt(other.accessToken))).toBe(false); + }); + + it('a token minted after logout-all is accepted again', async () => { + const before = await issueTokenPair(TEST_WALLET); + await revokeAllAccessTokens(TEST_WALLET, 'logout'); + + // Wait past the marker's `revokedBefore` second so the new token's `iat` + // is strictly greater. + const after = await issueTokenPair(TEST_WALLET); + const afterPayload = verifyJwt(after.accessToken); + const beforePayload = verifyJwt(before.accessToken); + + expect(await isAccessTokenRevoked(beforePayload)).toBe(true); + expect(afterPayload.iat).toBeGreaterThanOrEqual(beforePayload.iat); + }); + }); + + // ─── Store behaviour ─────────────────────────────────────────────────────── + + describe('InMemoryRevocationStore', () => { + let store: InMemoryRevocationStore; + + beforeEach(() => { + store = new InMemoryRevocationStore(); + }); + + it('honours per-token and wallet-wide revocation independently', async () => { + const now = Date.now(); + await store.revoke({ + tokenId: 'token-1', + walletAddress: 'wallet-a', + revokedAt: now, + reason: 'logout', + expiresAt: now + 60_000, + }); + + expect(await store.isRevoked('token-1')).toBe(true); + expect(await store.isRevoked('token-2')).toBe(false); + expect(await store.isWalletRevokedBefore('wallet-a', now - 1)).toBe(false); + }); + + it('wallet marker matches tokens issued at or before it', async () => { + const store2 = new InMemoryRevocationStore(); + const cut = Date.now(); + + await store2.revokeWalletBefore('wallet-b', cut, 'compromised'); + + expect(await store2.isWalletRevokedBefore('wallet-b', cut - 1000)).toBe(true); + expect(await store2.isWalletRevokedBefore('wallet-b', cut)).toBe(true); + expect(await store2.isWalletRevokedBefore('wallet-b', cut + 60_000)).toBe(false); + expect(await store2.isWalletRevokedBefore('wallet-c', cut)).toBe(false); + }); + + it('never moves the wallet marker backwards', async () => { + const later = Date.now(); + await store.revokeWalletBefore('wallet-d', later, 'compromised'); + await store.revokeWalletBefore('wallet-d', later - 60_000, 'logout'); + + // The earlier (narrower) revocation must not widen the existing one. + expect(await store.isWalletRevokedBefore('wallet-d', later + 1)).toBe(false); + }); + + it('revokeAllForWallet records the revocation instead of erasing it', async () => { + const now = Date.now(); + await store.revoke({ + tokenId: 'token-3', + walletAddress: 'wallet-e', + revokedAt: now, + reason: 'logout', + expiresAt: now + 60_000, + }); + + const count = await store.revokeAllForWallet('wallet-e', 'logout'); + expect(count).toBe(1); + // The individual record is gone, but the wallet marker still rejects it. + expect(await store.isRevoked('token-3')).toBe(false); + expect(await store.isWalletRevokedBefore('wallet-e', now)).toBe(true); + }); + + it('drops expired entries', async () => { + await store.revoke({ + tokenId: 'stale', + walletAddress: 'wallet-f', + revokedAt: Date.now() - 120_000, + reason: 'logout', + expiresAt: Date.now() - 60_000, + }); + + expect(await store.isRevoked('stale')).toBe(false); + }); + + it('clear() resets both maps', async () => { + const now = Date.now(); + await store.revoke({ + tokenId: 'token-4', + walletAddress: 'wallet-g', + revokedAt: now, + reason: 'logout', + expiresAt: now + 60_000, + }); + await store.revokeWalletBefore('wallet-g', now, 'logout'); + + await store.clear(); + + expect(await store.isRevoked('token-4')).toBe(false); + expect(await store.isWalletRevokedBefore('wallet-g', now)).toBe(false); + }); + }); + + describe('RedisRevocationStore falls back to the in-process store', () => { + function brokenRedis(): never { + return { + setex: () => Promise.reject(new Error('redis down')), + set: () => Promise.reject(new Error('redis down')), + get: () => Promise.reject(new Error('redis down')), + exists: () => Promise.reject(new Error('redis down')), + del: () => Promise.reject(new Error('redis down')), + sadd: () => Promise.reject(new Error('redis down')), + smembers: () => Promise.reject(new Error('redis down')), + keys: () => Promise.reject(new Error('redis down')), + pipeline: () => { + throw new Error('redis down'); + }, + } as never; + } + + it('still enforces revocation when Redis errors', async () => { + const store: RevocationStore = new RedisRevocationStore(brokenRedis()); + const now = Date.now(); + + await store.revoke({ + tokenId: 'redis-token', + walletAddress: 'wallet-r', + revokedAt: now, + reason: 'logout', + expiresAt: now + 60_000, + }); + + expect(await store.isRevoked('redis-token')).toBe(true); + }); + + it('still enforces the wallet marker when Redis errors', async () => { + const store: RevocationStore = new RedisRevocationStore(brokenRedis()); + const now = Date.now(); + + await store.revokeWalletBefore('wallet-r2', now, 'logout'); + + expect(await store.isWalletRevokedBefore('wallet-r2', now)).toBe(true); + expect(await store.isWalletRevokedBefore('wallet-r2', now + 1_000_000)).toBe(false); + }); + + it('revokeAllForWallet falls back rather than returning 0 silently', async () => { + const store: RevocationStore = new RedisRevocationStore(brokenRedis()); + const now = Date.now(); + + await store.revoke({ + tokenId: 'redis-token-2', + walletAddress: 'wallet-r3', + revokedAt: now, + reason: 'logout', + expiresAt: now + 60_000, + }); + + const count = await store.revokeAllForWallet('wallet-r3', 'logout'); + expect(count).toBe(1); + expect(await store.isWalletRevokedBefore('wallet-r3', now)).toBe(true); + }); + + it('clear() does not throw', async () => { + const store: RevocationStore = new RedisRevocationStore(brokenRedis()); + await expect(store.clear()).resolves.toBeUndefined(); + }); + }); +}); diff --git a/backend/src/__tests__/vaultsListLimits.test.ts b/backend/src/__tests__/vaultsListLimits.test.ts new file mode 100644 index 000000000..9611473f0 --- /dev/null +++ b/backend/src/__tests__/vaultsListLimits.test.ts @@ -0,0 +1,333 @@ +/** + * Regression tests for Issue #1430 — `GET /vaults` did not enforce a maximum + * `limit`, so `?limit=100000` reached Prisma as `take: 100000` and OOM-killed + * the API container. + * + * The load-bearing assertion is the last group: it spies on the Prisma + * client and proves that **no** query is ever issued with a `take` above the + * published ceiling, regardless of what the caller asked for. + */ + +import request from 'supertest'; +import app from '../index'; +import { getPrismaClient } from '../prismaClient'; +import { + DEFAULT_PAGE_SIZE, + MAX_PAGE, + MAX_PAGE_SIZE, + LIMIT_EXCEEDED_CODE, + PaginationLimitError, + enforcePaginationLimits, + resolvePagination, +} from '../middleware/paginationGuard'; +import { DEFAULT_PAGINATION_CONFIG, parsePaginationQuery } from '../pagination'; +import { specs } from '../swagger'; +import { VaultListResponseSchema, validateResponseAgainstSchema } from '../apiContractSnapshots'; +import { prisma } from '../prisma'; +import * as fs from 'fs'; +import * as path from 'path'; + +const VAULTS_PATH = '/api/v1/vaults'; + +/** Smallest `take` the route may ever issue: one lookahead row. */ +const MAX_ALLOWED_TAKE = MAX_PAGE_SIZE + 1; + +const prismaClient = getPrismaClient(); + +async function seedVaults(count: number): Promise { + const existing = await prisma.vault.count(); + if (existing >= count) return; + + for (let i = existing; i < count; i += 1) { + await prismaClient.vault.create({ + data: { + tenantId: 'tenant-1430', + aum: 1000 + i, + tvlUsd: `${1000 + i}.00`, + }, + }); + } +} + +describe('Issue #1430 — GET /api/v1/vaults pagination limits', () => { + beforeAll(async () => { + await seedVaults(55); + }); + + afterAll(async () => { + await prisma.vault.deleteMany({ where: { tenantId: 'tenant-1430' } }); + }); + + // ─── Defaults and accepted values ────────────────────────────────────────── + + it('defaults to 20 items per page', async () => { + const res = await request(app).get(VAULTS_PATH); + + expect(res.status).toBe(200); + expect(res.body.pagination.limit).toBe(DEFAULT_PAGE_SIZE); + expect(res.body.data.length).toBeLessThanOrEqual(DEFAULT_PAGE_SIZE); + }); + + it('honours the maximum permitted limit', async () => { + const res = await request(app).get(`${VAULTS_PATH}?limit=${MAX_PAGE_SIZE}`); + + expect(res.status).toBe(200); + expect(res.body.pagination.limit).toBe(MAX_PAGE_SIZE); + }); + + it('falls back to the default when limit is not a positive integer', async () => { + for (const raw of ['abc', '0', '-5', '1.5', '']) { + const res = await request(app).get(`${VAULTS_PATH}?limit=${encodeURIComponent(raw)}`); + expect(res.status).toBe(200); + expect(res.body.pagination.limit).toBe(DEFAULT_PAGE_SIZE); + } + }); + + it('rejects an oversized limit with 400 LIMIT_EXCEEDED', async () => { + const res = await request(app).get(`${VAULTS_PATH}?limit=100000`); + + expect(res.status).toBe(400); + expect(res.body.code).toBe(LIMIT_EXCEEDED_CODE); + expect(res.body.error).toBe('Bad Request'); + expect(res.body.retryable).toBe(false); + expect(res.body.message).toContain(String(MAX_PAGE_SIZE)); + expect(res.body.details).toMatchObject({ + field: 'limit', + requested: 100000, + maxLimit: MAX_PAGE_SIZE, + }); + }); + + it('rejects every limit above the ceiling, including the boundary + 1', async () => { + const res = await request(app).get(`${VAULTS_PATH}?limit=${MAX_PAGE_SIZE + 1}`); + + expect(res.status).toBe(400); + expect(res.body.code).toBe(LIMIT_EXCEEDED_CODE); + }); + + it('never issues a Prisma read with take above the ceiling', async () => { + const findMany = jest.spyOn(prismaClient.vault, 'findMany'); + const count = jest.spyOn(prismaClient.vault, 'count'); + + try { + const requests = [ + `${VAULTS_PATH}?limit=100000`, + `${VAULTS_PATH}?limit=999999`, + `${VAULTS_PATH}?limit=100000&page=3`, + `${VAULTS_PATH}?limit=51`, + `${VAULTS_PATH}`, + `${VAULTS_PATH}?limit=50`, + `${VAULTS_PATH}?limit=1`, + ]; + + for (const target of requests) { + await request(app).get(target); + } + + const takes = findMany.mock.calls.map(([args]) => (args as { take?: number }).take); + expect(takes.length).toBeGreaterThan(0); + for (const take of takes) { + expect(typeof take).toBe('number'); + expect(take!).toBeLessThanOrEqual(MAX_ALLOWED_TAKE); + } + } finally { + findMany.mockRestore(); + count.mockRestore(); + } + }); + + it('rejects the oversized limit before touching the database at all', async () => { + const findMany = jest.spyOn(prismaClient.vault, 'findMany'); + const count = jest.spyOn(prismaClient.vault, 'count'); + + try { + const res = await request(app).get(`${VAULTS_PATH}?limit=100000`); + + expect(res.status).toBe(400); + expect(findMany).not.toHaveBeenCalled(); + expect(count).not.toHaveBeenCalled(); + } finally { + findMany.mockRestore(); + count.mockRestore(); + } + }); + + // ─── Page clamping ───────────────────────────────────────────────────────── + + it('clamps page into 1..1000 instead of rejecting it', async () => { + const high = await request(app).get(`${VAULTS_PATH}?page=100000`); + expect(high.status).toBe(200); + expect(high.body.pagination.currentPage).toBe(MAX_PAGE); + + const low = await request(app).get(`${VAULTS_PATH}?page=-1`); + expect(low.status).toBe(200); + expect(low.body.pagination.currentPage).toBe(1); + }); + + it('caps the database offset derived from page', async () => { + const findMany = jest.spyOn(prismaClient.vault, 'findMany'); + + try { + await request(app).get(`${VAULTS_PATH}?limit=50&page=100000`); + const [args] = findMany.mock.calls[0] as [{ take: number; skip: number }]; + + expect(args.take).toBeLessThanOrEqual(MAX_ALLOWED_TAKE); + expect(args.skip).toBe((MAX_PAGE - 1) * 50); + } finally { + findMany.mockRestore(); + } + }); + + // ─── Response contract ───────────────────────────────────────────────────── + + it('matches the committed vault-list contract snapshot', async () => { + const res = await request(app).get(`${VAULTS_PATH}?limit=5`); + + expect(res.status).toBe(200); + const parsed = VaultListResponseSchema.safeParse(res.body); + expect(parsed.success).toBe(true); + expect(validateResponseAgainstSchema('GET /api/v1/vaults', res.body).success).toBe(true); + }); + + it('never leaks tenantId', async () => { + const res = await request(app).get(VAULTS_PATH); + + for (const vault of res.body.data as Array>) { + expect(vault).not.toHaveProperty('tenantId'); + } + }); + + // ─── Guard unit behaviour ────────────────────────────────────────────────── + + describe('resolvePagination', () => { + it('applies the documented defaults', () => { + expect(resolvePagination({})).toEqual({ limit: DEFAULT_PAGE_SIZE, page: 1 }); + }); + + it('throws PaginationLimitError above the ceiling', () => { + expect(() => resolvePagination({ limit: '100000' })).toThrow(PaginationLimitError); + }); + + it('honours per-route overrides', () => { + expect(resolvePagination({ limit: '75' }, { maxLimit: 100 })).toEqual({ + limit: 75, + page: 1, + }); + expect(() => resolvePagination({ limit: '75' }, { maxLimit: 50 })).toThrow( + PaginationLimitError, + ); + }); + + it('clamps page and ignores non-integer input', () => { + expect(resolvePagination({ page: '0' }).page).toBe(1); + expect(resolvePagination({ page: 'nope' }).page).toBe(1); + expect(resolvePagination({ page: '100000' }).page).toBe(MAX_PAGE); + expect(resolvePagination({ page: '3' }).page).toBe(3); + }); + + it('rejects repeated query parameters rather than guessing', () => { + expect(resolvePagination({ limit: ['10', '20'] }).limit).toBe(DEFAULT_PAGE_SIZE); + }); + }); + + describe('enforcePaginationLimits middleware', () => { + function runMiddleware(query: Record) { + const req = { query } as never; + const json = jest.fn(); + const res = { + status: jest.fn(() => ({ json })), + setHeader: jest.fn(), + } as never; + const next = jest.fn(); + + enforcePaginationLimits()(req, res, next); + return { req: req as { resolvedPagination?: { limit: number; page: number } }, json, next }; + } + + it('attaches the resolved pagination and continues', () => { + const { req, json, next } = runMiddleware({ limit: '25', page: '4' }); + + expect(next).toHaveBeenCalledTimes(1); + expect(json).not.toHaveBeenCalled(); + expect(req.resolvedPagination).toEqual({ limit: 25, page: 4 }); + }); + + it('sends 400 LIMIT_EXCEEDED and does not continue', () => { + const { json, next } = runMiddleware({ limit: '100000' }); + + expect(next).not.toHaveBeenCalled(); + expect(json).toHaveBeenCalledWith( + expect.objectContaining({ status: 400, code: LIMIT_EXCEEDED_CODE }), + ); + }); + }); + + describe('shared parsePaginationQuery', () => { + function parse(query: Record) { + return parsePaginationQuery({ query } as never); + } + + it('clamps page to maxPage but keeps the default limit of 20', () => { + expect(DEFAULT_PAGINATION_CONFIG.defaultLimit).toBe(20); + expect(parse({})).toMatchObject({ limit: 20 }); + expect(parse({ page: '1000000' })).toMatchObject({ page: MAX_PAGE }); + expect(parse({ page: '-1' })).toMatchObject({ page: 1 }); + }); + }); + + // ─── OpenAPI contract ────────────────────────────────────────────────────── + + describe('OpenAPI documentation of the pagination ceiling', () => { + const spec = specs as unknown as { + paths: Record> } }>; + components: { + parameters: Record< + string, + { name: string; schema: { maximum?: number; minimum?: number; default?: number } } + >; + }; + }; + + it('publishes maxLimit and the page ceiling as reusable parameters', () => { + expect(spec.components.parameters.pageSize).toMatchObject({ + name: 'limit', + schema: { minimum: 1, maximum: MAX_PAGE_SIZE, default: DEFAULT_PAGE_SIZE }, + }); + expect(spec.components.parameters.pageNumber).toMatchObject({ + name: 'page', + schema: { minimum: 1, maximum: MAX_PAGE }, + }); + }); + + it('references those parameters from GET /api/v1/vaults', () => { + const parameters = spec.paths[VAULTS_PATH]?.get?.parameters ?? []; + const refs = parameters.map((parameter) => parameter.$ref); + + expect(refs).toContain('#/components/parameters/pageSize'); + expect(refs).toContain('#/components/parameters/pageNumber'); + }); + + it('documents the LIMIT_EXCEEDED failure mode', () => { + const specWithResponses = specs as unknown as { + paths: Record< + string, + { get?: { responses?: Record }> } } + >; + }; + const example = + specWithResponses.paths[VAULTS_PATH]?.get?.responses?.['400']?.content?.[ + 'application/json' + ]?.example; + + expect(example?.code).toBe(LIMIT_EXCEEDED_CODE); + }); + + it('keeps the committed openapi.json in sync with the spec source', () => { + const committed = JSON.parse( + fs.readFileSync(path.join(__dirname, '..', '..', 'openapi.json'), 'utf8'), + ) as typeof specs; + + expect(JSON.stringify(committed)).toBe(JSON.stringify(specs)); + }); + }); +}); diff --git a/backend/src/apiContractSnapshots.ts b/backend/src/apiContractSnapshots.ts index acaa02eb4..cd93e22a6 100644 --- a/backend/src/apiContractSnapshots.ts +++ b/backend/src/apiContractSnapshots.ts @@ -17,6 +17,7 @@ export const CRITICAL_ENDPOINTS = [ 'GET /ready', 'GET /api/v1/vault/summary', 'GET /api/v1/transactions', + 'GET /api/v1/vaults', ] as const; export type CriticalEndpoint = (typeof CRITICAL_ENDPOINTS)[number]; @@ -110,11 +111,34 @@ export const TransactionsListResponseSchema = z }) .strict(); +/** + * Public projection of a vault row. `tenantId` is intentionally absent — the + * list route never exposes it (Issue #1430). + */ +export const VaultItemSchema = z + .object({ + id: z.string(), + aum: z.number(), + tvlUsd: z.string().nullable(), + createdAt: z.string(), + updatedAt: z.string(), + }) + .strict(); + +export const VaultListResponseSchema = z + .object({ + data: z.array(VaultItemSchema), + pagination: PaginationMetaSchema, + timestamp: z.string(), + }) + .strict(); + export const ENDPOINT_SCHEMAS: Record = { 'GET /health': HealthResponseSchema, 'GET /ready': ReadyResponseSchema, 'GET /api/v1/vault/summary': VaultSummaryResponseSchema, 'GET /api/v1/transactions': TransactionsListResponseSchema, + 'GET /api/v1/vaults': VaultListResponseSchema, }; export function endpointToFilename(endpoint: CriticalEndpoint): string { @@ -249,7 +273,7 @@ export function diffSchemaShapes( } for (const key of baselineRequired) { - if (!(key in baseline.properties ?? {})) { + if (!(key in baselineProps)) { issues.push({ path: at(key), message: 'required field missing from snapshot properties (orphaned reference)' }); } if (!(key in currentProps)) { @@ -260,8 +284,10 @@ export function diffSchemaShapes( } for (const key of currentRequired) { - if (!(key in current.properties ?? {})) { + if (!(key in (current.properties ?? {}))) { issues.push({ path: at(key), message: 'required field missing from live schema properties (invalid schema)' }); + continue; + } if (!baselineRequired.has(key)) { issues.push({ path: at(key), message: 'field is now required — regenerate snapshots with npm run snapshots:write' }); } diff --git a/backend/src/auth.ts b/backend/src/auth.ts index a000da84c..46c4b510a 100644 --- a/backend/src/auth.ts +++ b/backend/src/auth.ts @@ -48,6 +48,12 @@ import { type WalletAction, } from './walletNonce'; import { buildWalletSignMessage } from './walletSignature'; +import { + getRevocationStore, + setRevocationStore, + RedisRevocationStore, + type RevocationReason, +} from './tokenRevocation'; // ─── Config ─────────────────────────────────────────────────────────────────── @@ -155,6 +161,8 @@ export interface JwtPayload { iat: number; // issued-at (unix seconds) exp: number; // expiry (unix seconds) jti: string; // JWT ID (unique per token) + /** Optional "not before" (unix seconds). Tolerated by CLOCK_SKEW_TOLERANCE_SECONDS. */ + nbf?: number; } interface RefreshTokenEntry { @@ -287,6 +295,24 @@ function createRefreshTokenStore(): IRefreshTokenStore { const refreshTokenStore: IRefreshTokenStore = createRefreshTokenStore(); +// The revocation list is consulted on every authenticated request, so in a +// multi-instance deployment it has to be shared: a logout handled by pod A +// must be visible to pod B or the "stolen token" window re-opens. Redis is +// used when REDIS_URL is configured; otherwise the in-process store applies +// (single-instance / local development). +function initRevocationStore(): void { + const redisUrl = process.env.REDIS_URL; + if (!redisUrl) return; + + const redis = new Redis(redisUrl, { lazyConnect: true, enableOfflineQueue: false }); + redis.on('error', (err) => { + logger.log('error', 'Redis revocation store error', { error: err.message }); + }); + setRevocationStore(new RedisRevocationStore(redis)); +} + +initRevocationStore(); + // ─── HS256 JWT Helpers ──────────────────────────────────────────────────────── function base64UrlEncode(input: string | Buffer): string { @@ -315,6 +341,65 @@ function signJwt(payload: JwtPayload): string { return `${signingInput}.${sig}`; } +// ─── Time-claim validation (Issue #1431) ────────────────────────────────────────────────── + +/** + * Tolerance, in seconds, allowed on `nbf`/`iat`. + * + * This exists purely to absorb clock skew between the API pod and whatever + * minted the token. It is deliberately NOT applied to `exp`: a token whose + * lifetime has run out must stop working the instant it does, otherwise a + * stolen bearer token keeps authorising `POST /vault/:id/withdraw` for an + * extra minute after the user logs out. + */ +export const CLOCK_SKEW_TOLERANCE_SECONDS = 5; + +/** Thrown when a token's `exp` has passed. */ +export class TokenExpiredError extends Error { + constructor(message = 'JWT has expired') { + super(message); + this.name = 'TokenExpiredError'; + } +} + +/** Thrown when a token is not valid yet beyond the tolerated clock skew. */ +export class TokenNotYetValidError extends Error { + constructor(message = 'JWT is not yet valid') { + super(message); + this.name = 'TokenNotYetValidError'; + } +} + +/** + * Validates the time-based claims of an already signature-verified payload. + * + * `exp` uses **zero** tolerance and is inclusive of the expiry second: RFC 7519 + * says a token is invalid at and after `exp`, so `exp <= now` is rejected. + * `nbf`/`iat` get {@link CLOCK_SKEW_TOLERANCE_SECONDS} of slack. + */ +export function assertTimeClaims(payload: JwtPayload, nowSeconds = Math.floor(Date.now() / 1000)): void { + if (typeof payload.exp !== 'number' || !Number.isFinite(payload.exp)) { + throw new Error('Malformed JWT payload'); + } + + // Strict: no tolerance whatsoever, and the expiry second itself is too late. + if (payload.exp <= nowSeconds) { + throw new TokenExpiredError(); + } + + if (typeof payload.nbf === 'number' && Number.isFinite(payload.nbf)) { + if (payload.nbf > nowSeconds + CLOCK_SKEW_TOLERANCE_SECONDS) { + throw new TokenNotYetValidError(); + } + } + + if (typeof payload.iat === 'number' && Number.isFinite(payload.iat)) { + if (payload.iat > nowSeconds + CLOCK_SKEW_TOLERANCE_SECONDS) { + throw new TokenNotYetValidError('JWT was issued in the future'); + } + } +} + /** * Verifies a JWT string. * Returns the decoded payload on success. @@ -346,8 +431,7 @@ export function verifyJwt(token: string): JwtPayload { throw new Error('Malformed JWT payload'); } - const now = Math.floor(Date.now() / 1000); - if (payload.exp < now) throw new Error('JWT has expired'); + assertTimeClaims(payload); return payload; } @@ -391,6 +475,83 @@ export async function issueTokenPair(walletAddress: string, familyId?: string): // ─── Session Revocation ──────────────────────────────────────────────────────── +/** + * Adds a specific access token to the revocation list. + * + * The entry is retained until the token's own `exp`, at which point the token + * is rejected by the strict expiry check anyway, so the list can never grow + * unbounded. + */ +export async function revokeAccessToken( + payload: JwtPayload, + reason: RevocationReason = 'logout' +): Promise { + if (!payload?.jti || !payload?.sub) { + return; + } + + await getRevocationStore().revoke({ + tokenId: payload.jti, + walletAddress: payload.sub, + revokedAt: Date.now(), + reason, + expiresAt: payload.exp * 1000, + }); + + logger.log('info', 'Access token revoked', { + reason, + wallet: payload.sub.slice(0, 8) + '…', + jti: payload.jti, + }); +} + +/** + * Revokes every access token currently issued to `walletAddress`. + * This is used for /auth/logout-all and for "this device is compromised". + */ +export async function revokeAllAccessTokens( + walletAddress: string, + reason: RevocationReason = 'logout' +): Promise { + const normalizedAddress = normalizeWalletAddress(walletAddress); + const store = getRevocationStore(); + + await store.revokeWalletBefore(normalizedAddress, Date.now(), reason); + const removed = await store.revokeAllForWallet(normalizedAddress, reason); + + logger.log('info', 'All access tokens revoked for wallet', { + reason, + wallet: normalizedAddress.slice(0, 8) + '…', + revokedCount: removed, + }); + + return removed; +} + +/** + * True when the token has been revoked, either individually (`jti`) or because + * the whole wallet was revoked at or after the token was issued. + */ +export async function isAccessTokenRevoked(payload: JwtPayload): Promise { + if (!payload?.jti) return false; + + const store = getRevocationStore(); + + if (await store.isRevoked(payload.jti)) { + return true; + } + + // `iat` is unix seconds; the wallet marker uses unix milliseconds. The + // wallet is normalised on both sides so a lower-case `sub` cannot dodge a + // revocation written for the canonical (upper-case) address. + if (typeof payload.iat === 'number' && Number.isFinite(payload.iat)) { + return store.isWalletRevokedBefore(normalizeWalletAddress(payload.sub), payload.iat * 1000); + } + + return false; +} + + /** * Revokes the current session (all tokens in the same family). * This is used for /auth/logout. @@ -536,11 +697,21 @@ export interface AuthenticatedRequest extends Request { jwtPayload?: JwtPayload; } +/** Error code returned when a structurally valid token has been revoked. */ +export const TOKEN_REVOKED_CODE = 'TOKEN_REVOKED'; + /** * Express middleware that validates the Bearer access token from the * Authorization header and attaches the decoded payload to req.jwtPayload. * - * Returns 401 for missing / invalid / expired tokens. + * Two independent checks run on **every** request (Issue #1431): + * + * 1. Cryptographic verification plus time-claim validation, where `exp` is + * enforced with zero tolerance (`assertTimeClaims`). + * 2. A revocation-list lookup covering both the individual token id and any + * wallet-wide revocation issued at or after the token was minted. + * + * Returns 401 for missing / invalid / expired / revoked tokens. */ export function requireAuth( req: AuthenticatedRequest, @@ -560,9 +731,9 @@ export function requireAuth( return; } + let payload: JwtPayload; try { - req.jwtPayload = verifyJwt(match[1]); - next(); + payload = verifyJwt(match[1]); } catch (err) { sendApiError(req, res, { status: 401, @@ -570,7 +741,42 @@ export function requireAuth( message: err instanceof Error ? err.message : 'Invalid token', retryable: false, }); + return; } + + // The revocation lookup is async, so the rest of the chain runs from a + // continuation. Express ignores middleware return values, so returning + // before `next()` is safe for route usage; the single direct caller + // (authenticateTransactionExport) passes a synchronous next() and is + // unaffected. + isAccessTokenRevoked(payload).then( + (revoked) => { + if (revoked) { + sendApiError(req, res, { + status: 401, + code: TOKEN_REVOKED_CODE, + message: 'Token has been revoked. Please sign in again.', + retryable: false, + }); + return; + } + req.jwtPayload = payload; + next(); + }, + (err) => { + logger.log('error', 'Revocation lookup failed', { + error: err instanceof Error ? err.message : String(err), + }); + // Fail closed: an unreachable revocation store must not silently + // downgrade to "token is fine". + sendApiError(req, res, { + status: 503, + code: 'AUTH_REVOCATION_UNAVAILABLE', + message: 'Unable to verify session state. Please retry.', + retryable: true, + }); + }, + ); } // ─── Auth Route Handlers ────────────────────────────────────────────────────── diff --git a/backend/src/idempotency.ts b/backend/src/idempotency.ts index 1ab4b0448..fcc20199c 100644 --- a/backend/src/idempotency.ts +++ b/backend/src/idempotency.ts @@ -26,6 +26,22 @@ import { prisma } from './prisma'; import { logger } from './middleware/structuredLogging'; import type { Request, Response, NextFunction } from 'express'; +/** + * The replay cache used by money-moving operations (deposits, withdrawals, + * transfers) lives in `idempotencyStore.ts`; it is re-exported here so callers + * keep a single import surface for the whole idempotency feature. + */ +export { + idempotencyStore, + IdempotencyStore, + IdempotencyConflictError, + buildIdempotencyFingerprint, + getIdempotencyHashThreshold, + type IdempotentOperationResult, + type IdempotencyKeyInfo, + type IdempotencyMetrics, +} from './idempotencyStore'; + // ─── Types ────────────────────────────────────────────────────────────────── export interface IdempotencyRecord { diff --git a/backend/src/idempotencyStore.ts b/backend/src/idempotencyStore.ts new file mode 100644 index 000000000..41e4e6d3e --- /dev/null +++ b/backend/src/idempotencyStore.ts @@ -0,0 +1,360 @@ +/** + * @file idempotencyStore.ts + * Idempotency key store backed by Redis (when available) with NodeCache in-process fallback. + * + * Issue #811: Multi-instance deployments previously lost idempotency guarantees on pod + * recycle because responses were stored only in NodeCache (in-process memory). This revision + * persists completed responses to Redis using SET … EX so all replicas share the same store. + * + * Behavior when Redis is unavailable: + * - Falls back to NodeCache automatically (fail-open). + * - A warning is logged so operators are aware of the degraded guarantee. + * + * Existing observability API (inspectKeys, deleteKey, clear, getMetrics) is preserved; + * operations apply to whichever backend holds the key. + * + * Scope: this module owns the *replay cache* used by money-moving operations + * (`idempotencyStore.execute`). The durable, per-key record lifecycle used by the + * `enforceIdempotency` middleware lives in `idempotency.ts` (Prisma-backed). + */ + +import crypto from 'crypto'; +import NodeCache from 'node-cache'; +import { redisClientManager } from './rateLimiter'; + +// ─── Public Types ───────────────────────────────────────────────────────────── + +export interface IdempotentOperationResult { + statusCode: number; + body: T; +} + +/** Metadata attached to every idempotency key entry. */ +export interface IdempotencyKeyMetadata { + /** ISO-8601 timestamp when the key was first stored. */ + createdAt: string; + /** ISO-8601 timestamp of the most recent access (read or write). */ + lastAccessedAt: string; + /** Number of times this key has been replayed (returned cached result). */ + replayCount: number; + /** Current state of the entry. */ + status: 'pending' | 'completed'; +} + +/** Summary returned by GET /admin/idempotency/keys. */ +export interface IdempotencyKeyInfo { + key: string; + metadata: IdempotencyKeyMetadata; +} + +/** Snapshot of store-wide observability counters. */ +export interface IdempotencyMetrics { + hits: number; + conflicts: number; + evictions: number; + activeKeys: number; + pendingKeys: number; +} + +// ─── Internal Types ─────────────────────────────────────────────────────────── + +interface StoredResponse extends IdempotentOperationResult { + fingerprint: string; + metadata: IdempotencyKeyMetadata; +} + +interface PendingOperation { + fingerprint: string; + promise: Promise>; + metadata: IdempotencyKeyMetadata; +} + +// ─── Errors ─────────────────────────────────────────────────────────────────── + +export class IdempotencyConflictError extends Error { + constructor(message = 'Idempotency key already used for a different request body') { + super(message); + this.name = 'IdempotencyConflictError'; + } +} + +// ─── Redis key prefix ───────────────────────────────────────────────────────── + +const REDIS_PREFIX = 'idempotency:'; + +// ─── Store ──────────────────────────────────────────────────────────────────── + +export class IdempotencyStore { + /** Fallback in-process store used when Redis is unavailable. */ + private readonly localCache: NodeCache; + private readonly pendingResponses = new Map>(); + + // Observability counters + private _hits = 0; + private _conflicts = 0; + private _evictions = 0; + + constructor(private readonly ttlMs = 24 * 60 * 60 * 1000) { + const ttlSeconds = Math.max(1, Math.ceil(this.ttlMs / 1000)); + this.localCache = new NodeCache({ stdTTL: ttlSeconds, checkperiod: ttlSeconds }); + this.localCache.on('expired', () => { this._evictions++; }); + } + + // ─── Redis helpers ───────────────────────────────────────────────────────── + + private redisKey(key: string): string { + return `${REDIS_PREFIX}${key}`; + } + + private get redis() { + const client = redisClientManager.getClient(); + return redisClientManager.isReady() && client ? client : null; + } + + private async redisGet(key: string): Promise | null> { + const r = this.redis; + if (!r) return null; + try { + const raw = await r.get(this.redisKey(key)); + return raw ? (JSON.parse(raw) as StoredResponse) : null; + } catch { + return null; + } + } + + private async redisSet(key: string, value: StoredResponse): Promise { + const r = this.redis; + if (!r) return; + try { + const ttlSeconds = Math.max(1, Math.ceil(this.ttlMs / 1000)); + await r.set(this.redisKey(key), JSON.stringify(value), 'EX', ttlSeconds); + } catch (err) { + console.log(JSON.stringify({ level: 'warn', event: 'idempotency_redis_write_fail', key, reason: (err as Error).message })); + } + } + + private async redisDel(key: string): Promise { + const r = this.redis; + if (!r) return false; + try { + return (await r.del(this.redisKey(key))) > 0; + } catch { + return false; + } + } + + // ─── Core execute ────────────────────────────────────────────────────────── + + async execute( + key: string, + fingerprint: string, + operation: () => Promise> + ): Promise<{ result: IdempotentOperationResult; replayed: boolean }> { + const now = new Date().toISOString(); + + // 1. Check Redis first, then local cache + let completed = await this.redisGet(key); + if (!completed) { + completed = this.localCache.get>(key) ?? null; + } + + if (completed) { + if (completed.fingerprint !== fingerprint) { + this._conflicts++; + throw new IdempotencyConflictError(); + } + this._hits++; + completed.metadata.lastAccessedAt = now; + completed.metadata.replayCount++; + // Refresh in both backends; errors are non-fatal + await this.redisSet(key, completed); + this.localCache.set(key, completed); + return { result: { statusCode: completed.statusCode, body: completed.body }, replayed: true }; + } + + // 2. Currently in-flight + const pendingOperation = this.pendingResponses.get(key) as PendingOperation | undefined; + if (pendingOperation) { + if (pendingOperation.fingerprint !== fingerprint) { + this._conflicts++; + throw new IdempotencyConflictError(); + } + this._hits++; + pendingOperation.metadata.lastAccessedAt = now; + pendingOperation.metadata.replayCount++; + const replayed = await pendingOperation.promise; + return { result: { statusCode: replayed.statusCode, body: replayed.body }, replayed: true }; + } + + // 3. First execution + const metadata: IdempotencyKeyMetadata = { createdAt: now, lastAccessedAt: now, replayCount: 0, status: 'pending' }; + + const operationPromise = (async () => { + const result = await operation(); + const stored: StoredResponse = { + ...result, + fingerprint, + metadata: { ...metadata, status: 'completed', lastAccessedAt: new Date().toISOString() }, + }; + // Persist to Redis (primary) and local cache (fallback/fast-path) + await this.redisSet(key, stored); + this.localCache.set(key, stored, this.ttlMs / 1000); + return stored; + })(); + + this.pendingResponses.set(key, { fingerprint, promise: operationPromise, metadata }); + + try { + const stored = await operationPromise; + return { result: { statusCode: stored.statusCode, body: stored.body }, replayed: false }; + } finally { + this.pendingResponses.delete(key); + } + } + + // ─── Inspection ──────────────────────────────────────────────────────────── + + inspectKeys(prefix?: string): IdempotencyKeyInfo[] { + const results: IdempotencyKeyInfo[] = []; + for (const key of this.localCache.keys()) { + if (prefix && !key.startsWith(prefix)) continue; + const entry = this.localCache.get>(key); + if (entry) results.push({ key, metadata: { ...entry.metadata } }); + } + for (const [key, pending] of this.pendingResponses.entries()) { + if (prefix && !key.startsWith(prefix)) continue; + if (!results.some((r) => r.key === key)) { + results.push({ key, metadata: { ...pending.metadata } }); + } + } + return results; + } + + // ─── Targeted deletion ───────────────────────────────────────────────────── + + async deleteKey(key: string): Promise { + const deletedLocal = this.localCache.del(key) > 0; + const deletedPending = this.pendingResponses.delete(key); + const deletedRedis = await this.redisDel(key); + if (deletedLocal || deletedPending || deletedRedis) { + this._evictions++; + return true; + } + return false; + } + + // ─── Global clear (admin only) ───────────────────────────────────────────── + + clear(): void { + const count = this.localCache.keys().length + this.pendingResponses.size; + this._evictions += count; + this.localCache.flushAll(); + this.pendingResponses.clear(); + // Note: Redis keys are prefixed with REDIS_PREFIX; a full Redis FLUSHDB is intentionally + // not issued here to avoid clearing unrelated keys. Use deleteKey() per-key when needed. + } + + // ─── Observability ───────────────────────────────────────────────────────── + + getMetrics(): IdempotencyMetrics { + return { + hits: this._hits, + conflicts: this._conflicts, + evictions: this._evictions, + activeKeys: this.localCache.keys().length, + pendingKeys: this.pendingResponses.size, + }; + } + + // ─── Retention cleanup ───────────────────────────────────────────────────── + + async pruneStaleKeys( + retentionMs: number, + dryRun = false, + ): Promise<{ pruned: number; localPruned: number; redisPruned: number }> { + const cutoff = Date.now() - retentionMs; + let localPruned = 0; + let redisPruned = 0; + + for (const key of this.localCache.keys()) { + const entry = this.localCache.get>(key); + if (!entry) continue; + const createdAt = Date.parse(entry.metadata.createdAt); + if (Number.isNaN(createdAt) || createdAt >= cutoff) continue; + if (!dryRun) { + this.localCache.del(key); + this._evictions++; + } + localPruned++; + } + + const r = this.redis; + if (r) { + let cursor = '0'; + do { + const [nextCursor, keys] = await r.scan(cursor, 'MATCH', `${REDIS_PREFIX}*`, 'COUNT', 100); + cursor = nextCursor; + for (const redisKey of keys) { + try { + const raw = await r.get(redisKey); + if (!raw) continue; + const entry = JSON.parse(raw) as StoredResponse; + const createdAt = Date.parse(entry.metadata?.createdAt ?? ''); + const ttl = await r.ttl(redisKey); + const isStale = (!Number.isNaN(createdAt) && createdAt < cutoff) || ttl === 0; + if (!isStale) continue; + if (!dryRun) { + await r.del(redisKey); + this._evictions++; + } + redisPruned++; + } catch { + if (!dryRun) { + await r.del(redisKey); + this._evictions++; + } + redisPruned++; + } + } + } while (cursor !== '0'); + } + + return { pruned: localPruned + redisPruned, localPruned, redisPruned }; + } +} + +// ─── Singleton ──────────────────────────────────────────────────────────────── + +export const idempotencyStore = new IdempotencyStore( + parseInt(process.env.IDEMPOTENCY_KEY_TTL_MS || '86400000', 10) +); + +// ─── Fingerprint helper ─────────────────────────────────────────────────────── + +export function getIdempotencyHashThreshold(): number { + return parseInt(process.env.IDEMPOTENCY_HASH_THRESHOLD_BYTES || '4096', 10); +} + +export function buildIdempotencyFingerprint(payload: unknown): string { + const stable = stableStringify(payload); + const byteLength = Buffer.byteLength(stable, 'utf-8'); + if (byteLength > getIdempotencyHashThreshold()) { + return `hashv1:${crypto.createHash('sha256').update(stable).digest('hex')}`; + } + return stable; +} + +function stableStringify(value: unknown): string { + if (value === null) return 'null'; + if (value instanceof Date) return JSON.stringify(value.toISOString()); + if (typeof value !== 'object') return JSON.stringify(value); + + if (Array.isArray(value)) { + return `[${value.map((item) => stableStringify(item)).join(',')}]`; + } + + const record = value as Record; + const keys = Object.keys(record).sort(); + const serialized = keys.map((key) => `${JSON.stringify(key)}:${stableStringify(record[key])}`); + return `{${serialized.join(',')}}`; +} diff --git a/backend/src/index.ts b/backend/src/index.ts index c0fd1750b..e957cd9de 100644 --- a/backend/src/index.ts +++ b/backend/src/index.ts @@ -10,7 +10,17 @@ initTracing(); import express, { Express, Request, Response, NextFunction, ErrorRequestHandler } from 'express'; import NodeCache from 'node-cache'; -import { loginHandler, nonceHandler, refreshHandler, requireAuth, verifyJwt } from './auth'; +import { + loginHandler, + nonceHandler, + refreshHandler, + requireAuth, + verifyJwt, + revokeAccessToken, + revokeAllAccessTokens, + revokeCurrentSession, + revokeAllSessions, +} from './auth'; import { authLimiter, authIpLimiter, @@ -103,6 +113,7 @@ import { createVersionDiscoveryRouter } from './routes/apiVersions'; import { GracefulShutdownHandler } from './gracefulShutdown'; import { db } from './database'; import vaultRouter from './vaultEndpoints'; +import vaultsListRouter from './routes/vaults'; import walletAliasRouter from './walletAliasEndpoints'; import { walletAliasMappingService } from './walletAliasService'; import transactionRouter from './transactionEndpoints'; @@ -911,6 +922,7 @@ app.use('/api', createVersionDiscoveryRouter()); // Mount routers under /api/v1 apiV1.use('/vault', vaultRouter); +apiV1.use('/vaults', vaultsListRouter); apiV1.use('/wallet-aliases', walletAliasRouter); apiV1.use('/referrals', referralRouter); apiV1.use('/transactions', transactionRouter); @@ -979,16 +991,42 @@ app.use('/admin', validateApiKey, adminRbacMiddleware); /** * POST /api/v1/auth/logout - * Revokes the current session. Requires Bearer token. + * + * Revokes the presented access token (`jti`) on the revocation list, so the + * very next request with it is rejected 401 `TOKEN_REVOKED` (Issue #1431). + * When a `refreshToken` is also supplied the whole refresh family is revoked, + * so the session cannot be resurrected by rotation. */ -apiV1.post('/auth/logout', readsLimiter, requireAuth, (req: Request, res: Response) => { +apiV1.post('/auth/logout', readsLimiter, requireAuth, async (req: Request, res: Response) => { + const authReq = req as import('./auth').AuthenticatedRequest; + const payload = authReq.jwtPayload; + const walletAddress = payload?.sub; + + if (!payload || !walletAddress) { + res.status(500).json({ + error: 'Internal Server Error', + status: 500, + message: 'Unable to determine authenticated wallet', + }); + return; + } + try { - const authReq = req as import('./auth').AuthenticatedRequest; - const walletAddress = authReq.jwtPayload?.sub; - if (!walletAddress) throw new Error('Unable to determine authenticated wallet'); + await revokeAccessToken(payload, 'logout'); + + const refreshToken = + typeof (req.body as { refreshToken?: unknown } | undefined)?.refreshToken === 'string' + ? ((req.body as { refreshToken: string }).refreshToken) + : undefined; + if (refreshToken) { + await revokeCurrentSession(refreshToken); + } + res.status(200).json({ message: 'Session revoked successfully', - walletAddress: walletAddress.slice(0, 8) + '…', + walletAddress: walletAddress.slice(0, 8) + '…', + revokedAccessToken: true, + refreshSessionRevoked: Boolean(refreshToken), timestamp: new Date().toISOString(), }); } catch (err) { @@ -1002,17 +1040,34 @@ apiV1.post('/auth/logout', readsLimiter, requireAuth, (req: Request, res: Respon /** * POST /api/v1/auth/logout-all - * Revokes all active sessions for the authenticated wallet. + * + * Revokes every access token currently issued to the wallet by writing a + * wallet-wide high-water mark, so tokens minted before now stop working even + * though we hold no list of them (Issue #1431). Also revokes all refresh + * families for the wallet. */ -apiV1.post('/auth/logout-all', readsLimiter, requireAuth, (req: Request, res: Response) => { +apiV1.post('/auth/logout-all', readsLimiter, requireAuth, async (req: Request, res: Response) => { + const authReq = req as import('./auth').AuthenticatedRequest; + const walletAddress = authReq.jwtPayload?.sub; + + if (!walletAddress) { + res.status(500).json({ + error: 'Internal Server Error', + status: 500, + message: 'Unable to determine authenticated wallet', + }); + return; + } + try { - const authReq = req as import('./auth').AuthenticatedRequest; - const walletAddress = authReq.jwtPayload?.sub; - if (!walletAddress) throw new Error('Unable to determine authenticated wallet'); + const revokedAccessTokens = await revokeAllAccessTokens(walletAddress, 'logout'); + const revokedRefreshTokens = await revokeAllSessions(walletAddress); + res.status(200).json({ message: 'All sessions revoked successfully', - walletAddress: walletAddress.slice(0, 8) + '…', - revokedCount: 1, + walletAddress: walletAddress.slice(0, 8) + '…', + revokedCount: revokedRefreshTokens, + revokedAccessTokens, timestamp: new Date().toISOString(), }); } catch (err) { @@ -5404,6 +5459,7 @@ app.use((req: Request, res: Response) => { status: 404, code: 'ROUTE_NOT_FOUND', message: `Cannot ${req.method} ${req.originalUrl}`, + path: req.originalUrl, details: { path: req.originalUrl }, retryable: false, }); diff --git a/backend/src/middleware/apiError.ts b/backend/src/middleware/apiError.ts index a59b60aa5..469fe5dd2 100644 --- a/backend/src/middleware/apiError.ts +++ b/backend/src/middleware/apiError.ts @@ -10,6 +10,19 @@ export interface ApiErrorOptions { retryable?: boolean; retryAfterSeconds?: number | null; error?: string; + /** + * Short, stable headline for the failure. Part of the published error + * contract (see `ErrorEnvelope` in openapi.json) alongside `message`. + */ + summary?: string; + /** + * Field-level failures, surfaced at the top level for clients that read + * `errors` directly. `details` carries the same payload for the generic + * envelope. + */ + errors?: unknown[]; + /** Request path, echoed on routing failures so clients can log it directly. */ + path?: string; } export function sendApiError( @@ -30,6 +43,9 @@ export function sendApiError( code: options.code, message: options.message, retryable: options.retryable ?? options.status >= 500, + ...(options.summary !== undefined ? { summary: options.summary } : {}), + ...(options.errors !== undefined ? { errors: options.errors } : {}), + ...(options.path !== undefined ? { path: options.path } : {}), ...(options.details !== undefined ? { details: options.details } : {}), ...(correlationId ? { correlationId } : {}), ...(traceId ? { traceId } : {}), diff --git a/backend/src/middleware/paginationGuard.ts b/backend/src/middleware/paginationGuard.ts new file mode 100644 index 000000000..ab626d5a1 --- /dev/null +++ b/backend/src/middleware/paginationGuard.ts @@ -0,0 +1,160 @@ +/** + * @file paginationGuard.ts + * Hard caps for caller-controlled pagination on database-backed list routes. + * + * Issue #1430: `GET /vaults` forwarded `req.query.limit` straight into + * Prisma's `take`, so an unauthenticated caller could ask for + * `?limit=100000` and force the API to materialise 100k rows plus their + * relations — a memory spike large enough to OOM-kill the 512 MB + * container. + * + * Chosen behaviour: **reject, do not silently clamp.** A caller that asks + * for `limit=100000` has a bug (or is probing), and quietly returning 50 + * rows hides that behind a paginated response that looks complete. So an + * over-max `limit` fails fast with `400` and `code: 'LIMIT_EXCEEDED'`, and + * the response body names the ceiling. `page` is different: an out-of-range + * page number is a benign mistake, so it is clamped to `1..MAX_PAGE` and + * the effective value is echoed back in the pagination envelope. + * + * Anything that turns a request parameter into a Prisma `take`/`skip` + * should sit behind `enforcePaginationLimits()`. + */ + +import type { Request, Response, NextFunction } from 'express'; +import { sendApiError } from './apiError'; + +/** Page size used when the caller does not ask for one. */ +export const DEFAULT_PAGE_SIZE = 20; + +/** + * Hard ceiling on rows read per request. Sized so the largest permitted + * page stays comfortably inside the API container's memory budget. + */ +export const MAX_PAGE_SIZE = 50; + +/** + * Hard ceiling on offset pagination depth. `page * limit` is what reaches + * Prisma as `skip`, so an unbounded page number is an unbounded database + * offset scan. + */ +export const MAX_PAGE = 1000; + +export interface PaginationLimitOptions { + /** Rows per page when `limit` is absent. Defaults to {@link DEFAULT_PAGE_SIZE}. */ + defaultLimit?: number; + /** Largest permitted `limit`. Defaults to {@link MAX_PAGE_SIZE}. */ + maxLimit?: number; + /** Largest permitted `page`. Defaults to {@link MAX_PAGE}. */ + maxPage?: number; +} + +export interface ResolvedPagination { + limit: number; + page: number; +} + +/** Error code returned for a `limit` above the configured ceiling. */ +export const LIMIT_EXCEEDED_CODE = 'LIMIT_EXCEEDED'; + +/** + * Parses a raw query-string value that is expected to hold a base-10 + * integer. Returns `null` for anything that is not an exact integer so + * callers can tell "absent" from "garbage". + */ +function parseIntegerQueryValue(raw: unknown): number | null { + if (raw === undefined || raw === null) return null; + if (Array.isArray(raw)) return null; + if (typeof raw !== 'string') return null; + if (!/^-?\d+$/.test(raw.trim())) return null; + const parsed = Number.parseInt(raw.trim(), 10); + return Number.isSafeInteger(parsed) ? parsed : null; +} + +/** + * Resolves the effective `limit` and `page` for a request. + * + * Throws {@link PaginationLimitError} when `limit` exceeds the ceiling, so + * the middleware can turn it into a `400` instead of running a query. + */ +export function resolvePagination( + query: Record, + options: PaginationLimitOptions = {} +): ResolvedPagination { + const defaultLimit = options.defaultLimit ?? DEFAULT_PAGE_SIZE; + const maxLimit = options.maxLimit ?? MAX_PAGE_SIZE; + const maxPage = options.maxPage ?? MAX_PAGE; + + const rawLimit = parseIntegerQueryValue(query.limit); + if (rawLimit !== null && rawLimit > maxLimit) { + throw new PaginationLimitError(rawLimit, maxLimit); + } + + const limit = + rawLimit !== null && rawLimit > 0 ? Math.min(rawLimit, maxLimit) : defaultLimit; + + const rawPage = parseIntegerQueryValue(query.page); + const page = + rawPage === null || rawPage < 1 ? 1 : Math.min(rawPage, maxPage); + + return { limit, page }; +} + +/** Thrown when a caller asks for more rows per page than the route allows. */ +export class PaginationLimitError extends Error { + readonly requested: number; + readonly maxLimit: number; + + constructor(requested: number, maxLimit: number) { + super( + `limit must not exceed ${maxLimit} (received ${requested}). ` + + 'Request at most ' + + `${maxLimit} items per page and use the page/cursor parameters to walk the rest of the collection.`, + ); + this.name = 'PaginationLimitError'; + this.requested = requested; + this.maxLimit = maxLimit; + } +} + +/** + * Express middleware that rejects an over-max `limit` with + * `400` / `code: 'LIMIT_EXCEEDED'` before any database work happens, and + * pins the sanitised `limit`/`page` onto `req.resolvedPagination` for the + * handler. + */ +export function enforcePaginationLimits(options: PaginationLimitOptions = {}) { + const defaultLimit = options.defaultLimit ?? DEFAULT_PAGE_SIZE; + const maxLimit = options.maxLimit ?? MAX_PAGE_SIZE; + const maxPage = options.maxPage ?? MAX_PAGE; + + return (req: Request, res: Response, next: NextFunction): void => { + try { + const resolved = resolvePagination(req.query as Record, { + defaultLimit, + maxLimit, + maxPage, + }); + + req.resolvedPagination = resolved; + next(); + } catch (err) { + if (err instanceof PaginationLimitError) { + sendApiError(req, res, { + status: 400, + code: LIMIT_EXCEEDED_CODE, + message: err.message, + retryable: false, + details: { + field: 'limit', + requested: err.requested, + maxLimit: err.maxLimit, + defaultLimit, + maxPage, + }, + }); + return; + } + next(err); + } + }; +} diff --git a/backend/src/middleware/validate.ts b/backend/src/middleware/validate.ts index 76e48ecf9..69a2a3b5e 100644 --- a/backend/src/middleware/validate.ts +++ b/backend/src/middleware/validate.ts @@ -248,6 +248,8 @@ export function validate(schemas: ValidateTargets) { status: 400, code: 'VALIDATION_ERROR', message: formatZodError(issues), + summary: 'Request validation failed', + errors: details, details, retryable: false, }); diff --git a/backend/src/pagination.ts b/backend/src/pagination.ts index da4d1b0ee..71f3e7865 100644 --- a/backend/src/pagination.ts +++ b/backend/src/pagination.ts @@ -76,6 +76,12 @@ export interface PaginationConfig { defaultLimit: number; /** Maximum allowed items per page. */ maxLimit: number; + /** + * Highest page number a caller may request (Issue #1430). The effective + * offset handed to the database is `(page - 1) * limit`, so an unbounded + * page number turns into an unbounded `skip` scan. + */ + maxPage: number; /** Whether to include total count in response. */ includeTotal: boolean; /** Default sort field. */ @@ -89,6 +95,7 @@ export interface PaginationConfig { export const DEFAULT_PAGINATION_CONFIG: PaginationConfig = { defaultLimit: 20, maxLimit: 100, + maxPage: 1000, includeTotal: true, defaultSortOrder: 'desc', }; @@ -123,11 +130,11 @@ export function parsePaginationQuery( query.cursor = req.query.cursor; } - // Parse page (1-based) + // Parse page (1-based, clamped to 1..maxPage — Issue #1430) if (req.query.page !== undefined) { const page = parseInt(req.query.page as string, 10); if (!isNaN(page) && page > 0) { - query.page = page; + query.page = Math.min(page, mergedConfig.maxPage); } else { query.page = 1; } diff --git a/backend/src/routes/vaults.ts b/backend/src/routes/vaults.ts new file mode 100644 index 000000000..fb9c376f4 --- /dev/null +++ b/backend/src/routes/vaults.ts @@ -0,0 +1,156 @@ +/** + * @file routes/vaults.ts + * `GET /api/v1/vaults` — paginated, public listing of active vaults. + * + * Pagination limits (Issue #1430) + * ------------------------------- + * This is an unauthenticated read that maps a caller-supplied `limit` + * straight onto Prisma's `take`, so it is exactly the shape of request that + * let `?limit=100000` OOM the API container. The route therefore sits behind + * `enforcePaginationLimits()`, which: + * + * - defaults `limit` to 20, + * - rejects `limit > 50` with `400` / `code: 'LIMIT_EXCEEDED'` *before* + * any database call (rejecting rather than clamping is deliberate: a + * silent clamp hides a broken or probing client behind a response that + * looks like a complete page), + * - clamps `page` into `1..1000` because the effective offset handed to + * the database is `(page - 1) * limit`. + * + * `take` is always `limit + 1` — one extra row is fetched to learn whether a + * next page exists, so the largest possible read is `MAX_PAGE_SIZE + 1`. + */ + +import { Router, Request, Response } from 'express'; +import { readsLimiter } from '../rateLimiter'; +import { createTimeoutFor } from '../middleware/timeoutMiddleware'; +import { + enforcePaginationLimits, + DEFAULT_PAGE_SIZE, + MAX_PAGE_SIZE, +} from '../middleware/paginationGuard'; +import { + createPaginatedResponse, + createPaginationEnvelope, + encodeCursor, + type PaginatedResponse, +} from '../pagination'; +import { getPrismaClient } from '../prismaClient'; +import { withSpan } from '../tracing'; +import { logger } from '../middleware/structuredLogging'; + +const router = Router(); + +/** Public projection of a vault row. `tenantId` is never exposed. */ +export interface VaultListItem { + id: string; + aum: number; + tvlUsd: string | null; + createdAt: string; + updatedAt: string; +} + +interface VaultRow { + id: string; + aum: number; + tvlUsd: string | null; + createdAt: Date; + updatedAt: Date; +} + +function toVaultListItem(row: VaultRow): VaultListItem { + return { + id: row.id, + aum: row.aum, + tvlUsd: row.tvlUsd, + createdAt: row.createdAt.toISOString(), + updatedAt: row.updatedAt.toISOString(), + }; +} + +/** + * Reads one page of active vaults. + * + * Exported separately from the route so unit tests can exercise the query + * construction without an HTTP round trip. + */ +export async function buildVaultsResponse( + limit: number, + page: number +): Promise> { + const prisma = getPrismaClient(); + const take = Math.min(limit, MAX_PAGE_SIZE) + 1; + const skip = (Math.max(1, page) - 1) * limit; + + const where = { deletedAt: null }; + + const [total, rows] = await Promise.all([ + prisma.vault.count({ where }), + prisma.vault.findMany({ + where, + // `id` breaks ties so pages stay stable when `createdAt` collides. + orderBy: [{ createdAt: 'desc' }, { id: 'asc' }], + take, + skip, + }), + ]); + + const hasNextPage = rows.length > limit; + const pageRows = hasNextPage ? rows.slice(0, limit) : rows; + const data = pageRows.map((row) => toVaultListItem(row as unknown as VaultRow)); + + const pagination = createPaginationEnvelope({ + count: data.length, + limit, + total, + currentPage: page, + totalPages: Math.max(1, Math.ceil(total / limit)), + hasNextPage, + hasPrevPage: page > 1, + nextCursor: + hasNextPage && data.length > 0 ? encodeCursor(data[data.length - 1].id) : null, + }); + + return createPaginatedResponse(data, pagination); +} + +/** + * GET /api/v1/vaults + * + * Query parameters: + * - `limit`: items per page, 1..50 (default 20). `> 50` → 400 LIMIT_EXCEEDED. + * - `page`: 1-based page number, clamped to 1..1000. + */ +router.get( + '/', + readsLimiter, + enforcePaginationLimits({ defaultLimit: DEFAULT_PAGE_SIZE, maxLimit: MAX_PAGE_SIZE }), + createTimeoutFor.read(), + async (req: Request, res: Response) => { + const { limit, page } = req.resolvedPagination ?? { + limit: DEFAULT_PAGE_SIZE, + page: 1, + }; + + return withSpan('vaults.list', async (span) => { + span.setAttributes({ 'vaults.limit': limit, 'vaults.page': page }); + + try { + const response = await buildVaultsResponse(limit, page); + res.status(200).json(response); + } catch (error) { + logger.log('error', 'Failed to list vaults', { + error: error instanceof Error ? error.message : String(error), + }); + res.status(500).json({ + error: 'Internal Server Error', + status: 500, + code: 'VAULTS_LIST_FAILED', + message: 'Failed to fetch vaults', + }); + } + }); + }, +); + +export default router; diff --git a/backend/src/schemaSnapshot.ts b/backend/src/schemaSnapshot.ts index fa2fca475..19a5f8882 100644 --- a/backend/src/schemaSnapshot.ts +++ b/backend/src/schemaSnapshot.ts @@ -74,7 +74,7 @@ export function extractSchemaFromZod(zodSchema: z.ZodType): SchemaDefinitio // Handle object schemas if (zodSchema instanceof z.ZodObject) { - const shape = (zodSchema as z.ZodObject)._shape; + const shape = (zodSchema as z.ZodObject).shape; const properties: Record = {}; const required: string[] = []; diff --git a/backend/src/sessionAudit.ts b/backend/src/sessionAudit.ts index 295bb6cdc..ef94cc50a 100644 --- a/backend/src/sessionAudit.ts +++ b/backend/src/sessionAudit.ts @@ -59,7 +59,7 @@ export async function recordSessionEvent(entry: SessionAuditEntry): Promise { it('should reject invalid Idempotency-Key format', () => { const req = createMockRequest({ - get: (header: string) => (header === 'Idempotency-Key' ? 'invalid' : undefined), + get: ((header: string) => + header === 'Idempotency-Key' ? 'invalid' : undefined) as Request['get'], }); const res = createMockResponse(); const next = jest.fn(); @@ -187,7 +188,8 @@ describe('Idempotency', () => { it('should accept valid UUID format', () => { const validUuid = '550e8400-e29b-41d4-a716-446655440000'; const req = createMockRequest({ - get: (header: string) => (header === 'Idempotency-Key' ? validUuid : undefined), + get: ((header: string) => + header === 'Idempotency-Key' ? validUuid : undefined) as Request['get'], }); const res = createMockResponse(); const next = jest.fn(); @@ -202,7 +204,8 @@ describe('Idempotency', () => { it('should accept valid hex nonce', () => { const validHex = 'a'.repeat(32); const req = createMockRequest({ - get: (header: string) => (header === 'Idempotency-Key' ? validHex : undefined), + get: ((header: string) => + header === 'Idempotency-Key' ? validHex : undefined) as Request['get'], }); const res = createMockResponse(); const next = jest.fn(); diff --git a/backend/src/tests/tenantBoundary.test.ts b/backend/src/tests/tenantBoundary.test.ts index c7477b131..2fd8a7964 100644 --- a/backend/src/tests/tenantBoundary.test.ts +++ b/backend/src/tests/tenantBoundary.test.ts @@ -21,10 +21,10 @@ function createMockRequest(overrides?: Partial): Partial { authApiKeyRole: 'viewer', authApiKeyHash: 'hash123', ip: '127.0.0.1', - get: (header: string) => { + get: ((header: string) => { if (header === 'user-agent') return 'test-agent'; return undefined; - }, + }) as Request['get'], ...overrides, }; } diff --git a/backend/src/tokenRevocation.ts b/backend/src/tokenRevocation.ts index 86cd13330..491a7cea1 100644 --- a/backend/src/tokenRevocation.ts +++ b/backend/src/tokenRevocation.ts @@ -9,6 +9,15 @@ * - ROTATION: Token was rotated during refresh * - SUSPICIOUS: Suspicious activity detected * - COMPROMISED: Token was compromised + * + * Issue #1431: this store is consulted on **every** authenticated request, not + * only on refresh. Two kinds of revocation are therefore tracked: + * + * - per token id (`jti`) — the exact access token that was logged out, and + * - per wallet, as a high-water mark (`revokedBefore`) — every access token + * issued at or before that instant, which is what `/auth/logout-all` and a + * "this token is compromised" event need. A timestamp rather than an + * explicit id list keeps the store O(1) per wallet instead of O(tokens). */ import Redis from 'ioredis'; @@ -24,18 +33,45 @@ export interface RevocationRecord { expiresAt: number; // When to remove from store } +export interface WalletRevocation { + /** Access tokens issued at or before this instant (unix ms) are rejected. */ + revokedBefore: number; + reason: RevocationReason; + /** When the wallet marker itself expires from the store. */ + expiresAt: number; +} + export interface RevocationStore { revoke(record: RevocationRecord): Promise; isRevoked(tokenId: string): Promise; revokeAllForWallet(walletAddress: string, reason: RevocationReason): Promise; clear(): Promise; + /** + * Reject every access token for `walletAddress` issued at or before + * `issuedAtMs` (the token's `iat` claim, in unix milliseconds). + */ + revokeWalletBefore( + walletAddress: string, + issuedAtMs: number, + reason: RevocationReason, + ): Promise; + /** True when the wallet has been revoked at or after `issuedAtMs`. */ + isWalletRevokedBefore(walletAddress: string, issuedAtMs: number): Promise; } +/** + * How long a wallet-wide revocation marker is retained. No live access token + * can predate the marker once the longest access-token TTL has elapsed, so a + * marker older than that can never match and is safe to drop. + */ +const WALLET_REVOCATION_TTL_MS = 24 * 60 * 60 * 1000; + /** * In-memory revocation store for single-instance deployments. */ export class InMemoryRevocationStore implements RevocationStore { private revoked = new Map(); + private walletRevocations = new Map(); async revoke(record: RevocationRecord): Promise { this.revoked.set(record.tokenId, record); @@ -57,6 +93,9 @@ export class InMemoryRevocationStore implements RevocationStore { } async revokeAllForWallet(walletAddress: string, reason: RevocationReason): Promise { + // Drop the per-token entries for this wallet: the wallet marker written + // below now decides. Previously this method only deleted records without + // recording anything, which silently un-revoked every one of them. let count = 0; for (const [tokenId, record] of this.revoked.entries()) { if (record.walletAddress === walletAddress) { @@ -64,11 +103,42 @@ export class InMemoryRevocationStore implements RevocationStore { count++; } } + await this.revokeWalletBefore(walletAddress, Date.now(), reason); return count; } + async revokeWalletBefore( + walletAddress: string, + issuedAtMs: number, + reason: RevocationReason, + ): Promise { + const existing = this.walletRevocations.get(walletAddress); + if (existing && existing.revokedBefore >= issuedAtMs) { + // A broader (later) revocation always wins; never move the marker back. + return; + } + this.walletRevocations.set(walletAddress, { + revokedBefore: issuedAtMs, + reason, + expiresAt: issuedAtMs + WALLET_REVOCATION_TTL_MS, + }); + } + + async isWalletRevokedBefore(walletAddress: string, issuedAtMs: number): Promise { + const marker = this.walletRevocations.get(walletAddress); + if (!marker) return false; + + if (marker.expiresAt < Date.now()) { + this.walletRevocations.delete(walletAddress); + return false; + } + + return issuedAtMs <= marker.revokedBefore; + } + async clear(): Promise { this.revoked.clear(); + this.walletRevocations.clear(); } private cleanup(): void { @@ -78,6 +148,11 @@ export class InMemoryRevocationStore implements RevocationStore { this.revoked.delete(tokenId); } } + for (const [walletAddress, marker] of this.walletRevocations.entries()) { + if (marker.expiresAt < now) { + this.walletRevocations.delete(walletAddress); + } + } } } @@ -87,6 +162,8 @@ export class InMemoryRevocationStore implements RevocationStore { * Key schema: * - `revocation:token:{tokenId}` → JSON revocation record (with TTL = expiresAt) * - `revocation:wallet:{walletAddress}` → set of revoked token IDs + * - `revocation:wallet-revoked-before:{walletAddress}` → JSON + * {@link WalletRevocation} high-water mark (Issue #1431) */ export class RedisRevocationStore implements RevocationStore { private readonly keyPrefix = 'revocation:'; @@ -145,19 +222,22 @@ export class RedisRevocationStore implements RevocationStore { try { const walletKey = `${this.keyPrefix}wallet:${walletAddress}`; const tokenIds = await this.redis.smembers(walletKey); - - if (tokenIds.length === 0) return 0; - // Remove each token - const pipeline = this.redis.pipeline(); - for (const tokenId of tokenIds) { - const key = `${this.keyPrefix}token:${tokenId}`; - pipeline.del(key); + // Remove each token: the wallet high-water mark written below now + // decides, so keeping per-token entries would only add read cost. + if (tokenIds.length > 0) { + const pipeline = this.redis.pipeline(); + for (const tokenId of tokenIds) { + const key = `${this.keyPrefix}token:${tokenId}`; + pipeline.del(key); + } + pipeline.del(walletKey); + + await pipeline.exec(); } - pipeline.del(walletKey); - - await pipeline.exec(); - + + await this.revokeWalletBefore(walletAddress, Date.now(), reason); + logger.info('All tokens revoked for wallet', { walletAddress, reason, @@ -170,7 +250,59 @@ export class RedisRevocationStore implements RevocationStore { walletAddress, error: err instanceof Error ? err.message : String(err), }); - return 0; + return this.fallback.revokeAllForWallet(walletAddress, reason); + } + } + + async revokeWalletBefore( + walletAddress: string, + issuedAtMs: number, + reason: RevocationReason, + ): Promise { + try { + const key = `${this.keyPrefix}wallet-revoked-before:${walletAddress}`; + const existingRaw = await this.redis.get(key); + + if (existingRaw) { + const existing = JSON.parse(existingRaw) as WalletRevocation; + if (existing.revokedBefore >= issuedAtMs) { + // A broader (later) revocation always wins; never move it back. + return; + } + } + + const marker: WalletRevocation = { + revokedBefore: issuedAtMs, + reason, + expiresAt: issuedAtMs + WALLET_REVOCATION_TTL_MS, + }; + const ttl = Math.max(1, Math.ceil((marker.expiresAt - Date.now()) / 1000)); + await this.redis.set(key, JSON.stringify(marker), 'EX', ttl); + + logger.debug('Wallet access tokens revoked', { walletAddress, reason, issuedAtMs }); + } catch (err) { + logger.error('Failed to revoke wallet access tokens', { + walletAddress, + error: err instanceof Error ? err.message : String(err), + }); + await this.fallback.revokeWalletBefore(walletAddress, issuedAtMs, reason); + } + } + + async isWalletRevokedBefore(walletAddress: string, issuedAtMs: number): Promise { + try { + const key = `${this.keyPrefix}wallet-revoked-before:${walletAddress}`; + const raw = await this.redis.get(key); + if (!raw) return false; + + const marker = JSON.parse(raw) as WalletRevocation; + return issuedAtMs <= marker.revokedBefore; + } catch (err) { + logger.error('Failed to check wallet revocation', { + walletAddress, + error: err instanceof Error ? err.message : String(err), + }); + return this.fallback.isWalletRevokedBefore(walletAddress, issuedAtMs); } } @@ -185,6 +317,7 @@ export class RedisRevocationStore implements RevocationStore { logger.error('Failed to clear revocation store', { error: err instanceof Error ? err.message : String(err), }); + await this.fallback.clear(); } } } diff --git a/backend/src/tracing.ts b/backend/src/tracing.ts index a7d784c66..fc29dc015 100644 --- a/backend/src/tracing.ts +++ b/backend/src/tracing.ts @@ -8,9 +8,9 @@ * OTEL_ENABLED - Set to "false" to disable (default: true) */ -import { NodeSDK } from '@opentelemetry/sdk-node'; +import { NodeSDK, type NodeSDKConfiguration } from '@opentelemetry/sdk-node'; import { OTLPTraceExporter } from '@opentelemetry/exporter-trace-otlp-http'; -import { resourceFromAttributes } from '@opentelemetry/resources'; +import * as otelResources from '@opentelemetry/resources'; import { ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION } from '@opentelemetry/semantic-conventions'; import { HttpInstrumentation } from '@opentelemetry/instrumentation-http'; import { ExpressInstrumentation } from '@opentelemetry/instrumentation-express'; @@ -28,6 +28,26 @@ const SERVICE_NAME = process.env.OTEL_SERVICE_NAME || 'yieldvault-backend'; const OTLP_ENDPOINT = process.env.OTEL_EXPORTER_OTLP_ENDPOINT || 'http://localhost:4318'; +/** + * Builds the SDK resource from a flat attribute bag. + * + * `@opentelemetry/resources` renamed its factory between the v1 and v2 lines + * (`new Resource(attributes)` → `resourceFromAttributes(attributes)`), and the + * two OpenTelemetry major lines coexist in this workspace's lockfile. Resolve + * whichever factory the installed version exposes so the SDK boots on both. + */ +function buildResource(attributes: Record): NodeSDKConfiguration['resource'] { + const resources = otelResources as unknown as { + resourceFromAttributes?: (attrs: Record) => unknown; + Resource: new (attrs: Record) => unknown; + }; + + if (typeof resources.resourceFromAttributes === 'function') { + return resources.resourceFromAttributes(attributes) as NodeSDKConfiguration['resource']; + } + return new resources.Resource(attributes) as NodeSDKConfiguration['resource']; +} + let sdk: NodeSDK | null = null; export function initTracing(): void { @@ -62,7 +82,7 @@ export function initTracing(): void { const exporter = new OTLPTraceExporter({ url: `${OTLP_ENDPOINT}/v1/traces` }); sdk = new NodeSDK({ - resource: resourceFromAttributes({ + resource: buildResource({ [ATTR_SERVICE_NAME]: SERVICE_NAME, [ATTR_SERVICE_VERSION]: process.env.npm_package_version || '1.0.0', }), diff --git a/backend/src/types/express.d.ts b/backend/src/types/express.d.ts index 1bdeca878..d710c84ec 100644 --- a/backend/src/types/express.d.ts +++ b/backend/src/types/express.d.ts @@ -1,10 +1,17 @@ declare global { + import type { ResolvedPagination } from '../middleware/paginationGuard'; + namespace Express { interface Request { authApiKeyHash?: string; authApiKeyRole?: string; apiVersion?: 'v1' | 'v2'; apiVersionSource?: 'path' | 'legacy' | 'default'; + /** + * Effective, clamped pagination parameters for the current request. + * Populated by `enforcePaginationLimits()`. + */ + resolvedPagination?: ResolvedPagination; } } } diff --git a/backend/src/types/validation.ts b/backend/src/types/validation.ts index 20041be10..88180a744 100644 --- a/backend/src/types/validation.ts +++ b/backend/src/types/validation.ts @@ -26,7 +26,10 @@ export const PaginationQuerySchema = z .object({ limit: z.string().regex(/^\d+$/, 'limit must be a positive integer').optional(), cursor: z.string().optional(), - page: z.string().regex(/^\d+$/, 'page must be a positive integer').optional(), + // Signed so an out-of-range page (`page=-1`, `page=1000000`) reaches + // `parsePaginationQuery`, which clamps it into `1..maxPage`, instead of + // being rejected outright (Issue #1430). + page: z.string().regex(/^-?\d+$/, 'page must be an integer').optional(), sortBy: z.string().optional(), sortOrder: z.string().optional(), dryRun: z.enum(['true', 'false', '1', '0']).optional(), diff --git a/backend/src/vaultEndpoints.ts b/backend/src/vaultEndpoints.ts index abbbdf8cb..0b58a694f 100644 --- a/backend/src/vaultEndpoints.ts +++ b/backend/src/vaultEndpoints.ts @@ -3,7 +3,7 @@ import { emailService } from './emailService'; import { logger } from './middleware/structuredLogging'; import { allowlistMiddleware } from './middleware/allowlist'; import { triggerCacheInvalidation, registerInvalidationHook } from './middleware/cache'; -import { depositsLimiter, depositsUserLimiter } from './rateLimiter'; +import { depositsLimiter, depositsUserLimiter, readsLimiter } from './rateLimiter'; import { cacheMiddleware } from './middleware/cache'; import { idempotencyStore, @@ -769,7 +769,7 @@ router.post('/strategy', depositsLimiter, requireFlag('strategy-selection'), val }); }); - res.status(200).json({ message: 'Strategy selection endpoint (v2 preview)' }); + return res.status(200).json({ message: 'Strategy selection endpoint (v2 preview)' }); }); /** @@ -900,7 +900,7 @@ router.get('/receipts', readsLimiter, async (req: Request, res: Response) => { const transactions = await prisma.transaction.findMany({ where, - orderBy: { createdAt: 'desc' }, + orderBy: { timestamp: 'desc' }, take: limit + 1, ...(cursor ? { cursor: { id: cursor }, skip: 1 } : {}), }); @@ -916,7 +916,7 @@ router.get('/receipts', readsLimiter, async (req: Request, res: Response) => { status: tx.status, walletAddress: tx.user, explorerUrl: `${EXPLORER_BASE_URL}/${tx.id}`, - timestamp: tx.createdAt.toISOString(), + timestamp: tx.timestamp.toISOString(), })); res.status(200).json({