Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 0 additions & 1 deletion package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

8 changes: 4 additions & 4 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -9,11 +9,11 @@
"node": ">=20.0.0"
},
"scripts": {
"build": "node_modules/.bin/nest build",
"build": "nest build",
"format": "prettier --write \"src/**/*.ts\" \"test/**/*.ts\"",
"start": "node_modules/.bin/nest start",
"start:dev": "node_modules/.bin/nest start --watch",
"start:debug": "node_modules/.bin/nest start --debug --watch",
"start": "nest start",
"start:dev": "nest start --watch",
"start:debug": "nest start --debug --watch",
"start:prod": "node dist/main.js",
"lint": "eslint \"{src,apps,libs,test}/**/*.ts\"",
"typecheck": "tsc --noEmit",
Expand Down
2 changes: 2 additions & 0 deletions src/app.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import { AuditModule } from './modules/audit/audit.module';
import { AiModule } from './modules/ai/ai.module';
import { HealthModule } from './modules/health/health.module';
import { MetricsModule } from './modules/metrics/metrics.module';
import { AdminModule } from './modules/admin/admin.module';
import { RequestMetricsMiddleware } from './modules/metrics/metrics.middleware';
import { DeadLetterModule } from './modules/dead-letter/dead-letter.module';
import { AgentTraceInterceptor } from './common/interceptors/agent-trace.interceptor';
Expand Down Expand Up @@ -119,6 +120,7 @@ import { AgentTraceInterceptor } from './common/interceptors/agent-trace.interce
HealthModule,
MetricsModule,
DeadLetterModule,
AdminModule,
],
providers: [
{ provide: APP_GUARD, useClass: JwtAuthGuard },
Expand Down
16 changes: 16 additions & 0 deletions src/modules/admin/admin.module.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
import { Module } from '@nestjs/common';
import { QueueModule } from '../../queues/queue.module';
import { DlqController } from './dlq.controller';
import { DlqService } from './dlq.service';

/**
* Administrative module providing secured operator controls over
* background queues, dead-letter processing, and forensic error recovery.
*/
@Module({
imports: [QueueModule],
controllers: [DlqController],
providers: [DlqService],
exports: [DlqService],
})
export class AdminModule {}
159 changes: 159 additions & 0 deletions src/modules/admin/dlq.controller.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
import { describe, it, expect, vi, beforeEach } from 'vitest';
import { Reflector } from '@nestjs/core';
import { UserRole } from '@prisma/client';
import { DlqController } from './dlq.controller';
import { DlqService } from './dlq.service';
import { ROLES_KEY } from '../../common/decorators/roles.decorator';
import { Queues } from '../../queues/queues.constants';

describe('DlqController', () => {
let controller: DlqController;
let service: Record<string, ReturnType<typeof vi.fn>>;

beforeEach(() => {
service = {
listFailedJobs: vi.fn().mockResolvedValue({
items: [],
total: 0,
page: 1,
limit: 20,
}),
getQueueStats: vi.fn().mockResolvedValue([
{
queue: Queues.DeadLetter,
failed: 0,
active: 0,
waiting: 0,
delayed: 0,
completed: 0,
paused: 0,
},
]),
getJobDetails: vi.fn().mockResolvedValue({
id: 'job-1',
queue: Queues.DeadLetter,
name: 'test',
data: {},
opts: {},
attemptsMade: 1,
timestamp: 123456,
}),
retryJob: vi.fn().mockResolvedValue({
jobId: 'job-1',
queue: Queues.DeadLetter,
retried: true,
message: 'Job retried',
}),
retryAllFailedJobs: vi.fn().mockResolvedValue({
retriedCount: 5,
queues: [Queues.DeadLetter],
}),
removeJob: vi.fn().mockResolvedValue({
jobId: 'job-1',
queue: Queues.DeadLetter,
removed: true,
}),
purgeQueue: vi.fn().mockResolvedValue({
purgedCount: 10,
removedJobIds: ['job-1', 'job-2'],
queues: [Queues.DeadLetter],
}),
};

controller = new DlqController(service as unknown as DlqService);
});

describe('RBAC Roles Guard Configuration', () => {
it('has @Roles(UserRole.OWNER, UserRole.ADMIN) defined at the class level', () => {
const reflector = new Reflector();
const roles = reflector.get<UserRole[]>(ROLES_KEY, DlqController);

expect(roles).toBeDefined();
expect(roles).toContain(UserRole.OWNER);
expect(roles).toContain(UserRole.ADMIN);
expect(roles).toHaveLength(2);
});
});

describe('Endpoints', () => {
it('listFailedJobs delegates query to dlqService', async () => {
const query = { page: 2, limit: 10, queue: Queues.Webhooks };
const res = await controller.listFailedJobs(query);

expect(service.listFailedJobs).toHaveBeenCalledWith(query);
expect(res.page).toBe(1);
});

it('getQueueStats delegates to dlqService', async () => {
const res = await controller.getQueueStats();

expect(service.getQueueStats).toHaveBeenCalledOnce();
expect(res).toHaveLength(1);
});

it('getJobDetails delegates queue and id to dlqService', async () => {
await controller.getJobDetails(Queues.Webhooks, 'wh-job-1');

expect(service.getJobDetails).toHaveBeenCalledWith(Queues.Webhooks, 'wh-job-1');
});

it('getDlqJobDetails defaults to DeadLetter queue', async () => {
await controller.getDlqJobDetails('dlq-job-1');

expect(service.getJobDetails).toHaveBeenCalledWith(Queues.DeadLetter, 'dlq-job-1');
});

it('retryJob delegates queue and id to dlqService', async () => {
await controller.retryJob(Queues.Webhooks, 'wh-job-1');

expect(service.retryJob).toHaveBeenCalledWith(Queues.Webhooks, 'wh-job-1');
});

it('retryDlqJob delegates to dlqService with default queue', async () => {
await controller.retryDlqJob('dlq-job-1');

expect(service.retryJob).toHaveBeenCalledWith(Queues.DeadLetter, 'dlq-job-1');
});

it('retryAllJobs delegates to dlqService with optional queue', async () => {
await controller.retryAllJobs(Queues.Webhooks);

expect(service.retryAllFailedJobs).toHaveBeenCalledWith(Queues.Webhooks);
});

it('retryQueueAllJobs delegates to dlqService', async () => {
await controller.retryQueueAllJobs(Queues.Transactions);

expect(service.retryAllFailedJobs).toHaveBeenCalledWith(Queues.Transactions);
});

it('removeJob delegates queue and id to dlqService', async () => {
await controller.removeJob(Queues.Transactions, 'tx-job-1');

expect(service.removeJob).toHaveBeenCalledWith(Queues.Transactions, 'tx-job-1');
});

it('removeDlqJob defaults to DeadLetter queue', async () => {
await controller.removeDlqJob('dlq-job-1');

expect(service.removeJob).toHaveBeenCalledWith(Queues.DeadLetter, 'dlq-job-1');
});

it('purgeQueue delegates query to dlqService', async () => {
const purgeDto = { queue: Queues.Reports, gracePeriodMs: 1000, limit: 50 };
await controller.purgeQueue(purgeDto);

expect(service.purgeQueue).toHaveBeenCalledWith(purgeDto);
});

it('purgeSpecificQueue merges path parameter into purgeDto', async () => {
const purgeDto = { gracePeriodMs: 0, limit: 100 };
await controller.purgeSpecificQueue(Queues.RiskAnalysis, purgeDto);

expect(service.purgeQueue).toHaveBeenCalledWith({
...purgeDto,
queue: Queues.RiskAnalysis,
});
});
});
});
150 changes: 150 additions & 0 deletions src/modules/admin/dlq.controller.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
import {
Controller,
Get,
Post,
Delete,
Param,
Query,
HttpCode,
HttpStatus,
} from '@nestjs/common';
import { ApiOperation, ApiTags, ApiResponse } from '@nestjs/swagger';
import { UserRole } from '@prisma/client';
import { Roles } from '../../common/decorators/roles.decorator';
import { ZodValidationPipe } from '../../common/pipes/zod-validation.pipe';
import { DlqService } from './dlq.service';
import {
ListDlqJobsQuery,
listDlqJobsQuerySchema,
PurgeDlqDto,
purgeDlqSchema,
DlqJobDetails,
QueueJobCounts,
} from './dto/dlq.dto';
import { Queues } from '../../queues/queues.constants';

/**
* Administrative Dead-Letter Queue (DLQ) controller.
* Restricted strictly to system administrators (OWNER and ADMIN roles).
*/
@ApiTags('admin-dlq')
@Controller('admin/dlq')
@Roles(UserRole.OWNER, UserRole.ADMIN)
export class DlqController {
constructor(private readonly dlqService: DlqService) {}

@Get()
@ApiOperation({ summary: 'List failed jobs across queues or for a specific queue' })
@ApiResponse({ status: 200, description: 'List of failed / dead-lettered jobs' })
async listFailedJobs(
@Query(new ZodValidationPipe(listDlqJobsQuerySchema)) query: ListDlqJobsQuery,
) {
return this.dlqService.listFailedJobs(query);
}

@Get('stats')
@ApiOperation({ summary: 'Get queue job counts and DLQ health stats' })
@ApiResponse({ status: 200, description: 'Summary counts across all BullMQ queues' })
async getQueueStats(): Promise<QueueJobCounts[]> {
return this.dlqService.getQueueStats();
}

@Post('retry-all')
@HttpCode(HttpStatus.OK)
@ApiOperation({ summary: 'Retry all failed jobs across all queues or a specified queue' })
@ApiResponse({ status: 200, description: 'Results of batch retry operation' })
async retryAllJobs(@Query('queue') queue?: string) {
return this.dlqService.retryAllFailedJobs(queue);
}

@Delete('purge')
@HttpCode(HttpStatus.OK)
@ApiOperation({ summary: 'Purge failed jobs across all queues or a specified queue' })
@ApiResponse({ status: 200, description: 'Results of purge operation' })
async purgeQueue(
@Query(new ZodValidationPipe(purgeDlqSchema)) query: PurgeDlqDto,
) {
return this.dlqService.purgeQueue(query);
}

@Post(':queue/retry-all')
@HttpCode(HttpStatus.OK)
@ApiOperation({ summary: 'Retry all failed jobs in a specific queue' })
async retryQueueAllJobs(@Param('queue') queue: string) {
return this.dlqService.retryAllFailedJobs(queue);
}

@Delete(':queue/purge')
@HttpCode(HttpStatus.OK)
@ApiOperation({ summary: 'Purge failed jobs in a specific queue' })
async purgeSpecificQueue(
@Param('queue') queue: string,
@Query(new ZodValidationPipe(purgeDlqSchema)) query: PurgeDlqDto,
) {
return this.dlqService.purgeQueue({ ...query, queue });
}

@Get(':queue/:id')
@ApiOperation({ summary: 'Inspect a specific failed job and error payload in a named queue' })
@ApiResponse({ status: 200, description: 'Job inspection details' })
async getJobDetails(
@Param('queue') queue: string,
@Param('id') id: string,
): Promise<DlqJobDetails> {
return this.dlqService.getJobDetails(queue, id);
}

@Get(':id')
@ApiOperation({ summary: 'Inspect a specific failed job in the default Dead-Letter Queue' })
@ApiResponse({ status: 200, description: 'Job inspection details' })
async getDlqJobDetails(
@Param('id') id: string,
@Query('queue') queue?: string,
): Promise<DlqJobDetails> {
return this.dlqService.getJobDetails(queue ?? Queues.DeadLetter, id);
}

@Post(':queue/:id/retry')
@HttpCode(HttpStatus.OK)
@ApiOperation({ summary: 'Retry a specific failed job in a named queue' })
@ApiResponse({ status: 200, description: 'Job retry confirmation' })
async retryJob(
@Param('queue') queue: string,
@Param('id') id: string,
) {
return this.dlqService.retryJob(queue, id);
}

@Post(':id/retry')
@HttpCode(HttpStatus.OK)
@ApiOperation({ summary: 'Retry a specific failed job in the default Dead-Letter Queue' })
@ApiResponse({ status: 200, description: 'Job retry confirmation' })
async retryDlqJob(
@Param('id') id: string,
@Query('queue') queue?: string,
) {
return this.dlqService.retryJob(queue ?? Queues.DeadLetter, id);
}

@Delete(':queue/:id')
@HttpCode(HttpStatus.OK)
@ApiOperation({ summary: 'Delete/remove a specific failed job from a named queue' })
@ApiResponse({ status: 200, description: 'Job removal confirmation' })
async removeJob(
@Param('queue') queue: string,
@Param('id') id: string,
) {
return this.dlqService.removeJob(queue, id);
}

@Delete(':id')
@HttpCode(HttpStatus.OK)
@ApiOperation({ summary: 'Delete/remove a specific failed job from the default Dead-Letter Queue' })
@ApiResponse({ status: 200, description: 'Job removal confirmation' })
async removeDlqJob(
@Param('id') id: string,
@Query('queue') queue?: string,
) {
return this.dlqService.removeJob(queue ?? Queues.DeadLetter, id);
}
}
Loading
Loading