diff --git a/src/modules/queues/dlq/README.md b/src/modules/queues/dlq/README.md new file mode 100644 index 0000000..49d790c --- /dev/null +++ b/src/modules/queues/dlq/README.md @@ -0,0 +1,63 @@ +# Dead Letter Queue (DLQ) Module + +The Dead Letter Queue (DLQ) module captures, stores, monitors, and recovers failed jobs across Sentinel queues, preventing silent job loss and enabling controlled remediation. + +## Overview + +When background workers encounter terminal failures or exceed maximum retry thresholds, jobs are routed into the DLQ. This enables engineering and security teams to: +- **Inspect**: Review full payloads, stack traces, and execution attempt counts without logs rolling off. +- **Retry**: Re-dispatch jobs manually or through automated remediation handlers. +- **Audit & Discard**: Mark unrecoverable or invalid jobs as discarded with documented reasons. +- **Monitor**: Surface pending failure backlogs and error rates by queue. + +## Architecture + +- **`interfaces/dlq.interface.ts`**: Types for job payloads, statuses (`failed`, `retrying`, `retried`, `discarded`), filters, and metrics. +- **`dlq.service.ts`**: In-memory and extensible storage engine implementing `enqueueFailedJob`, `listJobs`, `retryJob`, `discardJob`, `getMetrics`, and `purge`. +- **`dlq.module.ts`**: NestJS module providing `DlqService` for injection across queue producers and consumers. + +## Usage + +### Enqueue a Failed Job + +```typescript +import { DlqService } from './dlq.service'; + +dlqService.enqueueFailedJob({ + queueName: 'transaction-monitoring', + jobName: 'audit-block-txs', + payload: { txHash: '0x123...' }, + error: new Error('RPC endpoint unavailable'), + attempts: 3, + maxAttempts: 3, + metadata: { workerId: 'worker-a' }, +}); +``` + +### Retry Controls + +```typescript +// Retry with custom execution logic +const result = await dlqService.retryJob(jobId, async (job) => { + return await dispatchToWorker(job.payload); +}); + +if (result.success) { + console.log('Job recovered successfully'); +} +``` + +### Metrics & Monitoring + +```typescript +const metrics = dlqService.getMetrics(); +console.log(`Pending failures: ${metrics.pending}, Total: ${metrics.totalFailed}`); +``` + +## Testing + +Run unit tests via Jest: + +```bash +npx jest src/modules/queues/dlq/dlq.service.spec.ts +``` diff --git a/src/modules/queues/dlq/dlq.module.ts b/src/modules/queues/dlq/dlq.module.ts new file mode 100644 index 0000000..e5363cc --- /dev/null +++ b/src/modules/queues/dlq/dlq.module.ts @@ -0,0 +1,8 @@ +import { Module } from '@nestjs/common'; +import { DlqService } from './dlq.service'; + +@Module({ + providers: [DlqService], + exports: [DlqService], +}) +export class DlqModule {} diff --git a/src/modules/queues/dlq/dlq.service.spec.ts b/src/modules/queues/dlq/dlq.service.spec.ts new file mode 100644 index 0000000..06ddc8d --- /dev/null +++ b/src/modules/queues/dlq/dlq.service.spec.ts @@ -0,0 +1,240 @@ +import { DlqService } from './dlq.service'; + +describe('DlqService', () => { + let service: DlqService; + + beforeEach(() => { + service = new DlqService(); + }); + + afterEach(() => { + service.clear(); + }); + + describe('enqueueFailedJob', () => { + it('captures failed jobs with payload, error metadata, and attempts', () => { + const job = service.enqueueFailedJob({ + queueName: 'transaction-monitoring', + jobName: 'audit-block-txs', + payload: { txHash: '0xabc123', ledgerSeq: 10452 }, + error: new Error('Database connection timeout'), + attempts: 3, + maxAttempts: 3, + metadata: { node: 'worker-1' }, + }); + + expect(job.id).toBeDefined(); + expect(job.queueName).toBe('transaction-monitoring'); + expect(job.jobName).toBe('audit-block-txs'); + expect(job.status).toBe('failed'); + expect(job.error.message).toBe('Database connection timeout'); + expect(job.error.stack).toBeDefined(); + expect(job.attempts).toBe(3); + expect(job.maxAttempts).toBe(3); + expect(job.retryCount).toBe(0); + expect(job.metadata).toEqual({ node: 'worker-1' }); + }); + + it('handles non-Error objects as error payload safely', () => { + const job = service.enqueueFailedJob({ + queueName: 'threat-feed', + jobName: 'sync-indicators', + payload: { feedId: 'misp-1' }, + error: { message: 'Network unreachable' }, + }); + + expect(job.error.message).toBe('Network unreachable'); + expect(job.attempts).toBe(1); + expect(job.maxAttempts).toBe(3); + }); + }); + + describe('getJob and listJobs', () => { + beforeEach(() => { + service.enqueueFailedJob({ + queueName: 'alerts-queue', + jobName: 'dispatch-slack-alert', + payload: { alertId: 'alt-1' }, + error: new Error('Rate limit exceeded'), + }); + service.enqueueFailedJob({ + queueName: 'analytics-queue', + jobName: 'compute-risk-score', + payload: { account: 'GABC' }, + error: new Error('Out of memory'), + }); + }); + + it('retrieves an existing job by ID', () => { + const all = service.listJobs(); + expect(all).toHaveLength(2); + + const job = service.getJob(all[0].id); + expect(job).toBeDefined(); + expect(job?.id).toBe(all[0].id); + }); + + it('filters jobs by queueName', () => { + const alertJobs = service.listJobs({ queueName: 'alerts-queue' }); + expect(alertJobs).toHaveLength(1); + expect(alertJobs[0].queueName).toBe('alerts-queue'); + + const analyticsJobs = service.listJobs({ queueName: 'analytics-queue' }); + expect(analyticsJobs).toHaveLength(1); + expect(analyticsJobs[0].queueName).toBe('analytics-queue'); + }); + + it('paginates results using limit and offset', () => { + const page1 = service.listJobs({ limit: 1, offset: 0 }); + expect(page1).toHaveLength(1); + + const page2 = service.listJobs({ limit: 1, offset: 1 }); + expect(page2).toHaveLength(1); + expect(page1[0].id).not.toBe(page2[0].id); + }); + }); + + describe('retryJob', () => { + it('marks job as retried when external handler succeeds', async () => { + const job = service.enqueueFailedJob({ + queueName: 'indexing', + jobName: 'index-contract', + payload: { contractId: 'CCONTRACT_ABC' }, + error: new Error('Transient 503 error'), + }); + + const mockRetryHandler = jest.fn().mockResolvedValue(true); + const result = await service.retryJob(job.id, mockRetryHandler); + + expect(result.success).toBe(true); + expect(result.job.status).toBe('retried'); + expect(result.job.retryCount).toBe(1); + expect(result.job.lastRetriedAt).toBeDefined(); + expect(mockRetryHandler).toHaveBeenCalledWith(result.job); + }); + + it('keeps job as failed when retry handler returns false', async () => { + const job = service.enqueueFailedJob({ + queueName: 'indexing', + jobName: 'index-contract', + payload: { contractId: 'CCONTRACT_ABC' }, + error: new Error('Service unavailable'), + }); + + const mockRetryHandler = jest.fn().mockResolvedValue(false); + const result = await service.retryJob(job.id, mockRetryHandler); + + expect(result.success).toBe(false); + expect(result.job.status).toBe('failed'); + expect(result.job.retryCount).toBe(1); + expect(result.error).toBe('Retry handler returned false'); + }); + + it('handles exceptions thrown by retry handler gracefully', async () => { + const job = service.enqueueFailedJob({ + queueName: 'indexing', + jobName: 'index-contract', + payload: { contractId: 'CCONTRACT_ABC' }, + error: new Error('Initial failure'), + }); + + const mockRetryHandler = jest.fn().mockRejectedValue(new Error('Fatal crash during retry')); + const result = await service.retryJob(job.id, mockRetryHandler); + + expect(result.success).toBe(false); + expect(result.job.status).toBe('failed'); + expect(result.job.error.message).toBe('Fatal crash during retry'); + }); + + it('throws NotFoundException for non-existent job ID', async () => { + await expect(service.retryJob('non-existent-id')).rejects.toThrow(); + }); + }); + + describe('discardJob', () => { + it('marks job as discarded and records audit reason', () => { + const job = service.enqueueFailedJob({ + queueName: 'webhook-delivery', + jobName: 'dispatch-webhook', + payload: { endpoint: 'http://malformed.url' }, + error: new Error('Malformed URL'), + }); + + const discarded = service.discardJob(job.id, 'Invalid endpoint URL cannot be delivered'); + expect(discarded.status).toBe('discarded'); + expect(discarded.discardReason).toBe('Invalid endpoint URL cannot be delivered'); + }); + }); + + describe('getMetrics', () => { + it('computes accurate status counts and per-queue breakdown', async () => { + const j1 = service.enqueueFailedJob({ + queueName: 'q1', + jobName: 'j1', + payload: {}, + error: new Error('e1'), + }); + const j2 = service.enqueueFailedJob({ + queueName: 'q1', + jobName: 'j2', + payload: {}, + error: new Error('e2'), + }); + service.enqueueFailedJob({ + queueName: 'q2', + jobName: 'j3', + payload: {}, + error: new Error('e3'), + }); + + await service.retryJob(j1.id); + service.discardJob(j2.id, 'Unrecoverable'); + + const metrics = service.getMetrics(); + expect(metrics.totalFailed).toBe(3); + expect(metrics.pending).toBe(1); + expect(metrics.retried).toBe(1); + expect(metrics.discarded).toBe(1); + + expect(metrics.byQueue['q1']).toEqual({ + total: 2, + pending: 0, + retried: 1, + discarded: 1, + }); + expect(metrics.byQueue['q2']).toEqual({ + total: 1, + pending: 1, + retried: 0, + discarded: 0, + }); + }); + }); + + describe('purge', () => { + it('purges retried and discarded jobs older than the cutoff threshold', () => { + const oldDate = new Date(Date.now() - 10 * 24 * 60 * 60 * 1000); // 10 days ago + const job1 = service.enqueueFailedJob({ + queueName: 'q1', + jobName: 'j1', + payload: {}, + error: new Error('e1'), + }); + job1.failedAt = oldDate; + job1.status = 'retried'; + + const job2 = service.enqueueFailedJob({ + queueName: 'q1', + jobName: 'j2', + payload: {}, + error: new Error('e2'), + }); + job2.failedAt = new Date(); // Recent + + const purged = service.purge(7 * 24 * 60 * 60 * 1000); // 7 days + expect(purged).toBe(1); + expect(service.getJob(job1.id)).toBeUndefined(); + expect(service.getJob(job2.id)).toBeDefined(); + }); + }); +}); diff --git a/src/modules/queues/dlq/dlq.service.ts b/src/modules/queues/dlq/dlq.service.ts new file mode 100644 index 0000000..c3a055c --- /dev/null +++ b/src/modules/queues/dlq/dlq.service.ts @@ -0,0 +1,219 @@ +import { Injectable, Logger, NotFoundException } from '@nestjs/common'; +import { randomUUID } from 'crypto'; +import { + DeadLetterJob, + DlqFilter, + DlqMetrics, + EnqueueFailedJobOptions, + RetryResult, +} from './interfaces/dlq.interface'; + +@Injectable() +export class DlqService { + private readonly logger = new Logger(DlqService.name); + private readonly jobs = new Map(); + + /** + * Capture a failed job into the Dead Letter Queue for investigation and recovery. + */ + enqueueFailedJob(options: EnqueueFailedJobOptions): DeadLetterJob { + const id = options.id || randomUUID(); + const errObj = + options.error instanceof Error + ? { + message: options.error.message, + name: options.error.name, + stack: options.error.stack, + } + : { + message: options.error.message || 'Unknown failure', + name: options.error.name, + stack: options.error.stack, + }; + + const job: DeadLetterJob = { + id, + queueName: options.queueName, + jobName: options.jobName, + payload: options.payload, + error: errObj, + attempts: options.attempts || 1, + maxAttempts: options.maxAttempts || 3, + failedAt: new Date(), + status: 'failed', + retryCount: 0, + metadata: options.metadata || {}, + }; + + this.jobs.set(id, job); + this.logger.warn( + `Job [${job.id}] from queue [${job.queueName}] routed to DLQ. Reason: ${job.error.message}`, + ); + + return job; + } + + /** + * Retrieve a single dead-letter job by ID. + */ + getJob(id: string): DeadLetterJob | undefined { + return this.jobs.get(id); + } + + /** + * List dead-letter jobs with optional filtering and pagination. + */ + listJobs(filter?: DlqFilter): DeadLetterJob[] { + let result = Array.from(this.jobs.values()); + + if (filter?.queueName) { + result = result.filter((j) => j.queueName === filter.queueName); + } + + if (filter?.status) { + result = result.filter((j) => j.status === filter.status); + } + + // Sort newest first + result.sort((a, b) => b.failedAt.getTime() - a.failedAt.getTime()); + + const offset = filter?.offset || 0; + const limit = filter?.limit || result.length; + + return result.slice(offset, offset + limit); + } + + /** + * Execute retry controls on a failed job. + * Invokes the optional handler; on success, marks the job as retried. + */ + async retryJob( + id: string, + retryHandler?: (job: DeadLetterJob) => Promise, + ): Promise { + const job = this.jobs.get(id); + if (!job) { + throw new NotFoundException(`Dead Letter Job with id ${id} not found`); + } + + job.status = 'retrying'; + job.retryCount += 1; + job.lastRetriedAt = new Date(); + + try { + if (retryHandler) { + const success = await retryHandler(job); + if (success) { + job.status = 'retried'; + this.logger.log(`Job [${job.id}] successfully retried and recovered.`); + return { success: true, job }; + } else { + job.status = 'failed'; + this.logger.warn(`Job [${job.id}] retry handler returned failure.`); + return { success: false, job, error: 'Retry handler returned false' }; + } + } else { + // Without an external handler, mark as retried for manual dispatch + job.status = 'retried'; + return { success: true, job }; + } + } catch (err: any) { + job.status = 'failed'; + job.error = { + message: err.message || 'Retry failed', + name: err.name, + stack: err.stack, + }; + this.logger.error(`Job [${job.id}] retry threw error: ${err.message}`); + return { success: false, job, error: err.message }; + } + } + + /** + * Explicitly discard a failed job from active retry consideration with an audit reason. + */ + discardJob(id: string, reason: string = 'Manually discarded'): DeadLetterJob { + const job = this.jobs.get(id); + if (!job) { + throw new NotFoundException(`Dead Letter Job with id ${id} not found`); + } + + job.status = 'discarded'; + job.discardReason = reason; + this.logger.log(`Job [${job.id}] marked as discarded. Reason: ${reason}`); + return job; + } + + /** + * Aggregate monitoring and observability metrics across queues. + */ + getMetrics(): DlqMetrics { + const metrics: DlqMetrics = { + totalFailed: 0, + pending: 0, + retried: 0, + discarded: 0, + byQueue: {}, + }; + + for (const job of this.jobs.values()) { + metrics.totalFailed += 1; + + if (!metrics.byQueue[job.queueName]) { + metrics.byQueue[job.queueName] = { + total: 0, + pending: 0, + retried: 0, + discarded: 0, + }; + } + + const q = metrics.byQueue[job.queueName]; + q.total += 1; + + if (job.status === 'failed' || job.status === 'retrying') { + metrics.pending += 1; + q.pending += 1; + } else if (job.status === 'retried') { + metrics.retried += 1; + q.retried += 1; + } else if (job.status === 'discarded') { + metrics.discarded += 1; + q.discarded += 1; + } + } + + return metrics; + } + + /** + * Purge resolved or discarded jobs older than the retention threshold. + */ + purge(olderThanMs: number = 7 * 24 * 60 * 60 * 1000): number { + const cutoff = Date.now() - olderThanMs; + let purged = 0; + + for (const [id, job] of this.jobs.entries()) { + if ( + (job.status === 'retried' || job.status === 'discarded') && + job.failedAt.getTime() < cutoff + ) { + this.jobs.delete(id); + purged += 1; + } + } + + if (purged > 0) { + this.logger.log(`Purged ${purged} expired jobs from DLQ storage.`); + } + + return purged; + } + + /** + * Clear all jobs (useful for testing teardown). + */ + clear(): void { + this.jobs.clear(); + } +} diff --git a/src/modules/queues/dlq/interfaces/dlq.interface.ts b/src/modules/queues/dlq/interfaces/dlq.interface.ts new file mode 100644 index 0000000..919fb9e --- /dev/null +++ b/src/modules/queues/dlq/interfaces/dlq.interface.ts @@ -0,0 +1,60 @@ +export type JobStatus = 'failed' | 'retrying' | 'retried' | 'discarded'; + +export interface DeadLetterJob { + id: string; + queueName: string; + jobName: string; + payload: T; + error: { + message: string; + name?: string; + stack?: string; + }; + attempts: number; + maxAttempts: number; + failedAt: Date; + status: JobStatus; + retryCount: number; + lastRetriedAt?: Date; + discardReason?: string; + metadata?: Record; +} + +export interface EnqueueFailedJobOptions { + id?: string; + queueName: string; + jobName: string; + payload: T; + error: Error | { message: string; name?: string; stack?: string }; + attempts?: number; + maxAttempts?: number; + metadata?: Record; +} + +export interface DlqFilter { + queueName?: string; + status?: JobStatus; + limit?: number; + offset?: number; +} + +export interface QueueMetricBreakdown { + total: number; + pending: number; + retried: number; + discarded: number; +} + +export interface DlqMetrics { + totalFailed: number; + pending: number; + retried: number; + discarded: number; + byQueue: Record; +} + +export interface RetryResult { + success: boolean; + job: DeadLetterJob; + error?: string; +}