Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion src/indexer/eventParser.js
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
10 changes: 10 additions & 0 deletions src/repositories/deliveryRepository.js
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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 = {
Expand Down
32 changes: 26 additions & 6 deletions src/startup/cacheWarm.js
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand Down
19 changes: 18 additions & 1 deletion src/ws/priceWebSocket.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 };