diff --git a/backend/.env.example b/backend/.env.example index eccf39c3..6295ebb9 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -77,6 +77,15 @@ CORS_ALLOWED_ORIGINS=http://localhost:3000 CREATE_PAYMENT_RATE_LIMIT_MAX=50 CREATE_PAYMENT_RATE_LIMIT_WINDOW_MS=60000 +# Payment session persistence retry (exponential backoff with full jitter, #1449) +# Total attempts including the first (clamped 1-6), base delay and delay cap (ms, max 10000) +PAYMENT_SESSION_RETRY_MAX_ATTEMPTS=3 +PAYMENT_SESSION_RETRY_BASE_DELAY_MS=100 +PAYMENT_SESSION_RETRY_MAX_DELAY_MS=2000 + +# Per-(merchant, Idempotency-Key) session lock TTL in ms (clamped 1000-120000, #1450) +PAYMENT_SESSION_LOCK_TTL_MS=30000 + # Resend API key for sending payment receipt emails to merchants RESEND_API_KEY=your_resend_api_key diff --git a/backend/PAYMENT_SESSION_VALIDATOR_SECURITY_AUDIT.md b/backend/PAYMENT_SESSION_VALIDATOR_SECURITY_AUDIT.md new file mode 100644 index 00000000..76c67292 --- /dev/null +++ b/backend/PAYMENT_SESSION_VALIDATOR_SECURITY_AUDIT.md @@ -0,0 +1,83 @@ +# Payment Session Validator & Merchant Settings Security Audit + +**Modules:** +- `backend/src/lib/payment-session-retry.js` (new) +- `backend/src/lib/payment-session-lock.js` (new) +- `backend/src/lib/merchant-payload-validation.js` (new) +- `backend/src/routes/payments.js`, `backend/src/services/paymentService.js` +- `backend/src/routes/merchants.js`, `backend/src/services/merchantService.js` +- `backend/src/lib/idempotency.js`, `backend/src/lib/request-schemas.js`, `backend/src/lib/webhooks.js`, `backend/src/lib/merchant-settings.js` + +**Issues:** #1449 (retry/backoff), #1450 (distributed locking), #1451 (integration & stress tests), #1482 (payload sanitization & strict validation) + +## Threat Model + +| Threat | Mitigation | Status | +|--------|------------|--------| +| Transient DB/network blip fails session creation | Persistence retried with full-jitter exponential backoff; bounded attempts/delay (#1449) | ✅ Implemented | +| Retry turns a rejected session into an accepted one | Only transport/5xx/429/connection-class Postgres errors retried; 4xx and constraint violations never retried | ✅ Implemented | +| Retry creates duplicate rows after a lost ack | Session id generated before first attempt; `23505` on a *retry* treated as already-persisted, on the *first* attempt surfaced | ✅ Implemented | +| Retry storm / request hang | Attempts clamped to ≤ 6, delay ≤ 10 s, full jitter spreads load | ✅ Implemented | +| Concurrent same-key requests create two sessions | `SET NX PX` lock per (merchant, Idempotency-Key); losers get `409 PAYMENT_SESSION_IN_PROGRESS` (#1450) | ✅ Implemented | +| Twin passes idempotency middleware before cache write | Idempotency cache re-checked *inside* the lock; replayed if present | ✅ Implemented | +| Releasing someone else's lock after TTL expiry | Random per-acquire token + compare-and-delete Lua release | ✅ Implemented | +| Lock key injection / cross-merchant collision | Key = `lock:payment-session::` | ✅ Implemented | +| Redis outage silently disables locking | No-op Redis client detected (`isOpen:false`); in-process lock fallback | ✅ Implemented | +| Prototype pollution via settings payloads | `__proto__`/`constructor`/`prototype` keys dropped recursively before validation (#1482) | ✅ Implemented | +| Mass assignment (`api_key`, `merchant_id`, `webhook_secret` in body) | `.strict()` schemas reject unknown keys with 400 | ✅ Implemented | +| Control characters / terminal escapes stored in profile fields | NUL and C0 controls (except `\t\n\r`) stripped | ✅ Implemented | +| Resource exhaustion via deep/wide payloads | Depth ≤ 6, ≤ 50 keys/object, ≤ 100 items/array, ≤ 4096 chars/string | ✅ Implemented | +| Webhook header injection (CR/LF) | Header values must be printable ASCII; send-time filter also drops CR/LF/NUL for legacy rows | ✅ Implemented | +| Signature/timestamp/transport header spoofing | Reserved header list enforced at write and send time (incl. `Host`, `Content-Length`, `Transfer-Encoding`) | ✅ Implemented | +| API key expiry lock-out or never-expiring keys | Expiry must be ≥ 1 minute ahead and ≤ 365 days; normalized to UTC | ✅ Implemented | +| Invalid rotate/expiry bodies returned 500 | Validation moved to `validateRequest` → structured 400 | ✅ Fixed | +| Non-HTTP callers bypass route validation | Service-level guards in `merchantService.rotateApiKey` / `setApiKeyExpiry` | ✅ Implemented | + +## Design Notes + +### Retry (#1449) + +`withSessionRetry(fn, opts)` runs `fn({ attempt })` and retries when `isRetryableSessionError(err)` is true. `delay(n) = random(0, min(maxDelayMs, baseDelayMs · 2ⁿ))`. The last error is rethrown with `retryAttempts` set, so upstream status mapping is unchanged. + +On-chain issuer verification is **not** wrapped again: `AssetIssuerErrorRecovery.verifyIssuerOnChain` already retries with its own circuit breaker, and nesting the two would multiply attempts. + +### Locking (#1450) + +Only requests that carry an `Idempotency-Key` are serialized. Requests without a key are independent by definition, and serializing them would only add latency. + +The idempotency middleware and the lock share one Redis client. The cached response is written (in `res.json`) *before* the lock is released on the same connection, and Redis runs a connection's commands in order. So a request that gets the lock after its twin finishes is guaranteed to see the cached response. + +**Degraded mode (Redis unavailable):** the in-process lock still prevents overlapping creations on one instance. There is no idempotency cache, so a request that arrives *after* its twin completed creates a new session. This matches the existing middleware's fail-open behavior. + +### Validation (#1482) + +Order on every settings / API-key route: auth → rate limit → `sanitizeMerchantPayload` → `validateRequest(strictSchema)` → handler. + +`POST /api/register-merchant` stays lenient at the top level for client compatibility, but `merchant_settings` is strict and `metadata` is sanitized. + +## Configuration + +| Variable | Default | Bounds | +|----------|---------|--------| +| `PAYMENT_SESSION_RETRY_MAX_ATTEMPTS` | 3 | 1–6 | +| `PAYMENT_SESSION_RETRY_BASE_DELAY_MS` | 100 | 0–10000 | +| `PAYMENT_SESSION_RETRY_MAX_DELAY_MS` | 2000 | base–10000 | +| `PAYMENT_SESSION_LOCK_TTL_MS` | 30000 | 1000–120000 | + +## API Behavior Changes + +- `POST /api/sessions`, `/api/create-payment`, `/api/support-transactions`: may return **409** `{ code: "PAYMENT_SESSION_IN_PROGRESS" }` when a same-key request is in flight. Clients should retry after a short delay. +- `POST /api/merchants/rotate-api-key`, `PUT /api/merchants/set-api-key-expiry`, `POST /api/merchants/rotate-webhook-secret`, `PUT /api/webhook-settings`: unknown body fields now return **400** (previously ignored, or a 500 on type errors). +- `PUT /api/merchants/set-api-key-expiry`: past dates and dates more than 365 days ahead now return **400**. +- `PUT /api/webhook-settings`: reserved names, duplicate names (case-insensitive), more than 20 headers, and values with line breaks now return **400**. + +## Test Coverage + +| Suite | Scope | +|-------|-------| +| `src/lib/payment-session-retry.test.js` | Error classification, option clamping, backoff math, retry/exhaustion, lost-ack recovery | +| `src/lib/payment-session-lock.test.js` | Key hashing, SET NX/PX, compare-and-delete, TTL expiry, local fallback, 409 conflict | +| `src/routes/payment-session-validator.integration.test.js` | HTTP end-to-end: rules, retry, concurrency (single and multi-node), replay, Redis-down fallback, stress (200 concurrent at 30% fault rate; 20 keys × 10 duplicates) | +| `src/lib/merchant-payload-validation.test.js` | Sanitizer limits, prototype pollution, strict schemas, header rules, expiry bounds | +| `src/routes/merchant-settings-validation.test.js` | HTTP-level 400s for mass assignment, header injection, expiry abuse, pollution | +| `src/services/merchantService.validation.test.js` | Service-level guards | diff --git a/backend/src/lib/idempotency.js b/backend/src/lib/idempotency.js index aedd2a46..d75619bd 100644 --- a/backend/src/lib/idempotency.js +++ b/backend/src/lib/idempotency.js @@ -6,6 +6,43 @@ import { getRedisClient } from "./redis.js"; */ const IDEMPOTENCY_TTL_SECONDS = 24 * 60 * 60; // 24 hours +export function buildIdempotencyRedisKey(merchantId, idempotencyKey) { + return `idempotency:${merchantId}:${idempotencyKey}`; +} + +export function hashIdempotencyPayload(body) { + return crypto + .createHash("sha256") + .update(JSON.stringify(body || {})) + .digest("hex"); +} + +/** + * Look up a previously cached idempotent response. + * + * Used by the payment session handler after it acquires the per-key session + * lock (issue #1450): a request that passed the middleware while a concurrent + * twin was still in flight must replay that twin's response instead of + * creating a second session. + * + * `payloadHash` must be the hash the middleware computed from the RAW body + * (exposed as `req.idempotency.payloadHash`), because downstream Zod schemas + * may transform `req.body`. + * + * @returns {Promise} + */ +export async function lookupIdempotentResponse({ redisClient, merchantId, idempotencyKey, payloadHash }) { + if (!redisClient || !merchantId || !idempotencyKey || !payloadHash) return null; + const cachedValue = await redisClient.get(buildIdempotencyRedisKey(merchantId, idempotencyKey)); + if (!cachedValue) return null; + + const { hash, response } = JSON.parse(cachedValue); + if (hash !== payloadHash) { + return { status: "mismatch" }; + } + return { status: "hit", response }; +} + /** * Idempotency middleware that checks and enforces idempotent requests. * Tracks the key tied to the payload hash and response. @@ -40,13 +77,11 @@ export async function idempotencyMiddleware(req, res, next) { } const redisClient = getRedisClient(); - const redisKey = `idempotency:${merchantId}:${idempotencyKey}`; - + const redisKey = buildIdempotencyRedisKey(merchantId, idempotencyKey); + // Calculate hash of payload to ensure consistency - const payloadHash = crypto - .createHash("sha256") - .update(JSON.stringify(req.body || {})) - .digest("hex"); + const payloadHash = hashIdempotencyPayload(req.body); + req.idempotency = { key: idempotencyKey, payloadHash }; try { const cachedValue = await redisClient.get(redisKey); diff --git a/backend/src/lib/merchant-payload-validation.js b/backend/src/lib/merchant-payload-validation.js new file mode 100644 index 00000000..c41c3a15 --- /dev/null +++ b/backend/src/lib/merchant-payload-validation.js @@ -0,0 +1,333 @@ +/** + * merchant-payload-validation.js + * + * Payload sanitization and strict validation for the Merchant Settings & + * API Key Service (issue #1482). + * + * Two layers, applied in this order on every settings / API-key route: + * + * 1. sanitizeMerchantPayload (middleware) + * - drops prototype-pollution keys (__proto__, constructor, prototype) + * - strips NUL and non-printable control characters from strings + * - enforces depth, key-count, array-length and string-length ceilings + * so a single request cannot force deep recursion or huge DB writes + * + * 2. Strict Zod schemas (via validateRequest) + * - `.strict()` objects: unknown keys are REJECTED, not silently ignored, + * so typos and mass-assignment attempts (e.g. `api_key`, `merchant_id`) + * fail loudly with 400 instead of being dropped or persisted + * - bounded numbers/dates (API-key expiry must be in the future and + * within MAX_API_KEY_LIFETIME_DAYS) + * - webhook custom headers reject CR/LF (header injection) and reserved + * system header names (signature/timestamp spoofing) + */ + +import { z } from "zod"; +import { logger } from "./logger.js"; + +export const SANITIZE_LIMITS = Object.freeze({ + maxDepth: 6, + maxKeysPerObject: 50, + maxArrayLength: 100, + maxStringLength: 4096, +}); + +const FORBIDDEN_KEYS = new Set(["__proto__", "constructor", "prototype"]); + +// C0 controls except \t \n \r, plus DEL. \r and \n are kept here so free-text +// fields remain usable; header values are rejected separately below. +// eslint-disable-next-line no-control-regex +const UNSAFE_CONTROL_CHARS = /[\u0000-\u0008\u000B\u000C\u000E-\u001F\u007F]/g; + +export class PayloadSanitizationError extends Error { + constructor(message, path) { + super(message); + this.name = "PayloadSanitizationError"; + this.status = 400; + this.path = path; + } +} + +function formatPath(path) { + return path.length === 0 ? "body" : path.join("."); +} + +/** + * Recursively sanitize an untrusted JSON payload. + * + * Returns a NEW value; the input is never mutated. Throws + * PayloadSanitizationError when a structural limit is exceeded. + * + * @param {unknown} value + * @param {Partial} [limits] + * @param {{ droppedKeys?: string[] }} [report] Collects dropped key paths + */ +export function sanitizePayload(value, limits = {}, report = {}) { + const opts = { ...SANITIZE_LIMITS, ...limits }; + const dropped = report.droppedKeys ?? (report.droppedKeys = []); + + function walk(node, depth, path) { + if (typeof node === "string") { + if (node.length > opts.maxStringLength) { + throw new PayloadSanitizationError( + `String at ${formatPath(path)} exceeds ${opts.maxStringLength} characters`, + formatPath(path), + ); + } + return node.replace(UNSAFE_CONTROL_CHARS, ""); + } + + if (node === null || typeof node !== "object") { + if (typeof node === "number" && !Number.isFinite(node)) { + throw new PayloadSanitizationError( + `Number at ${formatPath(path)} must be finite`, + formatPath(path), + ); + } + return node; + } + + if (depth >= opts.maxDepth) { + throw new PayloadSanitizationError( + `Payload nesting exceeds maximum depth of ${opts.maxDepth}`, + formatPath(path), + ); + } + + if (Array.isArray(node)) { + if (node.length > opts.maxArrayLength) { + throw new PayloadSanitizationError( + `Array at ${formatPath(path)} exceeds ${opts.maxArrayLength} items`, + formatPath(path), + ); + } + return node.map((item, index) => walk(item, depth + 1, [...path, index])); + } + + const keys = Object.keys(node); + if (keys.length > opts.maxKeysPerObject) { + throw new PayloadSanitizationError( + `Object at ${formatPath(path)} exceeds ${opts.maxKeysPerObject} keys`, + formatPath(path), + ); + } + + const out = {}; + for (const key of keys) { + if (FORBIDDEN_KEYS.has(key)) { + dropped.push(formatPath([...path, key])); + continue; + } + const cleanKey = key.replace(UNSAFE_CONTROL_CHARS, ""); + if (cleanKey !== key || cleanKey.length === 0) { + dropped.push(formatPath([...path, JSON.stringify(key)])); + continue; + } + out[cleanKey] = walk(node[key], depth + 1, [...path, cleanKey]); + } + return out; + } + + return walk(value, 0, []); +} + +/** + * Express middleware: sanitize req.body in place before schema validation. + * Structural violations are rejected with 400; dropped keys are logged + * (without values) so abuse attempts are observable. + */ +export function sanitizeMerchantPayload(req, res, next) { + if (req.body === undefined || req.body === null) { + return next(); + } + + const report = {}; + try { + req.body = sanitizePayload(req.body, {}, report); + } catch (err) { + if (err instanceof PayloadSanitizationError) { + logger.warn( + { merchantId: req.merchant?.id, path: err.path, route: req.originalUrl }, + "Rejected merchant payload: sanitization limit exceeded", + ); + return res.status(400).json({ error: "Validation failed", message: err.message }); + } + return next(err); + } + + if (report.droppedKeys.length > 0) { + logger.warn( + { merchantId: req.merchant?.id, droppedKeys: report.droppedKeys, route: req.originalUrl }, + "Dropped unsafe keys from merchant payload", + ); + } + return next(); +} + +// --------------------------------------------------------------------------- +// Strict schemas +// --------------------------------------------------------------------------- + +export const MAX_GRACE_PERIOD_HOURS = 168; +export const MAX_API_KEY_LIFETIME_DAYS = 365; +/** Expiry must be at least this far in the future to avoid instant lock-out. */ +export const MIN_API_KEY_EXPIRY_LEAD_MS = 60 * 1000; + +const gracePeriodHoursSchema = z + .number({ invalid_type_error: "grace_period_hours must be a number" }) + .int("grace_period_hours must be an integer") + .min(0, "grace_period_hours must be >= 0") + .max(MAX_GRACE_PERIOD_HOURS, `grace_period_hours must be <= ${MAX_GRACE_PERIOD_HOURS}`); + +export const rotateApiKeySchema = z + .object({ grace_period_hours: gracePeriodHoursSchema.optional() }) + .strict(); + +export const rotateWebhookSecretSchema = z + .object({ grace_period_hours: gracePeriodHoursSchema.optional() }) + .strict(); + +/** + * Validate an API key expiry timestamp and return it normalized to ISO-8601 + * UTC. Throws a 400 error on failure. Shared by the route schema and the + * service layer (defense in depth for non-HTTP callers). + * + * @param {unknown} value + * @param {number} [now=Date.now()] + * @returns {string} + */ +export function normalizeApiKeyExpiry(value, now = Date.now()) { + const fail = (message) => { + const err = new Error(message); + err.status = 400; + throw err; + }; + + if (typeof value !== "string" || value.trim() === "") { + fail("expires_at must be an ISO 8601 datetime string"); + } + const parsed = z.string().datetime({ offset: true }).safeParse(value.trim()); + if (!parsed.success) { + fail("expires_at must be an ISO 8601 datetime string"); + } + + const ts = Date.parse(value.trim()); + if (!Number.isFinite(ts)) { + fail("expires_at must be a valid datetime"); + } + if (ts < now + MIN_API_KEY_EXPIRY_LEAD_MS) { + fail("expires_at must be at least 1 minute in the future"); + } + if (ts > now + MAX_API_KEY_LIFETIME_DAYS * 24 * 60 * 60 * 1000) { + fail(`expires_at must be within ${MAX_API_KEY_LIFETIME_DAYS} days`); + } + return new Date(ts).toISOString(); +} + +export const setApiKeyExpirySchema = z + .object({ + expires_at: z.string({ required_error: "expires_at is required" }), + }) + .strict() + .transform((body, ctx) => { + try { + return { expires_at: normalizeApiKeyExpiry(body.expires_at) }; + } catch (err) { + ctx.addIssue({ code: z.ZodIssueCode.custom, path: ["expires_at"], message: err.message }); + return z.NEVER; + } + }); + +export const merchantSettingsSchema = z + .object({ + send_success_emails: z.boolean({ + invalid_type_error: "send_success_emails must be a boolean", + }).optional(), + }) + .strict(); + +const HEADER_NAME_RE = /^[A-Za-z0-9\-_]{1,64}$/; +// Printable ASCII + tab only. Rejects CR/LF (response splitting / header +// injection) and non-ASCII that some HTTP stacks mangle. +const HEADER_VALUE_RE = /^[\t\x20-\x7E]{1,1024}$/; +export const MAX_CUSTOM_HEADERS = 20; + +/** + * Header names merchants may not set: signature/timestamp headers would let + * a merchant spoof or confuse verification; transport headers could corrupt + * the outbound request. Kept in sync with sanitizeCustomHeaders() in + * webhooks.js, which silently drops them at send time. + */ +export const RESERVED_WEBHOOK_HEADERS = new Set([ + "content-type", + "content-length", + "transfer-encoding", + "connection", + "host", + "user-agent", + "pluto-signature", + "stellar-signature", + "pluto-timestamp", + "stellar-timestamp", +]); + +export const customHeadersSchema = z + .record( + z.string().regex( + HEADER_NAME_RE, + "Header names must be 1-64 characters: alphanumeric, hyphens, or underscores", + ), + z.string().regex( + HEADER_VALUE_RE, + "Header values must be 1-1024 printable ASCII characters with no line breaks", + ), + ) + .superRefine((headers, ctx) => { + const names = Object.keys(headers); + if (names.length > MAX_CUSTOM_HEADERS) { + ctx.addIssue({ + code: z.ZodIssueCode.custom, + message: `At most ${MAX_CUSTOM_HEADERS} custom headers are allowed`, + }); + } + const seen = new Set(); + for (const name of names) { + const lower = name.toLowerCase(); + if (RESERVED_WEBHOOK_HEADERS.has(lower)) { + ctx.addIssue({ + code: z.ZodIssueCode.custom, + path: [name], + message: `"${name}" is a reserved header and cannot be overridden`, + }); + } + if (seen.has(lower)) { + ctx.addIssue({ + code: z.ZodIssueCode.custom, + path: [name], + message: `Duplicate header name "${name}" (header names are case-insensitive)`, + }); + } + seen.add(lower); + } + }); + +const ASSET_CODE_RE = /^[A-Za-z0-9]{1,12}$/; + +export const paymentLimitsSchema = z + .record( + z.string().regex(ASSET_CODE_RE, "Asset codes must be 1-12 alphanumeric characters"), + z + .object({ + min: z.number().positive().finite().optional(), + max: z.number().positive().finite().optional(), + }) + .strict() + .refine( + (limits) => + limits.min === undefined || limits.max === undefined || limits.min <= limits.max, + { message: "min must be less than or equal to max" }, + ), + ) + .refine((limits) => Object.keys(limits).length <= 50, { + message: "At most 50 asset limits may be configured", + }); diff --git a/backend/src/lib/merchant-payload-validation.test.js b/backend/src/lib/merchant-payload-validation.test.js new file mode 100644 index 00000000..b55260d6 --- /dev/null +++ b/backend/src/lib/merchant-payload-validation.test.js @@ -0,0 +1,312 @@ +import { describe, it, expect, vi, beforeEach } from "vitest"; + +vi.mock("./logger.js", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +import { logger } from "./logger.js"; +import { + MAX_CUSTOM_HEADERS, + PayloadSanitizationError, + customHeadersSchema, + merchantSettingsSchema, + normalizeApiKeyExpiry, + paymentLimitsSchema, + rotateApiKeySchema, + rotateWebhookSecretSchema, + sanitizeMerchantPayload, + sanitizePayload, + setApiKeyExpirySchema, +} from "./merchant-payload-validation.js"; +import { resolveMerchantSettings } from "./merchant-settings.js"; +import { sanitizeCustomHeaders } from "./webhooks.js"; + +const DAY = 24 * 60 * 60 * 1000; + +beforeEach(() => vi.clearAllMocks()); + +describe("sanitizePayload (issue #1482)", () => { + it("returns a structurally equal copy for clean input without mutating it", () => { + const input = { a: 1, b: "x", c: [true, null, { d: 2 }] }; + const out = sanitizePayload(input); + expect(out).toEqual(input); + expect(out).not.toBe(input); + expect(out.c).not.toBe(input.c); + }); + + it("drops prototype-pollution keys at every depth and reports them", () => { + const input = JSON.parse( + '{"__proto__":{"polluted":true},"constructor":{"prototype":{"x":1}},"nested":{"__proto__":{"y":1},"ok":1},"list":[{"prototype":1}]}', + ); + const report = {}; + const out = sanitizePayload(input, {}, report); + + expect(out).toEqual({ nested: { ok: 1 }, list: [{}] }); + expect(Object.hasOwn(out, "__proto__")).toBe(false); + expect({}.polluted).toBeUndefined(); + expect(Object.prototype.polluted).toBeUndefined(); + expect(report.droppedKeys).toEqual( + expect.arrayContaining(["__proto__", "constructor", "nested.__proto__", "list.0.prototype"]), + ); + }); + + it("strips NUL and non-printable control characters but keeps tabs/newlines", () => { + const out = sanitizePayload({ name: "Acme\u0000 Corp\u0007\u001b[31m", note: "line1\nline2\tend" }); + expect(out.name).toBe("Acme Corp[31m"); + expect(out.note).toBe("line1\nline2\tend"); + }); + + it("drops keys that contain control characters", () => { + const report = {}; + const out = sanitizePayload({ "bad\u0000key": 1, good: 2 }, {}, report); + expect(out).toEqual({ good: 2 }); + expect(report.droppedKeys).toHaveLength(1); + }); + + it("rejects payloads nested beyond maxDepth", () => { + let deep = { v: 1 }; + for (let i = 0; i < 10; i += 1) deep = { n: deep }; + expect(() => sanitizePayload(deep)).toThrow(PayloadSanitizationError); + }); + + it("rejects objects with too many keys", () => { + const wide = Object.fromEntries(Array.from({ length: 51 }, (_, i) => [`k${i}`, i])); + expect(() => sanitizePayload(wide)).toThrow(/exceeds 50 keys/); + }); + + it("rejects oversized arrays and strings", () => { + expect(() => sanitizePayload({ a: new Array(101).fill(0) })).toThrow(/exceeds 100 items/); + expect(() => sanitizePayload({ s: "x".repeat(4097) })).toThrow(/exceeds 4096 characters/); + }); + + it("rejects non-finite numbers", () => { + expect(() => sanitizePayload({ n: Infinity })).toThrow(/must be finite/); + expect(() => sanitizePayload({ n: NaN })).toThrow(/must be finite/); + }); + + it("passes primitives through", () => { + expect(sanitizePayload(null)).toBeNull(); + expect(sanitizePayload(42)).toBe(42); + expect(sanitizePayload(false)).toBe(false); + }); + + it("marks errors as HTTP 400 with the offending path", () => { + try { + sanitizePayload({ a: { b: "x".repeat(5000) } }); + throw new Error("expected to throw"); + } catch (err) { + expect(err.status).toBe(400); + expect(err.path).toBe("a.b"); + } + }); +}); + +describe("sanitizeMerchantPayload middleware", () => { + function run(body) { + const req = { body, merchant: { id: "m1" }, originalUrl: "/api/x" }; + const res = { status: vi.fn().mockReturnThis(), json: vi.fn().mockReturnThis() }; + const next = vi.fn(); + sanitizeMerchantPayload(req, res, next); + return { req, res, next }; + } + + it("replaces req.body with the sanitized copy and continues", () => { + const { req, next } = run(JSON.parse('{"__proto__":{"x":1},"a":"b\\u0000"}')); + expect(req.body).toEqual({ a: "b" }); + expect(next).toHaveBeenCalledWith(); + expect(logger.warn).toHaveBeenCalledWith( + expect.objectContaining({ droppedKeys: ["__proto__"], merchantId: "m1" }), + expect.any(String), + ); + }); + + it("responds 400 for structural violations without calling next", () => { + const { res, next } = run({ s: "x".repeat(10_000) }); + expect(res.status).toHaveBeenCalledWith(400); + expect(res.json).toHaveBeenCalledWith( + expect.objectContaining({ error: "Validation failed" }), + ); + expect(next).not.toHaveBeenCalled(); + }); + + it("is a no-op for missing bodies", () => { + const { next } = run(undefined); + expect(next).toHaveBeenCalledWith(); + }); + + it("does not log values of dropped keys (no secret leakage)", () => { + run(JSON.parse('{"__proto__":{"api_key":"sk_live_secret"}}')); + expect(JSON.stringify(logger.warn.mock.calls)).not.toContain("sk_live_secret"); + }); +}); + +describe("rotateApiKeySchema / rotateWebhookSecretSchema", () => { + it.each([rotateApiKeySchema, rotateWebhookSecretSchema])("accepts empty and bounded bodies", (schema) => { + expect(schema.parse({})).toEqual({}); + expect(schema.parse({ grace_period_hours: 0 })).toEqual({ grace_period_hours: 0 }); + expect(schema.parse({ grace_period_hours: 168 })).toEqual({ grace_period_hours: 168 }); + }); + + it.each([ + { grace_period_hours: -1 }, + { grace_period_hours: 169 }, + { grace_period_hours: 1.5 }, + { grace_period_hours: "24" }, + { grace_period_hours: null }, + { grace_period_hours: 24, api_key: "sk_attacker" }, + { merchant_id: "someone-else" }, + ])("rejects %j", (body) => { + expect(rotateApiKeySchema.safeParse(body).success).toBe(false); + expect(rotateWebhookSecretSchema.safeParse(body).success).toBe(false); + }); +}); + +describe("normalizeApiKeyExpiry / setApiKeyExpirySchema", () => { + const now = Date.parse("2026-01-01T00:00:00Z"); + + it("normalizes offsets to UTC ISO-8601", () => { + expect(normalizeApiKeyExpiry("2026-02-01T01:00:00+01:00", now)).toBe("2026-02-01T00:00:00.000Z"); + }); + + it.each([ + ["past timestamp", "2025-12-31T00:00:00Z", /at least 1 minute in the future/], + ["now (instant lock-out)", "2026-01-01T00:00:00Z", /at least 1 minute/], + ["beyond 365 days", "2027-01-02T00:00:01Z", /within 365 days/], + ["date without time", "2026-02-01", /ISO 8601/], + ["garbage", "tomorrow", /ISO 8601/], + ["empty", " ", /ISO 8601/], + ])("rejects %s", (_label, value, message) => { + expect(() => normalizeApiKeyExpiry(value, now)).toThrow(message); + }); + + it("rejects non-strings with status 400", () => { + for (const value of [undefined, null, 123, {}, []]) { + try { + normalizeApiKeyExpiry(value, now); + throw new Error("expected to throw"); + } catch (err) { + expect(err.status).toBe(400); + } + } + }); + + it("schema accepts a valid future expiry and returns the normalized value", () => { + const future = new Date(Date.now() + 30 * DAY).toISOString(); + expect(setApiKeyExpirySchema.parse({ expires_at: future })).toEqual({ + expires_at: new Date(future).toISOString(), + }); + }); + + it("schema rejects unknown keys and missing expires_at", () => { + const future = new Date(Date.now() + DAY).toISOString(); + expect(setApiKeyExpirySchema.safeParse({ expires_at: future, api_key: "x" }).success).toBe(false); + expect(setApiKeyExpirySchema.safeParse({}).success).toBe(false); + }); + + it("schema reports the expiry rule on the expires_at path", () => { + const res = setApiKeyExpirySchema.safeParse({ expires_at: "2000-01-01T00:00:00Z" }); + expect(res.success).toBe(false); + expect(res.error.issues[0].path).toEqual(["expires_at"]); + }); +}); + +describe("merchantSettingsSchema", () => { + it("accepts known settings", () => { + expect(merchantSettingsSchema.parse({ send_success_emails: false })).toEqual({ + send_success_emails: false, + }); + expect(merchantSettingsSchema.parse({})).toEqual({}); + }); + + it("rejects unknown keys and wrong types", () => { + expect(merchantSettingsSchema.safeParse({ send_success_emails: "yes" }).success).toBe(false); + expect(merchantSettingsSchema.safeParse({ is_admin: true }).success).toBe(false); + }); +}); + +describe("resolveMerchantSettings hardening", () => { + it("ignores inherited (prototype) values", () => { + const proto = { send_success_emails: false }; + const input = Object.create(proto); + expect(resolveMerchantSettings(input)).toEqual({ send_success_emails: true }); + }); + + it("ignores arrays and unknown keys", () => { + expect(resolveMerchantSettings([false])).toEqual({ send_success_emails: true }); + expect(resolveMerchantSettings({ send_success_emails: false, extra: 1 })).toEqual({ + send_success_emails: false, + }); + }); +}); + +describe("customHeadersSchema", () => { + it("accepts safe headers", () => { + expect(customHeadersSchema.parse({ "X-Tenant": "acme", Authorization: "Bearer abc" })).toEqual({ + "X-Tenant": "acme", + Authorization: "Bearer abc", + }); + }); + + it.each([ + ["CRLF header injection", { "X-A": "ok\r\nX-Injected: 1" }], + ["bare LF", { "X-A": "a\nb" }], + ["non-ASCII value", { "X-A": "café" }], + ["empty value", { "X-A": "" }], + ["value too long", { "X-A": "x".repeat(1025) }], + ["name with colon", { "X-A:": "v" }], + ["name with space", { "X A": "v" }], + ["name too long", { ["X".repeat(65)]: "v" }], + ["reserved signature header", { "Pluto-Signature": "forged" }], + ["reserved timestamp header", { "stellar-timestamp": "0" }], + ["reserved host header", { Host: "evil.example" }], + ["reserved content-length", { "Content-Length": "0" }], + ["case-insensitive duplicates", { "X-Dup": "a", "x-dup": "b" }], + ["non-string value", { "X-A": 1 }], + ])("rejects %s", (_label, headers) => { + expect(customHeadersSchema.safeParse(headers).success).toBe(false); + }); + + it(`rejects more than ${MAX_CUSTOM_HEADERS} headers`, () => { + const many = Object.fromEntries( + Array.from({ length: MAX_CUSTOM_HEADERS + 1 }, (_, i) => [`X-H${i}`, "v"]), + ); + expect(customHeadersSchema.safeParse(many).success).toBe(false); + }); +}); + +describe("sanitizeCustomHeaders (send-time defense for legacy rows)", () => { + it("drops CR/LF/NUL values and reserved transport headers", () => { + expect( + sanitizeCustomHeaders({ + "X-Ok": "fine", + "X-Inject": "a\r\nX-Evil: 1", + "X-Nul": "a\u0000b", + Host: "evil.example", + "Transfer-Encoding": "chunked", + "PLUTO-Signature": "forged", + }), + ).toEqual({ "X-Ok": "fine" }); + }); +}); + +describe("paymentLimitsSchema", () => { + it("accepts valid limits", () => { + expect(paymentLimitsSchema.parse({ USDC: { min: 1, max: 100 }, XLM: { max: 5 } })).toBeTruthy(); + }); + + it.each([ + ["min greater than max", { USDC: { min: 10, max: 1 } }], + ["negative min", { USDC: { min: -1 } }], + ["unknown limit field", { USDC: { min: 1, cap: 5 } }], + ["invalid asset code", { "US DC": { min: 1 } }], + ["asset code too long", { ABCDEFGHIJKLM: { min: 1 } }], + ["infinite max", { USDC: { max: Infinity } }], + ])("rejects %s", (_label, limits) => { + expect(paymentLimitsSchema.safeParse(limits).success).toBe(false); + }); + + it("rejects more than 50 configured assets", () => { + const many = Object.fromEntries(Array.from({ length: 51 }, (_, i) => [`A${i}`, { min: 1 }])); + expect(paymentLimitsSchema.safeParse(many).success).toBe(false); + }); +}); diff --git a/backend/src/lib/merchant-settings.js b/backend/src/lib/merchant-settings.js index 5c97fbb8..ff135540 100644 --- a/backend/src/lib/merchant-settings.js +++ b/backend/src/lib/merchant-settings.js @@ -2,12 +2,20 @@ export const DEFAULT_MERCHANT_SETTINGS = Object.freeze({ send_success_emails: true, }); +/** + * Resolve stored/submitted settings into the canonical shape. Only OWN, + * correctly-typed properties are honoured, so inherited (prototype-polluted) + * values and unknown keys never leak into persisted settings (issue #1482). + */ export function resolveMerchantSettings(rawSettings) { const input = - rawSettings && typeof rawSettings === "object" ? rawSettings : {}; + rawSettings && typeof rawSettings === "object" && !Array.isArray(rawSettings) + ? rawSettings + : {}; return { send_success_emails: + Object.hasOwn(input, "send_success_emails") && typeof input.send_success_emails === "boolean" ? input.send_success_emails : DEFAULT_MERCHANT_SETTINGS.send_success_emails, diff --git a/backend/src/lib/payment-session-lock.js b/backend/src/lib/payment-session-lock.js new file mode 100644 index 00000000..bbb28b42 --- /dev/null +++ b/backend/src/lib/payment-session-lock.js @@ -0,0 +1,228 @@ +/** + * payment-session-lock.js + * + * Distributed concurrency control for the Payment Session Validator + * (issue #1450). + * + * Problem + * ------- + * The idempotency middleware checks Redis for a cached response and only + * writes the cache AFTER the handler responds. Two concurrent requests that + * carry the same Idempotency-Key (client retry storms, double-clicks, load + * balancer replays) can therefore both miss the cache, both pass validation + * and both create a payment session — on the same or on different instances. + * + * Solution + * -------- + * A short-lived mutual-exclusion lock per (merchant, Idempotency-Key): + * + * - Redis: SET NX PX — atomic acquire across instances + * - Release: compare-and-delete Lua script — a request can only release the + * lock it owns, never one that expired and was re-acquired by another. + * - Fallback: when Redis is unavailable, an in-process lock keeps a single + * instance correct instead of silently failing open. + * + * A request that loses the race gets HTTP 409 (never a second session); the + * client simply retries and then receives the cached idempotent response. + * + * Keys are namespaced per merchant and the client-supplied Idempotency-Key is + * SHA-256 hashed, so it cannot inject Redis key separators, collide across + * merchants, or create unbounded key sizes. + */ + +import { createHash, randomUUID } from "node:crypto"; +import { connectRedisClient } from "./redis.js"; +import { logger } from "./logger.js"; + +export const DEFAULT_LOCK_TTL_MS = 30_000; +const MIN_LOCK_TTL_MS = 1_000; +const MAX_LOCK_TTL_MS = 120_000; +const LOCK_KEY_PREFIX = "lock:payment-session"; + +/** Error code for a request rejected because its twin holds the lock. */ +export const PAYMENT_SESSION_IN_PROGRESS = "PAYMENT_SESSION_IN_PROGRESS"; + +/** + * Compare-and-delete: only the owner (matching token) may release the lock. + */ +export const RELEASE_LOCK_SCRIPT = + 'if redis.call("get", KEYS[1]) == ARGV[1] then return redis.call("del", KEYS[1]) else return 0 end'; + +/** In-process fallback: key -> { token, expiresAt } */ +const localLocks = new Map(); + +export function resolveLockTtlMs(value, env = process.env) { + const raw = Number(value ?? env.PAYMENT_SESSION_LOCK_TTL_MS); + if (!Number.isFinite(raw)) return DEFAULT_LOCK_TTL_MS; + return Math.min(Math.max(Math.trunc(raw), MIN_LOCK_TTL_MS), MAX_LOCK_TTL_MS); +} + +/** + * Build a collision-free, injection-safe lock key. + * + * @param {string} merchantId + * @param {string} idempotencyKey + */ +export function buildSessionLockKey(merchantId, idempotencyKey) { + if (!merchantId || typeof merchantId !== "string") { + throw new TypeError("merchantId is required to build a session lock key"); + } + if (!idempotencyKey || typeof idempotencyKey !== "string") { + throw new TypeError("idempotencyKey is required to build a session lock key"); + } + const digest = createHash("sha256").update(idempotencyKey).digest("hex"); + return `${LOCK_KEY_PREFIX}:${merchantId}:${digest}`; +} + +function acquireLocalLock(key, token, ttlMs, now = Date.now()) { + const existing = localLocks.get(key); + if (existing && existing.expiresAt > now) { + return false; + } + localLocks.set(key, { token, expiresAt: now + ttlMs }); + return true; +} + +function releaseLocalLock(key, token) { + const existing = localLocks.get(key); + if (existing && existing.token === token) { + localLocks.delete(key); + return true; + } + return false; +} + +async function resolveRedis(getRedis) { + try { + const client = await getRedis(); + // connectRedisClient() returns a no-op client (isOpen:false) when Redis is + // unreachable. Its SET would report success for everyone, so treat it as + // "no distributed backend" and use the local lock instead. + if (!client || client.isOpen === false || typeof client.sendCommand !== "function") { + return null; + } + return client; + } catch (err) { + logger.warn({ err: err?.message }, "Payment session lock: Redis unavailable, using local lock"); + return null; + } +} + +/** + * Try to acquire the session lock without waiting. + * + * @param {string} key From buildSessionLockKey() + * @param {object} [options] + * @param {number} [options.ttlMs] + * @param {() => Promise} [options.getRedis] Injected for tests + * @returns {Promise Promise }>} + * `null` when another request currently holds the lock. + */ +export async function acquireSessionLock(key, { ttlMs, getRedis = connectRedisClient } = {}) { + const ttl = resolveLockTtlMs(ttlMs); + const token = randomUUID(); + const client = await resolveRedis(getRedis); + + if (client) { + try { + const reply = await client.sendCommand(["SET", key, token, "NX", "PX", String(ttl)]); + if (reply !== "OK") { + return null; + } + return { + key, + token, + backend: "redis", + ttlMs: ttl, + release: async () => { + try { + const released = await client.sendCommand(["EVAL", RELEASE_LOCK_SCRIPT, "1", key, token]); + if (Number(released) !== 1) { + logger.warn( + { key, ttlMs: ttl }, + "Payment session lock expired before release; operation outlived lock TTL", + ); + return false; + } + return true; + } catch (err) { + // The lock will still expire via PX, so a failed release is not fatal. + logger.error({ err: err?.message, key }, "Payment session lock release failed"); + return false; + } + }, + }; + } catch (err) { + logger.warn( + { err: err?.message, key }, + "Payment session lock: Redis SET NX failed, falling back to local lock", + ); + } + } + + if (!acquireLocalLock(key, token, ttl)) { + return null; + } + return { + key, + token, + backend: "local", + ttlMs: ttl, + release: async () => releaseLocalLock(key, token), + }; +} + +/** + * Build the error returned when a concurrent request holds the lock. + */ +export function createLockConflictError() { + const error = new Error( + "A request with this Idempotency-Key is already being processed. Retry shortly.", + ); + error.status = 409; + error.code = PAYMENT_SESSION_IN_PROGRESS; + error.retryable = false; + return error; +} + +/** + * Run `fn` while holding the (merchant, Idempotency-Key) session lock. + * + * Requests without an Idempotency-Key are NOT serialized: two independent + * sessions with identical bodies are legitimate, and there is nothing to + * de-duplicate against. + * + * @template T + * @param {{ merchantId:string, idempotencyKey?:string|null }} scope + * @param {(lock: object|null) => Promise} fn + * @param {object} [options] Forwarded to acquireSessionLock + * @returns {Promise} + * @throws 409 error (code PAYMENT_SESSION_IN_PROGRESS) when the lock is held + */ +export async function withPaymentSessionLock(scope, fn, options = {}) { + const { merchantId, idempotencyKey } = scope ?? {}; + if (!idempotencyKey || !merchantId) { + return fn(null); + } + + const key = buildSessionLockKey(merchantId, idempotencyKey); + const lock = await acquireSessionLock(key, options); + if (!lock) { + logger.warn({ merchantId }, "Concurrent payment session request rejected: lock held"); + throw createLockConflictError(); + } + + const startedAt = Date.now(); + try { + return await fn(lock); + } finally { + const heldMs = Date.now() - startedAt; + await lock.release(); + logger.debug?.({ merchantId, backend: lock.backend, heldMs }, "Payment session lock released"); + } +} + +/** Test helper: clear the in-process fallback lock table. */ +export function resetLocalSessionLocksForTests() { + localLocks.clear(); +} diff --git a/backend/src/lib/payment-session-lock.test.js b/backend/src/lib/payment-session-lock.test.js new file mode 100644 index 00000000..d12e5698 --- /dev/null +++ b/backend/src/lib/payment-session-lock.test.js @@ -0,0 +1,260 @@ +import { describe, it, expect, vi, beforeEach } from "vitest"; + +vi.mock("./logger.js", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +vi.mock("./redis.js", () => ({ + connectRedisClient: vi.fn(async () => ({ isOpen: false })), +})); + +import { logger } from "./logger.js"; +import { createFakeRedis } from "../../tests/helpers/fake-redis.js"; +import { + DEFAULT_LOCK_TTL_MS, + PAYMENT_SESSION_IN_PROGRESS, + RELEASE_LOCK_SCRIPT, + acquireSessionLock, + buildSessionLockKey, + resetLocalSessionLocksForTests, + resolveLockTtlMs, + withPaymentSessionLock, +} from "./payment-session-lock.js"; + +beforeEach(() => { + vi.clearAllMocks(); + resetLocalSessionLocksForTests(); +}); + +describe("buildSessionLockKey (issue #1450)", () => { + it("namespaces by merchant and hashes the idempotency key", () => { + const key = buildSessionLockKey("merchant-1", "order-42"); + expect(key).toMatch(/^lock:payment-session:merchant-1:[a-f0-9]{64}$/); + expect(key).not.toContain("order-42"); + }); + + it("is deterministic", () => { + expect(buildSessionLockKey("m", "k")).toBe(buildSessionLockKey("m", "k")); + }); + + it("isolates merchants using the same idempotency key", () => { + expect(buildSessionLockKey("m1", "k")).not.toBe(buildSessionLockKey("m2", "k")); + }); + + it("neutralises key-injection attempts and bounds key size", () => { + const hostile = "x:*:lock:payment-session:other-merchant\r\n" + "A".repeat(100_000); + const key = buildSessionLockKey("m1", hostile); + expect(key.length).toBeLessThan(120); + expect(key).not.toMatch(/[\r\n*]/); + expect(key.startsWith("lock:payment-session:m1:")).toBe(true); + }); + + it("rejects missing inputs", () => { + expect(() => buildSessionLockKey("", "k")).toThrow(TypeError); + expect(() => buildSessionLockKey("m", "")).toThrow(TypeError); + expect(() => buildSessionLockKey("m", null)).toThrow(TypeError); + }); +}); + +describe("resolveLockTtlMs", () => { + it("defaults to 30s", () => { + expect(resolveLockTtlMs(undefined, {})).toBe(DEFAULT_LOCK_TTL_MS); + }); + it("reads env and clamps to [1s, 120s]", () => { + expect(resolveLockTtlMs(undefined, { PAYMENT_SESSION_LOCK_TTL_MS: "5000" })).toBe(5000); + expect(resolveLockTtlMs(10, {})).toBe(1000); + expect(resolveLockTtlMs(10_000_000, {})).toBe(120_000); + expect(resolveLockTtlMs("nope", {})).toBe(DEFAULT_LOCK_TTL_MS); + }); +}); + +describe("acquireSessionLock — Redis backend", () => { + it("acquires with SET NX PX and a unique token", async () => { + const redis = createFakeRedis(); + const lock = await acquireSessionLock("k1", { getRedis: async () => redis, ttlMs: 5000 }); + expect(lock).toMatchObject({ key: "k1", backend: "redis", ttlMs: 5000 }); + expect(lock.token).toMatch(/[0-9a-f-]{36}/); + expect(redis.store.get("k1").value).toBe(lock.token); + }); + + it("refuses a second holder while the lock is live", async () => { + const redis = createFakeRedis(); + const getRedis = async () => redis; + const first = await acquireSessionLock("k1", { getRedis }); + const second = await acquireSessionLock("k1", { getRedis }); + expect(first).not.toBeNull(); + expect(second).toBeNull(); + }); + + it("releases via compare-and-delete", async () => { + const redis = createFakeRedis(); + const send = vi.spyOn(redis, "sendCommand"); + const lock = await acquireSessionLock("k1", { getRedis: async () => redis }); + await expect(lock.release()).resolves.toBe(true); + expect(redis.store.has("k1")).toBe(false); + expect(send).toHaveBeenLastCalledWith(["EVAL", RELEASE_LOCK_SCRIPT, "1", "k1", lock.token]); + }); + + it("never deletes a lock that expired and was re-acquired by someone else", async () => { + const redis = createFakeRedis(); + const getRedis = async () => redis; + const stale = await acquireSessionLock("k1", { getRedis, ttlMs: 1000 }); + // Simulate expiry, then another request acquiring the key. + redis.store.delete("k1"); + const fresh = await acquireSessionLock("k1", { getRedis }); + + await expect(stale.release()).resolves.toBe(false); + expect(redis.store.get("k1").value).toBe(fresh.token); + expect(logger.warn).toHaveBeenCalledWith( + expect.objectContaining({ key: "k1" }), + expect.stringMatching(/expired before release/), + ); + }); + + it("lets the lock expire naturally via PX", async () => { + vi.useFakeTimers(); + try { + const redis = createFakeRedis(); + const getRedis = async () => redis; + await acquireSessionLock("k1", { getRedis, ttlMs: 1000 }); + expect(await acquireSessionLock("k1", { getRedis })).toBeNull(); + vi.advanceTimersByTime(1001); + expect(await acquireSessionLock("k1", { getRedis })).not.toBeNull(); + } finally { + vi.useRealTimers(); + } + }); + + it("treats a failed release as non-fatal", async () => { + const redis = createFakeRedis(); + const lock = await acquireSessionLock("k1", { getRedis: async () => redis }); + redis.sendCommand = vi.fn().mockRejectedValue(new Error("conn lost")); + await expect(lock.release()).resolves.toBe(false); + expect(logger.error).toHaveBeenCalled(); + }); +}); + +describe("acquireSessionLock — local fallback", () => { + it("uses the local lock when Redis is the no-op client (isOpen:false)", async () => { + const lock = await acquireSessionLock("k1", { getRedis: async () => ({ isOpen: false }) }); + expect(lock.backend).toBe("local"); + // No-op Redis would say "OK" to everyone; local lock must still exclude. + expect(await acquireSessionLock("k1", { getRedis: async () => ({ isOpen: false }) })).toBeNull(); + }); + + it("uses the local lock when connecting to Redis throws", async () => { + const lock = await acquireSessionLock("k1", { + getRedis: async () => { + throw new Error("ECONNREFUSED"); + }, + }); + expect(lock.backend).toBe("local"); + }); + + it("uses the local lock when SET NX itself errors", async () => { + const redis = createFakeRedis(); + redis.sendCommand = vi.fn().mockRejectedValue(new Error("READONLY")); + const lock = await acquireSessionLock("k1", { getRedis: async () => redis }); + expect(lock.backend).toBe("local"); + }); + + it("uses the default connectRedisClient when none is injected", async () => { + const lock = await acquireSessionLock("k-default"); + expect(lock.backend).toBe("local"); + }); + + it("releases only the owner's local lock", async () => { + const getRedis = async () => null; + const lock = await acquireSessionLock("k1", { getRedis }); + await expect(lock.release()).resolves.toBe(true); + await expect(lock.release()).resolves.toBe(false); + expect(await acquireSessionLock("k1", { getRedis })).not.toBeNull(); + }); + + it("expires local locks after the TTL", async () => { + vi.useFakeTimers(); + try { + const getRedis = async () => null; + await acquireSessionLock("k1", { getRedis, ttlMs: 1000 }); + expect(await acquireSessionLock("k1", { getRedis })).toBeNull(); + vi.advanceTimersByTime(1001); + expect(await acquireSessionLock("k1", { getRedis })).not.toBeNull(); + } finally { + vi.useRealTimers(); + } + }); +}); + +describe("withPaymentSessionLock", () => { + it("runs without locking when there is no Idempotency-Key", async () => { + const getRedis = vi.fn(); + const fn = vi.fn().mockResolvedValue("done"); + await expect( + withPaymentSessionLock({ merchantId: "m1", idempotencyKey: null }, fn, { getRedis }), + ).resolves.toBe("done"); + expect(fn).toHaveBeenCalledWith(null); + expect(getRedis).not.toHaveBeenCalled(); + }); + + it("runs without locking when there is no merchant", async () => { + const fn = vi.fn().mockResolvedValue("done"); + await withPaymentSessionLock({ idempotencyKey: "k" }, fn); + expect(fn).toHaveBeenCalledWith(null); + }); + + it("holds the lock during fn and releases it afterwards", async () => { + const redis = createFakeRedis(); + const getRedis = async () => redis; + let heldDuring = false; + await withPaymentSessionLock( + { merchantId: "m1", idempotencyKey: "k" }, + async (lock) => { + heldDuring = redis.store.has(lock.key); + }, + { getRedis }, + ); + expect(heldDuring).toBe(true); + expect(redis.keys("lock:")).toEqual([]); + }); + + it("releases the lock even when fn throws", async () => { + const redis = createFakeRedis(); + const boom = new Error("handler failed"); + await expect( + withPaymentSessionLock( + { merchantId: "m1", idempotencyKey: "k" }, + async () => { + throw boom; + }, + { getRedis: async () => redis }, + ), + ).rejects.toBe(boom); + expect(redis.keys("lock:")).toEqual([]); + }); + + it("rejects a concurrent twin with a 409 PAYMENT_SESSION_IN_PROGRESS error", async () => { + const redis = createFakeRedis(); + const getRedis = async () => redis; + let releaseFirst; + const first = withPaymentSessionLock( + { merchantId: "m1", idempotencyKey: "k" }, + () => new Promise((resolve) => (releaseFirst = resolve)), + { getRedis }, + ); + await vi.waitFor(() => expect(releaseFirst).toBeTypeOf("function")); + + const twin = withPaymentSessionLock( + { merchantId: "m1", idempotencyKey: "k" }, + vi.fn(), + { getRedis }, + ); + await expect(twin).rejects.toMatchObject({ + status: 409, + code: PAYMENT_SESSION_IN_PROGRESS, + retryable: false, + }); + + releaseFirst("ok"); + await expect(first).resolves.toBe("ok"); + }); +}); diff --git a/backend/src/lib/payment-session-retry.js b/backend/src/lib/payment-session-retry.js new file mode 100644 index 00000000..6444a674 --- /dev/null +++ b/backend/src/lib/payment-session-retry.js @@ -0,0 +1,274 @@ +/** + * payment-session-retry.js + * + * Automated retry with exponential backoff for the Payment Session Validator + * (issue #1449). + * + * Session creation performs two kinds of I/O that can fail transiently: + * + * - on-chain issuer verification (Horizon) + * - persisting the session row (Supabase / Postgres) + * + * A single dropped connection or 503 used to fail the whole request. This + * module retries ONLY failures that are safe and meaningful to retry + * (network errors, 5xx, 429, connection-class Postgres errors). Validation + * failures, auth failures and constraint violations are never retried, so a + * retry can never turn a rejected session into an accepted one. + * + * Delay schedule: "full jitter" exponential backoff + * delay(n) = random(0, min(maxDelayMs, baseDelayMs * 2^n)) + * which spreads retries from many instances and avoids a thundering herd + * against a recovering dependency. + */ + +import { logger } from "./logger.js"; + +export const DEFAULT_RETRY_OPTIONS = Object.freeze({ + maxAttempts: 3, + baseDelayMs: 100, + maxDelayMs: 2000, +}); + +/** Hard ceilings so env/config can never make a request hang indefinitely. */ +const MAX_ATTEMPTS_CEILING = 6; +const MAX_DELAY_CEILING_MS = 10_000; + +const RETRYABLE_NETWORK_CODES = new Set([ + "ECONNRESET", + "ECONNREFUSED", + "ECONNABORTED", + "ETIMEDOUT", + "EPIPE", + "EAI_AGAIN", + "ENOTFOUND", + "EHOSTUNREACH", + "ENETUNREACH", + "UND_ERR_SOCKET", + "UND_ERR_CONNECT_TIMEOUT", + "UND_ERR_HEADERS_TIMEOUT", +]); + +/** + * Postgres SQLSTATE classes that indicate a transient server-side condition: + * 08xxx connection exception, 40001 serialization failure, + * 40P01 deadlock, 53xxx insufficient resources, 57P01-03 shutdown/cannot connect. + */ +function isRetryablePgCode(code) { + if (typeof code !== "string") return false; + return ( + code.startsWith("08") || + code === "40001" || + code === "40P01" || + code.startsWith("53") || + code === "57P01" || + code === "57P02" || + code === "57P03" + ); +} + +/** Postgres unique_violation. */ +export const UNIQUE_VIOLATION = "23505"; + +function getStatus(err) { + const status = err?.status ?? err?.statusCode ?? err?.response?.status; + return typeof status === "number" ? status : null; +} + +/** + * Decide whether an error is transient and safe to retry. + * + * @param {unknown} err + * @returns {boolean} + */ +export function isRetryableSessionError(err) { + if (!err || typeof err !== "object") return false; + + // Explicit opt-out wins over every heuristic below. + if (err.retryable === false) return false; + if (err.retryable === true) return true; + + const status = getStatus(err); + if (status !== null) { + if (status === 408 || status === 429) return true; + if (status >= 500 && status !== 501) return true; + // Every other status (4xx validation/auth/not-found) is deterministic. + if (status >= 400) return false; + } + + if (RETRYABLE_NETWORK_CODES.has(err.code)) return true; + if (isRetryablePgCode(err.code)) return true; + + const message = typeof err.message === "string" ? err.message : ""; + if (/fetch failed|socket hang up|network error|timed? ?out/i.test(message)) { + return true; + } + + return false; +} + +function clampInt(value, fallback, min, max) { + const n = Number(value); + if (!Number.isFinite(n)) return fallback; + return Math.min(Math.max(Math.trunc(n), min), max); +} + +/** + * Normalize user/env supplied options into safe bounds. + */ +export function resolveRetryOptions(options = {}, env = process.env) { + const maxAttempts = clampInt( + options.maxAttempts ?? env.PAYMENT_SESSION_RETRY_MAX_ATTEMPTS, + DEFAULT_RETRY_OPTIONS.maxAttempts, + 1, + MAX_ATTEMPTS_CEILING, + ); + const baseDelayMs = clampInt( + options.baseDelayMs ?? env.PAYMENT_SESSION_RETRY_BASE_DELAY_MS, + DEFAULT_RETRY_OPTIONS.baseDelayMs, + 0, + MAX_DELAY_CEILING_MS, + ); + const maxDelayMs = clampInt( + options.maxDelayMs ?? env.PAYMENT_SESSION_RETRY_MAX_DELAY_MS, + DEFAULT_RETRY_OPTIONS.maxDelayMs, + baseDelayMs, + MAX_DELAY_CEILING_MS, + ); + return { maxAttempts, baseDelayMs, maxDelayMs }; +} + +/** + * Full-jitter exponential backoff delay for a zero-based retry index. + * + * @param {number} retryIndex 0 for the first retry, 1 for the second, ... + * @param {{baseDelayMs:number, maxDelayMs:number}} opts + * @param {() => number} [random=Math.random] + */ +export function computeBackoffDelay(retryIndex, { baseDelayMs, maxDelayMs }, random = Math.random) { + const ceiling = Math.min(maxDelayMs, baseDelayMs * 2 ** retryIndex); + return Math.floor(random() * ceiling); +} + +const defaultSleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); + +/** + * Run `operation` with retry + exponential backoff. + * + * `operation` receives `{ attempt }` (1-based) so callers can adapt, e.g. + * treat a duplicate-key error on a retried insert as "already persisted". + * + * On exhaustion the LAST error is rethrown, annotated with `retryAttempts`, + * so upstream error handling and status codes are unchanged. + * + * @template T + * @param {(ctx:{attempt:number}) => Promise} operation + * @param {object} [options] + * @param {string} [options.label] Log label for the operation + * @param {number} [options.maxAttempts] Total attempts, including the first + * @param {number} [options.baseDelayMs] + * @param {number} [options.maxDelayMs] + * @param {(err:unknown) => boolean} [options.shouldRetry] + * @param {(ms:number) => Promise} [options.sleep] Injected for tests + * @param {() => number} [options.random] Injected for tests + * @param {(info:object) => void} [options.onRetry] Hook for metrics + * @param {object} [options.context] Extra structured log fields + * @returns {Promise} + */ +export async function withSessionRetry(operation, options = {}) { + const { + label = "payment-session-operation", + shouldRetry = isRetryableSessionError, + sleep = defaultSleep, + random = Math.random, + onRetry, + context = {}, + } = options; + const resolved = resolveRetryOptions(options); + + let lastError; + for (let attempt = 1; attempt <= resolved.maxAttempts; attempt += 1) { + try { + return await operation({ attempt }); + } catch (err) { + lastError = err; + const retryable = shouldRetry(err); + const exhausted = attempt >= resolved.maxAttempts; + + if (!retryable || exhausted) { + if (err && typeof err === "object") { + err.retryAttempts = attempt; + } + if (retryable && exhausted) { + logger.error( + { ...context, label, attempts: attempt, err: err?.message, code: err?.code }, + "Payment session operation failed after exhausting retries", + ); + } + throw err; + } + + const delayMs = computeBackoffDelay(attempt - 1, resolved, random); + logger.warn( + { + ...context, + label, + attempt, + maxAttempts: resolved.maxAttempts, + delayMs, + err: err?.message, + code: err?.code, + status: getStatus(err), + }, + "Payment session operation failed, retrying with backoff", + ); + onRetry?.({ label, attempt, delayMs, error: err }); + await sleep(delayMs); + } + } + + // Unreachable: the loop either returns or throws. + throw lastError; +} + +/** + * Insert a payment session row with retry. + * + * The session id is generated by the server BEFORE the first attempt, so the + * insert is naturally idempotent: if attempt N committed but its response was + * lost, attempt N+1 hits a unique_violation on the primary key. That case is + * treated as success ONLY on a retry, never on the first attempt, so a real + * duplicate is still surfaced. + * + * @param {object} supabase Supabase client + * @param {object} payload Row to insert (must include server-generated `id`) + * @param {object} [options] Forwarded to withSessionRetry + * @returns {Promise<{ error: null, recoveredDuplicate: boolean }>} + * @throws the Supabase error (non-retryable or exhausted) + */ +export async function insertPaymentSessionWithRetry(supabase, payload, options = {}) { + let recoveredDuplicate = false; + + await withSessionRetry( + async ({ attempt }) => { + const { error } = await supabase.from("payments").insert(payload); + if (!error) return; + + if (attempt > 1 && error.code === UNIQUE_VIOLATION) { + recoveredDuplicate = true; + logger.warn( + { paymentId: payload.id, attempt }, + "Payment session already persisted by an earlier attempt; treating retry as success", + ); + return; + } + throw error; + }, + { + label: "payment-session-insert", + ...options, + context: { paymentId: payload.id, merchantId: payload.merchant_id, ...options.context }, + }, + ); + + return { error: null, recoveredDuplicate }; +} diff --git a/backend/src/lib/payment-session-retry.test.js b/backend/src/lib/payment-session-retry.test.js new file mode 100644 index 00000000..872782b3 --- /dev/null +++ b/backend/src/lib/payment-session-retry.test.js @@ -0,0 +1,304 @@ +import { describe, it, expect, vi, beforeEach } from "vitest"; + +vi.mock("./logger.js", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +import { logger } from "./logger.js"; +import { + DEFAULT_RETRY_OPTIONS, + UNIQUE_VIOLATION, + computeBackoffDelay, + insertPaymentSessionWithRetry, + isRetryableSessionError, + resolveRetryOptions, + withSessionRetry, +} from "./payment-session-retry.js"; + +const noSleep = vi.fn(async () => {}); + +function errWith(props) { + return Object.assign(new Error(props.message ?? "boom"), props); +} + +beforeEach(() => { + vi.clearAllMocks(); +}); + +describe("isRetryableSessionError (issue #1449)", () => { + it.each([ + ["HTTP 500", { status: 500 }], + ["HTTP 502", { status: 502 }], + ["HTTP 503", { status: 503 }], + ["HTTP 504", { statusCode: 504 }], + ["HTTP 429", { status: 429 }], + ["HTTP 408", { status: 408 }], + ["axios-style 503", { response: { status: 503 } }], + ["ECONNRESET", { code: "ECONNRESET" }], + ["ETIMEDOUT", { code: "ETIMEDOUT" }], + ["EAI_AGAIN", { code: "EAI_AGAIN" }], + ["undici socket", { code: "UND_ERR_SOCKET" }], + ["pg connection failure 08006", { code: "08006" }], + ["pg serialization failure", { code: "40001" }], + ["pg deadlock", { code: "40P01" }], + ["pg too many connections", { code: "53300" }], + ["pg admin shutdown", { code: "57P01" }], + ["supabase fetch failed", { message: "TypeError: fetch failed", code: "" }], + ["socket hang up", { message: "socket hang up" }], + ["explicit retryable flag", { retryable: true, status: 400 }], + ])("retries %s", (_label, props) => { + expect(isRetryableSessionError(errWith(props))).toBe(true); + }); + + it.each([ + ["HTTP 400 validation", { status: 400 }], + ["HTTP 401", { status: 401 }], + ["HTTP 403", { status: 403 }], + ["HTTP 404", { status: 404 }], + ["HTTP 409", { status: 409 }], + ["HTTP 422", { status: 422 }], + ["HTTP 501 not implemented", { status: 501 }], + ["pg unique violation", { code: UNIQUE_VIOLATION }], + ["pg check violation", { code: "23514" }], + ["pg undefined column", { code: "42703" }], + ["opt-out flag beats 503", { retryable: false, status: 503 }], + ["plain error", { message: "something unexpected" }], + ])("does NOT retry %s", (_label, props) => { + expect(isRetryableSessionError(errWith(props))).toBe(false); + }); + + it("does not retry non-object throwables", () => { + expect(isRetryableSessionError(null)).toBe(false); + expect(isRetryableSessionError(undefined)).toBe(false); + expect(isRetryableSessionError("ECONNRESET")).toBe(false); + }); +}); + +describe("resolveRetryOptions", () => { + it("uses defaults when nothing is configured", () => { + expect(resolveRetryOptions({}, {})).toEqual({ ...DEFAULT_RETRY_OPTIONS }); + }); + + it("reads env overrides", () => { + expect( + resolveRetryOptions( + {}, + { + PAYMENT_SESSION_RETRY_MAX_ATTEMPTS: "5", + PAYMENT_SESSION_RETRY_BASE_DELAY_MS: "50", + PAYMENT_SESSION_RETRY_MAX_DELAY_MS: "500", + }, + ), + ).toEqual({ maxAttempts: 5, baseDelayMs: 50, maxDelayMs: 500 }); + }); + + it("clamps hostile or nonsensical values to safe bounds", () => { + const opts = resolveRetryOptions( + { maxAttempts: 10_000, baseDelayMs: -5, maxDelayMs: 9e9 }, + {}, + ); + expect(opts.maxAttempts).toBe(6); + expect(opts.baseDelayMs).toBe(0); + expect(opts.maxDelayMs).toBe(10_000); + }); + + it("never lets maxDelay fall below baseDelay and forces at least one attempt", () => { + const opts = resolveRetryOptions({ maxAttempts: 0, baseDelayMs: 400, maxDelayMs: 10 }, {}); + expect(opts.maxAttempts).toBe(1); + expect(opts.maxDelayMs).toBe(400); + }); + + it("falls back to defaults for non-numeric input", () => { + expect(resolveRetryOptions({ maxAttempts: "abc" }, {}).maxAttempts).toBe(3); + }); +}); + +describe("computeBackoffDelay", () => { + const opts = { baseDelayMs: 100, maxDelayMs: 1000 }; + + it("grows the ceiling exponentially", () => { + const max = () => 0.999999; + expect(computeBackoffDelay(0, opts, max)).toBe(99); + expect(computeBackoffDelay(1, opts, max)).toBe(199); + expect(computeBackoffDelay(2, opts, max)).toBe(399); + expect(computeBackoffDelay(3, opts, max)).toBe(799); + }); + + it("caps at maxDelayMs", () => { + expect(computeBackoffDelay(10, opts, () => 0.999999)).toBe(999); + }); + + it("applies full jitter (can be zero)", () => { + expect(computeBackoffDelay(3, opts, () => 0)).toBe(0); + }); + + it("always stays within [0, ceiling] for random inputs", () => { + for (let i = 0; i < 500; i += 1) { + const retry = i % 8; + const d = computeBackoffDelay(retry, opts); + expect(d).toBeGreaterThanOrEqual(0); + expect(d).toBeLessThanOrEqual(Math.min(1000, 100 * 2 ** retry)); + } + }); +}); + +describe("withSessionRetry", () => { + it("returns immediately on success without sleeping", async () => { + const op = vi.fn().mockResolvedValue("ok"); + await expect(withSessionRetry(op, { sleep: noSleep })).resolves.toBe("ok"); + expect(op).toHaveBeenCalledTimes(1); + expect(op).toHaveBeenCalledWith({ attempt: 1 }); + expect(noSleep).not.toHaveBeenCalled(); + }); + + it("retries transient failures and succeeds", async () => { + const op = vi + .fn() + .mockRejectedValueOnce(errWith({ code: "ECONNRESET" })) + .mockRejectedValueOnce(errWith({ status: 503 })) + .mockResolvedValue("ok"); + const onRetry = vi.fn(); + + await expect( + withSessionRetry(op, { sleep: noSleep, random: () => 0.5, maxAttempts: 3, onRetry }), + ).resolves.toBe("ok"); + + expect(op).toHaveBeenCalledTimes(3); + expect(op.mock.calls.map(([ctx]) => ctx.attempt)).toEqual([1, 2, 3]); + expect(noSleep).toHaveBeenCalledTimes(2); + // base 100 → ceilings 100, 200; random 0.5 → 50, 100 + expect(noSleep.mock.calls.map(([ms]) => ms)).toEqual([50, 100]); + expect(onRetry).toHaveBeenCalledTimes(2); + expect(logger.warn).toHaveBeenCalledTimes(2); + }); + + it("does not retry deterministic failures and annotates the error", async () => { + const validation = errWith({ status: 400, message: "bad input" }); + const op = vi.fn().mockRejectedValue(validation); + + await expect(withSessionRetry(op, { sleep: noSleep })).rejects.toBe(validation); + expect(op).toHaveBeenCalledTimes(1); + expect(validation.retryAttempts).toBe(1); + expect(noSleep).not.toHaveBeenCalled(); + expect(logger.error).not.toHaveBeenCalled(); + }); + + it("rethrows the LAST error after exhausting attempts and logs it", async () => { + const first = errWith({ status: 503, message: "first" }); + const last = errWith({ status: 503, message: "last" }); + const op = vi + .fn() + .mockRejectedValueOnce(first) + .mockRejectedValueOnce(first) + .mockRejectedValueOnce(last); + + await expect(withSessionRetry(op, { sleep: noSleep, maxAttempts: 3 })).rejects.toBe(last); + expect(op).toHaveBeenCalledTimes(3); + expect(last.retryAttempts).toBe(3); + expect(logger.error).toHaveBeenCalledTimes(1); + }); + + it("honours a custom shouldRetry predicate", async () => { + const op = vi.fn().mockRejectedValue(errWith({ status: 503 })); + await expect( + withSessionRetry(op, { sleep: noSleep, shouldRetry: () => false }), + ).rejects.toThrow(); + expect(op).toHaveBeenCalledTimes(1); + }); + + it("includes context fields in retry logs without leaking payload data", async () => { + const op = vi + .fn() + .mockRejectedValueOnce(errWith({ status: 503 })) + .mockResolvedValue("ok"); + await withSessionRetry(op, { + sleep: noSleep, + label: "unit", + context: { paymentId: "p-1" }, + }); + const [fields] = logger.warn.mock.calls[0]; + expect(fields).toMatchObject({ label: "unit", paymentId: "p-1", attempt: 1, status: 503 }); + }); +}); + +describe("insertPaymentSessionWithRetry", () => { + function supabaseWith(insertImpl) { + const insert = vi.fn(insertImpl); + return { insert, client: { from: vi.fn(() => ({ insert })) } }; + } + const payload = { id: "pay-1", merchant_id: "m-1", amount: 10 }; + + it("inserts once when the first attempt succeeds", async () => { + const { insert, client } = supabaseWith(async () => ({ error: null })); + await expect( + insertPaymentSessionWithRetry(client, payload, { sleep: noSleep }), + ).resolves.toEqual({ error: null, recoveredDuplicate: false }); + expect(client.from).toHaveBeenCalledWith("payments"); + expect(insert).toHaveBeenCalledTimes(1); + expect(insert).toHaveBeenCalledWith(payload); + }); + + it("retries Supabase transport errors returned in { error }", async () => { + const { insert, client } = supabaseWith( + vi + .fn() + .mockResolvedValueOnce({ error: { message: "TypeError: fetch failed", code: "" } }) + .mockResolvedValueOnce({ error: { message: "upstream", status: 503 } }) + .mockResolvedValue({ error: null }), + ); + const res = await insertPaymentSessionWithRetry(client, payload, { sleep: noSleep }); + expect(res.recoveredDuplicate).toBe(false); + expect(insert).toHaveBeenCalledTimes(3); + }); + + it("treats a duplicate key on a RETRY as already persisted (lost ack)", async () => { + const { insert, client } = supabaseWith( + vi + .fn() + .mockResolvedValueOnce({ error: { message: "timeout", code: "ETIMEDOUT" } }) + .mockResolvedValueOnce({ error: { message: "dup", code: UNIQUE_VIOLATION } }), + ); + const res = await insertPaymentSessionWithRetry(client, payload, { sleep: noSleep }); + expect(res).toEqual({ error: null, recoveredDuplicate: true }); + expect(insert).toHaveBeenCalledTimes(2); + }); + + it("surfaces a duplicate key on the FIRST attempt (real conflict, not retried)", async () => { + const dup = { message: "dup", code: UNIQUE_VIOLATION }; + const { insert, client } = supabaseWith(async () => ({ error: dup })); + await expect( + insertPaymentSessionWithRetry(client, payload, { sleep: noSleep }), + ).rejects.toBe(dup); + expect(insert).toHaveBeenCalledTimes(1); + }); + + it("does not retry constraint/validation errors", async () => { + const check = { message: "violates check constraint", code: "23514" }; + const { insert, client } = supabaseWith(async () => ({ error: check })); + await expect( + insertPaymentSessionWithRetry(client, payload, { sleep: noSleep }), + ).rejects.toBe(check); + expect(insert).toHaveBeenCalledTimes(1); + }); + + it("gives up after maxAttempts on a persistent outage", async () => { + const down = { message: "Service Unavailable", status: 503 }; + const { insert, client } = supabaseWith(async () => ({ error: down })); + await expect( + insertPaymentSessionWithRetry(client, payload, { sleep: noSleep, maxAttempts: 4 }), + ).rejects.toBe(down); + expect(insert).toHaveBeenCalledTimes(4); + expect(down.retryAttempts).toBe(4); + }); + + it("retries thrown (not returned) network exceptions", async () => { + const { insert, client } = supabaseWith( + vi + .fn() + .mockRejectedValueOnce(errWith({ code: "ECONNREFUSED" })) + .mockResolvedValue({ error: null }), + ); + await insertPaymentSessionWithRetry(client, payload, { sleep: noSleep }); + expect(insert).toHaveBeenCalledTimes(2); + }); +}); diff --git a/backend/src/lib/request-schemas.js b/backend/src/lib/request-schemas.js index f1a9c37d..2b939102 100644 --- a/backend/src/lib/request-schemas.js +++ b/backend/src/lib/request-schemas.js @@ -7,6 +7,10 @@ import { validateMemo, } from "./stellar.js"; import { resolveAssetIssuer } from "../constants/assetConstants.js"; +import { + merchantSettingsSchema, + customHeadersSchema, +} from "./merchant-payload-validation.js"; const VALID_MEMO_TYPES = ["text", "id", "hash", "return"]; @@ -214,11 +218,8 @@ export const registerMerchantZodSchema = z.object({ logo_url: z.string().trim().optional(), }) .optional(), - merchant_settings: z - .object({ - send_success_emails: z.boolean().optional(), - }) - .optional(), + // Issue #1482: strict — unknown settings keys are rejected, not dropped. + merchant_settings: merchantSettingsSchema.optional(), metadata: z.record(z.string(), z.unknown()).optional(), }); @@ -252,8 +253,6 @@ export const paymentSessionZodSchema = paymentBaseSchema export const v2PaymentSessionSchema = paymentSessionZodSchema; -const SAFE_HEADER_NAME_RE = /^[a-zA-Z0-9\-_]+$/; - export const VALID_WEBHOOK_EVENTS = [ "payment.confirmed", "payment.failed", @@ -276,14 +275,8 @@ export const webhookSettingsSchema = z.object({ .refine((val) => val.startsWith("https://"), "webhook_url must use HTTPS") .optional(), ), - custom_headers: z - .record(z.string(), z.string().min(1, "Header value must not be empty")) - .refine( - (obj) => Object.keys(obj).every((k) => SAFE_HEADER_NAME_RE.test(k)), - "Header names must contain only alphanumeric characters, hyphens, or underscores", - ) - .optional() - .nullable(), + // Issue #1482: bounded count/length, no CR/LF, no reserved header names. + custom_headers: customHeadersSchema.optional().nullable(), subscribed_events: z .array( z.string().refine( @@ -293,9 +286,7 @@ export const webhookSettingsSchema = z.object({ ) .optional() .nullable(), -}); - - +}).strict(); /** * Helper to parse and validate payment body for session creation. diff --git a/backend/src/lib/webhooks.js b/backend/src/lib/webhooks.js index 1cd9e1d9..40cb9b58 100644 --- a/backend/src/lib/webhooks.js +++ b/backend/src/lib/webhooks.js @@ -8,6 +8,7 @@ import { webhookDispatchDuration, webhookDispatchRetriesTotal, } from "./metrics.js"; +import { RESERVED_WEBHOOK_HEADERS } from "./merchant-payload-validation.js"; let supabaseClientPromise; @@ -260,8 +261,8 @@ function scheduleRetries(url, payload, headers, paymentId) { * * Accepted: plain object whose keys are safe ASCII header names and whose * values are non-empty strings. - * Reserved system headers (Content-Type, User-Agent, PLUTO-Signature) are - * silently dropped to prevent merchants from overriding security controls. + * Reserved system headers (see RESERVED_WEBHOOK_HEADERS) and values containing + * CR/LF/NUL are silently dropped to prevent merchants from overriding security controls. * * @param {unknown} raw The value stored in merchants.webhook_custom_headers. * @returns {Record} A safe subset of the supplied headers. @@ -270,20 +271,15 @@ export function sanitizeCustomHeaders(raw) { if (!raw || typeof raw !== "object" || Array.isArray(raw)) return {}; const SAFE_HEADER_NAME = /^[a-zA-Z0-9\-_]+$/; - const RESERVED = new Set([ - "content-type", - "user-agent", - "pluto-signature", - "stellar-signature", - "pluto-timestamp", - "stellar-timestamp", - ]); const result = {}; for (const [key, value] of Object.entries(raw)) { if (!SAFE_HEADER_NAME.test(key)) continue; - if (RESERVED.has(key.toLowerCase())) continue; + if (RESERVED_WEBHOOK_HEADERS.has(key.toLowerCase())) continue; if (typeof value !== "string" || value.trim() === "") continue; + // Values persisted before issue #1482 were not CR/LF checked; never let + // one reach the outbound request (header injection). + if (/[\r\n\0]/.test(value)) continue; result[key] = value; } return result; diff --git a/backend/src/routes/merchant-settings-validation.test.js b/backend/src/routes/merchant-settings-validation.test.js new file mode 100644 index 00000000..272f3cc1 --- /dev/null +++ b/backend/src/routes/merchant-settings-validation.test.js @@ -0,0 +1,264 @@ +/** + * HTTP-level tests for payload sanitization & strict validation on the + * Merchant Settings & API Key Service routes (issue #1482). + */ +import express from "express"; +import request from "supertest"; +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const { db } = vi.hoisted(() => ({ + db: { updates: [], inserts: [], lookup: null }, +})); + +function builder(table) { + let lastPayload = null; + const b = { + select: () => b, + eq: () => b, + is: () => b, + update: (payload) => { + lastPayload = payload; + db.updates.push({ table, payload }); + return b; + }, + insert: (payload) => { + lastPayload = payload; + db.inserts.push({ table, payload }); + return b; + }, + maybeSingle: async () => ({ data: db.lookup, error: null }), + single: async () => ({ + data: { id: "merchant-1", metadata: {}, ...(lastPayload ?? {}) }, + error: null, + }), + then: (resolve) => resolve({ data: null, error: null }), + }; + return b; +} + +vi.mock("../lib/supabase.js", () => ({ + supabase: { from: vi.fn((table) => builder(table)) }, +})); + +vi.mock("../lib/auth.js", () => ({ + requireApiKeyAuth: () => (req, _res, next) => { + req.merchant = { id: "merchant-1", webhook_secret: "whsec_old" }; + next(); + }, + requireSessionAuth: () => (req, _res, next) => { + req.merchant = { id: "merchant-1" }; + next(); + }, + hashPassword: vi.fn(async () => "hashed"), +})); + +vi.mock("../lib/rate-limit.js", () => ({ + createMerchantSecurityActionRateLimit: () => (_req, _res, next) => next(), +})); + +vi.mock("../lib/sep10-auth.js", () => ({ + generateSessionToken: vi.fn(() => "session-token"), +})); + +vi.mock("../lib/api-usage.js", () => ({ getMerchantApiUsage: vi.fn() })); + +vi.mock("../lib/logger.js", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +import createMerchantsRouter from "./merchants.js"; + +const DAY = 24 * 60 * 60 * 1000; + +function buildApp() { + const app = express(); + app.use(express.json()); + app.use( + "/api", + createMerchantsRouter({ + merchantRegistrationRateLimit: (_req, _res, next) => next(), + merchantSecurityActionRateLimit: (_req, _res, next) => next(), + }), + ); + // Mirror app.js: errors carry .status + // eslint-disable-next-line no-unused-vars + app.use((err, _req, res, _next) => { + res.status(err.status || 500).json({ error: err.message }); + }); + return app; +} + +beforeEach(() => { + vi.clearAllMocks(); + db.updates = []; + db.inserts = []; + // Row returned by `.maybeSingle()` lookups (e.g. current API key). + db.lookup = { api_key: "sk_old" }; +}); + +describe("POST /api/merchants/rotate-api-key", () => { + it("rotates with a valid grace period", async () => { + const res = await request(buildApp()) + .post("/api/merchants/rotate-api-key") + .send({ grace_period_hours: 12 }); + expect(res.status).toBe(200); + expect(res.body.grace_period_hours).toBe(12); + expect(res.body.api_key).toMatch(/^sk_[a-f0-9]{48}$/); + }); + + it("accepts an empty body (default grace period)", async () => { + const res = await request(buildApp()).post("/api/merchants/rotate-api-key"); + expect(res.status).toBe(200); + expect(res.body.grace_period_hours).toBe(24); + }); + + it.each([ + ["string grace period", { grace_period_hours: "24" }], + ["negative grace period", { grace_period_hours: -5 }], + ["grace period over one week", { grace_period_hours: 500 }], + ["mass-assignment of api_key", { api_key: "sk_attacker_chosen" }], + ["unknown field", { grace: 1 }], + ])("returns 400 (not 500) for %s and never writes", async (_label, body) => { + const res = await request(buildApp()).post("/api/merchants/rotate-api-key").send(body); + expect(res.status).toBe(400); + expect(res.body.error).toBe("Validation failed"); + expect(db.updates).toHaveLength(0); + }); +}); + +describe("PUT /api/merchants/set-api-key-expiry", () => { + it("stores a normalized UTC expiry", async () => { + const future = new Date(Date.now() + 30 * DAY); + const withOffset = future.toISOString().replace("Z", "+00:00"); + const res = await request(buildApp()) + .put("/api/merchants/set-api-key-expiry") + .send({ expires_at: withOffset }); + expect(res.status).toBe(200); + expect(res.body.api_key_expires_at).toBe(future.toISOString()); + expect(db.updates[0].payload).toEqual({ api_key_expires_at: future.toISOString() }); + }); + + it.each([ + ["past expiry (self lock-out)", { expires_at: "2020-01-01T00:00:00Z" }], + ["expiry beyond 365 days", { expires_at: new Date(Date.now() + 400 * DAY).toISOString() }], + ["non-ISO string", { expires_at: "next tuesday" }], + ["missing expires_at", {}], + [ + "extra field", + { expires_at: new Date(Date.now() + DAY).toISOString(), api_key_old: "sk_x" }, + ], + ])("rejects %s with 400", async (_label, body) => { + const res = await request(buildApp()).put("/api/merchants/set-api-key-expiry").send(body); + expect(res.status).toBe(400); + expect(db.updates).toHaveLength(0); + }); +}); + +describe("POST /api/merchants/rotate-webhook-secret", () => { + it("rejects unknown fields", async () => { + const res = await request(buildApp()) + .post("/api/merchants/rotate-webhook-secret") + .send({ grace_period_hours: 1, webhook_secret: "whsec_mine" }); + expect(res.status).toBe(400); + expect(db.updates).toHaveLength(0); + }); + + it("rotates with a valid body", async () => { + const res = await request(buildApp()) + .post("/api/merchants/rotate-webhook-secret") + .send({ grace_period_hours: 1 }); + expect(res.status).toBe(200); + expect(res.body.grace_period_hours).toBe(1); + }); +}); + +describe("PUT /api/webhook-settings", () => { + it("accepts a valid URL with safe custom headers", async () => { + const res = await request(buildApp()) + .put("/api/webhook-settings") + .send({ + webhook_url: "https://hooks.example.com/pluto", + custom_headers: { "X-Tenant": "acme" }, + }); + expect(res.status).toBe(200); + expect(db.updates[0].payload.webhook_custom_headers).toEqual({ "X-Tenant": "acme" }); + }); + + it.each([ + ["CRLF header injection", { custom_headers: { "X-A": "v\r\nX-Evil: 1" } }], + ["signature header spoofing", { custom_headers: { "PLUTO-Signature": "forged" } }], + ["host override", { custom_headers: { Host: "internal.local" } }], + ["unknown top-level field", { webhook_url: "https://a.example.com", webhook_secret: "x" }], + ["plain http URL", { webhook_url: "http://a.example.com" }], + ])("rejects %s with 400", async (_label, body) => { + const res = await request(buildApp()).put("/api/webhook-settings").send(body); + expect(res.status).toBe(400); + expect(db.updates).toHaveLength(0); + }); + + it("rejects oversized payloads before they reach the database", async () => { + const res = await request(buildApp()) + .put("/api/webhook-settings") + .send({ webhook_url: "https://a.example.com/" + "x".repeat(5000) }); + expect(res.status).toBe(400); + expect(db.updates).toHaveLength(0); + }); +}); + +describe("POST /api/register-merchant", () => { + const base = { email: "owner@example.com", password: "correct horse battery" }; + + beforeEach(() => { + db.lookup = null; // no existing merchant with this email + }); + + it("persists only canonical merchant_settings", async () => { + const res = await request(buildApp()) + .post("/api/register-merchant") + .send({ ...base, merchant_settings: { send_success_emails: false } }); + expect(res.status).toBe(201); + expect(db.inserts[0].payload.merchant_settings).toEqual({ send_success_emails: false }); + }); + + it("rejects unknown merchant_settings keys", async () => { + const res = await request(buildApp()) + .post("/api/register-merchant") + .send({ ...base, merchant_settings: { send_success_emails: true, is_admin: true } }); + expect(res.status).toBe(400); + expect(db.inserts).toHaveLength(0); + }); + + it("strips prototype-pollution keys from metadata and does not pollute Object.prototype", async () => { + const raw = + '{"email":"owner@example.com","password":"correct horse battery",' + + '"metadata":{"industry":"retail","__proto__":{"isAdmin":true}},' + + '"__proto__":{"polluted":true}}'; + const res = await request(buildApp()) + .post("/api/register-merchant") + .set("Content-Type", "application/json") + .send(raw); + + expect(res.status).toBe(201); + expect(db.inserts[0].payload.metadata).toEqual({ industry: "retail" }); + expect({}.polluted).toBeUndefined(); + expect({}.isAdmin).toBeUndefined(); + }); + + it("strips control characters from business_name", async () => { + const res = await request(buildApp()) + .post("/api/register-merchant") + .send({ ...base, business_name: "Acme\u0000\u001b Ltd" }); + expect(res.status).toBe(201); + expect(db.inserts[0].payload.business_name).toBe("Acme Ltd"); + }); + + it("rejects deeply nested metadata", async () => { + let deep = { v: 1 }; + for (let i = 0; i < 10; i += 1) deep = { n: deep }; + const res = await request(buildApp()) + .post("/api/register-merchant") + .send({ ...base, metadata: deep }); + expect(res.status).toBe(400); + expect(db.inserts).toHaveLength(0); + }); +}); diff --git a/backend/src/routes/merchants.js b/backend/src/routes/merchants.js index 9b4710b4..0b5f8dc6 100644 --- a/backend/src/routes/merchants.js +++ b/backend/src/routes/merchants.js @@ -5,7 +5,6 @@ import { supabase } from "../lib/supabase.js"; import { requireApiKeyAuth, requireSessionAuth, hashPassword } from "../lib/auth.js"; import { generateSessionToken } from "../lib/sep10-auth.js"; import { getMerchantApiUsage } from "../lib/api-usage.js"; -import { z } from "zod"; import { validateRequest } from "../lib/validation.js"; import { createMerchantSecurityActionRateLimit } from "../lib/rate-limit.js"; import { @@ -17,6 +16,12 @@ import { } from "../lib/request-schemas.js"; import { merchantService } from "../services/merchantService.js"; import { resolveMerchantSettings } from "../lib/merchant-settings.js"; +import { + sanitizeMerchantPayload, + rotateApiKeySchema, + rotateWebhookSecretSchema, + setApiKeyExpirySchema, +} from "../lib/merchant-payload-validation.js"; import { renderReceiptEmail } from "../lib/email-templates.js"; import { createWebhookDomainVerificationState, @@ -35,14 +40,6 @@ const defaultMerchantRegistrationRateLimit = rateLimit({ const defaultMerchantSecurityActionRateLimit = createMerchantSecurityActionRateLimit(); -const rotateApiKeySchema = z.object({ - grace_period_hours: z.number().int().min(0).max(168).optional(), -}); - -const setApiKeyExpirySchema = z.object({ - expires_at: z.string().datetime({ offset: true }).or(z.string().datetime()), -}); - function createMerchantsRouter({ @@ -53,10 +50,6 @@ function createMerchantsRouter({ const DEFAULT_WEBHOOK_SECRET_ROTATION_GRACE_HOURS = 24; - const rotateWebhookSecretSchema = z.object({ - grace_period_hours: z.number().int().min(0).max(168).optional(), - }); - function resolveWebhookSecretRotationGraceHours(requestValue) { if (typeof requestValue === "number") { return requestValue; @@ -123,6 +116,7 @@ function createMerchantsRouter({ router.post( "/register-merchant", merchantRegistrationRateLimit, + sanitizeMerchantPayload, validateRequest({ body: registerMerchantZodSchema }), async (req, res, next) => { try { @@ -281,6 +275,7 @@ function createMerchantsRouter({ "/merchants/rotate-webhook-secret", requireApiKeyAuth({ requireSignature: true }), merchantSecurityActionRateLimit, + sanitizeMerchantPayload, validateRequest({ body: rotateWebhookSecretSchema }), async (req, res, next) => { try { @@ -345,6 +340,7 @@ function createMerchantsRouter({ router.put( "/merchant-branding", + sanitizeMerchantPayload, validateRequest({ body: sessionBrandingSchema }), async (req, res, next) => { try { @@ -397,6 +393,7 @@ function createMerchantsRouter({ router.post( "/preview-receipt", requireApiKeyAuth(), + sanitizeMerchantPayload, validateRequest({ body: sessionBrandingSchema }), async (req, res) => { try { @@ -504,15 +501,6 @@ function createMerchantsRouter({ }, ); - const paymentLimitsSchema = z - .record( - z.string().min(1), - z.object({ - min: z.number().positive().optional(), - max: z.number().positive().optional(), - }), - ) - .optional(); /** * @swagger * /api/webhook-settings: @@ -588,6 +576,7 @@ function createMerchantsRouter({ "/webhook-settings", requireApiKeyAuth(), merchantSecurityActionRateLimit, + sanitizeMerchantPayload, validateRequest({ body: webhookSettingsSchema }), async (req, res, next) => { try { @@ -794,9 +783,11 @@ function createMerchantsRouter({ "/merchants/rotate-api-key", requireApiKeyAuth({ requireSignature: true }), merchantSecurityActionRateLimit, + sanitizeMerchantPayload, + validateRequest({ body: rotateApiKeySchema }), async (req, res, next) => { try { - const body = rotateApiKeySchema.parse(req.body || {}); + const body = req.body; const result = await merchantService.rotateApiKey( req.merchant.id, body.grace_period_hours @@ -848,15 +839,17 @@ function createMerchantsRouter({ * expires_at: * type: string * format: date-time - * description: ISO 8601 datetime when the API key expires + * description: ISO 8601 datetime (with timezone) when the API key expires. Must be at least 1 minute in the future and within 365 days. Unknown body fields are rejected. */ router.put( "/merchants/set-api-key-expiry", requireApiKeyAuth({ requireSignature: true }), merchantSecurityActionRateLimit, + sanitizeMerchantPayload, + validateRequest({ body: setApiKeyExpirySchema }), async (req, res, next) => { try { - const body = setApiKeyExpirySchema.parse(req.body); + const body = req.body; const result = await merchantService.setApiKeyExpiry( req.merchant.id, body.expires_at diff --git a/backend/src/routes/payment-session-validator.integration.test.js b/backend/src/routes/payment-session-validator.integration.test.js new file mode 100644 index 00000000..de9cf74e --- /dev/null +++ b/backend/src/routes/payment-session-validator.integration.test.js @@ -0,0 +1,645 @@ +/** + * Integration & stress suite for the Payment Session Validator (issue #1451). + * + * Exercises the REAL session pipeline end-to-end over HTTP: + * + * idempotencyMiddleware → createSession → withPaymentSessionLock + * → payment-session-rules (issuer / limits / allowlist) + * → insertPaymentSessionWithRetry → 201 / 4xx / 5xx + * + * Only true process boundaries are faked: Supabase (an in-memory "payments" + * table with injectable faults) and Redis (tests/helpers/fake-redis.js, shared + * between app instances to model multiple API nodes). + */ +import express from "express"; +import request from "supertest"; +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; + +const { state } = vi.hoisted(() => ({ + state: { + redis: null, + db: null, + }, +})); + +vi.mock("../lib/redis.js", () => ({ + connectRedisClient: vi.fn(async () => state.redis), + getRedisClient: vi.fn(() => state.redis), + getCachedPayment: vi.fn(), + setCachedPayment: vi.fn(), + invalidatePaymentCache: vi.fn(), +})); + +vi.mock("../lib/supabase-client.js", () => ({ + getSupabaseClient: vi.fn(async () => state.db.client), +})); + +vi.mock("../lib/supabase.js", () => ({ + supabase: { from: vi.fn(() => state.db.client.from("payments")) }, +})); + +vi.mock("../lib/stellar.js", () => ({ + findMatchingPayment: vi.fn(), + findAnyRecentPayment: vi.fn(), + findStrictReceivePaths: vi.fn(), + getNetworkFeeStats: vi.fn(), + isValidStellarPublicKey: vi.fn( + (value) => typeof value === "string" && /^G[A-Z2-7]{55}$/.test(value), + ), + validateMemo: vi.fn(() => ({ valid: true })), + verifyTransactionSignature: vi.fn(), +})); + +vi.mock("../constants/assetConstants.js", async (importOriginal) => ({ + ...(await importOriginal()), + resolveAssetIssuer: vi.fn((_asset, issuer) => issuer || null), +})); + +vi.mock("../lib/create-payment-rate-limit.js", () => ({ + createCreatePaymentRateLimit: () => (_req, _res, next) => next(), +})); +vi.mock("../lib/rate-limit.js", () => ({ + createVerifyPaymentRateLimit: () => (_req, _res, next) => next(), +})); +vi.mock("../lib/recaptcha.js", () => ({ + recaptchaMiddleware: () => (_req, _res, next) => next(), +})); +// Request-shape validation is covered by request-schemas tests; here we +// target the business-rule validator, so bypass the Zod layer. +vi.mock("../lib/validation.js", () => ({ + validateRequest: () => (_req, _res, next) => next(), +})); +vi.mock("../lib/sanitize-metadata.js", () => ({ + sanitizeMetadataMiddleware: (_req, _res, next) => next(), +})); +vi.mock("../lib/logger.js", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); +vi.mock("../lib/webhooks.js", () => ({ + sendWebhook: vi.fn(), + isEventSubscribed: vi.fn(() => false), +})); +vi.mock("../lib/email.js", () => ({ sendReceiptEmail: vi.fn() })); +vi.mock("../lib/email-templates.js", () => ({ renderReceiptEmail: vi.fn(() => "") })); +vi.mock("../webhooks/resolver.js", () => ({ getPayloadForVersion: vi.fn(() => ({})) })); +vi.mock("../lib/stream-manager.js", () => ({ + streamManager: { notify: vi.fn(), addClient: vi.fn() }, +})); +vi.mock("../lib/metrics.js", () => ({ + paymentCreatedCounter: { inc: vi.fn() }, + paymentConfirmedCounter: { inc: vi.fn() }, + paymentConfirmationLatency: { observe: vi.fn() }, + paymentFailedCounter: { inc: vi.fn() }, + exchangeRateQuoteRequests: { inc: vi.fn() }, + exchangeRateQuoteDuration: { observe: vi.fn() }, + exchangeRateSlippageApplied: { inc: vi.fn() }, +})); +vi.mock("../services/paymentService.js", () => ({ paymentService: {} })); + +import createPaymentsRouter from "./payments.js"; +import { idempotencyMiddleware } from "../lib/idempotency.js"; +import { resetLocalSessionLocksForTests } from "../lib/payment-session-lock.js"; +import { createFakeRedis } from "../../tests/helpers/fake-redis.js"; + +const VALID_ISSUER = "GBBD47IF6LWK7P7MDEVSCWR7DPUWV3NY3DTQEVFL4NAT4AQH3ZLLFLA5"; +const OTHER_ISSUER = "GA5XIGA5C7FBPTVQ3CWHKNC7D2ZBHB24G3KUJG5WZ6S4EYWSSBFVL45T"; +const RECIPIENT = "GDRXE2BQUC3AZNPVFSCEZ76NJ3WWL25FYFK6RGZGIEKWE4SOOHSUJUJ6"; + +/** + * In-memory "payments" table. `faults` is a queue of results consumed by + * successive insert calls: "transient" | "outage" | "constraint" | + * "lost-ack" (row IS written but the client sees a timeout). + */ +function createFakeDb({ latencyMs = 0 } = {}) { + const rows = new Map(); + const faults = []; + const stats = { insertCalls: 0, inFlight: 0, maxInFlight: 0 }; + let randomFaultRate = 0; + + const tick = () => + latencyMs > 0 + ? new Promise((r) => setTimeout(r, Math.random() * latencyMs)) + : Promise.resolve(); + + async function insert(row) { + stats.insertCalls += 1; + stats.inFlight += 1; + stats.maxInFlight = Math.max(stats.maxInFlight, stats.inFlight); + try { + return await doInsert(row); + } finally { + stats.inFlight -= 1; + } + } + + async function doInsert(row) { + await tick(); + let fault = faults.shift(); + if (!fault && randomFaultRate > 0 && Math.random() < randomFaultRate) { + fault = Math.random() < 0.5 ? "transient" : "lost-ack"; + } + switch (fault) { + case "transient": + return { error: { message: "TypeError: fetch failed", code: "" } }; + case "outage": + return { error: { message: "Service Unavailable", status: 503 } }; + case "constraint": + return { error: { message: "violates check constraint", code: "23514" } }; + case "lost-ack": + if (!rows.has(row.id)) rows.set(row.id, structuredClone(row)); + return { error: { message: "timeout", code: "ETIMEDOUT" } }; + default: + break; + } + if (rows.has(row.id)) { + return { error: { message: "duplicate key value", code: "23505" } }; + } + rows.set(row.id, structuredClone(row)); + return { error: null }; + } + + return { + rows, + faults, + stats, + setRandomFaultRate(rate) { + randomFaultRate = rate; + }, + client: { from: () => ({ insert }) }, + }; +} + +const defaultMerchant = { + id: "merchant-1", + payment_limits: { USDC: { min: 1, max: 1000 } }, + allowed_issuers: [VALID_ISSUER], + branding_config: null, +}; + +const openServers = []; + +/** + * One "API node", listening on its own ephemeral port. Several nodes can + * share state.redis / state.db. A long-lived server per node (rather than + * supertest's server-per-request) keeps the stress tests about the code + * under test, not about socket churn in the test harness. + */ +async function buildApp() { + const server = createAppInstance().listen(0, "127.0.0.1"); + openServers.push(server); + await new Promise((resolve) => server.once("listening", resolve)); + return server; +} + +function createAppInstance() { + const app = express(); + app.use(express.json()); + app.use((req, _res, next) => { + req.merchant = { + ...defaultMerchant, + id: req.get("x-test-merchant") || defaultMerchant.id, + }; + next(); + }); + app.use("/api/sessions", idempotencyMiddleware); + app.use("/api/create-payment", idempotencyMiddleware); + app.use("/api", createPaymentsRouter()); + return app; +} + +/** Locks are released in `finally`, right after the response is flushed. */ +async function expectLocksDrained() { + await vi.waitFor(() => expect(state.redis.keys("lock:")).toEqual([]), { timeout: 5000 }); +} + +function sessionBody(overrides = {}) { + return { + amount: 25, + asset: "USDC", + asset_issuer: VALID_ISSUER, + recipient: RECIPIENT, + ...overrides, + }; +} + +beforeAll(() => { + process.env.PAYMENT_SESSION_RETRY_BASE_DELAY_MS = "1"; + process.env.PAYMENT_SESSION_RETRY_MAX_DELAY_MS = "5"; + process.env.PAYMENT_SESSION_RETRY_MAX_ATTEMPTS = "3"; +}); + +afterAll(() => { + delete process.env.PAYMENT_SESSION_RETRY_BASE_DELAY_MS; + delete process.env.PAYMENT_SESSION_RETRY_MAX_DELAY_MS; + delete process.env.PAYMENT_SESSION_RETRY_MAX_ATTEMPTS; +}); + +afterEach(async () => { + await Promise.all( + openServers.splice(0).map((server) => new Promise((resolve) => server.close(resolve))), + ); +}); + +beforeEach(() => { + vi.clearAllMocks(); + resetLocalSessionLocksForTests(); + state.redis = createFakeRedis(); + state.db = createFakeDb(); +}); + +// --------------------------------------------------------------------------- +// Validation rules over HTTP +// --------------------------------------------------------------------------- + +describe("Payment Session Validator — business rules (integration)", () => { + it("creates a session for a valid request and persists exactly one row", async () => { + const res = await request(await buildApp()).post("/api/sessions").send(sessionBody()); + + expect(res.status).toBe(201); + expect(res.body).toMatchObject({ status: "pending", sandbox: false }); + expect(res.body.payment_id).toMatch(/^[0-9a-f-]{36}$/); + expect(state.db.rows.size).toBe(1); + const row = state.db.rows.get(res.body.payment_id); + expect(row).toMatchObject({ + merchant_id: "merchant-1", + asset: "USDC", + asset_issuer: VALID_ISSUER, + amount: 25, + status: "pending", + }); + }); + + it("accepts native XLM without an issuer", async () => { + const res = await request(await buildApp()) + .post("/api/sessions") + .send(sessionBody({ asset: "XLM", asset_issuer: undefined, amount: 5 })); + expect(res.status).toBe(201); + }); + + it.each([ + [ + "missing issuer for non-native asset", + { asset_issuer: undefined }, + "asset_issuer is required for non-native assets", + ], + [ + "malformed issuer", + { asset_issuer: "not-a-stellar-key" }, + "asset_issuer must be a valid Stellar public key", + ], + [ + "issuer outside the merchant allowlist", + { asset_issuer: OTHER_ISSUER }, + "asset_issuer is not in the merchant's list of allowed issuers", + ], + ["amount below the per-asset minimum", { amount: 0.5 }, "Amount is below the minimum for USDC"], + ["amount above the per-asset maximum", { amount: 5000 }, "Amount exceeds the maximum for USDC"], + ])("rejects %s with 400 and never touches the database", async (_label, overrides, message) => { + const res = await request(await buildApp()).post("/api/sessions").send(sessionBody(overrides)); + expect(res.status).toBe(400); + expect(res.body.error).toBe(message); + expect(state.db.stats.insertCalls).toBe(0); + }); + + it("returns limit deltas for below-minimum rejections", async () => { + const res = await request(await buildApp()).post("/api/sessions").send(sessionBody({ amount: 0.25 })); + expect(res.body).toMatchObject({ min: 1, delta: 0.75 }); + }); + + it("applies identical rules on /create-payment", async () => { + const app = await buildApp(); + const ok = await request(app).post("/api/create-payment").send(sessionBody()); + const bad = await request(app) + .post("/api/create-payment") + .send(sessionBody({ asset_issuer: OTHER_ISSUER })); + expect(ok.status).toBe(201); + expect(bad.status).toBe(400); + }); +}); + +// --------------------------------------------------------------------------- +// Retry with exponential backoff (issue #1449) +// --------------------------------------------------------------------------- + +describe("Payment Session Validator — retry & backoff (integration)", () => { + it("recovers from transient persistence failures", async () => { + state.db.faults.push("transient", "outage"); + const res = await request(await buildApp()).post("/api/sessions").send(sessionBody()); + expect(res.status).toBe(201); + expect(state.db.stats.insertCalls).toBe(3); + expect(state.db.rows.size).toBe(1); + }); + + it("does not create a duplicate when an insert committed but its ack was lost", async () => { + state.db.faults.push("lost-ack"); + const res = await request(await buildApp()).post("/api/sessions").send(sessionBody()); + expect(res.status).toBe(201); + expect(state.db.stats.insertCalls).toBe(2); + expect(state.db.rows.size).toBe(1); + expect(state.db.rows.has(res.body.payment_id)).toBe(true); + }); + + it("returns 500 after exhausting retries on a sustained outage", async () => { + state.db.faults.push("outage", "outage", "outage", "outage"); + const res = await request(await buildApp()).post("/api/sessions").send(sessionBody()); + expect(res.status).toBe(500); + expect(state.db.stats.insertCalls).toBe(3); + expect(state.db.rows.size).toBe(0); + }); + + it("fails fast (no retry) on constraint violations", async () => { + state.db.faults.push("constraint"); + const res = await request(await buildApp()).post("/api/sessions").send(sessionBody()); + expect(res.status).toBe(500); + expect(state.db.stats.insertCalls).toBe(1); + }); + + it("does not cache failed responses under the Idempotency-Key", async () => { + state.db.faults.push("outage", "outage", "outage"); + const app = await buildApp(); + const failed = await request(app) + .post("/api/sessions") + .set("Idempotency-Key", "retry-after-500") + .send(sessionBody()); + expect(failed.status).toBe(500); + + const retried = await request(app) + .post("/api/sessions") + .set("Idempotency-Key", "retry-after-500") + .send(sessionBody()); + expect(retried.status).toBe(201); + expect(state.db.rows.size).toBe(1); + }); +}); + +// --------------------------------------------------------------------------- +// Distributed concurrency control (issue #1450) +// --------------------------------------------------------------------------- + +describe("Payment Session Validator — concurrency control (integration)", () => { + it("creates exactly ONE session for a burst of identical Idempotency-Key requests", async () => { + state.redis = createFakeRedis({ latencyMs: 3 }); + state.db = createFakeDb({ latencyMs: 10 }); + const app = await buildApp(); + + const responses = await Promise.all( + Array.from({ length: 25 }, () => + request(app).post("/api/sessions").set("Idempotency-Key", "burst-1").send(sessionBody()), + ), + ); + + const created = responses.filter((r) => r.status === 201); + const conflicts = responses.filter((r) => r.status === 409); + expect(created.length + conflicts.length).toBe(25); + expect(state.db.rows.size).toBe(1); + expect(state.db.stats.maxInFlight).toBe(1); + + // Every 201 (original or replay) must reference the same session. + const ids = new Set(created.map((r) => r.body.payment_id)); + expect(ids.size).toBe(1); + for (const c of conflicts) { + expect(c.body.code).toBe("PAYMENT_SESSION_IN_PROGRESS"); + } + }); + + it("serializes across multiple API nodes sharing Redis", async () => { + state.redis = createFakeRedis({ latencyMs: 3 }); + state.db = createFakeDb({ latencyMs: 10 }); + const nodes = [await buildApp(), await buildApp(), await buildApp()]; + + const responses = await Promise.all( + Array.from({ length: 30 }, (_, i) => + request(nodes[i % nodes.length]) + .post("/api/sessions") + .set("Idempotency-Key", "multi-node") + .send(sessionBody()), + ), + ); + + expect(state.db.rows.size).toBe(1); + expect(responses.every((r) => r.status === 201 || r.status === 409)).toBe(true); + }); + + it("replays the original response after the in-flight request completes", async () => { + const app = await buildApp(); + const first = await request(app) + .post("/api/sessions") + .set("Idempotency-Key", "replay-1") + .send(sessionBody()); + const second = await request(app) + .post("/api/sessions") + .set("Idempotency-Key", "replay-1") + .send(sessionBody()); + + expect(first.status).toBe(201); + expect(second.status).toBe(201); + expect(second.body.payment_id).toBe(first.body.payment_id); + expect(state.db.rows.size).toBe(1); + }); + + it("replays inside the lock when a twin passed the middleware before the cache was written", async () => { + // Model the race precisely: request B has already passed the idempotency + // middleware (cache miss) when A commits. B must replay, not re-create. + state.db = createFakeDb({ latencyMs: 25 }); + const app = await buildApp(); + const a = request(app).post("/api/sessions").set("Idempotency-Key", "race").send(sessionBody()); + const b = new Promise((resolve) => setTimeout(resolve, 5)).then(() => + request(app).post("/api/sessions").set("Idempotency-Key", "race").send(sessionBody()), + ); + const [ra, rb] = await Promise.all([a, b]); + + expect(ra.status).toBe(201); + expect([201, 409]).toContain(rb.status); + expect(state.db.rows.size).toBe(1); + }); + + it("rejects reuse of an Idempotency-Key with a different payload", async () => { + const app = await buildApp(); + await request(app).post("/api/sessions").set("Idempotency-Key", "k-mismatch").send(sessionBody()); + const res = await request(app) + .post("/api/sessions") + .set("Idempotency-Key", "k-mismatch") + .send(sessionBody({ amount: 999 })); + expect(res.status).toBe(400); + expect(state.db.rows.size).toBe(1); + }); + + it("does not let one merchant's Idempotency-Key block another merchant", async () => { + state.db = createFakeDb({ latencyMs: 10 }); + const app = await buildApp(); + const [m1, m2] = await Promise.all([ + request(app) + .post("/api/sessions") + .set("x-test-merchant", "merchant-A") + .set("Idempotency-Key", "shared") + .send(sessionBody()), + request(app) + .post("/api/sessions") + .set("x-test-merchant", "merchant-B") + .set("Idempotency-Key", "shared") + .send(sessionBody()), + ]); + expect(m1.status).toBe(201); + expect(m2.status).toBe(201); + expect(m1.body.payment_id).not.toBe(m2.body.payment_id); + expect(state.db.rows.size).toBe(2); + }); + + it("does NOT serialize requests without an Idempotency-Key", async () => { + state.db = createFakeDb({ latencyMs: 5 }); + const app = await buildApp(); + const responses = await Promise.all( + Array.from({ length: 10 }, () => request(app).post("/api/sessions").send(sessionBody())), + ); + expect(responses.every((r) => r.status === 201)).toBe(true); + expect(state.db.rows.size).toBe(10); + await expectLocksDrained(); + }); + + it("releases the lock after validation failures so the client can correct and retry", async () => { + const app = await buildApp(); + const bad = await request(app) + .post("/api/sessions") + .set("Idempotency-Key", "fix-and-retry") + .send(sessionBody({ amount: 99999 })); + expect(bad.status).toBe(400); + await expectLocksDrained(); + }); + + it("still enforces mutual exclusion when Redis is down (local fallback)", async () => { + // Degraded mode: without Redis there is no idempotency cache to replay + // from, so the guarantee is "never two in-flight creations for the same + // key" — overlapping twins get 409 instead of racing into the database. + state.redis = { isOpen: false, get: async () => null, set: async () => "OK" }; + state.db = createFakeDb({ latencyMs: 40 }); + const app = await buildApp(); + const responses = await Promise.all( + Array.from({ length: 10 }, () => + request(app).post("/api/sessions").set("Idempotency-Key", "no-redis").send(sessionBody()), + ), + ); + expect(responses.every((r) => r.status === 201 || r.status === 409)).toBe(true); + expect(responses.some((r) => r.status === 409)).toBe(true); + expect(state.db.stats.maxInFlight).toBe(1); + expect(state.db.rows.size).toBe(responses.filter((r) => r.status === 201).length); + }); + + it("hashes hostile Idempotency-Keys instead of embedding them in Redis keys", async () => { + const hostile = "x".repeat(2000) + ":lock:payment-session:merchant-2:*"; + const res = await request(await buildApp()) + .post("/api/sessions") + .set("Idempotency-Key", hostile) + .send(sessionBody()); + expect(res.status).toBe(201); + const lockCommands = state.redis.commandLog.filter( + ([cmd, key]) => cmd === "SET" && String(key).startsWith("lock:"), + ); + expect(lockCommands).toHaveLength(1); + expect(lockCommands[0][1]).toMatch(/^lock:payment-session:merchant-1:[a-f0-9]{64}$/); + }); +}); + +// --------------------------------------------------------------------------- +// Stress +// --------------------------------------------------------------------------- + +describe("Payment Session Validator — stress", { timeout: 30_000 }, () => { + it("handles 200 concurrent distinct sessions under a 30% fault rate without loss or duplication", async () => { + state.redis = createFakeRedis({ latencyMs: 2 }); + state.db = createFakeDb({ latencyMs: 4 }); + state.db.setRandomFaultRate(0.3); + process.env.PAYMENT_SESSION_RETRY_MAX_ATTEMPTS = "6"; + const nodes = [await buildApp(), await buildApp(), await buildApp(), await buildApp()]; + + try { + const started = Date.now(); + const responses = await Promise.all( + Array.from({ length: 200 }, (_, i) => + request(nodes[i % nodes.length]) + .post("/api/sessions") + .set("Idempotency-Key", `stress-${i}`) + .send(sessionBody({ amount: 1 + (i % 900) })), + ), + ); + const elapsed = Date.now() - started; + + const statuses = responses.map((r) => r.status); + const created = responses.filter((r) => r.status === 201); + // 0.3^6 ≈ 0.07% per request, so essentially all must succeed. + expect(created.length).toBeGreaterThanOrEqual(198); + expect(statuses.every((s) => s === 201 || s === 500)).toBe(true); + + const ids = new Set(created.map((r) => r.body.payment_id)); + expect(ids.size).toBe(created.length); + for (const id of ids) { + expect(state.db.rows.has(id)).toBe(true); + } + // Every persisted row belongs to a request (no orphans beyond lost-acks + // of requests that ultimately failed). + expect(state.db.rows.size).toBeLessThanOrEqual(200); + expect(state.db.rows.size).toBeGreaterThanOrEqual(created.length); + // All locks released (release runs just after the response flushes). + await expectLocksDrained(); + // Bounded latency: backoff is capped, so the whole burst stays fast. + expect(elapsed).toBeLessThan(10_000); + } finally { + process.env.PAYMENT_SESSION_RETRY_MAX_ATTEMPTS = "3"; + } + }); + + it("keeps exactly-once semantics for 20 keys × 10 concurrent duplicates each", async () => { + state.redis = createFakeRedis({ latencyMs: 2 }); + state.db = createFakeDb({ latencyMs: 5 }); + const nodes = [await buildApp(), await buildApp()]; + + const responses = await Promise.all( + Array.from({ length: 200 }, (_, i) => { + const key = `dup-${i % 20}`; + return request(nodes[i % 2]) + .post("/api/sessions") + .set("Idempotency-Key", key) + .send(sessionBody()) + .then((res) => ({ key, res })); + }), + ); + + expect(state.db.rows.size).toBe(20); + const byKey = new Map(); + for (const { key, res } of responses) { + expect([201, 409]).toContain(res.status); + if (res.status === 201) { + const ids = byKey.get(key) ?? new Set(); + ids.add(res.body.payment_id); + byKey.set(key, ids); + } + } + expect(byKey.size).toBe(20); + for (const ids of byKey.values()) { + expect(ids.size).toBe(1); + } + await expectLocksDrained(); + }); + + it("rejects a flood of invalid sessions without a single database write", async () => { + const app = await buildApp(); + const responses = await Promise.all( + Array.from({ length: 100 }, (_, i) => + request(app) + .post("/api/sessions") + .send( + sessionBody( + [ + { asset_issuer: undefined }, + { asset_issuer: "bad" }, + { asset_issuer: OTHER_ISSUER }, + { amount: 0.1 }, + { amount: 10_000 }, + ][i % 5], + ), + ), + ), + ); + expect(responses.every((r) => r.status === 400)).toBe(true); + expect(state.db.stats.insertCalls).toBe(0); + }); +}); diff --git a/backend/src/routes/payments.js b/backend/src/routes/payments.js index fe2c5744..a8a0c17e 100644 --- a/backend/src/routes/payments.js +++ b/backend/src/routes/payments.js @@ -42,6 +42,12 @@ import { import { sanitizeMetadataMiddleware } from "../lib/sanitize-metadata.js"; import { validatePaymentSession } from "../lib/payment-session-validator.js"; import { getSupabaseClient } from "../lib/supabase-client.js"; +import { insertPaymentSessionWithRetry } from "../lib/payment-session-retry.js"; +import { + withPaymentSessionLock, + PAYMENT_SESSION_IN_PROGRESS, +} from "../lib/payment-session-lock.js"; +import { lookupIdempotentResponse } from "../lib/idempotency.js"; import { paymentProcessorSessionsTotal, paymentProcessorSessionDuration, @@ -245,12 +251,52 @@ function createPaymentsRouter({ * type: string * 400: * description: Validation error or invalid Idempotency-Key + * 409: + * description: A concurrent request with the same Idempotency-Key is still being processed (code PAYMENT_SESSION_IN_PROGRESS). Retry shortly to receive the cached response. * 429: * description: Too many requests */ + /** + * Issue #1450: serialize concurrent requests that share an Idempotency-Key + * (per merchant, across instances) so they can never create two sessions. + * The loser of the race gets 409; a request that acquires the lock after + * its twin finished replays the cached idempotent response. + */ async function createSession(req, res, next) { - const sessionStart = Date.now(); + const idempotencyKey = req.idempotency?.key ?? req.get?.("Idempotency-Key") ?? null; try { + await withPaymentSessionLock( + { merchantId: req.merchant?.id, idempotencyKey }, + async (lock) => { + if (lock && req.idempotency?.payloadHash) { + const replay = await lookupIdempotentResponse({ + redisClient: await connectRedisClient(), + merchantId: req.merchant.id, + idempotencyKey, + payloadHash: req.idempotency.payloadHash, + }).catch((err) => { + logger.warn({ err: err?.message }, "Idempotency replay lookup failed; continuing"); + return null; + }); + if (replay?.status === "hit") { + return res.status(201).json(replay.response); + } + if (replay?.status === "mismatch") { + return res.status(400).json({ + error: "Idempotency-Key already used with a different request payload", + }); + } + } + return createSessionUnlocked(req, res); + }, + ); + } catch (err) { + if (err.code === PAYMENT_SESSION_IN_PROGRESS) { + return res.status(409).json({ error: err.message, code: err.code }); + } + logger.error({ err, merchantId: req.merchant?.id }, "DEBUG: createSession error"); + if (err.status === 400 && err.details) { + return res.status(400).json({ error: err.message, ...err.details }); const supabase = await getSupabaseClient(); logger.info({ merchantId: req.merchant?.id, amount: req.body?.amount, asset: req.body?.asset }, "DEBUG: createSession started"); @@ -278,86 +324,140 @@ function createPaymentsRouter({ ...(rejection.rule === "limits" ? rejection.details : {}), }); } + next(err); + } + } + async function createSessionUnlocked(req, res) { + const sessionStart = Date.now(); + const supabase = await getSupabaseClient(); + const body = req.body; + const asset = body.asset?.toUpperCase(); + logger.info({ merchantId: req.merchant?.id, amount: body.amount, asset: body.asset }, "DEBUG: createSession started"); + + // Shared business-rule validation (issue #1087) — issuer presence/format. + const { assetIssuer, rejection: issuerRejection } = resolveAndValidateIssuer( + asset, + body.asset_issuer, + ); + if (issuerRejection) { + paymentFailedCounter.inc({ asset: body.asset, reason: issuerRejection.reason }); + paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "validation_failed" }); + paymentProcessorSessionDuration.observe( + { asset: body.asset, outcome: "validation_failed" }, + (Date.now() - sessionStart) / 1000, + ); + return res.status(400).json({ error: issuerRejection.message }); + } const body = validation.payload; const { asset, assetIssuer } = validation; - const isSandbox = body.sandbox === true; - const baseId = randomUUID(); - const paymentId = isSandbox ? `test_${baseId}` : baseId; - const now = new Date().toISOString(); - const paymentLinkBase = - process.env.PAYMENT_LINK_BASE || "http://localhost:3000"; - const paymentLink = `${paymentLinkBase}/pay/${paymentId}`; - const resolvedBrandingConfig = resolveBrandingConfig({ - merchantBranding: req.merchant.branding_config, - brandingOverrides: body.branding_overrides, + // Shared business-rule validation (issue #1087) — per-asset limits (#153). + const limitRejection = validatePerAssetLimits({ + rawAsset: body.asset, + amount: body.amount, + paymentLimits: req.merchant.payment_limits, + }); + if (limitRejection) { + paymentFailedCounter.inc({ asset: body.asset, reason: limitRejection.reason }); + paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "validation_failed" }); + paymentProcessorSessionDuration.observe( + { asset: body.asset, outcome: "validation_failed" }, + (Date.now() - sessionStart) / 1000, + ); + return res.status(400).json({ + error: limitRejection.message, + ...limitRejection.details, }); + } - const metadata = - body.metadata && typeof body.metadata === "object" - ? { ...body.metadata } - : {}; - metadata.branding_config = resolvedBrandingConfig; - - const payload = { - id: paymentId, - merchant_id: req.merchant.id, - amount: body.amount, - asset, - asset_issuer: assetIssuer || null, - recipient: body.recipient, - description: body.description || null, - memo: body.message || body.memo || null, - memo_type: body.message ? "text" : (body.memo_type || null), - webhook_url: body.webhook_url || null, - client_id: body.client_id || null, - status: "pending", - tx_id: null, - metadata, - sandbox: isSandbox, - created_at: now, - }; - - const { error: insertError } = await supabase - .from("payments") - .insert(payload); + // Shared business-rule validation (issue #1087) — allowed-issuers check: + // if the merchant has configured a non-empty allowlist, only those + // issuer addresses may be used. + const allowedIssuerRejection = validateAllowedIssuers({ + asset, + assetIssuer, + allowedIssuers: req.merchant.allowed_issuers, + }); + if (allowedIssuerRejection) { + paymentFailedCounter.inc({ asset: body.asset, reason: "invalid_issuer" }); + paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "validation_failed" }); + paymentProcessorSessionDuration.observe( + { asset: body.asset, outcome: "validation_failed" }, + (Date.now() - sessionStart) / 1000, + ); + return res.status(400).json({ error: allowedIssuerRejection.message }); + } - if (insertError) { - insertError.status = 500; - paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "persistence_failed" }); - paymentProcessorSessionDuration.observe( - { asset: body.asset, outcome: "persistence_failed" }, - (Date.now() - sessionStart) / 1000, - ); - throw insertError; - } + const isSandbox = body.sandbox === true; + const baseId = randomUUID(); + const paymentId = isSandbox ? `test_${baseId}` : baseId; + const now = new Date().toISOString(); + const paymentLinkBase = + process.env.PAYMENT_LINK_BASE || "http://localhost:3000"; + const paymentLink = `${paymentLinkBase}/pay/${paymentId}`; + const resolvedBrandingConfig = resolveBrandingConfig({ + merchantBranding: req.merchant.branding_config, + brandingOverrides: body.branding_overrides, + }); - // Only record production metrics for non-sandbox payments. - if (!isSandbox) { - paymentCreatedCounter.inc({ asset: body.asset }); - paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "created" }); - paymentProcessorSessionDuration.observe( - { asset: body.asset, outcome: "created" }, - (Date.now() - sessionStart) / 1000, - ); - } + const metadata = + body.metadata && typeof body.metadata === "object" + ? { ...body.metadata } + : {}; + metadata.branding_config = resolvedBrandingConfig; + + const payload = { + id: paymentId, + merchant_id: req.merchant.id, + amount: body.amount, + asset, + asset_issuer: assetIssuer || null, + recipient: body.recipient, + description: body.description || null, + memo: body.message || body.memo || null, + memo_type: body.message ? "text" : (body.memo_type || null), + webhook_url: body.webhook_url || null, + client_id: body.client_id || null, + status: "pending", + tx_id: null, + metadata, + sandbox: isSandbox, + created_at: now, + }; + + // Issue #1449: transient persistence failures are retried with + // exponential backoff; validation/constraint errors are not. + try { + await insertPaymentSessionWithRetry(supabase, payload); + } catch (insertError) { + insertError.status = 500; + paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "persistence_failed" }); + paymentProcessorSessionDuration.observe( + { asset: body.asset, outcome: "persistence_failed" }, + (Date.now() - sessionStart) / 1000, + ); + throw insertError; + } - logger.info({ paymentId: paymentId }, "DEBUG: createSession success"); - res.status(201).json({ - payment_id: paymentId, - payment_link: paymentLink, - status: "pending", - sandbox: isSandbox, - branding_config: resolvedBrandingConfig, - }); - } catch (err) { - logger.error({ err, merchantId: req.merchant?.id }, "DEBUG: createSession error"); - if (err.status === 400 && err.details) { - return res.status(400).json({ error: err.message, ...err.details }); - } - next(err); + // Only record production metrics for non-sandbox payments. + if (!isSandbox) { + paymentCreatedCounter.inc({ asset: body.asset }); + paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "created" }); + paymentProcessorSessionDuration.observe( + { asset: body.asset, outcome: "created" }, + (Date.now() - sessionStart) / 1000, + ); } + + logger.info({ paymentId: paymentId }, "DEBUG: createSession success"); + res.status(201).json({ + payment_id: paymentId, + payment_link: paymentLink, + status: "pending", + sandbox: isSandbox, + branding_config: resolvedBrandingConfig, + }); } router.post("/create-payment", createPaymentRateLimit, recaptchaMiddleware(), validateRequest({ body: paymentSessionZodSchema }), sanitizeMetadataMiddleware, createSession); diff --git a/backend/src/services/merchantService.js b/backend/src/services/merchantService.js index 54ac8594..1fcd0310 100644 --- a/backend/src/services/merchantService.js +++ b/backend/src/services/merchantService.js @@ -2,6 +2,10 @@ import { randomBytes } from "crypto"; import { supabase } from "../lib/supabase.js"; import { resolveBrandingConfig } from "../lib/branding.js"; import { resolveMerchantSettings } from "../lib/merchant-settings.js"; +import { + MAX_GRACE_PERIOD_HOURS, + normalizeApiKeyExpiry, +} from "../lib/merchant-payload-validation.js"; import { sendWebhook } from "../lib/webhooks.js"; import { getPayloadForVersion } from "../webhooks/resolver.js"; @@ -26,6 +30,30 @@ function resolveWebhookSecretRotationGraceHours(requestValue) { return Math.min(parsed, 168); } +/** + * Issue #1482: service-level guards so non-HTTP callers (jobs, scripts, other + * services) get the same guarantees as the validated routes. + */ +function assertMerchantId(merchantId) { + if (typeof merchantId !== "string" || merchantId.trim() === "") { + const err = new Error("merchantId is required"); + err.status = 400; + throw err; + } +} + +function resolveApiKeyGraceHours(value) { + if (value === undefined || value === null) { + return DEFAULT_API_KEY_ROTATION_GRACE_HOURS; + } + if (!Number.isInteger(value)) { + const err = new Error("grace_period_hours must be an integer"); + err.status = 400; + throw err; + } + return Math.min(Math.max(value, 0), MAX_GRACE_PERIOD_HOURS); +} + export const merchantService = { async registerMerchant(body) { const { email } = body; @@ -85,6 +113,9 @@ export const merchantService = { }, async rotateApiKey(merchantId, gracePeriodHours = DEFAULT_API_KEY_ROTATION_GRACE_HOURS) { + assertMerchantId(merchantId); + const graceHours = resolveApiKeyGraceHours(gracePeriodHours); + // Get current merchant to preserve old key const { data: merchant, error: fetchError } = await supabase .from("merchants") @@ -105,7 +136,6 @@ export const merchantService = { const newApiKey = `sk_${randomBytes(24).toString("hex")}`; const now = Date.now(); - const graceHours = Math.min(Math.max(gracePeriodHours, 0), 168); // Clamp between 0 and 168 hours (1 week) const oldKeyExpiry = new Date(now + graceHours * 60 * 60 * 1000).toISOString(); const { error } = await supabase @@ -131,9 +161,12 @@ export const merchantService = { }, async setApiKeyExpiry(merchantId, expiresAt) { + assertMerchantId(merchantId); + const normalizedExpiry = normalizeApiKeyExpiry(expiresAt); + const { error } = await supabase .from("merchants") - .update({ api_key_expires_at: expiresAt }) + .update({ api_key_expires_at: normalizedExpiry }) .eq("id", merchantId); if (error) { @@ -141,7 +174,7 @@ export const merchantService = { throw error; } - return { api_key_expires_at: expiresAt }; + return { api_key_expires_at: normalizedExpiry }; }, async getApiKeyStatus(merchantId) { diff --git a/backend/src/services/merchantService.validation.test.js b/backend/src/services/merchantService.validation.test.js new file mode 100644 index 00000000..e8ea3e87 --- /dev/null +++ b/backend/src/services/merchantService.validation.test.js @@ -0,0 +1,91 @@ +/** + * Service-level guards for the Merchant Settings & API Key Service + * (issue #1482) — defense in depth for callers that bypass HTTP validation. + */ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const { mockUpdate, mockFrom } = vi.hoisted(() => { + const mockUpdate = vi.fn(); + const chain = { + select: vi.fn(() => chain), + eq: vi.fn(() => chain), + maybeSingle: vi.fn(async () => ({ data: { api_key: "sk_old" }, error: null })), + update: vi.fn((payload) => { + mockUpdate(payload); + return { eq: vi.fn(async () => ({ error: null })) }; + }), + }; + return { mockUpdate, mockFrom: vi.fn(() => chain) }; +}); + +vi.mock("../lib/supabase.js", () => ({ supabase: { from: mockFrom } })); +vi.mock("../lib/webhooks.js", () => ({ sendWebhook: vi.fn() })); +vi.mock("../webhooks/resolver.js", () => ({ getPayloadForVersion: vi.fn() })); +vi.mock("../lib/logger.js", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +import { merchantService } from "./merchantService.js"; + +const DAY = 24 * 60 * 60 * 1000; + +beforeEach(() => vi.clearAllMocks()); + +describe("merchantService.setApiKeyExpiry", () => { + it("normalizes and persists a valid expiry", async () => { + const future = new Date(Date.now() + 10 * DAY); + const result = await merchantService.setApiKeyExpiry( + "merchant-1", + future.toISOString().replace("Z", "+00:00"), + ); + expect(result).toEqual({ api_key_expires_at: future.toISOString() }); + expect(mockUpdate).toHaveBeenCalledWith({ api_key_expires_at: future.toISOString() }); + }); + + it.each([ + ["past", "2001-01-01T00:00:00Z"], + ["too far", new Date(Date.now() + 1000 * DAY).toISOString()], + ["garbage", "soon"], + ["non-string", 12345], + ])("rejects a %s expiry with 400 and does not write", async (_label, value) => { + await expect(merchantService.setApiKeyExpiry("merchant-1", value)).rejects.toMatchObject({ + status: 400, + }); + expect(mockUpdate).not.toHaveBeenCalled(); + }); + + it("rejects a missing merchantId", async () => { + const future = new Date(Date.now() + DAY).toISOString(); + await expect(merchantService.setApiKeyExpiry("", future)).rejects.toMatchObject({ status: 400 }); + await expect(merchantService.setApiKeyExpiry(undefined, future)).rejects.toMatchObject({ + status: 400, + }); + }); +}); + +describe("merchantService.rotateApiKey", () => { + it("clamps the grace period to one week", async () => { + const result = await merchantService.rotateApiKey("merchant-1", 10_000); + expect(result.grace_period_hours).toBe(168); + }); + + it("uses the default grace period when none is supplied", async () => { + const result = await merchantService.rotateApiKey("merchant-1", undefined); + expect(result.grace_period_hours).toBe(24); + }); + + it.each([["string", "24"], ["float", 1.5], ["NaN", Number.NaN]])( + "rejects a %s grace period with 400 instead of crashing with a 500", + async (_label, value) => { + await expect(merchantService.rotateApiKey("merchant-1", value)).rejects.toMatchObject({ + status: 400, + }); + expect(mockUpdate).not.toHaveBeenCalled(); + }, + ); + + it("rejects a missing merchantId before touching the database", async () => { + await expect(merchantService.rotateApiKey(null)).rejects.toMatchObject({ status: 400 }); + expect(mockFrom).not.toHaveBeenCalled(); + }); +}); diff --git a/backend/src/services/paymentService.js b/backend/src/services/paymentService.js index 3d99b51a..ee52a3a6 100644 --- a/backend/src/services/paymentService.js +++ b/backend/src/services/paymentService.js @@ -29,6 +29,7 @@ import { paymentSignatureVerifier } from "../lib/payment-signature-verification. import { logger } from "../lib/logger.js"; import { validatePaymentSession } from "../lib/payment-session-validator.js"; import { getSupabaseClient } from "../lib/supabase-client.js"; +import { insertPaymentSessionWithRetry } from "../lib/payment-session-retry.js"; import { paymentProcessorSessionsTotal, paymentProcessorSessionDuration, @@ -551,9 +552,10 @@ export const paymentService = { created_at: now, }; - const { error: insertError } = await supabase.from("payments").insert(payload); - - if (insertError) { + // Issue #1449: retry transient persistence failures with backoff. + try { + await insertPaymentSessionWithRetry(supabase, payload); + } catch (insertError) { insertError.status = 500; recordSessionOutcome("persistence_failed"); throw insertError; diff --git a/backend/src/services/paymentService.test.js b/backend/src/services/paymentService.test.js index d1c5b019..796cd642 100644 --- a/backend/src/services/paymentService.test.js +++ b/backend/src/services/paymentService.test.js @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; const { mockQueryWithRetry, @@ -206,6 +206,52 @@ describe("paymentService", () => { ); }); + describe("createPaymentSession persistence retry (issue #1449)", () => { + const merchant = { + id: "merchant-1", + allowed_issuers: [], + payment_limits: {}, + branding_config: {}, + }; + const body = { amount: 5, asset: "XLM", recipient: "GRECIPIENT" }; + + beforeEach(() => { + process.env.PAYMENT_SESSION_RETRY_BASE_DELAY_MS = "0"; + }); + + afterEach(() => { + delete process.env.PAYMENT_SESSION_RETRY_BASE_DELAY_MS; + }); + + it("retries a transient insert failure and creates the session", async () => { + const insert = vi + .fn() + .mockResolvedValueOnce({ error: { message: "TypeError: fetch failed", code: "" } }) + .mockResolvedValue({ error: null }); + mockSupabaseFrom.mockReturnValue({ insert }); + + const result = await paymentService.createPaymentSession(merchant, body); + + expect(result.status).toBe("pending"); + expect(insert).toHaveBeenCalledTimes(2); + // The same server-generated id is reused on retry (idempotent insert). + expect(insert.mock.calls[0][0].id).toBe(insert.mock.calls[1][0].id); + }); + + it("surfaces a non-retryable insert error as 500 without retrying", async () => { + const insert = vi + .fn() + .mockResolvedValue({ error: { message: "check violation", code: "23514" } }); + mockSupabaseFrom.mockReturnValue({ insert }); + + await expect(paymentService.createPaymentSession(merchant, body)).rejects.toMatchObject({ + status: 500, + code: "23514", + }); + expect(insert).toHaveBeenCalledTimes(1); + }); + }); + it("falls back to Supabase when the pooler exhausts retryable errors", async () => { const poolError = new Error("connection terminated"); poolError.code = "57P01"; diff --git a/backend/tests/helpers/fake-redis.js b/backend/tests/helpers/fake-redis.js index 6de8dbc0..74a9b7ff 100644 --- a/backend/tests/helpers/fake-redis.js +++ b/backend/tests/helpers/fake-redis.js @@ -1,4 +1,24 @@ /** + * Minimal in-memory Redis stand-in for concurrency tests. + * + * Supports the subset used by the idempotency middleware and the payment + * session lock: GET / SET (with EX/PX) / DEL, plus sendCommand for + * `SET key value NX PX ttl` and the compare-and-delete EVAL script. + * + * Each command awaits a FIFO latency tick and then executes its + * check-and-mutate step synchronously, which matches Redis' single-threaded + * atomicity while still letting concurrent callers interleave between + * commands — exactly the window real race conditions live in. + */ +export function createFakeRedis({ latencyMs = 0 } = {}) { + const store = new Map(); + const commandLog = []; + + const now = () => Date.now(); + const live = (key) => { + const entry = store.get(key); + if (!entry) return null; + if (entry.expiresAt !== null && entry.expiresAt <= now()) { * In-memory stand-in for the subset of Redis used by the exchange-rate * coordinator (issues #1445, #1446): SET [PX ms] [NX], GET, DEL and the * compare-and-delete EVAL script. Several coordinators can share one @@ -22,6 +42,89 @@ export function createFakeRedis({ latencyMs = 0 } = {}) { } return entry; }; + // Commands complete in the order they were issued (FIFO), like a single + // Redis connection: the idempotency middleware and the session lock share + // one client, and the lock's correctness relies on that ordering. + let queue = Promise.resolve(); + // With no latency, stay on the microtask queue so tests using + // vi.useFakeTimers() are not blocked on a timer that never fires. + const tick = () => { + if (latencyMs <= 0) return Promise.resolve(); + queue = queue.then( + () => new Promise((resolve) => setTimeout(resolve, Math.random() * latencyMs)), + ); + return queue; + }; + + function setEntry(key, value, { ex, px } = {}) { + let expiresAt = null; + if (ex) expiresAt = now() + Number(ex) * 1000; + if (px) expiresAt = now() + Number(px); + store.set(key, { value: String(value), expiresAt }); + } + + const client = { + isOpen: true, + store, + commandLog, + async get(key) { + await tick(); + commandLog.push(["GET", key]); + return live(key)?.value ?? null; + }, + async set(key, value, opts = {}) { + await tick(); + commandLog.push(["SET", key]); + if (opts.NX && live(key)) return null; + setEntry(key, value, { ex: opts.EX, px: opts.PX }); + return "OK"; + }, + async del(key) { + await tick(); + commandLog.push(["DEL", key]); + return store.delete(key) ? 1 : 0; + }, + async sendCommand(args) { + await tick(); + const [cmd, ...rest] = args; + commandLog.push([String(cmd).toUpperCase(), rest[0]]); + switch (String(cmd).toUpperCase()) { + case "SET": { + const [key, value, ...flags] = rest; + const upper = flags.map((f) => String(f).toUpperCase()); + if (upper.includes("NX") && live(key)) return null; + const pxIdx = upper.indexOf("PX"); + const exIdx = upper.indexOf("EX"); + setEntry(key, value, { + px: pxIdx >= 0 ? flags[pxIdx + 1] : undefined, + ex: exIdx >= 0 ? flags[exIdx + 1] : undefined, + }); + return "OK"; + } + case "EVAL": { + const [script, numKeys, key, token] = rest; + if (!/redis\.call\("get", KEYS\[1\]\) == ARGV\[1\]/.test(script) || numKeys !== "1") { + throw new Error("fake-redis: unsupported EVAL script"); + } + const entry = live(key); + if (entry && entry.value === token) { + store.delete(key); + return 1; + } + return 0; + } + case "GET": + return live(rest[0])?.value ?? null; + default: + throw new Error(`fake-redis: unsupported command ${cmd}`); + } + }, + keys(prefix = "") { + return [...store.keys()].filter((k) => k.startsWith(prefix) && live(k)); + }, + }; + + return client; const execute = (args) => { const [cmd, ...rest] = args;