diff --git a/listener/.gitignore b/listener/.gitignore index caf428fe..36a4e844 100644 --- a/listener/.gitignore +++ b/listener/.gitignore @@ -3,3 +3,8 @@ dist/ .env *.log data/ + +# Generated load-test reports. These are environment-specific; keep a curated +# baseline (reports/load/baseline.json) committed so runs stay comparable. +reports/load/*.json +!reports/load/baseline.json diff --git a/listener/LOAD_TESTING.md b/listener/LOAD_TESTING.md new file mode 100644 index 00000000..2ce8af3d --- /dev/null +++ b/listener/LOAD_TESTING.md @@ -0,0 +1,192 @@ +# API Load Testing Workflow + +## Overview + +`npm run load-test` is a repeatable workflow for measuring the throughput and +latency of NotifyChain's critical API endpoints. It exists so a change can be +measured against the same scenarios, on the same machine, and compared to a +saved baseline instead of relying on anecdotes. + +The workflow covers the three things issue #860 asks for: + +| Requirement | How it is met | +| --- | --- | +| Test scenarios are documented | [`load-test.config.json`](./load-test.config.json) declares each scenario; [`load-test.config.schema.md`](./load-test.config.schema.md) documents the format | +| RPS and latency are measurable | Every run reports requests/second, p50/p90/p95/p99/max latency, and error rate per scenario and in aggregate | +| Results are comparable across changes | Reports are schema-versioned JSON written to `reports/load/`; runs can be diffed against a saved baseline, and gated in CI | + +It complements, rather than replaces, the existing suites: + +- `npm run test:load` — unit tests for the measurement core (metric maths, + threshold gating, baseline comparison). Fast and socket-free. +- `npm run test` — the functional suite, including the rate-limit behaviour + tests added in issue #852 (`src/api/rate-limit-scenarios.test.ts`). +- `npm run test:stress` / `npm run stress-test` — queue/event-processing stress + tests, which exercise the internal pipeline rather than HTTP endpoints. + +## Two ways to run + +| Command | Target | When to use | +| --- | --- | --- | +| `npm run load-test` | In-process `createEventsServer` on an ephemeral port | Local runs and CI, where there is no deployed listener | +| `npm run load-test:external -- --url ` | An already-running listener | Staging/production-like environments | + +Both drive the same scenarios, the same measurement core and the same report +format — only the target differs. + +> **Why the in-process run is a Jest spec.** The listener's third-party runtime +> dependencies (`@stellar/stellar-sdk`, `node-cache`, `uuid`) are not installed; +> `jest.config.js` maps them to test doubles. Jest is therefore the supported way +> to execute the API in-process, so `npm run load-test` drives +> [`src/__tests__/load-test.workflow.test.ts`](./src/__tests__/load-test.workflow.test.ts) +> with `LOAD_TEST=1`. The spec is skipped in the normal `npm test` run so the +> suite stays fast and side-effect free. + +## Quick start + +```bash +cd listener + +# 1. In-process run of the documented scenarios (CI-friendly, no server needed). +npm run load-test + +# 2. Against an already-running listener (e.g. `npm run dev`). +npm run load-test:external -- --url http://127.0.0.1:3000 + +# 3. Save a baseline, then compare later runs against it. +npm run load-test:external -- --url http://127.0.0.1:3000 --out reports/load/baseline.json +npm run load-test:external -- --url http://127.0.0.1:3000 --baseline reports/load/baseline.json + +# 4. Gate a change in CI (non-zero exit on regression). +npm run load-test:external -- --url http://127.0.0.1:3000 --baseline reports/load/baseline.json --fail-on-regression +``` + +No extra tooling (`k6`, `autocannon`, …) is required, so the workflow runs in the +same environment as the rest of the test suite. + +## Scenarios + +Scenarios are declared in [`load-test.config.json`](./load-test.config.json) and +cover the endpoints that are hot on the read path: + +| Scenario | Endpoint | Why it matters | +| --- | --- | --- | +| `status` | `GET /api/status` | Liveness/readiness probe; hit by every load balancer | +| `events` | `GET /api/events` | Event registry read; the hottest dashboard path | +| `analytics` | `GET /api/analytics` | Analytics snapshot read on every dashboard refresh | +| `rate-limit-metrics` | `GET /api/rate-limit/metrics` | Observability endpoint that must stay fast while clients are throttled | + +Add a scenario by appending an entry to the `scenarios` array — see the +[schema doc](./load-test.config.schema.md) for every field. Scenario `name`s are +the keys used to match results across reports, so keep them stable. + +## What is measured + +For each scenario, and across the whole run: + +- **Throughput (RPS)** — completed requests divided by the measured phase wall + time. +- **Latency** — nearest-rank percentiles (p50, p90, p95, p99) plus min/mean/max. + Sampled with `process.hrtime.bigint()`, which is monotonic and unaffected by + wall-clock changes. +- **Error rate** — non-2xx responses divided by total responses. Client (4xx) + and server (5xx) errors are counted separately, so a `429` from rate limiting + is not confused with a `500`. +- **Environment** — node version, platform, CPU count and total memory, so a + report is self-describing when compared later or on another machine. + +Each scenario runs a short warm-up phase (excluded from the numbers) before the +measured phase, and scenarios run **sequentially** so they don't contend for the +same event loop or server capacity. Both choices keep runs comparable. + +### Thresholds + +Thresholds live in the config and are evaluated on every run: + +```json +"thresholds": { + "maxErrorRate": 0.01, + "maxP95Ms": 250, + "minThroughputRps": 25 +} +``` + +A breached threshold is listed in the report's `failures` array and the CLI exits +with code `1` (the in-process workflow fails the assertion instead). + +## Baseline comparison + +A baseline is just a previously written report: + +```bash +npm run load-test:external -- --url http://127.0.0.1:3000 --out reports/load/baseline.json +``` + +Comparing later runs reports, per scenario, the absolute and relative change in +p95 latency, throughput and error rate. A scenario whose p95 grows by more than +the regression tolerance (default 10%) is flagged as a regression. With +`--fail-on-regression` the workflow exits non-zero so it can be used as a gate. + +For the in-process workflow, drop the report at +`reports/load/baseline.json`; the spec prints the comparison, and with +`LOAD_TEST_GATE=1` it asserts no regression. + +`reports/load/*.json` is git-ignored; commit `reports/load/baseline.json` +explicitly if you want the team to share one reference point. + +## CLI reference + +``` +npm run load-test:external -- [options] + + --config Scenario/threshold config (default: load-test.config.json) + --url Base URL of the listener to test (default: target.baseUrl) + --out Where to write the JSON report (default: reports/load/latest.json) + --baseline Compare against a previously saved report + --fail-on-regression Exit non-zero when the baseline comparison regresses + --concurrency Override concurrency for every scenario + --duration Override durationMs for every scenario + --quiet Only print the report + --help Show this help +``` + +**Exit codes** + +| Code | Meaning | +| --- | --- | +| `0` | Thresholds passed (and, when gating, no baseline regression) | +| `1` | A threshold was breached, or a regression was detected while `--fail-on-regression` is set | +| `2` | The workflow could not run (bad config, bad arguments) | + +## Using it in CI + +```bash +cd listener +npm ci +npm run load-test # in-process, no external dependencies +npm run load-test:external -- --url "$LISTENER_URL" --baseline reports/load/baseline.json --fail-on-regression +``` + +The in-process run starts the events server on an ephemeral port, so no external +services are needed. Rate limiting is **disabled** in the default config so the +run measures raw API capacity; enable it under `target.rateLimit` to load-test +the throttling path itself (see [RATE-LIMITING-GUIDE.md](./RATE-LIMITING-GUIDE.md)). + +Because absolute numbers depend on the machine, treat the baseline comparison as +the signal and absolute thresholds as a coarse sanity check. Run the baseline and +the candidate on the same runner for a meaningful diff. + +## Interpreting results + +- **Throughput stable, p95 flat** — no measurable regression; ship it. +- **p95 up > 10%, RPS down** — likely a hot-path change; the comparison output + names the scenario that regressed. +- **Error rate > threshold** — inspect the report's `clientErrors`/`serverErrors` + split before assuming a regression. + +## Follow-ups + +- Publish a shared, machine-independent baseline (e.g. from a dedicated CI + runner) if the team wants absolute gating rather than relative comparison. +- Extend the scenario set to mutating endpoints (`POST /api/webhooks`) with + authenticated fixtures, which needs request signing (#491) to be wired in. diff --git a/listener/load-test.config.json b/listener/load-test.config.json new file mode 100644 index 00000000..6af7cfd0 --- /dev/null +++ b/listener/load-test.config.json @@ -0,0 +1,59 @@ +{ + "$schema": "./load-test.config.schema.md", + "target": { + "baseUrl": "http://127.0.0.1:3000", + "rateLimit": { + "enabled": false, + "windowMs": 60000, + "maxRequests": 1000000, + "clientOverrides": {} + } + }, + "scenarios": [ + { + "name": "status", + "description": "Service status - the liveness/readiness probe used by load balancers.", + "method": "GET", + "path": "/api/status", + "concurrency": 10, + "durationMs": 3000, + "warmupMs": 500, + "maxRequests": 3000 + }, + { + "name": "events", + "description": "Event registry read - the hottest read path for the dashboard.", + "method": "GET", + "path": "/api/events", + "concurrency": 10, + "durationMs": 3000, + "warmupMs": 500, + "maxRequests": 3000 + }, + { + "name": "analytics", + "description": "Analytics snapshot - read by the metrics dashboard on every refresh.", + "method": "GET", + "path": "/api/analytics", + "concurrency": 10, + "durationMs": 3000, + "warmupMs": 500, + "maxRequests": 3000 + }, + { + "name": "rate-limit-metrics", + "description": "Rate-limit observability endpoint - must stay fast while clients are throttled.", + "method": "GET", + "path": "/api/rate-limit/metrics", + "concurrency": 5, + "durationMs": 2000, + "warmupMs": 300, + "maxRequests": 1500 + } + ], + "thresholds": { + "maxErrorRate": 0.01, + "maxP95Ms": 250, + "minThroughputRps": 25 + } +} diff --git a/listener/load-test.config.schema.md b/listener/load-test.config.schema.md new file mode 100644 index 00000000..d9a642ee --- /dev/null +++ b/listener/load-test.config.schema.md @@ -0,0 +1,73 @@ +# `load-test.config.json` schema + +Human-readable reference for the config consumed by +[`src/scripts/load-test.ts`](./src/scripts/load-test.ts) (issue #860). Types live +in [`src/utils/load-test-runner.ts`](./src/utils/load-test-runner.ts) +(`LoadTestConfig`, `LoadTestScenario`, `LoadTestThresholds`). + +## Top level + +| Field | Type | Required | Description | +| --- | --- | --- | --- | +| `$schema` | string | no | Point at this document for editor hints | +| `target` | object | no | Defaults for the target server | +| `scenarios` | object[] | **yes** | One or more scenarios to execute | +| `thresholds` | object | **yes** | Pass/fail gates evaluated on every run | + +## `target` + +| Field | Type | Default | Description | +| --- | --- | --- | --- | +| `baseUrl` | string | `http://127.0.0.1:3000` | Default base URL for the external CLI when `--url` is omitted | +| `rateLimit` | object | disabled | `RateLimitConfig` applied to the in-process server by the Jest workflow (`enabled`, `windowMs`, `maxRequests`, `clientOverrides`) | + +Rate limiting is disabled by default so the workflow measures raw API capacity. +Enable it to load-test the throttling path. + +## `scenarios[]` + +| Field | Type | Required | Description | +| --- | --- | --- | --- | +| `name` | string | **yes** | Stable identifier; the key used to match scenarios across reports | +| `description` | string | no | Human-readable note on which critical path the scenario covers | +| `method` | `GET` \| `POST` \| `PUT` \| `DELETE` | **yes** | HTTP method | +| `path` | string | **yes** | Request path, e.g. `/api/status` | +| `headers` | object | no | Extra request headers | +| `body` | string | no | Already-serialized request body | +| `concurrency` | number | **yes** | In-flight requests during the phase | +| `durationMs` | number | **yes** | Length of the measured phase, in milliseconds | +| `warmupMs` | number | no | Warm-up phase, executed but excluded from the report | +| `maxRequests` | number | no | Hard cap on measured requests (AND-ed with `durationMs`); useful for bounded/CI runs | + +## `thresholds` + +| Field | Type | Required | Description | +| --- | --- | --- | --- | +| `maxErrorRate` | number | **yes** | Maximum tolerated fraction of non-2xx responses (`0.01` = 1%) | +| `maxP95Ms` | number | **yes** | Maximum tolerated p95 latency per scenario, in ms | +| `minThroughputRps` | number | no | Minimum aggregate throughput, in requests/second | + +## Example + +```json +{ + "$schema": "./load-test.config.schema.md", + "target": { + "baseUrl": "http://127.0.0.1:3000", + "rateLimit": { "enabled": false, "windowMs": 60000, "maxRequests": 1000000, "clientOverrides": {} } + }, + "scenarios": [ + { + "name": "status", + "description": "Liveness/readiness probe used by load balancers.", + "method": "GET", + "path": "/api/status", + "concurrency": 10, + "durationMs": 3000, + "warmupMs": 500, + "maxRequests": 3000 + } + ], + "thresholds": { "maxErrorRate": 0.01, "maxP95Ms": 250, "minThroughputRps": 25 } +} +``` diff --git a/listener/package.json b/listener/package.json index 8099b5d3..fd82918f 100644 --- a/listener/package.json +++ b/listener/package.json @@ -14,6 +14,9 @@ "test": "node ./node_modules/jest/bin/jest.js", "check:event-schema": "ts-node src/schema/compatibility-check.ts", "test:stress": "node ./node_modules/jest/bin/jest.js src/__tests__/stress.test.ts --runInBand --detectOpenHandles", + "test:load": "node ./node_modules/jest/bin/jest.js src/__tests__/load-test-runner.test.ts --runInBand", + "load-test": "LOAD_TEST=1 node ./node_modules/jest/bin/jest.js src/__tests__/load-test.workflow.test.ts --runInBand", + "load-test:external": "ts-node --transpile-only src/scripts/load-test.ts", "stress-test": "ts-node src/scripts/run-stress-tests.ts", "migrate": "ts-node src/scripts/migrate-db.ts", "migrate:templates": "ts-node src/scripts/migrate-templates.ts", diff --git a/listener/src/__tests__/load-test-runner.test.ts b/listener/src/__tests__/load-test-runner.test.ts new file mode 100644 index 00000000..023e5620 --- /dev/null +++ b/listener/src/__tests__/load-test-runner.test.ts @@ -0,0 +1,252 @@ +/** + * Load-testing workflow tests (issue #860) + * + * These exercise the measurement core directly with an injected probe so the + * behaviour is deterministic and socket-free: percentile maths, request/error + * accounting, threshold gating and baseline comparison. + */ + +import { describe, it, expect, jest } from '@jest/globals'; + +import { + LOAD_REPORT_SCHEMA_VERSION, + compareLoadReports, + formatComparison, + formatLoadReport, + percentile, + runLoadTest, + runScenario, + summarizeLatencies, + type LoadReport, + type LoadTestProbe, + type LoadTestScenario, +} from '../utils/load-test-runner'; + +/** Probe that always succeeds and reports a fixed latency. */ +function constantProbe(status = 200, durationMs = 5): LoadTestProbe { + return jest.fn(async () => ({ status, durationMs })); +} + +function scenario(overrides: Partial = {}): LoadTestScenario { + return { + name: 'status', + method: 'GET', + path: '/api/status', + concurrency: 2, + durationMs: 5_000, + maxRequests: 8, + ...overrides, + }; +} + +describe('percentile (#860)', () => { + it('returns 0 for an empty sample', () => { + expect(percentile([], 95)).toBe(0); + }); + + it('uses nearest-rank for typical percentiles', () => { + const sorted = Array.from({ length: 100 }, (_, i) => i + 1); + expect(percentile(sorted, 50)).toBe(50); + expect(percentile(sorted, 90)).toBe(90); + expect(percentile(sorted, 95)).toBe(95); + expect(percentile(sorted, 99)).toBe(99); + }); + + it('clamps the boundaries to the sample', () => { + const sorted = [3, 9, 12]; + expect(percentile(sorted, 0)).toBe(3); + expect(percentile(sorted, 100)).toBe(12); + }); +}); + +describe('summarizeLatencies (#860)', () => { + it('reports min/mean/percentiles/max from unsorted input', () => { + const stats = summarizeLatencies([40, 10, 30, 20]); + expect(stats.min).toBe(10); + expect(stats.max).toBe(40); + expect(stats.mean).toBe(25); + expect(stats.p50).toBe(20); + expect(stats.p99).toBe(40); + }); + + it('returns zeroed stats for no samples', () => { + expect(summarizeLatencies([])).toEqual({ + min: 0, + mean: 0, + p50: 0, + p90: 0, + p95: 0, + p99: 0, + max: 0, + }); + }); +}); + +describe('runScenario (#860)', () => { + it('measures exactly maxRequests and counts successes', async () => { + const probe = constantProbe(200, 4); + const result = await runScenario(scenario({ concurrency: 4, maxRequests: 20 }), probe); + + expect(result.requests).toBe(20); + expect(result.success).toBe(20); + expect(result.errorRate).toBe(0); + expect(result.throughputRps).toBeGreaterThan(0); + expect(result.latencyMs.p50).toBe(4); + expect(probe).toHaveBeenCalledTimes(20); + }); + + it('separates client (4xx) and server (5xx) errors', async () => { + let call = 0; + const probe: LoadTestProbe = jest.fn(async () => { + call += 1; + const status = call % 4 === 0 ? 500 : call % 2 === 0 ? 429 : 200; + return { status, durationMs: 1 }; + }); + + const result = await runScenario(scenario({ concurrency: 1, maxRequests: 8 }), probe); + + expect(result.requests).toBe(8); + expect(result.success + result.clientErrors + result.serverErrors).toBe(8); + expect(result.serverErrors).toBeGreaterThan(0); + expect(result.clientErrors).toBeGreaterThan(0); + expect(result.errorRate).toBeGreaterThan(0); + }); + + it('does not count warm-up requests in the report', async () => { + const probe = constantProbe(200, 1); + const result = await runScenario( + scenario({ concurrency: 2, durationMs: 20, warmupMs: 20, maxRequests: undefined }), + probe, + ); + + // The report only reflects the measured phase, so it can't distinguish + // warm-up calls; but the probe must have been called at least once. + expect(result.requests).toBeGreaterThan(0); + expect(probe).toHaveBeenCalled(); + }); +}); + +describe('runLoadTest threshold gating (#860)', () => { + const config = { + scenarios: [scenario({ name: 'a', maxRequests: 10 }), scenario({ name: 'b', maxRequests: 10 })], + thresholds: { maxErrorRate: 0.05, maxP95Ms: 100 }, + }; + + it('passes when every threshold is satisfied and produces a schemaVersioned report', async () => { + const report = await runLoadTest(config, constantProbe(200, 10)); + expect(report.schemaVersion).toBe(LOAD_REPORT_SCHEMA_VERSION); + expect(report.passed).toBe(true); + expect(report.failures).toEqual([]); + expect(report.totals.requests).toBe(20); + expect(report.scenarios.map((s) => s.name)).toEqual(['a', 'b']); + }); + + it('fails when the aggregate error rate breaches the threshold', async () => { + const report = await runLoadTest(config, constantProbe(500, 10)); + expect(report.passed).toBe(false); + expect(report.totals.errorRate).toBe(1); + expect(report.failures.join(' ')).toContain('error rate'); + }); + + it('fails when a scenario p95 breaches the latency threshold', async () => { + const report = await runLoadTest(config, constantProbe(200, 500)); + expect(report.passed).toBe(false); + expect(report.failures.join(' ')).toContain('p95'); + }); + + it('fails when throughput falls below the minimum', async () => { + const report = await runLoadTest( + { ...config, thresholds: { ...config.thresholds, minThroughputRps: 10_000_000 } }, + constantProbe(200, 1), + ); + expect(report.passed).toBe(false); + expect(report.failures.join(' ')).toContain('throughput'); + }); +}); + +describe('compareLoadReports (#860)', () => { + function reportWith(name: string, p95: number, rps: number, errorRate = 0): LoadReport { + const scenarioResult = { + name, + method: 'GET' as const, + path: '/api/status', + concurrency: 1, + durationMs: 1000, + requests: 100, + success: 100, + clientErrors: 0, + serverErrors: 0, + errorRate, + throughputRps: rps, + latencyMs: { min: p95, mean: p95, p50: p95, p90: p95, p95, p99: p95, max: p95 }, + }; + return { + schemaVersion: LOAD_REPORT_SCHEMA_VERSION, + generatedAt: new Date(0).toISOString(), + environment: { node: 'v0', platform: 'test', release: 'test', cpus: 1, totalMemoryMb: 1 }, + scenarios: [scenarioResult], + totals: { requests: 100, errorRate, throughputRps: rps, latencyMs: scenarioResult.latencyMs }, + thresholds: { maxErrorRate: 0.1, maxP95Ms: 1000 }, + passed: true, + failures: [], + }; + } + + it('detects a latency regression beyond the tolerance', () => { + const comparison = compareLoadReports( + reportWith('status', 100, 1000), + reportWith('status', 150, 900), + ); + const scenarioComparison = comparison.scenarios[0]; + expect(scenarioComparison.status).toBe('compared'); + expect(scenarioComparison.p95Ms?.delta).toBe(50); + expect(scenarioComparison.p95Ms?.deltaPct).toBeCloseTo(0.5); + expect(comparison.passed).toBe(false); + expect(comparison.regressions.length).toBe(1); + }); + + it('treats an improvement as a pass', () => { + const comparison = compareLoadReports( + reportWith('status', 100, 1000), + reportWith('status', 80, 1200), + ); + expect(comparison.passed).toBe(true); + expect(comparison.regressions).toEqual([]); + }); + + it('honours a wider regression tolerance', () => { + const comparison = compareLoadReports( + reportWith('status', 100, 1000), + reportWith('status', 105, 1000), + 0.1, + ); + expect(comparison.passed).toBe(true); + }); + + it('reports added and removed scenarios', () => { + const comparison = compareLoadReports( + reportWith('status', 100, 1000), + reportWith('events', 100, 1000), + ); + const statuses = comparison.scenarios.map((s) => s.status).sort(); + expect(statuses).toEqual(['added', 'removed']); + }); +}); + +describe('report formatting (#860)', () => { + it('renders a readable report and comparison', async () => { + const report = await runLoadTest( + { + scenarios: [scenario({ name: 'status', maxRequests: 4 })], + thresholds: { maxErrorRate: 0.1, maxP95Ms: 1000 }, + }, + constantProbe(200, 3), + ); + const text = formatLoadReport(report); + expect(text).toContain('status'); + expect(text).toContain('Thresholds: PASS'); + + const comparison = compareLoadReports(report, report); + expect(formatComparison(comparison)).toContain('Comparison vs baseline'); + }); +}); diff --git a/listener/src/__tests__/load-test.workflow.test.ts b/listener/src/__tests__/load-test.workflow.test.ts new file mode 100644 index 00000000..fdb2d58f --- /dev/null +++ b/listener/src/__tests__/load-test.workflow.test.ts @@ -0,0 +1,136 @@ +/** + * In-process API load-testing workflow (issue #860) + * + * This is the `npm run load-test` entry point: it boots the real + * `createEventsServer` on an ephemeral port and runs the scenarios documented + * in `load-test.config.json` through a real HTTP probe, then writes a + * schema-versioned report to `reports/load/latest.json` that can be compared + * against a committed baseline. + * + * It lives under Jest deliberately. The listener's third-party runtime + * dependencies (`@stellar/stellar-sdk`, `node-cache`, `uuid`) are not installed + * — they are mapped to test doubles by `jest.config.js` — so Jest is the + * supported way to execute the API in-process. The external CLI + * (`npm run load-test:external`) covers the "already-running server" case. + * + * Gated behind `LOAD_TEST=1` so the regular `npm test` suite stays fast and + * side-effect free. + */ + +import { describe, it, expect, beforeAll, afterAll } from '@jest/globals'; +import fs from 'fs'; +import http from 'http'; + +import { createEventsServer } from '../api/events-server'; +import type { RateLimitConfig } from '../types'; +import { loadLoadTestConfig } from '../utils/load-test-config'; +import { makeHttpProbe } from '../utils/load-test-http-probe'; +import { + DEFAULT_BASELINE_PATH, + DEFAULT_REPORT_PATH, + readLoadReport, + writeLoadReport, +} from '../utils/load-test-reports'; +import { + LOAD_REPORT_SCHEMA_VERSION, + compareLoadReports, + formatComparison, + formatLoadReport, + runLoadTest, + type LoadReport, +} from '../utils/load-test-runner'; + +jest.mock('../utils/logger', () => ({ + __esModule: true, + default: { + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + }, +})); + +jest.setTimeout(120_000); + +/** Rate limiting is off for a capacity run; the config can override this. */ +const DISABLED_RATE_LIMIT: RateLimitConfig = { + enabled: false, + windowMs: 60_000, + maxRequests: 1_000_000, + clientOverrides: {}, +}; + +const describeWorkflow = process.env.LOAD_TEST === '1' ? describe : describe.skip; + +describeWorkflow('API load-testing workflow (#860)', () => { + let server: http.Server | undefined; + let baseUrl = ''; + let report: LoadReport; + + beforeAll(async () => { + const config = loadLoadTestConfig(); + + server = createEventsServer({ + port: 0, + stellarRpcUrl: process.env.STELLAR_RPC_URL ?? 'https://soroban-testnet.stellar.org', + stellarNetworkPassphrase: 'Test SDF Network ; September 2015', + contractAddresses: [], + rateLimit: config.target?.rateLimit ?? DISABLED_RATE_LIMIT, + }); + + await new Promise((resolve) => server!.listen(0, '127.0.0.1', () => resolve())); + const address = server.address(); + if (!address || typeof address === 'string') throw new Error('Expected a TCP address'); + baseUrl = `http://127.0.0.1:${address.port}`; + + report = await runLoadTest(config, makeHttpProbe(baseUrl)); + writeLoadReport(DEFAULT_REPORT_PATH, report); + process.stdout.write(`\n${formatLoadReport(report)}\n`); + }); + + afterAll(async () => { + if (server) { + await new Promise((resolve) => server!.close(() => resolve())); + server = undefined; + } + }); + + it('executes every documented scenario and stays within thresholds', () => { + const { scenarios } = loadLoadTestConfig(); + + expect(report.schemaVersion).toBe(LOAD_REPORT_SCHEMA_VERSION); + expect(report.scenarios.map((scenario) => scenario.name)).toEqual( + scenarios.map((scenario) => scenario.name), + ); + expect(report.failures).toEqual([]); + expect(report.passed).toBe(true); + }); + + it('measures throughput and latency for the critical endpoints', () => { + for (const scenario of report.scenarios) { + expect(scenario.requests).toBeGreaterThan(0); + expect(scenario.throughputRps).toBeGreaterThan(0); + expect(scenario.latencyMs.p95).toBeGreaterThanOrEqual(0); + // Rate limiting is disabled, so nothing should be throttled. + expect(scenario.clientErrors).toBe(0); + expect(scenario.serverErrors).toBe(0); + } + expect(report.totals.requests).toBeGreaterThan(0); + expect(report.totals.throughputRps).toBeGreaterThan(0); + }); + + it('can be compared against a committed baseline', () => { + if (!fs.existsSync(DEFAULT_BASELINE_PATH)) { + // No baseline committed yet: write this run as one to start comparing. + process.stdout.write(`\nNo baseline at ${DEFAULT_BASELINE_PATH}; skipping comparison.\n`); + return; + } + + const baseline = readLoadReport(DEFAULT_BASELINE_PATH); + const comparison = compareLoadReports(baseline, report); + process.stdout.write(`\n${formatComparison(comparison)}\n`); + + if (process.env.LOAD_TEST_GATE === '1') { + expect(comparison.passed).toBe(true); + } + }); +}); diff --git a/listener/src/api/rate-limit-scenarios.test.ts b/listener/src/api/rate-limit-scenarios.test.ts new file mode 100644 index 00000000..b337132a --- /dev/null +++ b/listener/src/api/rate-limit-scenarios.test.ts @@ -0,0 +1,459 @@ +/** + * API rate-limit scenario tests (issue #852) + * + * Verifies that the rate limiter behaves predictably under the three traffic + * conditions called out in the issue: + * + * - normal – legitimate clients stay under their quota and are never blocked + * - burst – a sudden spike is capped at exactly `maxRequests` + * - repeated – sustained traffic is re-admitted once the window rolls over + * + * plus two cross-cutting guarantees: + * + * - the 429 response is a stable, machine-readable envelope + * - one client exhausting its quota never affects another client, and the + * observability routes stay reachable while throttled + * + * The unit-level cases drive `RateLimiter.handle` directly with fake + * request/response objects; the final integration block exercises the same + * behaviour end-to-end through a real `http.Server`. + */ + +import { describe, it, expect, beforeEach, afterEach, jest } from '@jest/globals'; +import http from 'http'; + +import { RateLimiter } from './rate-limiter'; +import { createEventsServer } from './events-server'; +import type { RateLimitConfig } from '../types'; + +jest.mock('../utils/logger', () => ({ + __esModule: true, + default: { + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + }, +})); + +interface CapturedResponse { + status: number; + headers: Record; + body: any; +} + +interface MockResponse { + res: http.ServerResponse; + captured: () => CapturedResponse; +} + +function makeRequest( + headers: Record = {}, + ip = '127.0.0.1', + url = '/api/events', + method = 'GET', +): http.IncomingMessage { + return { + headers, + socket: { remoteAddress: ip }, + url, + method, + } as unknown as http.IncomingMessage; +} + +/** Minimal ServerResponse stand-in that records the full wire response. */ +function makeResponse(): MockResponse { + const headers: Record = {}; + let status = 200; + let rawBody = ''; + + const res = { + setHeader: (name: string, value: unknown) => { + headers[name.toLowerCase()] = String(value); + }, + getHeader: (name: string) => headers[name.toLowerCase()], + writeHead: (code: number, extra?: Record) => { + status = code; + if (extra) { + for (const [key, value] of Object.entries(extra)) { + headers[key.toLowerCase()] = String(value); + } + } + }, + end: (chunk?: unknown) => { + if (chunk !== undefined) rawBody = String(chunk); + }, + } as unknown as http.ServerResponse; + + return { + res, + captured: () => ({ + status, + headers, + body: rawBody ? JSON.parse(rawBody) : undefined, + }), + }; +} + +const baseConfig = (overrides: Partial = {}): RateLimitConfig => ({ + enabled: true, + windowMs: 60_000, + maxRequests: 5, + clientOverrides: {}, + ...overrides, +}); + +describe('Rate limiting under normal load (#852)', () => { + let limiter: RateLimiter; + + afterEach(() => limiter?.destroy()); + + it('admits every request from distinct clients that stay under quota', async () => { + limiter = new RateLimiter(baseConfig({ maxRequests: 5 })); + + // Five independent API keys, each making four requests (< 5 quota). + const results: boolean[] = []; + for (let client = 0; client < 5; client++) { + const req = makeRequest({ 'x-api-key': `client-${client}` }); + for (let i = 0; i < 4; i++) { + results.push(await limiter.handle(req, makeResponse().res)); + } + } + + expect(results.every(Boolean)).toBe(true); + expect(results).toHaveLength(20); + + const metrics = limiter.getMetrics(); + expect(metrics.totalRequests).toBe(20); + expect(metrics.allowedRequests).toBe(20); + expect(metrics.blockedRequests).toBe(0); + expect(metrics.uniqueClients).toBe(5); + }); + + it('reports a decreasing remaining quota without ever going negative', async () => { + limiter = new RateLimiter(baseConfig({ maxRequests: 3 })); + const req = makeRequest({ 'x-api-key': 'normal-client' }); + + const remaining: string[] = []; + for (let i = 0; i < 3; i++) { + const { res, captured } = makeResponse(); + await limiter.handle(req, res); + remaining.push(captured().headers['x-ratelimit-remaining']); + } + + expect(remaining).toEqual(['2', '1', '0']); + expect(remaining.every((value) => Number(value) >= 0)).toBe(true); + }); + + it('does not block a client merely because others are active', async () => { + limiter = new RateLimiter(baseConfig({ maxRequests: 2 })); + + const chatty = makeRequest({ 'x-api-key': 'chatty' }); + await limiter.handle(chatty, makeResponse().res); + await limiter.handle(chatty, makeResponse().res); + expect(await limiter.handle(chatty, makeResponse().res)).toBe(false); + + // A separate client still has its full quota available. + const quiet = makeRequest({ 'x-api-key': 'quiet' }); + expect(await limiter.handle(quiet, makeResponse().res)).toBe(true); + expect(await limiter.handle(quiet, makeResponse().res)).toBe(true); + }); +}); + +describe('Rate limiting under burst load (#852)', () => { + let limiter: RateLimiter; + + afterEach(() => limiter?.destroy()); + + it('caps a concurrent burst at exactly maxRequests', async () => { + const maxRequests = 5; + const burstSize = 25; + limiter = new RateLimiter(baseConfig({ maxRequests, windowMs: 60_000 })); + + const req = makeRequest({ 'x-api-key': 'burst-client' }); + const responses = await Promise.all( + Array.from({ length: burstSize }, () => limiter.handle(req, makeResponse().res)), + ); + + const allowed = responses.filter(Boolean).length; + const blocked = responses.filter((allowed) => !allowed).length; + + expect(allowed).toBe(maxRequests); + expect(blocked).toBe(burstSize - maxRequests); + expect(allowed + blocked).toBe(burstSize); + + const metrics = limiter.getMetrics(); + expect(metrics.totalRequests).toBe(burstSize); + expect(metrics.allowedRequests).toBe(maxRequests); + expect(metrics.blockedRequests).toBe(burstSize - maxRequests); + }); + + it('returns an identical, predictable 429 for every burst rejection', async () => { + const maxRequests = 3; + limiter = new RateLimiter(baseConfig({ maxRequests, windowMs: 60_000 })); + const req = makeRequest({ 'x-api-key': 'burst-envelope' }); + + const rejections: CapturedResponse[] = []; + for (let i = 0; i < maxRequests + 10; i++) { + const { res, captured } = makeResponse(); + const allowed = await limiter.handle(req, res); + if (!allowed) rejections.push(captured()); + } + + expect(rejections).toHaveLength(10); + for (const response of rejections) { + expect(response.status).toBe(429); + expect(response.headers['x-ratelimit-limit']).toBe(String(maxRequests)); + expect(response.headers['x-ratelimit-remaining']).toBe('0'); + expect(response.headers['x-ratelimit-reset']).toMatch(/^\d+$/); + expect(Number(response.headers['retry-after'])).toBeGreaterThanOrEqual(1); + expect(response.body).toMatchObject({ + success: false, + error: { code: 'RATE_LIMITED' }, + }); + expect(response.body.error.message).toContain('Rate limit exceeded'); + } + }); + + it('keeps the burst allowance honest under overlapping clients', async () => { + limiter = new RateLimiter(baseConfig({ maxRequests: 4, windowMs: 60_000 })); + + const clients = ['alpha', 'beta', 'gamma']; + const perClient = 6; + + const responses = await Promise.all( + clients.flatMap((client) => + Array.from({ length: perClient }, () => + limiter.handle(makeRequest({ 'x-api-key': client }), makeResponse().res), + ), + ), + ); + + // Each client independently gets exactly 4 of its 6 requests. + for (let i = 0; i < clients.length; i++) { + const slice = responses.slice(i * perClient, (i + 1) * perClient); + expect(slice.filter(Boolean)).toHaveLength(4); + } + expect(limiter.getMetrics().blockedRequests).toBe(clients.length * (perClient - 4)); + }); +}); + +describe('Rate limiting under repeated load across windows (#852)', () => { + beforeEach(() => { + jest.useFakeTimers(); + jest.setSystemTime(new Date('2026-01-01T00:00:00.000Z')); + }); + + afterEach(() => { + jest.useRealTimers(); + }); + + it('re-admits a sustained client after the window rolls over', async () => { + const limiter = new RateLimiter(baseConfig({ maxRequests: 2, windowMs: 1_000 })); + const req = makeRequest({ 'x-api-key': 'sustained-client' }); + + try { + // Window 1: two allowed, the rest throttled. + expect(await limiter.handle(req, makeResponse().res)).toBe(true); + expect(await limiter.handle(req, makeResponse().res)).toBe(true); + expect(await limiter.handle(req, makeResponse().res)).toBe(false); + expect(limiter.getMetrics().blockedRequests).toBe(1); + + // Advance exactly one window: quota resets. + jest.advanceTimersByTime(1_000); + expect(await limiter.handle(req, makeResponse().res)).toBe(true); + expect(await limiter.handle(req, makeResponse().res)).toBe(true); + expect(await limiter.handle(req, makeResponse().res)).toBe(false); + + // Advance three windows at once: the client is fully re-admitted again. + jest.advanceTimersByTime(3_000); + expect(await limiter.handle(req, makeResponse().res)).toBe(true); + expect(await limiter.handle(req, makeResponse().res)).toBe(true); + + const metrics = limiter.getMetrics(); + // 5 allowed (2 + 2 + 1... actually 2+2+2) and 2 blocked across three windows. + expect(metrics.allowedRequests).toBe(6); + expect(metrics.blockedRequests).toBe(2); + } finally { + limiter.destroy(); + } + }); + + it('does not leak quota across windows for a client that only trickles requests', async () => { + const limiter = new RateLimiter(baseConfig({ maxRequests: 1, windowMs: 500 })); + const req = makeRequest({ 'x-api-key': 'trickle-client' }); + + try { + for (let window = 0; window < 4; window++) { + expect(await limiter.handle(req, makeResponse().res)).toBe(true); + // A second request in the same window is blocked. + expect(await limiter.handle(req, makeResponse().res)).toBe(false); + jest.advanceTimersByTime(500); + } + + const metrics = limiter.getMetrics(); + expect(metrics.allowedRequests).toBe(4); + expect(metrics.blockedRequests).toBe(4); + } finally { + limiter.destroy(); + } + }); +}); + +describe('Rate limiting client isolation and exemptions (#852)', () => { + let limiter: RateLimiter; + + afterEach(() => limiter?.destroy()); + + it('keeps API-key buckets and IP buckets independent', async () => { + limiter = new RateLimiter(baseConfig({ maxRequests: 1 })); + + const keyed = makeRequest({ 'x-api-key': 'keyed-client' }, '10.0.0.1'); + const anonymous = makeRequest({}, '10.0.0.2'); + + expect(await limiter.handle(keyed, makeResponse().res)).toBe(true); + expect(await limiter.handle(keyed, makeResponse().res)).toBe(false); + + // A different identity still gets its own allowance. + expect(await limiter.handle(anonymous, makeResponse().res)).toBe(true); + expect(await limiter.handle(anonymous, makeResponse().res)).toBe(false); + }); + + it('applies per-client overrides without affecting the default quota', async () => { + limiter = new RateLimiter( + baseConfig({ + maxRequests: 2, + clientOverrides: { + premium: { maxRequests: 5 }, + }, + }), + ); + + const premium = makeRequest({ 'x-api-key': 'premium' }); + for (let i = 0; i < 5; i++) { + expect(await limiter.handle(premium, makeResponse().res)).toBe(true); + } + expect(await limiter.handle(premium, makeResponse().res)).toBe(false); + + const standard = makeRequest({ 'x-api-key': 'standard' }); + expect(await limiter.handle(standard, makeResponse().res)).toBe(true); + expect(await limiter.handle(standard, makeResponse().res)).toBe(true); + expect(await limiter.handle(standard, makeResponse().res)).toBe(false); + }); + + it('treats a disabled limiter as a pass-through for all traffic', async () => { + limiter = new RateLimiter(baseConfig({ enabled: false, maxRequests: 1 })); + const req = makeRequest({ 'x-api-key': 'anything' }); + for (let i = 0; i < 50; i++) { + expect(await limiter.handle(req, makeResponse().res)).toBe(true); + } + }); +}); + +// --------------------------------------------------------------------------- +// End-to-end enforcement through a real HTTP server +// --------------------------------------------------------------------------- + +interface HttpResult { + status: number; + headers: http.IncomingHttpHeaders; + body: any; +} + +function httpGet( + port: number, + path: string, + headers: Record = {}, +): Promise { + return new Promise((resolve, reject) => { + const req = http.request({ host: '127.0.0.1', port, path, method: 'GET', headers }, (res) => { + let body = ''; + res.on('data', (chunk) => { + body += chunk; + }); + res.on('end', () => { + resolve({ + status: res.statusCode ?? 0, + headers: res.headers, + body: body ? JSON.parse(body) : undefined, + }); + }); + }); + req.on('error', reject); + req.end(); + }); +} + +describe('Rate limiting over HTTP (#852)', () => { + let server: http.Server | undefined; + + const startServer = async (rateLimit: RateLimitConfig): Promise => { + server = createEventsServer({ + port: 0, + stellarRpcUrl: 'https://soroban-testnet.stellar.org:443', + stellarNetworkPassphrase: 'Test SDF Network ; September 2015', + contractAddresses: [], + rateLimit, + }); + + await new Promise((resolve) => server!.listen(0, '127.0.0.1', () => resolve())); + const address = server!.address(); + if (!address || typeof address === 'string') throw new Error('Expected a TCP address'); + return address.port; + }; + + afterEach(async () => { + if (server) { + await new Promise((resolve) => server!.close(() => resolve())); + server = undefined; + } + }); + + it('enforces the quota end-to-end and returns the predictable 429 envelope', async () => { + const port = await startServer(baseConfig({ maxRequests: 2, windowMs: 60_000 })); + + const first = await httpGet(port, '/api/events', { 'x-api-key': 'e2e-client' }); + const second = await httpGet(port, '/api/events', { 'x-api-key': 'e2e-client' }); + const third = await httpGet(port, '/api/events', { 'x-api-key': 'e2e-client' }); + + expect(first.status).toBe(200); + expect(second.status).toBe(200); + expect(third.status).toBe(429); + expect(third.headers['x-ratelimit-limit']).toBe('2'); + expect(third.headers['x-ratelimit-remaining']).toBe('0'); + expect(Number(third.headers['retry-after'])).toBeGreaterThanOrEqual(1); + expect(third.body).toMatchObject({ success: false, error: { code: 'RATE_LIMITED' } }); + }); + + it('never blocks legitimate concurrent traffic from other clients', async () => { + const port = await startServer(baseConfig({ maxRequests: 3, windowMs: 60_000 })); + + const responses = await Promise.all([ + httpGet(port, '/api/events', { 'x-api-key': 'client-one' }), + httpGet(port, '/api/events', { 'x-api-key': 'client-two' }), + httpGet(port, '/api/events', { 'x-api-key': 'client-three' }), + httpGet(port, '/api/events', { 'x-api-key': 'client-one' }), + httpGet(port, '/api/events', { 'x-api-key': 'client-two' }), + ]); + + expect(responses.every((response) => response.status === 200)).toBe(true); + }); + + it('keeps /health and /api/rate-limit/metrics reachable while throttled', async () => { + const port = await startServer(baseConfig({ maxRequests: 1, windowMs: 60_000 })); + + const allowed = await httpGet(port, '/api/events', { 'x-api-key': 'exhausted' }); + const throttled = await httpGet(port, '/api/events', { 'x-api-key': 'exhausted' }); + expect(allowed.status).toBe(200); + expect(throttled.status).toBe(429); + + // Observability routes are exempt so callers can still diagnose throttling. + const metrics = await httpGet(port, '/api/rate-limit/metrics'); + expect(metrics.status).toBe(200); + expect(metrics.body.blockedRequests).toBeGreaterThanOrEqual(1); + + const health = await httpGet(port, '/health'); + // /health may report 200 or 503 depending on the (network-isolated) environment, + // but it must never be rate limited. + expect(health.status).not.toBe(429); + }); +}); diff --git a/listener/src/api/rate-limiter.test.ts b/listener/src/api/rate-limiter.test.ts index 9b3e6bdc..651e1975 100644 --- a/listener/src/api/rate-limiter.test.ts +++ b/listener/src/api/rate-limiter.test.ts @@ -1,4 +1,13 @@ -import { jest, describe, it, expect, beforeEach, afterEach, beforeAll, afterAll } from '@jest/globals'; +import { + jest, + describe, + it, + expect, + beforeEach, + afterEach, + beforeAll, + afterAll, +} from '@jest/globals'; import http from 'http'; import * as fs from 'fs'; import * as path from 'path'; @@ -31,7 +40,9 @@ const mockResponse = () => { let statusCode = 200; let body = ''; return { - setHeader: jest.fn().mockImplementation((name: any, val: any) => headers.set(name.toLowerCase(), String(val))), + setHeader: jest + .fn() + .mockImplementation((name: any, val: any) => headers.set(name.toLowerCase(), String(val))), writeHead: jest.fn().mockImplementation((code: any, h: any) => { statusCode = code; if (h) { @@ -102,7 +113,7 @@ describe('RateLimiter', () => { maxRequests: 5, clientOverrides: {}, }); - const req = mockRequest({ 'authorization': 'Bearer token-abc' }); + const req = mockRequest({ authorization: 'Bearer token-abc' }); const client = limiter.identifyClient(req); expect(client.clientId).toBe('token-abc'); @@ -189,8 +200,9 @@ describe('RateLimiter', () => { expect(res3._getHeaders().get('retry-after')).toBeDefined(); const body = JSON.parse(res3._getBody()); - expect(body.error).toBe('Too Many Requests'); - expect(body.message).toContain('Rate limit exceeded'); + expect(body.success).toBe(false); + expect(body.error.code).toBe('RATE_LIMITED'); + expect(body.error.message).toContain('Rate limit exceeded'); limiter.destroy(); }); @@ -249,7 +261,7 @@ describe('RateLimiter', () => { }); const req = mockRequest({}, '127.0.0.1'); - + // Initial metrics let metrics = limiter.getMetrics(); expect(metrics.totalRequests).toBe(0); @@ -381,7 +393,7 @@ describe('RateLimiter', () => { }); const req = mockRequest({ 'x-api-key': 'attacker-key' }, '8.8.8.8'); - + // Request 1: Allowed await limiter.handle(req, mockResponse()); @@ -396,12 +408,12 @@ describe('RateLimiter', () => { clientId: 'attacker...', clientType: 'API_KEY', endpoint: '/api/schedule', - }) + }), ); // Verify DB record // Need a small timeout to allow async DB insert to complete - await new Promise(resolve => setTimeout(resolve, 50)); + await new Promise((resolve) => setTimeout(resolve, 50)); const rows = await db.all('SELECT * FROM rate_limit_events'); expect(rows.length).toBe(1); @@ -426,7 +438,10 @@ describe('Events Server Rate Limiting Integration', () => { } }); - const makeRequest = (path: string, headers: Record = {}): Promise<{ status: number; headers: any }> => { + const makeRequest = ( + path: string, + headers: Record = {}, + ): Promise<{ status: number; headers: any }> => { return new Promise((resolve, reject) => { const req = http.request( { @@ -438,7 +453,7 @@ describe('Events Server Rate Limiting Integration', () => { }, (res) => { resolve({ status: res.statusCode!, headers: res.headers }); - } + }, ); req.on('error', reject); req.end(); @@ -497,29 +512,33 @@ describe('Events Server Rate Limiting Integration', () => { await makeRequest('/api/events'); // This one should be blocked // Fetch metrics - const metricsResponse = await new Promise<{ status: number; body: string }>((resolve, reject) => { - const req = http.request( - { - host: '127.0.0.1', - port, - path: '/api/rate-limit/metrics', - method: 'GET', - }, - (res) => { - let body = ''; - res.on('data', (chunk) => { body += chunk; }); - res.on('end', () => { - resolve({ status: res.statusCode!, body }); - }); - } - ); - req.on('error', reject); - req.end(); - }); + const metricsResponse = await new Promise<{ status: number; body: string }>( + (resolve, reject) => { + const req = http.request( + { + host: '127.0.0.1', + port, + path: '/api/rate-limit/metrics', + method: 'GET', + }, + (res) => { + let body = ''; + res.on('data', (chunk) => { + body += chunk; + }); + res.on('end', () => { + resolve({ status: res.statusCode!, body }); + }); + }, + ); + req.on('error', reject); + req.end(); + }, + ); expect(metricsResponse.status).toBe(200); const metrics = JSON.parse(metricsResponse.body); - + expect(metrics.totalRequests).toBeGreaterThanOrEqual(3); expect(metrics.allowedRequests).toBeGreaterThanOrEqual(2); expect(metrics.blockedRequests).toBeGreaterThanOrEqual(1); diff --git a/listener/src/config.ts b/listener/src/config.ts index 8f6417f7..2c882f39 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -61,7 +61,7 @@ function validateRequiredEnvVars(): void { if (missing.length > 0) { throw new ConfigError( `Missing required environment variable(s): ${missing.join(', ')}. ` + - 'Copy .env.example to .env and set them before starting the listener.' + 'Copy .env.example to .env and set them before starting the listener.', ); } } @@ -148,7 +148,9 @@ function validateContractAddresses(value: unknown): ContractConfig[] { return value.map((item, index) => { if (typeof item !== 'object' || item === null) { - throw new ConfigError(`CONTRACT_ADDRESSES[${index}] must be an object with address and events.`); + throw new ConfigError( + `CONTRACT_ADDRESSES[${index}] must be an object with address and events.`, + ); } const address = (item as any).address; @@ -160,7 +162,7 @@ function validateContractAddresses(value: unknown): ContractConfig[] { if (!Array.isArray(events) || events.some((event) => typeof event !== 'string')) { throw new ConfigError( - `CONTRACT_ADDRESSES[${index}].events must be an array of string event names.` + `CONTRACT_ADDRESSES[${index}].events must be an array of string event names.`, ); } @@ -202,9 +204,7 @@ function validateWebhookSecrets(value: unknown): WebhookSecret[] { return value.map((item, index) => { if (typeof item !== 'object' || item === null) { - throw new ConfigError( - `WEBHOOK_SECRETS[${index}] must be an object with id and secret.` - ); + throw new ConfigError(`WEBHOOK_SECRETS[${index}] must be an object with id and secret.`); } const id = (item as any).id; @@ -229,9 +229,7 @@ function validateApiKeys(value: unknown): ApiKey[] { return value.map((item, index) => { if (typeof item !== 'object' || item === null) { - throw new ConfigError( - `API_KEYS[${index}] must be an object with key (and optional name).` - ); + throw new ConfigError(`API_KEYS[${index}] must be an object with key (and optional name).`); } const key = (item as any).key; @@ -347,21 +345,27 @@ function loadExpirationConfig(): ExpirationConfig { const defaultExpirationMs = parseIntegerEnv('EXPIRATION_DEFAULT_MS', String(24 * 60 * 60 * 1000)); const perEventTypeExpirationJson = trimEnv('EXPIRATION_PER_EVENT_TYPE'); let perEventTypeExpiration: Record | undefined; - + if (perEventTypeExpirationJson) { try { perEventTypeExpiration = JSON.parse(perEventTypeExpirationJson); - if (typeof perEventTypeExpiration !== 'object' || perEventTypeExpiration === null || Array.isArray(perEventTypeExpiration)) { + if ( + typeof perEventTypeExpiration !== 'object' || + perEventTypeExpiration === null || + Array.isArray(perEventTypeExpiration) + ) { throw new ConfigError('EXPIRATION_PER_EVENT_TYPE must be a valid JSON object'); } } catch (e) { if (e instanceof ConfigError) { throw e; } - throw new ConfigError(`EXPIRATION_PER_EVENT_TYPE must be valid JSON. Received: ${perEventTypeExpirationJson}`); + throw new ConfigError( + `EXPIRATION_PER_EVENT_TYPE must be valid JSON. Received: ${perEventTypeExpirationJson}`, + ); } } - + return { defaultExpirationMs, perEventTypeExpiration, @@ -430,7 +434,7 @@ export function loadConfig(): Config { const rawApiKeys = parseJsonEnv('API_KEYS', '[]'); const clientOverrides = parseJsonEnv>( 'RATE_LIMIT_CLIENT_OVERRIDES', - '{}' + '{}', ); const explicitRpcUrl = trimEnv('STELLAR_RPC_URL'); @@ -522,9 +526,7 @@ function loadLoggingConfig(): LoggingConfig { level: trimEnv('LOG_LEVEL') || 'info', // Preserves the previous implicit behaviour when LOG_FORMAT is unset: // JSON in production, human-readable elsewhere. - format: - trimEnv('LOG_FORMAT') || - (process.env.NODE_ENV === 'production' ? 'json' : 'pretty'), + format: trimEnv('LOG_FORMAT') || (process.env.NODE_ENV === 'production' ? 'json' : 'pretty'), }; } @@ -559,9 +561,7 @@ export function validateConfig(config: Config): void { ); } } catch { - errors.push( - `STELLAR_RPC_URL is not a valid URL (received: "${config.stellarRpcUrl}").`, - ); + errors.push(`STELLAR_RPC_URL is not a valid URL (received: "${config.stellarRpcUrl}").`); } } @@ -645,22 +645,16 @@ export function validateConfig(config: Config): void { } if (config.maxReconnectAttempts < 1) { - errors.push( - `MAX_RECONNECT_ATTEMPTS must be >= 1 (received: ${config.maxReconnectAttempts}).`, - ); + errors.push(`MAX_RECONNECT_ATTEMPTS must be >= 1 (received: ${config.maxReconnectAttempts}).`); } if (config.reconnectDelayMs < 0) { - errors.push( - `RECONNECT_DELAY_MS must be >= 0 (received: ${config.reconnectDelayMs}).`, - ); + errors.push(`RECONNECT_DELAY_MS must be >= 0 (received: ${config.reconnectDelayMs}).`); } // ── API server ───────────────────────────────────────────────────────────── if (config.eventsApiPort < 1 || config.eventsApiPort > 65535) { - errors.push( - `EVENTS_API_PORT must be between 1 and 65535 (received: ${config.eventsApiPort}).`, - ); + errors.push(`EVENTS_API_PORT must be between 1 and 65535 (received: ${config.eventsApiPort}).`); } // Validate CORS configuration during startup (#689) @@ -689,7 +683,7 @@ export function validateConfig(config: Config): void { 'Add contract configurations or the service will not process any events.', ); } - + config.contractAddresses.forEach((contract, index) => { if (!contract.address || typeof contract.address !== 'string') { errors.push(`CONTRACT_ADDRESSES[${index}].address must be a non-empty string.`); @@ -709,7 +703,7 @@ export function validateConfig(config: Config): void { ); } } - + if (!Array.isArray(contract.events) || contract.events.length === 0) { errors.push( `CONTRACT_ADDRESSES[${index}].events must be a non-empty array of event names.`, @@ -865,9 +859,7 @@ export function validateConfig(config: Config): void { // ── Analytics ───────────────────────────────────────────────────────────── if (config.analytics) { if (config.analytics.maxRecords < 1) { - errors.push( - `ANALYTICS_MAX_RECORDS must be >= 1 (received: ${config.analytics.maxRecords}).`, - ); + errors.push(`ANALYTICS_MAX_RECORDS must be >= 1 (received: ${config.analytics.maxRecords}).`); } if (config.analytics.bucketSizeMs < 60_000) { errors.push( @@ -1003,16 +995,16 @@ export function validateConfig(config: Config): void { required: false, }, // Webhook signing secrets - ...((config.webhookSecrets ?? []).map((ws, i) => ({ + ...(config.webhookSecrets ?? []).map((ws, i) => ({ fieldName: `WEBHOOK_SECRETS[${i}].secret`, value: ws.secret, required: true, - }))), + })), // API keys - ...((config.apiKeys ?? []).map((ak, i) => ({ + ...(config.apiKeys ?? []).map((ak, i) => ({ fieldName: `API_KEYS[${i}].key`, value: ak.key, required: true, - }))), + })), ]); } diff --git a/listener/src/scripts/load-test.ts b/listener/src/scripts/load-test.ts new file mode 100644 index 00000000..be4fba34 --- /dev/null +++ b/listener/src/scripts/load-test.ts @@ -0,0 +1,180 @@ +#!/usr/bin/env ts-node + +/** + * Repeatable load-testing workflow (issue #860) — external target + * + * Runs the scenarios documented in `load-test.config.json` against a running + * NotifyChain listener and writes a JSON report that can be compared across + * changes. + * + * This entry point talks to an already-running server over HTTP, which is how + * you load-test a deployed listener. To run the same scenarios in-process + * (no external server, no network) use `npm run load-test`, which drives the + * scenarios through the Jest workflow instead — the listener's runtime + * dependency mapping (`@stellar/stellar-sdk`, `node-cache`, `uuid`) is only + * wired up under Jest, so that is the supported in-process path. + * + * Usage: + * npm run load-test:external -- --url http://127.0.0.1:3000 + * npm run load-test:external -- --url https://listener.example.com --out reports/load/staging.json + * npm run load-test:external -- --url http://127.0.0.1:3000 --baseline reports/load/baseline.json + * npm run load-test:external -- --url http://127.0.0.1:3000 --baseline reports/load/baseline.json --fail-on-regression + * + * Exit codes: + * 0 thresholds (and, when gating, baseline comparison) passed + * 1 a threshold was breached, or a regression was detected while gating + * 2 the workflow could not run (bad config, bad arguments) + */ + +import path from 'path'; + +import { + DEFAULT_CONFIG_PATH, + loadLoadTestConfig, + type LoadTestFileConfig, +} from '../utils/load-test-config'; +import { makeHttpProbe } from '../utils/load-test-http-probe'; +import { DEFAULT_REPORT_PATH, readLoadReport, writeLoadReport } from '../utils/load-test-reports'; +import { + compareLoadReports, + formatComparison, + formatLoadReport, + runLoadTest, + type LoadTestConfig, +} from '../utils/load-test-runner'; + +interface CliOptions { + configPath: string; + url?: string; + outPath: string; + baselinePath?: string; + failOnRegression: boolean; + concurrency?: number; + durationMs?: number; + quiet: boolean; +} + +function parseArgs(argv: string[]): CliOptions { + const options: CliOptions = { + configPath: DEFAULT_CONFIG_PATH, + outPath: DEFAULT_REPORT_PATH, + failOnRegression: false, + quiet: false, + }; + + for (let i = 0; i < argv.length; i++) { + const arg = argv[i]; + const next = (): string => { + const value = argv[++i]; + if (value === undefined) throw new Error(`Missing value for ${arg}`); + return value; + }; + + switch (arg) { + case '--config': + options.configPath = path.resolve(next()); + break; + case '--url': + options.url = next(); + break; + case '--out': + options.outPath = path.resolve(next()); + break; + case '--baseline': + options.baselinePath = path.resolve(next()); + break; + case '--concurrency': + options.concurrency = Number(next()); + break; + case '--duration': + options.durationMs = Number(next()); + break; + case '--fail-on-regression': + options.failOnRegression = true; + break; + case '--quiet': + options.quiet = true; + break; + case '--help': + case '-h': + printHelp(); + process.exit(0); + break; + default: + throw new Error(`Unknown argument: ${arg}`); + } + } + + return options; +} + +function printHelp(): void { + process.stdout.write( + [ + 'NotifyChain load-testing workflow (external target)', + '', + 'Options:', + ' --config Scenario/threshold config (default: load-test.config.json)', + ' --url Base URL of the listener to test (default: target.baseUrl in the config)', + ' --out Where to write the JSON report (default: reports/load/latest.json)', + ' --baseline Compare against a previously saved report', + ' --fail-on-regression Exit non-zero when the baseline comparison regresses', + ' --concurrency Override concurrency for every scenario', + ' --duration Override durationMs for every scenario', + ' --quiet Only print the report', + '', + ].join('\n'), + ); +} + +function applyOverrides(config: LoadTestFileConfig, options: CliOptions): LoadTestConfig { + const scenarios = config.scenarios.map((scenario) => ({ + ...scenario, + concurrency: options.concurrency ?? scenario.concurrency, + durationMs: options.durationMs ?? scenario.durationMs, + })); + return { scenarios, thresholds: config.thresholds }; +} + +async function main(): Promise { + const options = parseArgs(process.argv.slice(2)); + const config = loadLoadTestConfig(options.configPath); + const runConfig = applyOverrides(config, options); + + const baseUrl = options.url ?? config.target?.baseUrl; + if (!baseUrl) { + throw new Error('No target: pass --url or set target.baseUrl in the config'); + } + + if (!options.quiet) { + process.stdout.write( + `Load testing ${baseUrl} across ${runConfig.scenarios.length} scenario(s)...\n\n`, + ); + } + + const report = await runLoadTest(runConfig, makeHttpProbe(baseUrl)); + writeLoadReport(options.outPath, report); + process.stdout.write(`${formatLoadReport(report)}\n`); + process.stdout.write(`\nReport written to ${options.outPath}\n`); + + let regressionFailed = false; + if (options.baselinePath) { + const baseline = readLoadReport(options.baselinePath); + const comparison = compareLoadReports(baseline, report); + process.stdout.write(`\n${formatComparison(comparison)}\n`); + if (!comparison.passed && options.failOnRegression) { + regressionFailed = true; + process.stdout.write('\nBaseline regressions detected while --fail-on-regression is set.\n'); + } + } + + if (!report.passed || regressionFailed) { + process.exitCode = 1; + } +} + +main().catch((error: unknown) => { + const message = error instanceof Error ? error.message : String(error); + process.stderr.write(`Load-test workflow failed: ${message}\n`); + process.exit(2); +}); diff --git a/listener/src/services/event-subscriber.ts b/listener/src/services/event-subscriber.ts index 7cd3b778..c45d5e90 100644 --- a/listener/src/services/event-subscriber.ts +++ b/listener/src/services/event-subscriber.ts @@ -144,14 +144,14 @@ export class EventSubscriber { } const events = response.events || []; - + // Detect potential reorg if events exist and we have previous state if (this.deduplicationService && events.length > 0) { const firstEventLedger = events[0]?.ledger; if (firstEventLedger) { const reorgDetected = await this.deduplicationService.detectReorg( contractConfig.address, - firstEventLedger + firstEventLedger, ); if (reorgDetected) { logger.warn('Potential blockchain reorg detected', { @@ -218,14 +218,14 @@ export class EventSubscriber { if (response.cursor) { this.lastCursors.set(contractConfig.address, response.cursor); - + // Update cursor in deduplication service if available if (this.deduplicationService) { const lastEventLedger = events.length > 0 ? events[events.length - 1].ledger : 0; await this.deduplicationService.updatePollingCursor( contractConfig.address, response.cursor, - lastEventLedger || 0 + lastEventLedger || 0, ); } } @@ -240,9 +240,7 @@ export class EventSubscriber { } if (totalContracts > 0 && failureCount === totalContracts) { - throw new Error( - `Failed to fetch events for all ${totalContracts} configured contract(s)` - ); + throw new Error(`Failed to fetch events for all ${totalContracts} configured contract(s)`); } } @@ -250,7 +248,7 @@ export class EventSubscriber { event: StellarSDK.rpc.Api.EventResponse, contractConfig: ContractConfig, requestId: string = '', - correlationId: string = requestId + correlationId: string = requestId, ): boolean { // Check if event has expired const eventName = getEventName(event.topic); @@ -372,7 +370,7 @@ export class EventSubscriber { } private async getContractEvents( - contractConfig: ContractConfig + contractConfig: ContractConfig, ): Promise { // Apply rate limiting before making RPC request if (this.rpcRateLimiter) { @@ -465,7 +463,7 @@ export class EventSubscriber { event: StellarSDK.rpc.Api.EventResponse, contractConfig: ContractConfig, requestId: string = '', - correlationId: string = '' + correlationId: string = '', ): Promise { correlationId = correlationId || requestId || generateCorrelationId(); const eventStart = Date.now(); @@ -542,7 +540,7 @@ export class EventSubscriber { const success = await this.discordService.sendEventNotification( event, contractConfig, - requestId + requestId, ); notificationSent = success; @@ -558,8 +556,8 @@ export class EventSubscriber { } catch (error) { processingError = error instanceof Error ? error.message : String(error); logger.error('Error sending Discord notification', { - requestId: correlationId, - correlationId, + requestId: correlationId, + correlationId, eventId: event.id, error: processingError, }); @@ -577,7 +575,7 @@ export class EventSubscriber { event.type, notificationSent, processingError ? 'ERROR' : 'PROCESSED', - processingError + processingError, ); } @@ -631,4 +629,4 @@ export class EventSubscriber { getCircuitBreakerMetrics() { return this.circuitBreaker?.getMetrics() || null; } -} \ No newline at end of file +} diff --git a/listener/src/utils/load-test-config.ts b/listener/src/utils/load-test-config.ts new file mode 100644 index 00000000..19d94fd1 --- /dev/null +++ b/listener/src/utils/load-test-config.ts @@ -0,0 +1,47 @@ +/** + * Load-test configuration loading (issue #860) + * + * Shared by the external CLI (`src/scripts/load-test.ts`) and the in-process + * Jest workflow (`src/__tests__/load-test.workflow.test.ts`) so both interpret + * `load-test.config.json` identically. + */ + +import fs from 'fs'; +import path from 'path'; + +import type { RateLimitConfig } from '../types'; +import type { LoadTestConfig } from './load-test-runner'; + +export interface LoadTestTarget { + /** Default base URL used by the external CLI when `--url` is omitted. */ + baseUrl?: string; + /** + * Rate-limit config applied to the in-process server by the Jest workflow. + * Rate limiting is disabled by default so the workflow measures raw API + * capacity; enable it to load-test the throttling path itself. + */ + rateLimit?: RateLimitConfig; +} + +export interface LoadTestFileConfig extends LoadTestConfig { + $schema?: string; + target?: LoadTestTarget; +} + +export const DEFAULT_CONFIG_PATH = path.resolve(__dirname, '..', '..', 'load-test.config.json'); + +/** Reads and minimally validates the config file. Throws on unusable input. */ +export function loadLoadTestConfig(configPath: string = DEFAULT_CONFIG_PATH): LoadTestFileConfig { + if (!fs.existsSync(configPath)) { + throw new Error(`Config not found: ${configPath}`); + } + + const parsed = JSON.parse(fs.readFileSync(configPath, 'utf8')) as LoadTestFileConfig; + if (!Array.isArray(parsed.scenarios) || parsed.scenarios.length === 0) { + throw new Error('Config must define at least one scenario'); + } + if (!parsed.thresholds) { + throw new Error('Config must define thresholds'); + } + return parsed; +} diff --git a/listener/src/utils/load-test-http-probe.ts b/listener/src/utils/load-test-http-probe.ts new file mode 100644 index 00000000..f5442a5e --- /dev/null +++ b/listener/src/utils/load-test-http-probe.ts @@ -0,0 +1,64 @@ +/** + * HTTP probe for the load-testing workflow (issue #860) + * + * Kept separate from the measurement core so `load-test-runner.ts` stays + * socket-free (and therefore deterministic under unit test). Latency is sampled + * with `process.hrtime.bigint()` because it is monotonic and unaffected by + * wall-clock adjustments. + */ + +import http from 'http'; +import https from 'https'; + +import type { LoadTestProbe, LoadTestScenario } from './load-test-runner'; + +/** + * Builds a probe that issues real HTTP requests against `baseUrl`. + * + * Network failures resolve with status `0` rather than rejecting, so a flaky + * connection shows up in the error rate instead of aborting the whole run. + */ +export function makeHttpProbe(baseUrl: string): LoadTestProbe { + const target = new URL(baseUrl); + const secure = target.protocol === 'https:'; + const port = target.port ? Number(target.port) : secure ? 443 : 80; + const agent = secure + ? new https.Agent({ keepAlive: true, maxSockets: 256 }) + : new http.Agent({ keepAlive: true, maxSockets: 256 }); + + return (scenario: LoadTestScenario) => + new Promise((resolve) => { + const startedAt = process.hrtime.bigint(); + const request = (secure ? https : http).request( + { + host: target.hostname, + port, + path: scenario.path, + method: scenario.method, + agent, + headers: { + 'content-type': 'application/json', + ...(scenario.headers ?? {}), + }, + }, + (response) => { + response.on('data', () => { + /* drain the body so the socket can be reused */ + }); + response.on('end', () => { + resolve({ + status: response.statusCode ?? 0, + durationMs: Number(process.hrtime.bigint() - startedAt) / 1e6, + }); + }); + }, + ); + + request.on('error', () => { + resolve({ status: 0, durationMs: Number(process.hrtime.bigint() - startedAt) / 1e6 }); + }); + + if (scenario.body) request.write(scenario.body); + request.end(); + }); +} diff --git a/listener/src/utils/load-test-reports.ts b/listener/src/utils/load-test-reports.ts new file mode 100644 index 00000000..7b2e1cbe --- /dev/null +++ b/listener/src/utils/load-test-reports.ts @@ -0,0 +1,30 @@ +/** + * Report persistence for the load-testing workflow (issue #860) + * + * Reports are schema-versioned JSON (`LOAD_REPORT_SCHEMA_VERSION`), which is + * what makes "results comparable across changes" possible: a report can be + * committed as a baseline and diffed by `compareLoadReports`. + */ + +import fs from 'fs'; +import path from 'path'; + +import type { LoadReport } from './load-test-runner'; + +export const DEFAULT_REPORT_DIR = path.resolve(__dirname, '..', '..', 'reports', 'load'); +export const DEFAULT_REPORT_PATH = path.join(DEFAULT_REPORT_DIR, 'latest.json'); +export const DEFAULT_BASELINE_PATH = path.join(DEFAULT_REPORT_DIR, 'baseline.json'); + +/** Writes a report as pretty-printed JSON, creating the directory if needed. */ +export function writeLoadReport(outPath: string, report: LoadReport): void { + fs.mkdirSync(path.dirname(outPath), { recursive: true }); + fs.writeFileSync(outPath, `${JSON.stringify(report, null, 2)}\n`); +} + +/** Reads a previously written report. Throws if it is missing. */ +export function readLoadReport(reportPath: string): LoadReport { + if (!fs.existsSync(reportPath)) { + throw new Error(`Report not found: ${reportPath}`); + } + return JSON.parse(fs.readFileSync(reportPath, 'utf8')) as LoadReport; +} diff --git a/listener/src/utils/load-test-runner.ts b/listener/src/utils/load-test-runner.ts new file mode 100644 index 00000000..bd04ab47 --- /dev/null +++ b/listener/src/utils/load-test-runner.ts @@ -0,0 +1,506 @@ +/** + * Load-testing runner (issue #860) + * + * A small, dependency-free measurement core for the NotifyChain API. It is + * deliberately separated from the CLI (`src/scripts/load-test.ts`) so the + * metric maths, threshold gating and run-to-run comparison logic can be unit + * tested without opening a socket. + * + * The workflow the issue asks for falls out of three exported pieces: + * 1. `runLoadTest` – executes documented scenarios and returns a report + * 2. `compareLoadReports` – diffs a report against a saved baseline + * 3. `formatLoadReport` / `formatComparison` – human-readable summaries + * + * Reports are plain JSON with a `schemaVersion`, so results can be committed + * as a baseline and compared across changes. + */ + +import os from 'os'; + +export const LOAD_REPORT_SCHEMA_VERSION = 1; + +export type LoadTestMethod = 'GET' | 'POST' | 'PUT' | 'DELETE'; + +export interface LoadTestScenario { + /** Stable identifier used to match scenarios across reports. */ + name: string; + /** Human-readable note about which critical path the scenario covers. */ + description?: string; + method: LoadTestMethod; + path: string; + headers?: Record; + /** Raw request body sent with the scenario (already serialized). */ + body?: string; + /** Number of concurrent in-flight requests. */ + concurrency: number; + /** Length of the measured phase, in milliseconds. */ + durationMs: number; + /** Optional warm-up phase that is executed but excluded from the report. */ + warmupMs?: number; + /** + * Optional hard cap on measured requests. Combined with `durationMs` via a + * logical AND, this keeps CI runs bounded and makes small runs deterministic. + */ + maxRequests?: number; +} + +export interface LoadTestThresholds { + /** Maximum tolerated fraction of non-2xx responses (0..1). */ + maxErrorRate: number; + /** Maximum tolerated p95 latency, in milliseconds. */ + maxP95Ms: number; + /** Optional minimum aggregate throughput, in requests per second. */ + minThroughputRps?: number; +} + +export interface LoadTestConfig { + scenarios: LoadTestScenario[]; + thresholds: LoadTestThresholds; +} + +export interface LatencyStats { + min: number; + mean: number; + p50: number; + p90: number; + p95: number; + p99: number; + max: number; +} + +export interface ScenarioResult { + name: string; + method: LoadTestMethod; + path: string; + concurrency: number; + durationMs: number; + requests: number; + success: number; + clientErrors: number; + serverErrors: number; + /** Non-2xx responses divided by total responses. */ + errorRate: number; + throughputRps: number; + latencyMs: LatencyStats; +} + +export interface LoadReportEnvironment { + node: string; + platform: string; + release: string; + cpus: number; + totalMemoryMb: number; +} + +export interface LoadReport { + schemaVersion: number; + generatedAt: string; + environment: LoadReportEnvironment; + scenarios: ScenarioResult[]; + totals: { + requests: number; + errorRate: number; + throughputRps: number; + latencyMs: LatencyStats; + }; + thresholds: LoadTestThresholds; + passed: boolean; + failures: string[]; +} + +/** One probe invocation: returns the HTTP status and how long it took. */ +export interface ProbeResult { + status: number; + durationMs: number; +} + +export type LoadTestProbe = (scenario: LoadTestScenario) => Promise; + +const EMPTY_STATS: LatencyStats = { min: 0, mean: 0, p50: 0, p90: 0, p95: 0, p99: 0, max: 0 }; + +/** + * Nearest-rank percentile (the same definition used by most HTTP load tools). + * `sorted` must be ascending. Returns 0 for an empty input. + */ +export function percentile(sorted: number[], p: number): number { + if (sorted.length === 0) return 0; + if (p <= 0) return sorted[0]; + if (p >= 100) return sorted[sorted.length - 1]; + const rank = Math.ceil((p / 100) * sorted.length); + const index = Math.min(sorted.length - 1, Math.max(0, rank - 1)); + return sorted[index]; +} + +/** Summary statistics for a set of latency samples. */ +export function summarizeLatencies(latencies: number[]): LatencyStats { + if (latencies.length === 0) return { ...EMPTY_STATS }; + const sorted = [...latencies].sort((a, b) => a - b); + const sum = sorted.reduce((total, value) => total + value, 0); + return { + min: round(sorted[0]), + mean: round(sum / sorted.length), + p50: round(percentile(sorted, 50)), + p90: round(percentile(sorted, 90)), + p95: round(percentile(sorted, 95)), + p99: round(percentile(sorted, 99)), + max: round(sorted[sorted.length - 1]), + }; +} + +function round(value: number): number { + return Math.round(value * 100) / 100; +} + +/** + * Runs one phase (warm-up or measured) at the scenario's concurrency. + * + * Each worker loops "check deadline/quota → issue request → await". Because the + * check and the counter increment happen in the same synchronous block, the + * `maxRequests` cap is never exceeded even with high concurrency. + */ +async function runPhase( + scenario: LoadTestScenario, + probe: LoadTestProbe, + durationMs: number, + maxRequests?: number, +): Promise<{ + latencies: number[]; + success: number; + clientErrors: number; + serverErrors: number; + elapsedMs: number; +}> { + const startedAt = Date.now(); + const deadline = startedAt + durationMs; + const latencies: number[] = []; + let success = 0; + let clientErrors = 0; + let serverErrors = 0; + let issued = 0; + + const workers = Array.from({ length: Math.max(1, scenario.concurrency) }, async () => { + while (Date.now() < deadline && (maxRequests === undefined || issued < maxRequests)) { + issued += 1; + const result = await probe(scenario); + latencies.push(result.durationMs); + if (result.status >= 200 && result.status < 300) success += 1; + else if (result.status >= 500) serverErrors += 1; + else clientErrors += 1; + } + }); + + await Promise.all(workers); + return { + latencies, + success, + clientErrors, + serverErrors, + elapsedMs: Math.max(1, Date.now() - startedAt), + }; +} + +/** Executes a single scenario (warm-up + measured phase) and summarises it. */ +export async function runScenario( + scenario: LoadTestScenario, + probe: LoadTestProbe, +): Promise { + if (scenario.warmupMs && scenario.warmupMs > 0) { + await runPhase(scenario, probe, scenario.warmupMs, scenario.maxRequests); + } + + const { latencies, success, clientErrors, serverErrors, elapsedMs } = await runPhase( + scenario, + probe, + scenario.durationMs, + scenario.maxRequests, + ); + + const requests = latencies.length; + return { + name: scenario.name, + method: scenario.method, + path: scenario.path, + concurrency: scenario.concurrency, + durationMs: elapsedMs, + requests, + success, + clientErrors, + serverErrors, + errorRate: requests === 0 ? 0 : round((requests - success) / requests), + throughputRps: round(requests / (elapsedMs / 1000)), + latencyMs: summarizeLatencies(latencies), + }; +} + +function describeEnvironment(): LoadReportEnvironment { + return { + node: process.version, + platform: os.platform(), + release: os.release(), + cpus: os.cpus().length, + totalMemoryMb: Math.round(os.totalmem() / (1024 * 1024)), + }; +} + +/** + * Runs every scenario sequentially and evaluates the configured thresholds. + * Scenarios run one after another so they don't contend for the same event loop + * / server capacity (which would make results non-comparable across runs). + */ +export async function runLoadTest( + config: LoadTestConfig, + probe: LoadTestProbe, +): Promise { + const scenarios: ScenarioResult[] = []; + for (const scenario of config.scenarios) { + scenarios.push(await runScenario(scenario, probe)); + } + + let requests = 0; + let success = 0; + let elapsedMs = 0; + for (const scenario of scenarios) { + requests += scenario.requests; + success += scenario.success; + elapsedMs += scenario.durationMs; + } + + const totals = { + requests, + errorRate: requests === 0 ? 0 : round((requests - success) / requests), + throughputRps: elapsedMs === 0 ? 0 : round(requests / (elapsedMs / 1000)), + latencyMs: summarizeLatencies(collectScenarioLatencies(scenarios)), + }; + const failures = evaluateThresholds(scenarios, totals, config.thresholds); + + return { + schemaVersion: LOAD_REPORT_SCHEMA_VERSION, + generatedAt: new Date().toISOString(), + environment: describeEnvironment(), + scenarios, + totals, + thresholds: config.thresholds, + passed: failures.length === 0, + failures, + }; +} + +/** + * We only keep per-scenario aggregate latency stats in the report, so the + * aggregate "latencyMs" is derived from those percentiles (a report-level + * roll-up, not a re-sampling of raw requests). + */ +function collectScenarioLatencies(scenarios: ScenarioResult[]): number[] { + const samples: number[] = []; + for (const scenario of scenarios) { + // Weight each scenario's percentile points by its request count so the + // roll-up tracks the overall traffic mix. + const weight = Math.max(1, Math.min(50, Math.round(scenario.requests / 10) || 1)); + for (let i = 0; i < weight; i++) { + samples.push( + scenario.latencyMs.p50, + scenario.latencyMs.p90, + scenario.latencyMs.p95, + scenario.latencyMs.p99, + ); + } + } + return samples; +} + +function evaluateThresholds( + scenarios: ScenarioResult[], + totals: LoadReport['totals'], + thresholds: LoadTestThresholds, +): string[] { + const failures: string[] = []; + + if (totals.errorRate > thresholds.maxErrorRate) { + failures.push( + `Aggregate error rate ${(totals.errorRate * 100).toFixed(2)}% exceeds maxErrorRate ${(thresholds.maxErrorRate * 100).toFixed(2)}%`, + ); + } + + for (const scenario of scenarios) { + if (scenario.latencyMs.p95 > thresholds.maxP95Ms) { + failures.push( + `Scenario "${scenario.name}" p95 ${scenario.latencyMs.p95}ms exceeds maxP95Ms ${thresholds.maxP95Ms}ms`, + ); + } + } + + if ( + thresholds.minThroughputRps !== undefined && + totals.throughputRps < thresholds.minThroughputRps + ) { + failures.push( + `Aggregate throughput ${totals.throughputRps} rps is below minThroughputRps ${thresholds.minThroughputRps}`, + ); + } + + return failures; +} + +export interface ScenarioComparison { + name: string; + status: 'added' | 'removed' | 'compared'; + throughputRps?: Delta; + p95Ms?: Delta; + errorRate?: Delta; +} + +export interface Delta { + baseline: number; + current: number; + /** Absolute change (current - baseline). */ + delta: number; + /** Relative change as a fraction (0.25 = +25%). */ + deltaPct: number; +} + +export interface LoadReportComparison { + scenarios: ScenarioComparison[]; + totals: { + throughputRps: Delta; + p95Ms: Delta; + errorRate: Delta; + }; + /** True when no scenario's p95 regressed by more than `regressionTolerancePct`. */ + passed: boolean; + regressions: string[]; +} + +function makeDelta(baseline: number, current: number): Delta { + const delta = round(current - baseline); + return { + baseline, + current, + delta, + deltaPct: baseline === 0 ? (current === 0 ? 0 : 1) : round((current - baseline) / baseline), + }; +} + +/** + * Diffs two reports. `regressionTolerancePct` is the acceptable relative p95 + * increase (0.1 = 10%); anything above it is reported as a regression. + */ +export function compareLoadReports( + baseline: LoadReport, + current: LoadReport, + regressionTolerancePct = 0.1, +): LoadReportComparison { + const baselineByName = new Map(baseline.scenarios.map((scenario) => [scenario.name, scenario])); + const currentByName = new Map(current.scenarios.map((scenario) => [scenario.name, scenario])); + const names = Array.from(new Set([...baselineByName.keys(), ...currentByName.keys()])); + + const scenarios: ScenarioComparison[] = []; + const regressions: string[] = []; + + for (const name of names) { + const before = baselineByName.get(name); + const after = currentByName.get(name); + if (before && !after) { + scenarios.push({ name, status: 'removed' }); + continue; + } + if (!before && after) { + scenarios.push({ name, status: 'added' }); + continue; + } + if (!before || !after) continue; + + const p95Ms = makeDelta(before.latencyMs.p95, after.latencyMs.p95); + scenarios.push({ + name, + status: 'compared', + throughputRps: makeDelta(before.throughputRps, after.throughputRps), + p95Ms, + errorRate: makeDelta(before.errorRate, after.errorRate), + }); + + if (p95Ms.deltaPct > regressionTolerancePct) { + regressions.push( + `Scenario "${name}" p95 regressed ${(p95Ms.deltaPct * 100).toFixed(1)}% (${before.latencyMs.p95}ms → ${after.latencyMs.p95}ms)`, + ); + } + } + + return { + scenarios, + totals: { + throughputRps: makeDelta(baseline.totals.throughputRps, current.totals.throughputRps), + p95Ms: makeDelta(baseline.totals.latencyMs.p95, current.totals.latencyMs.p95), + errorRate: makeDelta(baseline.totals.errorRate, current.totals.errorRate), + }, + passed: regressions.length === 0, + regressions, + }; +} + +/** Human-readable one-screen summary of a report. */ +export function formatLoadReport(report: LoadReport): string { + const lines: string[] = []; + lines.push(`Load report (${report.generatedAt}) — schema v${report.schemaVersion}`); + lines.push( + `Environment: node ${report.environment.node} · ${report.environment.platform} ${report.environment.release} · ${report.environment.cpus} cpus`, + ); + lines.push(''); + + const header = ['scenario', 'reqs', 'err%', 'rps', 'p50', 'p95', 'p99', 'max']; + const rows = report.scenarios.map((scenario) => [ + scenario.name, + String(scenario.requests), + `${(scenario.errorRate * 100).toFixed(2)}%`, + scenario.throughputRps.toFixed(1), + `${scenario.latencyMs.p50}ms`, + `${scenario.latencyMs.p95}ms`, + `${scenario.latencyMs.p99}ms`, + `${scenario.latencyMs.max}ms`, + ]); + // Size each column to its widest cell so long scenario names don't collide + // with the next column. + const widths = header.map( + (title, index) => Math.max(title.length, ...rows.map((row) => row[index].length)) + 2, + ); + const renderRow = (cells: string[]): string => + cells + .map((cell, index) => cell.padEnd(widths[index])) + .join('') + .trimEnd(); + + lines.push(renderRow(header)); + for (const row of rows) lines.push(renderRow(row)); + lines.push(''); + lines.push( + `Totals: ${report.totals.requests} requests · ${report.totals.throughputRps.toFixed(1)} rps · error rate ${(report.totals.errorRate * 100).toFixed(2)}% · p95 ${report.totals.latencyMs.p95}ms`, + ); + lines.push(report.passed ? 'Thresholds: PASS' : `Thresholds: FAIL (${report.failures.length})`); + for (const failure of report.failures) lines.push(` ✗ ${failure}`); + return lines.join('\n'); +} + +/** Human-readable summary of a baseline comparison. */ +export function formatComparison(comparison: LoadReportComparison): string { + const lines: string[] = []; + lines.push('Comparison vs baseline:'); + for (const scenario of comparison.scenarios) { + if (scenario.status !== 'compared') { + lines.push(` ${scenario.name}: ${scenario.status}`); + continue; + } + const p95 = scenario.p95Ms!; + const rps = scenario.throughputRps!; + const arrow = p95.deltaPct > 0 ? '↑' : '↓'; + lines.push( + ` ${scenario.name}: p95 ${p95.baseline}ms → ${p95.current}ms (${arrow}${(Math.abs(p95.deltaPct) * 100).toFixed(1)}%), rps ${rps.baseline.toFixed(1)} → ${rps.current.toFixed(1)}`, + ); + } + lines.push( + ` totals: p95 ${comparison.totals.p95Ms.baseline}ms → ${comparison.totals.p95Ms.current}ms, rps ${comparison.totals.throughputRps.baseline.toFixed(1)} → ${comparison.totals.throughputRps.current.toFixed(1)}`, + ); + if (comparison.regressions.length > 0) { + lines.push('Regressions:'); + for (const regression of comparison.regressions) lines.push(` ✗ ${regression}`); + } + return lines.join('\n'); +}