From d39e8a8861dc82a685bd652e300098ba3d0a6e8f Mon Sep 17 00:00:00 2001 From: davidjhonsonmate Date: Sat, 26 Sep 2026 02:10:57 +0100 Subject: [PATCH] fix: add WebSocket health check, fix inFlight promise sharing, fix leader election Redis error - priceWebSocket.js: add getHealth() function to expose WS server status (#418) - index.js: add WebSocket health check to /health endpoint (#418) - priceOracle.js: use wrapper promise pattern to prevent concurrent callers from sharing rejections (#417) - leaderElection.js: return false on Redis error instead of stale leader flag (#416) --- src/index.js | 10 ++++++++-- src/services/leaderElection.js | 10 +++++++--- src/services/priceOracle.js | 23 +++++++++++++++++++---- src/ws/priceWebSocket.js | 21 +++++++++++++++++++-- 4 files changed, 53 insertions(+), 11 deletions(-) diff --git a/src/index.js b/src/index.js index 4a3c485..7ed5310 100644 --- a/src/index.js +++ b/src/index.js @@ -79,6 +79,7 @@ app.get('/health', async (req, res) => { const webhookWorkerHealth = wrappedWebhookRetryWorker.getHealth(); const airdropExpiryHealth = wrappedAirdropExpiryJob.getHealth(); const database = await checkDatabase(); + const wsHealth = priceWebSocket.getHealth(); // Compute overall status: // unhealthy – Redis is down, or a job is stalled past its grace period @@ -90,11 +91,11 @@ app.get('/health', async (req, res) => { // work. The health check distinguishes "not leader" from "stalled" via the // `leader` field. let status = 'ok'; - if (!redisConnected || !priceRefreshHealth.healthy || !webhookWorkerHealth.healthy || database.status === 'error') { + if (!redisConnected || !priceRefreshHealth.healthy || !webhookWorkerHealth.healthy || database.status === 'error' || !wsHealth.healthy) { const jobsDegraded = (!priceRefreshHealth.healthy && !priceRefreshHealth.stalled) || (!webhookWorkerHealth.healthy && !webhookWorkerHealth.stalled); - status = (!redisConnected || priceRefreshHealth.stalled || webhookWorkerHealth.stalled || database.status === 'error') + status = (!redisConnected || priceRefreshHealth.stalled || webhookWorkerHealth.stalled || database.status === 'error' || !wsHealth.healthy) ? 'unhealthy' : jobsDegraded ? 'degraded' : 'unhealthy'; } @@ -108,6 +109,11 @@ app.get('/health', async (req, res) => { redis: { connected: redisConnected, }, + websocket: { + healthy: wsHealth.healthy, + connections: wsHealth.connections, + error: wsHealth.error, + }, jobs: { price_refresh: { healthy: priceRefreshHealth.healthy, diff --git a/src/services/leaderElection.js b/src/services/leaderElection.js index d288f7b..8d72351 100644 --- a/src/services/leaderElection.js +++ b/src/services/leaderElection.js @@ -151,9 +151,13 @@ function createLeaderElection(jobName, opts = {}) { lockKey, error: err.message, }); - // Don't clear leader flag on transient Redis errors — the lease may - // still be valid. We'll retry on the next renewal cycle. - return leader; + // Return false on Redis errors to prevent the instance from + // continuing to act as leader when it cannot verify its lease. + // The instance will attempt to re-acquire leadership on the next cycle. + leader = false; + acquiredAt = null; + lastRenewedAt = null; + return false; } } diff --git a/src/services/priceOracle.js b/src/services/priceOracle.js index 52cfd97..070ba8e 100644 --- a/src/services/priceOracle.js +++ b/src/services/priceOracle.js @@ -288,6 +288,11 @@ async function doFetchFreshPrice(assetCode, issuer = null, redisUnavailable = fa * Single-flight wrapper around doFetchFreshPrice. Concurrent calls for the * same assetCode:issuer pair while a fetch is already in-flight will await * the same promise instead of each independently hitting all upstream sources. + * + * Uses a wrapper promise pattern to ensure concurrent callers receive + * independent promise instances. This prevents the issue where a rejection + * from the original promise would propagate to all waiting callers, causing + * them to share the same rejection error. */ async function fetchFreshPrice(assetCode, issuer = null, redisUnavailable = false) { const normalisedIssuer = issuer || null; @@ -301,14 +306,24 @@ async function fetchFreshPrice(assetCode, issuer = null, redisUnavailable = fals issuer: normalisedIssuer, coalesced: coalescedCount, }); - return existing; + return existing.promise; } - const promise = doFetchFreshPrice(assetCode, issuer, redisUnavailable) + let resolve, reject; + const wrapperPromise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + + const actualPromise = doFetchFreshPrice(assetCode, issuer, redisUnavailable); + + // Forward resolution/rejection to the wrapper, then clean up + actualPromise + .then(resolve, reject) .finally(() => inFlight.delete(key)); - inFlight.set(key, promise); - return promise; + inFlight.set(key, { promise: wrapperPromise, actual: actualPromise }); + return wrapperPromise; } async function refreshAllCachedPrices() { diff --git a/src/ws/priceWebSocket.js b/src/ws/priceWebSocket.js index 5095054..577cbdd 100644 --- a/src/ws/priceWebSocket.js +++ b/src/ws/priceWebSocket.js @@ -5,6 +5,8 @@ const logger = require('../logger'); const apiKeys = require('../services/apiKeys'); const subscriptionManager = require('./PriceSubscriptionManager'); +let wss = null; + function extractBearerToken(header) { if (!header || typeof header !== 'string') return null; const match = header.match(/^Bearer\s+(.+)$/i); @@ -37,7 +39,7 @@ function authenticateUpgrade(info, callback) { * Clients connect at ws:///ws */ function attach(httpServer) { - const wss = new WebSocketServer({ + wss = new WebSocketServer({ server: httpServer, path: '/ws', verifyClient: authenticateUpgrade, @@ -58,4 +60,19 @@ function attach(httpServer) { return wss; } -module.exports = { attach }; +/** + * Returns health information about the WebSocket server. + * Used by the /health endpoint to check if the WS server is accepting connections. + */ +function getHealth() { + if (!wss) { + return { healthy: false, connections: 0, error: 'WebSocket server not initialized' }; + } + return { + healthy: wss.readyState === undefined || wss.readyState === 0, // 0 = CONNECTING (normal for server) + connections: wss.clients.size, + error: null, + }; +} + +module.exports = { attach, getHealth };