Skip to content
Closed
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: 4 additions & 0 deletions src/indexer/eventPoller.js
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,10 @@ class EventPoller {
let response;
let lastError;
for (let attempt = 1; attempt <= RPC_MAX_RETRIES; attempt++) {
if (this.stopped) {
throw new Error('Poll stopped during RPC retry');
}

try {
response = await this.rpcBreaker.call(() =>
this.server.getEvents({
Expand Down
39 changes: 25 additions & 14 deletions src/services/idempotency.js
Original file line number Diff line number Diff line change
Expand Up @@ -63,32 +63,43 @@ async function storeIdempotencyResponse(key, statusCode, responseBody) {
function idempotencyMiddleware(resourceType = 'resource') {
return async (req, res, next) => {
const idempotencyKey = req.get('Idempotency-Key');


// Validate idempotency key format before processing
if (idempotencyKey && typeof idempotencyKey !== 'string') {
return res.status(400).json({ error: 'Invalid Idempotency-Key header' });
}

if (idempotencyKey && idempotencyKey.length === 0) {
return res.status(400).json({ error: 'Idempotency-Key header cannot be empty' });
}

// Check if this idempotency key was already processed before any other processing
if (idempotencyKey) {
const cached = await getIdempotencyResponse(idempotencyKey);
if (cached) {
// Return the cached response
res.set('Idempotency-Replay', 'true');
return res.status(cached.statusCode).json(cached.body);
}
}

// Store the original json() method
const originalJson = res.json.bind(res);
// Override json() to capture and cache the response

// Override json() to capture and cache the response only after request is processed successfully
res.json = function(data) {
if (idempotencyKey && res.statusCode >= 200 && res.statusCode < 300) {
// Only cache successful responses
storeIdempotencyResponse(idempotencyKey, res.statusCode, data);
}
return originalJson(data);
};
// Check if this idempotency key was already processed

// Mark that we're processing this key
if (idempotencyKey) {
const cached = await getIdempotencyResponse(idempotencyKey);
if (cached) {
// Return the cached response
res.set('Idempotency-Replay', 'true');
return res.status(cached.statusCode).json(cached.body);
}

// Mark that we're processing this key
res.set('Idempotency-Key', idempotencyKey);
}

return next();
};
}
Expand Down
32 changes: 14 additions & 18 deletions src/services/webhook.js
Original file line number Diff line number Diff line change
Expand Up @@ -108,25 +108,21 @@ async function probeReachability(webhookUrl, options = {}) {
}

async function deliver(webhookUrl, secret, payload) {
try {
const result = await sendSignedRequest(webhookUrl, secret, payload);
if (result.ok) {
logger.info('Webhook delivered', { alert_id: payload.alert_id, url: webhookUrl });
return;
}

logger.warn('Webhook delivery failed', {
alert_id: payload.alert_id,
url: webhookUrl,
status: result.status,
});
} catch (err) {
logger.warn('Webhook delivery failed', {
alert_id: payload.alert_id,
url: webhookUrl,
error: err.message,
});
const result = await sendSignedRequest(webhookUrl, secret, payload);
if (result.ok) {
logger.info('Webhook delivered', { alert_id: payload.alert_id, url: webhookUrl });
return;
}

const error = new Error(`Webhook delivery failed with status ${result.status}`);
error.statusCode = result.status;
error.duration_ms = result.duration_ms;
logger.warn('Webhook delivery failed', {
alert_id: payload.alert_id,
url: webhookUrl,
status: result.status,
});
throw error;
}

module.exports = {
Expand Down
Loading