From 980d8df92958498007b7de8ffccffbdac470023a Mon Sep 17 00:00:00 2001 From: chukwuemekarita Date: Fri, 25 Sep 2026 22:15:29 +0100 Subject: [PATCH] fix: WS graceful shutdown, cacheWarm success count, eventParser flat data, webhook validation - #414: Add shutdown() to priceWebSocket.js that drains connections via PriceSubscriptionManager.drain() before closing the server - #413: Track succeeded count incrementally in cacheWarm.js so the timeout handler reports real progress instead of hardcoding succeeded: 0 - #412: Return value directly instead of { value } wrapper for non-array non-object decoded values in eventParser.dataFromValue - #411: Validate webhook_id exists via webhookRepository.findById before creating delivery records to prevent orphaned deliveries --- src/indexer/eventParser.js | 4 +++- src/repositories/deliveryRepository.js | 10 ++++++++ src/startup/cacheWarm.js | 32 +++++++++++++++++++++----- src/ws/priceWebSocket.js | 19 ++++++++++++++- 4 files changed, 57 insertions(+), 8 deletions(-) diff --git a/src/indexer/eventParser.js b/src/indexer/eventParser.js index 6fead34..c4c059d 100644 --- a/src/indexer/eventParser.js +++ b/src/indexer/eventParser.js @@ -50,7 +50,9 @@ function dataFromValue(eventName, value, topicHintCount = 0) { return Object.fromEntries(valueFields.map((field, index) => [field, value[index] ?? null])); } - return { value }; + // For non-array non-object values (strings, numbers, bigints, etc.), + // return the value directly instead of wrapping in { value } (#412). + return value; } function mergeTopicHints(eventName, data, topics) { diff --git a/src/repositories/deliveryRepository.js b/src/repositories/deliveryRepository.js index c230eaf..d7f1bf3 100644 --- a/src/repositories/deliveryRepository.js +++ b/src/repositories/deliveryRepository.js @@ -38,6 +38,7 @@ const crypto = require('crypto'); const cache = require('../services/cache'); const logger = require('../logger'); +const webhookRepository = require('./webhookRepository'); const RETRY_QUEUE_KEY = 'webhooks:retries'; const RECENT_DELIVERIES_LIMIT = 100; @@ -79,6 +80,15 @@ function generateTraceId() { } async function create({ webhook_id, event_id, event_type, trace_id, request_id }) { + // Validate webhook exists before creating delivery to prevent orphaned records (#411). + if (!webhook_id) { + throw new Error('deliveryRepository.create: webhook_id is required'); + } + const webhook = await webhookRepository.findById(webhook_id); + if (!webhook) { + throw new Error(`deliveryRepository.create: webhook '${webhook_id}' does not exist`); + } + const id = generateId(); const now = new Date().toISOString(); const record = { diff --git a/src/startup/cacheWarm.js b/src/startup/cacheWarm.js index d6dde7f..1bbbcdf 100644 --- a/src/startup/cacheWarm.js +++ b/src/startup/cacheWarm.js @@ -68,22 +68,42 @@ async function warmCache( let timedOut = false; let timeoutId; - const warming = runWarmCache(allAssets, oracle, abortController.signal).then((summary) => { - if (!timedOut) { - log.info('Cache warm complete', summary); + // Wrap each fetch to track successes incrementally, so the timeout + // handler can snapshot the real count instead of hardcoding 0 (#413). + let succeededSoFar = 0; + const trackedResults = allAssets.map(({ code, issuer }) => { + if (abortController.signal.aborted) { + return Promise.reject(new Error('Warming cancelled')); } - return summary; + return Promise.resolve() + .then(() => oracle.fetchFreshPrice(code, issuer || null)) + .then((value) => { + if (isWarmSuccess({ status: 'fulfilled', value })) { + succeededSoFar++; + } + return { status: 'fulfilled', value }; + }) + .catch((reason) => ({ status: 'rejected', reason })); }); + const warming = Promise.all(trackedResults).then((results) => ({ + total: allAssets.length, + succeeded: results.filter(isWarmSuccess).length, + failed: allAssets.length - results.filter(isWarmSuccess).length, + timedOut: false, + durationMs: 0, + })); + const timeout = new Promise((resolve) => { timeoutId = setTimeout(() => { timedOut = true; // Signal abort to cancel in-flight fetches (#402) abortController.abort(); + const failed = allAssets.length - succeededSoFar; const summary = { total: allAssets.length, - succeeded: 0, - failed: allAssets.length, + succeeded: succeededSoFar, + failed, timedOut: true, durationMs: timeoutMs, }; diff --git a/src/ws/priceWebSocket.js b/src/ws/priceWebSocket.js index 5095054..d7ee37a 100644 --- a/src/ws/priceWebSocket.js +++ b/src/ws/priceWebSocket.js @@ -58,4 +58,21 @@ function attach(httpServer) { return wss; } -module.exports = { attach }; +/** + * Gracefully close all WebSocket connections and stop the heartbeat. + * Call this during process shutdown to avoid abrupt connection drops. + */ +async function shutdown(wss, drainTimeoutMs = 5000) { + if (!wss) return; + + await subscriptionManager.drain(drainTimeoutMs); + + await new Promise((resolve) => { + wss.close(() => { + logger.info('WebSocket server closed'); + resolve(); + }); + }); +} + +module.exports = { attach, shutdown };