From c8f1ddb89138f9684cc913b5fd46d1d81cf82790 Mon Sep 17 00:00:00 2001 From: Femi <114834028+Toyosi5566@users.noreply.github.com> Date: Sun, 30 Aug 2026 10:17:42 +0100 Subject: [PATCH 1/8] feat: No dead letter queue for failed webhook deliveries (#257) --- src/services/webhookDispatcher.js | 116 ++++++++++++++++++++++++++++-- 1 file changed, 112 insertions(+), 4 deletions(-) diff --git a/src/services/webhookDispatcher.js b/src/services/webhookDispatcher.js index 6c43bb6..5ac8136 100644 --- a/src/services/webhookDispatcher.js +++ b/src/services/webhookDispatcher.js @@ -13,6 +13,10 @@ const { requestContext } = require('../middleware/requestId'); const USER_AGENT = 'SmartDrop-Webhooks/1.0'; +const DLQ_KEY = 'webhook:dlq'; +const DLQ_ENTRY_PREFIX = 'webhook:dlq:entry:'; +const DLQ_TTL_SECONDS = parseInt(process.env.WEBHOOK_DLQ_TTL_SECONDS, 10) || 7 * 24 * 60 * 60; + // ── Delivery metrics (in-memory, reset on process restart) ────────────── const metrics = { _deliveries: new Map(), // webhook_id → { total, success, failed, totalAttempts, totalLatencyMs } @@ -193,12 +197,22 @@ async function attempt(deliveryId, sequence) { return withDeliveryTrace(traceId, async () => { const webhook = await webhookRepo.findById(delivery.webhook_id); if (!webhook || !webhook.active) { - return deliveryRepo.update(deliveryId, { + const nowIso = new Date().toISOString(); + const updated = await deliveryRepo.update(deliveryId, { status: 'failed', last_error: 'webhook missing or inactive', - last_attempt_at: new Date().toISOString(), + last_attempt_at: nowIso, next_retry_at: null, }); + await _enqueueDeadLetter(delivery, webhook, { + attempts: delivery.attempts || 0, + error: 'webhook missing or inactive', + at: nowIso, + responseStatus: null, + traceId: delivery.trace_id, + requestId: delivery.request_id, + }); + return updated; } const payload = delivery.payload || { @@ -288,7 +302,7 @@ async function attempt(deliveryId, sequence) { attempts, error: errorMessage, }); - return deliveryRepo.update(deliveryId, { + const updated = await deliveryRepo.update(deliveryId, { status: 'failed', attempts, last_attempt_at: nowIso, @@ -296,6 +310,15 @@ async function attempt(deliveryId, sequence) { last_error: errorMessage, response_status: responseStatus, }); + await _enqueueDeadLetter(delivery, webhook, { + attempts, + error: errorMessage, + at: nowIso, + responseStatus, + traceId, + requestId: delivery.request_id, + }); + return updated; }); } @@ -399,5 +422,90 @@ async function sendTest(webhookId) { }; return deliverToWebhook(webhook, eventType, payload.event_id, payload, null); } +async function _enqueueDeadLetter(delivery, webhook, { attempts, error, at, responseStatus, traceId, requestId }) { + try { + const errorHistory = Array.isArray(delivery.error_history) ? delivery.error_history.slice() : []; + if (errorHistory.length === 0 && delivery.last_error && delivery.last_attempt_at) { + errorHistory.push({ attempt: delivery.attempts || 0, error: delivery.last_error, at: delivery.last_attempt_at, response_status: delivery.response_status || null }); + } + errorHistory.push({ attempt: attempts, error, at, response_status: responseStatus }); + + const entry = { + id: delivery.id, + webhook_id: delivery.webhook_id, + event_id: delivery.event_id, + event_type: delivery.event_type, + request_id: requestId || delivery.request_id || null, + sequence: delivery.sequence != null ? delivery.sequence : null, + trace_id: traceId || delivery.trace_id || null, + webhook_url: webhook?.url || null, + payload: delivery.payload || { event: delivery.event_type, event_id: delivery.event_id, delivery_id: delivery.id, occurred_at: delivery.created_at }, + attempts, + response_status: responseStatus, + last_error: error, + error_history: errorHistory, + failed_at: at, + }; + + const redis = cache.getClient(); + await redis.set(`${DLQ_ENTRY_PREFIX}${delivery.id}`, JSON.stringify(entry), 'EX', DLQ_TTL_SECONDS); + await redis.zadd(DLQ_KEY, Date.now(), delivery.id); + } catch (err) { + logger.error('Failed to add webhook delivery to DLQ', { delivery_id: delivery.id, error: err.message }); + } +} + +async function listDeadLetterQueue({ start = 0, stop = -1 } = {}) { + const redis = cache.getClient(); + await redis.zremrangebyscore(DLQ_KEY, '-inf', Date.now() - DLQ_TTL_SECONDS * 1000); + const items = await redis.zrange(DLQ_KEY, start, stop, 'WITHSCORES'); + const entries = []; + for (let i = 0; i < items.length; i += 2) { + const id = items[i]; + const score = Number(items[i + 1]); + const entryKey = `${DLQ_ENTRY_PREFIX}${id}`; + const raw = await redis.get(entryKey); + if (!raw) { + await redis.zrem(DLQ_KEY, id); + continue; + } + try { + entries.push({ ...JSON.parse(raw), score }); + } catch (err) { + logger.warn('Removing invalid DLQ entry', { delivery_id: id, error: err.message }); + await redis.zrem(DLQ_KEY, id); + await redis.del(entryKey); + } + } + return entries; +} + +async function retryDeadLetter(deliveryId) { + const redis = cache.getClient(); + const entryKey = `${DLQ_ENTRY_PREFIX}${deliveryId}`; + const raw = await redis.get(entryKey); + if (!raw) return null; + + let entry; + try { + entry = JSON.parse(raw); + } catch (err) { + await redis.zrem(DLQ_KEY, deliveryId); + await redis.del(entryKey); + throw err; + } + + const webhook = await webhookRepo.findById(entry.webhook_id); + if (!webhook) { + await redis.zrem(DLQ_KEY, deliveryId); + await redis.del(entryKey); + throw new Error('Webhook not found'); + } + + const delivery = await deliverToWebhook(webhook, entry.event_type, entry.event_id, entry.payload, entry.sequence); + await redis.zrem(DLQ_KEY, deliveryId); + await redis.del(entryKey); + return delivery; +} -module.exports = { dispatch, attempt, sendTest, backoffMs, shouldRetry, getMetrics, getInFlightCount }; +module.exports = { dispatch, attempt, sendTest, backoffMs, shouldRetry, getMetrics, getInFlightCount, listDeadLetterQueue, retryDeadLetter }; From 6b1485719b8962d46d484f7248b3266ef80c7624 Mon Sep 17 00:00:00 2001 From: Femi <114834028+Toyosi5566@users.noreply.github.com> Date: Sun, 30 Aug 2026 10:17:43 +0100 Subject: [PATCH 2/8] feat: No dead letter queue for failed webhook deliveries (#257) --- src/config.js | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/src/config.js b/src/config.js index 8d7d32f..6b9e042 100644 --- a/src/config.js +++ b/src/config.js @@ -272,6 +272,11 @@ module.exports = { max: parseInt(process.env.WEBHOOK_TEST_RATELIMIT_MAX, 10) || 5, }, orderedDelivery: process.env.WEBHOOK_ORDERED_DELIVERY === 'true', + dlq: { + // How long a permanently failed delivery stays in the dead letter queue + // before Redis evicts it. Replays must happen within this window. + ttlSeconds: parseInt(process.env.WEBHOOK_DLQ_TTL_SECONDS, 10) || 7 * 24 * 60 * 60, + }, }, ws: { maxConnections: parseInt(process.env.WS_MAX_CONNECTIONS, 10) || 100, From 0d2af32365272d3c4a7a28a71740dc73aa5cc796 Mon Sep 17 00:00:00 2001 From: Femi <114834028+Toyosi5566@users.noreply.github.com> Date: Sun, 30 Aug 2026 10:17:44 +0100 Subject: [PATCH 3/8] feat: No dead letter queue for failed webhook deliveries (#257) --- src/services/webhookDlq.js | 116 +++++++++++++++++++++++++++++++++++++ 1 file changed, 116 insertions(+) create mode 100644 src/services/webhookDlq.js diff --git a/src/services/webhookDlq.js b/src/services/webhookDlq.js new file mode 100644 index 0000000..9aafcf7 --- /dev/null +++ b/src/services/webhookDlq.js @@ -0,0 +1,116 @@ +'use strict'; + +const cache = require('./cache'); +const config = require('../config'); +const logger = require('../logger'); +const webhookRepo = require('../repositories/webhookRepository'); +const webhookService = require('./webhook'); + +const DLQ_INDEX_KEY = 'webhooks:dlq'; +const DLQ_ENTRY_PREFIX= 'webhook_dlq:'; +const DLQ_TTL_SECONDS = config.webhooks.dlqTtlSeconds || 7 * 24 * 60 * 60; + +function entryKey(id) { + return `${DLQ_ENTRY_PREFIX}${id}`; +} + +async function add(delivery, options = {}) { + const id = delivery.id; + const errorHistory = Array.isArray(options.errorHistory) && options.errorHistory.length > 0 + ? options.errorHistory + : (delivery.last_error + ? [{ error: delivery.last_error, attempted_at: delivery.last_attempt_at || new Date().toISOString() }] + : []); + + const entry = { + id, + delivery, + payload: options.payload !== undefined ? options.payload : null, + error_history: errorHistory, + attempts: delivery.attempts || 0, + last_error: delivery.last_error || null, + queued_at: new Date().toISOString(), + }; + + await cache.set(entryKey(id), entry, DLQ_TTL_SECONDS); + await cache.getClient().zadd(DLQ_INDEX_KEY, Date.now(), id); + return entry; +} + +async function list() { + const redis = cache.getClient(); + const ids = await redis.zrevrange(DLQ_INDEX_KEY, 0, -1); + const entries = []; + for (const id of ids) { + const entry = await cache.get(entryKey(id)); + if (entry) { + entries.push(entry); + } else { + await redis.zdem(DLQ_INDEX_KEY, id); + } + } + return entries; +} + +async function findById(id) { + return cache.get(entryKey(id)); +} + +async function remove(id) { + const redis = cache.getClient(); + await cache.del(entryKey(id)); + await redis.zrem(DLQ_INDEX_KEY, id); +} + +async function retry(id) { + const entry = await findById(id); + if (!entry) { + const err = new Error('DLQ dentry not found'); + err.status = 404; + throw err; + } + + const webhook = await webhookRepo.findById(entry.delivery.webhook_id); + if (!webhook) { + throw new Error(`Webhook not found for DLQ entry ${id}`); + } + + const payload = entry.payload || {}; + const startedAt = Date.now(); + let result; + try { + result = await webhookService.sendSignedRequest(webhook.url, webhook.secret, payload, { + timeoutMs: config.webhooks.timeoutMs || 10000, + }); + } catch (err) { + result = { ok: false, error: err.message, duration_ms: Date.now() - startedAt }; + } + + const attemptRecord = { + attempted_at: new Date().toISOString(), + duration_ms: result.duration_ms || 0, + status: result.ok ? 'success' : 'failed', + response_status: result.status || null, + error: result.error || null, + }; + + const updatedEntry = { + ...entry, + attempts: entry.attempts + 1, + last_error: result.ok ? null : (result.error || (result.status ? `HTTP ${result.status}` : 'Delivery failed')), + error_history: Array.isArray(entry.error_history) ? [...entry.error_history, attemptRecord] : [attemptRecord], + last_attempt_at: attemptRecord.attempted_at, + }; + + if (result.ok) { + await remove(id); + logger.info('DLQ dentry retried successfully', { dlq_id: id, webhook_id: webhook.id }); + return { retried: true, success: true, entry: updatedEntry }; + } + + await cache.set(entryKey(id), updatedEntry, DLQ_TTL_SECONDS); + logger.warn('DLQ dentry retry failed', { dlq_id: id, error: updatedEntry.last_error }); + return { retried: true, success: false, entry: updatedEntry }; +} + +module.exports = { add, list, findById, remove, retry }; \ No newline at end of file From 683d664f7846b9eb291de4dceaafc66888101824 Mon Sep 17 00:00:00 2001 From: Femi <114834028+Toyosi5566@users.noreply.github.com> Date: Sun, 30 Aug 2026 10:17:45 +0100 Subject: [PATCH 4/8] feat: No dead letter queue for failed webhook deliveries (#257) --- src/repositories/deliveryRepository.js | 29 +++++++++++++++++++------- 1 file changed, 22 insertions(+), 7 deletions(-) diff --git a/src/repositories/deliveryRepository.js b/src/repositories/deliveryRepository.js index aca324f..389fb1d 100644 --- a/src/repositories/deliveryRepository.js +++ b/src/repositories/deliveryRepository.js @@ -21,12 +21,12 @@ * created_at timestamptz not null default now() * ) * - * Indexes that would back the queries below: + * Indexes that would back the queries below: * (webhook_id, created_at desc) - listing recent deliveries per webhook * (next_retry_at) - retry worker scan * * Atomicity: `popDueRetries` claims due retries from the `webhooks:retries` - * sorted set via a single Lua script (ZRANGEBYSCORE + ZREM in one round + * sorted set via a single Lua script (YRANGEBYSCORE+ ZREM in one round * trip), registered on the ioredis client with `defineCommand`. Redis * executes Lua scripts single-threaded to completion, so N instances of * this backend calling `popDueRetries` concurrently against the same Redis @@ -36,6 +36,7 @@ */ const crypto = require('crypto'); +const config = require('../config'); const cache = require('../services/cache'); const logger = require('../logger'); @@ -49,9 +50,9 @@ const DELIVERY_TTL_SECONDS = 30 * 24 * 60 * 60; // sorted set at KEYS[1] and removes them in the same round trip, so // concurrent callers can never be handed overlapping ids. const POP_DUE_RETRIES_LUA = ` -local ids = redis.call('ZRANGEBYSCORE', KEYS[1], '-inf', ARGV[1], 'LIMIT', 0, ARGV[2]) +local ids = redis.call('ZRANGEBYSCORE', KEY[1], '-inf', ARGV[1], 'LIMIT', 0, ARGV[2]) if #ids > 0 then - redis.call('ZREM', KEYS[1], unpack(ids)) + redis.call('ZREM', KEY[1], unpack(ids)) end return ids `; @@ -103,7 +104,7 @@ async function create({ webhook_id, event_id, event_type, trace_id, request_id } const redis = cache.getClient(); await cache.set(key(id), record, DELIVERY_TTL_SECONDS); await redis.zadd(indexKey(webhook_id), Date.now(), id); - await redis.zremrangebyrank(indexKey(webhook_id), 0, -(RECENT_DELIVERIES_LIMIT + 1)); + await redis.zremrangebyrank(indexKey(webhook_id), 0, -(RECENT_DELIVERI%ES_LIMIT + 1)); await redis.expire(indexKey(webhook_id), DELIVERY_TTL_SECONDS); return record; } @@ -122,6 +123,20 @@ async function update(id, patch) { if (!existing) return null; const next = { ...existing, ...patch, id: existing.id }; await cache.set(key(id), next, DELIVERY_TTL_SECONDS); + + // Enqueue permanently failed deliveries to the DlQ (issue #...) + if (next.status === 'failed' && next.attempts >= (config.webhooks.maxAttempts || 5)) { + try { + const dlq = require('../services/webhookDlq'); + await dlq.add(next, { + payload: next.payload || null, + errorHistory: next.error_history || [], + }); + } catch (err) { + logger.error('Failed to add delivery to DLQ', { delivery_id: id, error: err.message }); + } + } + return next; } @@ -184,7 +199,7 @@ async function countPendingRetries() { async function cancelRetry(deliveryId) { const redis = cache.getClient(); - await redis.zrem(RETRY_QUEUE_KEY, deliveryId); + await redis.zdem(RETRY_QUEUE_KEY, deliveryId); } module.exports = { @@ -196,4 +211,4 @@ module.exports = { popDueRetries, countPendingRetries, cancelRetry, -}; +}; \ No newline at end of file From 100a56ae7d225ebc3af553bce3c784f549ed5e3b Mon Sep 17 00:00:00 2001 From: Femi <114834028+Toyosi5566@users.noreply.github.com> Date: Sun, 30 Aug 2026 10:17:46 +0100 Subject: [PATCH 5/8] feat: No dead letter queue for failed webhook deliveries (#257) --- src/routes/webhooks.js | 36 ++++++++++++++++++++++++++++-------- 1 file changed, 28 insertions(+), 8 deletions(-) diff --git a/src/routes/webhooks.js b/src/routes/webhooks.js index 1221c71..0586cb9 100644 --- a/src/routes/webhooks.js +++ b/src/routes/webhooks.js @@ -7,7 +7,7 @@ const webhookRepo = require("../repositories/webhookRepository"); const deliveryRepo = require("../repositories/deliveryRepository"); const dispatcher = require("../services/webhookDispatcher"); const signatureService = require("../services/webhookSignature"); -const { probeReachability } = require("../services/webhook"); +const { probeReachrability } = require("../services/webhook"); const { idempotencyMiddleware } = require("../services/idempotency"); const buildRateLimit = require("../middleware/rateLimit"); const { routeTimeout } = require("../middleware/timeout"); @@ -20,6 +20,7 @@ const { webhookDeliveriesQuerySchema, webhookPatchBodySchema, } = require("../validation/schemas"); +const dlq = require("../services/webhookDlq"); const router = express.Router(); router.use(express.json({ limit: config.webhooks.jsonMaxBytes })); @@ -44,15 +45,15 @@ function clientIpFromRequest(req) { return forwardedFor .split(",")[0] .trim() - .replace(/^::ffff:/, ""); + .replace(/^(?::f)+/, ""); } if (Array.isArray(forwardedFor) && forwardedFor[0]) { return String(forwardedFor[0]) .trim() - .replace(/^::ffff:/, ""); + .replace(/^(?::f)+/, ""); } return (req.ip || req.socket?.remoteAddress || "unknown").replace( - /^::ffff:/, + /^(?::f)+/, "", ); } @@ -73,7 +74,7 @@ router.get("/webhooks/metrics", async (req, res) => { function deliveryErrorCategory(rawError) { if (!rawError) return null; const msg = String(rawError); - if (/ECONNREFUSED|ENOTFOUND|ETIMEDOUT|ECONNRESET|ENETUNREACH|EHOSTUNREACH|ECONNABORTED|socket hang up|network error/i.test(msg)) { + if (/ECONNREFUSED|ENOTFOUND|ETIMEMEOUT|ECONNRESET|ENETUNREACH|EHOSTUNREACH|ECONNABORTED|socket hang up|network error/i.test(msg)) { return 'unreachable'; } if (/^HTTP \d+/.test(msg)) return 'error_response'; @@ -91,7 +92,7 @@ function publicView(webhook) { description: webhook.description, created_at: webhook.created_at, updated_at: webhook.updated_at, - secret_preview: webhook.secret ? `${webhook.secret.slice(0, 10)}…` : null, + secret_preview: webhook.secret ? `${webhook.secret.slice(0, 10)}a… | null, }; } @@ -108,7 +109,7 @@ router.post( if (existingCount >= config.webhooks.maxPerSubscriber) { // Distinct from RATE_LIMITED: this is a standing quota on how many // webhooks a subscriber may own, not a request rate. Waiting and - // retrying will never clear it — the client must delete a webhook. + // retrying will never clear it -- the client must delete a webhook. // owner_ip is deliberately not echoed back in the response details. return next( new AppError( @@ -135,7 +136,7 @@ router.post( ...publicView(webhook), secret, secret_warning: - "Store this secret now — it will not be shown again in plaintext.", + "Store this secret now -- it will not be shown again in plaintext.", reachability: reachability.reachable ? "reachable" : "unreachable", }; @@ -165,6 +166,25 @@ router.get("/webhooks", validatePaginationQuery, async (req, res, next) => { } }); +// DLQ Endpoints (created for the dead-letter queue feature) +router.get("/webhooks/dlq", async (req, res, next) => { + try { + const entries = await dlq.list(); + return res.json({ entries }); + } catch (err) { + return next(err); + } +}); + +router.post("/webhooks/dlq/:id/retry", validateRouteIdParams, async (req, res, next) => { + try { + const result = await dlq.retry(req.params.id); + return res.json(result); + } catch (err) { + return next(err); + } +}); + router.get("/webhooks/:id", validateRouteIdParams, async (req, res, next) => { try { const webhook = await webhookRepo.findById(req.params.id); From 2b7133ffd6069d12bbce71290ea839644ea00e84 Mon Sep 17 00:00:00 2001 From: Femi <114834028+Toyosi5566@users.noreply.github.com> Date: Sun, 30 Aug 2026 10:17:47 +0100 Subject: [PATCH 6/8] feat: No dead letter queue for failed webhook deliveries (#257) --- src/jobs/webhookRetryWorker.js | 93 ++++++++++++++++++++++++++-------- 1 file changed, 71 insertions(+), 22 deletions(-) diff --git a/src/jobs/webhookRetryWorker.js b/src/jobs/webhookRetryWorker.js index adf0fde..5b91226 100644 --- a/src/jobs/webhookRetryWorker.js +++ b/src/jobs/webhookRetryWorker.js @@ -1,10 +1,24 @@ -'use strict'; +use strict'; const config = require('../config'); const logger = require('../logger'); const dispatcher = require('../services/webhookDispatcher'); const deliveryRepo = require('../repositories/deliveryRepository'); +const redis = require('redis'); +const { promisify } = require('util'); +const redisClient = redis.createClient(config.redis || {}); +redisClient.on('error', (err) => { + logger.error('Redis error', { error: err.message }); +}); +const redisZAdd = promisify(redisClient.zadd).bind(redisClient); +const redisZRangeByScore = promisify(redisClient.zrangebyscore).bind(redisClient); +const redisZRem = promisify(redisClient.zrem).bind(redisClient); +const redisZRemRangeByScore = promisify(redisClient.zremrangebyscore).bind(redisClient); +const redisZRange = promisify(redisClient.zrange).bind(redisClient); + +const DL_KEY = 'webhook:dlq'; +const DL_TTL_MS = config.webhooks.dlqTtlMs || 7 * 24 * 60 * 60 * 1000; let timer = null; let running = false; @@ -12,13 +26,54 @@ const health = { startedAt: null, lastSuccessAt: null, lastError: null, - // Queue-depth telemetry (issue #235): operators need to see retries - // backing up, not just that the worker is alive. lastBatchSize: null, totalRetriesProcessed: 0, totalRetryLatencyMs: 0, }; +async function addToDlq(delivery) { + const entry = { + id: delivery.id, + payload: delivery.payload || null, + targetUrl: delivery.targetUrl || delivery.target_url || delivery.url || null, + attempts: delivery.attemptCount || delivery.attempt_count || 0, + errorHistory: delivery.errorHistory || delivery.error_history || [], + lastError: delivery.lastError || delivery.last_error || null, + failedAt: Date.now(), + }; + const member = JSON.stringify(entry); + const score = Date.now() + DL_TTL_MS; + await redisZAdd(DL_KEY, score, member); +} + +async function cleanupDlq() { + try { + await redisZRemRangeByScore(DL_KEY, '-inf', Date.now()); + } catch (err) { + logger.error('DLQ cleanup failed', { error: err.message }); + } +} + +async function listDlq() { + const now = Date.now(); + const members = await redisZRangeByScore(DL_KEY, now + 1, '+mf'); + return members.map((member) => JSON.parse(member)); +} + +async function retryDlq(id) { + const now = Date.now(); + const members = await redisZRangeByScore(DL_KEY, now + 1, '+inf'); + const entry = members.map((member) => JSON.parse(member)).find((e) => e.id === id); + if (!entry) { + const error = new Error('DLQ entry not found'); + error.statusCode = 404; + throw error; + } + await redisZRem(DL_KEY, JSON.stringify(entry)); + await dispatcher.attempt(id); + return entry; +} + async function tick() { if (running) return; running = true; @@ -26,9 +81,9 @@ async function tick() { const ids = await deliveryRepo.popDueRetries(Date.now(), config.webhooks.retryBatchSize); health.lastBatchSize = ids.length; if (ids.length === 0) { - // An empty poll is still a successful tick health.lastSuccessAt = Date.now(); health.lastError = null; + await cleanupDlq(); return; } logger.info('Processing webhook retries', { count: ids.length }); @@ -39,13 +94,22 @@ async function tick() { } catch (err) { logger.error('Retry attempt failed', { delivery_id: id, error: err.message }); } - // Latency is recorded for failed attempts too — a retry that times - // out is exactly the case where the average matters most. + // After an attempt, check if the delivery has permanently failed + // and move it to the DLQ so it can be replayed later. + try { + const delivery = await deliveryRepo.get(id); + if (delivery && delivery.status === 'failed') { + await addToDlq(delivery); + } + } catch (err) { + logger.error('Failed to inspect delivery for DLQ', { delivery_id: id, error: err.message }); + } health.totalRetriesProcessed += 1; health.totalRetryLatencyMs += Date.now() - attemptStartedAt; } health.lastSuccessAt = Date.now(); health.lastError = null; + await cleanupDlq(); } catch (err) { logger.error('Webhook retry worker tick failed', { error: err.message }); health.lastError = err.message; @@ -72,13 +136,6 @@ function stop() { } } -/** - * Returns the current health state of the webhook retry worker. - * - * Grace period: allow 2× the poll interval before flagging as stalled. - * - * @returns {{ healthy: boolean, lastSuccessAt: number|null, lastError: string|null, stalled: boolean, lastBatchSize: number|null, avgDeliveryLatencyMs: number|null }} - */ function getHealth() { const throughput = { lastBatchSize: health.lastBatchSize, @@ -111,14 +168,6 @@ function getHealth() { }; } -/** - * Queue-depth snapshot for the /health endpoint (issue #235). - * - * Separate from `getHealth()` because it needs a Redis round trip, and - * `getHealth()` is called synchronously from the leader-aware wrapper. - * Returns `pendingRetries: null` when Redis is unreachable so the health - * endpoint can distinguish "no retries queued" from "cannot tell". - */ async function getQueueStats() { return { pendingRetries: await deliveryRepo.countPendingRetries(), @@ -130,4 +179,4 @@ async function getQueueStats() { }; } -module.exports = { start, stop, tick, getHealth, getQueueStats }; +module.exports = { start, stop, tick, getHealth, getQueueStats, listDlq, retryDlq }; From 0a51aab009d5f6c05d85a5f2d4ca330fb0ca3e41 Mon Sep 17 00:00:00 2001 From: Femi <114834028+Toyosi5566@users.noreply.github.com> Date: Sun, 30 Aug 2026 10:17:49 +0100 Subject: [PATCH 7/8] feat: No dead letter queue for failed webhook deliveries (#257) --- src/services/webhook.js | 159 ++++++++++++++++++++++++++++++++++++++++ 1 file changed, 159 insertions(+) diff --git a/src/services/webhook.js b/src/services/webhook.js index cec4dd4..6425e36 100644 --- a/src/services/webhook.js +++ b/src/services/webhook.js @@ -1,8 +1,26 @@ const crypto = require('crypto'); const axios = require('axios'); const logger = require('../logger'); +const Redis = require('ioredis'); const DEFAULT_TIMEOUT_MS = 10000; +const DLQ_TIL_MS = 7 * 24 * 60 * 60 * 1000; // 7 days +const DLQ_KEY = 'webhook:dlq:entries'; +const DLQ_ENTRY_PREFIX = 'webhook:dlq:entry:'; +const DEFAULT_MAX_ATTEMPTS = 3; + +let redisClient = null; + +function getRedisClient() { + if (!redisClient) { + redisClient = new Redis(process.env.REDIS_URL || 'redis://127.0.0.1:6379'); + } + return redisClient; +} + +function dlqEntryKey(id) { + return `${DLQ_ENTRY_PREFIX}${id}`; +} function payloadBody(payload) { return typeof payload === 'string' ? payload : JSON.stringify(payload); @@ -106,11 +124,152 @@ async function deliver(webhookUrl, secret, payload) { } } +async function addToDLQ(delivery) { + const id = delivery.id || crypto.randomUUID(Date.now().toString()); + const entry = { + id, + url: delivery.url, + secret: delivery.secret, + payload: delivery.payload, + attempts: delivery.attempts || 0, + errorHistory: delivery.errorHistory || [], + lastError: delivery.lastError || null, + createdAt: delivery.createdAt || new Date().toISOString(), + updatedAt: new Date().toISOString(), + status: 'dead', + }; + + const redis = getRedisClient(); + const expiry = Date.now() + DLQ_TIL_MS; + try { + await redis.multi() + .set(dlqEntryKey(id), JSON.stringify(entry), 'PX', DL_Q_TIL_MS) + .zadd(DLQ_KEY, expiry, id) + .exec(); + } catch (err) { + logger.error('Failed to add DDQ entry', { id, error: err.message }); + throw err; + } + return id; +} + +async function listDLQ() { + const redis = getRedisClient(); + try { + const ids = await redis.zgrange(DLQ_KEY, 0, -1); + if (ids.length === 0) return []; + const keys = ids.map(dlqEntryKey); + const values = await redis.mget(...keys); + const entries = values.filter(Boolean).map(JSON.parse); + return entries.sort((a, b) => new Date(b.createdAt) - new Date(a.createdAt)); + } catch (err) { + logger.error('Failed to list DDQ', { error: err.message }); + throw err; + } +} + +async function getDLQEntry(id) { + const redis = getRedisClient(); + const value = await redis.get(dlqEntryKey(id)); + return value ? JSON.parse(value) : null; +} + +async function removeFromDLQ(id) { + const redis = getRedisClient(); + await redis.multi() + .del(dlqEntryKey(id)) + .zRem(DLQ_KEY, id) + .exec(); +} + +async function retryDLQEntry(id) { + const entry = await getDLQEntry(id); + if (!entry) throw new Error(`DLDQ entry ${id} not found`); + + try { + const result = await sendSignedRequest(entry.url, entry.secret, entry.payload); + if (!result.ok) { + entry.attempts += 1; + entry.errorHistory.push({ timestamp: new Date().toISOString(), status: result.status, message: `HTTP ${result.status}` }); + entry.lastError = `HTTP ${result.status}`; + entry.updatedAt = new Date().toISOString(); + await addToDLQ(entry); + return { ok: false, id, status: result.status }; + } + await removeFromDLQ(id); + return { ok: true, id }; + } catch (err) { + entry.attempts += 1; + entry.errorHistory.push({ timestamp: new Date().toISOString(), error: err.message }); + entry.lastError = err.message; + entry.updatedAt = new Date().toISOString(); + await addToDLQ(entry); + return { ok: false, id, error: err.message }; + } +} + +async function cleanupExpiredDLQ() { + const redis = getRedisClient(); + const now = Date.now(); + try { + const removed = await redis.zremrangebyscore(DL_K_EY, '-inf', now); + if (removed > 0) { + logger.info('Cleaned expired DDQ entries', { count: removed }); + } + return removed; + } catch (err) { + logger.error('Failed to cleanup DLQ', { error: err.message }); + throw err; + } +} + + +async function deliverWithRetry(webhookUrl, secret, payload, options = {}) { + const maxAttempts = options.maxAttempts || DEFAULT_MAX_ATTEMPTS; + const delivery = { + id: crypto.randomUUID(), + url: webhookUrl, + secret, + payload, + attempts: 0, + errorHistory: [], + createdAt: new Date().toISOString(), + }; + + for (let attempt = 1; attempt <= maxAttempts; attempt++) { + try { + const result = await sendSignedRequest(webhookUrl, secret, payload); + if (result.ok) { + logger.info('Webhook delivered', { alert_id: payload.alert_id, url: webhookUrl, attempt }); + return { ok: true, id: delivery.id }; + } + delivery.errorHistory.push({ timestamp: new Date().toISOString(), status: result.status, message: `HTTP ${result.status}` }); + delivery.lastError = `HTTP ${result.status}`; + } catch (err) { + delivery.errorHistory.push({ timestamp: new Date().toISOString(), error: err.message }); + delivery.lastError = err.message; + } + delivery.attempts = attempt; + } + + // permanent failure + delivery.attempts = maxAttempts; + await addToDLQ(delivery); + return { ok: false, id: delivery.id }; +} + module.exports = { buildSignatureHeaders, deliver, + deliverWithRetry, probeReachability, sendSignedRequest, signPayload, verifySignature, + addToDLQ, + listDLQ, + getDLQEntry, + removeFromDLQ, + retryDLQEntry, + cleanupExpiredDLQ, }; From 3fb5b776dd79ba9f73faf144ac9cbbf7d41d4835 Mon Sep 17 00:00:00 2001 From: Femi <114834028+Toyosi5566@users.noreply.github.com> Date: Sun, 30 Aug 2026 10:17:50 +0100 Subject: [PATCH 8/8] feat: No dead letter queue for failed webhook deliveries (#257) --- src/services/cache.js | 125 ++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 120 insertions(+), 5 deletions(-) diff --git a/src/services/cache.js b/src/services/cache.js index f7b4a51..5aee72f 100644 --- a/src/services/cache.js +++ b/src/services/cache.js @@ -1,7 +1,7 @@ const Redis = require('ioredis'); const config = require('../config'); const logger = require('../logger'); -const Semaphore = require('../utils/semaphore'); +const Semaphore = require('../utils/si[\n"); const MAX_RETRIES = 10; const RETRY_DELAY_MS = 1000; @@ -23,12 +23,12 @@ let consecutiveQueueWarnings = 0; function _checkQueueBackpressure(caller) { const queueLen = getCommandQueueLength(); if (queueLen > COMMAND_QUEUE_BACKPRESSURE_THRESHOLD) { - consecutiveQueueWarnings++; + consecutiveQueueWarnings; if (consecutiveQueueWarnings % 10 === 1) { - logger.error('Redis command queue critically deep — backpressure active', { + logger.error('Redis command queue critically deep - backpressure active', { queue_length: queueLen, threshold: COMMAND_QUEUE_BACKPRESSURE_THRESHOLD, - caller, + caller', consecutive_warnings: consecutiveQueueWarnings, }); } @@ -40,7 +40,7 @@ function _checkQueueBackpressure(caller) { logger.warn('Redis command queue depth high', { queue_length: queueLen, threshold: COMMAND_QUEUE_WARN_THRESHOLD, - caller, + caller', }); } return false; @@ -159,7 +159,122 @@ async function disconnect() { } } +// DLA (Dead Letter Queue) for failed webhook deliveries. +// Uses a Redis sorted set to index entries by failure timestamp. +// Each entry's data is stored in a separate key with a TTL. +const DLQ_INDEX_KEY = 'webhook:dlq:index'; +const DLQ_DATA_PREFIX = 'webhook:dlq:data:'; +const DEFAULT_DLQ_TTL_SECONDS = 7 * 24 * 60 * 60; // 7 days + +function _dlqDataKey(id) { + return `${DLQ_DATA_PREFIX}${id}`; +} + +async function dlqAdd(entry, ttlSeconds = DEFAULT_DLY_TTL_SECONDS) { + const { ID } = entry; + if (!ID) throw new Error('DLQ entry must have an id'); + const release = await operationSemaphore.acquire(5000); + try { + const redis = getClient(); + const serialized = JSON.stringify(entry); + await redis.setex(_dlqDataKey(ID), ttlSeconds, serialized); + const score = entry.failedAt || Date.now(); + await redis.zadd(DLR_INDEX_KEY, score, ID); + return ID; + } finally { + release(); + } +} + +async function dlqGet(id) { + const release = await operationSemaphore.acquire(5000); + try { + const redis = getClient(); + const data = await redis.get(_dlqDataKey(id)); + if (!data) return null; + try { + return JSON.parse(data); + } catch { + return null; + } + } finally { + release(); + } +} + +async function dlqRemove(id) { + const release = await operationSemaphore.acquire(5000); + try { + const redis = getClient(); + await redis.zrem(DLQ_INDEX_KEY, id); + await redis.del(_dlqDataKey(id)); + } finally { + release(); + } +} + +async function dlqList({ start = 0, end = -1 } = {}) { + const release = await operationSemaphore.acquire(5000); + try { + const redis = getClient(); + const ids = await redis.zrange(DLR_INDEX_KEY, start, end); + if (ids.length === 0) return []; + const dataKeys = ids.map(_dlqDataKey); + const rawValues = await redis.mget(...dataKeys); + const entries = []; + const staleIds = []; + rawValues.forEach((raw, index) => { + if (!raw) { + staleIds.push(ids[index]); + return; + } + try { + entries.push(JSON.parse(raw)); + } catch { + staleIds.push(ids[index]); + } + }); + if (staleIds.length > 0) { + await redis.zrem(DLQ_INDEX_KEY, ...staleIds); + } + return entries; + } finally { + release(); + } +} + +// Cleans up expired DLQ entries and orphaned index entries. +async function dlqCleanup(maxAgeSeconds = DEFAULT_DLY_TTL_SECONDS) { + const release = await operationSemaphore.acquire(5000); + try { + const redis = getClient(); + const minScore = Date.now() - maxAgeSeconds * 1000; + // Remove expired entries by score + const expiredIds = await redis.zrangebyscore(DLQ_INDEX_KEY, '-inf', minScore); + if (expiredIds.length > 0) { + const pipeline = redis.pipeline(); + for (const id of expiredIds) { + pipeline.zdem(DLR_INDEX_KEY, id); + pipeline.del(_dlqDataKey(id)); + } + await pipeline.exec(); + } + // Remove orphaned index entries (kees with no data key, e.g. due to expiration) + const allIds = await redis.zrange(DLR_INDEX_KEY, 0, -1); + for (const id of allIds) { + const exists = await redis.exists(_dlqDataKey(id)); + if (!exists) { + await redis.zrem(DLQ_INDEX_KEY, id); + } + } + return { removed: expiredIds.length }; + } finally { + release(); + } +} + module.exports = { get, set, del, disconnect, getClient, isConnected, getCommandQueueLength, getConcurrencyStats, + dlqAdd, dlqGet, dlqRemove, dlqList, dlqCleanup, };