From c31b92059a002a851feaff27977550b5b781c8d1 Mon Sep 17 00:00:00 2001 From: woahwhattheheck Date: Thu, 24 Sep 2026 18:07:21 -0400 Subject: [PATCH 1/6] feat(health): add dependency-aware readiness and recovery diagnostics Separate process liveness from traffic readiness. Ready probes now check store, payments (Stellar), and FX under a per-dependency timeout, return redacted reason codes on failure, and re-evaluate every request so recovery does not require a restart. Liveness stays dependency-free so outages do not flap orchestrator restarts. Closes #134 --- .env.example | 5 + CHANGELOG.md | 6 + README.md | 9 +- src/config/index.js | 7 + src/controllers/healthController.js | 30 ++- src/services/dependencyHealthService.js | 319 ++++++++++++++++++++++++ src/services/rateService.js | 17 ++ src/services/stellarService.js | 17 ++ test/dependencyHealth.test.js | 241 ++++++++++++++++++ test/smoke.test.js | 7 +- 10 files changed, 644 insertions(+), 14 deletions(-) create mode 100644 src/services/dependencyHealthService.js create mode 100644 test/dependencyHealth.test.js diff --git a/.env.example b/.env.example index dcba576..4b8c2fa 100644 --- a/.env.example +++ b/.env.example @@ -56,3 +56,8 @@ PAGINATION_MAX_SCAN=10000 # Valid scopes: transfers:read, transfers:write, users:read, users:write, audit:read # Example (single line, properly escaped for shell): # API_TOKENS={"my-secret-token":["transfers:read","transfers:write","users:read","users:write","audit:read"]} + +# Health / readiness probes +# Per-dependency time budget for GET /api/health/ready. Kept well below +# REQUEST_TIMEOUT_MS so a hung dependency cannot stall the probe. +HEALTH_CHECK_TIMEOUT_MS=1000 diff --git a/CHANGELOG.md b/CHANGELOG.md index a8e4e0a..2ea9e8c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,6 +18,12 @@ When preparing a new release: ### Added +- Dependency-aware readiness diagnostics: `GET /api/health/ready` now probes + store, payments (Stellar), and FX under a per-check timeout + (`HEALTH_CHECK_TIMEOUT_MS`), returns redacted reason codes on failure, and + recovers without a process restart. Liveness (`GET /api/health/live`) stays + process-only so orchestrators do not flap restarts during dependency outages. + - 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..46486a2 100644 --- a/README.md +++ b/README.md @@ -227,9 +227,12 @@ The API implements Cache-Control response headers for security and efficiency: ### Health -- `GET /api/health` — service health snapshot (version, uptime, env). -- `GET /api/health/live` — liveness probe. -- `GET /api/health/ready` — readiness probe with dependency checks. +- `GET /api/health` — service health snapshot (version, uptime, env). Not a dependency gate. +- `GET /api/health/live` — liveness probe. Process-only; stays responsive during dependency outages. +- `GET /api/health/ready` — readiness probe with bounded checks for store, payments (Stellar), and FX. + Returns `200` when ready and `503` with redacted reason codes when a dependency is down or times out. + Probes re-run on every request so recovery does not require a restart. + Tune the per-check budget with `HEALTH_CHECK_TIMEOUT_MS` (default `1000`). - `GET /api/version` — service name and version. ### Rates & quotes diff --git a/src/config/index.js b/src/config/index.js index 41f9a71..ff29378 100644 --- a/src/config/index.js +++ b/src/config/index.js @@ -69,6 +69,13 @@ const config = { maxScan: parseInt(process.env.PAGINATION_MAX_SCAN, 10) || 10000, }, + + health: { + // Per-dependency budget for readiness probes. Must stay well below the + // request timeout so a hung dependency cannot stall the readiness route. + checkTimeoutMs: parseInt(process.env.HEALTH_CHECK_TIMEOUT_MS, 10) || 1000, + }, + apiTokens: (() => { try { if (process.env.API_TOKENS) { diff --git a/src/controllers/healthController.js b/src/controllers/healthController.js index 0e41aa5..4c91e07 100644 --- a/src/controllers/healthController.js +++ b/src/controllers/healthController.js @@ -2,14 +2,21 @@ const config = require('../config'); const { name, version } = require('../../package.json'); +const dependencyHealth = require('../services/dependencyHealthService'); /** * Health controller. + * + * Liveness answers "is the process up?" and must stay cheap and dependency- + * free so orchestrators do not restart a process that is merely waiting on + * a degraded dependency. Readiness answers "can this instance serve traffic?" + * by running bounded, redacted dependency probes that recover automatically + * once the dependency is healthy again. */ /** * GET /api/health - * Reports basic liveness information. + * Reports basic process information (not a dependency gate). */ function getHealth(req, res) { res.json({ @@ -36,7 +43,8 @@ function getVersion(req, res) { /** * GET /api/health/live - * Liveness probe: confirms the process is up and responding. + * Liveness probe: confirms the process is up and responding. Intentionally + * ignores store / payment / FX state so outages do not flap restarts. */ function getLiveness(req, res) { res.json({ status: 'alive', timestamp: new Date().toISOString() }); @@ -44,15 +52,19 @@ function getLiveness(req, res) { /** * GET /api/health/ready - * Readiness probe: confirms dependencies needed to serve traffic are up. - * The demo store is always in-memory, so readiness simply mirrors liveness. + * Readiness probe: bounded checks against store, payments, and FX. + * Returns 200 when every dependency is healthy and 503 otherwise. Reason + * codes are redacted; recovery is automatic on the next successful probe. */ -function getReadiness(req, res) { - res.json({ - status: 'ready', - checks: { store: 'ok' }, +async function getReadiness(req, res) { + const result = await dependencyHealth.evaluateReadiness(); + const payload = { + status: result.status, + checks: result.checks, + timeoutMs: result.timeoutMs, timestamp: new Date().toISOString(), - }); + }; + res.status(result.ready ? 200 : 503).json(payload); } module.exports = { diff --git a/src/services/dependencyHealthService.js b/src/services/dependencyHealthService.js new file mode 100644 index 0000000..4efeed0 --- /dev/null +++ b/src/services/dependencyHealthService.js @@ -0,0 +1,319 @@ +'use strict'; + +const config = require('../config'); +const { store } = require('../store'); +const rateService = require('./rateService'); +const stellarService = require('./stellarService'); + +/** + * Dependency-aware readiness diagnostics. + * + * Separates process liveness from traffic readiness by probing the store + * (database stand-in), payment provider (Stellar), and FX rate table with + * a per-check time budget. Failures surface as stable, redacted reason + * codes — never raw messages, stacks, or connection material — and every + * probe is re-evaluated on the next request so recovery does not require + * a process restart. + */ + +/** Stable dependency names exposed on the readiness payload. */ +const DEPENDENCIES = Object.freeze(['store', 'payments', 'fx']); + +/** Reason codes returned to callers. Keep these stable for monitors. */ +const REASON = Object.freeze({ + STORE_UNAVAILABLE: 'STORE_UNAVAILABLE', + STORE_TIMEOUT: 'STORE_TIMEOUT', + PAYMENTS_UNAVAILABLE: 'PAYMENTS_UNAVAILABLE', + PAYMENTS_TIMEOUT: 'PAYMENTS_TIMEOUT', + FX_UNAVAILABLE: 'FX_UNAVAILABLE', + FX_TIMEOUT: 'FX_TIMEOUT', + CHECK_ERROR: 'CHECK_ERROR', +}); + +/** + * Test-only overrides. A Map of dependency name -> override descriptor: + * { mode: 'fail', reason?: string } + * { mode: 'timeout', delayMs?: number } + * { mode: 'throw', message?: string } // message is redacted from responses + * Cleared between tests so forced failures never stick across process life. + * @type {Map} + */ +const forcedStates = new Map(); + +/** + * Default probe implementations. Each returns a Promise that resolves on + * success or rejects with an Error carrying a `reasonCode`. + */ +const defaultProbes = Object.freeze({ + async store() { + if (!store || !(store.users instanceof Map) || !(store.transfers instanceof Map)) { + const err = new Error('store unavailable'); + err.reasonCode = REASON.STORE_UNAVAILABLE; + throw err; + } + // Touch the maps so a corrupted store surfaces as unavailable. + void store.users.size; + void store.transfers.size; + return { ok: true }; + }, + + async payments() { + const result = stellarService.ping(); + if (!result || result.ok !== true) { + const err = new Error('payments unavailable'); + err.reasonCode = REASON.PAYMENTS_UNAVAILABLE; + throw err; + } + return { ok: true }; + }, + + async fx() { + const result = rateService.ping(); + if (!result || result.ok !== true) { + const err = new Error('fx unavailable'); + err.reasonCode = REASON.FX_UNAVAILABLE; + throw err; + } + return { ok: true }; + }, +}); + +/** Active probes — start as defaults; tests may replace individual ones. */ +const probes = { + store: defaultProbes.store, + payments: defaultProbes.payments, + fx: defaultProbes.fx, +}; + +/** + * Resolve the per-check timeout budget in milliseconds. + * @returns {number} + */ +function checkTimeoutMs() { + const raw = config.health && config.health.checkTimeoutMs; + const n = Number(raw); + return Number.isFinite(n) && n > 0 ? n : 1000; +} + +/** + * Race a probe against a hard deadline. The original promise is not + * cancelled (Node has no Abort for plain promises), but readiness never + * waits longer than the budget. + * @param {Promise} promise + * @param {number} ms + * @param {string} timeoutReason + * @returns {Promise} + */ +function withTimeout(promise, ms, timeoutReason) { + let timer; + const timeoutPromise = new Promise((_, reject) => { + timer = setTimeout(() => { + const err = new Error('dependency check timed out'); + err.reasonCode = timeoutReason; + err.code = 'TIMEOUT'; + reject(err); + }, ms); + // Do not keep the process alive solely for a readiness timer. + if (typeof timer.unref === 'function') { + timer.unref(); + } + }); + + return Promise.race([promise, timeoutPromise]).finally(() => { + clearTimeout(timer); + }); +} + +/** + * Map an arbitrary failure to a redacted reason code. Raw messages and + * stacks never leave this function. + * @param {string} name + * @param {unknown} err + * @returns {string} + */ +function redactReason(name, err) { + const code = err && typeof err === 'object' ? err.reasonCode : undefined; + if (typeof code === 'string' && Object.values(REASON).includes(code)) { + return code; + } + + const timeouts = { + store: REASON.STORE_TIMEOUT, + payments: REASON.PAYMENTS_TIMEOUT, + fx: REASON.FX_TIMEOUT, + }; + if (err && typeof err === 'object' && err.code === 'TIMEOUT') { + return timeouts[name] || REASON.CHECK_ERROR; + } + + const unavailable = { + store: REASON.STORE_UNAVAILABLE, + payments: REASON.PAYMENTS_UNAVAILABLE, + fx: REASON.FX_UNAVAILABLE, + }; + return unavailable[name] || REASON.CHECK_ERROR; +} + +/** + * Apply a test override before the real probe, if one is set. + * @param {string} name + * @returns {Promise|null} a substitute promise, or null to run the real probe + */ +function applyForcedState(name) { + const forced = forcedStates.get(name); + if (!forced) return null; + + if (forced.mode === 'fail') { + const err = new Error('forced failure'); + err.reasonCode = + forced.reason || + ({ + store: REASON.STORE_UNAVAILABLE, + payments: REASON.PAYMENTS_UNAVAILABLE, + fx: REASON.FX_UNAVAILABLE, + }[name] || REASON.CHECK_ERROR); + return Promise.reject(err); + } + + if (forced.mode === 'timeout') { + const delay = Number(forced.delayMs) > 0 ? Number(forced.delayMs) : checkTimeoutMs() * 5; + return new Promise((resolve) => { + const t = setTimeout(resolve, delay); + if (typeof t.unref === 'function') t.unref(); + }); + } + + if (forced.mode === 'throw') { + // Deliberately include sensitive-looking content so redaction tests can + // assert it never reaches the HTTP response. + const err = new Error( + forced.message || + 'postgres://user:super-secret@db.internal:5432/remitflow leaked' + ); + return Promise.reject(err); + } + + return null; +} + +/** + * Run a single named dependency check under the shared time budget. + * @param {string} name + * @param {number} [timeoutMs] + * @returns {Promise<{name: string, status: 'ok'|'error', reason?: string, latencyMs: number}>} + */ +async function runCheck(name, timeoutMs = checkTimeoutMs()) { + const started = Date.now(); + const timeoutReason = { + store: REASON.STORE_TIMEOUT, + payments: REASON.PAYMENTS_TIMEOUT, + fx: REASON.FX_TIMEOUT, + }[name] || REASON.CHECK_ERROR; + + try { + const forced = applyForcedState(name); + const probe = typeof probes[name] === 'function' ? probes[name] : null; + if (!probe && !forced) { + const err = new Error('unknown dependency'); + err.reasonCode = REASON.CHECK_ERROR; + throw err; + } + await withTimeout(forced || probe(), timeoutMs, timeoutReason); + return { + name, + status: 'ok', + latencyMs: Date.now() - started, + }; + } catch (err) { + return { + name, + status: 'error', + reason: redactReason(name, err), + latencyMs: Date.now() - started, + }; + } +} + +/** + * Evaluate every dependency. Safe to call on every readiness request — + * nothing is cached as permanently failed. + * @param {object} [options] + * @param {number} [options.timeoutMs] + * @returns {Promise<{ + * ready: boolean, + * status: 'ready'|'not_ready', + * checks: Record, + * timeoutMs: number + * }>} + */ +async function evaluateReadiness(options = {}) { + const timeoutMs = Number(options.timeoutMs) > 0 ? Number(options.timeoutMs) : checkTimeoutMs(); + const results = await Promise.all(DEPENDENCIES.map((name) => runCheck(name, timeoutMs))); + + /** @type {Record} */ + const checks = {}; + let ready = true; + for (const result of results) { + const entry = { status: result.status, latencyMs: result.latencyMs }; + if (result.reason) entry.reason = result.reason; + checks[result.name] = entry; + if (result.status !== 'ok') ready = false; + } + + return { + ready, + status: ready ? 'ready' : 'not_ready', + checks, + timeoutMs, + }; +} + +/** + * Force a dependency into a known bad state for tests. Cleared with + * `clearForcedStates` so readiness can recover without a restart. + * @param {string} name + * @param {{mode: 'fail'|'timeout'|'throw', reason?: string, delayMs?: number, message?: string}} state + */ +function forceDependencyState(name, state) { + if (!DEPENDENCIES.includes(name)) { + throw new Error(`unknown dependency: ${name}`); + } + forcedStates.set(name, state); +} + +/** Remove every test override so subsequent probes hit the real path. */ +function clearForcedStates() { + forcedStates.clear(); +} + +/** + * Replace a probe implementation (tests only). Pass `null` to restore default. + * @param {string} name + * @param {null|(() => Promise)} fn + */ +function setProbeForTests(name, fn) { + if (!DEPENDENCIES.includes(name)) { + throw new Error(`unknown dependency: ${name}`); + } + probes[name] = typeof fn === 'function' ? fn : defaultProbes[name]; +} + +/** Restore default probes and clear forced states (tests only). */ +function resetForTests() { + for (const name of DEPENDENCIES) { + probes[name] = defaultProbes[name]; + } + forcedStates.clear(); +} + +module.exports = { + DEPENDENCIES, + REASON, + evaluateReadiness, + runCheck, + forceDependencyState, + clearForcedStates, + setProbeForTests, + resetForTests, + checkTimeoutMs, +}; diff --git a/src/services/rateService.js b/src/services/rateService.js index 95c7112..29486dc 100644 --- a/src/services/rateService.js +++ b/src/services/rateService.js @@ -76,10 +76,27 @@ function getPair(from, to) { }; } + +/** + * Lightweight FX probe used by readiness checks. + * Confirms the rate table is loaded and returns at least one currency. + * @returns {{ ok: true, currencies: number }} + */ +function ping() { + const rates = listRates(); + if (!Array.isArray(rates) || rates.length === 0) { + const err = new Error('fx rate table empty'); + err.reasonCode = 'FX_UNAVAILABLE'; + throw err; + } + return { ok: true, currencies: rates.length }; +} + module.exports = { listRates, isSupported, getRate, convert, getPair, + ping, }; diff --git a/src/services/stellarService.js b/src/services/stellarService.js index e77cea9..99ab210 100644 --- a/src/services/stellarService.js +++ b/src/services/stellarService.js @@ -11,6 +11,22 @@ const logger = require('../utils/logger'); * rest of the app can pretend a settlement happened. */ +/** + * Lightweight payment-provider probe used by readiness checks. + * Validates that the Stellar network configuration is present so a + * misconfigured process fails readiness instead of reporting healthy. + * @returns {{ ok: true, network: string }} + */ +function ping() { + const network = config.stellar && config.stellar.network; + if (typeof network !== 'string' || network.trim() === '') { + const err = new Error('stellar network not configured'); + err.reasonCode = 'PAYMENTS_UNAVAILABLE'; + throw err; + } + return { ok: true, network }; +} + /** * Pretend to submit a payment to the Stellar network. * @param {object} params @@ -36,6 +52,7 @@ function createClaimableBalanceId() { } module.exports = { + ping, submitPayment, createClaimableBalanceId, }; diff --git a/test/dependencyHealth.test.js b/test/dependencyHealth.test.js new file mode 100644 index 0000000..7e952ec --- /dev/null +++ b/test/dependencyHealth.test.js @@ -0,0 +1,241 @@ +'use strict'; + +const { test, before, after, beforeEach, afterEach } = require('node:test'); +const assert = require('node:assert/strict'); + +process.env.NODE_ENV = 'test'; + +const createApp = require('../src/app'); +const config = require('../src/config'); +const dependencyHealth = require('../src/services/dependencyHealthService'); + +let server; +let baseUrl; +let originalTimeout; + +before(() => { + originalTimeout = config.health.checkTimeoutMs; + // Keep probe budgets tight so timeout tests stay fast. + config.health.checkTimeoutMs = 50; + + const app = createApp(); + return new Promise((resolve) => { + server = app.listen(0, () => { + const { port } = server.address(); + baseUrl = `http://127.0.0.1:${port}`; + resolve(); + }); + }); +}); + +after(() => { + config.health.checkTimeoutMs = originalTimeout; + dependencyHealth.resetForTests(); + if (server) { + server.close(); + } +}); + +beforeEach(() => { + dependencyHealth.resetForTests(); + config.health.checkTimeoutMs = 50; +}); + +afterEach(() => { + dependencyHealth.resetForTests(); +}); + +async function fetchJson(path) { + const res = await fetch(`${baseUrl}${path}`); + const body = await res.json(); + return { status: res.status, body }; +} + +// ─── Happy path ────────────────────────────────────────────────────────────── + +test('readiness is ready when store, payments, and fx are healthy', async () => { + const { status, body } = await fetchJson('/api/health/ready'); + assert.equal(status, 200); + assert.equal(body.status, 'ready'); + for (const name of dependencyHealth.DEPENDENCIES) { + assert.equal(body.checks[name].status, 'ok'); + assert.equal(body.checks[name].reason, undefined); + assert.ok(typeof body.checks[name].latencyMs === 'number'); + } +}); + +test('evaluateReadiness reports ready=true with all ok checks', async () => { + const result = await dependencyHealth.evaluateReadiness({ timeoutMs: 50 }); + assert.equal(result.ready, true); + assert.equal(result.status, 'ready'); + assert.deepEqual( + Object.keys(result.checks).sort(), + [...dependencyHealth.DEPENDENCIES].sort() + ); +}); + +// ─── Dependency failure (original failure mode regression) ─────────────────── + +test('readiness returns 503 when the payment provider is unavailable', async () => { + dependencyHealth.forceDependencyState('payments', { + mode: 'fail', + reason: dependencyHealth.REASON.PAYMENTS_UNAVAILABLE, + }); + + const { status, body } = await fetchJson('/api/health/ready'); + assert.equal(status, 503); + assert.equal(body.status, 'not_ready'); + assert.equal(body.checks.payments.status, 'error'); + assert.equal(body.checks.payments.reason, 'PAYMENTS_UNAVAILABLE'); + assert.equal(body.checks.store.status, 'ok'); + assert.equal(body.checks.fx.status, 'ok'); +}); + +test('readiness returns 503 when the store dependency fails', async () => { + dependencyHealth.forceDependencyState('store', { + mode: 'fail', + reason: dependencyHealth.REASON.STORE_UNAVAILABLE, + }); + + const { status, body } = await fetchJson('/api/health/ready'); + assert.equal(status, 503); + assert.equal(body.status, 'not_ready'); + assert.equal(body.checks.store.reason, 'STORE_UNAVAILABLE'); +}); + +test('readiness returns 503 when the FX dependency fails', async () => { + dependencyHealth.forceDependencyState('fx', { + mode: 'fail', + reason: dependencyHealth.REASON.FX_UNAVAILABLE, + }); + + const { status, body } = await fetchJson('/api/health/ready'); + assert.equal(status, 503); + assert.equal(body.checks.fx.reason, 'FX_UNAVAILABLE'); +}); + +// ─── Liveness stays responsive during outages ──────────────────────────────── + +test('liveness remains 200 while readiness is not_ready', async () => { + dependencyHealth.forceDependencyState('payments', { mode: 'fail' }); + dependencyHealth.forceDependencyState('fx', { mode: 'fail' }); + dependencyHealth.forceDependencyState('store', { mode: 'fail' }); + + const live = await fetchJson('/api/health/live'); + assert.equal(live.status, 200); + assert.equal(live.body.status, 'alive'); + + const ready = await fetchJson('/api/health/ready'); + assert.equal(ready.status, 503); + assert.equal(ready.body.status, 'not_ready'); +}); + +test('base /api/health stays ok during dependency outages', async () => { + dependencyHealth.forceDependencyState('payments', { mode: 'fail' }); + const { status, body } = await fetchJson('/api/health'); + assert.equal(status, 200); + assert.equal(body.status, 'ok'); +}); + +// ─── Timeout: dependency checks cannot hang ────────────────────────────────── + +test('readiness returns 503 with *_TIMEOUT when a dependency hangs', async () => { + dependencyHealth.forceDependencyState('fx', { + mode: 'timeout', + delayMs: 5000, + }); + + const started = Date.now(); + const { status, body } = await fetchJson('/api/health/ready'); + const elapsed = Date.now() - started; + + assert.equal(status, 503); + assert.equal(body.checks.fx.status, 'error'); + assert.equal(body.checks.fx.reason, 'FX_TIMEOUT'); + // Must finish near the 50ms budget, not the 5s hang. + assert.ok(elapsed < 1000, `readiness hung for ${elapsed}ms`); +}); + +test('runCheck maps a hanging probe to a timeout reason within budget', async () => { + dependencyHealth.setProbeForTests('store', () => new Promise(() => {})); + const started = Date.now(); + const result = await dependencyHealth.runCheck('store', 40); + const elapsed = Date.now() - started; + + assert.equal(result.status, 'error'); + assert.equal(result.reason, 'STORE_TIMEOUT'); + assert.ok(elapsed < 500, `check hung for ${elapsed}ms`); +}); + +// ─── Recovery without restart ──────────────────────────────────────────────── + +test('readiness recovers after a failed dependency becomes healthy again', async () => { + dependencyHealth.forceDependencyState('payments', { + mode: 'fail', + reason: dependencyHealth.REASON.PAYMENTS_UNAVAILABLE, + }); + + const down = await fetchJson('/api/health/ready'); + assert.equal(down.status, 503); + assert.equal(down.body.checks.payments.reason, 'PAYMENTS_UNAVAILABLE'); + + // Clear the forced failure — no process restart, no module reload. + dependencyHealth.clearForcedStates(); + + const up = await fetchJson('/api/health/ready'); + assert.equal(up.status, 200); + assert.equal(up.body.status, 'ready'); + assert.equal(up.body.checks.payments.status, 'ok'); + assert.equal(up.body.checks.payments.reason, undefined); +}); + +// ─── Status codes ──────────────────────────────────────────────────────────── + +test('status codes: 200 when ready, 503 when any dependency fails', async () => { + const ok = await fetchJson('/api/health/ready'); + assert.equal(ok.status, 200); + + dependencyHealth.forceDependencyState('store', { mode: 'fail' }); + const bad = await fetchJson('/api/health/ready'); + assert.equal(bad.status, 503); + + dependencyHealth.clearForcedStates(); + const recovered = await fetchJson('/api/health/ready'); + assert.equal(recovered.status, 200); +}); + +// ─── Redaction ─────────────────────────────────────────────────────────────── + +test('readiness redacts raw error messages and secrets from responses', async () => { + const secret = 'postgres://user:super-secret@db.internal:5432/remitflow'; + dependencyHealth.forceDependencyState('store', { + mode: 'throw', + message: `${secret} password=hunter2 apiKey=sk_live_abc`, + }); + + const { status, body } = await fetchJson('/api/health/ready'); + assert.equal(status, 503); + assert.equal(body.checks.store.status, 'error'); + assert.equal(body.checks.store.reason, 'STORE_UNAVAILABLE'); + + const serialized = JSON.stringify(body); + assert.equal(serialized.includes('super-secret'), false); + assert.equal(serialized.includes('hunter2'), false); + assert.equal(serialized.includes('sk_live_abc'), false); + assert.equal(serialized.includes('postgres://'), false); + assert.equal(serialized.includes('password='), false); + // Only the stable reason code is present — no free-form message field. + assert.equal(body.checks.store.message, undefined); + assert.equal(body.checks.store.stack, undefined); +}); + +test('unit redact path never leaks custom throw messages', async () => { + dependencyHealth.forceDependencyState('payments', { + mode: 'throw', + message: 'Authorization: Bearer sk_live_should_not_leak', + }); + const result = await dependencyHealth.runCheck('payments', 50); + assert.equal(result.status, 'error'); + assert.equal(result.reason, 'PAYMENTS_UNAVAILABLE'); + assert.equal(JSON.stringify(result).includes('sk_live'), false); +}); diff --git a/test/smoke.test.js b/test/smoke.test.js index c7b039b..eb5c486 100644 --- a/test/smoke.test.js +++ b/test/smoke.test.js @@ -57,11 +57,14 @@ test('health liveness probe returns alive', async () => { assert.equal(body.status, 'alive'); }); -test('health readiness probe returns ready', async () => { +test('health readiness probe returns ready with dependency checks', async () => { const { status, body } = await fetchJson('/api/health/ready'); assert.equal(status, 200); assert.equal(body.status, 'ready'); - assert.deepEqual(body.checks, { store: 'ok' }); + assert.equal(body.checks.store.status, 'ok'); + assert.equal(body.checks.payments.status, 'ok'); + assert.equal(body.checks.fx.status, 'ok'); + assert.ok(typeof body.timeoutMs === 'number'); }); test('version endpoint returns name and version', async () => { From 47d93f81617b1416d8babc13defab1f32bbc9919 Mon Sep 17 00:00:00 2001 From: woahwhattheheck Date: Sat, 3 Oct 2026 03:45:12 -0400 Subject: [PATCH 2/6] fix(health): reject unusable advertised FX rates Validate the actual loaded rate values for every advertised currency so a missing or corrupt FX table cannot report ready while rate requests fail. Keep the existing success payload, FX_UNAVAILABLE code and corridors. Exercise table loss, liveness and same-process recovery through the actual HTTP app, plus partial invalid-rate states. The new data-loss regression failed against the previous source with 200 instead of 503. Validation: npm test passed all 271 tests; no tests skipped. --- src/services/rateService.js | 7 +++-- test/dependencyHealth.test.js | 54 +++++++++++++++++++++++++++++++++++ 2 files changed, 58 insertions(+), 3 deletions(-) diff --git a/src/services/rateService.js b/src/services/rateService.js index 29486dc..f37770f 100644 --- a/src/services/rateService.js +++ b/src/services/rateService.js @@ -79,13 +79,14 @@ function getPair(from, to) { /** * Lightweight FX probe used by readiness checks. - * Confirms the rate table is loaded and returns at least one currency. + * Confirms every advertised currency has a usable rate in the loaded table. * @returns {{ ok: true, currencies: number }} */ function ping() { const rates = listRates(); - if (!Array.isArray(rates) || rates.length === 0) { - const err = new Error('fx rate table empty'); + if (!Array.isArray(rates) || rates.length === 0 || + rates.some(({ rateToUsd }) => !Number.isFinite(rateToUsd) || rateToUsd <= 0)) { + const err = new Error('fx rate table unavailable'); err.reasonCode = 'FX_UNAVAILABLE'; throw err; } diff --git a/test/dependencyHealth.test.js b/test/dependencyHealth.test.js index 7e952ec..822538b 100644 --- a/test/dependencyHealth.test.js +++ b/test/dependencyHealth.test.js @@ -7,6 +7,7 @@ process.env.NODE_ENV = 'test'; const createApp = require('../src/app'); const config = require('../src/config'); +const { RATES_TO_USD } = require('../src/config/rates'); const dependencyHealth = require('../src/services/dependencyHealthService'); let server; @@ -114,6 +115,59 @@ test('readiness returns 503 when the FX dependency fails', async () => { assert.equal(body.checks.fx.reason, 'FX_UNAVAILABLE'); }); +test('readiness detects missing FX data and recovers when the table is restored', async () => { + const originalRates = { ...RATES_TO_USD }; + const originalPair = await fetchJson('/api/rates/USD-NGN'); + assert.equal(originalPair.status, 200); + + try { + for (const currency of Object.keys(RATES_TO_USD)) { + delete RATES_TO_USD[currency]; + } + + const unavailablePair = await fetchJson('/api/rates/USD-NGN'); + assert.equal(unavailablePair.status, 400); + + const down = await fetchJson('/api/health/ready'); + assert.equal(down.status, 503); + assert.equal(down.body.status, 'not_ready'); + assert.equal(down.body.checks.fx.status, 'error'); + assert.equal(down.body.checks.fx.reason, 'FX_UNAVAILABLE'); + assert.equal(down.body.checks.store.status, 'ok'); + assert.equal(down.body.checks.payments.status, 'ok'); + + const live = await fetchJson('/api/health/live'); + assert.equal(live.status, 200); + assert.equal(live.body.status, 'alive'); + } finally { + Object.assign(RATES_TO_USD, originalRates); + } + + const recovered = await fetchJson('/api/health/ready'); + assert.equal(recovered.status, 200); + assert.equal(recovered.body.checks.fx.status, 'ok'); + assert.equal(recovered.body.checks.fx.reason, undefined); + assert.deepEqual(await fetchJson('/api/rates/USD-NGN'), originalPair); +}); + +test('readiness rejects unusable rates in an advertised FX corridor', async () => { + const originalRate = RATES_TO_USD.NGN; + try { + for (const rate of [undefined, 0, -1, NaN, Infinity]) { + RATES_TO_USD.NGN = rate; + const down = await fetchJson('/api/health/ready'); + assert.equal(down.status, 503, `NGN rate ${String(rate)} must not be ready`); + assert.equal(down.body.checks.fx.reason, 'FX_UNAVAILABLE'); + } + } finally { + RATES_TO_USD.NGN = originalRate; + } + + const recovered = await fetchJson('/api/health/ready'); + assert.equal(recovered.status, 200); + assert.equal(recovered.body.checks.fx.status, 'ok'); +}); + // ─── Liveness stays responsive during outages ──────────────────────────────── test('liveness remains 200 while readiness is not_ready', async () => { From 7cf06801b387cc57f7cc2c2368d86aecd3a6ab19 Mon Sep 17 00:00:00 2001 From: woahwhattheheck Date: Sat, 3 Oct 2026 13:07:07 -0400 Subject: [PATCH 3/6] fix(health): keep liveness outside the API quota Serve only GET/HEAD liveness before the global API budget, preserving common middleware and the existing health controller. Readiness, business traffic and unmatched routes stay limited. Native createApp before/after: 17 local HTTP requests per phase prove availability after quota exhaustion and that liveness polls do not spend the business budget. Existing FX outage/recovery behavior is preserved. All 273 tests from the unchanged npm test selection pass on Node 24.19.0 with --test-concurrency=1 and retained packages matching all 76 lockfile versions. Focused health checks: 17 pass; exact parent with identical final tests: 2 fail / 15 pass. Syntax and whitespace checks pass. CI uses Node 22 and remains a separate hosted gate; services retain their in-memory/mock boundary. --- README.md | 5 +- src/app.js | 6 +++ src/routes/healthRoutes.js | 3 -- test/dependencyHealth.test.js | 95 +++++++++++++++++++++++++++++++++++ 4 files changed, 105 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index 46486a2..369ebed 100644 --- a/README.md +++ b/README.md @@ -228,7 +228,10 @@ The API implements Cache-Control response headers for security and efficiency: ### Health - `GET /api/health` — service health snapshot (version, uptime, env). Not a dependency gate. -- `GET /api/health/live` — liveness probe. Process-only; stays responsive during dependency outages. +- `GET /api/health/live` (also `HEAD`) — liveness probe. Process-only; stays responsive during dependency outages + and exhausted API quotas. Liveness polls do not consume the business API budget. Common security, cache, + parsing, timeout, request-ID and logging middleware still apply. Readiness, health snapshots and business + routes retain the API quota; unmatched methods or subpaths continue through the existing route stack. - `GET /api/health/ready` — readiness probe with bounded checks for store, payments (Stellar), and FX. Returns `200` when ready and `503` with redacted reason codes when a dependency is down or times out. Probes re-run on every request so recovery does not require a restart. diff --git a/src/app.js b/src/app.js index bd5b065..2c4f2e7 100644 --- a/src/app.js +++ b/src/app.js @@ -6,6 +6,7 @@ const morgan = require('morgan'); const config = require('./config'); const routes = require('./routes'); +const healthController = require('./controllers/healthController'); const securityHeaders = require('./middleware/securityHeaders'); const cacheControl = require('./middleware/cacheControl'); const requestTimeout = require('./middleware/requestTimeout'); @@ -46,6 +47,11 @@ function createApp() { } app.use(requestLogger); + // Process liveness must not consume or depend on the business API budget. + // Express also serves HEAD through this GET route; other paths/methods + // continue through the normal API middleware below. + app.get('/api/health/live', healthController.getLiveness); + // Basic abuse protection on the API surface. app.use('/api', rateLimit(config.rateLimit)); diff --git a/src/routes/healthRoutes.js b/src/routes/healthRoutes.js index b60f184..678aa48 100644 --- a/src/routes/healthRoutes.js +++ b/src/routes/healthRoutes.js @@ -9,9 +9,6 @@ const router = express.Router(); // GET /api/health router.get('/', asyncHandler(healthController.getHealth)); -// GET /api/health/live -router.get('/live', asyncHandler(healthController.getLiveness)); - // GET /api/health/ready router.get('/ready', asyncHandler(healthController.getReadiness)); diff --git a/test/dependencyHealth.test.js b/test/dependencyHealth.test.js index 822538b..7ec17a4 100644 --- a/test/dependencyHealth.test.js +++ b/test/dependencyHealth.test.js @@ -52,6 +52,101 @@ async function fetchJson(path) { return { status: res.status, body }; } +async function withSmallApiBudget(run) { + const originalMax = config.rateLimit.max; + const originalWindowMs = config.rateLimit.windowMs; + let budgetServer; + try { + config.rateLimit.max = 2; + config.rateLimit.windowMs = 60000; + budgetServer = createApp().listen(0, '127.0.0.1'); + } finally { + config.rateLimit.max = originalMax; + config.rateLimit.windowMs = originalWindowMs; + } + + try { + await new Promise((resolve, reject) => { + budgetServer.once('listening', resolve); + budgetServer.once('error', reject); + }); + await run(`http://127.0.0.1:${budgetServer.address().port}`); + } finally { + budgetServer.closeAllConnections(); + await new Promise((resolve) => budgetServer.close(resolve)); + } +} + +test('liveness GET and HEAD stay available after the business API quota is exhausted', async () => { + await withSmallApiBudget(async (url) => { + for (const status of [200, 200, 429]) { + const response = await fetch(`${url}/api/version`); + assert.equal(response.status, status); + await response.text(); + } + + for (const [method, path] of [ + ['GET', '/api/health/live'], + ['HEAD', '/api/health/live'], + ['GET', '/api/health/live/'], + ['GET', '/api/health/live?probe=quota'], + ]) { + const response = await fetch(url + path, { + method, + headers: { 'X-Request-Id': 'liveness-quota-test' }, + }); + assert.equal(response.status, 200, `${method} ${path}`); + assert.equal(response.headers.get('x-request-id'), 'liveness-quota-test'); + assert.equal(response.headers.get('x-content-type-options'), 'nosniff'); + assert.match(response.headers.get('cache-control'), /no-store/); + assert.equal(response.headers.get('x-ratelimit-limit'), null); + if (method === 'HEAD') { + assert.equal(await response.text(), ''); + } else { + assert.equal((await response.json()).status, 'alive'); + } + } + + for (const [method, path] of [ + ['GET', '/api/health/ready'], + ['GET', '/api/health'], + ['POST', '/api/health/live'], + ['GET', '/api/health/live/extra'], + ]) { + const response = await fetch(url + path, { method }); + assert.equal(response.status, 429, `${method} ${path} remains limited`); + assert.equal(response.headers.get('x-ratelimit-remaining'), '0'); + await response.text(); + } + }); +}); + +test('liveness polling during an FX outage does not consume the business API quota', async () => { + const originalRates = { ...RATES_TO_USD }; + try { + for (const currency of Object.keys(RATES_TO_USD)) delete RATES_TO_USD[currency]; + await withSmallApiBudget(async (url) => { + for (const method of ['GET', 'HEAD', 'GET']) { + const response = await fetch(`${url}/api/health/live`, { method }); + assert.equal(response.status, 200); + await response.text(); + } + + const ready = await fetch(`${url}/api/health/ready`); + assert.equal(ready.status, 503); + assert.equal((await ready.json()).checks.fx.reason, 'FX_UNAVAILABLE'); + + for (const status of [200, 429]) { + const response = await fetch(`${url}/api/version`); + assert.equal(response.status, status); + await response.text(); + } + }); + } finally { + Object.assign(RATES_TO_USD, originalRates); + } +}); + // ─── Happy path ────────────────────────────────────────────────────────────── test('readiness is ready when store, payments, and fx are healthy', async () => { From 67c4e52d8bcb6b286922104702e717c75cdfa053 Mon Sep 17 00:00:00 2001 From: woahwhattheheck Date: Sat, 3 Oct 2026 18:56:31 -0400 Subject: [PATCH 4/6] fix(health): await payment and FX readiness adapters --- README.md | 6 + src/services/dependencyHealthService.js | 4 +- test/dependencyHealth.test.js | 194 ++++++++++++++++++++++++ 3 files changed, 202 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 369ebed..db5e4ac 100644 --- a/README.md +++ b/README.md @@ -236,6 +236,12 @@ The API implements Cache-Control response headers for security and efficiency: Returns `200` when ready and `503` with redacted reason codes when a dependency is down or times out. Probes re-run on every request so recovery does not require a restart. Tune the per-check budget with `HEALTH_CHECK_TIMEOUT_MS` (default `1000`). + Payment and FX `ping()` adapters may return their result immediately or as a Promise; + the default probes await that result before checking `ok`. Rejections remain attached + to the readiness error path, including after a timeout. The deadline bounds waiting, + but does not cancel the underlying Promise or interrupt synchronous blocking work. + The shipped adapters remain synchronous mocks: Stellar configuration and loaded FX + data checks do not establish live provider reachability. - `GET /api/version` — service name and version. ### Rates & quotes diff --git a/src/services/dependencyHealthService.js b/src/services/dependencyHealthService.js index 4efeed0..5a361b2 100644 --- a/src/services/dependencyHealthService.js +++ b/src/services/dependencyHealthService.js @@ -58,7 +58,7 @@ const defaultProbes = Object.freeze({ }, async payments() { - const result = stellarService.ping(); + const result = await stellarService.ping(); if (!result || result.ok !== true) { const err = new Error('payments unavailable'); err.reasonCode = REASON.PAYMENTS_UNAVAILABLE; @@ -68,7 +68,7 @@ const defaultProbes = Object.freeze({ }, async fx() { - const result = rateService.ping(); + const result = await rateService.ping(); if (!result || result.ok !== true) { const err = new Error('fx unavailable'); err.reasonCode = REASON.FX_UNAVAILABLE; diff --git a/test/dependencyHealth.test.js b/test/dependencyHealth.test.js index 7ec17a4..559b0b8 100644 --- a/test/dependencyHealth.test.js +++ b/test/dependencyHealth.test.js @@ -388,3 +388,197 @@ test('unit redact path never leaks custom throw messages', async () => { assert.equal(result.reason, 'PAYMENTS_UNAVAILABLE'); assert.equal(JSON.stringify(result).includes('sk_live'), false); }); + +// Exercise the default wrappers by replacing adapter methods, not health probes. +const { execFile: execHealthChild } = require('node:child_process'); +const healthTestRepoRoot = require('node:path').resolve(__dirname, '..'); +const asyncHealthAdapters = [ + { name: 'payments', reasonPrefix: 'PAYMENTS', service: require('../src/services/stellarService') }, + { name: 'fx', reasonPrefix: 'FX', service: require('../src/services/rateService') }, +]; + +// Detached provider rejections must fail only this child on an unrepaired head. +// No unhandledRejection listener is installed in either process. +const healthRejectionChildSource = '(' + (async function healthRejectionChild() { + process.env.NODE_ENV = 'test'; + const createChildApp = require('./src/app'); + const childConfig = require('./src/config'); + const name = process.argv[1]; + const mode = process.argv[2]; + const provider = name === 'payments' + ? require('./src/services/stellarService') + : require('./src/services/rateService'); + const originalPing = provider.ping; + const secret = 'RF140_SYNTHETIC_ASYNC_PROVIDER_SECRET'; + const timeoutMs = 40; + childConfig.health.checkTimeoutMs = timeoutMs; + + let childServer; + let rejectionTimer; + let result; + let markRejectionDelivered; + const rejectionDelivered = new Promise((resolve) => { + markRejectionDelivered = resolve; + }); + + try { + childServer = createChildApp().listen(0, '127.0.0.1'); + await new Promise((resolve, reject) => { + childServer.once('listening', resolve); + childServer.once('error', reject); + }); + const url = 'http://127.0.0.1:' + childServer.address().port; + async function requestJson(route) { + const response = await fetch(url + route, { signal: AbortSignal.timeout(2000) }); + return { status: response.status, body: await response.json() }; + } + + provider.ping = () => new Promise((resolve, reject) => { + const fail = () => { + reject(new Error(secret)); + markRejectionDelivered(); + }; + if (mode === 'late') { + rejectionTimer = setTimeout(fail, timeoutMs * 3); + } else { + fail(); + } + }); + + const started = Date.now(); + const down = await requestJson('/api/health/ready'); + const elapsedMs = Date.now() - started; + const live = await requestJson('/api/health/live'); + await rejectionDelivered; + // Let a detached rejection take its normal fatal path before reporting. + await new Promise((resolve) => setImmediate(resolve)); + + provider.ping = originalPing; + const recovered = await requestJson('/api/health/ready'); + result = { name, mode, timeoutMs, elapsedMs, down, live, recovered }; + } finally { + provider.ping = originalPing; + clearTimeout(rejectionTimer); + if (childServer) { + childServer.closeAllConnections(); + await new Promise((resolve) => childServer.close(resolve)); + } + } + + process.stdout.write('RF140_CHILD_RESULT ' + JSON.stringify(result) + '\n'); +}).toString() + ')().catch((error) => { console.error(error); process.exitCode = 1; });'; + +async function observeHealthRejectionChild(name, mode) { + const child = await new Promise((resolve) => { + execHealthChild( + process.execPath, + ['--max-old-space-size=64', '--unhandled-rejections=strict', '-e', healthRejectionChildSource, name, mode], + { + cwd: healthTestRepoRoot, + env: { ...process.env, NODE_ENV: 'test' }, + timeout: 10000, + maxBuffer: 128 * 1024, + }, + (error, stdout, stderr) => resolve({ + exitCode: error ? (error.code ?? 'PROCESS_ERROR') : 0, + signal: error ? error.signal : null, + stdout, + stderr, + }) + ); + }); + + assert.equal( + child.exitCode, + 0, + name + ' ' + mode + ' child exit=' + child.exitCode + + ' signal=' + child.signal + '\n' + child.stderr + ); + const marker = 'RF140_CHILD_RESULT '; + const resultLine = child.stdout.split('\n').find((line) => line.startsWith(marker)); + assert.ok(resultLine, 'child did not report completed HTTP observations:\n' + child.stdout); + return JSON.parse(resultLine.slice(marker.length)); +} + +for (const { name, reasonPrefix, service } of asyncHealthAdapters) { + test(name + ' default adapter awaits async healthy and negative results and recovers', async () => { + const originalPing = service.ping; + try { + service.ping = async () => ({ ok: true }); + const healthy = await fetchJson('/api/health/ready'); + assert.equal(healthy.status, 200); + assert.equal(healthy.body.checks[name].status, 'ok'); + + service.ping = async () => ({ ok: false }); + const down = await fetchJson('/api/health/ready'); + assert.equal(down.status, 503); + assert.equal(down.body.checks[name].status, 'error'); + assert.equal(down.body.checks[name].reason, reasonPrefix + '_UNAVAILABLE'); + + service.ping = async () => ({ ok: true }); + const recovered = await fetchJson('/api/health/ready'); + assert.equal(recovered.status, 200); + assert.equal(recovered.body.checks[name].status, 'ok'); + assert.equal(recovered.body.checks[name].reason, undefined); + } finally { + service.ping = originalPing; + } + }); + + test(name + ' default adapter bounds a pending ping while liveness and recovery remain available', async () => { + const originalPing = service.ping; + try { + service.ping = () => new Promise(() => {}); + const started = Date.now(); + const [down, live] = await Promise.all([ + fetchJson('/api/health/ready'), + fetchJson('/api/health/live'), + ]); + const elapsedMs = Date.now() - started; + + assert.equal(down.status, 503); + assert.equal(down.body.checks[name].status, 'error'); + assert.equal(down.body.checks[name].reason, reasonPrefix + '_TIMEOUT'); + assert.ok(elapsedMs < 1000, 'pending ' + name + ' check took ' + elapsedMs + 'ms'); + assert.equal(live.status, 200); + assert.equal(live.body.status, 'alive'); + + service.ping = originalPing; + const recovered = await fetchJson('/api/health/ready'); + assert.equal(recovered.status, 200); + assert.equal(recovered.body.checks[name].status, 'ok'); + assert.equal(recovered.body.checks[name].reason, undefined); + } finally { + service.ping = originalPing; + } + }); + + test(name + ' default adapter catches and redacts an immediate rejection without killing the process', async () => { + const observed = await observeHealthRejectionChild(name, 'immediate'); + assert.equal(observed.down.status, 503); + assert.equal(observed.down.body.checks[name].status, 'error'); + assert.equal(observed.down.body.checks[name].reason, reasonPrefix + '_UNAVAILABLE'); + assert.equal(JSON.stringify(observed.down.body).includes('RF140_SYNTHETIC_ASYNC_PROVIDER_SECRET'), false); + assert.equal(observed.down.body.checks[name].message, undefined); + assert.equal(observed.down.body.checks[name].stack, undefined); + assert.equal(observed.live.status, 200); + assert.equal(observed.live.body.status, 'alive'); + assert.equal(observed.recovered.status, 200); + assert.equal(observed.recovered.body.checks[name].status, 'ok'); + assert.equal(observed.recovered.body.checks[name].reason, undefined); + }); + + test(name + ' default adapter survives a rejection after its deadline and recovers', async () => { + const observed = await observeHealthRejectionChild(name, 'late'); + assert.equal(observed.down.status, 503); + assert.equal(observed.down.body.checks[name].status, 'error'); + assert.equal(observed.down.body.checks[name].reason, reasonPrefix + '_TIMEOUT'); + assert.ok(observed.elapsedMs < 1000, 'late ' + name + ' check took ' + observed.elapsedMs + 'ms'); + assert.equal(JSON.stringify(observed.down.body).includes('RF140_SYNTHETIC_ASYNC_PROVIDER_SECRET'), false); + assert.equal(observed.live.status, 200); + assert.equal(observed.live.body.status, 'alive'); + assert.equal(observed.recovered.status, 200); + assert.equal(observed.recovered.body.checks[name].status, 'ok'); + assert.equal(observed.recovered.body.checks[name].reason, undefined); + }); +} From b4c9505ee1ac758ebb3fff659bbe39747febaa6b Mon Sep 17 00:00:00 2001 From: woahwhattheheck Date: Sun, 4 Oct 2026 09:54:27 +0000 Subject: [PATCH 5/6] fix: bound unfinished readiness probes across overlapping requests [skip ci] --- src/services/dependencyHealthService.js | 32 +++++-- src/services/inFlightProbe.js | 52 +++++++++++ test/readinessProbeSharing.test.js | 118 ++++++++++++++++++++++++ 3 files changed, 193 insertions(+), 9 deletions(-) create mode 100644 src/services/inFlightProbe.js create mode 100644 test/readinessProbeSharing.test.js diff --git a/src/services/dependencyHealthService.js b/src/services/dependencyHealthService.js index 5a361b2..f1d574f 100644 --- a/src/services/dependencyHealthService.js +++ b/src/services/dependencyHealthService.js @@ -4,6 +4,8 @@ const config = require('../config'); const { store } = require('../store'); const rateService = require('./rateService'); const stellarService = require('./stellarService'); +const { createProbePool } = require('./inFlightProbe'); +const probePool = createProbePool(); /** * Dependency-aware readiness diagnostics. @@ -11,8 +13,8 @@ const stellarService = require('./stellarService'); * Separates process liveness from traffic readiness by probing the store * (database stand-in), payment provider (Stellar), and FX rate table with * a per-check time budget. Failures surface as stable, redacted reason - * codes — never raw messages, stacks, or connection material — and every - * probe is re-evaluated on the next request so recovery does not require + * codes — never raw messages, stacks, or connection material — and completed + * probes are re-evaluated on the next request so recovery does not require * a process restart. */ @@ -57,8 +59,8 @@ const defaultProbes = Object.freeze({ return { ok: true }; }, - async payments() { - const result = await stellarService.ping(); + async payments(ping = stellarService.ping) { + const result = await ping.call(stellarService); if (!result || result.ok !== true) { const err = new Error('payments unavailable'); err.reasonCode = REASON.PAYMENTS_UNAVAILABLE; @@ -67,8 +69,8 @@ const defaultProbes = Object.freeze({ return { ok: true }; }, - async fx() { - const result = await rateService.ping(); + async fx(ping = rateService.ping) { + const result = await ping.call(rateService); if (!result || result.ok !== true) { const err = new Error('fx unavailable'); err.reasonCode = REASON.FX_UNAVAILABLE; @@ -218,7 +220,18 @@ async function runCheck(name, timeoutMs = checkTimeoutMs()) { err.reasonCode = REASON.CHECK_ERROR; throw err; } - await withTimeout(forced || probe(), timeoutMs, timeoutReason); + // Replacing an adapter must not keep joining its predecessor's hung call. + // Capture that method now, including its receiver, before deferred execution. + const adapter = probe === defaultProbes.payments ? stellarService.ping + : probe === defaultProbes.fx ? rateService.ping : probe; + const invoke = probe === defaultProbes.payments || probe === defaultProbes.fx + ? () => probe(adapter) : probe; + const subscription = forced ? null : probePool.acquire(name, invoke, adapter); + try { + await withTimeout(forced || subscription.promise, timeoutMs, timeoutReason); + } finally { + if (subscription) subscription.release(); + } return { name, status: 'ok', @@ -235,8 +248,8 @@ async function runCheck(name, timeoutMs = checkTimeoutMs()) { } /** - * Evaluate every dependency. Safe to call on every readiness request — - * nothing is cached as permanently failed. + * Evaluate every dependency. Overlapping requests share only unfinished + * probes; each caller keeps its own deadline and completed checks are not cached. * @param {object} [options] * @param {number} [options.timeoutMs] * @returns {Promise<{ @@ -300,6 +313,7 @@ function setProbeForTests(name, fn) { /** Restore default probes and clear forced states (tests only). */ function resetForTests() { + probePool.clear(); for (const name of DEPENDENCIES) { probes[name] = defaultProbes[name]; } diff --git a/src/services/inFlightProbe.js b/src/services/inFlightProbe.js new file mode 100644 index 0000000..5dfabe9 --- /dev/null +++ b/src/services/inFlightProbe.js @@ -0,0 +1,52 @@ +'use strict'; + +/** + * Share only unfinished provider work. Each caller owns a detachable waiter, + * so a timed-out readiness request leaves no callback on a hung provider. + * The three dependency names bound retained operations in normal use. + */ +function createProbePool() { + const pending = new Map(); + + function acquire(name, probe, identity = probe) { + let entry = pending.get(name); + if (!entry || entry.identity !== identity) { + entry = { identity, waiters: new Set() }; + pending.set(name, entry); + const current = entry; + const finish = (ok, value) => { + // Replaced adapters may finish after their successor started. + if (pending.get(name) === current) pending.delete(name); + for (const waiter of current.waiters) { + if (ok) waiter.resolve(value); + else waiter.reject(value); + } + current.waiters.clear(); + }; + // One fulfillment/rejection pair per actual operation, not per request. + // Synchronous throws and thenables follow the same failure path. + Promise.resolve().then(probe).then( + (value) => finish(true, value), + (error) => finish(false, error) + ); + } + + let waiter; + const promise = new Promise((resolve, reject) => { + waiter = { resolve, reject }; + entry.waiters.add(waiter); + }); + return { + promise, + release() { entry.waiters.delete(waiter); }, + }; + } + + return { + acquire, + // Test adapter resets cannot let an old operation remove its replacement. + clear() { pending.clear(); }, + }; +} + +module.exports = { createProbePool }; diff --git a/test/readinessProbeSharing.test.js b/test/readinessProbeSharing.test.js new file mode 100644 index 0000000..6148d44 --- /dev/null +++ b/test/readinessProbeSharing.test.js @@ -0,0 +1,118 @@ +'use strict'; + +const { test, before, after, beforeEach, afterEach } = require('node:test'); +const assert = require('node:assert/strict'); +const { setImmediate: tick } = require('node:timers/promises'); + +process.env.NODE_ENV = 'test'; +const createApp = require('../src/app'); +const config = require('../src/config'); +const health = require('../src/services/dependencyHealthService'); + +let server; +let base; +let originalTimeout; + +before(async () => { + originalTimeout = config.health.checkTimeoutMs; + server = createApp().listen(0, '127.0.0.1'); + await new Promise((resolve, reject) => { + server.once('listening', resolve); + server.once('error', reject); + }); + base = `http://127.0.0.1:${server.address().port}`; +}); +beforeEach(() => { health.resetForTests(); config.health.checkTimeoutMs = 5; }); +afterEach(() => health.resetForTests()); +after(async () => { + config.health.checkTimeoutMs = originalTimeout; + health.resetForTests(); + server.closeAllConnections(); + await new Promise((resolve) => server.close(resolve)); +}); + +test('repeated readiness timeouts do not multiply unfinished provider work', async () => { + let calls = 0; + let finish; + const hung = new Promise((resolve) => { finish = resolve; }); + health.setProbeForTests('payments', () => { calls++; return hung; }); + try { + for (let i = 0; i < 40; i++) { + const response = await fetch(`${base}/api/health/ready`); + assert.equal(response.status, 503); + assert.equal((await response.json()).checks.payments.reason, 'PAYMENTS_TIMEOUT'); + } + const live = await fetch(`${base}/api/health/live`); + assert.equal(live.status, 200); + assert.equal((await live.json()).status, 'alive'); + assert.equal(calls, 1, '40 timed-out requests must share one unfinished provider call'); + } finally { + finish({ ok: true }); + await tick(); + } + // The settled result is not a health cache: the next request calls again. + const response = await fetch(`${base}/api/health/ready`); + assert.equal(response.status, 200); + assert.equal((await response.json()).checks.payments.status, 'ok'); + assert.equal(calls, 2); +}); + +test('each waiter retains its own deadline and a longer waiter can recover', async () => { + let calls = 0; + let finish; + const pending = new Promise((resolve) => { finish = resolve; }); + health.setProbeForTests('fx', () => { calls++; return pending; }); + const short = health.runCheck('fx', 10); + const long = health.runCheck('fx', 1000); + try { + assert.equal((await short).reason, 'FX_TIMEOUT'); + finish({ ok: true }); + assert.equal((await long).status, 'ok'); + assert.equal(calls, 1); + assert.equal((await health.runCheck('fx', 1000)).status, 'ok'); + assert.equal(calls, 2, 'completed checks must be probed anew'); + } finally { finish({ ok: true }); await long; } +}); + +test('shared rejection is redacted for every waiter and does not stick', async () => { + let calls = 0; + let fail; + const pending = new Promise((_, reject) => { fail = reject; }); + health.setProbeForTests('payments', () => { calls++; return pending; }); + const first = health.runCheck('payments', 1000); + const second = health.runCheck('payments', 1000); + await tick(); + fail(new Error('Authorization: Bearer private-test-secret')); + const results = await Promise.all([first, second]); + assert.equal(calls, 1); + for (const result of results) { + assert.equal(result.reason, 'PAYMENTS_UNAVAILABLE'); + assert.equal(JSON.stringify(result).includes('private-test-secret'), false); + } + health.setProbeForTests('payments', async () => ({ ok: true })); + assert.equal((await health.runCheck('payments', 1000)).status, 'ok'); +}); + +test('an old completion cannot remove a replacement adapter operation', async () => { + let finishOld; + let finishNew; + let newCalls = 0; + const oldPending = new Promise((resolve) => { finishOld = resolve; }); + const newPending = new Promise((resolve) => { finishNew = resolve; }); + health.setProbeForTests('store', () => oldPending); + const old = health.runCheck('store', 1000); + await tick(); + health.setProbeForTests('store', () => { newCalls++; return newPending; }); + const replacement = health.runCheck('store', 1000); + try { + await tick(); + finishOld({ ok: true }); + assert.equal((await old).status, 'ok'); + const joined = health.runCheck('store', 1000); + await tick(); + finishNew({ ok: true }); + assert.equal((await replacement).status, 'ok'); + assert.equal((await joined).status, 'ok'); + assert.equal(newCalls, 1); + } finally { finishOld({ ok: true }); finishNew({ ok: true }); } +}); From 4fd4284ccead15abede008a95779c680aab2682b Mon Sep 17 00:00:00 2001 From: woahwhattheheck Date: Sun, 4 Oct 2026 05:57:56 -0400 Subject: [PATCH 6/6] docs: record bounded readiness work and recovery evidence [skip ci] --- docs/READINESS_PROBE_SHARING.md | 34 +++++++++++++++++++++++++++++++++ 1 file changed, 34 insertions(+) create mode 100644 docs/READINESS_PROBE_SHARING.md diff --git a/docs/READINESS_PROBE_SHARING.md b/docs/READINESS_PROBE_SHARING.md new file mode 100644 index 0000000..f8b6fe1 --- /dev/null +++ b/docs/READINESS_PROBE_SHARING.md @@ -0,0 +1,34 @@ +# Bounded readiness probe work + +Follow-up to the existing dependency-readiness contribution, PR #140 / issue #134. + +## Behavior + +A request deadline previously stopped waiting without stopping the underlying operation. Repeated readiness requests could therefore accumulate unfinished calls to a hung payment or FX adapter. The service now shares only an unfinished operation for each dependency and adapter generation. Every request retains its own timeout; a timed-out request removes its subscriber instead of leaving another callback attached to the hung operation. Completed results are discarded, not cached as healthy or failed. + +The payment and FX method references and their receiver objects are captured before execution. Replacing an adapter does not inherit the old adapter's pending operation; an old completion cannot evict its successor. Existing stable reason codes, redaction, forced test states and the liveness/business-API rate-limit distinction are unchanged. + +**Limit:** this bounds outstanding work for an unchanged adapter; it does not cancel an arbitrary Promise. A truly never-settling adapter stays not-ready rather than spawning more calls. Recovery occurs when the operation settles or the adapter is replaced. A real external driver still needs transport-level cancellation/timeouts. Deliberately replacing adapters repeatedly can leave their old uncancellable operations outstanding; this is not a global provider quota or fleet-wide limiter. + +## Executed evidence — October 4, 2026 + +Source commit `b4c9505ee1ac758ebb3fff659bbe39747febaa6b`, tree `9a49ce0645d3b95eed16aa8d2671dbc19682e989`, sole parent `67c4e52d8bcb6b286922104702e717c75cdfa053`. + +- Production service blob: `f1d574ff23df4440aa350ca57ace9df51b88080b`. +- Shared-operation helper: `5dfabe9ed6437b00e336b0eb8e3dca58e9bc5279`. +- New regression file: `6148d44d1e1f46d05a2920531ae0bf9986028b44`. + +[Final run 37193635966](https://github.com/woahwhattheheck/RemitFlow-Backend/actions/runs/37193635966), Node 24.21.0 on Ubuntu, executed: + +```sh +npm ci --no-audit --no-fund +node --unhandled-rejections=strict --test --test-reporter=tap test/dependencyHealth.test.js test/readinessProbeSharing.test.js +``` + +**29 tests passed; zero failures, skips or cancellations.** Existing health tests and dependency manifests stayed byte-identical. This includes actual Express HTTP requests and default adapter wrappers, with controlled local provider promises; no live Stellar/FX/database connection is claimed. + +The repeated-timeout regression sent 40 readiness requests. Before the repair it made 40 provider calls; afterward it made one unfinished provider call, retained liveness HTTP 200 and recovered after completion. The other new cases cover independent caller deadlines, shared redacted rejection and a late predecessor completion. The earlier default-adapter tests also prove replacement recovery. + +The initial baseline run [37192889105](https://github.com/woahwhattheheck/RemitFlow-Backend/actions/runs/37192889105) reproduced all four new failures. Its wrapper then stopped because it expected TAP while Node emitted the spec reporter; that run is not represented as a repair pass. The first repaired candidate [37193040714](https://github.com/woahwhattheheck/RemitFlow-Backend/actions/runs/37193040714) passed 27 tests but failed two existing adapter-replacement recovery cases. Capturing and keying the actual adapter generation repaired those failures without changing any existing assertion. The original baseline was not rerun. + +[Final artifact 11299902454](https://github.com/woahwhattheheck/RemitFlow-Backend/actions/runs/37193635966/artifacts/11299902454) retains TAP output, commands, source diff and identities. ZIP SHA-256: `a08f9dad1374587e4b202b190fa97b897a602522d20888a670d45dd726995c5f`. Earlier baseline and failed-candidate artifacts remain attached to their runs. This result does not establish full-suite, deployment, live-provider, bounty-acceptance or payment status.