From 1299e576ddf5e84a60f90efb77fb332dd5795729 Mon Sep 17 00:00:00 2001 From: aaronmanuel309-bot Date: Sat, 29 Aug 2026 14:18:57 +0000 Subject: [PATCH] feat(webhook): add dead letter queue for failed message processing Permanently failed deliveries (max retries exhausted) and SSRF-blocked jobs are now routed to a bounded dead letter queue instead of being silently dropped, preventing silent data loss during downstream outages. Adds inspection, single-entry redelivery, and removal endpoints under /deadletter, Prometheus DLQ metrics, and architecture/runbook documentation. Closes #124 --- .gitignore | 1 + docs/WEBHOOK_ARCHITECTURE.md | 16 ++ docs/WEBHOOK_RUNBOOK.md | 39 +++ .../src/deadLetterQueue.ts | 160 ++++++++++++ webhook-delivery-service/src/delivery.ts | 46 +++- webhook-delivery-service/src/index.ts | 81 +++++- webhook-delivery-service/src/metrics.ts | 73 ++++++ .../tests/deadLetterQueue.test.ts | 247 ++++++++++++++++++ 8 files changed, 660 insertions(+), 3 deletions(-) create mode 100644 webhook-delivery-service/src/deadLetterQueue.ts create mode 100644 webhook-delivery-service/tests/deadLetterQueue.test.ts diff --git a/.gitignore b/.gitignore index 2f8e346..22dccc3 100644 --- a/.gitignore +++ b/.gitignore @@ -11,6 +11,7 @@ AGENTS.md agents/ issue.md node_modules/ +coverage/ # Windows build artifacts *.exe diff --git a/docs/WEBHOOK_ARCHITECTURE.md b/docs/WEBHOOK_ARCHITECTURE.md index 8652feb..94c8f0c 100644 --- a/docs/WEBHOOK_ARCHITECTURE.md +++ b/docs/WEBHOOK_ARCHITECTURE.md @@ -92,6 +92,18 @@ Transient network drops, rate limits (HTTP 429), and short-term receiver outages - **Full Jitter**: Prevents "thundering herd" issues by introducing randomized delay ($t_{jitter} = \text{random}(0, t_{backoff})$). - **Max Retries**: Defaulted to **5 attempts** before a webhook is classified as failed. +### 3.6 Dead Letter Queue (`/deadletter`) + +Rather than silently dropping a message that can never be delivered, the service routes permanently failed jobs into a bounded **Dead Letter Queue (DLQ)** for inspection and operator-driven redelivery: + +- **What is dead-lettered**: Deliveries that exhaust their maximum retry budget (`MAX_ATTEMPTS_EXHAUSTED`) and jobs rejected by the SSRF shield (`SSRF_BLOCKED`). +- **Retention**: A bounded, in-memory store (default **1,000 entries**, FIFO). When full, the oldest dead letter is evicted and counted as discarded for alerting. +- **Inspection**: `GET /deadletter` lists entries (newest first); `GET /deadletter/:id` fetches a single entry including the failed payload, attempt history, reason, and last error. +- **Redelivery**: `POST /deadletter/:id/requeue` reconstructs a fresh delivery job from the stored entry (with a fresh retry budget) and pushes it back onto the active queue. Requeued deliveries are re-signed and pass through the full security + retry pipeline again. +- **Removal**: `DELETE /deadletter/:id` removes a single entry; `DELETE /deadletter?confirm=true` purges the whole queue. Purging requires an explicit confirmation query parameter to prevent accidental data loss. + +The DLQ is intentionally in-memory to mirror the service's ingestion pipeline and keep inspection/redelivery on the fast path. For deployments requiring cross-restart durability, the out-of-process persistent queue pattern (Redis/RabbitMQ) noted in the runbook should be substituted. + --- ## 4. Monitoring & Metrics @@ -101,3 +113,7 @@ The service registers Prometheus counters and histograms to measure health indic - `webhook_delivery_duration_seconds`: Histogram of endpoint response latency. - `webhook_queue_size_current`: Gauge representing current queue occupancy. - `webhook_failures_total`: Total dropped or exhausted delivery alerts. +- `webhook_dead_letter_queue_size_current`: Gauge of the current DLQ occupancy. +- `webhook_dead_letter_enqueued_total`: Counter of messages entering the DLQ, labeled by `reason` (`MAX_ATTEMPTS_EXHAUSTED` | `SSRF_BLOCKED`). +- `webhook_dead_letter_requeued_total`: Counter of dead letters pushed back onto the active queue. +- `webhook_dead_letter_discarded_total`: Counter of dead letters evicted, purged, or manually removed. diff --git a/docs/WEBHOOK_RUNBOOK.md b/docs/WEBHOOK_RUNBOOK.md index db4b737..dd5ac1b 100644 --- a/docs/WEBHOOK_RUNBOOK.md +++ b/docs/WEBHOOK_RUNBOOK.md @@ -79,3 +79,42 @@ To prevent breaking integrations during rotation: To ensure the safety of the off-chain system, the SSRF (Server-Side Request Forgery) engine must be audited after any networking or DNS upgrades: 1. Verify that the URL parser correctly flags subnets by running integration tests. 2. Inspect server firewalls, ensuring egress traffic is strictly barred from routing to cloud provider private IP ranges and internal Kubernetes API service accounts. + +--- + +## 5. Dead Letter Queue (DLQ) Operations + +When a webhook permanently fails after exhausting its retry budget, or is rejected by the SSRF shield, it is moved to the **dead letter queue** instead of being dropped. `webhook_dead_letter_queue_size_current` climbing or `webhook_dead_letter_enqueued_total` increasing indicates persistent downstream failures. + +### Step 1: Inspect the dead letter queue +```bash +curl -s http://webhook-service.internal/deadletter +# {"count": 3, "deadLetters": [ { "id": "...", "reason": "MAX_ATTEMPTS_EXHAUSTED", ... } ]} + +# Inspect a single entry to see the failure reason and last error +curl -s http://webhook-service.internal/deadletter/ +``` + +### Step 2: Confirm the root cause before redelivering +1. Verify the downstream endpoint is healthy (`curl` / `GET /health` on the receiver). +2. Confirm the stored `errorMessage` is a transient failure (5xx, timeout) and not a request you should not re-send (e.g. 4xx contract violations). + +### Step 3: Redeliver the message +```bash +# Push a single dead letter back onto the active queue with a fresh retry budget +curl -X POST http://webhook-service.internal/deadletter//requeue +# { "status": "REQUEUED", "jobId": "" } +``` +The requeued message is re-signed and passes through the full security + retry pipeline again. + +### Step 4: Discard dead letters +```bash +# Remove a single entry +curl -X DELETE http://webhook-service.internal/deadletter/ +# Purge the entire queue (requires explicit confirmation) +curl -X DELETE http://webhook-service.internal/deadletter?confirm=true +``` + +### Operational Notes +- **Bounded retention**: The DLQ holds up to 1,000 entries in memory; the oldest entry is evicted (and counted as `webhook_dead_letter_discarded_total`) when capacity is exceeded. +- **Durability**: The DLQ is in-memory. For transactions that must survive service restarts, deploy the out-of-process persistent queue pattern (Redis/RabbitMQ) as the backing store. diff --git a/webhook-delivery-service/src/deadLetterQueue.ts b/webhook-delivery-service/src/deadLetterQueue.ts new file mode 100644 index 0000000..ec0ce51 --- /dev/null +++ b/webhook-delivery-service/src/deadLetterQueue.ts @@ -0,0 +1,160 @@ +/** + * Dead Letter Queue (DLQ) for the Webhook Delivery Service. + * + * Webhook deliveries that permanently fail (max attempts exhausted) or that are + * rejected by the SSRF security shield are moved into the dead letter queue + * instead of being silently dropped. This preserves the failed messages for + * inspection, offline retry ("requeue"), and operational forensics, so that a + * downstream subscriber outage or a mis-configured endpoint is never the cause + * of silent data loss. + * + * The queue is backed by an in-memory bounded store (mirroring the rest of the + * service's in-memory delivery pipeline) and is drained through the `/deadletter` + * HTTP endpoints. A pluggable requeue handler is registered by the delivery + * module so that dead letters can be pushed back onto the active delivery queue + * with a fresh retry budget. + */ +import type { WebhookPayload } from './delivery'; +import { + trackDeadLetterCount, + trackDeadLetterEnqueued, + trackDeadLetterDiscarded, +} from './metrics'; + +export type DeadLetterReason = 'MAX_ATTEMPTS_EXHAUSTED' | 'SSRF_BLOCKED'; + +export interface DeadLetterEntry { + /** Identifier of the original webhook job. */ + id: string; + /** Destination endpoint that the webhook was targeting. */ + url: string; + /** Event payload that could not be delivered. */ + payload: WebhookPayload; + /** Shared secret used for HMAC signing (required for requeue). */ + secret: string; + /** Optional Ed25519 private key used for signing (required for requeue). */ + privateKey?: string; + /** Number of delivery attempts made before dead-lettering. */ + attempts: number; + /** Maximum number of attempts permitted for redelivery. */ + maxAttempts: number; + /** Why the message was dead-lettered. */ + reason: DeadLetterReason; + /** Human readable error description captured at failure time. */ + errorMessage: string; + /** HTTP status code observed on the last failed attempt, if any. */ + statusCode?: number; + /** Timestamp (ms) at which the message entered the dead letter queue. */ + deadLetteredAt: number; +} + +const DEFAULT_MAX_DEAD_LETTERS = 1000; + +const store = new Map(); +let maxDeadLetters = DEFAULT_MAX_DEAD_LETTERS; + +/* + * The requeue orchestration (dead letter -> active delivery queue) lives in the + * delivery module, which owns both the queue and this store. This module only + * exposes the primitive data operations; requeueing pops an entry and the + * delivery module reconstructs a fresh job from it. + */ + +/** + * Insert a permanently failed delivery into the dead letter queue. If the + * queue is at capacity the oldest entry is evicted (FIFO) so the queue stays + * bounded; evicted entries are tracked as discarded for alerting purposes. + */ +export function reportDeadLetter(entry: DeadLetterEntry): void { + if (!store.has(entry.id) && store.size >= maxDeadLetters) { + const oldest = oldestEntryId(); + if (oldest) { + store.delete(oldest); + trackDeadLetterDiscarded(); + } + } + store.set(entry.id, entry); + trackDeadLetterEnqueued(entry.reason); + trackDeadLetterCount(store.size); +} + +/** Return a snapshot of all dead letters, newest first. */ +export function getDeadLetters(): DeadLetterEntry[] { + return [...store.values()].sort((a, b) => b.deadLetteredAt - a.deadLetteredAt); +} + +/** Return a single dead letter by id. */ +export function getDeadLetter(id: string): DeadLetterEntry | undefined { + return store.get(id); +} + +/** Return the current number of dead letters. */ +export function getDeadLetterCount(): number { + return store.size; +} + +/** + * Atomically remove a dead letter from the queue and return it so the caller + * can re-enqueue it as a fresh delivery job. Returns `undefined` if the entry + * does not exist. + */ +export function popDeadLetter(id: string): DeadLetterEntry | undefined { + const entry = store.get(id); + if (!entry) { + return undefined; + } + store.delete(id); + trackDeadLetterCount(store.size); + return entry; +} + +/** + * Remove a single dead letter from the queue without requeueing it. + * Returns `true` if an entry was present and removed. + */ +export function removeDeadLetter(id: string): boolean { + const existed = store.delete(id); + if (existed) { + trackDeadLetterCount(store.size); + trackDeadLetterDiscarded(); + } + return existed; +} + +/** Remove all dead letters from the queue. Returns the number removed. */ +export function purgeDeadLetters(): number { + const count = store.size; + if (count > 0) { + store.clear(); + trackDeadLetterCount(0); + for (let i = 0; i < count; i++) { + trackDeadLetterDiscarded(); + } + } + return count; +} + +/** Clear the DLQ and reset internal state (primarily for tests). */ +export function resetDeadLetterQueue(): void { + store.clear(); + maxDeadLetters = DEFAULT_MAX_DEAD_LETTERS; + trackDeadLetterCount(0); +} + +/** Configure the maximum size of the dead letter queue (primarily for tests). */ +export function setMaxDeadLetters(size: number): void { + maxDeadLetters = size; +} + +/** Return the id of the oldest entry (the FIFO eviction candidate). */ +function oldestEntryId(): string | undefined { + let oldestId: string | undefined; + let oldestTime = Infinity; + for (const entry of store.values()) { + if (entry.deadLetteredAt < oldestTime) { + oldestTime = entry.deadLetteredAt; + oldestId = entry.id; + } + } + return oldestId; +} \ No newline at end of file diff --git a/webhook-delivery-service/src/delivery.ts b/webhook-delivery-service/src/delivery.ts index 87be392..2effbf6 100644 --- a/webhook-delivery-service/src/delivery.ts +++ b/webhook-delivery-service/src/delivery.ts @@ -1,6 +1,7 @@ import axios from 'axios'; import { generateSignatures, validateUrlForSsrf } from './security'; -import { trackDeliveryAttempt, trackQueueSize, trackFailure } from './metrics'; +import { trackDeliveryAttempt, trackQueueSize, trackFailure, trackDeadLetterRequeued } from './metrics'; +import { reportDeadLetter, popDeadLetter, resetDeadLetterQueue, DeadLetterEntry } from './deadLetterQueue'; import { logger, LogAttributes } from './logger'; // Structured logging is skipped in tests: Jest's console interception adds @@ -106,12 +107,28 @@ export function getQueueSize(): number { return queue.length; } +/** + * Push a dead letter back onto the active delivery queue as a fresh job with a + * fresh retry budget. Returns the new webhook job id, or `null` if the dead + * letter does not exist. + */ +export function requeueDeadLetter(id: string): string | null { + const entry = popDeadLetter(id); + if (!entry) { + return null; + } + const newJobId = enqueueWebhook(entry.payload, entry.url, entry.secret, entry.privateKey, entry.maxAttempts); + trackDeadLetterRequeued(); + return newJobId; +} + /** * Clear queue and logs (primarily for testing) */ export function clearQueueAndLogs(): void { queue.length = 0; deliveryLogs.length = 0; + resetDeadLetterQueue(); trackQueueSize(0); } @@ -191,6 +208,18 @@ async function deliverWebhook(job: WebhookJob) { if (!ssrfCheck.valid) { const errorMsg = `SSRF Prevention: ${ssrfCheck.reason}`; trackFailure(); + reportDeadLetter({ + id: job.id, + url: job.url, + payload: job.payload, + secret: job.secret, + privateKey: job.privateKey, + attempts: job.attempts, + maxAttempts: job.maxAttempts, + reason: 'SSRF_BLOCKED', + errorMessage: errorMsg, + deadLetteredAt: Date.now(), + }); logDelivery('warn', 'webhook delivery dropped by SSRF check', { 'webhook.id': job.id, 'webhook.event': job.payload.event, @@ -286,8 +315,21 @@ async function deliverWebhook(job: WebhookJob) { lastAttemptTime: Date.now(), }); } else { - // Max attempts exhausted + // Max attempts exhausted -> move to the dead letter queue trackFailure(); + reportDeadLetter({ + id: job.id, + url: job.url, + payload: job.payload, + secret: job.secret, + privateKey: job.privateKey, + attempts: job.attempts, + maxAttempts: job.maxAttempts, + reason: 'MAX_ATTEMPTS_EXHAUSTED', + errorMessage: `Max attempts (${job.maxAttempts}) exhausted. Last error: ${errorMessage}`, + statusCode, + deadLetteredAt: Date.now(), + }); logDelivery('error', 'webhook delivery failed permanently', { 'webhook.id': job.id, 'webhook.event': job.payload.event, diff --git a/webhook-delivery-service/src/index.ts b/webhook-delivery-service/src/index.ts index 65b0fd1..9a667b1 100644 --- a/webhook-delivery-service/src/index.ts +++ b/webhook-delivery-service/src/index.ts @@ -1,5 +1,12 @@ import express, { Request, Response } from 'express'; -import { enqueueWebhook, getDeliveryLogs, getQueueSize } from './delivery'; +import { enqueueWebhook, getDeliveryLogs, getQueueSize, requeueDeadLetter } from './delivery'; +import { + getDeadLetters, + getDeadLetter, + getDeadLetterCount, + removeDeadLetter, + purgeDeadLetters, +} from './deadLetterQueue'; import { getPrometheusMetrics, getStatsSummary } from './metrics'; import { logger } from './logger'; @@ -92,6 +99,77 @@ app.get('/logs', (req: Request, res: Response) => { return res.json(logs); }); +/** + * GET /deadletter + * List all messages held in the dead letter queue + */ +app.get('/deadletter', (req: Request, res: Response) => { + return res.json({ + count: getDeadLetterCount(), + deadLetters: getDeadLetters(), + }); +}); + +/** + * GET /deadletter/:id + * Fetch a single dead letter by id + */ +app.get('/deadletter/:id', (req: Request, res: Response) => { + const entry = getDeadLetter(req.params.id); + if (!entry) { + return res.status(404).json({ error: 'Dead letter not found.' }); + } + return res.json(entry); +}); + +/** + * POST /deadletter/:id/requeue + * Push a dead letter back onto the active delivery queue + */ +app.post('/deadletter/:id/requeue', (req: Request, res: Response) => { + const newJobId = requeueDeadLetter(req.params.id); + if (newJobId === null) { + return res.status(404).json({ error: 'Dead letter not found or cannot be requeued.' }); + } + return res.json({ + status: 'REQUEUED', + message: 'Dead letter pushed back onto the active delivery queue.', + jobId: newJobId, + }); +}); + +/** + * DELETE /deadletter/:id + * Permanently remove a single dead letter + */ +app.delete('/deadletter/:id', (req: Request, res: Response) => { + const removed = removeDeadLetter(req.params.id); + if (!removed) { + return res.status(404).json({ error: 'Dead letter not found.' }); + } + return res.json({ + status: 'DELETED', + message: 'Dead letter removed from the queue.', + }); +}); + +/** + * DELETE /deadletter + * Purge all dead letters (requires ?confirm=true) + */ +app.delete('/deadletter', (req: Request, res: Response) => { + if (req.query.confirm !== 'true') { + return res.status(400).json({ + error: 'Purging the dead letter queue requires a ?confirm=true query parameter.', + }); + } + const removed = purgeDeadLetters(); + return res.json({ + status: 'PURGED', + message: `Removed ${removed} dead letter(s) from the queue.`, + }); +}); + /** * GET /health * Basic system health indicator @@ -101,6 +179,7 @@ app.get('/health', (req: Request, res: Response) => { status: 'UP', timestamp: Date.now(), queueSize: getQueueSize(), + deadLetterQueueSize: getDeadLetterCount(), }); }); diff --git a/webhook-delivery-service/src/metrics.ts b/webhook-delivery-service/src/metrics.ts index 67d353a..947be22 100644 --- a/webhook-delivery-service/src/metrics.ts +++ b/webhook-delivery-service/src/metrics.ts @@ -31,12 +31,38 @@ const totalFailures = new Counter({ registers: [registry], }); +const deadLetterCount = new Gauge({ + name: 'webhook_dead_letter_queue_size_current', + help: 'Current number of messages held in the dead letter queue', + registers: [registry], +}); + +const deadLetterEnqueued = new Counter({ + name: 'webhook_dead_letter_enqueued_total', + help: 'Total number of messages moved into the dead letter queue', + labelNames: ['reason'], + registers: [registry], +}); + +const deadLetterRequeued = new Counter({ + name: 'webhook_dead_letter_requeued_total', + help: 'Total number of dead letters pushed back onto the active delivery queue', + registers: [registry], +}); + +const deadLetterDiscarded = new Counter({ + name: 'webhook_dead_letter_discarded_total', + help: 'Total number of dead letters discarded (evicted, purged, or removed)', + registers: [registry], +}); + // Cache variables for UI Dashboard querying without Prometheus integration const statCache = { successCount: 0, failureCount: 0, totalAttempts: 0, durations: [] as number[], + dlqCount: 0, lastReset: Date.now(), }; @@ -104,6 +130,51 @@ export function trackFailure(): void { } } +/** + * Tracks the current dead letter queue size + */ +export function trackDeadLetterCount(size: number): void { + statCache.dlqCount = size; + try { + deadLetterCount.set(size); + } catch { + // Ignored + } +} + +/** + * Tracks a message entering the dead letter queue + */ +export function trackDeadLetterEnqueued(reason: string): void { + try { + deadLetterEnqueued.inc({ reason }); + } catch { + // Ignored + } +} + +/** + * Tracks a dead letter being pushed back onto the active queue + */ +export function trackDeadLetterRequeued(): void { + try { + deadLetterRequeued.inc(); + } catch { + // Ignored + } +} + +/** + * Tracks a dead letter being permanently discarded + */ +export function trackDeadLetterDiscarded(): void { + try { + deadLetterDiscarded.inc(); + } catch { + // Ignored + } +} + /** * Reset metric caches (primarily for testing) */ @@ -112,6 +183,7 @@ export function resetMetricCache(): void { statCache.failureCount = 0; statCache.totalAttempts = 0; statCache.durations = []; + statCache.dlqCount = 0; statCache.lastReset = Date.now(); } @@ -141,6 +213,7 @@ export function getStatsSummary() { p95LatencyMs: Math.round(p95 * 1000), p99LatencyMs: Math.round(p99 * 1000), successRate: Math.round(successRate * 100) / 100, + dlqCount: statCache.dlqCount, uptimeSeconds: Math.floor((Date.now() - statCache.lastReset) / 1000), }; } diff --git a/webhook-delivery-service/tests/deadLetterQueue.test.ts b/webhook-delivery-service/tests/deadLetterQueue.test.ts new file mode 100644 index 0000000..ef06205 --- /dev/null +++ b/webhook-delivery-service/tests/deadLetterQueue.test.ts @@ -0,0 +1,247 @@ +import request from 'supertest'; +import app from '../src/index'; +import { DeadLetterEntry, DeadLetterReason } from '../src/deadLetterQueue'; +import * as dlq from '../src/deadLetterQueue'; +import * as delivery from '../src/delivery'; +import { resetMetricCache, getStatsSummary } from '../src/metrics'; +import axios from 'axios'; + +jest.mock('axios'); +const mockedAxios = axios as jest.Mocked; + +function makeEntry(overrides: Partial = {}): DeadLetterEntry { + return { + id: Math.random().toString(36).substring(2, 15), + url: 'https://api.example.com/hook', + payload: { event: 'low_balance', timestamp: Date.now(), data: { meter_id: 1 } }, + secret: 'secret', + attempts: 5, + maxAttempts: 5, + reason: 'MAX_ATTEMPTS_EXHAUSTED' as DeadLetterReason, + errorMessage: 'Max attempts (5) exhausted. Last error: Network Error', + deadLetteredAt: Date.now(), + ...overrides, + }; +} + +describe('Dead Letter Queue Module', () => { + beforeEach(async () => { + delivery.clearQueueAndLogs(); + resetMetricCache(); + await new Promise((resolve) => setTimeout(resolve, 30)); // let in-flight timers settle + }); + + describe('1. Reporting dead letters', () => { + test('should add an entry and reflect it in the queue', () => { + const entry = makeEntry({ id: 'job-1' }); + dlq.reportDeadLetter(entry); + + expect(dlq.getDeadLetterCount()).toBe(1); + expect(dlq.getDeadLetter('job-1')).toEqual(entry); + expect(dlq.getDeadLetters()[0].id).toBe('job-1'); + expect(getStatsSummary().dlqCount).toBe(1); + }); + + test('should return undefined for unknown ids', () => { + expect(dlq.getDeadLetter('nope')).toBeUndefined(); + expect(dlq.getDeadLetterCount()).toBe(0); + }); + + test('should return newest-first ordering', () => { + dlq.reportDeadLetter(makeEntry({ id: 'a', deadLetteredAt: 100 })); + dlq.reportDeadLetter(makeEntry({ id: 'b', deadLetteredAt: 300 })); + dlq.reportDeadLetter(makeEntry({ id: 'c', deadLetteredAt: 200 })); + + expect(dlq.getDeadLetters().map((e) => e.id)).toEqual(['b', 'c', 'a']); + }); + }); + + describe('2. FIFO eviction at capacity', () => { + test('should evict the oldest entry when the queue is full', () => { + dlq.setMaxDeadLetters(2); + dlq.reportDeadLetter(makeEntry({ id: 'first', deadLetteredAt: 100 })); + dlq.reportDeadLetter(makeEntry({ id: 'second', deadLetteredAt: 200 })); + dlq.reportDeadLetter(makeEntry({ id: 'third', deadLetteredAt: 300 })); + + expect(dlq.getDeadLetterCount()).toBe(2); + expect(dlq.getDeadLetter('first')).toBeUndefined(); // evicted + expect(dlq.getDeadLetter('second')).toBeDefined(); + expect(dlq.getDeadLetter('third')).toBeDefined(); + }); + + test('should not evict when reporting an existing id', () => { + dlq.setMaxDeadLetters(1); + dlq.reportDeadLetter(makeEntry({ id: 'a', deadLetteredAt: 100 })); + dlq.reportDeadLetter(makeEntry({ id: 'a', attempts: 6 })); + + expect(dlq.getDeadLetterCount()).toBe(1); + expect(dlq.getDeadLetter('a')!.attempts).toBe(6); + }); + }); + + describe('3. Popping dead letters for requeue', () => { + test('should pop the entry and remove it from the queue', () => { + dlq.reportDeadLetter(makeEntry({ id: 'job-1' })); + + const popped = dlq.popDeadLetter('job-1'); + + expect(popped).toBeDefined(); + expect(popped!.id).toBe('job-1'); + expect(dlq.getDeadLetter('job-1')).toBeUndefined(); + expect(dlq.getDeadLetterCount()).toBe(0); + }); + + test('should return undefined for an unknown id', () => { + expect(dlq.popDeadLetter('unknown')).toBeUndefined(); + expect(dlq.getDeadLetterCount()).toBe(0); + }); + }); + + describe('4. Removing and purging dead letters', () => { + test('should remove a single dead letter', () => { + dlq.reportDeadLetter(makeEntry({ id: 'job-1' })); + expect(dlq.removeDeadLetter('job-1')).toBe(true); + expect(dlq.getDeadLetterCount()).toBe(0); + expect(dlq.removeDeadLetter('job-1')).toBe(false); + }); + + test('should purge all dead letters and report the count', () => { + dlq.reportDeadLetter(makeEntry({ id: 'a' })); + dlq.reportDeadLetter(makeEntry({ id: 'b' })); + const removed = dlq.purgeDeadLetters(); + + expect(removed).toBe(2); + expect(dlq.getDeadLetterCount()).toBe(0); + expect(getStatsSummary().dlqCount).toBe(0); + }); + }); +}); + +describe('DLQ Integration with Delivery', () => { + beforeEach(async () => { + jest.clearAllMocks(); + delivery.clearQueueAndLogs(); + dlq.resetDeadLetterQueue(); + resetMetricCache(); + await new Promise((resolve) => setTimeout(resolve, 30)); + }); + + test('permanently failed deliveries (max attempts exhausted) enter the DLQ', async () => { + jest.spyOn(delivery, 'calculateRetryDelay').mockReturnValue(1); + mockedAxios.post.mockRejectedValue({ message: 'Network Error', response: { status: 500 } }); + + const jobId = delivery.enqueueWebhook( + { event: 'test', timestamp: Date.now(), data: {} }, + 'https://api.example.com/hook', + 'secret', + undefined, + 2 + ); + + await new Promise((resolve) => setTimeout(resolve, 500)); + + const entry = dlq.getDeadLetter(jobId); + expect(entry).toBeDefined(); + expect(entry!.reason).toBe('MAX_ATTEMPTS_EXHAUSTED'); + expect(entry!.attempts).toBe(2); + expect(entry!.maxAttempts).toBe(2); + expect(entry!.url).toBe('https://api.example.com/hook'); + expect(entry!.errorMessage).toContain('Max attempts (2) exhausted'); + }); + + test('SSRF-blocked deliveries enter the DLQ without an HTTP request', async () => { + const jobId = delivery.enqueueWebhook( + { event: 'test', timestamp: Date.now(), data: {} }, + 'http://127.0.0.1:9000/hooks', + 'secret' + ); + + await new Promise((resolve) => setTimeout(resolve, 100)); + + const entry = dlq.getDeadLetter(jobId); + expect(entry).toBeDefined(); + expect(entry!.reason).toBe('SSRF_BLOCKED'); + expect(entry!.errorMessage).toContain('SSRF Prevention'); + expect(mockedAxios.post).not.toHaveBeenCalled(); + }); + + test('requeueing a dead letter re-enqueues and can be delivered successfully', async () => { + jest.spyOn(delivery, 'calculateRetryDelay').mockReturnValue(1); + + // First: let the job fail permanently so it lands in the DLQ. + mockedAxios.post + .mockRejectedValueOnce({ message: 'Network Error', response: { status: 500 } }) + .mockRejectedValueOnce({ message: 'Network Error', response: { status: 500 } }); + + const jobId = delivery.enqueueWebhook( + { event: 'test', timestamp: Date.now(), data: {} }, + 'https://api.example.com/hook', + 'secret', + undefined, + 2 // fails after 2 attempts + ); + + await new Promise((resolve) => setTimeout(resolve, 500)); + expect(dlq.getDeadLetter(jobId)).toBeDefined(); + + // Second: requeue the dead letter and let it succeed this time. + mockedAxios.post.mockResolvedValue({ status: 200, data: {} }); + const newJobId = delivery.requeueDeadLetter(jobId); + expect(newJobId).not.toBeNull(); + + await new Promise((resolve) => setTimeout(resolve, 500)); + + const logs = delivery.getDeliveryLogs(); + const log = logs.find((l) => l.id === newJobId!); + expect(log).toBeDefined(); + expect(log!.status).toBe('SUCCESS'); + expect(dlq.getDeadLetter(jobId)).toBeUndefined(); // moved off the DLQ + expect(dlq.getDeadLetterCount()).toBe(0); + }); + + test('DLQ endpoints list, requeue, and delete dead letters over HTTP', async () => { + // Fail a webhook so it lands in the DLQ. + jest.spyOn(delivery, 'calculateRetryDelay').mockReturnValue(1); + mockedAxios.post.mockRejectedValue({ message: 'Network Error', response: { status: 500 } }); + const jobId = delivery.enqueueWebhook( + { event: 'test', timestamp: Date.now(), data: {} }, + 'https://api.example.com/hook', + 'secret', + undefined, + 2 + ); + await new Promise((resolve) => setTimeout(resolve, 500)); + + // List + let res = await request(app).get('/deadletter'); + expect(res.status).toBe(200); + expect(res.body.count).toBe(1); + expect(res.body.deadLetters[0].id).toBe(jobId); + + // Fetch single + res = await request(app).get(`/deadletter/${jobId}`); + expect(res.status).toBe(200); + expect(res.body.id).toBe(jobId); + + // Requeue (delivery now succeeds) -> removed from DLQ + mockedAxios.post.mockResolvedValue({ status: 200, data: {} }); + res = await request(app).post(`/deadletter/${jobId}/requeue`); + expect(res.status).toBe(200); + expect(res.body.status).toBe('REQUEUED'); + expect(dlq.getDeadLetterCount()).toBe(0); + + // 404 for unknown id + res = await request(app).get('/deadletter/does-not-exist'); + expect(res.status).toBe(404); + + // Purge requires confirmation + dlq.reportDeadLetter(makeEntry({ id: 'x' })); + res = await request(app).delete('/deadletter'); + expect(res.status).toBe(400); + + res = await request(app).delete('/deadletter?confirm=true'); + expect(res.status).toBe(200); + expect(res.body.message).toContain('1'); + expect(dlq.getDeadLetterCount()).toBe(0); + }); +}); \ No newline at end of file