diff --git a/docs/WEBHOOK_ARCHITECTURE.md b/docs/WEBHOOK_ARCHITECTURE.md index 7d92a6d..94c8f0c 100644 --- a/docs/WEBHOOK_ARCHITECTURE.md +++ b/docs/WEBHOOK_ARCHITECTURE.md @@ -92,16 +92,17 @@ 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 Distributed Job Scheduler with Lease-based Worker Claiming +### 3.6 Dead Letter Queue (`/deadletter`) -Queue processing runs on a **distributed job scheduler** (`jobScheduler.ts`) where multiple worker loops compete to claim due webhook jobs under short-lived **leases**: +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: -- **Claim protocol**: A job becomes claimable once its `runAt` time has elapsed. A worker acquires the job by claiming a lease from a shared lease registry (`LeaseStore`); while that lease is valid, no other worker can claim the same job, so concurrent replicas/workers can never double-deliver the same webhook. -- **Fencing across processes**: `LeaseStore.claim` is atomic (synchronous within the Node event loop and serialisable against a shared store such as Redis/etcd in production), which gives cross-process mutual exclusion. Each lease carries a monotonic fencing token. -- **Heartbeat / lease renewal**: while a job is executing, the owning worker renews its lease on a configurable interval, so a healthy long-running delivery is never stolen by a competing worker. -- **Crash recovery**: if a worker dies without renewing, its lease expires and another worker reclaims the job — exactly-once under normal operation, at-least-once on worker failure. -- **Retry via rescheduling**: a failed attempt with retries remaining calls `ctx.reschedule(nextAttemptTime)` (exponential backoff + jitter), returning the job to the claimable pool at a future time. -- **Horizontal scaling**: worker count is controlled by `WEBHOOK_WORKER_COUNT` (default `3`). Scaling replicas or raising the worker count increases delivery concurrency without risking duplicate deliveries. +- **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. --- @@ -112,9 +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_scheduler_workers_current`: Gauge of active worker loops. -- `webhook_scheduler_active_leases_current`: Gauge of jobs currently executing under a worker lease. -- `webhook_scheduler_jobs_submitted_total`: Counter of jobs submitted to the scheduler. -- `webhook_scheduler_jobs_processed_total`: Counter of jobs executed by workers. -- `webhook_scheduler_jobs_failed_total`: Counter of jobs whose execution threw. -- `webhook_scheduler_lease_reclaimed_total`: Counter of expired leases reclaimed by another worker (crash recovery events). +- `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 2279cc6..dd5ac1b 100644 --- a/docs/WEBHOOK_RUNBOOK.md +++ b/docs/WEBHOOK_RUNBOOK.md @@ -82,20 +82,39 @@ To ensure the safety of the off-chain system, the SSRF (Server-Side Request Forg --- -## 5. Scheduler & Worker Operations +## 5. Dead Letter Queue (DLQ) Operations -Queue processing runs on the distributed job scheduler, where workers claim jobs under short-lived leases. Diagnose scheduler health through the `webhook_scheduler_*` metrics on `/metrics` and the scheduler block on `/health`. +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. -### Health Indicators -- `webhook_scheduler_workers_current` **0** → worker loops are not running; the scheduler cannot drain the queue. Restart the service. -- `webhook_scheduler_active_leases_current` sustained at the worker count → all workers are blocked on slow deliveries; inspect downstream endpoint latency and consider scaling out. -- `webhook_scheduler_lease_reclaimed_total` climbing → workers are expiring mid-delivery (lease not renewed). Investigate event-loop blocking / GC pauses, or increase the lease duration. +### 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/ +``` -### Diagnosing duplicate or missed deliveries -1. Confirm workers are healthy: `curl -s http://webhook-service.internal/health` and check `scheduler.workers` is non-empty and `scheduler.pendingCount` is not climbing. -2. Confirm no lease thrash: `curl -s http://webhook-service.internal/metrics | grep webhook_scheduler_lease_reclaimed_total`. -3. If `pendingCount` climbs while `active_leases` stays low, a worker crash loop is likely; scale the deployment and inspect container restart counts. +### 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 +``` -### Tuning -- Delivery concurrency per instance: `WEBHOOK_WORKER_COUNT` (default `3`). -- Lease duration and heartbeat are configurable in `jobScheduler.ts` (`leaseDurationMs`, `leaseRenewIntervalMs`, `pollIntervalMs`). +### 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 192686f..199e0d1 100644 --- a/webhook-delivery-service/src/delivery.ts +++ b/webhook-delivery-service/src/delivery.ts @@ -1,7 +1,7 @@ import axios from 'axios'; import { generateSignatures, validateUrlForSsrf } from './security'; -import { trackDeliveryAttempt, trackQueueSize, trackFailure } from './metrics'; -import { JobScheduler, ExecuteContext } from './jobScheduler'; +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 @@ -121,12 +121,28 @@ export function getSchedulerStatus() { return scheduler.getStatus(); } +/** + * 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 { scheduler.clear(); deliveryLogs.length = 0; + resetDeadLetterQueue(); trackQueueSize(0); } @@ -170,6 +186,18 @@ async function deliverWebhook(job: WebhookJob, ctx: ExecuteContext) { 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, @@ -266,8 +294,21 @@ async function deliverWebhook(job: WebhookJob, ctx: ExecuteContext) { 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 bd428a9..f3d8951 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, getSchedulerStatus } 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'; @@ -95,6 +102,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 @@ -104,7 +182,7 @@ app.get('/health', (req: Request, res: Response) => { status: 'UP', timestamp: Date.now(), queueSize: getQueueSize(), - scheduler: getSchedulerStatus(), + deadLetterQueueSize: getDeadLetterCount(), }); }); diff --git a/webhook-delivery-service/src/metrics.ts b/webhook-delivery-service/src/metrics.ts index a7339ea..248cf95 100644 --- a/webhook-delivery-service/src/metrics.ts +++ b/webhook-delivery-service/src/metrics.ts @@ -74,12 +74,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(), }; @@ -225,6 +251,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) */ @@ -233,6 +304,7 @@ export function resetMetricCache(): void { statCache.failureCount = 0; statCache.totalAttempts = 0; statCache.durations = []; + statCache.dlqCount = 0; statCache.lastReset = Date.now(); } @@ -262,6 +334,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