diff --git a/src/index.js b/src/index.js index fc62d9a..8cb7f27 100644 --- a/src/index.js +++ b/src/index.js @@ -181,6 +181,11 @@ app.get("/health", healthRateLimit, async (req, res) => { draining: subscriptionManager.isDraining, drain_stats: subscriptionManager.drainStats, }, + 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 ac3587a..3468c4d 100644 --- a/src/services/leaderElection.js +++ b/src/services/leaderElection.js @@ -173,9 +173,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 0653102..9cc172f 100644 --- a/src/services/priceOracle.js +++ b/src/services/priceOracle.js @@ -335,6 +335,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; @@ -348,7 +353,7 @@ async function fetchFreshPrice(assetCode, issuer = null, redisUnavailable = fals issuer: normalisedIssuer, coalesced: coalescedCount, }); - return existing; + return existing.promise; } // Wrap the promise to handle rejections cleanly: on rejection, remove from @@ -363,8 +368,8 @@ async function fetchFreshPrice(assetCode, issuer = null, redisUnavailable = fals return result; }); - 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 d7ee37a..aa47668 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,