From c53a4c7fc31af2c8f9f0b5db53fce765b8d4f4c6 Mon Sep 17 00:00:00 2001 From: Ved-viraj Date: Mon, 28 Sep 2026 14:53:50 +0100 Subject: [PATCH 1/2] feat(backend): add prometheus alert metrics and health telemetry to Exchange Rate Oracle Cache Track quote-cache load outcomes and expose a rolling health snapshot so operators can alert on internal errors, timeouts, and stale lookups. --- backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md | 61 +++- .../exchange-rate-oracle-cache.rules.yml | 78 +++++ backend/src/app.js | 26 ++ backend/src/lib/exchange-rate-cache.js | 13 + .../src/lib/exchange-rate-oracle-telemetry.js | 317 ++++++++++++++++++ .../exchange-rate-oracle-telemetry.test.js | 212 ++++++++++++ backend/src/routes/payments.js | 98 ++---- backend/src/routes/prometheus.js | 6 +- backend/tests/helpers/fake-redis.js | 223 +++++------- .../integration/exchange-rate-cache.test.js | 17 + 10 files changed, 830 insertions(+), 221 deletions(-) create mode 100644 backend/docs/alerts/exchange-rate-oracle-cache.rules.yml create mode 100644 backend/src/lib/exchange-rate-oracle-telemetry.js create mode 100644 backend/src/lib/exchange-rate-oracle-telemetry.test.js diff --git a/backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md b/backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md index 7197058e..06da2e52 100644 --- a/backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md +++ b/backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md @@ -3,7 +3,8 @@ Caching and concurrency control for path-payment exchange-rate quotes (`GET /api/path-payment-quote/:id`). -Covers issues **#1445** (distributed concurrency control and locking) and +Covers issues **#1443** (Prometheus alert metrics and health telemetry), +**#1445** (distributed concurrency control and locking) and **#1446** (integration and stress test suite). --- @@ -16,6 +17,8 @@ Covers issues **#1445** (distributed concurrency control and locking) and | `src/lib/exchange-rate-coordinator.js` | Cross-instance coordination: Redis lock + shared quote store | | `src/services/exchangeRateService.js` | `getExchangeRateQuote()` composes both layers around the Horizon query | | `src/lib/path-payment-metrics.js` | Prometheus series (existing cache metrics + concurrency metrics) | +| `src/lib/exchange-rate-oracle-telemetry.js` | Alert metrics, rolling-window health, separate registry (#1443) | +| `docs/alerts/exchange-rate-oracle-cache.rules.yml` | Prometheus alert rules (#1443) | ``` getExchangeRateQuote(key) @@ -154,11 +157,65 @@ HTTP bursts wait until every request has joined the in-flight load, using Horizon response. supertest opens a separate server per request, so arrival order is otherwise not guaranteed. +## 8. Alert metrics and health (#1443) + +Lookups and loads are counted in `exchange-rate-oracle-telemetry.js`, separate +from the path-payment series so a scrape can alert on the cache itself. +`/metrics` merges the registry. Every label is from a fixed set (`hit`, +`miss`, `stale`, `success`, `error`, `timeout`, `not_found`). Asset codes, +issuers, amounts and cache keys are never labels. + +| Metric | Type | Labels | +|---|---|---| +| `exchange_rate_oracle_cache_lookups_total` | counter | `result` (hit/miss/stale) | +| `exchange_rate_oracle_cache_loads_total` | counter | `outcome` (success/error/timeout/not_found) | +| `exchange_rate_oracle_cache_load_duration_seconds` | histogram | `outcome` | +| `exchange_rate_oracle_cache_health_state` | gauge | 0 healthy, 1 degraded, 2 unhealthy | +| `exchange_rate_oracle_cache_error_ratio` | gauge | rolling window | +| `exchange_rate_oracle_cache_timeout_ratio` | gauge | rolling window | +| `exchange_rate_oracle_cache_stale_ratio` | gauge | rolling window | +| `exchange_rate_oracle_cache_last_load_timestamp_seconds` | gauge | none | + +`not_found` is a normal "Horizon has no path" result. It is counted, but it +does not move the error ratio and cannot mark the cache unhealthy. + +`GET /health/exchange-rate-oracle-cache` (public, no quote or account data): + +| Status | Condition | HTTP | +|---|---|---| +| `unhealthy` | ≥ `min_samples` loads and internal error ratio ≥ threshold | 503 | +| `degraded` | ≥ `min_samples` loads and timeout ratio ≥ threshold, **or** ≥ `min_samples` lookups and stale ratio ≥ threshold | 200 | +| `healthy` | otherwise, including no traffic | 200 | + +`GET /health` also reports `services.exchange_rate_oracle_cache`. That value +does not change `ok` or the status code. Gauges refresh at scrape time, so +they decay when traffic stops. The window is a fixed ring of 5-second buckets. + +| Variable | Default | +|---|---| +| `EXCHANGE_RATE_ORACLE_HEALTH_WINDOW_MS` | `300000` | +| `EXCHANGE_RATE_ORACLE_HEALTH_MIN_SAMPLES` | `20` | +| `EXCHANGE_RATE_ORACLE_ERROR_RATIO_THRESHOLD` | `0.05` | +| `EXCHANGE_RATE_ORACLE_TIMEOUT_RATIO_THRESHOLD` | `0.2` | +| `EXCHANGE_RATE_ORACLE_STALE_RATIO_THRESHOLD` | `0.5` | + +`docs/alerts/exchange-rate-oracle-cache.rules.yml` alerts on unhealthy state, +error ratio, repeated timeouts, a high stale ratio, and p99 load latency +above 5s. A unit test checks that every metric named in the rules file is +registered. + +Failed loads are logged at `warn` with outcome, duration, status and error +message. The cache key is not logged. A failure inside the metrics client is +swallowed so it cannot replace the loader's error. + +## 9. Tests + The suites were mutation-checked against `exchange-rate-cache.js`. Disabling single-flight fails 16 tests. Dropping the invalidation guard fails 4. ``` -npx vitest run src/lib/exchange-rate-cache.test.js \ +npx vitest run src/lib/exchange-rate-oracle-telemetry.test.js \ + src/lib/exchange-rate-cache.test.js \ src/lib/exchange-rate-coordinator.test.js \ src/services/exchangeRateService.test.js \ tests/integration/exchange-rate-cache.test.js diff --git a/backend/docs/alerts/exchange-rate-oracle-cache.rules.yml b/backend/docs/alerts/exchange-rate-oracle-cache.rules.yml new file mode 100644 index 00000000..9a2b3c80 --- /dev/null +++ b/backend/docs/alerts/exchange-rate-oracle-cache.rules.yml @@ -0,0 +1,78 @@ +# Prometheus alerting rules for the Exchange Rate Oracle Cache (issue #1443). +# +# Load with: rule_files: ["docs/alerts/exchange-rate-oracle-cache.rules.yml"] +# Validate: promtool check rules docs/alerts/exchange-rate-oracle-cache.rules.yml +# +# Metric reference: docs/EXCHANGE_RATE_ORACLE_CACHE.md + +groups: + - name: exchange-rate-oracle-cache + rules: + - alert: ExchangeRateOracleCacheUnhealthy + expr: max(exchange_rate_oracle_cache_health_state) >= 2 + for: 5m + labels: + severity: critical + annotations: + summary: Exchange rate oracle cache reports unhealthy + description: > + Rolling-window internal error ratio is above threshold. Quote loads + are failing. See GET /health/exchange-rate-oracle-cache for counts + and reasons. "No path" results are excluded from this ratio. + + - alert: ExchangeRateOracleCacheHighErrorRatio + expr: | + sum(rate(exchange_rate_oracle_cache_loads_total{outcome="error"}[10m])) + / + clamp_min(sum(rate(exchange_rate_oracle_cache_loads_total[10m])), 1e-9) + > 0.05 + and sum(increase(exchange_rate_oracle_cache_loads_total[10m])) >= 20 + for: 10m + labels: + severity: critical + annotations: + summary: More than 5% of exchange-rate oracle loads are internal errors + description: > + Horizon or the quote loader is failing. Check logs for + "Exchange rate oracle cache load failed". + + - alert: ExchangeRateOracleCacheLoadTimeouts + expr: sum(increase(exchange_rate_oracle_cache_loads_total{outcome="timeout"}[10m])) > 5 + for: 10m + labels: + severity: warning + annotations: + summary: Exchange-rate oracle loads are timing out + description: > + More than 5 loads hit the load timeout in 10m. Waiters receive 504. + A hung Horizon call is the usual cause. + + - alert: ExchangeRateOracleCacheHighStaleRatio + expr: | + sum(rate(exchange_rate_oracle_cache_lookups_total{result="stale"}[10m])) + / + clamp_min(sum(rate(exchange_rate_oracle_cache_lookups_total[10m])), 1e-9) + > 0.5 + and sum(increase(exchange_rate_oracle_cache_lookups_total[10m])) >= 20 + for: 10m + labels: + severity: warning + annotations: + summary: Most exchange-rate cache lookups are past the fresh TTL + description: > + Quotes are being refreshed constantly. The TTL may be shorter than + the request interval, or invalidation is too aggressive. + + - alert: ExchangeRateOracleCacheSlowLoads + expr: | + histogram_quantile(0.99, + sum by (le) (rate(exchange_rate_oracle_cache_load_duration_seconds_bucket[10m])) + ) > 5 + for: 10m + labels: + severity: warning + annotations: + summary: Exchange-rate oracle load p99 latency above 5s + description: > + Upstream quote latency is high. The HTTP load timeout is 15s by + default (EXCHANGE_RATE_LOAD_TIMEOUT_MS). diff --git a/backend/src/app.js b/backend/src/app.js index 503fc3cd..b330d308 100644 --- a/backend/src/app.js +++ b/backend/src/app.js @@ -48,6 +48,7 @@ import { import { versionDeprecationMiddleware } from "./lib/version-deprecation.js"; import oracleRouter from "./routes/oracle.js"; import { getPaymentSessionValidatorHealth } from "./lib/payment-session-validator.js"; +import { getExchangeRateOracleHealth } from "./lib/exchange-rate-oracle-telemetry.js"; import { configureExchangeRateCoordination } from "./services/exchangeRateService.js"; export async function createApp({ redisClient }) { @@ -260,6 +261,7 @@ export async function createApp({ redisClient }) { redis: redisAvailable ? "ok" : "unavailable", // Informational only — does not affect `ok` / the status code. payment_session_validator: getPaymentSessionValidatorHealth().status, + exchange_rate_oracle_cache: getExchangeRateOracleHealth().status, }, }); }); @@ -288,6 +290,30 @@ export async function createApp({ redisClient }) { res.status(health.status === "unhealthy" ? 503 : 200).json(health); }); + /** + * @swagger + * /health/exchange-rate-oracle-cache: + * get: + * summary: Exchange Rate Oracle Cache health telemetry + * description: > + * Rolling-window load and lookup counts plus derived status for the + * exchange-rate quote cache (issue #1443). Returns 503 only when the + * cache is unhealthy (internal error ratio above threshold). A + * degraded status (timeouts or stale lookups) still returns 200. + * The body contains no quotes, asset codes, or account ids. + * tags: [Health] + * security: [] + * responses: + * 200: + * description: Cache healthy or degraded + * 503: + * description: Cache unhealthy + */ + app.get("/health/exchange-rate-oracle-cache", (_req, res) => { + const health = getExchangeRateOracleHealth(); + res.status(health.status === "unhealthy" ? 503 : 200).json(health); + }); + const verifyPaymentRateLimit = createVerifyPaymentRateLimit({ store: redisAvailable ? createRedisRateLimitStore({ client: redisClient }) : undefined, }); diff --git a/backend/src/lib/exchange-rate-cache.js b/backend/src/lib/exchange-rate-cache.js index 2ef9b781..8013582a 100644 --- a/backend/src/lib/exchange-rate-cache.js +++ b/backend/src/lib/exchange-rate-cache.js @@ -36,6 +36,11 @@ import { exchangeRateCacheLoadTimeouts, exchangeRateCacheStaleWritesPrevented, } from './path-payment-metrics.js'; +import { + classifyOracleLoadError, + recordOracleLoad, + recordOracleLookup, +} from './exchange-rate-oracle-telemetry.js'; /** * Default cache metrics wired to the granular path-payment series (issue #1048) @@ -142,6 +147,7 @@ export class ExchangeRateCache { const entry = this.cache.get(key); if (!entry) { this.metrics?.miss?.inc?.({ cache: 'exchange_rate' }); + recordOracleLookup('miss'); return { hit: false, data: null, stale: false }; } @@ -150,11 +156,13 @@ export class ExchangeRateCache { if (age > this.staleToleranceMs) { this.cache.delete(key); this.metrics?.miss?.inc?.({ cache: 'exchange_rate' }); + recordOracleLookup('miss'); return { hit: false, data: null, stale: false }; } const stale = age > this.ttlMs; this.metrics?.hit?.inc?.({ cache: 'exchange_rate', stale: stale ? '1' : '0' }); + recordOracleLookup(stale ? 'stale' : 'hit'); // Refresh recency in LRU order this.cache.delete(key); @@ -204,6 +212,7 @@ export class ExchangeRateCache { const entry = { promise: null, invalidated: false }; entry.promise = (async () => { + const started = Date.now(); try { // Invoke the loader on a later microtask so the entry is registered // in `inflight` first — even a synchronously throwing loader then @@ -217,7 +226,11 @@ export class ExchangeRateCache { } else { this.set(key, data); } + recordOracleLoad('success', Date.now() - started); return data; + } catch (err) { + recordOracleLoad(classifyOracleLoadError(err), Date.now() - started, err); + throw err; } finally { // Only remove our own entry; an invalidation may already have // detached it and a newer load may occupy the slot. diff --git a/backend/src/lib/exchange-rate-oracle-telemetry.js b/backend/src/lib/exchange-rate-oracle-telemetry.js new file mode 100644 index 00000000..2fae7b5e --- /dev/null +++ b/backend/src/lib/exchange-rate-oracle-telemetry.js @@ -0,0 +1,317 @@ +/** + * Exchange Rate Oracle Cache telemetry (issue #1443). + * + * Prometheus series and a rolling-window health snapshot for the in-process + * quote cache (exchange-rate-cache.js). Label values are a fixed server-side + * set — never asset codes, issuers, or amounts — so quote traffic cannot + * inflate cardinality. + * + * The metrics live in their own registry so they can be tested in isolation. + * /metrics merges this registry with the others. + */ + +import client from 'prom-client'; +import { logger } from './logger.js'; + +const register = new client.Registry(); + +register.setDefaultLabels({ + app: 'stellar-payment-api', +}); + +/** result: hit | miss | stale */ +export const exchangeRateOracleLookupsTotal = new client.Counter({ + name: 'exchange_rate_oracle_cache_lookups_total', + help: 'Exchange-rate oracle cache lookups, by result', + labelNames: ['result'], +}); + +/** + * outcome: success | error | timeout | not_found + * not_found is a normal "no path" result and is not an internal error. + */ +export const exchangeRateOracleLoadsTotal = new client.Counter({ + name: 'exchange_rate_oracle_cache_loads_total', + help: 'Exchange-rate oracle cache loads that called the upstream quote source, by outcome', + labelNames: ['outcome'], +}); + +export const exchangeRateOracleLoadDuration = new client.Histogram({ + name: 'exchange_rate_oracle_cache_load_duration_seconds', + help: 'Wall time of an exchange-rate oracle cache load, including upstream time', + labelNames: ['outcome'], + buckets: [0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 15], +}); + +/** + * Rolling-window gauges. Refreshed at scrape time so they decay when idle. + * health_state: 0 = healthy, 1 = degraded, 2 = unhealthy + */ +export const exchangeRateOracleHealthState = new client.Gauge({ + name: 'exchange_rate_oracle_cache_health_state', + help: 'Exchange-rate oracle cache health (0 = healthy, 1 = degraded, 2 = unhealthy)', +}); + +export const exchangeRateOracleErrorRatio = new client.Gauge({ + name: 'exchange_rate_oracle_cache_error_ratio', + help: 'Share of oracle cache loads that failed with an internal error over the rolling window', +}); + +export const exchangeRateOracleTimeoutRatio = new client.Gauge({ + name: 'exchange_rate_oracle_cache_timeout_ratio', + help: 'Share of oracle cache loads that timed out over the rolling window', +}); + +export const exchangeRateOracleStaleRatio = new client.Gauge({ + name: 'exchange_rate_oracle_cache_stale_ratio', + help: 'Share of oracle cache lookups that found a stale-but-tolerable entry over the rolling window', +}); + +export const exchangeRateOracleLastLoadTimestamp = new client.Gauge({ + name: 'exchange_rate_oracle_cache_last_load_timestamp_seconds', + help: 'Unix time of the most recent exchange-rate oracle cache load', +}); + +const LOOKUP_RESULTS = new Set(['hit', 'miss', 'stale']); +const LOAD_OUTCOMES = new Set(['success', 'error', 'timeout', 'not_found']); + +export const HEALTH_STATE_VALUES = Object.freeze({ healthy: 0, degraded: 1, unhealthy: 2 }); + +function readNumberEnv(name, fallback) { + const raw = Number.parseFloat(process.env[name] ?? ''); + return Number.isFinite(raw) && raw >= 0 ? raw : fallback; +} + +export const DEFAULT_HEALTH_OPTIONS = Object.freeze({ + windowMs: readNumberEnv('EXCHANGE_RATE_ORACLE_HEALTH_WINDOW_MS', 300_000), + bucketMs: 5_000, + minSamples: readNumberEnv('EXCHANGE_RATE_ORACLE_HEALTH_MIN_SAMPLES', 20), + errorRatioThreshold: readNumberEnv('EXCHANGE_RATE_ORACLE_ERROR_RATIO_THRESHOLD', 0.05), + timeoutRatioThreshold: readNumberEnv('EXCHANGE_RATE_ORACLE_TIMEOUT_RATIO_THRESHOLD', 0.2), + staleRatioThreshold: readNumberEnv('EXCHANGE_RATE_ORACLE_STALE_RATIO_THRESHOLD', 0.5), +}); + +function emptyBucket() { + return { + slot: -1, + hits: 0, + misses: 0, + stale: 0, + success: 0, + errors: 0, + timeouts: 0, + notFound: 0, + }; +} + +/** + * Rolling-window health tracker. Counts live in a fixed ring of time buckets + * so memory stays constant regardless of traffic. + */ +export class ExchangeRateOracleHealthMonitor { + constructor(options = {}, now = () => Date.now()) { + this.options = { ...DEFAULT_HEALTH_OPTIONS, ...options }; + this.now = now; + this.bucketCount = Math.max(1, Math.ceil(this.options.windowMs / this.options.bucketMs)); + this.reset(); + } + + reset() { + this.buckets = Array.from({ length: this.bucketCount }, () => emptyBucket()); + this.lastLoadAt = null; + } + + _bucketFor(time) { + const slot = Math.floor(time / this.options.bucketMs); + const bucket = this.buckets[slot % this.bucketCount]; + if (bucket.slot !== slot) { + Object.assign(bucket, emptyBucket(), { slot }); + } + return bucket; + } + + /** @param {'hit'|'miss'|'stale'} result */ + recordLookup(result) { + const bucket = this._bucketFor(this.now()); + if (result === 'hit') bucket.hits += 1; + else if (result === 'stale') bucket.stale += 1; + else bucket.misses += 1; + } + + /** @param {'success'|'error'|'timeout'|'not_found'} outcome */ + recordLoad(outcome) { + const time = this.now(); + const bucket = this._bucketFor(time); + if (outcome === 'success') bucket.success += 1; + else if (outcome === 'timeout') bucket.timeouts += 1; + else if (outcome === 'not_found') bucket.notFound += 1; + else bucket.errors += 1; + this.lastLoadAt = time; + } + + snapshot() { + const time = this.now(); + const currentSlot = Math.floor(time / this.options.bucketMs); + const oldestSlot = currentSlot - this.bucketCount + 1; + const counts = { + hits: 0, + misses: 0, + stale: 0, + success: 0, + errors: 0, + timeouts: 0, + notFound: 0, + }; + + for (const bucket of this.buckets) { + if (bucket.slot < oldestSlot || bucket.slot > currentSlot) continue; + counts.hits += bucket.hits; + counts.misses += bucket.misses; + counts.stale += bucket.stale; + counts.success += bucket.success; + counts.errors += bucket.errors; + counts.timeouts += bucket.timeouts; + counts.notFound += bucket.notFound; + } + + const loads = counts.success + counts.errors + counts.timeouts + counts.notFound; + const lookups = counts.hits + counts.misses + counts.stale; + const errorRatio = loads > 0 ? counts.errors / loads : 0; + const timeoutRatio = loads > 0 ? counts.timeouts / loads : 0; + const staleRatio = lookups > 0 ? counts.stale / lookups : 0; + const { minSamples, errorRatioThreshold, timeoutRatioThreshold, staleRatioThreshold } = + this.options; + + const reasons = []; + let status = 'healthy'; + if (loads >= minSamples && errorRatio >= errorRatioThreshold) { + status = 'unhealthy'; + reasons.push('error_ratio_exceeded'); + } + if (loads >= minSamples && timeoutRatio >= timeoutRatioThreshold) { + if (status === 'healthy') status = 'degraded'; + reasons.push('timeout_ratio_exceeded'); + } + if (lookups >= minSamples && staleRatio >= staleRatioThreshold) { + if (status === 'healthy') status = 'degraded'; + reasons.push('stale_ratio_exceeded'); + } + + return { + status, + reasons, + window_ms: this.bucketCount * this.options.bucketMs, + loads, + success: counts.success, + errors: counts.errors, + timeouts: counts.timeouts, + not_found: counts.notFound, + lookups, + hits: counts.hits, + misses: counts.misses, + stale: counts.stale, + error_ratio: Number(errorRatio.toFixed(4)), + timeout_ratio: Number(timeoutRatio.toFixed(4)), + stale_ratio: Number(staleRatio.toFixed(4)), + last_load_at: this.lastLoadAt === null ? null : new Date(this.lastLoadAt).toISOString(), + thresholds: { + min_samples: minSamples, + error_ratio: errorRatioThreshold, + timeout_ratio: timeoutRatioThreshold, + stale_ratio: staleRatioThreshold, + }, + }; + } +} + +const healthMonitor = new ExchangeRateOracleHealthMonitor(); + +exchangeRateOracleHealthState.collect = function collectOracleHealth() { + const snap = healthMonitor.snapshot(); + this.set(HEALTH_STATE_VALUES[snap.status] ?? 0); + exchangeRateOracleErrorRatio.set(snap.error_ratio); + exchangeRateOracleTimeoutRatio.set(snap.timeout_ratio); + exchangeRateOracleStaleRatio.set(snap.stale_ratio); + if (healthMonitor.lastLoadAt !== null) { + exchangeRateOracleLastLoadTimestamp.set(healthMonitor.lastLoadAt / 1000); + } +}; + +/** Current cache health, for /health endpoints. No quote or account data. */ +export function getExchangeRateOracleHealth() { + return healthMonitor.snapshot(); +} + +/** Test-only: clear the rolling health window. Counters are left intact. */ +export function resetExchangeRateOracleHealth() { + healthMonitor.reset(); +} + +function safeRecord(fn) { + try { + fn(); + } catch (err) { + logger.error({ err: err?.message }, 'Failed to record exchange-rate oracle telemetry'); + } +} + +/** + * Classify a failed load without reading caller-controlled strings into labels. + * @param {unknown} err + * @returns {'timeout'|'not_found'|'error'} + */ +export function classifyOracleLoadError(err) { + if (!err || typeof err !== 'object') return 'error'; + if (err.name === 'CacheLoadTimeoutError') return 'timeout'; + const status = err.status ?? err.statusCode ?? err.response?.status ?? null; + if (status === 504 || status === 408) return 'timeout'; + if (status === 404 || err.name === 'NoPathFoundError') return 'not_found'; + return 'error'; +} + +/** @param {'hit'|'miss'|'stale'} result */ +export function recordOracleLookup(result) { + const safe = LOOKUP_RESULTS.has(result) ? result : 'miss'; + safeRecord(() => { + exchangeRateOracleLookupsTotal.inc({ result: safe }); + healthMonitor.recordLookup(safe); + }); +} + +/** + * @param {'success'|'error'|'timeout'|'not_found'} outcome + * @param {number} durationMs + * @param {unknown} [err] + */ +export function recordOracleLoad(outcome, durationMs, err) { + const safeOutcome = LOAD_OUTCOMES.has(outcome) ? outcome : 'error'; + const ms = Number.isFinite(durationMs) && durationMs > 0 ? durationMs : 0; + safeRecord(() => { + exchangeRateOracleLoadsTotal.inc({ outcome: safeOutcome }); + exchangeRateOracleLoadDuration.observe({ outcome: safeOutcome }, ms / 1000); + healthMonitor.recordLoad(safeOutcome); + }); + if (safeOutcome === 'error' || safeOutcome === 'timeout') { + logger.warn( + { + outcome: safeOutcome, + durationMs: ms, + err: err && typeof err === 'object' ? err.message : undefined, + code: err && typeof err === 'object' ? err.code : undefined, + status: err && typeof err === 'object' ? (err.status ?? err.statusCode) : undefined, + }, + 'Exchange rate oracle cache load failed', + ); + } +} + +register.registerMetric(exchangeRateOracleLookupsTotal); +register.registerMetric(exchangeRateOracleLoadsTotal); +register.registerMetric(exchangeRateOracleLoadDuration); +register.registerMetric(exchangeRateOracleHealthState); +register.registerMetric(exchangeRateOracleErrorRatio); +register.registerMetric(exchangeRateOracleTimeoutRatio); +register.registerMetric(exchangeRateOracleStaleRatio); +register.registerMetric(exchangeRateOracleLastLoadTimestamp); + +export { register as exchangeRateOracleRegister }; diff --git a/backend/src/lib/exchange-rate-oracle-telemetry.test.js b/backend/src/lib/exchange-rate-oracle-telemetry.test.js new file mode 100644 index 00000000..8534b55d --- /dev/null +++ b/backend/src/lib/exchange-rate-oracle-telemetry.test.js @@ -0,0 +1,212 @@ +import { readFileSync } from 'node:fs'; +import { describe, it, expect, beforeEach, vi } from 'vitest'; + +vi.mock('./logger.js', () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +import { logger } from './logger.js'; +import { ExchangeRateCache } from './exchange-rate-cache.js'; +import { + ExchangeRateOracleHealthMonitor, + classifyOracleLoadError, + exchangeRateOracleLoadsTotal, + exchangeRateOracleRegister, + getExchangeRateOracleHealth, + resetExchangeRateOracleHealth, +} from './exchange-rate-oracle-telemetry.js'; + +function make(overrides = {}) { + let clock = 1_700_000_000_000; + const monitor = new ExchangeRateOracleHealthMonitor( + { windowMs: 60_000, bucketMs: 5_000, minSamples: 10, ...overrides }, + () => clock, + ); + return { + monitor, + advance(ms) { + clock += ms; + }, + }; +} + +beforeEach(() => { + resetExchangeRateOracleHealth(); + vi.clearAllMocks(); +}); + +describe('classifyOracleLoadError', () => { + it('treats load timeouts and gateway timeouts as timeout', () => { + expect(classifyOracleLoadError({ name: 'CacheLoadTimeoutError', status: 504 })).toBe('timeout'); + expect(classifyOracleLoadError({ status: 504 })).toBe('timeout'); + expect(classifyOracleLoadError({ statusCode: 408 })).toBe('timeout'); + }); + + it('treats a missing path as not_found, not an internal error', () => { + expect(classifyOracleLoadError({ name: 'NoPathFoundError', statusCode: 404 })).toBe('not_found'); + expect(classifyOracleLoadError({ status: 404 })).toBe('not_found'); + }); + + it('classifies everything else, including non-objects, as error', () => { + expect(classifyOracleLoadError({ status: 502, message: 'down' })).toBe('error'); + expect(classifyOracleLoadError(null)).toBe('error'); + expect(classifyOracleLoadError('nope')).toBe('error'); + }); +}); + +describe('ExchangeRateOracleHealthMonitor', () => { + it('stays healthy with no traffic and below the sample floor', () => { + const { monitor } = make(); + expect(monitor.snapshot().status).toBe('healthy'); + monitor.recordLoad('error'); + expect(monitor.snapshot()).toMatchObject({ status: 'healthy', loads: 1, errors: 1 }); + }); + + it('goes unhealthy only on internal errors, not on not_found', () => { + const failing = make().monitor; + for (let i = 0; i < 10; i += 1) failing.recordLoad('error'); + expect(failing.snapshot()).toMatchObject({ + status: 'unhealthy', + reasons: ['error_ratio_exceeded'], + }); + + const empty = make().monitor; + for (let i = 0; i < 10; i += 1) empty.recordLoad('not_found'); + expect(empty.snapshot()).toMatchObject({ status: 'healthy', not_found: 10, error_ratio: 0 }); + }); + + it('goes degraded on timeouts or stale lookups and keeps every reason', () => { + const timeouts = make().monitor; + for (let i = 0; i < 8; i += 1) timeouts.recordLoad('success'); + for (let i = 0; i < 2; i += 1) timeouts.recordLoad('timeout'); + expect(timeouts.snapshot()).toMatchObject({ + status: 'degraded', + reasons: ['timeout_ratio_exceeded'], + }); + + const stale = make().monitor; + for (let i = 0; i < 5; i += 1) stale.recordLookup('hit'); + for (let i = 0; i < 5; i += 1) stale.recordLookup('stale'); + expect(stale.snapshot()).toMatchObject({ + status: 'degraded', + reasons: ['stale_ratio_exceeded'], + }); + }); + + it('lets unhealthy take precedence while still reporting degraded reasons', () => { + const { monitor } = make(); + for (let i = 0; i < 10; i += 1) monitor.recordLoad('error'); + for (let i = 0; i < 10; i += 1) monitor.recordLookup('stale'); + const snap = monitor.snapshot(); + expect(snap.status).toBe('unhealthy'); + expect(snap.reasons).toEqual(['error_ratio_exceeded', 'stale_ratio_exceeded']); + }); + + it('forgets events that fall out of the window', () => { + const { monitor, advance } = make(); + for (let i = 0; i < 10; i += 1) monitor.recordLoad('error'); + expect(monitor.snapshot().status).toBe('unhealthy'); + advance(61_000); + const snap = monitor.snapshot(); + expect(snap).toMatchObject({ loads: 0, status: 'healthy' }); + expect(snap.last_load_at).not.toBeNull(); + }); + + it('reuses a bucket slot without leaking the previous slot counts', () => { + let clock = 0; + const monitor = new ExchangeRateOracleHealthMonitor( + { windowMs: 5_000, bucketMs: 5_000, minSamples: 1 }, + () => clock, + ); + monitor.recordLoad('error'); + clock += 5_000; + monitor.recordLoad('success'); + expect(monitor.snapshot()).toMatchObject({ loads: 1, success: 1, errors: 0 }); + }); + + it('keeps a fixed bucket count under heavy load', () => { + let clock = 0; + const monitor = new ExchangeRateOracleHealthMonitor( + { windowMs: 60_000, bucketMs: 5_000, minSamples: 20 }, + () => clock, + ); + for (let i = 0; i < 20_000; i += 1) { + clock += 7; + monitor.recordLoad(i % 2 ? 'success' : 'error'); + } + expect(monitor.buckets).toHaveLength(12); + expect(monitor.snapshot().loads).toBeLessThanOrEqual(Math.ceil(60_000 / 7) + 1); + }); +}); + +describe('cache instrumentation', () => { + it('records fresh hits, misses, stale lookups, successes, not-found and timeouts', async () => { + vi.useFakeTimers(); + let snap; + try { + const cache = new ExchangeRateCache({ ttlMs: 100, staleToleranceMs: 500, maxEntries: 5 }); + cache.set('fresh', { rate: 1 }); + expect(cache.get('fresh').stale).toBe(false); + + vi.advanceTimersByTime(150); + expect(cache.get('fresh').stale).toBe(true); + expect(cache.get('missing').hit).toBe(false); + + await cache.getOrLoad('loaded', async () => ({ rate: 2 })); + const notFound = Object.assign(new Error('no path'), { name: 'NoPathFoundError', statusCode: 404 }); + await expect(cache.getOrLoad('none', async () => { throw notFound; })).rejects.toThrow('no path'); + + const pending = cache.getOrLoad('hung', () => new Promise(() => {}), { timeoutMs: 30 }); + const assertion = expect(pending).rejects.toMatchObject({ name: 'CacheLoadTimeoutError' }); + await vi.advanceTimersByTimeAsync(30); + await assertion; + snap = getExchangeRateOracleHealth(); + } finally { + vi.useRealTimers(); + } + expect(snap.hits).toBeGreaterThanOrEqual(1); + expect(snap.stale).toBeGreaterThanOrEqual(1); + expect(snap.misses).toBeGreaterThanOrEqual(1); + expect(snap.success).toBe(1); + expect(snap.not_found).toBe(1); + expect(snap.timeouts).toBe(1); + expect(snap.errors).toBe(0); + expect(logger.warn).toHaveBeenCalled(); + }); + + it('does not let a telemetry failure replace the loader error', async () => { + const cache = new ExchangeRateCache({ ttlMs: 1000 }); + const original = exchangeRateOracleLoadsTotal.inc; + exchangeRateOracleLoadsTotal.inc = () => { + throw new Error('metrics down'); + }; + try { + await expect(cache.getOrLoad('k', async () => { + throw Object.assign(new Error('horizon down'), { status: 503 }); + })).rejects.toThrow('horizon down'); + } finally { + exchangeRateOracleLoadsTotal.inc = original; + } + }); +}); + +describe('alert rules (issue #1443)', () => { + it('only reference metrics that the oracle cache registry exposes', () => { + const rules = readFileSync( + new URL('../../docs/alerts/exchange-rate-oracle-cache.rules.yml', import.meta.url), + 'utf8', + ); + const referenced = new Set( + [...rules.matchAll(/\b(exchange_rate_oracle_cache_[a-z0-9_]+)/g)].map(([, name]) => + name.replace(/_(bucket|sum|count)$/, ''), + ), + ); + const registered = new Set( + exchangeRateOracleRegister.getMetricsAsArray().map((metric) => metric.name), + ); + expect(referenced.size).toBeGreaterThan(0); + for (const name of referenced) { + expect(registered, `alert rule references unknown metric ${name}`).toContain(name); + } + }); +}); diff --git a/backend/src/routes/payments.js b/backend/src/routes/payments.js index a8a0c17e..266efa64 100644 --- a/backend/src/routes/payments.js +++ b/backend/src/routes/payments.js @@ -297,32 +297,6 @@ function createPaymentsRouter({ logger.error({ err, merchantId: req.merchant?.id }, "DEBUG: createSession error"); if (err.status === 400 && err.details) { return res.status(400).json({ error: err.message, ...err.details }); - const supabase = await getSupabaseClient(); - logger.info({ merchantId: req.merchant?.id, amount: req.body?.amount, asset: req.body?.asset }, "DEBUG: createSession started"); - - // Sanitization, strict payload checks and shared business rules - // (issues #1087, #1447) with validator metrics/health (#1448). - const validation = validatePaymentSession({ - body: req.body, - merchant: req.merchant, - source: "http", - }); - if (!validation.ok) { - const { rejection } = validation; - const assetLabel = req.body?.asset; - paymentFailedCounter.inc({ - asset: assetLabel, - reason: rejection.reason === "issuer_not_allowed" ? "invalid_issuer" : rejection.reason, - }); - paymentProcessorSessionsTotal.inc({ asset: assetLabel, outcome: "validation_failed" }); - paymentProcessorSessionDuration.observe( - { asset: assetLabel, outcome: "validation_failed" }, - (Date.now() - sessionStart) / 1000, - ); - return res.status(400).json({ - error: rejection.message, - ...(rejection.rule === "limits" ? rejection.details : {}), - }); } next(err); } @@ -331,63 +305,35 @@ function createPaymentsRouter({ async function createSessionUnlocked(req, res) { const sessionStart = Date.now(); const supabase = await getSupabaseClient(); - const body = req.body; - const asset = body.asset?.toUpperCase(); - logger.info({ merchantId: req.merchant?.id, amount: body.amount, asset: body.asset }, "DEBUG: createSession started"); - - // Shared business-rule validation (issue #1087) — issuer presence/format. - const { assetIssuer, rejection: issuerRejection } = resolveAndValidateIssuer( - asset, - body.asset_issuer, - ); - if (issuerRejection) { - paymentFailedCounter.inc({ asset: body.asset, reason: issuerRejection.reason }); - paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "validation_failed" }); - paymentProcessorSessionDuration.observe( - { asset: body.asset, outcome: "validation_failed" }, - (Date.now() - sessionStart) / 1000, - ); - return res.status(400).json({ error: issuerRejection.message }); - } - const body = validation.payload; - const { asset, assetIssuer } = validation; - - // Shared business-rule validation (issue #1087) — per-asset limits (#153). - const limitRejection = validatePerAssetLimits({ - rawAsset: body.asset, - amount: body.amount, - paymentLimits: req.merchant.payment_limits, + logger.info({ merchantId: req.merchant?.id, amount: req.body?.amount, asset: req.body?.asset }, "DEBUG: createSession started"); + + // Sanitization and shared business rules (issues #1087, #1447). + // The validator owns issuer, limit and allowlist checks; this route + // persists validation.payload rather than the raw body. + const validation = validatePaymentSession({ + body: req.body, + merchant: req.merchant, + source: "http", }); - if (limitRejection) { - paymentFailedCounter.inc({ asset: body.asset, reason: limitRejection.reason }); - paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "validation_failed" }); + if (!validation.ok) { + const { rejection } = validation; + const assetLabel = req.body?.asset; + paymentFailedCounter.inc({ + asset: assetLabel, + reason: rejection.reason === "issuer_not_allowed" ? "invalid_issuer" : rejection.reason, + }); + paymentProcessorSessionsTotal.inc({ asset: assetLabel, outcome: "validation_failed" }); paymentProcessorSessionDuration.observe( - { asset: body.asset, outcome: "validation_failed" }, + { asset: assetLabel, outcome: "validation_failed" }, (Date.now() - sessionStart) / 1000, ); return res.status(400).json({ - error: limitRejection.message, - ...limitRejection.details, + error: rejection.message, + ...(rejection.rule === "limits" ? rejection.details : {}), }); } - - // Shared business-rule validation (issue #1087) — allowed-issuers check: - // if the merchant has configured a non-empty allowlist, only those - // issuer addresses may be used. - const allowedIssuerRejection = validateAllowedIssuers({ - asset, - assetIssuer, - allowedIssuers: req.merchant.allowed_issuers, - }); - if (allowedIssuerRejection) { - paymentFailedCounter.inc({ asset: body.asset, reason: "invalid_issuer" }); - paymentProcessorSessionsTotal.inc({ asset: body.asset, outcome: "validation_failed" }); - paymentProcessorSessionDuration.observe( - { asset: body.asset, outcome: "validation_failed" }, - (Date.now() - sessionStart) / 1000, - ); - return res.status(400).json({ error: allowedIssuerRejection.message }); - } + const body = validation.payload; + const { asset, assetIssuer } = validation; const isSandbox = body.sandbox === true; const baseId = randomUUID(); diff --git a/backend/src/routes/prometheus.js b/backend/src/routes/prometheus.js index bd812212..4a9cdde3 100644 --- a/backend/src/routes/prometheus.js +++ b/backend/src/routes/prometheus.js @@ -12,6 +12,9 @@ import { pathPaymentRegister } from "../lib/path-payment-metrics.js"; // Payment Session Validator metrics live in their own registry (issue #1448) // and are merged into the scrape output below. import { paymentSessionValidatorRegister } from "../lib/payment-session-validator-metrics.js"; +// Exchange Rate Oracle Cache alert metrics live in their own registry (issue #1443) +// and are merged into the scrape output below. +import { exchangeRateOracleRegister } from "../lib/exchange-rate-oracle-telemetry.js"; const router = express.Router(); @@ -20,7 +23,7 @@ const router = express.Router(); * /metrics: * get: * summary: Expose Prometheus metrics - * description: Returns the current state of Prometheus metrics for the application, including granular payment processor, trustline manager, path payment and payment session validator metrics. + * description: Returns the current state of Prometheus metrics for the application, including granular payment processor, trustline manager, path payment, payment session validator and exchange-rate oracle cache metrics. * tags: [Monitoring] * responses: * 200: @@ -38,6 +41,7 @@ router.get("/metrics", async (req, res) => { trustlineManagerRegister.metrics(), pathPaymentRegister.metrics(), paymentSessionValidatorRegister.metrics(), + exchangeRateOracleRegister.metrics(), ]); res.set("Content-Type", register.contentType); res.end(scrapes.join("\n")); diff --git a/backend/tests/helpers/fake-redis.js b/backend/tests/helpers/fake-redis.js index 74a9b7ff..6c682755 100644 --- a/backend/tests/helpers/fake-redis.js +++ b/backend/tests/helpers/fake-redis.js @@ -1,182 +1,121 @@ /** - * Minimal in-memory Redis stand-in for concurrency tests. + * In-memory Redis stand-in shared by the exchange-rate coordinator and the + * payment-session lock / idempotency tests. * - * Supports the subset used by the idempotency middleware and the payment - * session lock: GET / SET (with EX/PX) / DEL, plus sendCommand for - * `SET key value NX PX ttl` and the compare-and-delete EVAL script. - * - * Each command awaits a FIFO latency tick and then executes its - * check-and-mutate step synchronously, which matches Redis' single-threaded - * atomicity while still letting concurrent callers interleave between - * commands — exactly the window real race conditions live in. - */ -export function createFakeRedis({ latencyMs = 0 } = {}) { - const store = new Map(); - const commandLog = []; - - const now = () => Date.now(); - const live = (key) => { - const entry = store.get(key); - if (!entry) return null; - if (entry.expiresAt !== null && entry.expiresAt <= now()) { - * In-memory stand-in for the subset of Redis used by the exchange-rate - * coordinator (issues #1445, #1446): SET [PX ms] [NX], GET, DEL and the - * compare-and-delete EVAL script. Several coordinators can share one - * instance to simulate multiple API processes against the same Redis. - * - * Commands resolve asynchronously (optionally after `latencyMs`) so they - * interleave the way network round-trips do. Expiry follows Date.now(), so - * it works with vi.useFakeTimers(). + * Supports GET, SET (EX/PX/NX), DEL and the compare-and-delete EVAL script, + * both as methods and through sendCommand. Commands run in FIFO order. + * Expiry follows Date.now(), so it works with vi.useFakeTimers(). + * setFailure() makes later commands reject, which the coordinator uses to + * exercise its fail-open path. */ export function createFakeRedis({ latencyMs = 0 } = {}) { const store = new Map(); // key -> { value, expiresAt } const calls = []; let failWith = null; + let queue = Promise.resolve(); const live = (key) => { const entry = store.get(key); if (!entry) return null; - if (entry.expiresAt !== null && Date.now() >= entry.expiresAt) { + if (entry.expiresAt !== null && entry.expiresAt <= Date.now()) { store.delete(key); return null; } return entry; }; - // Commands complete in the order they were issued (FIFO), like a single - // Redis connection: the idempotency middleware and the session lock share - // one client, and the lock's correctness relies on that ordering. - let queue = Promise.resolve(); - // With no latency, stay on the microtask queue so tests using - // vi.useFakeTimers() are not blocked on a timer that never fires. + const tick = () => { if (latencyMs <= 0) return Promise.resolve(); queue = queue.then( - () => new Promise((resolve) => setTimeout(resolve, Math.random() * latencyMs)), + () => new Promise((resolve) => setTimeout(resolve, latencyMs)), ); return queue; }; function setEntry(key, value, { ex, px } = {}) { let expiresAt = null; - if (ex) expiresAt = now() + Number(ex) * 1000; - if (px) expiresAt = now() + Number(px); + if (ex) expiresAt = Date.now() + Number(ex) * 1000; + if (px) expiresAt = Date.now() + Number(px); store.set(key, { value: String(value), expiresAt }); } - const client = { - isOpen: true, - store, - commandLog, - async get(key) { - await tick(); - commandLog.push(["GET", key]); - return live(key)?.value ?? null; - }, - async set(key, value, opts = {}) { - await tick(); - commandLog.push(["SET", key]); - if (opts.NX && live(key)) return null; - setEntry(key, value, { ex: opts.EX, px: opts.PX }); - return "OK"; - }, - async del(key) { - await tick(); - commandLog.push(["DEL", key]); - return store.delete(key) ? 1 : 0; - }, - async sendCommand(args) { - await tick(); - const [cmd, ...rest] = args; - commandLog.push([String(cmd).toUpperCase(), rest[0]]); - switch (String(cmd).toUpperCase()) { - case "SET": { - const [key, value, ...flags] = rest; - const upper = flags.map((f) => String(f).toUpperCase()); - if (upper.includes("NX") && live(key)) return null; - const pxIdx = upper.indexOf("PX"); - const exIdx = upper.indexOf("EX"); - setEntry(key, value, { - px: pxIdx >= 0 ? flags[pxIdx + 1] : undefined, - ex: exIdx >= 0 ? flags[exIdx + 1] : undefined, - }); - return "OK"; - } - case "EVAL": { - const [script, numKeys, key, token] = rest; - if (!/redis\.call\("get", KEYS\[1\]\) == ARGV\[1\]/.test(script) || numKeys !== "1") { - throw new Error("fake-redis: unsupported EVAL script"); - } - const entry = live(key); - if (entry && entry.value === token) { - store.delete(key); - return 1; - } - return 0; - } - case "GET": - return live(rest[0])?.value ?? null; - default: - throw new Error(`fake-redis: unsupported command ${cmd}`); - } - }, - keys(prefix = "") { - return [...store.keys()].filter((k) => k.startsWith(prefix) && live(k)); - }, - }; - - return client; + function isCompareAndDelete(script) { + const text = String(script); + return /redis\.call\(['"]get['"], KEYS\[1\]\)/.test(text) + && /redis\.call\(['"]del['"], KEYS\[1\]\)/.test(text); + } - const execute = (args) => { - const [cmd, ...rest] = args; - switch (String(cmd).toUpperCase()) { - case 'GET': - return live(rest[0])?.value ?? null; - case 'SET': { - const [key, value, ...opts] = rest; - const upper = opts.map((o) => String(o).toUpperCase()); - const pxIndex = upper.indexOf('PX'); - const ttl = pxIndex >= 0 ? Number(opts[pxIndex + 1]) : null; - if (upper.includes('NX') && live(key)) return null; - store.set(key, { value: String(value), expiresAt: ttl ? Date.now() + ttl : null }); - return 'OK'; - } - case 'DEL': - return rest.reduce((n, key) => n + (live(key) && store.delete(key) ? 1 : 0), 0); - case 'EVAL': { - // Only the compare-and-delete release script is supported. - const [, , key, token] = rest; - if (live(key)?.value === token) { - store.delete(key); - return 1; - } - return 0; - } - default: - throw new Error(`fake-redis: unsupported command ${cmd}`); - } - }; + async function run(args, fn) { + calls.push(args); + await tick(); + if (failWith) throw failWith; + return fn(); + } return { isOpen: true, - calls, store, - /** Make every subsequent command reject with `error` (null to heal). */ + calls, + // Payment-session tests read [command, key] pairs. Full sendCommand + // argument lists still destructure that way. + commandLog: calls, setFailure(error) { failWith = error; }, - async sendCommand(args) { - calls.push(args); - if (latencyMs > 0) { - await new Promise((resolve) => setTimeout(resolve, latencyMs)); - } else { - await Promise.resolve(); - } - if (failWith) throw failWith; - return execute(args); - }, - /** Test helper: count calls by command name. */ count(cmd) { - return calls.filter(([c]) => String(c).toUpperCase() === cmd).length; + return calls.filter(([name]) => String(name).toUpperCase() === String(cmd).toUpperCase()).length; + }, + keys(prefix = '') { + return [...store.keys()].filter((key) => key.startsWith(prefix) && live(key)); + }, + get(key) { + return run(['GET', key], () => live(key)?.value ?? null); + }, + set(key, value, opts = {}) { + return run(['SET', key, value], () => { + if (opts.NX && live(key)) return null; + setEntry(key, value, { ex: opts.EX, px: opts.PX }); + return 'OK'; + }); + }, + del(key) { + return run(['DEL', key], () => (live(key) && store.delete(key) ? 1 : 0)); + }, + sendCommand(args) { + const [cmd, ...rest] = args; + return run(args, () => { + switch (String(cmd).toUpperCase()) { + case 'GET': + return live(rest[0])?.value ?? null; + case 'SET': { + const [key, value, ...flags] = rest; + const upper = flags.map((flag) => String(flag).toUpperCase()); + if (upper.includes('NX') && live(key)) return null; + const pxIdx = upper.indexOf('PX'); + const exIdx = upper.indexOf('EX'); + setEntry(key, value, { + px: pxIdx >= 0 ? flags[pxIdx + 1] : undefined, + ex: exIdx >= 0 ? flags[exIdx + 1] : undefined, + }); + return 'OK'; + } + case 'DEL': + return rest.reduce((count, key) => count + (live(key) && store.delete(key) ? 1 : 0), 0); + case 'EVAL': { + const [script, numKeys, key, token] = rest; + if (!isCompareAndDelete(script) || String(numKeys) !== '1') { + throw new Error('fake-redis: unsupported EVAL script'); + } + if (live(key)?.value === token) { + store.delete(key); + return 1; + } + return 0; + } + default: + throw new Error(`fake-redis: unsupported command ${cmd}`); + } + }); }, }; } diff --git a/backend/tests/integration/exchange-rate-cache.test.js b/backend/tests/integration/exchange-rate-cache.test.js index f8ded940..d11a21b9 100644 --- a/backend/tests/integration/exchange-rate-cache.test.js +++ b/backend/tests/integration/exchange-rate-cache.test.js @@ -19,6 +19,7 @@ import { createApp } from '../../src/app.js'; import { closePool } from '../../src/lib/db.js'; import { findStrictReceivePaths } from '../../src/lib/stellar.js'; import { resetExchangeRateCache, generateRateCacheKey } from '../../src/lib/exchange-rate-cache.js'; +import { resetExchangeRateOracleHealth } from '../../src/lib/exchange-rate-oracle-telemetry.js'; import { configureExchangeRateCoordination, resetExchangeRateCoordination, @@ -150,6 +151,7 @@ describe('Exchange Rate Oracle Cache — HTTP integration', () => { for (let i = 1; i <= 5; i++) seedPayment(i, `${i}.0000000`); resetExchangeRateCache(); resetExchangeRateCoordination(); + resetExchangeRateOracleHealth(); findStrictReceivePaths.mockReset(); findStrictReceivePaths.mockImplementation(async ({ destAmount }) => horizonPath(destAmount)); }); @@ -254,6 +256,21 @@ describe('Exchange Rate Oracle Cache — HTTP integration', () => { expect(res.text).toContain('exchange_rate_cache_inflight_loads'); expect(res.text).toContain('exchange_rate_lock_acquisitions_total'); expect(res.text).toContain('exchange_rate_coordination_fallbacks_total'); + expect(res.text).toContain('exchange_rate_oracle_cache_health_state'); + expect(res.text).toContain('exchange_rate_oracle_cache_loads_total'); + }); + + it('reports oracle cache health without quote or account data', async () => { + expect((await quote(app)).status).toBe(200); + const res = await request(app).get('/health/exchange-rate-oracle-cache'); + expect(res.status).toBe(200); + expect(res.body).toMatchObject({ status: 'healthy', success: 1, errors: 0 }); + expect(res.body.thresholds.min_samples).toBeGreaterThan(0); + expect(JSON.stringify(res.body)).not.toContain(SOURCE_ACCOUNT); + expect(JSON.stringify(res.body)).not.toContain('USDC'); + + const health = await request(app).get('/health'); + expect(health.body.services.exchange_rate_oracle_cache).toBe('healthy'); }); describe('with Redis coordination', () => { From 7f09da2788a84adc9c38fb8bec332085526051b3 Mon Sep 17 00:00:00 2001 From: Ved-viraj Date: Tue, 29 Sep 2026 10:57:29 +0100 Subject: [PATCH 2/2] feat(backend): implement automated retry with exponential backoff in Exchange Rate Oracle Cache Retry transient Horizon quote failures inside the single-flight loader so a burst shares one backoff loop and deterministic errors are not retried. --- backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md | 37 +++- .../exchange-rate-oracle-cache.rules.yml | 11 ++ backend/src/lib/exchange-rate-oracle-retry.js | 173 +++++++++++++++++ .../lib/exchange-rate-oracle-retry.test.js | 175 ++++++++++++++++++ .../src/lib/exchange-rate-oracle-telemetry.js | 17 ++ backend/src/services/exchangeRateService.js | 26 ++- .../src/services/exchangeRateService.test.js | 31 +++- .../integration/exchange-rate-cache.test.js | 10 +- 8 files changed, 464 insertions(+), 16 deletions(-) create mode 100644 backend/src/lib/exchange-rate-oracle-retry.js create mode 100644 backend/src/lib/exchange-rate-oracle-retry.test.js diff --git a/backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md b/backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md index 06da2e52..5b29cd33 100644 --- a/backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md +++ b/backend/docs/EXCHANGE_RATE_ORACLE_CACHE.md @@ -4,6 +4,7 @@ Caching and concurrency control for path-payment exchange-rate quotes (`GET /api/path-payment-quote/:id`). Covers issues **#1443** (Prometheus alert metrics and health telemetry), +**#1444** (automated retry with exponential backoff), **#1445** (distributed concurrency control and locking) and **#1446** (integration and stress test suite). @@ -18,7 +19,8 @@ Covers issues **#1443** (Prometheus alert metrics and health telemetry), | `src/services/exchangeRateService.js` | `getExchangeRateQuote()` composes both layers around the Horizon query | | `src/lib/path-payment-metrics.js` | Prometheus series (existing cache metrics + concurrency metrics) | | `src/lib/exchange-rate-oracle-telemetry.js` | Alert metrics, rolling-window health, separate registry (#1443) | -| `docs/alerts/exchange-rate-oracle-cache.rules.yml` | Prometheus alert rules (#1443) | +| `src/lib/exchange-rate-oracle-retry.js` | Full-jitter exponential backoff around the Horizon read (#1444) | +| `docs/alerts/exchange-rate-oracle-cache.rules.yml` | Prometheus alert rules (#1443, #1444) | ``` getExchangeRateQuote(key) @@ -208,13 +210,44 @@ Failed loads are logged at `warn` with outcome, duration, status and error message. The cache key is not logged. A failure inside the metrics client is swallowed so it cannot replace the loader's error. -## 9. Tests +| Metric | Type | Labels | +|---|---|---| +| `exchange_rate_oracle_cache_retries_total` | counter | `result` (scheduled/recovered/exhausted) | + +## 9. Retry with exponential backoff (#1444) + +`getExchangeRateQuote()` wraps the Horizon read in `withOracleRetry()`. The +wrapper sits inside `ExchangeRateCache.getOrLoad()`, so concurrent callers +for the same key share one retry loop. Cache hits do not retry. Redis lock +polling is unchanged. + +Delay is full jitter: `random(0, min(maxDelayMs, baseDelayMs * 2^n))`. +Defaults are 3 attempts, 100ms base, 1000ms cap. Env and callers cannot +exceed 6 attempts or 10s of delay. + +| Variable | Default | +|---|---| +| `EXCHANGE_RATE_ORACLE_RETRY_MAX_ATTEMPTS` | `3` | +| `EXCHANGE_RATE_ORACLE_RETRY_BASE_DELAY_MS` | `100` | +| `EXCHANGE_RATE_ORACLE_RETRY_MAX_DELAY_MS` | `1000` | + +Retried: network errors, HTTP 408, 429, and 5xx except 501. +Not retried: `NoPathFoundError` (404), other 4xx, `CacheLoadTimeoutError`, +and any error with `retryable: false`. The last error is rethrown with +`retryAttempts` set, and it is not cached. The outer load timeout still +bounds how long HTTP waiters block. + +The quote is a read of public DEX data, so a retry cannot submit a payment. +`ExchangeRateOracleCacheRetriesExhausted` fires when a load gives up. + +## 10. Tests The suites were mutation-checked against `exchange-rate-cache.js`. Disabling single-flight fails 16 tests. Dropping the invalidation guard fails 4. ``` npx vitest run src/lib/exchange-rate-oracle-telemetry.test.js \ + src/lib/exchange-rate-oracle-retry.test.js \ src/lib/exchange-rate-cache.test.js \ src/lib/exchange-rate-coordinator.test.js \ src/services/exchangeRateService.test.js \ diff --git a/backend/docs/alerts/exchange-rate-oracle-cache.rules.yml b/backend/docs/alerts/exchange-rate-oracle-cache.rules.yml index 9a2b3c80..6dea339a 100644 --- a/backend/docs/alerts/exchange-rate-oracle-cache.rules.yml +++ b/backend/docs/alerts/exchange-rate-oracle-cache.rules.yml @@ -63,6 +63,17 @@ groups: Quotes are being refreshed constantly. The TTL may be shorter than the request interval, or invalidation is too aggressive. + - alert: ExchangeRateOracleCacheRetriesExhausted + expr: sum(increase(exchange_rate_oracle_cache_retries_total{result="exhausted"}[10m])) > 0 + for: 5m + labels: + severity: warning + annotations: + summary: Exchange-rate oracle retries were exhausted + description: > + A quote load kept failing after exponential backoff. The last error + is returned to the caller and is not cached. Check Horizon. + - alert: ExchangeRateOracleCacheSlowLoads expr: | histogram_quantile(0.99, diff --git a/backend/src/lib/exchange-rate-oracle-retry.js b/backend/src/lib/exchange-rate-oracle-retry.js new file mode 100644 index 00000000..4be75ea1 --- /dev/null +++ b/backend/src/lib/exchange-rate-oracle-retry.js @@ -0,0 +1,173 @@ +/** + * Automated retry with exponential backoff for Exchange Rate Oracle Cache + * loads (issue #1444). + * + * The quote fetch is a read of public DEX data. Retrying it cannot create a + * payment or change a balance. Only transient failures are retried: network + * errors, 408, 429 and 5xx (except 501). A missing path (404), validation + * failures and CacheLoadTimeoutError are never retried, so a deterministic + * rejection cannot be turned into a quote by repetition. + * + * The wrapper runs inside the single-flight loader, so a burst of identical + * requests shares one retry loop instead of one loop per caller. + * + * Delay schedule: full jitter + * delay(n) = random(0, min(maxDelayMs, baseDelayMs * 2^n)) + */ + +import { logger } from './logger.js'; +import { recordOracleRetry } from './exchange-rate-oracle-telemetry.js'; + +export const DEFAULT_ORACLE_RETRY_OPTIONS = Object.freeze({ + maxAttempts: 3, + baseDelayMs: 100, + maxDelayMs: 1000, +}); + +const MAX_ATTEMPTS_CEILING = 6; +const MAX_DELAY_CEILING_MS = 10_000; + +const RETRYABLE_NETWORK_CODES = new Set([ + 'ECONNRESET', + 'ECONNREFUSED', + 'ECONNABORTED', + 'ETIMEDOUT', + 'EPIPE', + 'EAI_AGAIN', + 'ENOTFOUND', + 'EHOSTUNREACH', + 'ENETUNREACH', + 'UND_ERR_SOCKET', + 'UND_ERR_CONNECT_TIMEOUT', + 'UND_ERR_HEADERS_TIMEOUT', +]); + +function getStatus(err) { + const status = err?.status ?? err?.statusCode ?? err?.response?.status; + return typeof status === 'number' ? status : null; +} + +/** + * @param {unknown} err + * @returns {boolean} + */ +export function isRetryableOracleError(err) { + if (!err || typeof err !== 'object') return false; + if (err.name === 'NoPathFoundError' || err.name === 'CacheLoadTimeoutError') return false; + if (err.retryable === false) return false; + if (err.retryable === true) return true; + + const status = getStatus(err); + if (status !== null) { + if (status === 408 || status === 429) return true; + if (status >= 500 && status !== 501) return true; + if (status >= 400) return false; + } + + if (RETRYABLE_NETWORK_CODES.has(err.code)) return true; + + const message = typeof err.message === 'string' ? err.message : ''; + if (/fetch failed|socket hang up|network error|timed? ?out/i.test(message)) return true; + return false; +} + +function clampInt(value, fallback, min, max) { + const n = Number(value); + if (!Number.isFinite(n)) return fallback; + return Math.min(Math.max(Math.trunc(n), min), max); +} + +export function resolveOracleRetryOptions(options = {}, env = process.env) { + const maxAttempts = clampInt( + options.maxAttempts ?? env.EXCHANGE_RATE_ORACLE_RETRY_MAX_ATTEMPTS, + DEFAULT_ORACLE_RETRY_OPTIONS.maxAttempts, + 1, + MAX_ATTEMPTS_CEILING, + ); + const baseDelayMs = clampInt( + options.baseDelayMs ?? env.EXCHANGE_RATE_ORACLE_RETRY_BASE_DELAY_MS, + DEFAULT_ORACLE_RETRY_OPTIONS.baseDelayMs, + 0, + MAX_DELAY_CEILING_MS, + ); + const maxDelayMs = clampInt( + options.maxDelayMs ?? env.EXCHANGE_RATE_ORACLE_RETRY_MAX_DELAY_MS, + DEFAULT_ORACLE_RETRY_OPTIONS.maxDelayMs, + baseDelayMs, + MAX_DELAY_CEILING_MS, + ); + return { maxAttempts, baseDelayMs, maxDelayMs }; +} + +/** + * @param {number} retryIndex 0 for the first retry + * @param {{baseDelayMs:number, maxDelayMs:number}} opts + * @param {() => number} [random] + */ +export function computeOracleBackoffDelay(retryIndex, { baseDelayMs, maxDelayMs }, random = Math.random) { + const ceiling = Math.min(maxDelayMs, baseDelayMs * 2 ** retryIndex); + if (ceiling <= 0) return 0; + return Math.floor(random() * ceiling); +} + +const defaultSleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)); + +/** + * @template T + * @param {(ctx:{attempt:number}) => Promise} operation + * @param {object} [options] + */ +export async function withOracleRetry(operation, options = {}) { + const { + label = 'exchange-rate-oracle-fetch', + shouldRetry = isRetryableOracleError, + sleep = defaultSleep, + random = Math.random, + context = {}, + } = options; + const resolved = resolveOracleRetryOptions(options); + + let lastError; + for (let attempt = 1; attempt <= resolved.maxAttempts; attempt += 1) { + try { + const result = await operation({ attempt }); + if (attempt > 1) recordOracleRetry('recovered'); + return result; + } catch (err) { + lastError = err; + const retryable = shouldRetry(err); + const exhausted = attempt >= resolved.maxAttempts; + + if (!retryable || exhausted) { + if (err && typeof err === 'object') err.retryAttempts = attempt; + if (retryable && exhausted) { + recordOracleRetry('exhausted'); + logger.error( + { ...context, label, attempts: attempt, err: err?.message, code: err?.code, status: getStatus(err) }, + 'Exchange rate oracle fetch failed after exhausting retries', + ); + } + throw err; + } + + const delayMs = computeOracleBackoffDelay(attempt - 1, resolved, random); + logger.warn( + { + ...context, + label, + attempt, + maxAttempts: resolved.maxAttempts, + delayMs, + err: err?.message, + code: err?.code, + status: getStatus(err), + }, + 'Exchange rate oracle fetch failed, retrying with backoff', + ); + recordOracleRetry('scheduled'); + await sleep(delayMs); + } + } + + throw lastError; +} diff --git a/backend/src/lib/exchange-rate-oracle-retry.test.js b/backend/src/lib/exchange-rate-oracle-retry.test.js new file mode 100644 index 00000000..3dd138d7 --- /dev/null +++ b/backend/src/lib/exchange-rate-oracle-retry.test.js @@ -0,0 +1,175 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest'; + +vi.mock('./logger.js', () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() }, +})); + +import { logger } from './logger.js'; +import { + DEFAULT_ORACLE_RETRY_OPTIONS, + computeOracleBackoffDelay, + isRetryableOracleError, + resolveOracleRetryOptions, + withOracleRetry, +} from './exchange-rate-oracle-retry.js'; +import { + exchangeRateOracleRegister, + resetExchangeRateOracleHealth, +} from './exchange-rate-oracle-telemetry.js'; + +const noSleep = vi.fn(async () => {}); + +function errWith(props) { + return Object.assign(new Error(props.message ?? 'boom'), props); +} + +async function retryCount(result) { + const metric = exchangeRateOracleRegister.getSingleMetric('exchange_rate_oracle_cache_retries_total'); + const snapshot = await metric.get(); + return snapshot.values + .filter((entry) => entry.labels.result === result) + .reduce((sum, entry) => sum + entry.value, 0); +} + +beforeEach(() => { + vi.clearAllMocks(); + resetExchangeRateOracleHealth(); +}); + +describe('isRetryableOracleError (issue #1444)', () => { + it.each([ + ['HTTP 500', { status: 500 }], + ['HTTP 502', { status: 502 }], + ['HTTP 503', { statusCode: 503 }], + ['HTTP 429', { status: 429 }], + ['HTTP 408', { status: 408 }], + ['ECONNRESET', { code: 'ECONNRESET' }], + ['ETIMEDOUT', { code: 'ETIMEDOUT' }], + ['fetch failed', { message: 'TypeError: fetch failed' }], + ['socket hang up', { message: 'socket hang up' }], + ['explicit retryable flag', { retryable: true, status: 400 }], + ])('retries %s', (_label, props) => { + expect(isRetryableOracleError(errWith(props))).toBe(true); + }); + + it.each([ + ['HTTP 400', { status: 400 }], + ['HTTP 401', { status: 401 }], + ['HTTP 404', { status: 404 }], + ['HTTP 409', { status: 409 }], + ['HTTP 422', { status: 422 }], + ['HTTP 501', { status: 501 }], + ['no path', { name: 'NoPathFoundError', statusCode: 404 }], + ['load timeout', { name: 'CacheLoadTimeoutError', status: 504 }], + ['opt-out beats 503', { retryable: false, status: 503 }], + ['plain error', { message: 'something unexpected' }], + ])('does not retry %s', (_label, props) => { + expect(isRetryableOracleError(errWith(props))).toBe(false); + }); + + it('does not retry non-object throwables', () => { + expect(isRetryableOracleError(null)).toBe(false); + expect(isRetryableOracleError(undefined)).toBe(false); + expect(isRetryableOracleError('ECONNRESET')).toBe(false); + }); +}); + +describe('resolveOracleRetryOptions', () => { + it('uses defaults when nothing is configured', () => { + expect(resolveOracleRetryOptions({}, {})).toEqual({ ...DEFAULT_ORACLE_RETRY_OPTIONS }); + }); + + it('reads env overrides and clamps hostile values', () => { + expect(resolveOracleRetryOptions({}, { + EXCHANGE_RATE_ORACLE_RETRY_MAX_ATTEMPTS: '4', + EXCHANGE_RATE_ORACLE_RETRY_BASE_DELAY_MS: '25', + EXCHANGE_RATE_ORACLE_RETRY_MAX_DELAY_MS: '400', + })).toEqual({ maxAttempts: 4, baseDelayMs: 25, maxDelayMs: 400 }); + + const hostile = resolveOracleRetryOptions( + { maxAttempts: 10_000, baseDelayMs: -5, maxDelayMs: 9e9 }, + {}, + ); + expect(hostile).toEqual({ maxAttempts: 6, baseDelayMs: 0, maxDelayMs: 10_000 }); + }); + + it('keeps maxDelay at least baseDelay and at least one attempt', () => { + const opts = resolveOracleRetryOptions({ maxAttempts: 0, baseDelayMs: 400, maxDelayMs: 10 }, {}); + expect(opts.maxAttempts).toBe(1); + expect(opts.maxDelayMs).toBe(400); + }); +}); + +describe('computeOracleBackoffDelay', () => { + const opts = { baseDelayMs: 100, maxDelayMs: 1000 }; + + it('grows the ceiling exponentially and caps it', () => { + const max = () => 0.999999; + expect(computeOracleBackoffDelay(0, opts, max)).toBe(99); + expect(computeOracleBackoffDelay(1, opts, max)).toBe(199); + expect(computeOracleBackoffDelay(2, opts, max)).toBe(399); + expect(computeOracleBackoffDelay(8, opts, max)).toBe(999); + }); + + it('can delay 0 under full jitter', () => { + expect(computeOracleBackoffDelay(3, opts, () => 0)).toBe(0); + }); +}); + +describe('withOracleRetry', () => { + it('returns the first success without retrying', async () => { + const operation = vi.fn(async () => 'quote'); + await expect(withOracleRetry(operation, { sleep: noSleep, maxAttempts: 3 })).resolves.toBe('quote'); + expect(operation).toHaveBeenCalledTimes(1); + expect(logger.warn).not.toHaveBeenCalled(); + }); + + it('retries a transient error then recovers', async () => { + const before = await retryCount('scheduled'); + const operation = vi.fn() + .mockRejectedValueOnce(errWith({ status: 503, message: 'unavailable' })) + .mockResolvedValueOnce('quote'); + + await expect(withOracleRetry(operation, { sleep: noSleep, random: () => 0, maxAttempts: 3 })).resolves.toBe('quote'); + expect(operation).toHaveBeenCalledTimes(2); + expect(logger.warn).toHaveBeenCalledTimes(1); + expect(await retryCount('scheduled')).toBe(before + 1); + expect(await retryCount('recovered')).toBeGreaterThan(0); + }); + + it('does not retry a missing path', async () => { + const operation = vi.fn(async () => { + throw errWith({ name: 'NoPathFoundError', statusCode: 404, message: 'no path' }); + }); + await expect(withOracleRetry(operation, { sleep: noSleep, maxAttempts: 3 })).rejects.toThrow('no path'); + expect(operation).toHaveBeenCalledTimes(1); + expect(logger.error).not.toHaveBeenCalled(); + }); + + it('stops at the attempt ceiling and annotates the last error', async () => { + const before = await retryCount('exhausted'); + const operation = vi.fn(async () => { + throw errWith({ status: 502, message: 'bad gateway' }); + }); + await expect(withOracleRetry(operation, { sleep: noSleep, maxAttempts: 3 })).rejects.toMatchObject({ + message: 'bad gateway', + retryAttempts: 3, + }); + expect(operation).toHaveBeenCalledTimes(3); + expect(logger.error).toHaveBeenCalledTimes(1); + expect(await retryCount('exhausted')).toBe(before + 1); + }); + + it('passes a 1-based attempt number to the operation', async () => { + const seen = []; + await withOracleRetry( + async ({ attempt }) => { + seen.push(attempt); + if (attempt < 2) throw errWith({ code: 'ECONNRESET' }); + return 'ok'; + }, + { sleep: noSleep, maxAttempts: 3 }, + ); + expect(seen).toEqual([1, 2]); + }); +}); diff --git a/backend/src/lib/exchange-rate-oracle-telemetry.js b/backend/src/lib/exchange-rate-oracle-telemetry.js index 2fae7b5e..f83d13ca 100644 --- a/backend/src/lib/exchange-rate-oracle-telemetry.js +++ b/backend/src/lib/exchange-rate-oracle-telemetry.js @@ -72,8 +72,16 @@ export const exchangeRateOracleLastLoadTimestamp = new client.Gauge({ help: 'Unix time of the most recent exchange-rate oracle cache load', }); +/** result: scheduled | recovered | exhausted (issue #1444) */ +export const exchangeRateOracleRetriesTotal = new client.Counter({ + name: 'exchange_rate_oracle_cache_retries_total', + help: 'Exchange-rate oracle fetch retries with exponential backoff, by result', + labelNames: ['result'], +}); + const LOOKUP_RESULTS = new Set(['hit', 'miss', 'stale']); const LOAD_OUTCOMES = new Set(['success', 'error', 'timeout', 'not_found']); +const RETRY_RESULTS = new Set(['scheduled', 'recovered', 'exhausted']); export const HEALTH_STATE_VALUES = Object.freeze({ healthy: 0, degraded: 1, unhealthy: 2 }); @@ -313,5 +321,14 @@ register.registerMetric(exchangeRateOracleErrorRatio); register.registerMetric(exchangeRateOracleTimeoutRatio); register.registerMetric(exchangeRateOracleStaleRatio); register.registerMetric(exchangeRateOracleLastLoadTimestamp); +register.registerMetric(exchangeRateOracleRetriesTotal); + +/** @param {'scheduled'|'recovered'|'exhausted'} result */ +export function recordOracleRetry(result) { + const safe = RETRY_RESULTS.has(result) ? result : 'exhausted'; + safeRecord(() => { + exchangeRateOracleRetriesTotal.inc({ result: safe }); + }); +} export { register as exchangeRateOracleRegister }; diff --git a/backend/src/services/exchangeRateService.js b/backend/src/services/exchangeRateService.js index 571cddb5..548f418a 100644 --- a/backend/src/services/exchangeRateService.js +++ b/backend/src/services/exchangeRateService.js @@ -15,6 +15,10 @@ * client, that load is further coordinated across instances by a * distributed lock + shared quote store (exchange-rate-coordinator.js). * Coordination fails open to a direct Horizon query. + * - The Horizon read itself is retried with full-jitter exponential backoff + * (issue #1444). Retries stay inside the single-flight loader, so a burst + * shares one retry loop. Missing paths and other deterministic failures + * are not retried. * * The route handler calls getExchangeRateQuote() and only handles HTTP concerns; * all exchange-rate logic lives here. @@ -26,6 +30,7 @@ import { generateRateCacheKey, } from '../lib/exchange-rate-cache.js'; import { ExchangeRateCoordinator } from '../lib/exchange-rate-coordinator.js'; +import { withOracleRetry } from '../lib/exchange-rate-oracle-retry.js'; import { logger } from '../lib/logger.js'; const DEFAULT_SLIPPAGE = parseFloat(process.env.PATH_PAYMENT_SLIPPAGE ?? '0.01'); @@ -108,15 +113,18 @@ export async function getExchangeRateQuote({ destAssetIssuer, ); - const fetchFromHorizon = () => fetchQuoteFromHorizon({ - sourceAssetCode, - sourceAssetIssuer, - destAssetCode, - destAssetIssuer, - destAmount, - sourceAccount, - slippage, - }); + const fetchFromHorizon = () => withOracleRetry( + () => fetchQuoteFromHorizon({ + sourceAssetCode, + sourceAssetIssuer, + destAssetCode, + destAssetIssuer, + destAmount, + sourceAccount, + slippage, + }), + { label: 'exchange-rate-oracle-fetch' }, + ); // Captured per call so a later reconfiguration cannot change an in-flight load. const activeCoordinator = coordinator; diff --git a/backend/src/services/exchangeRateService.test.js b/backend/src/services/exchangeRateService.test.js index a2b0b56f..9f355e27 100644 --- a/backend/src/services/exchangeRateService.test.js +++ b/backend/src/services/exchangeRateService.test.js @@ -34,11 +34,15 @@ describe('getExchangeRateQuote', () => { beforeEach(() => { resetExchangeRateCache(); vi.clearAllMocks(); + process.env.EXCHANGE_RATE_ORACLE_RETRY_BASE_DELAY_MS = '0'; + process.env.EXCHANGE_RATE_ORACLE_RETRY_MAX_DELAY_MS = '0'; findStrictReceivePaths.mockResolvedValue(MOCK_PATH); }); afterEach(() => { resetExchangeRateCache(); + delete process.env.EXCHANGE_RATE_ORACLE_RETRY_BASE_DELAY_MS; + delete process.env.EXCHANGE_RATE_ORACLE_RETRY_MAX_DELAY_MS; }); it('returns a quote with correct fields', async () => { @@ -96,13 +100,36 @@ describe('getExchangeRateQuote', () => { } }); - it('propagates Horizon errors as-is', async () => { + it('propagates Horizon errors as-is after retries are exhausted', async () => { const horizonErr = new Error('Horizon unavailable'); horizonErr.status = 503; findStrictReceivePaths.mockRejectedValue(horizonErr); await expect( getExchangeRateQuote({ sourceAssetCode: 'XLM', destAssetCode: 'USDC', destAmount: '1.0' }), ).rejects.toThrow('Horizon unavailable'); + expect(findStrictReceivePaths).toHaveBeenCalledTimes(3); + expect(horizonErr.retryAttempts).toBe(3); + }); + + it('retries a transient Horizon failure and then returns the quote', async () => { + findStrictReceivePaths + .mockRejectedValueOnce(Object.assign(new Error('unavailable'), { status: 503 })) + .mockResolvedValueOnce(MOCK_PATH); + const quote = await getExchangeRateQuote({ + sourceAssetCode: 'XLM', + destAssetCode: 'USDC', + destAmount: '1.0', + }); + expect(quote.sourceAmount).toBe('0.5000000'); + expect(findStrictReceivePaths).toHaveBeenCalledTimes(2); + }); + + it('does not retry a deterministic 400', async () => { + findStrictReceivePaths.mockRejectedValue(Object.assign(new Error('bad request'), { status: 400 })); + await expect( + getExchangeRateQuote({ sourceAssetCode: 'XLM', destAssetCode: 'USDC', destAmount: '1.0' }), + ).rejects.toThrow('bad request'); + expect(findStrictReceivePaths).toHaveBeenCalledTimes(1); }); it('applies custom slippage correctly', async () => { @@ -153,6 +180,8 @@ describe('concurrency control (issue #1445)', () => { resetExchangeRateCache(); resetExchangeRateCoordination(); vi.clearAllMocks(); + process.env.EXCHANGE_RATE_ORACLE_RETRY_BASE_DELAY_MS = '0'; + process.env.EXCHANGE_RATE_ORACLE_RETRY_MAX_DELAY_MS = '0'; findStrictReceivePaths.mockResolvedValue(MOCK_PATH); }); diff --git a/backend/tests/integration/exchange-rate-cache.test.js b/backend/tests/integration/exchange-rate-cache.test.js index d11a21b9..a81676ca 100644 --- a/backend/tests/integration/exchange-rate-cache.test.js +++ b/backend/tests/integration/exchange-rate-cache.test.js @@ -13,6 +13,8 @@ vi.hoisted(() => { process.env.PATH_PAYMENT_QUOTE_RATE_LIMIT_MAX = '100000'; // Wide enough that every request in a burst joins the load before it times out. process.env.EXCHANGE_RATE_LOAD_TIMEOUT_MS = '1000'; + process.env.EXCHANGE_RATE_ORACLE_RETRY_BASE_DELAY_MS = '0'; + process.env.EXCHANGE_RATE_ORACLE_RETRY_MAX_DELAY_MS = '0'; }); import { createApp } from '../../src/app.js'; @@ -210,11 +212,11 @@ describe('Exchange Rate Oracle Cache — HTTP integration', () => { }); it('fails every waiter on a Horizon error without caching it', async () => { - const horizon = gatedHorizon(); - const burst = await burstJoined(app, horizon, 10); - horizon.failAll(Object.assign(new Error('Horizon 503'), { status: 502 })); - const responses = await Promise.all(burst); + findStrictReceivePaths.mockRejectedValue(Object.assign(new Error('Horizon 503'), { status: 502 })); + const responses = await Promise.all(Array.from({ length: 10 }, () => quote(app))); expect(responses.every((r) => r.status === 502)).toBe(true); + // One single-flight load, retried to the attempt ceiling — not once per waiter. + expect(findStrictReceivePaths).toHaveBeenCalledTimes(3); findStrictReceivePaths.mockImplementation(async ({ destAmount }) => horizonPath(destAmount)); expect((await quote(app)).status).toBe(200);