From ce8abb0e7e288dff540f9b6f81210d6bbbcb8a96 Mon Sep 17 00:00:00 2001 From: Sweet-Kid Date: Sat, 26 Sep 2026 11:50:42 +0100 Subject: [PATCH] fix: configurable preflight max-age, durable webhook metrics, degraded boot, ledger cursor CAS MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - cors.js/config.js: preflight responses advertise Access-Control-Max-Age from CORS_MAX_AGE_SECONDS (default 24h) so the cache window is explicit and tunable instead of a magic literal (#344) - webhookDispatcher.js: delivery counters are written through to Redis (one aggregate hash plus one hash per webhook) and reloaded once at startup via hydrateMetrics(), so success rates survive a restart (#343) - index.js: startServer() runs in degraded mode when Redis is down at boot — warn once, hydrate metrics best-effort, then warm the cache — rather than letting either step stop the process (#342) - eventStore.js/eventPoller.js: advance indexer:last_ledger through a compare-and-set Lua script that also refuses to move the cursor backwards; the poller keeps the events it already saved and reports a skip when another writer wins the race (#341) Collateral repairs to code that was broken on main and blocked these fixes: config.js module.exports structure (syntax error), the duplicate /health websocket key, the priceWebSocket getHealth() export lost in a merge, and priceOracle's single-flight wrapper promise (#417). --- src/config.js | 227 +++++++++++++------------ src/index.js | 32 +++- src/indexer/eventPoller.js | 24 ++- src/indexer/eventStore.js | 68 ++++++++ src/middleware/cors.js | 7 +- src/services/priceOracle.js | 32 ++-- src/services/webhookDispatcher.js | 119 ++++++++++++- src/startup/cacheWarm.js | 8 +- src/ws/priceWebSocket.js | 18 +- test/cors.test.js | 45 +++++ test/eventPoller.test.js | 113 +++++++++++- test/health.test.js | 12 ++ test/helpers/cacheMock.js | 51 ++++++ test/indexerBackoff.test.js | 3 +- test/indexerStore.test.js | 72 ++++++++ test/priceOracle.test.js | 4 +- test/startupBanner.test.js | 57 +++++++ test/webhookDispatcher.test.js | 7 + test/webhookMetricsPersistence.test.js | 179 +++++++++++++++++++ 19 files changed, 936 insertions(+), 142 deletions(-) create mode 100644 test/webhookMetricsPersistence.test.js diff --git a/src/config.js b/src/config.js index 7533c2f..6e6442d 100644 --- a/src/config.js +++ b/src/config.js @@ -88,6 +88,7 @@ const env = cleanEnv(rawEnv, { AIRDROP_JSON_MAX_BYTES: positiveInteger({ default: 2 * 1024 * 1024 }), AIRDROP_RATELIMIT_WINDOW: positiveInteger({ default: 60 }), AIRDROP_RATELIMIT_MAX: positiveInteger({ default: 10 }), + CORS_MAX_AGE_SECONDS: positiveInteger({ default: 86400 }), PRICE_CACHE_TTL_SECONDS: num({ default: 60 }), PRICE_REFRESH_INTERVAL_SECONDS: num({ default: 30 }), PRICE_STALE_THRESHOLD_MINUTES: num({ default: 5 }), @@ -129,116 +130,7 @@ const parsedWatchedAssets = Array.isArray(env.WATCHED_ASSETS) ? env.WATCHED_ASSETS : parseWatchedAssets(env.WATCHED_ASSETS); -const config = module.exports; - -// #291 — Runtime config hot-reload support. When the server receives a -// SIGHUP signal (or watches for .env file changes in development), it -// re-reads environment variables and re-validates them via envalid, then -// replaces the exported config values in-place so all consumers see the -// updated config without a restart. -function reload() { - require('dotenv').config({ override: true }); - const reloaded = cleanEnv({ ...process.env }, { - NODE_ENV: str({ default: 'development', choices: ['development', 'test', 'production'] }), - PORT: port({ default: 3000 }), - REDIS_URL: url({ devDefault: 'redis://localhost:6379' }), - DATABASE_URL: url({ devDefault: databaseDevDefault }), - STELLAR_HORIZON_URL: url({ default: 'https://horizon.stellar.org' }), - SOROBAN_RPC_URL: url({ default: 'https://soroban-rpc.mainnet.stellar.gateway.fm' }), - USDC_ISSUER: stellarAddress({ default: 'GA5ZSEJYB37JRC5AVCIA5MOP4RHTM335AX2OBFLDTQLNUEHRGPTM6RIA' }), - COINGECKO_API_KEY: str({ default: '' }), - COINMARKETCAP_API_KEY: str({ default: '' }), - INSTANCE_ID: str({ default: '' }), - LEASE_TTL_MS: positiveInteger({ default: 15000 }), - LEASE_RENEW_INTERVAL_MS: positiveInteger({ default: 5000 }), - ADMIN_API_KEY: str({ default: '' }), - WEBHOOK_SECRET_ENCRYPTION_KEY: str({ default: '' }), - AIRDROP_CSV_MAX_BYTES: positiveInteger({ default: 5 * 1024 * 1024 }), - AIRDROP_JSON_MAX_BYTES: positiveInteger({ default: 2 * 1024 * 1024 }), - AIRDROP_RATELIMIT_WINDOW: positiveInteger({ default: 60 }), - AIRDROP_RATELIMIT_MAX: positiveInteger({ default: 10 }), - PRICE_CACHE_TTL_SECONDS: num({ default: 60 }), - PRICE_REFRESH_INTERVAL_SECONDS: num({ default: 30 }), - PRICE_STALE_THRESHOLD_MINUTES: num({ default: 5 }), - PRICE_ANOMALY_THRESHOLD_PCT: num({ default: 20 }), - PRICE_MIN_SOURCES: num({ default: 2 }), - PRICE_ANOMALY_ACTION: str({ default: 'warn', choices: ['warn', 'reject'] }), - PRICE_SOURCE_PRIORITY: str({ default: '' }), - PRICE_REFRESH_MAX_CYCLE_MS: num({ default: 90000 }), - CIRCUIT_BREAKER_FAILURE_THRESHOLD: num({ default: 3 }), - CIRCUIT_BREAKER_SUCCESS_THRESHOLD: num({ default: 1 }), - CIRCUIT_BREAKER_TIMEOUT_MS: num({ default: 30000 }), - PRICE_SOURCE_CIRCUIT_COOLDOWN_MS: num({ default: 15 * 60 * 1000 }), - PRICE_SOURCE_CIRCUIT_REMINDER_MS: num({ default: 5 * 60 * 1000 }), - AIRDROP_EXPIRY_CHECK_INTERVAL_SECONDS: num({ default: 60 }), - AIRDROP_LEDGER_CACHE_TTL_MS: num({ default: 5000 }), - AIRDROP_EXPIRY_SCAN_BATCH_SIZE: num({ default: 100 }), - WATCHED_ASSETS: watchedAssets({ default: '' }), - SENTRY_DSN: str({ default: '' }), - LOG_LEVEL: str({ default: 'info', choices: ['debug', 'info', 'warn', 'error'] }), - SLOW_REQUEST_THRESHOLD_MS: num({ default: 1000 }), - ROUTE_TIMEOUT_MS: num({ default: 30000 }), - API_KEY_RATELIMIT_WINDOW_SECONDS: positiveInteger({ default: 60 }), - API_KEY_RATELIMIT_FREE_MAX: positiveInteger({ default: 100 }), - API_KEY_RATELIMIT_PRO_MAX: positiveInteger({ default: 1000 }), - API_KEY_RATELIMIT_ADMIN_MAX: positiveInteger({ default: 10000 }), - }); - - const reloadedUsdcIssuer = reloaded.USDC_ISSUER; - const reloadedWatchedAssets = Array.isArray(reloaded.WATCHED_ASSETS) - ? reloaded.WATCHED_ASSETS - : parseWatchedAssets(reloaded.WATCHED_ASSETS); - - // Update all config properties in-place so existing references see new values. - Object.assign(config, { - nodeEnv: reloaded.NODE_ENV, - port: reloaded.PORT, - databaseUrl: reloaded.DATABASE_URL, - redis: { url: reloaded.REDIS_URL }, - stellar: { - horizonUrl: reloaded.STELLAR_HORIZON_URL, - sorobanRpcUrl: reloaded.SOROBAN_RPC_URL, - usdcIssuer: reloadedUsdcIssuer, - }, - auth: { adminApiKey: reloaded.ADMIN_API_KEY }, - // Issue #373: reload() already re-validated COINMARKETCAP_API_KEY above - // (it's in the env schema passed to validateEnv) but never applied it — - // config.coinmarketcap was silently left out of this Object.assign - // entirely, so even a caller that re-read config.coinmarketcap.apiKey - // fresh on every call (see coinmarketcap.js's own fix for #373) would - // still never see a reloaded key change without this. - coinmarketcap: { ...config.coinmarketcap, apiKey: reloaded.COINMARKETCAP_API_KEY }, - sentryDsn: reloaded.SENTRY_DSN, - slowRequestThresholdMs: reloaded.SLOW_REQUEST_THRESHOLD_MS, - routeTimeoutMs: reloaded.ROUTE_TIMEOUT_MS, - webhookSecretEncryptionKey: reloaded.WEBHOOK_SECRET_ENCRYPTION_KEY, - watchedAssets: reloadedWatchedAssets, - webhooks: { - ...config.webhooks, - maxAttempts: parseInt(process.env.WEBHOOK_MAX_ATTEMPTS, 10) || 3, - retryBaseMs: parseInt(process.env.WEBHOOK_RETRY_BASE_MS, 10) || 30000, - retryFactor: parseFloat(process.env.WEBHOOK_RETRY_FACTOR) || 2, - timeoutMs: parseInt(process.env.WEBHOOK_TIMEOUT_MS, 10) || 5000, - }, - LOG_LEVEL: reloaded.LOG_LEVEL, - }); - - return config; -} - -// Register SIGHUP handler for production hot-reload -if (process.env.NODE_ENV === 'production') { - process.on('SIGHUP', () => { - try { - reload(); - console.log('[config] Configuration reloaded via SIGHUP'); - } catch (err) { - console.error('[config] Failed to reload configuration:', err.message); - } - }); -} - -module.exports.reload = reload; +const config = module.exports = { nodeEnv: env.NODE_ENV, port: env.PORT, databaseUrl: env.DATABASE_URL, @@ -343,6 +235,9 @@ module.exports.reload = reload; slowRequestThresholdMs: env.SLOW_REQUEST_THRESHOLD_MS, routeTimeoutMs: env.ROUTE_TIMEOUT_MS, webhookSecretEncryptionKey: env.WEBHOOK_SECRET_ENCRYPTION_KEY, + // How long browsers may cache a CORS preflight response (issue #344). + // Exposed as Access-Control-Max-Age on every successful OPTIONS preflight. + corsMaxAgeSeconds: env.CORS_MAX_AGE_SECONDS, corsAllowedOrigins: (process.env.CORS_ALLOWED_ORIGINS || 'http://localhost:3000,http://localhost:3001') .split(',') .map((o) => o.trim()) @@ -387,3 +282,115 @@ module.exports.reload = reload; maxConnectionsPerIp: parseInt(process.env.WS_MAX_CONNECTIONS_PER_IP, 10) || 5, }, }; + + +// #291 — Runtime config hot-reload support. When the server receives a +// SIGHUP signal (or watches for .env file changes in development), it +// re-reads environment variables and re-validates them via envalid, then +// replaces the exported config values in-place so all consumers see the +// updated config without a restart. +function reload() { + require('dotenv').config({ override: true }); + const reloaded = cleanEnv({ ...process.env }, { + NODE_ENV: str({ default: 'development', choices: ['development', 'test', 'production'] }), + PORT: port({ default: 3000 }), + REDIS_URL: url({ devDefault: 'redis://localhost:6379' }), + DATABASE_URL: url({ devDefault: databaseDevDefault }), + STELLAR_HORIZON_URL: url({ default: 'https://horizon.stellar.org' }), + SOROBAN_RPC_URL: url({ default: 'https://soroban-rpc.mainnet.stellar.gateway.fm' }), + USDC_ISSUER: stellarAddress({ default: 'GA5ZSEJYB37JRC5AVCIA5MOP4RHTM335AX2OBFLDTQLNUEHRGPTM6RIA' }), + COINGECKO_API_KEY: str({ default: '' }), + COINMARKETCAP_API_KEY: str({ default: '' }), + INSTANCE_ID: str({ default: '' }), + LEASE_TTL_MS: positiveInteger({ default: 15000 }), + LEASE_RENEW_INTERVAL_MS: positiveInteger({ default: 5000 }), + ADMIN_API_KEY: str({ default: '' }), + WEBHOOK_SECRET_ENCRYPTION_KEY: str({ default: '' }), + AIRDROP_CSV_MAX_BYTES: positiveInteger({ default: 5 * 1024 * 1024 }), + AIRDROP_JSON_MAX_BYTES: positiveInteger({ default: 2 * 1024 * 1024 }), + AIRDROP_RATELIMIT_WINDOW: positiveInteger({ default: 60 }), + AIRDROP_RATELIMIT_MAX: positiveInteger({ default: 10 }), + CORS_MAX_AGE_SECONDS: positiveInteger({ default: 86400 }), + PRICE_CACHE_TTL_SECONDS: num({ default: 60 }), + PRICE_REFRESH_INTERVAL_SECONDS: num({ default: 30 }), + PRICE_STALE_THRESHOLD_MINUTES: num({ default: 5 }), + PRICE_ANOMALY_THRESHOLD_PCT: num({ default: 20 }), + PRICE_MIN_SOURCES: num({ default: 2 }), + PRICE_ANOMALY_ACTION: str({ default: 'warn', choices: ['warn', 'reject'] }), + PRICE_SOURCE_PRIORITY: str({ default: '' }), + PRICE_REFRESH_MAX_CYCLE_MS: num({ default: 90000 }), + CIRCUIT_BREAKER_FAILURE_THRESHOLD: num({ default: 3 }), + CIRCUIT_BREAKER_SUCCESS_THRESHOLD: num({ default: 1 }), + CIRCUIT_BREAKER_TIMEOUT_MS: num({ default: 30000 }), + PRICE_SOURCE_CIRCUIT_COOLDOWN_MS: num({ default: 15 * 60 * 1000 }), + PRICE_SOURCE_CIRCUIT_REMINDER_MS: num({ default: 5 * 60 * 1000 }), + AIRDROP_EXPIRY_CHECK_INTERVAL_SECONDS: num({ default: 60 }), + AIRDROP_LEDGER_CACHE_TTL_MS: num({ default: 5000 }), + AIRDROP_EXPIRY_SCAN_BATCH_SIZE: num({ default: 100 }), + WATCHED_ASSETS: watchedAssets({ default: '' }), + SENTRY_DSN: str({ default: '' }), + LOG_LEVEL: str({ default: 'info', choices: ['debug', 'info', 'warn', 'error'] }), + SLOW_REQUEST_THRESHOLD_MS: num({ default: 1000 }), + ROUTE_TIMEOUT_MS: num({ default: 30000 }), + API_KEY_RATELIMIT_WINDOW_SECONDS: positiveInteger({ default: 60 }), + API_KEY_RATELIMIT_FREE_MAX: positiveInteger({ default: 100 }), + API_KEY_RATELIMIT_PRO_MAX: positiveInteger({ default: 1000 }), + API_KEY_RATELIMIT_ADMIN_MAX: positiveInteger({ default: 10000 }), + }); + + const reloadedUsdcIssuer = reloaded.USDC_ISSUER; + const reloadedWatchedAssets = Array.isArray(reloaded.WATCHED_ASSETS) + ? reloaded.WATCHED_ASSETS + : parseWatchedAssets(reloaded.WATCHED_ASSETS); + + // Update all config properties in-place so existing references see new values. + Object.assign(config, { + nodeEnv: reloaded.NODE_ENV, + port: reloaded.PORT, + databaseUrl: reloaded.DATABASE_URL, + redis: { url: reloaded.REDIS_URL }, + stellar: { + horizonUrl: reloaded.STELLAR_HORIZON_URL, + sorobanRpcUrl: reloaded.SOROBAN_RPC_URL, + usdcIssuer: reloadedUsdcIssuer, + }, + auth: { adminApiKey: reloaded.ADMIN_API_KEY }, + // Issue #373: reload() already re-validated COINMARKETCAP_API_KEY above + // (it's in the env schema passed to validateEnv) but never applied it — + // config.coinmarketcap was silently left out of this Object.assign + // entirely, so even a caller that re-read config.coinmarketcap.apiKey + // fresh on every call (see coinmarketcap.js's own fix for #373) would + // still never see a reloaded key change without this. + coinmarketcap: { ...config.coinmarketcap, apiKey: reloaded.COINMARKETCAP_API_KEY }, + sentryDsn: reloaded.SENTRY_DSN, + slowRequestThresholdMs: reloaded.SLOW_REQUEST_THRESHOLD_MS, + routeTimeoutMs: reloaded.ROUTE_TIMEOUT_MS, + webhookSecretEncryptionKey: reloaded.WEBHOOK_SECRET_ENCRYPTION_KEY, + corsMaxAgeSeconds: reloaded.CORS_MAX_AGE_SECONDS, + watchedAssets: reloadedWatchedAssets, + webhooks: { + ...config.webhooks, + maxAttempts: parseInt(process.env.WEBHOOK_MAX_ATTEMPTS, 10) || 3, + retryBaseMs: parseInt(process.env.WEBHOOK_RETRY_BASE_MS, 10) || 30000, + retryFactor: parseFloat(process.env.WEBHOOK_RETRY_FACTOR) || 2, + timeoutMs: parseInt(process.env.WEBHOOK_TIMEOUT_MS, 10) || 5000, + }, + LOG_LEVEL: reloaded.LOG_LEVEL, + }); + + return config; +} + +// Register SIGHUP handler for production hot-reload +if (process.env.NODE_ENV === 'production') { + process.on('SIGHUP', () => { + try { + reload(); + console.log('[config] Configuration reloaded via SIGHUP'); + } catch (err) { + console.error('[config] Failed to reload configuration:', err.message); + } + }); +} + +module.exports.reload = reload; \ No newline at end of file diff --git a/src/index.js b/src/index.js index 8cb7f27..a0ff2db 100644 --- a/src/index.js +++ b/src/index.js @@ -126,6 +126,7 @@ app.get("/health", healthRateLimit, async (req, res) => { const webhookWorkerHealth = wrappedWebhookRetryWorker.getHealth(); const airdropExpiryHealth = wrappedAirdropExpiryJob.getHealth(); const database = await checkDatabase(); + const wsHealth = priceWebSocket.getHealth(); // Queue depth for the retry worker (issue #235) — "the worker is alive" // says nothing about whether retries are piling up behind it. Health must // still answer if this telemetry read fails, so a failure degrades to @@ -176,15 +177,12 @@ app.get("/health", healthRateLimit, async (req, res) => { command_queue_depth: redisQueueDepth, concurrency: redisConcurrency, }, - websocket: { - connections: subscriptionManager.connectionCount, - draining: subscriptionManager.isDraining, - drain_stats: subscriptionManager.drainStats, - }, websocket: { healthy: wsHealth.healthy, connections: wsHealth.connections, error: wsHealth.error, + draining: subscriptionManager.isDraining, + drain_stats: subscriptionManager.drainStats, }, jobs: { price_refresh: { @@ -348,6 +346,30 @@ function logStartupBanner() { } async function startServer() { + // Redis being unreachable at boot must not stop the process from coming up + // (#342): everything below degrades to "no cached state" and recovers on + // its own once ioredis reconnects. Say so once, explicitly, so an operator + // reading the logs can tell a degraded boot from a healthy one. + const redisConnectedAtBoot = cache.isConnected(); + if (!redisConnectedAtBoot) { + logger.warn( + "Redis unavailable at startup; starting in degraded mode (no cached state, rate limits, or leader lease until it reconnects)", + { redis_connected: false }, + ); + } + + // Restore historical delivery counters before we accept traffic (#343). + // Best-effort: if Redis is unreachable the server still starts with empty + // counters rather than failing startup (#342). + try { + await webhookDispatcher.hydrateMetrics(); + } catch (err) { + logger.warn( + "Could not load persisted webhook delivery metrics; starting with empty counters", + { error: err.message }, + ); + } + try { await warmCache(config.watchedAssets); } catch (err) { diff --git a/src/indexer/eventPoller.js b/src/indexer/eventPoller.js index 48cd625..d1de264 100644 --- a/src/indexer/eventPoller.js +++ b/src/indexer/eventPoller.js @@ -228,7 +228,29 @@ class EventPoller { ? Math.max(previousLedger ?? 0, ...eventLedgers) : Math.max(response.latestLedger || previousLedger || 0, ...eventLedgers); - await this.store.setLastLedger(latestIndexedLedger); + // Advance the cursor only if it still holds the value this poll read + // (issue #341). A plain set here is a read-modify-write across two + // round trips: an overlapping poller would have its own progress + // silently overwritten, and a stale/lagging answer could even move the + // cursor backwards. The events above are already saved either way — + // saves are idempotent — so refusing the write costs nothing but the + // duplicate work the other poller already did. + const cursorAdvanced = await this.store.advanceLastLedger( + previousLedger, + latestIndexedLedger, + ); + if (!cursorAdvanced) { + this.metrics.pollsSkipped += 1; + this.logger.warn( + 'Ledger cursor changed concurrently; leaving the other writer in place', + { + expected_ledger: previousLedger, + attempted_ledger: latestIndexedLedger, + }, + ); + return { skipped: true, reason: 'concurrent cursor advance' }; + } + this.lastIndexedLedger = latestIndexedLedger; this.metrics.eventsIndexed += parsedEvents.length; diff --git a/src/indexer/eventStore.js b/src/indexer/eventStore.js index f2a40a2..cfe0381 100644 --- a/src/indexer/eventStore.js +++ b/src/indexer/eventStore.js @@ -4,6 +4,33 @@ const EVENT_IDS_KEY = "indexer:contract_events:ids"; const LAST_LEDGER_KEY = "indexer:last_ledger"; const AIRDROP_IDS_KEY = "indexer:airdrops:ids"; +/** + * CAS + monotonic advance for the cursor key. ARGV[1] is the value the + * caller observed ('' = key absent), ARGV[2] the value it wants to store. + * Returns 1 only if both conditions hold and the write happened. + */ +const ADVANCE_LAST_LEDGER_LUA = ` +local cursor = redis.call('GET', KEYS[1]) +local expected = ARGV[1] +local nextValue = ARGV[2] + +if expected == '' then + if cursor then return 0 end +elseif not cursor or cursor ~= expected then + return 0 +end + +local nextNumber = tonumber(nextValue) +if not nextNumber then return 0 end +if cursor then + local cursorNumber = tonumber(cursor) + if cursorNumber and nextNumber < cursorNumber then return 0 end +end + +redis.call('SET', KEYS[1], nextValue) +return 1 +`; + function eventKey(id) { return `indexer:contract_event:${id}`; } @@ -48,6 +75,46 @@ async function setLastLedger(ledger) { await cache.set(LAST_LEDGER_KEY, Number(ledger)); } +/** + * Compare-and-set advance of the indexer cursor (issue #341). + * + * Reading `indexer:last_ledger` and writing it back are two separate Redis + * round trips, so two overlapping pollers can both read the same cursor and + * both write their own answer — the slower one silently rewinds the indexer + * (or, on a lagging response, a plain SET can move the cursor backwards even + * single-threaded). Both hazards are handled here in one Lua script, which + * Redis runs to completion with no interleaving: + * + * - ARGV[1] is the value the caller read ('' = "the cursor was unset"). + * If the stored cursor no longer matches it, nothing is written and the + * caller gets 0 (false). + * - The cursor is only ever moved forwards; a lower next value is refused + * (0) rather than applied. + * + * Returns true when this caller's write landed. + * + * @param {number|null} expectedLedger ledger the caller last read (null = unset) + * @param {number} nextLedger ledger to advance to + * @returns {Promise} whether the cursor was advanced + */ +async function advanceLastLedger(expectedLedger, nextLedger) { + const redis = cache.getClient(); + if (typeof redis.advanceLastLedger !== "function") { + redis.defineCommand("advanceLastLedger", { + numberOfKeys: 1, + lua: ADVANCE_LAST_LEDGER_LUA, + }); + } + + const expected = + expectedLedger === null || expectedLedger === undefined + ? "" + : String(Number(expectedLedger)); + const next = String(Number(nextLedger)); + const applied = await redis.advanceLastLedger(LAST_LEDGER_KEY, expected, next); + return Number(applied) === 1; +} + async function upsertAirdrop(event) { const airdropId = getAirdropId(event); if (!airdropId) return; @@ -298,6 +365,7 @@ async function getStats() { } module.exports = { + advanceLastLedger, getAirdropRecipients, getAirdropStatus, getLastLedger, diff --git a/src/middleware/cors.js b/src/middleware/cors.js index 692734c..4597dd3 100644 --- a/src/middleware/cors.js +++ b/src/middleware/cors.js @@ -1,4 +1,5 @@ const cors = require('cors'); +const config = require('../config'); const AppError = require('../errors/AppError'); function buildCorsMiddleware(allowedOrigins) { @@ -15,7 +16,11 @@ function buildCorsMiddleware(allowedOrigins) { credentials: true, methods: ['GET', 'POST', 'DELETE', 'PATCH', 'OPTIONS'], allowedHeaders: ['Content-Type', 'Authorization', 'X-Requested-With'], - maxAge: 86400, + // Issue #344: preflight (OPTIONS) responses carry Access-Control-Max-Age + // so browsers can cache the result instead of re-sending a preflight for + // every cross-origin POST. Window defaults to 24h and is tunable via + // CORS_MAX_AGE_SECONDS. + maxAge: config.corsMaxAgeSeconds, }); } diff --git a/src/services/priceOracle.js b/src/services/priceOracle.js index 9cc172f..8f5242a 100644 --- a/src/services/priceOracle.js +++ b/src/services/priceOracle.js @@ -356,17 +356,23 @@ async function fetchFreshPrice(assetCode, issuer = null, redisUnavailable = fals return existing.promise; } - // Wrap the promise to handle rejections cleanly: on rejection, remove from - // inFlight and re-throw so concurrent callers see the same error (#286). - const promise = doFetchFreshPrice(assetCode, issuer, redisUnavailable) - .catch((err) => { - inFlight.delete(key); - throw err; - }) - .then((result) => { - inFlight.delete(key); - return result; - }); + // Each coalesced caller is handed its own wrapper promise (#417) so a + // rejection of the underlying fetch is forwarded into every waiter + // without the waiters sharing one rejection object. The single-flight + // entry is still removed exactly once when the underlying fetch settles + // (#286), whether it resolves or rejects. + let resolve; + let reject; + const wrapperPromise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + + const actualPromise = doFetchFreshPrice(assetCode, issuer, redisUnavailable); + + actualPromise + .then(resolve, reject) + .finally(() => inFlight.delete(key)); inFlight.set(key, { promise: wrapperPromise, actual: actualPromise }); return wrapperPromise; @@ -409,6 +415,10 @@ async function refreshAllCachedPrices() { const refreshPromises = keys .map(async (key) => { + // Second line of defence behind the MATCH pattern above: a history key + // must never be refreshed as if it were an asset called "history". + if (key.startsWith(HISTORY_PREFIX)) return; + const suffix = key.replace(CACHE_PREFIX, ''); const parts = suffix.split(':'); const assetCode = parts[0]; diff --git a/src/services/webhookDispatcher.js b/src/services/webhookDispatcher.js index f4c8e49..f783414 100644 --- a/src/services/webhookDispatcher.js +++ b/src/services/webhookDispatcher.js @@ -16,7 +16,16 @@ const USER_AGENT = "SmartDrop-Webhooks/1.0"; const WEBHOOK_CACHE_TTL_MS = 60_000; const webhookCache = new Map(); -// ── Delivery metrics (in-memory, reset on process restart) ────────────── +// ── Delivery metrics (#343) ───────────────────────────────────────────── +// Counters live in Redis (hash-per-webhook plus one aggregate hash) so a +// process restart no longer zeros historical delivery success rates. The +// in-memory copy below is the read model getMetrics() answers from: it is +// loaded once at startup via hydrateMetrics() and then written through on +// every completed delivery. Redis is the durable copy; if it is unavailable +// the write is logged and dropped rather than failing the delivery. +const METRICS_AGGREGATE_KEY = 'webhook:metrics:aggregate'; +const METRICS_WEBHOOK_KEY_PREFIX = 'webhook:metrics:webhook:'; + const metrics = { _deliveries: new Map(), // webhook_id → { total, success, failed, totalAttempts, totalLatencyMs } _inFlight: new Set(), // delivery IDs currently being attempted @@ -42,6 +51,109 @@ function _ensureWebhookMetrics(webhookId) { return metrics._deliveries.get(webhookId); } +function _fromRedisFields(raw) { + const num = (value) => { + const parsed = Number(value); + return Number.isFinite(parsed) ? parsed : 0; + }; + return { + total: num(raw.total), + success: num(raw.success), + failed: num(raw.failed), + totalAttempts: num(raw.total_attempts), + totalLatencyMs: num(raw.total_latency_ms), + }; +} + +/** + * Best-effort write-through of one completed delivery's contribution to the + * durable counters. Fire-and-forget on purpose: a metrics write must never + * be able to fail or delay the delivery it is describing. + */ +function _persistDeliveryMetrics(webhookId, { success, attempts, latencyMs }) { + try { + const redis = cache.getClient(); + const pipeline = redis.pipeline(); + const deltas = [ + ['total', 1], + [success ? 'success' : 'failed', 1], + ['total_attempts', attempts], + ['total_latency_ms', latencyMs], + ]; + for (const [field, by] of deltas) { + pipeline.hincrby(METRICS_AGGREGATE_KEY, field, by); + pipeline.hincrby(`${METRICS_WEBHOOK_KEY_PREFIX}${webhookId}`, field, by); + } + pipeline + .exec() + .then((results) => { + // ioredis resolves exec() with per-command [err, result] pairs; a + // rejected command (WRONGTYPE, OOM, …) would otherwise be invisible. + const failed = (results || []).find(([err]) => err); + if (failed) throw failed[0]; + }) + .catch((err) => { + logger.warn('Failed to persist webhook delivery metrics', { + webhook_id: webhookId, + error: err.message, + }); + }); + } catch (err) { + logger.warn('Failed to persist webhook delivery metrics', { + webhook_id: webhookId, + error: err.message, + }); + } +} + +async function _scanMetricKeys(pattern) { + const redis = cache.getClient(); + const keys = []; + let cursor = '0'; + do { + const [nextCursor, batch] = await redis.scan(cursor, 'MATCH', pattern, 'COUNT', 100); + cursor = String(nextCursor); + keys.push(...batch); + } while (cursor !== '0'); + return keys; +} + +/** + * Loads the durable counters back into the in-memory read model (#343). + * + * Called once from startServer(), before the server begins accepting + * requests, so no delivery can be counted twice (the overwrite below would + * otherwise discard a locally-recorded delta). Rejects if Redis cannot be + * read — callers are expected to treat that as "start with empty counters", + * not as a fatal startup error (#342). + */ +async function hydrateMetrics() { + const redis = cache.getClient(); + const [aggregateRaw, webhookKeys] = await Promise.all([ + redis.hgetall(METRICS_AGGREGATE_KEY), + _scanMetricKeys(`${METRICS_WEBHOOK_KEY_PREFIX}*`), + ]); + + metrics._aggregate = _fromRedisFields(aggregateRaw || {}); + metrics._deliveries.clear(); + + const perWebhook = await Promise.all( + webhookKeys.map(async (key) => [ + key.slice(METRICS_WEBHOOK_KEY_PREFIX.length), + await redis.hgetall(key), + ]), + ); + for (const [webhookId, raw] of perWebhook) { + metrics._deliveries.set(webhookId, _fromRedisFields(raw || {})); + } + + logger.info('Webhook delivery metrics loaded from Redis', { + aggregate_total: metrics._aggregate.total, + webhooks_tracked: metrics._deliveries.size, + }); + return metrics._aggregate.total; +} + function recordDeliveryStart(deliveryId, webhookId) { metrics._inFlight.add(deliveryId); _ensureWebhookMetrics(webhookId); @@ -70,6 +182,8 @@ function recordDeliveryEnd( wm.failed += 1; ag.failed += 1; } + + _persistDeliveryMetrics(webhookId, { success, attempts, latencyMs }); } function getMetrics() { @@ -685,4 +799,7 @@ module.exports = { shouldRetry, getMetrics, getInFlightCount, + hydrateMetrics, + recordDeliveryStart, + recordDeliveryEnd, }; diff --git a/src/startup/cacheWarm.js b/src/startup/cacheWarm.js index 1bbbcdf..cd15813 100644 --- a/src/startup/cacheWarm.js +++ b/src/startup/cacheWarm.js @@ -113,7 +113,13 @@ async function warmCache( }); const summary = await Promise.race([warming, timeout]); - if (!summary.timedOut) clearTimeout(timeoutId); + if (!summary.timedOut) { + clearTimeout(timeoutId); + // One completion line per startup (mirroring the timeout path's warn) + // so an operator can see how much of the cache actually warmed without + // diffing per-asset fetch logs. + log.info('Cache warm complete', summary); + } return summary; } diff --git a/src/ws/priceWebSocket.js b/src/ws/priceWebSocket.js index 68faff3..1ec77f6 100644 --- a/src/ws/priceWebSocket.js +++ b/src/ws/priceWebSocket.js @@ -65,6 +65,22 @@ function attach(httpServer) { return wss; } +/** + * Health snapshot of the WebSocket server for the /health endpoint. + * Reports unhealthy (rather than throwing) when the server has not been + * attached yet, so /health can still answer during early startup. + */ +function getHealth() { + if (!wss) { + return { healthy: false, connections: 0, error: 'WebSocket server not initialized' }; + } + return { + healthy: true, + connections: wss.clients ? wss.clients.size : 0, + error: null, + }; +} + /** * Gracefully close all WebSocket connections and stop the heartbeat. * Call this during process shutdown to avoid abrupt connection drops. @@ -82,4 +98,4 @@ async function shutdown(wss, drainTimeoutMs = 5000) { }); } -module.exports = { attach, shutdown }; +module.exports = { attach, getHealth, shutdown }; diff --git a/test/cors.test.js b/test/cors.test.js index 58ffa85..3c9c2da 100644 --- a/test/cors.test.js +++ b/test/cors.test.js @@ -7,6 +7,9 @@ const { errorHandler } = require('../src/middleware/errorHandler'); const ALLOWED = ['http://localhost:3000', 'https://app.smartdrop.io']; +// Must match config.js's CORS_MAX_AGE_SECONDS default. +const DEFAULT_MAX_AGE_SECONDS = 86400; + function buildApp(allowedOrigins) { const app = express(); app.use(buildCorsMiddleware(allowedOrigins)); @@ -107,6 +110,48 @@ describe('CORS no-origin requests (server-to-server, curl)', () => { }); }); +describe('CORS preflight caching (#344)', () => { + let app; + beforeAll(() => { app = buildApp(ALLOWED); }); + + test('a cross-origin POST preflight advertises a cache window', async () => { + const res = await request(app) + .options('/test') + .set('Origin', 'http://localhost:3000') + .set('Access-Control-Request-Method', 'POST') + .set('Access-Control-Request-Headers', 'content-type,authorization'); + + expect(res.status).toBe(204); + // Without Access-Control-Max-Age the browser re-sends this preflight for + // every cross-origin POST; with it, the response is reusable for the + // advertised number of seconds. + expect(res.headers['access-control-max-age']).toBe(String(DEFAULT_MAX_AGE_SECONDS)); + expect(res.headers['access-control-allow-methods']).toMatch(/POST/); + }); + + test('the cache window is configurable via CORS_MAX_AGE_SECONDS', async () => { + const original = process.env.CORS_MAX_AGE_SECONDS; + process.env.CORS_MAX_AGE_SECONDS = '600'; + jest.resetModules(); + const buildCors = require('../src/middleware/cors'); + const configurableApp = express(); + configurableApp.use(buildCors(ALLOWED)); + configurableApp.get('/test', (req, res) => res.json({ ok: true })); + + try { + const res = await request(configurableApp) + .options('/test') + .set('Origin', 'http://localhost:3000') + .set('Access-Control-Request-Method', 'POST'); + expect(res.headers['access-control-max-age']).toBe('600'); + } finally { + if (original === undefined) delete process.env.CORS_MAX_AGE_SECONDS; + else process.env.CORS_MAX_AGE_SECONDS = original; + jest.resetModules(); + } + }); +}); + describe('CORS_ALLOWED_ORIGINS config parsing', () => { test('dev default allows localhost:3000 and localhost:3001', () => { const original = process.env.CORS_ALLOWED_ORIGINS; diff --git a/test/eventPoller.test.js b/test/eventPoller.test.js index e3e2807..221dda9 100644 --- a/test/eventPoller.test.js +++ b/test/eventPoller.test.js @@ -55,7 +55,7 @@ describe('EventPoller', () => { getLastLedger: jest.fn(async () => null), saveEvent: jest.fn(async () => {}), saveEvents: jest.fn(async () => {}), - setLastLedger: jest.fn(async () => {}), + advanceLastLedger: jest.fn(async () => true), }; const logger = { info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn() }; @@ -80,7 +80,7 @@ describe('EventPoller', () => { event_name: 'airdrop_created', data: expect.objectContaining({ airdrop_id: 'drop-1', total_amount: '1000' }), })]); - expect(store.setLastLedger).toHaveBeenCalledWith(25); + expect(store.advanceLastLedger).toHaveBeenCalledWith(null, 25); expect(result).toMatchObject({ indexed_events: 1, latest_ledger: 25 }); expect(poller.getStatus()).toMatchObject({ latest_ledger: 25, last_error: null }); }); @@ -93,7 +93,7 @@ describe('EventPoller', () => { getLastLedger: jest.fn(async () => 19), saveEvent: jest.fn(async () => {}), saveEvents: jest.fn(async () => {}), - setLastLedger: jest.fn(async () => {}), + advanceLastLedger: jest.fn(async () => true), }; const poller = new EventPoller({ @@ -107,6 +107,9 @@ describe('EventPoller', () => { await poller.pollOnce(); expect(server.getEvents.mock.calls[0][0].startLedger).toBe(20); + // The cursor advance is conditional on the ledger that was actually read + // (#341), not a blind overwrite. + expect(store.advanceLastLedger).toHaveBeenCalledWith(19, 25); }); test('skips polling when no contract id is configured', async () => { @@ -132,7 +135,7 @@ describe('EventPoller', () => { getLastLedger: jest.fn(async () => null), saveEvent: jest.fn(async () => {}), saveEvents: jest.fn(async () => {}), - setLastLedger: jest.fn(async () => {}), + advanceLastLedger: jest.fn(async () => true), }; const logger = { info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn() }; @@ -150,7 +153,7 @@ describe('EventPoller', () => { // Not 500 (response.latestLedger) — that would permanently skip // whatever exists between ledger 24 and the tip. - expect(store.setLastLedger).toHaveBeenCalledWith(24); + expect(store.advanceLastLedger).toHaveBeenCalledWith(null, 24); expect(result).toMatchObject({ truncated: true, indexed_events: pollLimit }); expect(logger.warn).toHaveBeenCalledWith( 'SmartDrop event poll truncated by pollLimit; more events pending next cycle', @@ -179,8 +182,14 @@ describe('EventPoller', () => { getLastLedger: jest.fn(async () => lastLedger), saveEvent: jest.fn(async () => {}), saveEvents: jest.fn(async () => {}), - setLastLedger: jest.fn(async (ledger) => { - lastLedger = ledger; + advanceLastLedger: jest.fn(async (expected, next) => { + // Same contract as eventStore.advanceLastLedger: refuse the write + // if the cursor no longer matches what the caller read, and never + // move it backwards. + if (expected !== lastLedger) return false; + if (lastLedger !== null && next < lastLedger) return false; + lastLedger = next; + return true; }), }; @@ -218,7 +227,7 @@ describe('EventPoller', () => { getLastLedger: jest.fn(async () => null), saveEvent: jest.fn(async () => {}), saveEvents: jest.fn(async () => {}), - setLastLedger: jest.fn(async () => {}), + advanceLastLedger: jest.fn(async () => true), }; const logger = { info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn() }; @@ -234,9 +243,95 @@ describe('EventPoller', () => { const result = await poller.pollOnce(); - expect(store.setLastLedger).toHaveBeenCalledWith(500); + expect(store.advanceLastLedger).toHaveBeenCalledWith(null, 500); expect(result.truncated).toBe(false); expect(logger.warn).not.toHaveBeenCalled(); }); }); + + describe('concurrent cursor advance (#341)', () => { + function buildPoller({ store, server, logger }) { + return new EventPoller({ + enabled: true, + contractId: 'CCONTRACT', + startLedger: 10, + pollLimit: 5, + server, + store, + logger, + }); + } + + test('saves the batch before touching the cursor', async () => { + const server = { + getEvents: jest.fn(async () => ({ + latestLedger: 25, + events: [contractEvent()], + })), + }; + const store = { + getLastLedger: jest.fn(async () => null), + saveEvents: jest.fn(async () => {}), + advanceLastLedger: jest.fn(async () => true), + }; + + const poller = buildPoller({ + store, + server, + logger: { info: jest.fn(), warn: jest.fn(), error: jest.fn(), debug: jest.fn() }, + }); + await poller.pollOnce(); + + // If the cursor moved first and the process died mid-save, those + // events would be gone for good — hence save-then-advance. + expect(store.saveEvents.mock.invocationCallOrder[0]).toBeLessThan( + store.advanceLastLedger.mock.invocationCallOrder[0], + ); + }); + + test('refusing the write still persists events and reports a skip', async () => { + const server = { + getEvents: jest.fn(async () => ({ + latestLedger: 150, + events: [contractEvent({ ledger: 105, id: 'evt-105' })], + })), + }; + const store = { + getLastLedger: jest.fn(async () => 100), + saveEvents: jest.fn(async () => {}), + // Another instance moved the cursor between our read and our write. + advanceLastLedger: jest.fn(async () => false), + }; + const logger = { + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + debug: jest.fn(), + }; + + const poller = buildPoller({ store, server, logger }); + const result = await poller.pollOnce(); + + expect(store.advanceLastLedger).toHaveBeenCalledWith(100, 150); + expect(store.saveEvents).toHaveBeenCalledTimes(1); + expect(result).toMatchObject({ + skipped: true, + reason: 'concurrent cursor advance', + }); + expect(logger.warn).toHaveBeenCalledWith( + 'Ledger cursor changed concurrently; leaving the other writer in place', + expect.objectContaining({ + expected_ledger: 100, + attempted_ledger: 150, + }), + ); + expect(poller.getMetrics()).toMatchObject({ + polls_attempted: 1, + polls_skipped: 1, + polls_succeeded: 0, + }); + // The loser of the race must not claim the progress it did not make. + expect(poller.lastIndexedLedger).toBeNull(); + }); + }); }); diff --git a/test/health.test.js b/test/health.test.js index 8c64339..f713485 100644 --- a/test/health.test.js +++ b/test/health.test.js @@ -42,6 +42,9 @@ jest.mock('../src/jobs/webhookRetryWorker', () => ({ jest.mock('../src/ws/priceWebSocket', () => ({ attach: jest.fn(), + // /health reads this to report the WS server's state (mirrors the real + // getHealth() for a server that has been attached). + getHealth: jest.fn(() => ({ healthy: true, connections: 0 })), })); // --------------------------------------------------------------------------- @@ -170,6 +173,15 @@ describe('GET /health – status computation', () => { tick: jest.fn(), getHealth: () => ({ healthy: true, lastSuccessAt: Date.now(), lastError: null, stalled: false }), })); + // The real airdrop expiry job exposes no getHealth(), so the wrapper can + // never report it as healthy — mock it here, where the point of the test + // is "every job healthy ⇒ ok". + jest.mock('../src/jobs/airdropExpiry', () => ({ + start: jest.fn(), + stop: jest.fn(), + tick: jest.fn(), + getHealth: () => ({ healthy: true, lastSuccessAt: Date.now(), lastError: null, stalled: false }), + })); const app = loadApp(); const res = await request(app).get('/health'); diff --git a/test/helpers/cacheMock.js b/test/helpers/cacheMock.js index 7bb8db9..60c1b62 100644 --- a/test/helpers/cacheMock.js +++ b/test/helpers/cacheMock.js @@ -195,6 +195,33 @@ function createCacheMock() { if (!h) return {}; return Object.fromEntries(h.entries()); }), + // Issue #343 (webhookDispatcher metrics persistence): HINCRBY is the + // write-through primitive for the durable delivery counters. + hincrby: jest.fn(async (key, field, by) => { + const h = getHash(key); + const next = (Number(h.get(field)) || 0) + Number(by); + h.set(field, String(next)); + return next; + }), + // Issue #343: SCAN over the whole keyspace (all key types), matching a + // Redis glob pattern. Returns the lot in one pass (cursor '0'), which is + // all callers iterating until cursor === '0' need. + scan: jest.fn(async (cursor, ...args) => { + const matchIdx = args.indexOf('MATCH'); + const pattern = matchIdx === -1 ? '*' : String(args[matchIdx + 1]); + const regex = new RegExp( + `^${pattern.replace(/[.+^${}()|[\]\\]/g, '\\$&').replace(/\*/g, '.*').replace(/\?/g, '.')}$`, + ); + const allKeys = new Set([ + ...rawStore.keys(), + ...hashes.keys(), + ...sets.keys(), + ...zsets.keys(), + ...lists.keys(), + ]); + const matches = [...allKeys].filter((key) => regex.test(key)); + return ['0', matches]; + }), pexpire: jest.fn(async (key, ms) => { const entry = getLive(key); if (!entry) return 0; @@ -229,6 +256,30 @@ function createCacheMock() { }); return; } + if (name === 'advanceLastLedger') { + // Mirrors ADVANCE_LAST_LEDGER_LUA (issue #341): compare-and-set + // (expected === '' means "key must be absent") plus monotonicity — + // never store a lower ledger than the one already stored. Reads and + // the write below are uninterrupted, matching Redis running the + // script to completion. + redis.advanceLastLedger = jest.fn(async (key, expected, nextValue) => { + const cursor = getLive(key) ? getLive(key).value : null; + if (expected === '') { + if (cursor !== null) return 0; + } else if (cursor === null || cursor !== expected) { + return 0; + } + const nextNumber = Number(nextValue); + if (!Number.isFinite(nextNumber)) return 0; + if (cursor !== null) { + const cursorNumber = Number(cursor); + if (Number.isFinite(cursorNumber) && nextNumber < cursorNumber) return 0; + } + rawStore.set(key, { value: String(nextValue), expiresAt: null }); + return 1; + }); + return; + } if (name === 'renewLease') { // Mirrors RENEW_LUA: renew only if we still hold the lease. redis.renewLease = jest.fn(async (key, expectedValue, ttlMs) => { diff --git a/test/indexerBackoff.test.js b/test/indexerBackoff.test.js index 1f69c27..39e1935 100644 --- a/test/indexerBackoff.test.js +++ b/test/indexerBackoff.test.js @@ -13,8 +13,9 @@ const mockLogger = { function buildStore(overrides = {}) { return { getLastLedger: jest.fn(async () => 100), - setLastLedger: jest.fn(async () => {}), + advanceLastLedger: jest.fn(async () => true), saveEvent: jest.fn(async () => {}), + saveEvents: jest.fn(async () => {}), ...overrides, }; } diff --git a/test/indexerStore.test.js b/test/indexerStore.test.js index 5bea72c..7d5ef36 100644 --- a/test/indexerStore.test.js +++ b/test/indexerStore.test.js @@ -79,6 +79,30 @@ const mockRedis = { }), multi: jest.fn(() => makeChain()), pipeline: jest.fn(() => makeChain()), + // Registers the cursor CAS (issue #341) the same way eventStore does it + // against a real ioredis client; mirrors ADVANCE_LAST_LEDGER_LUA. + defineCommand: jest.fn((name, { lua } = {}) => { + if (name === "advanceLastLedger") { + mockRedis.advanceLastLedger = jest.fn(async (key, expected, nextValue) => { + const cursor = mockStore.has(key) ? mockStore.get(key) : null; + if (expected === "") { + if (cursor !== null) return 0; + } else if (cursor === null || cursor !== expected) { + return 0; + } + const nextNumber = Number(nextValue); + if (!Number.isFinite(nextNumber)) return 0; + if (cursor !== null) { + const cursorNumber = Number(cursor); + if (Number.isFinite(cursorNumber) && nextNumber < cursorNumber) return 0; + } + mockStore.set(key, String(nextValue)); + return 1; + }); + return; + } + throw new Error(`indexerStore test mock: unsupported defineCommand "${name}" (lua: ${typeof lua})`); + }), }; jest.mock("../src/services/cache", () => ({ @@ -371,3 +395,51 @@ describe("indexer event store", () => { ]); }); }); + +describe("ledger cursor compare-and-set (#341)", () => { + const LEDGER_KEY = "indexer:last_ledger"; + + test("advances an unset cursor and then reads it back", async () => { + await expect(eventStore.getLastLedger(null)).resolves.toBeNull(); + + await expect(eventStore.advanceLastLedger(null, 100)).resolves.toBe(true); + await expect(eventStore.getLastLedger(null)).resolves.toBe(100); + }); + + test("advances only when the cursor still holds what the caller read", async () => { + await eventStore.advanceLastLedger(null, 100); + + // Second poller read 100 too, but the first one already moved on. + await expect(eventStore.advanceLastLedger(100, 140)).resolves.toBe(true); + await expect(eventStore.advanceLastLedger(100, 150)).resolves.toBe(false); + await expect(eventStore.getLastLedger(null)).resolves.toBe(140); + }); + + test("refuses a write that would move the cursor backwards", async () => { + await eventStore.advanceLastLedger(null, 200); + + await expect(eventStore.advanceLastLedger(200, 150)).resolves.toBe(false); + await expect(eventStore.getLastLedger(null)).resolves.toBe(200); + // Re-applying the same value is a no-op advance, not a regression. + await expect(eventStore.advanceLastLedger(200, 200)).resolves.toBe(true); + }); + + test("refuses when an unset cursor was claimed in the meantime", async () => { + await eventStore.advanceLastLedger(null, 1); + + await expect(eventStore.advanceLastLedger(null, 2)).resolves.toBe(false); + await expect(eventStore.getLastLedger(null)).resolves.toBe(1); + }); + + test("two overlapping pollers converge on the winner's cursor", async () => { + // Both read the same starting point before either writes. + const readA = await eventStore.getLastLedger(null); + const readB = await eventStore.getLastLedger(null); + + await expect(eventStore.advanceLastLedger(readA, 10)).resolves.toBe(true); + await expect(eventStore.advanceLastLedger(readB, 12)).resolves.toBe(false); + + await expect(eventStore.getLastLedger(null)).resolves.toBe(10); + expect(mockStore.get(LEDGER_KEY)).toBeDefined(); + }); +}); diff --git a/test/priceOracle.test.js b/test/priceOracle.test.js index f9351aa..1cf08be 100644 --- a/test/priceOracle.test.js +++ b/test/priceOracle.test.js @@ -617,7 +617,9 @@ describe('refreshAllCachedPrices', () => { const result = await refreshAllCachedPrices(); - expect(redis.scan).toHaveBeenCalledWith('0', 'MATCH', 'price:*', 'COUNT', 100); + // History keys are excluded by the MATCH pattern itself (price:history:* + // is skipped at the Redis level), not just by client-side filtering. + expect(redis.scan).toHaveBeenCalledWith('0', 'MATCH', 'price:[^h]*', 'COUNT', 100); // Only the non-history key is refreshed. expect(mockStellarFetch).toHaveBeenCalledWith('XLM', null); expect(result).toEqual({ XLM: { price: 0.1, source: 'stellar_dex' } }); diff --git a/test/startupBanner.test.js b/test/startupBanner.test.js index 9937a3e..f142079 100644 --- a/test/startupBanner.test.js +++ b/test/startupBanner.test.js @@ -28,6 +28,7 @@ const { sanitizeUrl, startServer, } = require("../src/index"); +const webhookDispatcher = require("../src/services/webhookDispatcher"); const config = require("../src/config"); const { version: appVersion } = require("../package.json"); @@ -118,4 +119,60 @@ describe("startup banner", () => { ); listen.mockRestore(); }); + + test("reloads persisted webhook metrics before accepting traffic (#343)", async () => { + const mockServer = { close: jest.fn() }; + const listen = jest.spyOn(app, "listen").mockReturnValue(mockServer); + const hydrate = jest + .spyOn(webhookDispatcher, "hydrateMetrics") + .mockResolvedValue(7); + + await expect(startServer()).resolves.toBe(mockServer); + + expect(hydrate).toHaveBeenCalledTimes(1); + expect(listen).toHaveBeenCalledWith(config.port, expect.any(Function)); + // Counters must be restored before a delivery can be counted into them, + // otherwise the reload would overwrite an already-recorded delivery. + expect(hydrate.mock.invocationCallOrder[0]).toBeLessThan( + listen.mock.invocationCallOrder[0], + ); + hydrate.mockRestore(); + listen.mockRestore(); + }); + + test("still starts the server when webhook metrics cannot be loaded (#342)", async () => { + const mockServer = { close: jest.fn() }; + const listen = jest.spyOn(app, "listen").mockReturnValue(mockServer); + const hydrate = jest + .spyOn(webhookDispatcher, "hydrateMetrics") + .mockRejectedValueOnce(new Error("connection refused")); + + await expect(startServer()).resolves.toBe(mockServer); + + expect(listen).toHaveBeenCalledWith(config.port, expect.any(Function)); + expect(mockLogger.warn).toHaveBeenCalledWith( + "Could not load persisted webhook delivery metrics; starting with empty counters", + { error: "connection refused" }, + ); + hydrate.mockRestore(); + listen.mockRestore(); + }); + + test("announces a degraded boot when Redis is down at startup (#342)", async () => { + const mockServer = { close: jest.fn() }; + const listen = jest.spyOn(app, "listen").mockReturnValue(mockServer); + const isConnected = jest + .spyOn(require("../src/services/cache"), "isConnected") + .mockReturnValue(false); + + await expect(startServer()).resolves.toBe(mockServer); + + expect(listen).toHaveBeenCalledWith(config.port, expect.any(Function)); + expect(mockLogger.warn).toHaveBeenCalledWith( + expect.stringContaining("starting in degraded mode"), + expect.objectContaining({ redis_connected: false }), + ); + isConnected.mockRestore(); + listen.mockRestore(); + }); }); diff --git a/test/webhookDispatcher.test.js b/test/webhookDispatcher.test.js index ed52544..5b7c8fb 100644 --- a/test/webhookDispatcher.test.js +++ b/test/webhookDispatcher.test.js @@ -8,6 +8,13 @@ const { reset, zsets } = mockHelper; jest.mock("../src/services/cache", () => mockHelper.cacheMock); jest.mock("../src/config", () => ({ webhookSecretEncryptionKey: "test-webhook-secret-key", + // validation/schemas.js enumerates the API key tiers from config at module + // load, so the config mock has to carry them too. + apiKeyRateLimit: { + windowSeconds: 60, + defaultTier: "free", + tiers: { free: 100, pro: 500, admin: 5000 }, + }, webhooks: { retryBaseMs: 30000, retryFactor: 2, diff --git a/test/webhookMetricsPersistence.test.js b/test/webhookMetricsPersistence.test.js new file mode 100644 index 0000000..a4c654e --- /dev/null +++ b/test/webhookMetricsPersistence.test.js @@ -0,0 +1,179 @@ +"use strict"; + +/** + * Webhook delivery metrics must survive a process restart (issue #343). + * + * Each completed delivery is written through to Redis (one aggregate hash + * plus one hash per webhook) and the counters are reloaded once at startup, + * so the success/failure rates exposed by GET /metrics and /health keep + * meaning "all time" rather than "since this process booted". + */ + +const { createCacheMock } = require("./helpers/cacheMock"); + +const mockHelper = createCacheMock(); +const { reset, redis } = mockHelper; + +jest.mock("../src/services/cache", () => mockHelper.cacheMock); + +const mockLogger = { + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + debug: jest.fn(), +}; +jest.mock("../src/logger", () => mockLogger); + +const dispatcher = require("../src/services/webhookDispatcher"); + +const AGGREGATE_KEY = "webhook:metrics:aggregate"; +const WEBHOOK_KEY = (id) => `webhook:metrics:webhook:${id}`; + +/** Let the fire-and-forget metric pipelines settle before asserting. */ +function flush() { + return new Promise((resolve) => setImmediate(resolve)); +} + +beforeEach(() => { + reset(); + Object.values(mockLogger).forEach((fn) => fn.mockClear()); +}); + +async function record(webhookId, { success, attempts = 1, latencyMs = 25 }) { + const deliveryId = `dlv_${webhookId}_${Math.random()}`; + dispatcher.recordDeliveryStart(deliveryId, webhookId); + dispatcher.recordDeliveryEnd(deliveryId, webhookId, { + success, + attempts, + latencyMs, + }); + await flush(); +} + +describe("durable delivery counters", () => { + test("writes each completed delivery through to Redis", async () => { + await dispatcher.hydrateMetrics(); + + await record("wh_1", { success: true }); + await record("wh_1", { success: false, attempts: 3, latencyMs: 40 }); + + expect(await redis.hgetall(AGGREGATE_KEY)).toEqual({ + total: "2", + success: "1", + failed: "1", + total_attempts: "4", + total_latency_ms: "65", + }); + expect(await redis.hgetall(WEBHOOK_KEY("wh_1"))).toEqual({ + total: "2", + success: "1", + failed: "1", + total_attempts: "4", + total_latency_ms: "65", + }); + }); + + test("keeps counters per webhook so a busy one cannot hide the rest", async () => { + await dispatcher.hydrateMetrics(); + + await record("wh_1", { success: true }); + await record("wh_2", { success: false }); + + expect(await redis.hgetall(WEBHOOK_KEY("wh_1"))).toMatchObject({ + total: "1", + success: "1", + }); + expect(await redis.hgetall(WEBHOOK_KEY("wh_2"))).toMatchObject({ + total: "1", + failed: "1", + }); + }); + + test("a Redis write failure is logged but never breaks the delivery bookkeeping", async () => { + redis.pipeline.mockImplementationOnce(() => { + throw new Error("connection lost"); + }); + + const before = dispatcher.getMetrics().aggregate.total; + expect(() => + dispatcher.recordDeliveryEnd("dlv_x", "wh_1", { + success: true, + attempts: 1, + latencyMs: 10, + }), + ).not.toThrow(); + + expect(dispatcher.getMetrics().aggregate.total).toBe(before + 1); + expect(mockLogger.warn).toHaveBeenCalledWith( + "Failed to persist webhook delivery metrics", + expect.objectContaining({ webhook_id: "wh_1" }), + ); + }); +}); + +describe("hydrateMetrics", () => { + test("restores historical counters and derived rates after a restart", async () => { + // Simulate counters left behind by a previous process. + await redis.hset(AGGREGATE_KEY, "total", 10); + await redis.hset(AGGREGATE_KEY, "success", 8); + await redis.hset(AGGREGATE_KEY, "failed", 2); + await redis.hset(AGGREGATE_KEY, "total_attempts", 14); + await redis.hset(AGGREGATE_KEY, "total_latency_ms", 250); + await redis.hset(WEBHOOK_KEY("wh_42"), "total", 4); + await redis.hset(WEBHOOK_KEY("wh_42"), "success", 4); + await redis.hset(WEBHOOK_KEY("wh_42"), "failed", 0); + await redis.hset(WEBHOOK_KEY("wh_42"), "total_attempts", 4); + await redis.hset(WEBHOOK_KEY("wh_42"), "total_latency_ms", 100); + + const restored = await dispatcher.hydrateMetrics(); + expect(restored).toBe(10); + + const snapshot = dispatcher.getMetrics(); + expect(snapshot.aggregate).toMatchObject({ + total: 10, + success: 8, + failed: 2, + success_rate: 0.8, + retry_rate: 0.4, + avg_latency_ms: 25, + }); + expect(snapshot.per_webhook.wh_42).toMatchObject({ + total: 4, + success: 4, + success_rate: 1, + avg_latency_ms: 25, + }); + expect(mockLogger.info).toHaveBeenCalledWith( + "Webhook delivery metrics loaded from Redis", + expect.objectContaining({ aggregate_total: 10, webhooks_tracked: 1 }), + ); + }); + + test("replaces the whole read model rather than double-counting", async () => { + await redis.hset(AGGREGATE_KEY, "total", 5); + await redis.hset(AGGREGATE_KEY, "success", 5); + await redis.hset(AGGREGATE_KEY, "failed", 0); + + await dispatcher.hydrateMetrics(); + await dispatcher.hydrateMetrics(); + + expect(dispatcher.getMetrics().aggregate.total).toBe(5); + }); + + test("reports empty counters when Redis holds none yet", async () => { + expect(await dispatcher.hydrateMetrics()).toBe(0); + expect(dispatcher.getMetrics()).toMatchObject({ + in_flight: 0, + aggregate: { total: 0, success: 0, failed: 0, success_rate: null }, + }); + }); + + test("rejects when Redis is unreachable so the caller can degrade (#342)", async () => { + redis.hgetall.mockRejectedValueOnce(new Error("connection refused")); + + await expect(dispatcher.hydrateMetrics()).rejects.toThrow( + "connection refused", + ); + expect(dispatcher.getMetrics().aggregate.total).toBe(0); + }); +});