diff --git a/.env.example b/.env.example index dcba576..a217342 100644 --- a/.env.example +++ b/.env.example @@ -20,6 +20,21 @@ CORS_ORIGIN=* # Rate limiting (per client IP, applied to /api) RATE_LIMIT_WINDOW_MS=60000 RATE_LIMIT_MAX=100 +RATE_LIMIT_MAX_KEYS=10000 + +# Honour X-Forwarded-For only behind a trusted reverse proxy +TRUST_PROXY=false + +# Per-actor mutation budgets (see docs/ABUSE_CONTROLS.md) +MUTATION_RATE_LIMIT_MAX_KEYS=10000 +MUTATION_RATE_LIMIT_TRANSFERS_WINDOW_MS=60000 +MUTATION_RATE_LIMIT_TRANSFERS_MAX=30 +MUTATION_RATE_LIMIT_USERS_WINDOW_MS=60000 +MUTATION_RATE_LIMIT_USERS_MAX=20 +MUTATION_RATE_LIMIT_QUOTE_WINDOW_MS=60000 +MUTATION_RATE_LIMIT_QUOTE_MAX=60 +MUTATION_RATE_LIMIT_ADMIN_WINDOW_MS=60000 +MUTATION_RATE_LIMIT_ADMIN_MAX=30 # Error tracking ERROR_TRACKING_ENABLED=true diff --git a/CHANGELOG.md b/CHANGELOG.md index a8e4e0a..7e7cd77 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,6 +18,13 @@ When preparing a new release: ### Added +- Route-family mutation rate limits and sanitized correlation IDs for abuse + control on transfer writes, user writes, quote, and admin diagnostics. + Actor keys are truncated token fingerprints (never raw secrets); limiter + tables are bounded by `RATE_LIMIT_MAX_KEYS` / + `MUTATION_RATE_LIMIT_MAX_KEYS`. `TRUST_PROXY` gates `X-Forwarded-For`. + See `docs/ABUSE_CONTROLS.md`. + - Cursor pagination for `GET /api/transfers` and `GET /api/audit`. Pass `?cursor=` (with optional `?order=asc|desc`) to page by an indexed position instead of a row offset; responses carry a `pageInfo` block with diff --git a/README.md b/README.md index 4791414..8de6f7b 100644 --- a/README.md +++ b/README.md @@ -37,6 +37,9 @@ The application is configured using environment variables (typically defined in | `CORS_ORIGIN` | Allowed CORS origin | `*` | | `RATE_LIMIT_WINDOW_MS` | Time window for rate limiting (ms) | `60000` | | `RATE_LIMIT_MAX` | Max requests per window | `100` | +| `RATE_LIMIT_MAX_KEYS` | Max distinct client identities retained by the global limiter | `10000` | +| `TRUST_PROXY` | Honour `X-Forwarded-For` when resolving client IP (`true`/`1`) | `false` | +| `MUTATION_RATE_LIMIT_*` | Per-actor budgets for transfer/user/quote/admin mutations (see `docs/ABUSE_CONTROLS.md`) | see docs | | `BODY_LIMIT` | Max JSON request body size | `100kb` | | `REQUEST_TIMEOUT_MS` | Request timeout before returning 503 (ms) | `15000` | | `DB_POOL_MIN` | Minimum database connections in pool | `2` | diff --git a/docs/ABUSE_CONTROLS.md b/docs/ABUSE_CONTROLS.md new file mode 100644 index 0000000..3eee0d0 --- /dev/null +++ b/docs/ABUSE_CONTROLS.md @@ -0,0 +1,109 @@ +# Abuse controls and correlation IDs + +RemitFlow Backend bounds high-volume retries and automated abuse on mutation +and provider-adjacent routes, and propagates a safe correlation id on every +request so incidents can be traced without leaking account secrets. + +## Global limit + +Every `/api/*` request is subject to an IP-keyed fixed-window limit: + +| Setting | Env | Default | +|---|---|---| +| Window | `RATE_LIMIT_WINDOW_MS` | `60000` | +| Max requests | `RATE_LIMIT_MAX` | `100` | +| Max tracked keys | `RATE_LIMIT_MAX_KEYS` | `10000` | + +Responses carry `X-RateLimit-Limit`, `X-RateLimit-Remaining`, +`X-RateLimit-Reset`, and `X-RateLimit-Policy`. Exhausted budgets answer +`429` with `Retry-After` and a JSON error that includes `requestId`. + +## Mutation / route-family limits + +Authenticated write paths use a **stricter actor-keyed** budget on top of the +global limit. The actor key is a truncated SHA-256 fingerprint of the API +token (or admin key) — never the raw secret. Public quote uses the client IP. + +| Family | Routes | Env (window / max) | Default max / min | +|---|---|---|---| +| `transfers` | `POST /api/transfers`, claim, cancel, archive, unarchive | `MUTATION_RATE_LIMIT_TRANSFERS_*` | 30 / 60s | +| `users` | `POST /api/users` | `MUTATION_RATE_LIMIT_USERS_*` | 20 / 60s | +| `quote` | `GET /api/quote` | `MUTATION_RATE_LIMIT_QUOTE_*` | 60 / 60s | +| `admin` | `GET /api/admin/diagnostics` | `MUTATION_RATE_LIMIT_ADMIN_*` | 30 / 60s | + +Shared cap: `MUTATION_RATE_LIMIT_MAX_KEYS` (default `10000`). + +Actors are isolated: one token burning its transfer budget does not exhaust +another token's budget. Limiter tables are bounded. Expired windows are pruned +first; when every key is still live, a new identity receives 429 with +`Retry-After` until a slot expires. This preserves active budgets under an +identity flood, at the cost of delaying new identities when the table is full. + +## Proxy trust + +`TRUST_PROXY=true` uses Express `trust proxy` for one proxy hop. Rate +limiting reads Express's resolved `req.ip`, so an attacker-controlled left-most +`X-Forwarded-For` value cannot rotate the identity when the trusted proxy +appends the connecting address. Enable this only when the deployment has +exactly one trusted reverse proxy; otherwise leave it off until the app's +trust proxy setting matches the deployment topology. + +## Correlation IDs + +On every request, the [request-ID middleware](../src/middleware/requestId.js): + +- Trims leading and trailing whitespace from each candidate header value. +- Accepts a nonempty value of at most 128 characters after trimming, matching + `^[A-Za-z0-9._:-]+$`. +- Uses a valid `X-Request-Id` first, otherwise a valid `X-Correlation-Id`. + If neither is valid, generates a fresh UUID. +- Echoes the selected value on both `X-Request-Id` and `X-Correlation-Id`. +- Exposes it as `req.id` / `req.correlationId` and on error envelopes as + `error.requestId`. + +Values outside this format are ignored after trimming. These checks do not +identify secrets or personal data, or hash/redact an accepted identifier. An +API token or account identifier made only of permitted characters can still +pass them. + +Supply a fresh opaque ID, such as a UUID, and keep credentials, account details +and other personal data out of both headers. Omit both headers when a +server-generated ID is sufficient. Accepted IDs are echoed and written by the +[request logger](../src/middleware/requestLogger.js), so treat them as values +that can appear in response headers, error envelopes and logs. + +## Test behaviour + +When `NODE_ENV=test`, rate limiters are no-ops unless +`ENABLE_RATE_LIMIT_IN_TEST=1` (or a limiter is constructed with +`forceInTest: true`). This keeps the functional suite independent of the +abuse budget while still allowing focused regression coverage. + +## Reproduce saturated-table HTTP load + +With this checkout's normal dependencies installed, run: + +```sh +node scripts/benchmark-rate-limit.cjs rate-limit-result.json > /dev/null +``` + +Use a new result filename. The script starts the real application on an ephemeral +loopback port, enables the limiter explicitly in test mode, and uses one trusted +proxy hop to generate distinct local identities. It admits 10,000 identities, +then measures three batches of 2,000 rejected newcomers at 32 concurrent +connections. Every measured response must be 429 with the global policy and +matching correlation ID; an exhausted original identity must remain blocked. +The script closes its server and connections when finished. + +JSON output retains each elapsed/CPU sample, Node and Express versions, and +SHA-256 hashes of the application and limiter source. CPU includes both the +HTTP client and server in the same process. Standard application logs go to +stdout, so keep the same redirection for both versions being compared. + +To compare another checkout with its dependencies already installed, pass its +path as the second argument. Alternate original/repaired/original/repaired runs +using the same script and distinct output filenames; report medians and ranges +because load on the host can vary. The workload uses generated local traffic and +does not measure deployed throughput or external provider behavior. Repeated +rejection avoids table scans until the next expiry; expiry-triggered cleanup +still takes time proportional to the number of tracked identities. diff --git a/scripts/benchmark-rate-limit.cjs b/scripts/benchmark-rate-limit.cjs new file mode 100644 index 0000000..9f3b948 --- /dev/null +++ b/scripts/benchmark-rate-limit.cjs @@ -0,0 +1,119 @@ +'use strict'; + +const assert = require('node:assert/strict'); +const { createHash } = require('node:crypto'); +const fs = require('node:fs'); +const http = require('node:http'); +const path = require('node:path'); +const { performance } = require('node:perf_hooks'); + +// Usage: node scripts/benchmark-rate-limit.cjs NEW-RESULT.json [APP-CHECKOUT] +const outputPath = process.argv[2]; +if (!outputPath || process.argv.length > 4) { + console.error('Usage: node scripts/benchmark-rate-limit.cjs NEW-RESULT.json [APP-CHECKOUT]'); + process.exitCode = 2; +} else { + run().catch((error) => { + console.error(error); + process.exitCode = 1; + }); +} + +async function run() { + if (fs.existsSync(outputPath)) throw new Error(`Result already exists: ${outputPath}`); + const appRoot = path.resolve(process.argv[3] || path.join(__dirname, '..')); + Object.assign(process.env, { + NODE_ENV: 'test', + ENABLE_RATE_LIMIT_IN_TEST: '1', + RATE_LIMIT_MAX_KEYS: '10000', + RATE_LIMIT_MAX: '1', + RATE_LIMIT_WINDOW_MS: '180000', + TRUST_PROXY: 'true', + ERROR_TRACKING_ENABLED: 'false', + }); + + const createApp = require(path.join(appRoot, 'src/app')); + const server = http.createServer(createApp()); + const agent = new http.Agent({ keepAlive: true, maxSockets: 32 }); + let port; + + function request(identity) { + const ip = `10.${(identity >>> 16) & 255}.${(identity >>> 8) & 255}.${identity & 255}`; + return new Promise((resolve, reject) => { + const req = http.get({ + host: '127.0.0.1', port, path: '/api/version', agent, + headers: { 'X-Forwarded-For': ip, 'X-Request-Id': `load-${identity}` }, + }, (res) => { + let body = ''; + res.on('data', (chunk) => { body += chunk; }); + res.on('error', reject); + res.on('end', () => { + try { + resolve({ status: res.statusCode, headers: res.headers, body: JSON.parse(body) }); + } catch (error) { + reject(error); + } + }); + }); + req.setTimeout(30_000, () => req.destroy(new Error('Local benchmark request timed out'))); + req.on('error', reject); + }); + } + + async function batch(start, count, expectedStatus) { + let next = 0; + let completed = 0; + await Promise.all(Array.from({ length: 32 }, async () => { + while (next < count) { + const identity = start + next++; + const result = await request(identity); + assert.equal(result.status, expectedStatus); + if (expectedStatus === 429) { + assert.equal(result.body.error.details.policy, 'global'); + assert.equal(result.headers['x-request-id'], `load-${identity}`); + } + completed += 1; + } + })); + return completed; + } + + try { + await new Promise((resolve, reject) => { + server.once('error', reject); + server.listen(0, '127.0.0.1', resolve); + }); + port = server.address().port; + const fill = await batch(1, 10_000, 200); + const samples = []; + for (let round = 0; round < 3; round += 1) { + const started = performance.now(); + const cpuStarted = process.cpuUsage(); + const count = await batch(20_000 + round * 2_000, 2_000, 429); + const cpu = process.cpuUsage(cpuStarted); + samples.push({ count, elapsedMs: performance.now() - started, cpuMs: (cpu.user + cpu.system) / 1_000 }); + } + + const original = await request(1); + assert.equal(original.status, 429); + assert.equal(original.body.error.message, 'Too many requests, please try again later'); + + const sourceSha256 = {}; + for (const file of ['src/app.js', 'src/middleware/rateLimit.js']) { + sourceSha256[file] = createHash('sha256').update(fs.readFileSync(path.join(appRoot, file))).digest('hex'); + } + const result = { + node: process.version, + express: require(path.join(appRoot, 'node_modules/express/package.json')).version, + sourceSha256, + path: '/api/version', concurrency: 32, maxKeys: 10_000, fill, samples, + activeBudgetPreserved: true, + scope: 'Real createApp HTTP on loopback; test mode with rate limiting enabled; one trusted proxy hop; error tracking disabled; CPU includes client and server in one process; generated local load, not deployed performance.', + }; + fs.writeFileSync(outputPath, `${JSON.stringify(result, null, 2)}\n`, { flag: 'wx' }); + } finally { + agent.destroy(); + server.closeAllConnections(); + if (server.listening) await new Promise((resolve) => server.close(resolve)); + } +} diff --git a/src/app.js b/src/app.js index bd5b065..82909b2 100644 --- a/src/app.js +++ b/src/app.js @@ -16,6 +16,7 @@ const maintenanceMode = require('./middleware/maintenanceMode'); const jsonError = require('./middleware/jsonError'); const notFound = require('./middleware/notFound'); const errorHandler = require('./middleware/errorHandler'); +const { resolveClientIp } = require('./utils/clientIdentity'); /** * Build and configure the Express application. @@ -26,16 +27,16 @@ const errorHandler = require('./middleware/errorHandler'); function createApp() { const app = express(); + // Only honour X-Forwarded-* when explicitly configured. Blind trust lets + // clients rotate forged IPs and evade the global abuse budget. + if (config.trustProxy) { + app.set('trust proxy', 1); + } + // Core middleware. app.use(securityHeaders); app.use(cacheControl({ policy: config.cache.defaultPolicy })); app.use(cors({ origin: config.corsOrigin })); - app.use(express.json({ limit: config.bodyLimit })); - app.use(express.urlencoded({ extended: false, limit: config.bodyLimit })); - app.use(jsonError); - - // Fail slow requests instead of hanging the connection. - app.use(requestTimeout({ ms: config.requestTimeoutMs })); // Assign/propagate a correlation id before logging. app.use(requestId); @@ -46,8 +47,27 @@ function createApp() { } app.use(requestLogger); - // Basic abuse protection on the API surface. - app.use('/api', rateLimit(config.rateLimit)); + // Apply the API budget before body parsing so invalid bodies cannot bypass it. + app.use( + '/api', + rateLimit({ + name: 'global', + windowMs: config.rateLimit.windowMs, + max: config.rateLimit.max, + maxKeys: config.rateLimit.maxKeys, + trustProxy: config.trustProxy, + keyGenerator(req) { + return resolveClientIp(req, { trustProxy: config.trustProxy }); + }, + }) + ); + + app.use(express.json({ limit: config.bodyLimit })); + app.use(express.urlencoded({ extended: false, limit: config.bodyLimit })); + app.use(jsonError); + + // Fail slow handlers instead of hanging the connection. + app.use(requestTimeout({ ms: config.requestTimeoutMs })); // Block all non-health API traffic while maintenance mode is active. app.use(maintenanceMode); diff --git a/src/config/index.js b/src/config/index.js index 41f9a71..320e04f 100644 --- a/src/config/index.js +++ b/src/config/index.js @@ -2,13 +2,27 @@ require('dotenv').config(); +/** + * Parse a positive integer env var with a fallback. + * @param {string|undefined} value + * @param {number} fallback + * @returns {number} + */ +function intEnv(value, fallback) { + const parsed = parseInt(value, 10); + return Number.isFinite(parsed) && parsed > 0 ? parsed : fallback; +} + /** * Centralized application configuration. * Values are read from environment variables with sensible defaults * so the app can boot even without a .env file present. */ +const env = process.env.NODE_ENV || 'development'; +const isTest = env === 'test'; + const config = { - env: process.env.NODE_ENV || 'development', + env, port: parseInt(process.env.PORT, 10) || 3000, baseCurrency: process.env.DEFAULT_BASE_CURRENCY || 'USD', @@ -34,9 +48,44 @@ const config = { // Per-request time budget before a 503 is returned. requestTimeoutMs: parseInt(process.env.REQUEST_TIMEOUT_MS, 10) || 15 * 1000, + /** + * Whether to trust `X-Forwarded-For` when resolving the client IP. + * Off by default so untrusted clients cannot rotate IPs to bypass limits. + * Enable only behind a reverse proxy that strips/forges the header safely. + */ + trustProxy: process.env.TRUST_PROXY === 'true' || process.env.TRUST_PROXY === '1', + rateLimit: { - windowMs: parseInt(process.env.RATE_LIMIT_WINDOW_MS, 10) || 60 * 1000, - max: parseInt(process.env.RATE_LIMIT_MAX, 10) || 100, + windowMs: intEnv(process.env.RATE_LIMIT_WINDOW_MS, 60 * 1000), + max: intEnv(process.env.RATE_LIMIT_MAX, isTest ? 10_000 : 100), + maxKeys: intEnv(process.env.RATE_LIMIT_MAX_KEYS, 10_000), + }, + + /** + * Stricter per-actor budgets for mutation / expensive routes. + * Defaults are deliberately higher than a single integration test suite + * needs, while still bounding automated abuse of provider-backed paths. + */ + mutationRateLimit: { + maxKeys: intEnv(process.env.MUTATION_RATE_LIMIT_MAX_KEYS, 10_000), + transfers: { + windowMs: intEnv(process.env.MUTATION_RATE_LIMIT_TRANSFERS_WINDOW_MS, 60 * 1000), + // Generous under test so the suite does not trip the abuse budget; production + // defaults stay tight enough to bound provider-quota exhaustion. + max: intEnv(process.env.MUTATION_RATE_LIMIT_TRANSFERS_MAX, isTest ? 10_000 : 30), + }, + users: { + windowMs: intEnv(process.env.MUTATION_RATE_LIMIT_USERS_WINDOW_MS, 60 * 1000), + max: intEnv(process.env.MUTATION_RATE_LIMIT_USERS_MAX, isTest ? 10_000 : 20), + }, + quote: { + windowMs: intEnv(process.env.MUTATION_RATE_LIMIT_QUOTE_WINDOW_MS, 60 * 1000), + max: intEnv(process.env.MUTATION_RATE_LIMIT_QUOTE_MAX, isTest ? 10_000 : 60), + }, + admin: { + windowMs: intEnv(process.env.MUTATION_RATE_LIMIT_ADMIN_WINDOW_MS, 60 * 1000), + max: intEnv(process.env.MUTATION_RATE_LIMIT_ADMIN_MAX, isTest ? 10_000 : 30), + }, }, errorTracking: { diff --git a/src/controllers/adminController.js b/src/controllers/adminController.js index 870e10e..ea41def 100644 --- a/src/controllers/adminController.js +++ b/src/controllers/adminController.js @@ -30,7 +30,9 @@ function getDiagnostics(req, res) { fee: config.fee, maxTransferAmount: config.maxTransferAmount, stellar: config.stellar, + trustProxy: config.trustProxy, rateLimit: config.rateLimit, + mutationRateLimit: config.mutationRateLimit, errorTrackingEnabled: config.errorTracking.enabled, }, stats: { diff --git a/src/middleware/adminAuth.js b/src/middleware/adminAuth.js index 7828ad7..5217b65 100644 --- a/src/middleware/adminAuth.js +++ b/src/middleware/adminAuth.js @@ -23,6 +23,8 @@ function adminAuth(req, res, next) { return next(new ApiError(401, 'Unauthorized')); } + // Expose for downstream actor-keyed rate limiting (never log the raw value). + req.adminToken = token; next(); } diff --git a/src/middleware/mutationRateLimit.js b/src/middleware/mutationRateLimit.js new file mode 100644 index 0000000..c4f6ccb --- /dev/null +++ b/src/middleware/mutationRateLimit.js @@ -0,0 +1,55 @@ +'use strict'; + +const config = require('../config'); +const rateLimit = require('./rateLimit'); +const { resolveActorKey } = require('../utils/clientIdentity'); + +/** + * Route-family mutation rate limiters. + * + * Applied after authentication so the actor key is a token fingerprint + * (never the raw secret). Quote stays IP-keyed because it is public. + * Each limiter is independently bounded so one hot route cannot starve + * another family's budget, and the shared maxKeys cap keeps memory finite. + */ + +function buildLimiter(family, overrides = {}) { + const settings = config.mutationRateLimit[family] || {}; + const trustProxy = config.trustProxy; + + return rateLimit({ + name: `mutation:${family}`, + windowMs: overrides.windowMs || settings.windowMs, + max: overrides.max || settings.max, + maxKeys: overrides.maxKeys || config.mutationRateLimit.maxKeys, + trustProxy, + forceInTest: Boolean(overrides.forceInTest), + keyGenerator(req) { + return `${family}:${resolveActorKey(req, { trustProxy })}`; + }, + }); +} + +const transfers = buildLimiter('transfers'); +const users = buildLimiter('users'); +const quote = buildLimiter('quote'); +const admin = buildLimiter('admin'); + +/** + * Reset every mutation limiter. Intended for tests only. + */ +function resetAll() { + transfers.reset(); + users.reset(); + quote.reset(); + admin.reset(); +} + +module.exports = { + transfers, + users, + quote, + admin, + resetAll, + buildLimiter, +}; diff --git a/src/middleware/rateLimit.js b/src/middleware/rateLimit.js index b7fa7bd..6748ce8 100644 --- a/src/middleware/rateLimit.js +++ b/src/middleware/rateLimit.js @@ -1,52 +1,140 @@ 'use strict'; const ApiError = require('../utils/ApiError'); +const { resolveClientIp } = require('../utils/clientIdentity'); /** - * Minimal in-memory rate limiter. - * Tracks request counts per client (by IP) within a fixed time window. - * Counters live only in process memory, which is fine for a single-node - * demo; a real deployment would use a shared store such as Redis. + * In-memory fixed-window rate limiter with a bounded key table. + * + * Counters live only in process memory (fine for a single-node demo). A + * production deployment should back this with a shared store such as Redis. + * The map is capped by `maxKeys`: expired entries are pruned first. When + * all slots are live, new identities receive 429 until a slot expires rather + * than evicting an active budget and allowing repeat attempts. + * + * Under `NODE_ENV=test` the limiter is a no-op unless + * `ENABLE_RATE_LIMIT_IN_TEST=1`, so the rest of the suite is not coupled to + * the abuse budget. Dedicated limiter and abuse-control tests opt back in. * * @param {object} [options] * @param {number} [options.windowMs] - length of the window in milliseconds. - * @param {number} [options.max] - max requests allowed per window per client. - * @returns {import('express').RequestHandler} + * @param {number} [options.max] - max requests allowed per window per key. + * @param {number} [options.maxKeys] - hard cap on tracked identities. + * @param {(req: import('express').Request) => string} [options.keyGenerator] + * @param {string} [options.name] - label included in 429 details (no secrets). + * @param {boolean} [options.trustProxy] - whether to honour X-Forwarded-For. + * @param {boolean} [options.forceInTest] - enforce even when NODE_ENV=test. + * @returns {import('express').RequestHandler & { reset: Function, size: Function }} */ function rateLimit(options = {}) { const windowMs = options.windowMs || 60 * 1000; const max = options.max || 100; + const maxKeys = options.maxKeys || 10_000; + const name = options.name || 'default'; + const trustProxy = Boolean(options.trustProxy); + const forceInTest = Boolean(options.forceInTest); + const keyGenerator = + options.keyGenerator || + ((req) => resolveClientIp(req, { trustProxy })); + + /** @type {Map} */ const hits = new Map(); + let nextResetAt = Infinity; + + function pruneExpired(now) { + // A full table should reject new identities without rescanning live budgets. + if (now < nextResetAt) return; + nextResetAt = Infinity; + for (const [key, entry] of hits) { + if (now >= entry.resetAt) { + hits.delete(key); + } else { + nextResetAt = Math.min(nextResetAt, entry.resetAt); + } + } + } + + function hasCapacity(now) { + // Also prune before renewing an expired identity below capacity, so the + // cached deadline remains exact if the wall clock later moves backward. + pruneExpired(now); + return hits.size < maxKeys; + } + + function capacityRetryAfter(now) { + return Math.max(1, Math.ceil((nextResetAt - now) / 1000)); + } + + function rateLimitMiddleware(req, res, next) { + const skipForTest = + process.env.NODE_ENV === 'test' && + process.env.ENABLE_RATE_LIMIT_IN_TEST !== '1' && + !forceInTest; + if (skipForTest) { + return next(); + } - return function rateLimitMiddleware(req, res, next) { const now = Date.now(); - const key = req.ip || 'unknown'; + const key = keyGenerator(req) || 'unknown'; let entry = hits.get(key); if (!entry || now >= entry.resetAt) { - entry = { count: 0, resetAt: now + windowMs }; + if (!hasCapacity(now)) { + const retryAfter = capacityRetryAfter(now); + res.set('X-RateLimit-Limit', String(max)); + res.set('X-RateLimit-Remaining', '0'); + res.set('X-RateLimit-Policy', name); + res.set('Retry-After', String(retryAfter)); + return next( + ApiError.tooManyRequests('Rate limit capacity reached, please try again later', { + retryAfter, + limit: max, + windowMs, + policy: name, + }) + ); + } + entry = { count: 0, resetAt: now + windowMs, touchedAt: now }; + hits.delete(key); hits.set(key, entry); + nextResetAt = Math.min(nextResetAt, entry.resetAt); } entry.count += 1; + entry.touchedAt = now; const remaining = Math.max(0, max - entry.count); res.set('X-RateLimit-Limit', String(max)); res.set('X-RateLimit-Remaining', String(remaining)); res.set('X-RateLimit-Reset', String(Math.ceil(entry.resetAt / 1000))); + res.set('X-RateLimit-Policy', name); if (entry.count > max) { - const retryAfter = Math.ceil((entry.resetAt - now) / 1000); + const retryAfter = Math.max(1, Math.ceil((entry.resetAt - now) / 1000)); res.set('Retry-After', String(retryAfter)); return next( ApiError.tooManyRequests('Too many requests, please try again later', { retryAfter, + limit: max, + windowMs, + policy: name, }) ); } return next(); + } + + rateLimitMiddleware.reset = function reset() { + hits.clear(); + nextResetAt = Infinity; }; + + rateLimitMiddleware.size = function size() { + return hits.size; + }; + + return rateLimitMiddleware; } module.exports = rateLimit; diff --git a/src/middleware/requestId.js b/src/middleware/requestId.js index 91cecea..028658a 100644 --- a/src/middleware/requestId.js +++ b/src/middleware/requestId.js @@ -2,18 +2,51 @@ const { newId } = require('../utils/ids'); +/** Max length accepted for an inbound correlation / request id. */ +const MAX_CORRELATION_ID_LENGTH = 128; + +/** + * Safe correlation ids: printable, non-secret, and short enough to log. + * Rejects anything that looks like it could carry a token or free-form PII + * dump (spaces, quotes, control chars, oversized values). + * @param {*} value + * @returns {string|null} + */ +function sanitizeCorrelationId(value) { + if (typeof value !== 'string') { + return null; + } + const trimmed = value.trim(); + if (!trimmed || trimmed.length > MAX_CORRELATION_ID_LENGTH) { + return null; + } + if (!/^[A-Za-z0-9._:-]+$/.test(trimmed)) { + return null; + } + return trimmed; +} + /** - * Attach a unique identifier to every request. - * Honours an inbound `X-Request-Id` header when present so a caller can - * correlate logs across services; otherwise a fresh id is generated. - * The id is exposed on `req.id` and echoed back in the response header. + * Attach a unique correlation identifier to every request. + * + * Honour inbound `X-Request-Id` or `X-Correlation-Id` when they pass + * sanitization so callers can stitch logs across services. Otherwise a + * fresh id is generated. The same value is exposed as `req.id` and + * `req.correlationId`, and echoed on both response headers so either + * convention works for clients. */ function requestId(req, res, next) { - const incoming = req.get('X-Request-Id'); - const id = incoming && incoming.trim() ? incoming.trim() : newId(); + const incoming = + sanitizeCorrelationId(req.get('X-Request-Id')) || + sanitizeCorrelationId(req.get('X-Correlation-Id')); + const id = incoming || newId(); req.id = id; + req.correlationId = id; res.set('X-Request-Id', id); + res.set('X-Correlation-Id', id); next(); } module.exports = requestId; +module.exports.sanitizeCorrelationId = sanitizeCorrelationId; +module.exports.MAX_CORRELATION_ID_LENGTH = MAX_CORRELATION_ID_LENGTH; diff --git a/src/routes/adminRoutes.js b/src/routes/adminRoutes.js index 8fe8f4f..472a6b4 100644 --- a/src/routes/adminRoutes.js +++ b/src/routes/adminRoutes.js @@ -3,6 +3,7 @@ const express = require('express'); const asyncHandler = require('../utils/asyncHandler'); const adminAuth = require('../middleware/adminAuth'); +const mutationRateLimit = require('../middleware/mutationRateLimit'); const adminController = require('../controllers/adminController'); const router = express.Router(); @@ -11,6 +12,7 @@ const router = express.Router(); router.get( '/diagnostics', adminAuth, + mutationRateLimit.admin, asyncHandler(adminController.getDiagnostics) ); diff --git a/src/routes/rateRoutes.js b/src/routes/rateRoutes.js index 6f0fc53..1b18982 100644 --- a/src/routes/rateRoutes.js +++ b/src/routes/rateRoutes.js @@ -5,6 +5,7 @@ const config = require('../config'); const cacheControl = require('../middleware/cacheControl'); const asyncHandler = require('../utils/asyncHandler'); const validate = require('../middleware/validate'); +const mutationRateLimit = require('../middleware/mutationRateLimit'); const rateController = require('../controllers/rateController'); const { validateQuoteQuery } = require('../validators/quoteValidator'); @@ -25,8 +26,11 @@ router.get( ); // GET /api/quote?amount=&from=&to= +// Quote is public but provider-adjacent; apply a dedicated IP budget so bursts +// cannot exhaust FX / settlement quotas shared with write paths. router.get( '/quote', + mutationRateLimit.quote, validate(validateQuoteQuery), asyncHandler(rateController.getQuote) ); diff --git a/src/routes/transferRoutes.js b/src/routes/transferRoutes.js index 58eac20..2181954 100644 --- a/src/routes/transferRoutes.js +++ b/src/routes/transferRoutes.js @@ -4,6 +4,7 @@ const express = require('express'); const asyncHandler = require('../utils/asyncHandler'); const validate = require('../middleware/validate'); const requireScope = require('../middleware/requireScope'); +const mutationRateLimit = require('../middleware/mutationRateLimit'); const transferController = require('../controllers/transferController'); const { validateCreateTransfer } = require('../validators/transferValidator'); @@ -13,6 +14,7 @@ const router = express.Router(); router.post( '/', requireScope(['transfers:write']), + mutationRateLimit.transfers, validate(validateCreateTransfer), asyncHandler(transferController.createTransfer) ); @@ -27,15 +29,35 @@ router.get('/stats', requireScope(['transfers:read']), asyncHandler(transferCont router.get('/:id', requireScope(['transfers:read']), asyncHandler(transferController.getTransfer)); // POST /api/transfers/:id/claim -router.post('/:id/claim', requireScope(['transfers:write']), asyncHandler(transferController.claimTransfer)); +router.post( + '/:id/claim', + requireScope(['transfers:write']), + mutationRateLimit.transfers, + asyncHandler(transferController.claimTransfer) +); // POST /api/transfers/:id/cancel -router.post('/:id/cancel', requireScope(['transfers:write']), asyncHandler(transferController.cancelTransfer)); +router.post( + '/:id/cancel', + requireScope(['transfers:write']), + mutationRateLimit.transfers, + asyncHandler(transferController.cancelTransfer) +); // POST /api/transfers/:id/archive -router.post('/:id/archive', requireScope(['transfers:write']), asyncHandler(transferController.archiveTransfer)); +router.post( + '/:id/archive', + requireScope(['transfers:write']), + mutationRateLimit.transfers, + asyncHandler(transferController.archiveTransfer) +); // POST /api/transfers/:id/unarchive -router.post('/:id/unarchive', requireScope(['transfers:write']), asyncHandler(transferController.unarchiveTransfer)); +router.post( + '/:id/unarchive', + requireScope(['transfers:write']), + mutationRateLimit.transfers, + asyncHandler(transferController.unarchiveTransfer) +); module.exports = router; diff --git a/src/routes/userRoutes.js b/src/routes/userRoutes.js index 1e82221..00b85ac 100644 --- a/src/routes/userRoutes.js +++ b/src/routes/userRoutes.js @@ -4,6 +4,7 @@ const express = require('express'); const asyncHandler = require('../utils/asyncHandler'); const validate = require('../middleware/validate'); const requireScope = require('../middleware/requireScope'); +const mutationRateLimit = require('../middleware/mutationRateLimit'); const userController = require('../controllers/userController'); const { validateCreateUser } = require('../validators/userValidator'); @@ -19,6 +20,7 @@ router.get('/:id', requireScope(['users:read']), asyncHandler(userController.get router.post( '/', requireScope(['users:write']), + mutationRateLimit.users, validate(validateCreateUser), asyncHandler(userController.createUser) ); diff --git a/src/utils/clientIdentity.js b/src/utils/clientIdentity.js new file mode 100644 index 0000000..68c51d2 --- /dev/null +++ b/src/utils/clientIdentity.js @@ -0,0 +1,69 @@ +'use strict'; + +const crypto = require('crypto'); + +/** + * Client identity helpers for abuse controls. + * + * Keys and correlation values must never embed raw API tokens, admin keys, + * account numbers, or other secrets. Fingerprints are one-way and truncated + * so they are useful for isolation without being reversible. + */ + +const FINGERPRINT_CHARS = 16; + +/** + * One-way, truncated fingerprint of a secret string. + * @param {string} value + * @returns {string} + */ +function fingerprint(value) { + if (typeof value !== 'string' || value.length === 0) { + return 'anon'; + } + return crypto.createHash('sha256').update(value).digest('hex').slice(0, FINGERPRINT_CHARS); +} + +/** + * Resolve the client IP for rate limiting. + * + * When `trustProxy` is false (default), only the direct socket address is + * used. When true, Express resolves `req.ip` from its configured trusted + * proxy hops. The left-most forwarded value is not necessarily trustworthy. + * + * @param {import('express').Request} req + * @param {{ trustProxy?: boolean }} [options] + * @returns {string} + */ +function resolveClientIp(req, options = {}) { + const trustProxy = Boolean(options.trustProxy); + if (trustProxy && typeof req.ip === 'string' && req.ip) { + return req.ip; + } + return (req.socket && req.socket.remoteAddress) || req.ip || 'unknown'; +} + +/** + * Build a stable actor key for mutation rate limiting. + * Prefers an authenticated token fingerprint; falls back to client IP. + * Never returns the raw token or admin key. + * + * @param {import('express').Request} req + * @param {{ trustProxy?: boolean }} [options] + * @returns {string} + */ +function resolveActorKey(req, options = {}) { + if (typeof req.token === 'string' && req.token) { + return `actor:${fingerprint(req.token)}`; + } + if (typeof req.adminToken === 'string' && req.adminToken) { + return `admin:${fingerprint(req.adminToken)}`; + } + return `ip:${resolveClientIp(req, options)}`; +} + +module.exports = { + fingerprint, + resolveClientIp, + resolveActorKey, +}; diff --git a/test/abuseControls.test.js b/test/abuseControls.test.js new file mode 100644 index 0000000..04de561 --- /dev/null +++ b/test/abuseControls.test.js @@ -0,0 +1,498 @@ +'use strict'; + +const { test, beforeEach } = require('node:test'); +const assert = require('node:assert/strict'); +const http = require('node:http'); +const express = require('express'); + +process.env.NODE_ENV = 'test'; + +const rateLimit = require('../src/middleware/rateLimit'); +const requestId = require('../src/middleware/requestId'); +const { + fingerprint, + resolveClientIp, + resolveActorKey, +} = require('../src/utils/clientIdentity'); +const { sanitizeCorrelationId } = require('../src/middleware/requestId'); +const errorHandler = require('../src/middleware/errorHandler'); +const ApiError = require('../src/utils/ApiError'); + +function mockReq(overrides = {}) { + const headers = { ...(overrides.headers || {}) }; + return { + ip: overrides.ip || '127.0.0.1', + socket: { remoteAddress: overrides.remoteAddress || '127.0.0.1' }, + token: overrides.token, + adminToken: overrides.adminToken, + get(name) { + const key = Object.keys(headers).find( + (k) => k.toLowerCase() === String(name).toLowerCase() + ); + return key ? headers[key] : undefined; + }, + ...overrides, + }; +} + +function mockRes() { + const headers = {}; + return { + headers, + set(name, value) { + headers[name] = String(value); + }, + statusCode: 200, + status(code) { + this.statusCode = code; + return this; + }, + json(body) { + this.body = body; + return this; + }, + }; +} + +function run(middleware, req, res) { + return new Promise((resolve, reject) => { + middleware(req, res, (err) => { + if (err) { + resolve({ err, res }); + } else { + resolve({ err: null, res }); + } + }); + }); +} + +test('fingerprint is stable, truncated, and never echoes the secret', () => { + const a = fingerprint('test-token-admin'); + const b = fingerprint('test-token-admin'); + assert.equal(a, b); + assert.equal(a.length, 16); + assert.ok(!a.includes('test-token')); + assert.equal(fingerprint(''), 'anon'); +}); + +test('resolveClientIp ignores X-Forwarded-For unless trustProxy is on', () => { + const req = mockReq({ + remoteAddress: '10.0.0.5', + // Express trust proxy 1 resolves the nearest untrusted hop. + ip: '10.0.0.1', + headers: { 'X-Forwarded-For': '203.0.113.9, 10.0.0.1' }, + }); + assert.equal(resolveClientIp(req, { trustProxy: false }), '10.0.0.5'); + assert.equal(resolveClientIp(req, { trustProxy: true }), '10.0.0.1'); +}); + +test('resolveActorKey prefers token fingerprint over IP', () => { + const req = mockReq({ + token: 'test-token-admin', + remoteAddress: '10.0.0.5', + }); + const key = resolveActorKey(req); + assert.match(key, /^actor:/); + assert.ok(!key.includes('test-token-admin')); + assert.ok(!key.includes('10.0.0.5')); +}); + +test('sanitizeCorrelationId rejects oversized or unsafe values', () => { + assert.equal(sanitizeCorrelationId('abc-123'), 'abc-123'); + assert.equal(sanitizeCorrelationId(' ok_id '), 'ok_id'); + assert.equal(sanitizeCorrelationId('has space'), null); + assert.equal(sanitizeCorrelationId('bad"quote'), null); + assert.equal(sanitizeCorrelationId('x'.repeat(200)), null); + assert.equal(sanitizeCorrelationId(''), null); +}); + +test('burst over max returns 429 with Retry-After and policy details', async () => { + const limiter = rateLimit({ + name: 'mutation:transfers', + windowMs: 60_000, + max: 3, + forceInTest: true, + keyGenerator: () => 'actor:one', + }); + + for (let i = 0; i < 3; i += 1) { + const { err } = await run(limiter, mockReq(), mockRes()); + assert.equal(err, null); + } + + const { err, res } = await run(limiter, mockReq(), mockRes()); + assert.ok(err instanceof ApiError); + assert.equal(err.statusCode, 429); + assert.equal(err.details.policy, 'mutation:transfers'); + assert.equal(err.details.limit, 3); + assert.ok(err.details.retryAfter >= 1); + assert.equal(res.headers['Retry-After'], String(err.details.retryAfter)); + assert.equal(res.headers['X-RateLimit-Remaining'], '0'); +}); + +test('identity isolation: one actor bursting does not block another', async () => { + const limiter = rateLimit({ + name: 'mutation:transfers', + windowMs: 60_000, + max: 2, + forceInTest: true, + keyGenerator: (req) => resolveActorKey(req), + }); + + const actorA = mockReq({ token: 'token-a' }); + const actorB = mockReq({ token: 'token-b' }); + + assert.equal((await run(limiter, actorA, mockRes())).err, null); + assert.equal((await run(limiter, actorA, mockRes())).err, null); + const blocked = await run(limiter, actorA, mockRes()); + assert.equal(blocked.err.statusCode, 429); + + const other = await run(limiter, actorB, mockRes()); + assert.equal(other.err, null); +}); + +test('proxy trust: forged X-Forwarded-For cannot rotate identity when trustProxy is false', async () => { + const limiter = rateLimit({ + name: 'global', + windowMs: 60_000, + max: 2, + forceInTest: true, + trustProxy: false, + }); + + const first = mockReq({ + remoteAddress: '10.0.0.5', + headers: { 'X-Forwarded-For': '198.51.100.1' }, + }); + const second = mockReq({ + remoteAddress: '10.0.0.5', + headers: { 'X-Forwarded-For': '198.51.100.2' }, + }); + + assert.equal((await run(limiter, first, mockRes())).err, null); + assert.equal((await run(limiter, first, mockRes())).err, null); + const blocked = await run(limiter, second, mockRes()); + assert.equal(blocked.err.statusCode, 429); +}); + +test('proxy trust: forged left-most X-Forwarded-For cannot rotate a one-hop identity', async () => { + const app = express(); + app.set('trust proxy', 1); + app.use(rateLimit({ + name: 'global', + windowMs: 60_000, + max: 1, + forceInTest: true, + trustProxy: true, + })); + app.get('/ping', (_req, res) => res.json({ ok: true })); + app.use(errorHandler); + + const server = http.createServer(app); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const { port } = server.address(); + + try { + async function ping(forwarded) { + return fetch(`http://127.0.0.1:${port}/ping`, { + headers: { 'X-Forwarded-For': forwarded }, + }); + } + const first = await ping('198.51.100.1, 203.0.113.9'); + assert.equal(first.status, 200); + const forged = await ping('198.51.100.2, 203.0.113.9'); + assert.equal(forged.status, 429); + const other = await ping('198.51.100.2, 203.0.113.10'); + assert.equal(other.status, 200); + } finally { + await new Promise((resolve) => server.close(resolve)); + } +}); + +test('maxKeys capacity preserves active budgets under identity flood', async () => { + const limiter = rateLimit({ + name: 'global', + windowMs: 60_000, + max: 1, + maxKeys: 5, + forceInTest: true, + keyGenerator: (req) => req.actor, + }); + + for (let i = 0; i < 5; i += 1) { + const req = mockReq(); + req.actor = `actor-${i}`; + assert.equal((await run(limiter, req, mockRes())).err, null); + } + + const newcomer = mockReq(); + newcomer.actor = 'actor-new'; + const full = await run(limiter, newcomer, mockRes()); + assert.equal(full.err.statusCode, 429); + assert.ok(Number(full.res.headers['Retry-After']) >= 1); + assert.equal(limiter.size(), 5); + + const original = mockReq(); + original.actor = 'actor-0'; + assert.equal((await run(limiter, original, mockRes())).err.statusCode, 429); +}); + +test('an unchanged full table has bounded traversal work during a rejection burst', async (t) => { + t.mock.method(Date, 'now', () => 1_000); + const NativeMap = global.Map; + let visited = 0; + function* countEntries(iterator) { + for (const entry of iterator) { + visited += 1; + yield entry; + } + } + class CountingMap extends NativeMap { + [Symbol.iterator]() { return countEntries(super[Symbol.iterator]()); } + entries() { return countEntries(super.entries()); } + values() { return countEntries(super.values()); } + keys() { return countEntries(super.keys()); } + forEach(callback, thisArg) { + return super.forEach((value, key) => { + visited += 1; + callback.call(thisArg, value, key, this); + }); + } + } + + const maxKeys = 256; + let limiter; + // Observe collection work only in this limiter. Restore the constructor + // before any asynchronous request so unrelated code keeps the native Map. + try { + global.Map = CountingMap; + limiter = rateLimit({ + windowMs: 60_000, + max: 1, + maxKeys, + forceInTest: true, + keyGenerator: (req) => req.actor, + }); + } finally { + global.Map = NativeMap; + } + + for (let i = 0; i < maxKeys; i += 1) { + assert.equal((await run(limiter, { actor: `known-${i}` }, mockRes())).err, null); + } + visited = 0; + for (let i = 0; i < 128; i += 1) { + const { err, res } = await run(limiter, { actor: `new-${i}` }, mockRes()); + assert.equal(err.statusCode, 429); + assert.equal(res.headers['Retry-After'], '60'); + } + assert.equal(limiter.size(), maxKeys); + assert.equal((await run(limiter, { actor: 'known-0' }, mockRes())).err.statusCode, 429); + // Permit one lazy scan, but reject repeated work proportional to the table + // size for a burst that cannot release or consume any slot. + assert.ok(visited <= maxKeys, `${visited} entries visited during unchanged-table rejection`); +}); + +test('capacity retry deadlines round up and release only expired budgets', async (t) => { + let now = 1_000; + t.mock.method(Date, 'now', () => now); + const limiter = rateLimit({ + windowMs: 10_000, max: 1, maxKeys: 2, forceInTest: true, + keyGenerator: (req) => req.actor, + }); + const request = (actor) => run(limiter, { actor }, mockRes()); + assert.equal((await request('a')).err, null); + now = 5_000; + assert.equal((await request('b')).err, null); + + now = 9_500; + let blocked = await request('c'); + assert.equal(blocked.err.statusCode, 429); + assert.equal(blocked.res.headers['Retry-After'], '2'); + now = 10_001; + blocked = await request('c'); + assert.equal(blocked.res.headers['Retry-After'], '1'); + + now = 11_000; + assert.equal((await request('a')).err, null); + assert.equal(limiter.size(), 2); + blocked = await request('b'); + assert.equal(blocked.err.statusCode, 429); + assert.equal(blocked.res.headers['Retry-After'], '4'); + now = 15_000; + assert.equal((await request('c')).err, null); + assert.equal((await request('a')).err.statusCode, 429); + assert.equal(limiter.size(), 2); +}); + +test('renewal below capacity and reset keep retry deadlines correct after clock rollback', async (t) => { + let now = 1_000; + t.mock.method(Date, 'now', () => now); + const limiter = rateLimit({ + windowMs: 10_000, max: 1, maxKeys: 2, forceInTest: true, + keyGenerator: (req) => req.actor, + }); + const request = (actor) => run(limiter, { actor }, mockRes()); + assert.equal((await request('a')).err, null); + now = 11_000; + assert.equal((await request('a')).err, null); + assert.equal((await request('b')).err, null); + now = 10_000; + let blocked = await request('c'); + assert.equal(blocked.err.statusCode, 429); + assert.equal(blocked.res.headers['Retry-After'], '11'); + + now = 20_000; + limiter.reset(); + assert.equal(limiter.size(), 0); + assert.equal((await request('a')).err, null); + assert.equal((await request('b')).err, null); + blocked = await request('c'); + assert.equal(blocked.err.statusCode, 429); + assert.equal(blocked.res.headers['Retry-After'], '10'); +}); + +test('a newly admitted earlier window becomes the capacity deadline after clock rollback', async (t) => { + let now = 11_000; + t.mock.method(Date, 'now', () => now); + const limiter = rateLimit({ + windowMs: 10_000, max: 1, maxKeys: 2, forceInTest: true, + keyGenerator: (req) => req.actor, + }); + const request = (actor) => run(limiter, { actor }, mockRes()); + assert.equal((await request('a')).err, null); + now = 1_000; + assert.equal((await request('b')).err, null); + const blocked = await request('c'); + assert.equal(blocked.err.statusCode, 429); + assert.equal(blocked.res.headers['Retry-After'], '10'); + now = 11_000; + assert.equal((await request('c')).err, null); + assert.equal((await request('a')).err.statusCode, 429); +}); + +test('correlation id is echoed on success and on 429 without leaking tokens', async () => { + const app = express(); + app.use(requestId); + app.use( + rateLimit({ + name: 'mutation:quote', + windowMs: 60_000, + max: 1, + forceInTest: true, + keyGenerator: () => 'ip:test', + }) + ); + app.get('/quote', (req, res) => { + res.json({ ok: true, correlationId: req.correlationId }); + }); + app.use(errorHandler); + + const server = http.createServer(app); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const { port } = server.address(); + const base = `http://127.0.0.1:${port}`; + + try { + const first = await fetch(`${base}/quote`, { + headers: { 'X-Correlation-Id': 'client-corr-1' }, + }); + assert.equal(first.status, 200); + assert.equal(first.headers.get('x-correlation-id'), 'client-corr-1'); + assert.equal(first.headers.get('x-request-id'), 'client-corr-1'); + const firstBody = await first.json(); + assert.equal(firstBody.correlationId, 'client-corr-1'); + + const second = await fetch(`${base}/quote`, { + headers: { + 'X-Correlation-Id': 'client-corr-2', + Authorization: 'Bearer test-token-admin', + }, + }); + assert.equal(second.status, 429); + assert.equal(second.headers.get('x-correlation-id'), 'client-corr-2'); + assert.ok(second.headers.get('retry-after')); + const body = await second.json(); + assert.equal(body.error.status, 429); + assert.equal(body.error.requestId, 'client-corr-2'); + assert.equal(body.error.details.policy, 'mutation:quote'); + const dumped = JSON.stringify(body); + assert.ok(!dumped.includes('test-token-admin')); + } finally { + await new Promise((resolve) => server.close(resolve)); + } +}); + +test('unsafe inbound correlation id is replaced, not echoed', async () => { + const app = express(); + app.use(requestId); + app.get('/ping', (req, res) => res.json({ id: req.id })); + const server = http.createServer(app); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const { port } = server.address(); + + try { + const res = await fetch(`http://127.0.0.1:${port}/ping`, { + headers: { 'X-Request-Id': 'not a safe id' }, + }); + const body = await res.json(); + assert.notEqual(body.id, 'not a safe id'); + assert.equal(res.headers.get('x-request-id'), body.id); + assert.equal(res.headers.get('x-correlation-id'), body.id); + } finally { + await new Promise((resolve) => server.close(resolve)); + } +}); + +test('actual app counts parser failures in the global budget and correlates errors', async () => { + const config = require('../src/config'); + const createApp = require('../src/app'); + const originalMax = config.rateLimit.max; + const originalBodyLimit = config.bodyLimit; + const originalEnableInTest = process.env.ENABLE_RATE_LIMIT_IN_TEST; + let server; + + try { + config.rateLimit.max = 2; + config.bodyLimit = '32b'; + process.env.ENABLE_RATE_LIMIT_IN_TEST = '1'; + server = http.createServer(createApp()); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const base = `http://127.0.0.1:${server.address().port}`; + const malformed = '{"amount":'; + const oversized = JSON.stringify({ note: 'x'.repeat(64) }); + const responses = []; + + for (const [index, payload] of [malformed, oversized, malformed, oversized].entries()) { + const id = `parser-budget-${index}`; + const response = await fetch(`${base}/api/transfers`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', 'X-Request-Id': id }, + body: payload, + }); + responses.push({ id, response, body: await response.json() }); + } + + assert.deepEqual(responses.map(({ response }) => response.status), [400, 413, 429, 429]); + for (const { id, response, body } of responses) { + assert.equal(response.headers.get('x-request-id'), id); + assert.equal(response.headers.get('x-correlation-id'), id); + assert.equal(body.error.requestId, id); + assert.equal(response.headers.get('x-ratelimit-policy'), 'global'); + if (response.status === 429) { + assert.ok(Number(response.headers.get('retry-after')) >= 1); + assert.equal(response.headers.get('x-ratelimit-remaining'), '0'); + assert.equal(body.error.details.policy, 'global'); + } + } + } finally { + config.rateLimit.max = originalMax; + config.bodyLimit = originalBodyLimit; + if (originalEnableInTest === undefined) delete process.env.ENABLE_RATE_LIMIT_IN_TEST; + else process.env.ENABLE_RATE_LIMIT_IN_TEST = originalEnableInTest; + if (server) { + server.closeAllConnections(); + await new Promise((resolve) => server.close(resolve)); + } + } +});