diff --git a/backend/prisma/schema.prisma b/backend/prisma/schema.prisma index 20e8c2e41..e19d3f945 100644 --- a/backend/prisma/schema.prisma +++ b/backend/prisma/schema.prisma @@ -675,6 +675,8 @@ model WalletTenantAssociation { @@index([walletAddress, tenantId]) @@index([tenantId]) +} + model IdempotencyKey { id String @id @default(uuid()) keyId String diff --git a/backend/src/__tests__/endpointResponseContract.test.ts b/backend/src/__tests__/endpointResponseContract.test.ts new file mode 100644 index 000000000..a6b7a952c --- /dev/null +++ b/backend/src/__tests__/endpointResponseContract.test.ts @@ -0,0 +1,387 @@ +/** + * Response-shape contract tests for endpoints that are also reachable through + * an in-process snapshot builder. + * + * `GET /admin/impersonate/:wallet` does not proxy the HTTP route — it + * re-synthesizes each sub-response (summary, transactions, holdings, history, + * referral stats, referral code) in-process so an admin sees exactly what the + * wallet owner sees. That makes every one of those builders a second + * implementation of a wire contract, and nothing but a test stops the two from + * drifting. + * + * These tests therefore assert each endpoint's shape on *both* paths: the real + * HTTP response, and the proxied snapshot, for the success and error cases + * alike. Any builder that hand-rolls a partial body fails here. + * + * The referral routes also sit behind a response cache, so every test clears it + * first: a 200 recorded by an earlier case would be replayed instead of + * re-running the handler and would quietly mask the failure-path cases. + * + * Issues: #1319 (in-process snapshots must mirror the HTTP error envelope), + * #1322 (validate direct and proxied shapes for the same endpoint), + * #1318 (invalid pagination must be clamped, not rejected). + */ + +import request from 'supertest'; +import app from '../index'; +import { getPrismaClient, disconnectPrismaClient } from '../prismaClient'; +import { referralService } from '../referralService'; +import { normalizeWalletAddress } from '../walletUtils'; +import { clearAdminAuditLogsForTests } from '../adminAudit'; +import { registerApiKey } from '../middleware/apiKeyAuth'; +import { invalidateCache } from '../middleware/cache'; +import { THIRD_TEST_WALLET } from './setup'; + +const getPrisma = () => getPrismaClient(); + +const SUPER_ADMIN_KEY = 'contract-super-admin-key'; +const ADMIN_WALLET = 'GADMIN000000000000000000000000000000000000000000000009'; + +const REFERRER_WALLET = 'GCONTRACTREFERRER000000000000000000000000000000001'; +const REFERRED_WALLET = 'GCONTRACTREFERRED000000000000000000000000000000002'; +/** Wallet with no referral activity at all — exercises the 404 envelope. */ +const INACTIVE_WALLET = THIRD_TEST_WALLET; + +interface ImpersonationSnapshot { + walletAddress: string; + transactions: unknown; + portfolioHoldings: unknown; + vaultHistory: unknown; + referralStats: { statusCode: number; body: unknown }; + referralCode: { statusCode: number; body: unknown }; +} + +/** Assert a body carries the full canonical error envelope, not a partial one. */ +function expectCanonicalErrorEnvelope(body: unknown, status: number, context: string): void { + const envelope = (body ?? {}) as Record; + const required = ['error', 'status', 'code', 'message', 'retryable'] as const; + const missing = required.filter( + (field) => !Object.prototype.hasOwnProperty.call(envelope, field) + ); + + // Report the missing fields together with the call site so a drifting + // snapshot builder points straight at the field it dropped. + expect({ context, missing }).toEqual({ context, missing: [] }); + expect(envelope.status).toBe(status); + expect(typeof envelope.error).toBe('string'); + expect(typeof envelope.code).toBe('string'); + expect((envelope.code as string).length).toBeGreaterThan(0); + expect(typeof envelope.message).toBe('string'); + expect(typeof envelope.retryable).toBe('boolean'); +} + +const LIST_ENVELOPE_KEYS = ['data', 'pagination', 'timestamp'] as const; +const PAGINATION_META_KEYS = [ + 'count', + 'limit', + 'total', + 'nextCursor', + 'prevCursor', + 'currentPage', + 'totalPages', + 'hasNextPage', + 'hasPrevPage', +] as const; + +/** + * Assert a list body is a well-formed paginated envelope. + * + * The success counterpart of `expectCanonicalErrorEnvelope`: a snapshot builder + * that emits a partial list body still compares equal to the direct response + * only if both are broken, so each key set is asserted on its own merits too. + */ +function expectCanonicalListEnvelope(body: unknown, context: string): void { + const envelope = (body ?? {}) as Record; + const missing = LIST_ENVELOPE_KEYS.filter( + (field) => !Object.prototype.hasOwnProperty.call(envelope, field) + ); + + expect({ context, missing }).toEqual({ context, missing: [] }); + expect(Array.isArray(envelope.data)).toBe(true); + + const pagination = (envelope.pagination ?? {}) as Record; + const missingPagination = PAGINATION_META_KEYS.filter( + (field) => !Object.prototype.hasOwnProperty.call(pagination, field) + ); + + expect({ context: `${context}.pagination`, missing: missingPagination }).toEqual({ + context: `${context}.pagination`, + missing: [], + }); + expect(typeof pagination.count).toBe('number'); + expect(typeof pagination.limit).toBe('number'); + expect(typeof pagination.hasNextPage).toBe('boolean'); + expect(typeof pagination.hasPrevPage).toBe('boolean'); +} + +async function impersonate(wallet: string): Promise { + const response = await request(app) + .get(`/admin/impersonate/${wallet}`) + .set('Authorization', `ApiKey ${SUPER_ADMIN_KEY}`) + .set('x-admin-id', ADMIN_WALLET); + + expect(response.status).toBe(200); + return response.body as ImpersonationSnapshot; +} + +describe('Response shape contract: endpoints re-synthesized by /admin/impersonate', () => { + beforeAll(async () => { + const prisma = getPrisma(); + await prisma.referral.deleteMany(); + await prisma.referralCode.deleteMany(); + await prisma.sharePriceSnapshot.deleteMany(); + await prisma.transaction.deleteMany(); + + // Active referrer: a referral row with firstDepositAt set is what makes + // getReferralStats() return stats instead of null. + await referralService.createReferralCode(REFERRER_WALLET, 'CONTRACT1'); + await prisma.referral.create({ + data: { + referrerAddress: normalizeWalletAddress(REFERRER_WALLET), + referredAddress: normalizeWalletAddress(REFERRED_WALLET), + firstDepositAt: new Date(), + }, + }); + }); + + afterAll(async () => { + const prisma = getPrisma(); + await prisma.referral.deleteMany(); + await prisma.referralCode.deleteMany(); + await disconnectPrismaClient(); + }); + + beforeEach(() => { + clearAdminAuditLogsForTests(); + // The referral routes sit behind a 30s response cache, so a 200 recorded by + // an earlier test would be replayed instead of re-running the handler — + // which would silently hide the failure-path cases below. + invalidateCache(); + registerApiKey(SUPER_ADMIN_KEY, { role: 'super-admin' }); + }); + + afterEach(() => { + jest.restoreAllMocks(); + }); + + describe('error path (404)', () => { + it('returns the canonical error envelope over HTTP', async () => { + const direct = await request(app).get(`/api/v1/referrals/${INACTIVE_WALLET}`); + + expect(direct.status).toBe(404); + expectCanonicalErrorEnvelope(direct.body, 404, 'direct HTTP 404'); + }); + + it('returns the same envelope through the impersonation proxy', async () => { + const snapshot = await impersonate(INACTIVE_WALLET); + + expect(snapshot.referralStats.statusCode).toBe(404); + expectCanonicalErrorEnvelope(snapshot.referralStats.body, 404, 'proxied 404'); + }); + + it('matches the direct response byte-for-byte on both paths (#1319, #1322)', async () => { + const direct = await request(app).get(`/api/v1/referrals/${INACTIVE_WALLET}`); + const snapshot = await impersonate(INACTIVE_WALLET); + + expect(snapshot.referralStats).toEqual({ + statusCode: direct.status, + body: direct.body, + }); + }); + }); + + describe('success path (200)', () => { + it('returns the same stats body over HTTP and through the proxy', async () => { + const direct = await request(app).get(`/api/v1/referrals/${REFERRER_WALLET}`); + + expect(direct.status).toBe(200); + expect(direct.body).toMatchObject({ + referral_count: expect.any(Number), + total_reward_earned: expect.any(String), + }); + + const snapshot = await impersonate(REFERRER_WALLET); + + expect(snapshot.referralStats).toEqual({ + statusCode: direct.status, + body: direct.body, + }); + }); + }); + + describe('sibling endpoint: GET /api/v1/referrals/code/:wallet', () => { + it('matches the direct response through the proxy', async () => { + const direct = await request(app).get(`/api/v1/referrals/code/${REFERRER_WALLET}`); + + expect(direct.status).toBe(200); + + const snapshot = await impersonate(REFERRER_WALLET); + + expect(snapshot.referralCode).toEqual({ + statusCode: direct.status, + body: direct.body, + }); + }); + }); + + describe('failure path (500) — the snapshot must not reject (#1319)', () => { + it('mirrors the stats 500 envelope on both paths and still returns 200 overall', async () => { + jest + .spyOn(referralService, 'getReferralStats') + .mockRejectedValue(new Error('referral store unavailable')); + + const direct = await request(app).get(`/api/v1/referrals/${REFERRER_WALLET}`); + + expect(direct.status).toBe(500); + expectCanonicalErrorEnvelope(direct.body, 500, 'direct HTTP 500'); + + const snapshot = await impersonate(REFERRER_WALLET); + + // The proxy answers 200 with a per-sub-resource status, so a thrown + // rejection here would surface as a 500 on the whole impersonation + // response — and take every other synthesized sub-resource with it. + expect(snapshot.referralStats).toEqual({ + statusCode: direct.status, + body: direct.body, + }); + expectCanonicalListEnvelope(snapshot.transactions, 'proxied transactions'); + }); + + it('mirrors the referral-code 500 envelope on both paths', async () => { + jest + .spyOn(referralService, 'getOrCreateReferralCode') + .mockRejectedValue(new Error('referral code store unavailable')); + + const direct = await request(app).get(`/api/v1/referrals/code/${REFERRER_WALLET}`); + + expect(direct.status).toBe(500); + expectCanonicalErrorEnvelope(direct.body, 500, 'direct HTTP 500 (code)'); + + const snapshot = await impersonate(REFERRER_WALLET); + + expect(snapshot.referralCode).toEqual({ + statusCode: direct.status, + body: direct.body, + }); + }); + }); + + describe('every synthesized impersonation sub-resource (#1322)', () => { + /** The list endpoints `/admin/impersonate/:wallet` re-synthesizes. */ + const PROXIED_LIST_ENDPOINTS = [ + { + name: 'transactions', + path: (wallet: string) => `/api/v1/transactions?walletAddress=${wallet}`, + snapshotKey: 'transactions' as const, + }, + { + name: 'portfolio holdings', + path: (wallet: string) => `/api/v1/portfolio/holdings?walletAddress=${wallet}`, + snapshotKey: 'portfolioHoldings' as const, + }, + { + name: 'vault history', + path: (_wallet: string) => '/api/v1/vault/history', + snapshotKey: 'vaultHistory' as const, + }, + ]; + + it.each(PROXIED_LIST_ENDPOINTS)( + 'GET $name is a well-formed list envelope on both the direct and the proxied path', + async ({ path, snapshotKey }) => { + const direct = await request(app).get(path(INACTIVE_WALLET)); + + expect(direct.status).toBe(200); + expectCanonicalListEnvelope(direct.body, `direct ${snapshotKey}`); + + const snapshot = await impersonate(INACTIVE_WALLET); + + expectCanonicalListEnvelope(snapshot[snapshotKey], `proxied ${snapshotKey}`); + } + ); + + it.each(PROXIED_LIST_ENDPOINTS)( + 'the proxied $name payload matches what the endpoint serves', + async ({ path, snapshotKey }) => { + const direct = await request(app).get(path(INACTIVE_WALLET)); + const snapshot = await impersonate(INACTIVE_WALLET); + const proxied = snapshot[snapshotKey] as { data: unknown; pagination: unknown }; + + // Only data and the pagination metadata are compared: every envelope + // carries a per-call `timestamp`, so whole-body equality would fail for + // reasons that have nothing to do with the wire contract. + expect(proxied.data).toEqual(direct.body.data); + expect(proxied.pagination).toEqual(direct.body.pagination); + } + ); + }); +}); + +describe('Pagination contract: invalid page params are clamped, never rejected (#1318)', () => { + /** Every one of these must resolve to the first page, not a 400. */ + const INVALID_PAGE_VALUES = ['-1', '0', 'abc', '1.5', '-999']; + + it.each(INVALID_PAGE_VALUES)('resolves page=%s to the first page', async (page) => { + const response = await request(app).get(`/api/v1/transactions?page=${page}`); + + expect(response.status).toBe(200); + expect(response.body.pagination.currentPage).toBe(1); + }); + + it('resolves an invalid limit to the endpoint default', async () => { + const response = await request(app).get('/api/v1/transactions?limit=abc'); + + expect(response.status).toBe(200); + expect(response.body.pagination.limit).toBeGreaterThan(0); + }); + + it('clamps a limit above the endpoint maximum instead of rejecting it', async () => { + const response = await request(app).get('/api/v1/transactions?limit=100000'); + + expect(response.status).toBe(200); + expect(response.body.pagination.limit).toBeLessThanOrEqual(100); + }); + + it('still accepts a valid page unchanged', async () => { + const response = await request(app).get('/api/v1/transactions?page=1&limit=5'); + + expect(response.status).toBe(200); + expect(response.body.pagination.currentPage).toBe(1); + }); + + /** + * The other list endpoints share the same contract, so they are held to it + * here too: each one resolves an unusable `limit` to a valid page size + * instead of letting a NaN reach the data layer as `take` and 500ing. + */ + it.each([ + ['/api/v1/transactions', 'pagination.limit'], + ['/api/v1/portfolio/holdings', 'pagination.limit'], + ['/api/v1/vault/history', 'pagination.limit'], + ])('GET %s clamps limit=abc to a usable page size', async (path, limitField) => { + const response = await request(app).get(`${path}?limit=abc`); + + expect(response.status).toBe(200); + + const value = limitField + .split('.') + .reduce((node, key) => (node as Record)[key], response.body); + + expect(typeof value).toBe('number'); + expect(value as number).toBeGreaterThan(0); + }); + + it('GET /api/v1/vault/receipts clamps limit=abc instead of 500ing', async () => { + const response = await request(app).get('/api/v1/vault/receipts?limit=abc'); + + expect(response.status).toBe(200); + expect(Array.isArray(response.body.receipts)).toBe(true); + }); + + it('GET /api/v1/vault/receipts clamps a limit above the endpoint maximum', async () => { + const response = await request(app).get('/api/v1/vault/receipts?limit=100000'); + + expect(response.status).toBe(200); + expect(response.body.receipts.length).toBeLessThanOrEqual(100); + }); +}); diff --git a/backend/src/__tests__/governance.test.ts b/backend/src/__tests__/governance.test.ts index 1d67fe75a..1bddbf349 100644 --- a/backend/src/__tests__/governance.test.ts +++ b/backend/src/__tests__/governance.test.ts @@ -182,6 +182,23 @@ describe('Backend governance', () => { body: referralCode.body, }); + // The proxied referral stats are synthesized in-process, not proxied over + // HTTP, so assert the envelope explicitly rather than only via equality — + // a partial body must fail on its own merits (#1319, #1322). + for (const key of ['error', 'status', 'code', 'message', 'retryable'] as const) { + expect(Object.prototype.hasOwnProperty.call(response.body.referralStats.body, key)).toBe( + true + ); + } + expect(response.body.referralStats.statusCode).toBe(404); + expect(response.body.referralStats.body).toMatchObject({ + error: 'Not Found', + status: 404, + code: expect.any(String), + message: 'No referral activity found for this wallet', + retryable: false, + }); + const auditResponse = await request(app) .get('/admin/audit/logs') .query({ action: 'admin.impersonate', limit: 5 }) diff --git a/backend/src/__tests__/issues711.test.ts b/backend/src/__tests__/issues711.test.ts index 7feef6c6b..f2fb0bcc4 100644 --- a/backend/src/__tests__/issues711.test.ts +++ b/backend/src/__tests__/issues711.test.ts @@ -47,21 +47,19 @@ 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; + // diffSchemaShapes(baseline, current) compares a committed snapshot against + // the live schema, so the shape that is missing the field is the baseline + // and the shape carrying it is `current`. + const current = zodToJsonShape(HealthResponseSchema); + const baseline = JSON.parse(JSON.stringify(current)) as typeof current; // 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( + delete baseline.properties?.checks?.properties?.indexer; + baseline.properties!.checks!.required = (baseline.properties!.checks!.required ?? []).filter( (k: string) => k !== 'indexer', ); const issues = diffSchemaShapes(baseline, current, 'GET /health'); - expect( - issues.some( - (issue) => - issue.message.includes('new field added') || issue.message.includes('now required'), - ), - ).toBe(true); + expect(issues.some((issue) => issue.message.includes('now required'))).toBe(true); }); it('validates a conforming health payload', () => { diff --git a/backend/src/__tests__/requestValidation.test.ts b/backend/src/__tests__/requestValidation.test.ts index 39e28194d..cca89507e 100644 --- a/backend/src/__tests__/requestValidation.test.ts +++ b/backend/src/__tests__/requestValidation.test.ts @@ -3,10 +3,12 @@ import { ApyBackfillBodySchema, MaintenanceToggleSchema, PaginationQuerySchema, + TransactionListQuerySchema, WebhookVerifyBodySchema, DeadLetterIdsSchema, } from '../types/validation'; import { WEBHOOK_EVENT_TYPES } from '../types/webhooks'; +import { clampLimitNumber, clampPageNumber } from '../pagination'; describe('request validation schemas', () => { it('rejects an inverted APY backfill range', () => { @@ -19,9 +21,22 @@ describe('request validation schemas', () => { expect(result.enabled).toBe(true); }); - it('rejects a non-numeric pagination limit', () => { - const result = PaginationQuerySchema.safeParse({ limit: 'abc' }); - expect(result.success).toBe(false); + it('accepts an out-of-range pagination limit for the parser to clamp', () => { + // Contract tests require invalid pagination to resolve gracefully (200) via + // parsePaginationQuery's clamping, so the schema must not reject it first. + for (const limit of ['abc', '-1', '0', '100000']) { + expect(PaginationQuerySchema.safeParse({ limit }).success).toBe(true); + } + }); + + it('accepts an out-of-range page for the parser to clamp', () => { + for (const page of ['abc', '-1', '0', '1.5']) { + expect(PaginationQuerySchema.safeParse({ page }).success).toBe(true); + } + }); + + it('still rejects a non-string pagination value', () => { + expect(PaginationQuerySchema.safeParse({ page: ['1', '2'] }).success).toBe(false); }); it('requires a webhook verify secret', () => { diff --git a/backend/src/apiContractSnapshots.ts b/backend/src/apiContractSnapshots.ts index acaa02eb4..b0f2d37dd 100644 --- a/backend/src/apiContractSnapshots.ts +++ b/backend/src/apiContractSnapshots.ts @@ -249,7 +249,10 @@ export function diffSchemaShapes( } for (const key of baselineRequired) { - if (!(key in baseline.properties ?? {})) { + // `in` binds tighter than `??`, so the parentheses are required: without + // them this reads as `(key in baseline.properties) ?? {}` and throws a + // TypeError whenever a committed snapshot has no `properties` object. + if (!(key in (baseline.properties ?? {}))) { issues.push({ path: at(key), message: 'required field missing from snapshot properties (orphaned reference)' }); } if (!(key in currentProps)) { @@ -260,8 +263,9 @@ 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)' }); + } 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/idempotency.ts b/backend/src/idempotency.ts index 1ab4b0448..5f9a62b2f 100644 --- a/backend/src/idempotency.ts +++ b/backend/src/idempotency.ts @@ -22,7 +22,9 @@ */ import crypto from 'crypto'; +import NodeCache from 'node-cache'; import { prisma } from './prisma'; +import { redisClientManager } from './rateLimiter'; import { logger } from './middleware/structuredLogging'; import type { Request, Response, NextFunction } from 'express'; @@ -419,3 +421,346 @@ export function startIdempotencyCleanupTask(intervalMs = 3600000): NodeJS.Timer }); }, intervalMs); } + +// ─── In-flight store (Redis + NodeCache) ───────────────────────────────────── +// +// Issue #811: multi-instance deployments lost idempotency guarantees on pod +// recycle because responses were stored only in-process. This revision +// persists completed responses to Redis using SET … EX so all replicas share +// the same store, falling back to NodeCache when Redis is unavailable +// (fail-open, with a warning). +// +// The Prisma helpers above (enforceIdempotency / *IdempotencyRecord) are the +// newer, DB-backed design but are not yet wired into any route. The mutation +// endpoints (deposit/withdrawal/transfer), the admin key routes and the +// retention sweeper all still call the store below, so both surfaces live here +// until those call sites are migrated. Restoring this store is what keeps +// `idempotencyStore`, `IdempotencyStore` and `IdempotencyConflictError` +// resolvable for those consumers. + +export interface IdempotentOperationResult { + statusCode: number; + body: T; +} + +/** Metadata attached to every idempotency key entry. */ +export interface IdempotencyEntryMetadata { + /** 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: IdempotencyEntryMetadata; +} + +/** Snapshot of store-wide observability counters. */ +export interface IdempotencyMetrics { + hits: number; + conflicts: number; + evictions: number; + activeKeys: number; + pendingKeys: number; +} + +interface StoredResponse extends IdempotentOperationResult { + fingerprint: string; + metadata: IdempotencyEntryMetadata; +} + +interface PendingOperation { + fingerprint: string; + promise: Promise>; + metadata: IdempotencyEntryMetadata; +} + +export class IdempotencyConflictError extends Error { + constructor(message = 'Idempotency key already used for a different request body') { + super(message); + this.name = 'IdempotencyConflictError'; + } +} + +const REDIS_PREFIX = 'idempotency:'; + +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: IdempotencyEntryMetadata = { 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..a1246e3f7 100644 --- a/backend/src/index.ts +++ b/backend/src/index.ts @@ -111,10 +111,15 @@ import { buildTransactionExportArtifact, buildTransactionsResponse, buildVaultHistoryResponse, + TRANSACTION_PAGINATION_CONFIG, } from './listEndpoints'; import { createPaginatedResponse, createPaginationEnvelope, encodeCursor } from './pagination'; import listRouter from './listEndpoints'; -import referralRouter from './referralEndpoints'; +import referralRouter, { + buildReferralCodeErrorBody, + buildReferralStatsErrorBody, + buildReferralStatsNotFoundBody, +} from './referralEndpoints'; import auditLogRouter from './auditLogEndpoints'; import { referralService } from './referralService'; import { @@ -425,30 +430,67 @@ function sendStandardListEnvelope( async function buildReferralStatsSnapshot(wallet: string) { const normalizedWallet = normalizeWalletAddress(wallet); - const stats = await referralService.getReferralStats(normalizedWallet); - if (!stats) { + + try { + const stats = await referralService.getReferralStats(normalizedWallet); + + if (!stats) { + // Must be the exact envelope GET /api/v1/referrals/:wallet returns, not + // a hand-rolled stand-in — admin impersonation is compared against that + // endpoint field-for-field by the governance contract test. + return { + statusCode: 404, + body: buildReferralStatsNotFoundBody(), + }; + } + return { - statusCode: 404, - body: { - error: 'Not Found', - status: 404, - code: 'ROUTE_NOT_FOUND', - message: 'No referral activity found for this wallet', - retryable: false, - }, + statusCode: 200, + body: stats, + }; + } catch (error) { + // Same reasoning as the 404 above. The route reports a failed lookup as a + // 500 with the canonical envelope; letting the rejection escape instead + // would 500 the entire impersonation response with a different message, so + // one broken sub-resource would take the whole snapshot with it (#1319). + logger.log('error', 'Error building referral stats snapshot', { + error: error instanceof Error ? error.message : String(error), + wallet: normalizedWallet, + }); + return { + statusCode: 500, + body: buildReferralStatsErrorBody(), }; } +} - return { - statusCode: 200, - body: stats, - }; +async function buildReferralCodeSnapshot(wallet: string) { + const normalizedWallet = normalizeWalletAddress(wallet); + + try { + const code = await referralService.getOrCreateReferralCode(normalizedWallet); + return { + statusCode: 200, + body: { code }, + }; + } catch (error) { + logger.log('error', 'Error building referral code snapshot', { + error: error instanceof Error ? error.message : String(error), + wallet: normalizedWallet, + }); + return { + statusCode: 500, + body: buildReferralCodeErrorBody(), + }; + } } async function buildWalletTransactionsSnapshot(wallet: string) { const normalizedWallet = normalizeWalletAddress(wallet); const prisma = getPrismaClient(); - const limit = 20; + // Shared with GET /api/v1/transactions so the proxy can never drift to a + // different page size than the endpoint it stands in for (#1322). + const limit = TRANSACTION_PAGINATION_CONFIG.defaultLimit ?? 20; const where = { user: normalizedWallet }; const [total, transactions] = await Promise.all([ prisma.transaction.count({ where }), @@ -486,10 +528,7 @@ async function buildImpersonatedVaultState(wallet: string) { portfolioHoldings: buildPortfolioHoldingsResponse({ walletAddress: normalizedWallet }), vaultHistory: buildVaultHistoryResponse({}), referralStats: await buildReferralStatsSnapshot(normalizedWallet), - referralCode: { - statusCode: 200, - body: { code: await referralService.getOrCreateReferralCode(normalizedWallet) }, - }, + referralCode: await buildReferralCodeSnapshot(normalizedWallet), }; } @@ -5404,6 +5443,10 @@ app.use((req: Request, res: Response) => { status: 404, code: 'ROUTE_NOT_FOUND', message: `Cannot ${req.method} ${req.originalUrl}`, + // `path` is also surfaced at the top level because that is where the + // documented not-found envelope puts it; `details.path` stays for callers + // that already read the field from there. + path: req.originalUrl, details: { path: req.originalUrl }, retryable: false, }); diff --git a/backend/src/listEndpoints.ts b/backend/src/listEndpoints.ts index 5d255bd67..d263765c4 100644 --- a/backend/src/listEndpoints.ts +++ b/backend/src/listEndpoints.ts @@ -13,6 +13,8 @@ import { Router, Request, Response } from 'express'; import { Readable } from 'stream'; import { parsePaginationQuery, + clampLimitNumber, + clampPageNumber, paginateWithCursor, paginateWithOffset, sortItems, @@ -196,7 +198,7 @@ const MOCK_VAULT_HISTORY: VaultHistoryPoint[] = Array.from({ length: 365 }, (_, // ─── Pagination Configs ───────────────────────────────────────────────────── -const TRANSACTION_PAGINATION_CONFIG: Partial = { +export const TRANSACTION_PAGINATION_CONFIG: Partial = { defaultLimit: 20, maxLimit: 100, defaultSortBy: 'timestamp', @@ -457,7 +459,12 @@ export async function buildTransactionsResponse( query: WalletStateQuery ): Promise> { const prisma = getPrismaClient(); - const limit = query.limit ?? TRANSACTION_PAGINATION_CONFIG.defaultLimit ?? 20; + // Clamp rather than trust the caller: this builder is reachable from the + // routes, from the legacy /api/transactions handler and from the in-process + // impersonation snapshot, and a raw parseInt from any of them must not be + // able to put a NaN, zero or oversized limit in the envelope. + const limit = clampLimitNumber(query.limit, TRANSACTION_PAGINATION_CONFIG); + const page = clampPageNumber(query.page); const normalizedDateRange = parseDateRangeOrThrow({ from: query.from, to: query.to }); const sortBy = query.sortBy ?? TRANSACTION_PAGINATION_CONFIG.defaultSortBy ?? 'timestamp'; @@ -485,8 +492,8 @@ export async function buildTransactionsResponse( take: limit + 1, ...(query.cursor ? { cursor: { id: Buffer.from(query.cursor, 'base64url').toString('utf-8') }, skip: 1 } - : query.page && query.page > 1 - ? { skip: (query.page - 1) * limit } + : page && page > 1 + ? { skip: (page - 1) * limit } : {}), }), ]); @@ -500,10 +507,10 @@ export async function buildTransactionsResponse( limit, total, hasNextPage, - hasPrevPage: !!(query.cursor || (query.page && query.page > 1)), + hasPrevPage: !!(query.cursor || (page && page > 1)), nextCursor: hasNextPage && data.length > 0 ? encodeCursor(data[data.length - 1].id) : null, - currentPage: query.page ?? null, - totalPages: query.page ? Math.ceil(total / limit) : null, + currentPage: page ?? null, + totalPages: page ? Math.max(1, Math.ceil(total / limit)) : null, }); return createPaginatedResponse(mapped, pagination, { @@ -627,7 +634,7 @@ export function buildPortfolioHoldingsResponse( query: WalletStateQuery ): PaginatedResponse { const pagination = { - limit: query.limit ?? PORTFOLIO_PAGINATION_CONFIG.defaultLimit ?? 20, + limit: clampLimitNumber(query.limit, PORTFOLIO_PAGINATION_CONFIG), cursor: query.cursor, sortBy: query.sortBy ?? PORTFOLIO_PAGINATION_CONFIG.defaultSortBy, sortOrder: query.sortOrder ?? PORTFOLIO_PAGINATION_CONFIG.defaultSortOrder ?? 'desc', @@ -653,7 +660,7 @@ export function buildVaultHistoryResponse( query: Pick ): PaginatedResponse { const pagination = { - limit: query.limit ?? VAULT_HISTORY_PAGINATION_CONFIG.defaultLimit ?? 30, + limit: clampLimitNumber(query.limit, VAULT_HISTORY_PAGINATION_CONFIG), cursor: query.cursor, sortBy: query.sortBy ?? VAULT_HISTORY_PAGINATION_CONFIG.defaultSortBy, sortOrder: query.sortOrder ?? VAULT_HISTORY_PAGINATION_CONFIG.defaultSortOrder ?? 'desc', diff --git a/backend/src/middleware/apiError.ts b/backend/src/middleware/apiError.ts index a59b60aa5..2d168e163 100644 --- a/backend/src/middleware/apiError.ts +++ b/backend/src/middleware/apiError.ts @@ -10,6 +10,75 @@ export interface ApiErrorOptions { retryable?: boolean; retryAfterSeconds?: number | null; error?: string; + summary?: string; + errors?: unknown[]; + path?: string; +} + +/** + * Canonical wire shape for every error response. + * + * `error`/`status`/`code`/`message`/`retryable` are always present; `details`, + * `correlationId` and `traceId` are only emitted when known. + * + * `summary`, `errors` and `path` are additive and only set by the callers that + * have that information: `errors` mirrors `details` on validation failures so + * clients can read the field list under either key, and `path` is set by the + * catch-all route handler. + */ +export interface ApiErrorBody { + error: string; + status: number; + code: string; + message: string; + retryable: boolean; + details?: unknown; + correlationId?: string; + traceId?: string; + summary?: string; + errors?: unknown[]; + path?: string; +} + +export interface BuildApiErrorBodyOptions { + status?: number; + code?: string; + message: string; + details?: unknown; + retryable?: boolean; + error?: string; + correlationId?: string; + traceId?: string; + summary?: string; + errors?: unknown[]; + path?: string; +} + +/** + * Build the canonical error envelope for a status/message pair. + * + * Every code path that produces an error body must go through this helper — + * `sendApiError` for thrown/structured errors, `apiErrorContractMiddleware` + * for handlers that hand-roll `res.status(n).json({ error, message })`, and + * in-process snapshot builders that stand in for an HTTP endpoint. Synthesizing + * a response outside this helper is what produced bodies that matched a route + * in name but drifted from the real wire shape. + */ +export function buildApiErrorBody(options: BuildApiErrorBodyOptions): ApiErrorBody { + const status = options.status ?? 500; + return { + error: options.error ?? statusLabel(status), + status, + code: options.code ?? defaultErrorCode(status), + message: options.message, + retryable: options.retryable ?? status >= 500, + ...(options.details !== undefined ? { details: options.details } : {}), + ...(options.summary !== undefined ? { summary: options.summary } : {}), + ...(options.errors !== undefined ? { errors: options.errors } : {}), + ...(options.path !== undefined ? { path: options.path } : {}), + ...(options.correlationId ? { correlationId: options.correlationId } : {}), + ...(options.traceId ? { traceId: options.traceId } : {}), + }; } export function sendApiError( @@ -24,16 +93,21 @@ export function sendApiError( res.setHeader('Retry-After', String(options.retryAfterSeconds)); } - res.status(options.status).json({ - error: options.error ?? statusLabel(options.status), - status: options.status, - code: options.code, - message: options.message, - retryable: options.retryable ?? options.status >= 500, - ...(options.details !== undefined ? { details: options.details } : {}), - ...(correlationId ? { correlationId } : {}), - ...(traceId ? { traceId } : {}), - }); + res.status(options.status).json( + buildApiErrorBody({ + error: options.error, + status: options.status, + code: options.code, + message: options.message, + retryable: options.retryable, + details: options.details, + summary: options.summary, + errors: options.errors, + path: options.path, + correlationId, + ...(traceId ? { traceId } : {}), + }) + ); } export function apiErrorContractMiddleware( @@ -52,13 +126,20 @@ export function apiErrorContractMiddleware( return json(body); } + // Every field is passed explicitly so this normalization keeps its original + // defaults; the shared builder supplies the shape (key set and order), not + // the fallback values. return json({ ...errorBody, - error: typeof errorBody.error === 'string' ? errorBody.error : statusLabel(res.statusCode), - status: typeof errorBody.status === 'number' ? errorBody.status : res.statusCode, - code: typeof errorBody.code === 'string' ? errorBody.code : defaultErrorCode(res.statusCode), - message: typeof errorBody.message === 'string' ? errorBody.message : String(errorBody.error), - retryable: typeof errorBody.retryable === 'boolean' ? errorBody.retryable : res.statusCode >= 500, + ...buildApiErrorBody({ + error: typeof errorBody.error === 'string' ? errorBody.error : statusLabel(res.statusCode), + status: typeof errorBody.status === 'number' ? errorBody.status : res.statusCode, + code: typeof errorBody.code === 'string' ? errorBody.code : defaultErrorCode(res.statusCode), + message: + typeof errorBody.message === 'string' ? errorBody.message : String(errorBody.error), + retryable: typeof errorBody.retryable === 'boolean' ? errorBody.retryable : res.statusCode >= 500, + details: errorBody.details, + }), }); }) as Response['json']; diff --git a/backend/src/middleware/validate.ts b/backend/src/middleware/validate.ts index 76e48ecf9..3568203d0 100644 --- a/backend/src/middleware/validate.ts +++ b/backend/src/middleware/validate.ts @@ -244,11 +244,17 @@ export function validate(schemas: ValidateTargets) { message: e.message, })); + // `details` carries the field list for clients that read the canonical + // key; `errors` mirrors it because the webhook and admin validation + // contract documents the field list under `errors`, and `summary` is the + // stable human-readable label those clients assert on. sendApiError(req, res, { status: 400, code: 'VALIDATION_ERROR', + summary: 'Request validation failed', message: formatZodError(issues), details, + errors: details, retryable: false, }); return; diff --git a/backend/src/pagination.ts b/backend/src/pagination.ts index da4d1b0ee..00ba155db 100644 --- a/backend/src/pagination.ts +++ b/backend/src/pagination.ts @@ -95,6 +95,52 @@ export const DEFAULT_PAGINATION_CONFIG: PaginationConfig = { // ─── Query Parsing ────────────────────────────────────────────────────────── +/** + * Clamp a parsed 1-based page number to a usable value. + * + * Non-numeric, non-finite and non-positive values all become 1 so that + * out-of-range `page` params resolve to the first page instead of producing a + * negative offset or a negative `currentPage` in the response envelope. + * + * This is the single definition of page clamping: `parsePaginationQuery` and the + * list-response builders both route through it so the request parser and the + * response metadata can never disagree about what a given `page` means. + */ +export function clampPageNumber(page: number | undefined): number | undefined { + if (page === undefined) { + return undefined; + } + return Number.isFinite(page) && page > 0 ? Math.floor(page) : 1; +} + +/** + * Clamp a requested page size into the endpoint's usable range. + * + * Undefined, non-numeric, non-finite and non-positive values all fall back to + * the endpoint default; anything above `maxLimit` is capped at `maxLimit`. This + * never rejects — a `limit` the caller could not use becomes a valid page size + * instead of a 400, which is what the contract tests require of every list + * endpoint. + * + * Like `clampPageNumber`, this is the single definition of limit clamping: + * `parsePaginationQuery` and the list-response builders both route through it, + * so a hand-rolled `parseInt(req.query.limit)` in a route handler can no longer + * leak a `NaN`, zero or oversized `limit` into a pagination envelope. + */ +export function clampLimitNumber( + limit: number | undefined, + config: Partial = {} +): number { + const mergedConfig = { ...DEFAULT_PAGINATION_CONFIG, ...config }; + + if (limit === undefined || !Number.isFinite(limit)) { + return mergedConfig.defaultLimit; + } + + const floored = Math.floor(limit); + return floored > 0 ? Math.min(floored, mergedConfig.maxLimit) : mergedConfig.defaultLimit; +} + /** * Parse and validate pagination query parameters from request. * @@ -109,28 +155,20 @@ export function parsePaginationQuery( const mergedConfig = { ...DEFAULT_PAGINATION_CONFIG, ...config }; const query: PaginationQuery = {}; - // Parse limit + // Parse limit (clamped — never rejected) if (req.query.limit !== undefined) { - const limit = parseInt(req.query.limit as string, 10); - if (!isNaN(limit) && limit > 0) { - query.limit = Math.min(limit, mergedConfig.maxLimit); - } + query.limit = clampLimitNumber(parseInt(req.query.limit as string, 10), mergedConfig); } - query.limit = query.limit || mergedConfig.defaultLimit; + query.limit = query.limit ?? mergedConfig.defaultLimit; // Parse cursor (opaque string, no validation needed) if (req.query.cursor !== undefined && typeof req.query.cursor === 'string') { query.cursor = req.query.cursor; } - // Parse page (1-based) + // Parse page (1-based, clamped — never rejected) if (req.query.page !== undefined) { - const page = parseInt(req.query.page as string, 10); - if (!isNaN(page) && page > 0) { - query.page = page; - } else { - query.page = 1; - } + query.page = clampPageNumber(parseInt(req.query.page as string, 10)); } // Parse sortBy @@ -167,12 +205,13 @@ export function paginateWithCursor( query: PaginationQuery, getCursor: (item: T) => string ): { data: T[]; pagination: PaginationMeta } { - const limit = query.limit || DEFAULT_PAGINATION_CONFIG.defaultLimit; + const limit = clampLimitNumber(query.limit); + const page = clampPageNumber(query.page); let startIndex = 0; const invalidCursor = false; - if (query.page && query.page > 0) { - startIndex = (query.page - 1) * limit; + if (page && page > 1) { + startIndex = (page - 1) * limit; } // Find starting position based on cursor @@ -184,11 +223,11 @@ export function paginateWithCursor( pagination: createPaginationEnvelope({ count: 0, limit, - total: query.page ? items.length : null, + total: page ? items.length : null, hasNextPage: false, hasPrevPage: false, - currentPage: query.page ? Math.max(1, query.page) : null, - totalPages: query.page ? Math.max(1, Math.ceil(items.length / limit)) : null, + currentPage: page ?? null, + totalPages: page ? Math.max(1, Math.ceil(items.length / limit)) : null, }), }; } @@ -221,8 +260,8 @@ export function paginateWithCursor( total: items.length, hasNextPage: hasMore, hasPrevPage: startIndex > 0, - currentPage: query.page || null, - totalPages: query.page ? Math.max(1, Math.ceil(items.length / limit)) : null, + currentPage: page ?? null, + totalPages: page ? Math.max(1, Math.ceil(items.length / limit)) : null, }); if (hasMore && data.length > 0) { @@ -251,8 +290,8 @@ export function paginateWithOffset( items: T[], query: PaginationQuery ): { data: T[]; pagination: PaginationMeta } { - const limit = query.limit || DEFAULT_PAGINATION_CONFIG.defaultLimit; - const page = query.page || 1; + const limit = clampLimitNumber(query.limit); + const page = clampPageNumber(query.page) ?? 1; const startIndex = (page - 1) * limit; const endIndex = startIndex + limit; diff --git a/backend/src/referralEndpoints.ts b/backend/src/referralEndpoints.ts index 8fe6b6176..f4d26b9b7 100644 --- a/backend/src/referralEndpoints.ts +++ b/backend/src/referralEndpoints.ts @@ -3,10 +3,32 @@ import { referralService } from './referralService'; import { logger } from './middleware/structuredLogging'; import { normalizeWalletAddress } from './walletUtils'; import { cacheMiddleware } from './middleware/cache'; +import { buildApiErrorBody, type ApiErrorBody } from './middleware/apiError'; const router = Router(); const REFERRAL_CACHE_TTL_MS = parseInt(process.env.CACHE_LIST_ENDPOINTS_TTL_MS || '30000', 10); +const WALLET_REQUIRED_MESSAGE = 'Wallet address is required'; +const REFERRAL_STATS_NOT_FOUND_MESSAGE = 'No referral activity found for this wallet'; +const REFERRAL_STATS_ERROR_MESSAGE = 'Failed to fetch referral stats'; +const REFERRAL_CODE_ERROR_MESSAGE = 'Failed to get referral code'; + +export function buildReferralWalletRequiredBody(): ApiErrorBody { + return buildApiErrorBody({ status: 400, message: WALLET_REQUIRED_MESSAGE }); +} + +export function buildReferralStatsNotFoundBody(): ApiErrorBody { + return buildApiErrorBody({ status: 404, message: REFERRAL_STATS_NOT_FOUND_MESSAGE }); +} + +export function buildReferralStatsErrorBody(): ApiErrorBody { + return buildApiErrorBody({ status: 500, message: REFERRAL_STATS_ERROR_MESSAGE }); +} + +export function buildReferralCodeErrorBody(): ApiErrorBody { + return buildApiErrorBody({ status: 500, message: REFERRAL_CODE_ERROR_MESSAGE }); +} + /** * @openapi * /api/v1/referrals/{wallet}: @@ -39,11 +61,7 @@ router.get('/:wallet', cacheMiddleware({ ttl: REFERRAL_CACHE_TTL_MS }), async (r const { wallet } = req.params; if (!wallet) { - return res.status(400).json({ - error: 'Bad Request', - status: 400, - message: 'Wallet address is required', - }); + return res.status(400).json(buildReferralWalletRequiredBody()); } const normalizedWallet = normalizeWalletAddress(wallet); @@ -52,11 +70,7 @@ router.get('/:wallet', cacheMiddleware({ ttl: REFERRAL_CACHE_TTL_MS }), async (r const stats = await referralService.getReferralStats(normalizedWallet); if (!stats) { - return res.status(404).json({ - error: 'Not Found', - status: 404, - message: 'No referral activity found for this wallet', - }); + return res.status(404).json(buildReferralStatsNotFoundBody()); } return res.status(200).json(stats); @@ -65,11 +79,7 @@ router.get('/:wallet', cacheMiddleware({ ttl: REFERRAL_CACHE_TTL_MS }), async (r error: error instanceof Error ? error.message : String(error), wallet: normalizedWallet, }); - return res.status(500).json({ - error: 'Internal Server Error', - status: 500, - message: 'Failed to fetch referral stats', - }); + return res.status(500).json(buildReferralStatsErrorBody()); } }); @@ -102,11 +112,7 @@ router.get('/code/:wallet', cacheMiddleware({ ttl: REFERRAL_CACHE_TTL_MS }), asy const { wallet } = req.params; if (!wallet) { - return res.status(400).json({ - error: 'Bad Request', - status: 400, - message: 'Wallet address is required', - }); + return res.status(400).json(buildReferralWalletRequiredBody()); } const normalizedWallet = normalizeWalletAddress(wallet); @@ -119,11 +125,7 @@ router.get('/code/:wallet', cacheMiddleware({ ttl: REFERRAL_CACHE_TTL_MS }), asy error: error instanceof Error ? error.message : String(error), wallet: normalizedWallet, }); - return res.status(500).json({ - error: 'Internal Server Error', - status: 500, - message: 'Failed to get referral code', - }); + return res.status(500).json(buildReferralCodeErrorBody()); } }); diff --git a/backend/src/transactionEndpoints.ts b/backend/src/transactionEndpoints.ts index ef1a1728c..42d690936 100644 --- a/backend/src/transactionEndpoints.ts +++ b/backend/src/transactionEndpoints.ts @@ -19,7 +19,7 @@ import { parseTypeFilter, parseStatusFilter, } from './transactionQuery'; -import { buildTransactionsResponse } from './listEndpoints'; +import { buildTransactionsResponse, TRANSACTION_PAGINATION_CONFIG } from './listEndpoints'; import { cacheMiddleware } from './middleware/cache'; import { tenantGuard } from './middleware/tenantGuard'; import { Permission } from './middleware/rbac'; @@ -63,6 +63,12 @@ router.get('/', const from = req.query.from as string | undefined; const to = req.query.to as string | undefined; + // One parser for both branches below. Hand-rolling parseInt here is what + // let `limit=abc` through as NaN (Prisma `take: NaN` -> 500) and let + // `limit=100000` bypass the endpoint max, both of which the contract + // tests require to resolve gracefully to a clamped 200 (#1318). + const paginationQuery = parsePaginationQuery(req, TRANSACTION_PAGINATION_CONFIG); + if (!walletAddress) { // Validate type filter if provided const { error: typeError } = parseTypeFilter(typeof type === 'string' ? type : undefined); @@ -82,11 +88,7 @@ router.get('/', try { const response = await buildTransactionsResponse({ - limit: typeof req.query.limit === 'string' ? parseInt(req.query.limit, 10) : undefined, - cursor: typeof req.query.cursor === 'string' ? req.query.cursor : undefined, - page: typeof req.query.page === 'string' ? parseInt(req.query.page, 10) : undefined, - sortBy: typeof req.query.sortBy === 'string' ? req.query.sortBy : undefined, - sortOrder: req.query.sortOrder === 'asc' ? 'asc' : 'desc', + ...paginationQuery, type: typeof type === 'string' ? type : undefined, status: typeof status === 'string' ? status : undefined, from, @@ -126,13 +128,7 @@ router.get('/', return; } - // Parse pagination parameters - const paginationQuery = parsePaginationQuery(req, { - ...DEFAULT_PAGINATION_CONFIG, - defaultSortBy: 'timestamp', - defaultSortOrder: 'desc', - }); - + // Parse pagination parameters (shared with the unscoped branch above) // Resolve and validate the requested sort field against the allowlist. const sort = resolveTransactionSort(paginationQuery.sortBy, paginationQuery.sortOrder); if (!sort.valid) { diff --git a/backend/src/types/validation.ts b/backend/src/types/validation.ts index 20041be10..4d18bf63b 100644 --- a/backend/src/types/validation.ts +++ b/backend/src/types/validation.ts @@ -22,11 +22,19 @@ export const stellarWalletAddressField = z export const walletAddressField = z.string().trim().min(1, 'walletAddress is required'); +/** + * Pagination params are validated for *shape* only. Numeric ranges are + * deliberately not enforced here: `parsePaginationQuery` already clamps every + * out-of-range value gracefully (non-numeric / <=0 `page` -> 1, non-numeric / + * <=0 `limit` -> default, `limit` above the endpoint max -> max), and contract + * tests require those inputs to resolve to a 200 first page rather than a 400. + * Rejecting them at the schema layer made the two disagree for the same input. + */ export const PaginationQuerySchema = z .object({ - limit: z.string().regex(/^\d+$/, 'limit must be a positive integer').optional(), + limit: z.string().optional(), cursor: z.string().optional(), - page: z.string().regex(/^\d+$/, 'page must be a positive integer').optional(), + page: z.string().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..681a4e449 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, @@ -31,6 +31,7 @@ import crypto from 'crypto'; // crypto is still used below for generateFingerprint and body.id generation. import { tryAcquireWalletLock } from './walletLock'; import { normalizeWalletAddress } from './walletUtils'; +import { clampLimitNumber, type PaginationConfig } from './pagination'; import { recordVaultLifecycleEvent } from './vaultAuditLog'; import { registerWithdrawalPlan, @@ -45,6 +46,12 @@ const ZERO = new Decimal(0); const DEFAULT_SHARE_PRICE = new Decimal(1); const STRATEGY_CACHE_TTL_MS = parseInt(process.env.CACHE_STRATEGY_TTL_MS || '30000', 10); +/** Page size bounds for GET /receipts; see clampLimitNumber for the semantics. */ +const RECEIPTS_PAGINATION_CONFIG: Partial = { + defaultLimit: 50, + maxLimit: 100, +}; + // Register cache invalidation hooks for transaction state changes registerInvalidationHook((eventType) => { if (eventType.startsWith('transaction.')) { @@ -769,7 +776,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)' }); }); /** @@ -893,14 +900,23 @@ router.post( router.get('/receipts', readsLimiter, async (req: Request, res: Response) => { const prisma = getPrismaClient(); const wallet = req.query.wallet as string | undefined; - const limit = Math.min(parseInt(req.query.limit as string || '50', 10), 100); + // Clamp rather than parseInt: `limit=abc` used to become NaN and `limit=-5` + // stayed negative, both of which reach Prisma as `take` and 500 the request + // instead of serving a valid page (#1318). + const limit = clampLimitNumber( + req.query.limit === undefined ? undefined : parseInt(req.query.limit as string, 10), + RECEIPTS_PAGINATION_CONFIG + ); const cursor = req.query.cursor as string | undefined; const where = wallet ? { user: wallet } : {}; const transactions = await prisma.transaction.findMany({ where, - orderBy: { createdAt: 'desc' }, + // `Transaction` records their time in `timestamp`; ordering by `createdAt` + // made every call fail Prisma's validation, so this endpoint could only + // ever answer 500. + orderBy: { timestamp: 'desc' }, take: limit + 1, ...(cursor ? { cursor: { id: cursor }, skip: 1 } : {}), }); @@ -916,7 +932,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({