diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index eb86b946..a4868d56 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -100,3 +100,21 @@ jobs: - name: Run tests run: npm run test + + # Fast static verification of the committed migrations: checks that every + # migration directory contains a non-empty, structurally valid migration.sql + # and that directory names follow the Prisma naming convention. Runs on every + # pull request and needs no PostgreSQL connection, so it also catches broken + # migrations in containerized runners without database services. + migrations-static: + name: migrations (static) + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-node@v4 + with: + node-version: 22 + cache: npm + - run: npm ci --include=dev + - name: Verify migration files (no database required) + run: npm run db:verify:static diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index a092da49..5d8c26d0 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -42,8 +42,39 @@ its own microservice. 2. `npm run db:verify` passes (`scripts/verify-migrations.sh`). See [Database migrations](#database-migrations) below. 3. Database schema changes include a valid Prisma migration folder containing `migration.sql`. -4. New endpoints are documented with OpenAPI/Swagger decorators. -5. Cross-repo contracts (response envelope, entity/enum names) still match `astroid-web` and `astroid-sdk`. +4. `npm run db:verify:static` passes — see [Database migration verification](#database-migration-verification). +5. New endpoints are documented with OpenAPI/Swagger decorators. +6. Cross-repo contracts (response envelope, entity/enum names) still match `astroid-web` and `astroid-sdk`. + +## Database migration verification + +Schema changes in Prisma are strictly verified before they reach production so a +malformed or drifted migration can never break a deployment. CI runs the static +check on every pull request, and the database-backed verification runs whenever +`DATABASE_URL` is available (locally or in CI): + +| Job | Requires a database | What it checks | +|---|---|---| +| `migrations (static)` (`npm run db:verify:static`) | No | Every migration directory has a non-empty `migration.sql` containing executable SQL (not only comments), quotes/parentheses are balanced, and directory names follow the Prisma convention `_` | +| `migrations (database)` (`npm run db:verify`) | Yes | Migrations apply cleanly to a fresh PostgreSQL instance and the schema rebuilt from the migrations alone matches `prisma/schema.prisma` (no drift) | + +### Naming convention + +Migration directories **must** be named `_` +(e.g. `20260830174000_sync_schema`). CI fails the build otherwise. Create +migrations with `npm run prisma:migrate` — never by hand. + +### Locally + +```bash +npm run db:verify:static # fast static checks, no database needed +npm run db:verify # full verification (needs DATABASE_URL, and + # SHADOW_DATABASE_URL for the drift check) +``` + +The static mode is what containerized CI runners without an active PostgreSQL +connection use; the full script degrades to it automatically when `DATABASE_URL` +is unset. ## Database migrations diff --git a/package.json b/package.json index c593b3b9..d8ab08da 100644 --- a/package.json +++ b/package.json @@ -26,7 +26,8 @@ "prisma:seed": "ts-node prisma/seed.ts", "db:seed": "ts-node prisma/seed.ts", "db:verify": "bash scripts/verify-migrations.sh", - "db:migrate": "ts-node src/database/migrate.cli.ts" + "db:migrate": "ts-node src/database/migrate.cli.ts", + "db:verify:static": "DATABASE_URL= SHADOW_DATABASE_URL= bash scripts/verify-migrations.sh" }, "prisma": { "seed": "ts-node prisma/seed.ts" diff --git a/src/modules/audit/audit.module.ts b/src/modules/audit/audit.module.ts index 2088d627..7f3f52fa 100644 --- a/src/modules/audit/audit.module.ts +++ b/src/modules/audit/audit.module.ts @@ -20,6 +20,19 @@ import { AuditCleanupQueue } from './queues/audit-cleanup.queue'; db: redisConfig().db, }, }), + // The dedicated `audit` queue persists audit entries asynchronously via the + // AuditWorker (src/workers/audit.worker.ts). Bounded retries with + // exponential backoff ride out transient database outages without dropping + // entries; terminal failures land in the dead-letter queue. + BullModule.registerQueue({ + name: Queues.Audit, + defaultJobOptions: { + attempts: 5, + backoff: { type: 'exponential', delay: 1000 }, + removeOnComplete: { count: 1000 }, + removeOnFail: { age: 7 * 24 * 3600 }, + }, + }), BullModule.registerQueue({ name: Queues.AuditCleanup, defaultJobOptions: { diff --git a/src/modules/health/health.controller.ts b/src/modules/health/health.controller.ts index 44fc41c8..6bc23ccf 100644 --- a/src/modules/health/health.controller.ts +++ b/src/modules/health/health.controller.ts @@ -1,14 +1,20 @@ -import { Controller, Get, Res, HttpStatus } from '@nestjs/common'; -import { ApiOperation, ApiResponse, ApiTags } from '@nestjs/swagger'; +import { Controller, Get, HttpException, HttpStatus, Res, UseGuards } from '@nestjs/common'; +import { ApiOperation, ApiResponse, ApiTags, ApiBearerAuth } from '@nestjs/swagger'; import { HealthIndicatorResult } from '@nestjs/terminus'; import { SkipThrottle } from '@nestjs/throttler'; import { Response } from 'express'; import { Public } from '../../common/decorators/public.decorator'; import { SkipAudit } from '../../common/decorators/skip-audit.decorator'; +import { JwtAuthGuard } from '../../common/guards/jwt-auth.guard'; +import { RolesGuard } from '../../common/guards/roles.guard'; +import { Roles } from '../../common/decorators/roles.decorator'; +import { UserRole } from '@prisma/client'; +import { SkipPublicRateLimit } from '../../common/decorators/skip-public-rate-limit.decorator'; import { PrismaHealthIndicator } from './indicators/prisma.health'; import { RedisHealthIndicator } from './indicators/redis.health'; import { StellarHealthIndicator } from './indicators/stellar.health'; import { DatabaseMigrationHealthIndicator } from './indicators/database-migration.health'; +import { BullMQHealthIndicator, QueuesHealthReport } from './indicators/bullmq.health'; /** Per-dependency report shape returned under `services` in the readiness body. */ interface ReadinessServiceReport { @@ -98,7 +104,7 @@ export class HealthController { }; } - @Get(['ready', 'readiness']) + @Get('readiness') @ApiOperation({ summary: 'Application readiness check' }) @ApiResponse({ status: 200, description: 'Application is ready' }) @ApiResponse({ status: 503, description: 'Application is not ready' }) @@ -192,3 +198,52 @@ function unwrap(result: HealthIndicatorResult, key: string): ReadinessServiceRep ...(message ?? error ? { error: message ?? error } : {}), }; } + +/** + * BullMQ queue diagnostics. Separate from `HealthController` so the queue + * internals stay protected: this controller is NOT marked `@Public()` and + * requires an authenticated OWNER, ADMIN, DEVELOPER or AUDITOR. + */ +@ApiTags('health') +@Controller('health') +@SkipAudit() +@SkipThrottle({ api: true, auth: true }) +@SkipPublicRateLimit() +export class QueuesHealthController { + constructor(private readonly bullmqHealthIndicator: BullMQHealthIndicator) {} + + @Get('queues') + @UseGuards(JwtAuthGuard, RolesGuard) + @Roles(UserRole.OWNER, UserRole.ADMIN, UserRole.DEVELOPER, UserRole.AUDITOR) + @ApiBearerAuth('access-token') + @ApiOperation({ + summary: 'BullMQ queue health check', + description: + 'Inspects every registered BullMQ queue (notifications, webhooks, ' + + 'stellar-sync, analytics, reports, outbox-events, stellar-fee-bump, ' + + 'transactions, risk-analysis, dead-letter, audit-cleanup, audit) and ' + + 'returns waiting/active/failed/delayed/completed/paused job counts plus ' + + 'Redis connectivity. Used by Kubernetes probes and dashboards to monitor ' + + 'asynchronous worker health.', + }) + @ApiResponse({ status: 200, description: 'Per-queue job counts and Redis connectivity' }) + @ApiResponse({ status: 401, description: 'Not authenticated' }) + @ApiResponse({ status: 403, description: 'Insufficient permissions' }) + @ApiResponse({ status: 503, description: 'Redis unreachable or all queues failing' }) + async checkQueuesHealth(): Promise { + const report = await this.bullmqHealthIndicator.checkHealth(); + + if (report.status === 'down') { + throw new HttpException( + { + statusCode: 503, + message: 'Redis unreachable or all BullMQ queues failing health probes', + report, + }, + HttpStatus.SERVICE_UNAVAILABLE, + ); + } + + return report; + } +} diff --git a/src/modules/health/health.module.ts b/src/modules/health/health.module.ts index 72adb104..df8f9631 100644 --- a/src/modules/health/health.module.ts +++ b/src/modules/health/health.module.ts @@ -1,28 +1,31 @@ -import { Module } from '@nestjs/common'; -import { TerminusModule } from '@nestjs/terminus'; -import { HealthController } from './health.controller'; -import { PrismaHealthIndicator } from './indicators/prisma.health'; -import { RedisHealthIndicator } from './indicators/redis.health'; -import { StellarHealthIndicator } from './indicators/stellar.health'; -import { DatabaseMigrationHealthIndicator } from './indicators/database-migration.health'; -import { DatabaseModule } from '../../database/database.module'; - -@Module({ - // TerminusModule supplies `HealthCheckService` and the indicator base class - // used by PrismaHealthIndicator. - imports: [DatabaseModule, TerminusModule], - controllers: [HealthController], - providers: [ - PrismaHealthIndicator, - RedisHealthIndicator, - StellarHealthIndicator, - DatabaseMigrationHealthIndicator, - ], - exports: [ - PrismaHealthIndicator, - RedisHealthIndicator, - StellarHealthIndicator, - DatabaseMigrationHealthIndicator, - ], -}) -export class HealthModule {} +import { Module } from '@nestjs/common'; +import { TerminusModule } from '@nestjs/terminus'; +import { HealthController, QueuesHealthController } from './health.controller'; +import { PrismaHealthIndicator } from './indicators/prisma.health'; +import { RedisHealthIndicator } from './indicators/redis.health'; +import { StellarHealthIndicator } from './indicators/stellar.health'; +import { DatabaseMigrationHealthIndicator } from './indicators/database-migration.health'; +import { BullMQHealthIndicator } from './indicators/bullmq.health'; +import { DatabaseModule } from '../../database/database.module'; + +@Module({ + // TerminusModule supplies `HealthCheckService` and the indicator base class + // used by PrismaHealthIndicator. + imports: [DatabaseModule, TerminusModule], + controllers: [HealthController, QueuesHealthController], + providers: [ + PrismaHealthIndicator, + RedisHealthIndicator, + StellarHealthIndicator, + DatabaseMigrationHealthIndicator, + BullMQHealthIndicator, + ], + exports: [ + PrismaHealthIndicator, + RedisHealthIndicator, + StellarHealthIndicator, + DatabaseMigrationHealthIndicator, + BullMQHealthIndicator, + ], +}) +export class HealthModule {} diff --git a/src/modules/health/index.ts b/src/modules/health/index.ts index 10e6d76e..c28e3054 100644 --- a/src/modules/health/index.ts +++ b/src/modules/health/index.ts @@ -1,6 +1,7 @@ -export * from './health.module'; -export * from './health.controller'; -export * from './indicators/prisma.health'; -export * from './indicators/redis.health'; -export * from './indicators/stellar.health'; -export * from './indicators/database-migration.health'; +export * from './health.module'; +export * from './health.controller'; +export * from './indicators/prisma.health'; +export * from './indicators/redis.health'; +export * from './indicators/stellar.health'; +export * from './indicators/database-migration.health'; +export * from './indicators/bullmq.health'; diff --git a/src/modules/health/indicators/bullmq.health.spec.ts b/src/modules/health/indicators/bullmq.health.spec.ts new file mode 100644 index 00000000..f613ba0b --- /dev/null +++ b/src/modules/health/indicators/bullmq.health.spec.ts @@ -0,0 +1,159 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { Queue } from 'bullmq'; +import { BullMQHealthIndicator } from './bullmq.health'; + +const makeQueue = (overrides: Record = {}) => + ({ + getJobCounts: vi.fn().mockResolvedValue({ + waiting: 0, + active: 0, + completed: 0, + failed: 0, + delayed: 0, + paused: 0, + ...overrides, + }), + close: vi.fn().mockResolvedValue(undefined), + }) as unknown as Queue; + +describe('BullMQHealthIndicator', () => { + let indicator: BullMQHealthIndicator; + + beforeEach(() => { + vi.clearAllMocks(); + indicator = new BullMQHealthIndicator(100); + }); + + it('monitors the core BullMQ queues', () => { + const queues = indicator.monitoredQueues; + expect(queues).toContain('webhooks'); + expect(queues).toContain('risk-analysis'); + expect(queues).toContain('transactions'); + expect(queues).toContain('dead-letter'); + expect(queues).toContain('audit'); + }); + + it('reports UP with job counts for all healthy queues', async () => { + const handle = makeQueue({ waiting: 3, active: 1, completed: 50, failed: 2 }); + indicator.setQueueHandle('webhooks', handle); + // Every queue without an explicit handle falls back to the same healthy mock. + vi.spyOn( + indicator as unknown as { getQueue: (name: string) => Queue }, + 'getQueue', + ).mockImplementation(() => handle); + + const report = await indicator.checkHealth(); + + expect(report.status).toBe('up'); + expect(report.redis).toBe('up'); + const webhooks = report.queues.find((q) => q.queue === 'webhooks'); + expect(webhooks?.connection).toBe('up'); + expect(webhooks?.counts).toEqual({ + waiting: 3, + active: 1, + completed: 50, + failed: 2, + delayed: 0, + paused: 0, + }); + expect(report.queues.length).toBe(indicator.monitoredQueues.length); + }); + + it('reports DEGRADED when a subset of queues fail their probes', async () => { + vi.spyOn( + indicator as unknown as { getQueue: (name: string) => Queue }, + 'getQueue', + ).mockImplementation((name: string) => { + if (name === 'webhooks') { + return makeQueue(); + } + // Every other queue times out. + return { + getJobCounts: vi.fn(() => new Promise(() => undefined)), // never resolves + close: vi.fn().mockResolvedValue(undefined), + } as unknown as Queue; + }); + + const report = await indicator.checkHealth(); + + expect(report.status).toBe('degraded'); + expect(report.redis).toBe('up'); + const webhooks = report.queues.find((q) => q.queue === 'webhooks'); + expect(webhooks?.error).toBeUndefined(); + const timedOut = report.queues.filter((q) => q.error?.includes('timed out')); + expect(timedOut.length).toBe(indicator.monitoredQueues.length - 1); + expect(timedOut[0]?.connection).toBe('down'); + }); + + it('reports DOWN when every queue probe fails', async () => { + vi.spyOn( + indicator as unknown as { getQueue: (name: string) => Queue }, + 'getQueue', + ).mockImplementation(() => { + return { + getJobCounts: vi.fn().mockRejectedValue(new Error('Connection is closed.')), + close: vi.fn().mockResolvedValue(undefined), + } as unknown as Queue; + }); + + const report = await indicator.checkHealth(); + + expect(report.status).toBe('down'); + expect(report.redis).toBe('down'); + for (const q of report.queues) { + expect(q.connection).toBe('down'); + expect(q.error).toContain('Connection is closed'); + expect(q.counts).toEqual({ waiting: 0, active: 0, completed: 0, failed: 0, delayed: 0, paused: 0 }); + } + }); + + it('probeQueue captures Redis errors per queue without throwing', async () => { + const failing = { + getJobCounts: vi.fn().mockRejectedValue(new Error('Redis timeout')), + close: vi.fn().mockResolvedValue(undefined), + } as unknown as Queue; + indicator.setQueueHandle('webhooks', failing); + + const status = await indicator.probeQueue('webhooks'); + + expect(status.connection).toBe('down'); + expect(status.error).toContain('Redis timeout'); + expect(status.counts.failed).toBe(0); + }); + + it('probeQueue returns normalized counts when the probe succeeds', async () => { + const healthy = makeQueue({ failed: 7, delayed: 4, paused: 1 }); + indicator.setQueueHandle('transactions', healthy); + + const status = await indicator.probeQueue('transactions'); + + expect(status.connection).toBe('up'); + expect(status.counts).toEqual({ + waiting: 0, + active: 0, + completed: 0, + failed: 7, + delayed: 4, + paused: 1, + }); + }); + + it('probes respect the configured timeout', async () => { + vi.useFakeTimers(); + try { + const hanging = { + getJobCounts: vi.fn(() => new Promise(() => undefined)), + close: vi.fn().mockResolvedValue(undefined), + } as unknown as Queue; + indicator.setQueueHandle('webhooks', hanging); + + const promise = indicator.probeQueue('webhooks'); + vi.advanceTimersByTime(101); + const status = await promise; + + expect(status.error).toContain('timed out after 100ms'); + } finally { + vi.useRealTimers(); + } + }); +}); diff --git a/src/modules/health/indicators/bullmq.health.ts b/src/modules/health/indicators/bullmq.health.ts new file mode 100644 index 00000000..ac6a22ab --- /dev/null +++ b/src/modules/health/indicators/bullmq.health.ts @@ -0,0 +1,199 @@ +import { Injectable, Logger } from '@nestjs/common'; +import { Queue } from 'bullmq'; +import { redisConfig } from '../../../config/redis.config'; +import { Queues } from '../../../queues/queues.constants'; + +export interface QueueHealthStatus { + /** Registered BullMQ queue name. */ + queue: string; + /** Connection status of the queue's underlying Redis client. */ + connection: 'up' | 'down'; + /** Error message when the queue could not be reached or probed timed out. */ + error?: string; + counts: { + waiting: number; + active: number; + completed: number; + failed: number; + delayed: number; + paused: number; + }; +} + +export interface QueuesHealthReport { + status: 'up' | 'down' | 'degraded'; + timestamp: string; + redis: 'up' | 'down'; + queues: QueueHealthStatus[]; +} + +/** + * Default timeout for probing a single queue, in milliseconds. A Redis + * connection that drops mid-probe can hang the underlying command; the probe + * is raced against this timeout so the health endpoint always answers before + * a Kubernetes liveness/readiness probe deadline. + */ +const PROBE_TIMEOUT_MS = 2_000; + +/** + * BullMQ queue health indicator. + * + * Inspects the core BullMQ queues (notifications, webhooks, stellar-sync, + * analytics, reports, outbox-events, stellar-fee-bump, transactions, + * risk-analysis, dead-letter, audit-cleanup, audit) and gathers waiting / + * active / failed / delayed / completed / paused job counts plus Redis + * connectivity, so Kubernetes (or any orchestrator) can monitor asynchronous + * processing health through one JSON endpoint. + * + * The indicator constructs lightweight Queue handles on the shared Redis + * connection settings and closes them after probing — it never starts + * workers. Failures and timeouts degrade gracefully: a queue that cannot be + * probed is reported as `down` with an error instead of throwing, and the + * overall report status becomes `degraded` (or `down` when Redis itself is + * unreachable). + */ +@Injectable() +export class BullMQHealthIndicator { + private readonly logger = new Logger(BullMQHealthIndicator.name); + private readonly queueHandles = new Map(); + + constructor(timeoutMs: number = PROBE_TIMEOUT_MS) { + this.probeTimeoutMs = timeoutMs; + } + + private readonly probeTimeoutMs: number; + + /** Queues monitored by this indicator. */ + get monitoredQueues(): string[] { + return Object.values(Queues); + } + + /** + * Probes every monitored queue and assembles the overall health report. + * Never throws — all failures are captured per queue in the report. + */ + async checkHealth(): Promise { + const queueStatuses = await Promise.all( + this.monitoredQueues.map((queueName) => this.probeQueue(queueName)), + ); + + // Redis is considered reachable when at least one probe succeeded — a + // per-queue timeout with other successes indicates a queue-level issue, + // not a Redis outage. + const anyQueueUp = queueStatuses.some((q) => q.connection === 'up'); + const failingQueues = queueStatuses.filter((q) => q.error !== undefined).length; + + let status: QueuesHealthReport['status'] = 'up'; + if (!anyQueueUp) { + status = 'down'; + } else if (failingQueues > 0) { + status = 'degraded'; + } + + return { + status, + timestamp: new Date().toISOString(), + redis: anyQueueUp ? 'up' : 'down', + queues: queueStatuses, + }; + } + + /** + * Probes a single queue with a hard timeout. Returns a structured status + * even when Redis is unreachable or the queue is unresponsive. + */ + async probeQueue(queueName: string): Promise { + let queue: Queue | undefined; + try { + queue = this.getQueue(queueName); + const counts = await this.withTimeout(queue.getJobCounts( + 'waiting', + 'active', + 'completed', + 'failed', + 'delayed', + 'paused', + )); + + return { + queue: queueName, + connection: 'up', + counts: { + waiting: counts.waiting ?? 0, + active: counts.active ?? 0, + completed: counts.completed ?? 0, + failed: counts.failed ?? 0, + delayed: counts.delayed ?? 0, + paused: counts.paused ?? 0, + }, + }; + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + this.logger.warn(`Queue health probe failed for '${queueName}': ${message}`); + return { + queue: queueName, + connection: 'down', + error: message, + counts: { waiting: 0, active: 0, completed: 0, failed: 0, delayed: 0, paused: 0 }, + }; + } + } + + /** + * Racy timeout wrapper — resolves with the probe result or rejects when the + * probe exceeds {@link probeTimeoutMs}. Guarantees the health endpoint + * responds even when Redis drops mid-command. + */ + private withTimeout(promise: Promise): Promise { + return new Promise((resolve, reject) => { + const timer = setTimeout( + () => reject(new Error(`Queue probe timed out after ${this.probeTimeoutMs}ms`)), + this.probeTimeoutMs, + ); + promise.then( + (value) => { + clearTimeout(timer); + resolve(value); + }, + (err: unknown) => { + clearTimeout(timer); + reject(err instanceof Error ? err : new Error(String(err))); + }, + ); + }); + } + + /** + * Lazily creates (or returns a previously created test-injected) queue + * handle. In tests the handles map may be pre-populated with mocks. + */ + protected getQueue(queueName: string): Queue { + return this.queueHandles.get(queueName) ?? this.createQueue(queueName); + } + + /** Registers (replaces) the handle used for a queue — used by tests. */ + setQueueHandle(queueName: string, queue: Queue): void { + const existing = this.queueHandles.get(queueName); + if (existing && existing !== queue) { + void existing.close().catch(() => undefined); + } + this.queueHandles.set(queueName, queue); + } + + private createQueue(queueName: string): Queue { + const { host, port, password, db } = redisConfig(); + const queue = new Queue(queueName, { + connection: { host, port, password: password || undefined, db }, + }); + this.queueHandles.set(queueName, queue); + return queue; + } + + async onModuleDestroy(): Promise { + await Promise.all( + Array.from(this.queueHandles.values()).map((queue) => + queue.close().catch(() => undefined), + ), + ); + } +} diff --git a/src/modules/webhooks/index.ts b/src/modules/webhooks/index.ts index 0165f55e..f606835c 100644 --- a/src/modules/webhooks/index.ts +++ b/src/modules/webhooks/index.ts @@ -1,2 +1,3 @@ export * from './webhook.service'; export * from './webhook.module'; +export * from './services/webhook-circuit-breaker.service'; diff --git a/src/modules/webhooks/services/webhook-circuit-breaker.service.spec.ts b/src/modules/webhooks/services/webhook-circuit-breaker.service.spec.ts new file mode 100644 index 00000000..02e6d8c3 --- /dev/null +++ b/src/modules/webhooks/services/webhook-circuit-breaker.service.spec.ts @@ -0,0 +1,155 @@ +import { beforeEach, afterEach, describe, expect, it, vi } from 'vitest'; +import { + WebhookCircuitBreakerService, + WebhookCircuitState, +} from './webhook-circuit-breaker.service'; + +/** + * The breaker is exercised through its in-memory fallback path (Redis is not + * available in unit tests), which mirrors the Redis semantics with identical + * state transitions. + */ +describe('WebhookCircuitBreakerService', () => { + let breaker: WebhookCircuitBreakerService; + const URL_A = 'https://failing.example.com/hook'; + const URL_B = 'https://healthy.example.com/hook'; + + beforeEach(() => { + vi.useFakeTimers(); + // Threshold 3 keeps trip scenarios short; the production default is 5. + breaker = new WebhookCircuitBreakerService(3, 10_000); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + const fail = (url: string) => breaker.recordFailure(url, new Error('HTTP 500')); + const succeed = (url: string) => breaker.recordSuccess(url); + + describe('state transitions', () => { + it('starts CLOSED and allows delivery', async () => { + const report = await breaker.getReport(URL_A); + expect(report.state).toBe(WebhookCircuitState.CLOSED); + expect(await breaker.isDeliveryAllowed(URL_A)).toBe(true); + }); + + it('stays CLOSED below the failure threshold', async () => { + await fail(URL_A); + await fail(URL_A); + + const report = await breaker.getReport(URL_A); + expect(report.state).toBe(WebhookCircuitState.CLOSED); + expect(report.consecutiveFailures).toBe(2); + expect(await breaker.isDeliveryAllowed(URL_A)).toBe(true); + }); + + it('trips OPEN after the threshold of consecutive failures', async () => { + await fail(URL_A); + await fail(URL_A); + await fail(URL_A); + + const report = await breaker.getReport(URL_A); + expect(report.state).toBe(WebhookCircuitState.OPEN); + expect(report.consecutiveFailures).toBe(3); + expect(await breaker.isDeliveryAllowed(URL_A)).toBe(false); + expect(report.remainingOpenMs).toBeGreaterThan(0); + }); + + it('tracks domains independently', async () => { + await fail(URL_A); + await fail(URL_A); + await fail(URL_A); + + expect((await breaker.getReport(URL_A)).state).toBe(WebhookCircuitState.OPEN); + expect((await breaker.getReport(URL_B)).state).toBe(WebhookCircuitState.CLOSED); + expect(await breaker.isDeliveryAllowed(URL_B)).toBe(true); + }); + + it('normalizes hosts so paths do not split the circuit', async () => { + await fail('https://example.com/a'); + await fail('https://example.com/b'); + await fail('https://example.com/c'); + + expect((await breaker.getReport('https://example.com/anything')).state).toBe( + WebhookCircuitState.OPEN, + ); + }); + + it('a success resets the consecutive failure counter while CLOSED', async () => { + await fail(URL_A); + await fail(URL_A); + await succeed(URL_A); + await fail(URL_A); + await fail(URL_A); + + expect((await breaker.getReport(URL_A)).state).toBe(WebhookCircuitState.CLOSED); + }); + + it('half-opens after the open-state TTL elapses and allows a trial', async () => { + await fail(URL_A); + await fail(URL_A); + await fail(URL_A); + expect(await breaker.isDeliveryAllowed(URL_A)).toBe(false); + + // Advance past the 10s open window. + vi.advanceTimersByTime(10_001); + + expect(await breaker.isDeliveryAllowed(URL_A)).toBe(true); + expect((await breaker.getReport(URL_A)).state).toBe(WebhookCircuitState.HALF_OPEN); + }); + + it('a failed trial immediately reopens the circuit', async () => { + await fail(URL_A); + await fail(URL_A); + await fail(URL_A); + vi.advanceTimersByTime(10_001); + + // Trial delivery goes through, then fails again. + expect(await breaker.isDeliveryAllowed(URL_A)).toBe(true); + await fail(URL_A); + + expect((await breaker.getReport(URL_A)).state).toBe(WebhookCircuitState.OPEN); + expect(await breaker.isDeliveryAllowed(URL_A)).toBe(false); + }); + + it('closes again after consecutive successful trials while HALF_OPEN', async () => { + await fail(URL_A); + await fail(URL_A); + await fail(URL_A); + vi.advanceTimersByTime(10_001); + + await breaker.isDeliveryAllowed(URL_A); // transition to HALF_OPEN + await succeed(URL_A); + expect((await breaker.getReport(URL_A)).state).toBe(WebhookCircuitState.HALF_OPEN); + await succeed(URL_A); + expect((await breaker.getReport(URL_A)).state).toBe(WebhookCircuitState.CLOSED); + expect((await breaker.getReport(URL_A)).consecutiveFailures).toBe(0); + }); + + it('reset force-closes the circuit', async () => { + await fail(URL_A); + await fail(URL_A); + await fail(URL_A); + expect((await breaker.getReport(URL_A)).state).toBe(WebhookCircuitState.OPEN); + + await breaker.reset(URL_A); + expect((await breaker.getReport(URL_A)).state).toBe(WebhookCircuitState.CLOSED); + expect(await breaker.isDeliveryAllowed(URL_A)).toBe(true); + }); + + it('getAllReports lists tracked domains', async () => { + await fail(URL_A); + const reports = await breaker.getAllReports(); + expect(reports.some((r) => r.host === 'failing.example.com')).toBe(true); + }); + }); + + describe('host extraction edge cases', () => { + it('falls back to a stable key for malformed URLs', async () => { + await expect(breaker.isDeliveryAllowed('not-a-url')).resolves.toBe(true); + const report = await breaker.getReport('not-a-url'); + expect(report.host).toBe('unknown'); + }); + }); +}); diff --git a/src/modules/webhooks/services/webhook-circuit-breaker.service.ts b/src/modules/webhooks/services/webhook-circuit-breaker.service.ts new file mode 100644 index 00000000..84381adf --- /dev/null +++ b/src/modules/webhooks/services/webhook-circuit-breaker.service.ts @@ -0,0 +1,334 @@ +import { Injectable, Logger } from '@nestjs/common'; +import Redis from 'ioredis'; +import { redisConfig } from '../../../config/redis.config'; + +/** How long a tripped (OPEN) circuit stays open before half-opening, ms. */ +const DEFAULT_OPEN_STATE_TTL_MS = 60_000; + +/** Consecutive failures required to trip the circuit for one domain. */ +export const DEFAULT_FAILURE_THRESHOLD = 5; + +/** Consecutive successes (while HALF_OPEN) required to close the circuit again. */ +const RECOVERY_SUCCESS_THRESHOLD = 2; + +/** Circuit states tracked per webhook endpoint domain. */ +export enum WebhookCircuitState { + /** Deliveries flow through normally; failures are counted. */ + CLOSED = 'CLOSED', + /** Deliveries to this domain fail fast until the open-state TTL elapses. */ + OPEN = 'OPEN', + /** Trial deliveries are allowed to probe whether the endpoint recovered. */ + HALF_OPEN = 'HALF_OPEN', +} + +export interface WebhookCircuitReport { + host: string; + state: WebhookCircuitState; + consecutiveFailures: number; + consecutiveSuccesses: number; + openedAt: number | null; + remainingOpenMs: number; +} + +interface DomainBreakerSnapshot { + failures: number; + successes: number; + state: WebhookCircuitState; + openedAt: number | null; +} + +/** + * Per-domain circuit breaker for outbound webhook deliveries. + * + * When one subscriber endpoint starts failing (transient outage, 5xx storm), + * the BullMQ webhook processor would otherwise keep hammering it with retries + * from every queued job — compounding the downstream failure and wasting + * worker concurrency. This service tracks consecutive failures keyed by the + * endpoint's host and "trips" the circuit once a threshold is exceeded, + * pausing deliveries to that domain only. + * + * State is stored in Redis so all worker replicas share one view per domain + * (with atomic INCR/GET/TTL semantics and automatic expiry of the OPEN + * window); when Redis is unavailable the service degrades to per-process + * in-memory tracking so circuit protection never disappears entirely. + * + * Recovery follows the classic three-state machine: + * CLOSED → (N consecutive failures) → OPEN → (TTL elapsed) → HALF_OPEN → + * (M consecutive successes) → CLOSED. A failure while HALF_OPEN reopens. + */ +@Injectable() +export class WebhookCircuitBreakerService { + private readonly logger = new Logger(WebhookCircuitBreakerService.name); + private readonly redis: Redis | null; + private readonly failureThreshold: number; + private readonly openStateTtlMs: number; + + /** In-memory fallback state, keyed by host, when Redis is unavailable. */ + private readonly memoryState = new Map(); + + constructor(failureThreshold: number = DEFAULT_FAILURE_THRESHOLD, openStateTtlMs: number = DEFAULT_OPEN_STATE_TTL_MS) { + this.failureThreshold = failureThreshold; + this.openStateTtlMs = openStateTtlMs; + + // Redis is optional: without it the breaker still works per-process. + try { + const { host, port, password, db } = redisConfig(); + this.redis = new Redis({ + host, + port, + password: password || undefined, + db, + lazyConnect: true, + maxRetriesPerRequest: 1, + retryStrategy: (times: number) => (times > 2 ? null : Math.min(times * 200, 1_000)), + enableOfflineQueue: false, + }); + this.redis.connect().catch((err: Error) => { + this.logger.warn(`Webhook circuit breaker Redis unavailable, using in-memory state: ${err.message}`); + }); + this.redis.on('error', () => { + /* Logged once via connect() catch; keep the breaker silent afterwards. */ + }); + } catch { + this.logger.warn('Redis not configured — webhook circuit breaker using in-memory state'); + this.redis = null; + } + } + + /** + * Determines whether a delivery to the given URL is allowed right now. + * + * - CLOSED: allowed. + * - OPEN within the TTL: blocked (fail fast — the processor short-circuits + * and lets BullMQ retry later, when the breaker may have half-opened). + * - OPEN past the TTL: transitions to HALF_OPEN and allows one trial. + * - HALF_OPEN: a trial delivery is allowed. + */ + async isDeliveryAllowed(url: string): Promise { + const host = this.extractHost(url); + const snap = await this.loadSnapshot(host); + + if (snap.state === WebhookCircuitState.OPEN) { + const elapsed = snap.openedAt !== null ? Date.now() - snap.openedAt : this.openStateTtlMs; + if (elapsed < this.openStateTtlMs) { + return false; + } + // Stale OPEN state (TTL elapsed) — move to HALF_OPEN for a trial call. + await this.saveSnapshot(host, { + ...snap, + state: WebhookCircuitState.HALF_OPEN, + successes: 0, + }); + this.logger.log(`Circuit for ${host} half-opened; allowing trial delivery`); + return true; + } + + return true; + } + + /** + * Records a successful delivery against the endpoint's domain. + * While HALF_OPEN, enough consecutive successes close the circuit again. + */ + async recordSuccess(url: string): Promise { + const host = this.extractHost(url); + const snap = await this.loadSnapshot(host); + + if (snap.state === WebhookCircuitState.CLOSED && snap.failures === 0) { + return; // Nothing to reset — avoid a redundant write. + } + + if (snap.state === WebhookCircuitState.HALF_OPEN) { + const successes = snap.successes + 1; + if (successes >= RECOVERY_SUCCESS_THRESHOLD) { + this.logger.log(`Circuit for ${host} closed again after recovery`); + await this.saveSnapshot(host, { + failures: 0, + successes: 0, + state: WebhookCircuitState.CLOSED, + openedAt: null, + }); + return; + } + await this.saveSnapshot(host, { ...snap, successes }); + return; + } + + // CLOSED: a success resets the consecutive failure counter. + await this.saveSnapshot(host, { ...snap, failures: 0 }); + } + + /** + * Records a failed delivery against the endpoint's domain. + * Trips the circuit OPEN when consecutive failures reach the threshold. + */ + async recordFailure(url: string, error?: unknown): Promise { + const host = this.extractHost(url); + const snap = await this.loadSnapshot(host); + + const failures = + snap.state === WebhookCircuitState.HALF_OPEN + ? this.failureThreshold // A failed trial immediately reopens. + : snap.failures + 1; + + if (failures >= this.failureThreshold && snap.state !== WebhookCircuitState.OPEN) { + this.logger.warn( + `Circuit for ${host} OPEN after ${failures} consecutive failures` + + `${error ? `: ${(error as Error)?.message ?? String(error)}` : ''}`, + ); + await this.saveSnapshot(host, { + failures, + successes: 0, + state: WebhookCircuitState.OPEN, + openedAt: Date.now(), + }); + return; + } + + await this.saveSnapshot(host, { ...snap, failures, successes: 0 }); + } + + /** Force-closes the circuit for a domain (administrative override / tests). */ + async reset(url: string): Promise { + const host = this.extractHost(url); + await this.saveSnapshot(host, { + failures: 0, + successes: 0, + state: WebhookCircuitState.CLOSED, + openedAt: null, + }); + } + + /** Returns the current circuit report for a URL's host (observability). */ + async getReport(url: string): Promise { + const host = this.extractHost(url); + const snap = await this.loadSnapshot(host); + const remainingOpenMs = + snap.state === WebhookCircuitState.OPEN && snap.openedAt !== null + ? Math.max(this.openStateTtlMs - (Date.now() - snap.openedAt), 0) + : 0; + + return { + host, + state: snap.state, + consecutiveFailures: snap.failures, + consecutiveSuccesses: snap.successes, + openedAt: snap.openedAt, + remainingOpenMs, + }; + } + + /** Reports for every domain the breaker currently tracks. */ + async getAllReports(): Promise { + if (this.redis && this.redis.status === 'ready') { + try { + const keys = await this.redis.keys(`${this.keyPrefix()}*`); + const reports: WebhookCircuitReport[] = []; + for (const key of keys) { + const host = key.slice(this.keyPrefix().length); + const raw = await this.redis.get(key); + if (!raw) continue; + const snap = JSON.parse(raw) as DomainBreakerSnapshot; + reports.push({ + host, + state: snap.state, + consecutiveFailures: snap.failures, + consecutiveSuccesses: snap.successes, + openedAt: snap.openedAt, + remainingOpenMs: + snap.state === WebhookCircuitState.OPEN && snap.openedAt !== null + ? Math.max(this.openStateTtlMs - (Date.now() - snap.openedAt), 0) + : 0, + }); + } + return reports; + } catch { + /* fall through to memory */ + } + } + return Array.from(this.memoryState.entries()).map(([host, snap]) => ({ + host, + state: snap.state, + consecutiveFailures: snap.failures, + consecutiveSuccesses: snap.successes, + openedAt: snap.openedAt, + remainingOpenMs: + snap.state === WebhookCircuitState.OPEN && snap.openedAt !== null + ? Math.max(this.openStateTtlMs - (Date.now() - snap.openedAt), 0) + : 0, + })); + } + + onModuleDestroy(): void { + this.redis?.disconnect(); + } + + /** Strips the scheme, path and port-independence: tracks by hostname only. */ + private extractHost(url: string): string { + try { + const parsed = new URL(url); + // Include the port so distinct services on one host don't share a fuse, + // but normalize the default ports so http(s)://host traffic matches. + return parsed.port && parsed.port !== '80' && parsed.port !== '443' + ? `${parsed.hostname}:${parsed.port}` + : parsed.hostname; + } catch { + // Malformed URLs still get a stable (if coarse) key. + return 'unknown'; + } + } + + private keyPrefix(): string { + return 'webhook:circuit:'; + } + + private stateKey(host: string): string { + return `${this.keyPrefix()}${host}`; + } + + private async loadSnapshot(host: string): Promise { + if (this.redis && this.redis.status === 'ready') { + try { + const raw = await this.redis.get(this.stateKey(host)); + if (raw) { + const parsed = JSON.parse(raw) as DomainBreakerSnapshot; + // Guard against stale/foreign shapes in Redis. + if (typeof parsed?.failures === 'number' && parsed?.state) { + return parsed; + } + } + return this.emptySnapshot(); + } catch { + /* fall through to memory */ + } + } + return this.memoryState.get(host) ?? this.emptySnapshot(); + } + + private async saveSnapshot(host: string, snap: DomainBreakerSnapshot): Promise { + if (this.redis && this.redis.status === 'ready') { + try { + // Expire slightly after the OPEN window so stale states self-clean; + // CLOSED snapshots keep a short TTL to bound memory usage. + const ttlSeconds = + snap.state === WebhookCircuitState.OPEN + ? Math.ceil(this.openStateTtlMs / 1_000) + 30 + : 300; + await this.redis.set(this.stateKey(host), JSON.stringify(snap), 'EX', ttlSeconds); + return; + } catch { + /* fall through to memory */ + } + } + this.memoryState.set(host, snap); + } + + private emptySnapshot(): DomainBreakerSnapshot { + return { + failures: 0, + successes: 0, + state: WebhookCircuitState.CLOSED, + openedAt: null, + }; + } +} diff --git a/src/modules/webhooks/webhook.module.ts b/src/modules/webhooks/webhook.module.ts index 884316c9..66625c20 100644 --- a/src/modules/webhooks/webhook.module.ts +++ b/src/modules/webhooks/webhook.module.ts @@ -7,23 +7,24 @@ import { WebhookService } from './webhook.service'; import { WebhookRepository } from './webhook.repository'; import { WebhookDispatcher } from './webhook.dispatcher'; import { WebhookDeliveryService } from './services/webhook-delivery.service'; -import { WebhookAuditService } from './services/webhook-audit.service'; +import { WebhookCircuitBreakerService } from './services/webhook-circuit-breaker.service'; import { WebhooksProcessor } from './webhooks.processor'; -import { createWebhookQueueOptions } from '../../queues/webhook.queue'; +import { Queues } from '../../queues/queues.constants'; import { redisConfig } from '../../config/redis.config'; +import { webhookBackoffStrategy } from '../../utils/backoff.util'; import { MetricsModule } from '../metrics/metrics.module'; import { RawBodyMiddleware } from '../../common/middleware/raw-body.middleware'; import { WebhookSignatureGuard } from '../../common/guards/webhook-signature.guard'; import { SlidingWindowThrottlerGuard } from '../../common/guards/sliding-window-throttler.guard'; +import type { RegisterQueueOptions } from '@nestjs/bullmq'; /** - * Webhooks module. The dispatcher listens to domain events and queues the - * curated WEBHOOK_EVENTS set to subscribed external endpoints via BullMQ. + * Webhooks module. The dispatcher listens to domain events and queues + * the curated WEBHOOK_EVENTS set to subscribed external endpoints via BullMQ. * - * Retry policy (5 attempts, exponential backoff with jitter) and the - * dead-letter routing live in `@queues/webhook.queue`, so the API and the worker - * can never drift apart. `WebhookAuditService` records deliveries that exhaust - * their retries in the compliance audit trail. + * Uses a custom backoffStrategy with randomized jitter (20% of base delay) + * to prevent thundering herd problems when multiple webhook deliveries + * are retried simultaneously. */ @Module({ imports: [ @@ -35,7 +36,24 @@ import { SlidingWindowThrottlerGuard } from '../../common/guards/sliding-window- db: redisConfig().db, }, }), - BullModule.registerQueue(createWebhookQueueOptions()), + BullModule.registerQueue({ + name: Queues.Webhooks, + defaultJobOptions: { + attempts: 5, + backoff: { + type: 'exponential', + delay: 2000, + }, + removeOnComplete: { count: 1000 }, + removeOnFail: { age: 24 * 3600 }, + }, + // BullMQ reads queue.opts.settings.backoffStrategy at retry time. + // The AdvancedOptions type is not fully exposed by @nestjs/bullmq, so we + // cast to include the backoffStrategy field that BullMQ supports at runtime. + settings: { + backoffStrategy: webhookBackoffStrategy, + } as RegisterQueueOptions['settings'], + }), MetricsModule, ], controllers: [WebhookController, WebhookIngressController], @@ -44,13 +62,13 @@ import { SlidingWindowThrottlerGuard } from '../../common/guards/sliding-window- WebhookRepository, WebhookDispatcher, WebhookDeliveryService, - WebhookAuditService, + WebhookCircuitBreakerService, WebhooksProcessor, WebhookIngressService, WebhookSignatureGuard, SlidingWindowThrottlerGuard, ], - exports: [WebhookService], + exports: [WebhookService, WebhookCircuitBreakerService], }) export class WebhookModule implements NestModule { configure(consumer: MiddlewareConsumer): void { diff --git a/src/modules/webhooks/webhooks.processor.spec.ts b/src/modules/webhooks/webhooks.processor.spec.ts index 362e21e5..a56cae65 100644 --- a/src/modules/webhooks/webhooks.processor.spec.ts +++ b/src/modules/webhooks/webhooks.processor.spec.ts @@ -2,7 +2,10 @@ import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest'; import { Job, UnrecoverableError } from 'bullmq'; import { WebhooksProcessor } from './webhooks.processor'; import { WebhookJobData } from './types/webhook-job.types'; -import { verifyWebhookSignature } from './utils/signing'; +import { createHmac } from 'crypto'; +import { + WebhookCircuitBreakerService, +} from './services/webhook-circuit-breaker.service'; describe('WebhooksProcessor', () => { let processor: WebhooksProcessor; @@ -19,6 +22,7 @@ describe('WebhooksProcessor', () => { webhookId: WEBHOOK_ID, organizationId: ORG_ID, url: WEBHOOK_URL, + secret: WEBHOOK_SECRET, eventName: 'transaction.completed', payload: { event: 'transaction.completed', data: { transactionId: 'txn-123' } }, eventId: EVENT_ID, @@ -34,9 +38,7 @@ describe('WebhooksProcessor', () => { }) as unknown as Job; beforeEach(() => { - mockPrisma = { - webhook: { findFirst: vi.fn().mockResolvedValue({ secret: WEBHOOK_SECRET }) }, - }; + mockPrisma = {}; // Access private property via type assertion processor = new WebhooksProcessor(mockPrisma as never); fetchSpy = vi.fn(); @@ -70,13 +72,14 @@ describe('WebhooksProcessor', () => { expect(options.headers['x-astroid-delivery']).toBe(EVENT_ID); expect(options.headers['x-astroid-timestamp']).toBeDefined(); expect(options.headers['x-astroid-timestamp']).toMatch(/^\d+$/); - expect(options.headers['x-astroid-signature-version']).toBe('v1'); - const body = Buffer.from(options.body); + // Verify HMAC-SHA256 signature = HMAC(secret, timestamp + body) + const body = options.body; const timestamp = options.headers['x-astroid-timestamp']; - expect(options.headers['x-astroid-signature']).toMatch(/^v1=[0-9a-f]{64}$/); - expect(verifyWebhookSignature(WEBHOOK_SECRET, timestamp, EVENT_ID, body, options.headers['x-astroid-signature'])).toBe(true); - expect(job.data).not.toHaveProperty('secret'); + const expectedSignature = createHmac('sha256', WEBHOOK_SECRET) + .update(`${timestamp}.${body}`) + .digest('hex'); + expect(options.headers['x-astroid-signature']).toBe(expectedSignature); }); it('returns success result with status code', async () => { @@ -234,34 +237,92 @@ describe('WebhooksProcessor', () => { const job = createMockJob({ attemptsMade: 4 } as Partial>); await expect(processor.process(job)).rejects.toThrow('HTTP 503'); }); + }); - it('keeps the same event identity across retry attempts', async () => { - const upsert = vi.fn().mockResolvedValue({}); - mockPrisma.webhookDelivery = { upsert }; - fetchSpy.mockResolvedValue({ ok: true, status: 200, text: () => Promise.resolve('OK') }); - - const firstAttempt = createMockJob(); - const retryAttempt = createMockJob({ attemptsMade: 1 } as Partial>); - await processor.process(firstAttempt); - await processor.process(retryAttempt); - - expect(firstAttempt.data.eventId).toBe(retryAttempt.data.eventId); - expect(upsert.mock.calls[0][0].where).toEqual(upsert.mock.calls[1][0].where); + describe('per-domain circuit breaker (issue #219)', () => { + let breaker: WebhookCircuitBreakerService; + + beforeEach(() => { + breaker = new WebhookCircuitBreakerService(2, 10_000); + processor = new WebhooksProcessor( + mockPrisma as never, + undefined, + undefined, + breaker, + ); }); - it('does not include a downstream response body in failure messages or logs', async () => { - const secretEcho = `${WEBHOOK_SECRET}:${EVENT_ID}:payload`; - const warn = vi.spyOn(processor['logger'], 'warn').mockImplementation(() => undefined); - const error = vi.spyOn(processor['logger'], 'error').mockImplementation(() => undefined); + it('records failures and trips the circuit after consecutive HTTP 500s', async () => { fetchSpy.mockResolvedValue({ ok: false, status: 500, - text: () => Promise.resolve(secretEcho), + text: () => Promise.resolve('Internal Server Error'), }); await expect(processor.process(createMockJob())).rejects.toThrow('HTTP 500'); - expect(warn.mock.calls.flat().join(' ')).not.toContain(secretEcho); - expect(error.mock.calls.flat().join(' ')).not.toContain(secretEcho); + await expect(processor.process(createMockJob())).rejects.toThrow('HTTP 500'); + + const report = await breaker.getReport(WEBHOOK_URL); + expect(report.state).toBe('OPEN'); + expect(report.consecutiveFailures).toBe(2); + }); + + it('fail-fasts without calling fetch while the circuit is OPEN', async () => { + await breaker.recordFailure(WEBHOOK_URL, new Error('HTTP 500')); + await breaker.recordFailure(WEBHOOK_URL, new Error('HTTP 500')); + + await expect(processor.process(createMockJob())).rejects.toThrow('Circuit open'); + expect(fetchSpy).not.toHaveBeenCalled(); + }); + + it('records successes and keeps the circuit CLOSED on healthy delivery', async () => { + fetchSpy.mockResolvedValue({ ok: true, status: 200, text: () => Promise.resolve('OK') }); + + await processor.process(createMockJob()); + + const report = await breaker.getReport(WEBHOOK_URL); + expect(report.state).toBe('CLOSED'); + expect(report.consecutiveFailures).toBe(0); + }); + + it('allows delivery again after the circuit half-opens and delivery succeeds', async () => { + vi.useFakeTimers(); + try { + fetchSpy.mockResolvedValue({ + ok: false, + status: 500, + text: () => Promise.resolve('Internal Server Error'), + }); + await expect(processor.process(createMockJob())).rejects.toThrow('HTTP 500'); + await expect(processor.process(createMockJob())).rejects.toThrow('HTTP 500'); + expect((await breaker.getReport(WEBHOOK_URL)).state).toBe('OPEN'); + + // Advance past the open window, then deliver successfully. + vi.advanceTimersByTime(10_001); + fetchSpy.mockResolvedValue({ ok: true, status: 200, text: () => Promise.resolve('OK') }); + expect((await breaker.isDeliveryAllowed(WEBHOOK_URL)) as boolean).toBe(true); + await processor.process(createMockJob()); + + const report = await breaker.getReport(WEBHOOK_URL); + expect(['HALF_OPEN', 'CLOSED']).toContain(report.state); + } finally { + vi.useRealTimers(); + } + }); + + it('deliveries to a different domain are unaffected by an OPEN circuit', async () => { + fetchSpy.mockResolvedValue({ ok: true, status: 200, text: () => Promise.resolve('OK') }); + + await breaker.recordFailure(WEBHOOK_URL, new Error('HTTP 500')); + await breaker.recordFailure(WEBHOOK_URL, new Error('HTTP 500')); + + const otherUrlJob = createMockJob({ + data: createJobData({ url: 'https://other-domain.com/hook' }), + }); + const result = await processor.process(otherUrlJob); + expect(result.success).toBe(true); + expect(fetchSpy).toHaveBeenCalledOnce(); + expect(fetchSpy.mock.calls[0][0]).toBe('https://other-domain.com/hook'); }); }); }); diff --git a/src/modules/webhooks/webhooks.processor.ts b/src/modules/webhooks/webhooks.processor.ts index ef37c5d2..20e2d81c 100644 --- a/src/modules/webhooks/webhooks.processor.ts +++ b/src/modules/webhooks/webhooks.processor.ts @@ -1,20 +1,13 @@ import { Processor, WorkerHost } from '@nestjs/bullmq'; +import { ConfigService } from '@nestjs/config'; import { Inject, Logger, Optional } from '@nestjs/common'; import { Job, UnrecoverableError } from 'bullmq'; import { Queues } from '../../queues/queues.constants'; import { WebhookJobData, WebhookJobResult } from './types/webhook-job.types'; -import { signWebhookPayload, WEBHOOK_SIGNATURE_VERSION } from './utils/signing'; +import { signWebhookPayload } from './utils/signing'; import { PrismaService } from '../../database/prisma.service'; import { WorkerMetricsService } from '../../modules/metrics/worker-metrics.service'; -import { WebhookAuditService } from './services/webhook-audit.service'; -import { - WEBHOOK_DELIVERY_HEADER, - WEBHOOK_EVENT_HEADER, - WEBHOOK_EVENT_ID_HEADER, - WEBHOOK_SIGNATURE_HEADER, - WEBHOOK_SIGNATURE_VERSION_HEADER, - WEBHOOK_TIMESTAMP_HEADER, -} from '../../common/constants/headers'; +import { WebhookCircuitBreakerService } from './services/webhook-circuit-breaker.service'; /** * BullMQ job processor for webhook event delivery with exponential backoff + jitter. @@ -26,6 +19,11 @@ import { * - Persistent delivery status tracking (PENDING → RETRYING → FAILED/DELIVERED) * - Fail-safe: retry failures never crash the master process * + * Per-domain circuit breaker (issue #219): + * - Consecutive downstream failures (per endpoint host) trip a circuit that + * fail-fasts further deliveries to that domain until it half-opens again + * - Successes reset the failure counter; recovered domains resume delivery + * * Jitter is applied via a custom backoffStrategy configured on the BullMQ * queue registration (see webhook.module.ts). BullMQ reads the strategy from * queue.opts.settings.backoffStrategy at retry time. @@ -62,81 +60,67 @@ export class WebhooksProcessor extends WorkerHost implements OnModuleDestroy { constructor( @Optional() @Inject(PrismaService) private readonly prisma?: PrismaService, + @Optional() private readonly configService?: ConfigService, @Optional() private readonly workerMetrics?: WorkerMetricsService, - @Optional() private readonly webhookAudit?: WebhookAuditService, + @Optional() private readonly circuitBreaker?: WebhookCircuitBreakerService, ) { super(); } - /** - * Audit entry for a delivery that will not be retried again: an unrecoverable - * 4xx or the final attempt. `WebhookAuditService` swallows its own failures, so - * this can never mask the original delivery error. - */ - private async auditTerminalFailure( - job: Job, - failedReason: string, - responseStatus?: number, - ): Promise { - if (!this.webhookAudit) return; - try { - await this.webhookAudit.recordTerminalFailure({ - webhookId: job.data.webhookId, - organizationId: job.data.organizationId, - url: job.data.url, - eventName: job.data.eventName, - eventId: job.data.eventId, - attemptsMade: job.attemptsMade + 1, - failedReason, - responseStatus, - }); - } catch (error) { - // Never let compliance bookkeeping mask the original delivery failure. - this.logger.warn( - `Could not audit webhook ${job.data.webhookId} failure: ${(error as Error).message}`, - ); - } - } - - private async resolveSecret(webhookId: string, organizationId: string): Promise { - const client = this.prisma?.workerClient ?? this.prisma; - if (!client) throw new UnrecoverableError('Webhook signing secret is unavailable'); - const webhook = await client.webhook.findFirst({ - where: { id: webhookId, organizationId }, - select: { secret: true }, - }); - if (!webhook?.secret) throw new UnrecoverableError('Webhook signing secret is unavailable'); - return webhook.secret; + private resolveSecret(jobSecret?: string): string { + if (jobSecret) return jobSecret; + const fallback = + this.configService?.get('WEBHOOK_SECRET') ?? + this.configService?.get('STELLAR_WEBHOOK_SECRET') ?? + this.configService?.get('WEBHOOK_SIGNING_SECRET') ?? + ''; + return fallback; } async process(job: Job): Promise { const jobName = job.name ?? 'webhook-delivery'; const execute = async (): Promise => { - const { webhookId, organizationId, url, eventName, payload, eventId, metadata } = job.data; - const requestTrace = metadata?.requestId ? ` requestId=${metadata.requestId}` : ''; - this.logger.debug(`Processing webhook ${webhookId} event ${eventName} attempt ${job.attemptsMade + 1}/5${requestTrace}`); + const { webhookId, organizationId, url, secret, eventName, payload, eventId } = job.data; + + // Circuit breaker (issue #219): fail fast when this endpoint's domain is + // OPEN. The error surfaces as a regular (transient) failure so BullMQ + // retries with backoff — by which time the breaker may have half-opened. + if (this.circuitBreaker) { + const allowed = await this.circuitBreaker.isDeliveryAllowed(url); + if (!allowed) { + const report = await this.circuitBreaker.getReport(url); + this.logger.warn( + `Circuit OPEN for ${report.host} — skipping webhook ${webhookId} delivery ` + + `(${report.remainingOpenMs}ms until half-open trial)`, + ); + throw new Error( + `Circuit open for ${report.host}; delivery paused until recovery trial`, + ); + } + } + + this.logger.debug(`Processing webhook ${webhookId} event ${eventName} attempt ${job.attemptsMade + 1}/5`); let responseStatus: number | undefined; let errorMessage: string | undefined; let isNonTransient = false; try { - const body = Buffer.from(JSON.stringify(payload), 'utf8'); + const body = JSON.stringify(payload); const timestamp = Math.floor(Date.now() / 1000).toString(); - const effectiveSecret = await this.resolveSecret(webhookId, organizationId); - const signature = signWebhookPayload(effectiveSecret, timestamp, eventId, body); + const effectiveSecret = this.resolveSecret(secret); + const signature = signWebhookPayload(effectiveSecret, timestamp, body); const response = await fetch(url, { method: 'POST', headers: { 'content-type': 'application/json', - [WEBHOOK_SIGNATURE_HEADER]: signature, - [WEBHOOK_TIMESTAMP_HEADER]: timestamp, - [WEBHOOK_EVENT_ID_HEADER]: eventId, - [WEBHOOK_DELIVERY_HEADER]: eventId, - [WEBHOOK_EVENT_HEADER]: eventName, - [WEBHOOK_SIGNATURE_VERSION_HEADER]: WEBHOOK_SIGNATURE_VERSION, + 'x-astroid-signature': signature, + 'x-astroid-timestamp': timestamp, + 'x-astroid-delivery': eventId, + 'x-astroid-event': eventName, + 'x-astroid-event-id': eventId, 'user-agent': 'Astroid-Webhook-Bot/1.0', }, body, @@ -145,9 +129,10 @@ export class WebhooksProcessor extends WorkerHost implements OnModuleDestroy { responseStatus = response.status; if (!response.ok) { - errorMessage = `HTTP ${response.status}`; + const errorText = await response.text().catch(() => response.statusText); + errorMessage = `HTTP ${response.status}: ${errorText}`; isNonTransient = WebhooksProcessor.NON_TRANSIENT_STATUSES.has(response.status); - this.logger.warn(`Webhook ${webhookId} responded ${response.status}${requestTrace}`); + this.logger.warn(`Webhook ${webhookId} responded ${response.status}: ${errorText}`); if (isNonTransient) { await this.persistState({ webhookId, @@ -160,19 +145,16 @@ export class WebhooksProcessor extends WorkerHost implements OnModuleDestroy { lastError: errorMessage, responseStatus, }); - // Non-transient (4xx): record the abandoned delivery before BullMQ - // moves it straight to the failed set. - await this.auditTerminalFailure(job, errorMessage ?? 'HTTP error', responseStatus); throw new UnrecoverableError(errorMessage); } throw new Error(errorMessage); } - this.logger.debug(`Webhook ${webhookId} delivered successfully${requestTrace}`); + this.logger.debug(`Webhook ${webhookId} delivered successfully`); } catch (error) { if (error instanceof UnrecoverableError) throw error; - errorMessage = error instanceof Error ? error.message : 'Delivery attempt failed'; + errorMessage = (error as Error).message; const isLastAttempt = job.attemptsMade >= 4; - this.logger.error(`Webhook ${webhookId} failed attempt ${job.attemptsMade + 1}/5${requestTrace}`); + this.logger.error(`Webhook ${webhookId} failed attempt ${job.attemptsMade + 1}/5: ${errorMessage}`); await this.persistState({ webhookId, organizationId, @@ -185,14 +167,29 @@ export class WebhooksProcessor extends WorkerHost implements OnModuleDestroy { responseStatus, }); if (isLastAttempt) { - this.logger.error(`Webhook ${webhookId} exhausted all retry attempts${requestTrace}`); - // Retries are exhausted: the delivery is dead-lettered by the queue - // failure listener, so record it permanently in the audit trail. - await this.auditTerminalFailure(job, errorMessage ?? 'unknown error', responseStatus); + this.logger.error(`Webhook ${webhookId} exhausted all retry attempts`); + } + // Record the failure with the per-domain circuit breaker (issue #219). + if (this.circuitBreaker) { + try { + await this.circuitBreaker.recordFailure(url, error); + } catch (err) { + this.logger.warn(`Circuit breaker recordFailure failed: ${(err as Error).message}`); + } } throw error; } + // Record the outcome with the per-domain circuit breaker (issue #219). + // Best-effort: breaker bookkeeping failures must never affect delivery. + if (this.circuitBreaker) { + try { + await this.circuitBreaker.recordSuccess(url); + } catch (err) { + this.logger.warn(`Circuit breaker recordSuccess failed: ${(err as Error).message}`); + } + } + await this.persistState({ webhookId, organizationId, diff --git a/src/queues/queues.constants.ts b/src/queues/queues.constants.ts index 6cbd4d43..9db32765 100644 --- a/src/queues/queues.constants.ts +++ b/src/queues/queues.constants.ts @@ -25,6 +25,8 @@ export const Queues = { DeadLetter: 'dead-letter', /** Audit logs cleanup job. */ AuditCleanup: 'audit-cleanup', + /** Asynchronous audit log persistence (batched writes to the audit trail). */ + Audit: 'audit', } as const; export type QueueName = (typeof Queues)[keyof typeof Queues]; @@ -57,3 +59,23 @@ export interface DlqJobData { metadata?: Record; } +/** A single audit log entry queued for asynchronous persistence. */ +export interface AuditLogJobEntry { + organizationId: string; + userId?: string | null; + action: string; + entity: string; + entityId?: string | null; + oldValue?: unknown; + newValue?: unknown; + ipAddress?: string | null; + device?: string | null; + requestId?: string | null; +} + +/** Standard BullMQ job payload for the audit persistence queue. */ +export interface AuditJobData { + /** Batch of audit entries to persist in one transaction. */ + entries: AuditLogJobEntry[]; +} + diff --git a/src/workers/audit.worker.spec.ts b/src/workers/audit.worker.spec.ts new file mode 100644 index 00000000..5f1bfed0 --- /dev/null +++ b/src/workers/audit.worker.spec.ts @@ -0,0 +1,180 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { AuditWorker } from './audit.worker'; +import { AuditJobData } from '../queues/queues.constants'; +import { AuditHashService } from '../modules/audit/audit-hash.service'; + +type CreateArgs = { + organizationId: string; + action: string; + entity: string; + entityId?: string | null; + userId?: string | null; + requestId?: string | null; + previousHash?: string | null; + hash?: string | null; +}; + +const createJob = (data: AuditJobData, attemptsMade = 0, attempts = 5) => + ({ + id: 'job-1', + name: 'audit-persist', + data, + attemptsMade, + opts: { attempts }, + }) as never; + +describe('AuditWorker', () => { + let auditLogCreate: ReturnType; + let workerClient: { $transaction: ReturnType }; + let prisma: { workerClient: unknown }; + let hashService: { + getLatestHash: ReturnType; + computeEntryHash: ReturnType; + }; + let worker: AuditWorker; + + beforeEach(() => { + vi.clearAllMocks(); + + auditLogCreate = vi.fn().mockResolvedValue({}); + workerClient = { + $transaction: vi.fn(async (fn: (tx: unknown) => Promise) => + fn({ auditLog: { create: auditLogCreate } }), + ), + }; + prisma = { workerClient }; + + hashService = { + getLatestHash: vi.fn().mockResolvedValue('prev-hash'), + computeEntryHash: vi + .fn() + .mockImplementation((input: { organizationId: string }) => ({ + previousHash: 'prev-hash', + hash: `hash-of-${input.organizationId}`, + })), + }; + + worker = new AuditWorker( + prisma as never, + hashService as unknown as AuditHashService, + ); + }); + + it('persists a batch of entries in a single transaction', async () => { + const data: AuditJobData = { + entries: [ + { + organizationId: 'org-1', + action: 'POLICY_CHECK', + entity: 'Policy', + entityId: 'pol-1', + }, + { + organizationId: 'org-1', + action: 'AUTH_ATTEMPT', + entity: 'User', + entityId: 'user-1', + }, + { + organizationId: 'org-1', + action: 'RISK_EVALUATED', + entity: 'Transaction', + entityId: 'txn-1', + }, + ], + }; + + const result = await worker.process(createJob(data)); + + expect(workerClient.$transaction).toHaveBeenCalledOnce(); + expect(auditLogCreate).toHaveBeenCalledTimes(3); + expect(result).toEqual({ persisted: 3 }); + }); + + it('chains each entry to the previous hash within the batch', async () => { + const data: AuditJobData = { + entries: [ + { organizationId: 'org-1', action: 'A', entity: 'E1' }, + { organizationId: 'org-1', action: 'B', entity: 'E2' }, + ], + }; + + await worker.process(createJob(data)); + + const firstCall = auditLogCreate.mock.calls[0][0] as { data: CreateArgs }; + const secondCall = auditLogCreate.mock.calls[1][0] as { data: CreateArgs }; + + expect(firstCall.data.previousHash).toBe('prev-hash'); + expect(secondCall.data.previousHash).toBe('hash-of-org-1'); + }); + + it('records null hash fields when the hash service is unavailable', async () => { + const fallbackWorker = new AuditWorker(prisma as never); + const data: AuditJobData = { + entries: [{ organizationId: 'org-1', action: 'A', entity: 'E1' }], + }; + + await fallbackWorker.process(createJob(data)); + + const call = auditLogCreate.mock.calls[0][0] as { data: CreateArgs }; + expect(call.data.previousHash).toBeNull(); + expect(call.data.hash).toBeNull(); + }); + + it('normalizes optional fields to null on insert', async () => { + const data: AuditJobData = { + entries: [{ organizationId: 'org-1', action: 'A', entity: 'E1' }], + }; + + await worker.process(createJob(data)); + + const call = auditLogCreate.mock.calls[0][0] as { data: CreateArgs }; + expect(call.data.userId).toBeNull(); + expect(call.data.requestId).toBeNull(); + }); + + it('returns zero persisted without touching the database for an empty batch', async () => { + const result = await worker.process(createJob({ entries: [] })); + + expect(result).toEqual({ persisted: 0 }); + expect(workerClient.$transaction).not.toHaveBeenCalled(); + }); + + it('rethrows persistence failures so BullMQ retries with backoff', async () => { + workerClient.$transaction.mockRejectedValue(new Error('Connection terminated')); + + await expect( + worker.process( + createJob({ entries: [{ organizationId: 'org-1', action: 'A', entity: 'E1' }] }), + ), + ).rejects.toThrow('Connection terminated'); + }); + + it('survives (does not crash) when the database write fails on the final attempt', async () => { + workerClient.$transaction.mockRejectedValue(new Error('write failure')); + + // Final attempt: attemptsMade is the last of 5 attempts. The error still + // propagates to BullMQ, but process() itself must not throw anything other + // than the original error (no unhandled crash / process exit). + const promise = worker.process( + createJob({ entries: [{ organizationId: 'org-1', action: 'A', entity: 'E1' }] }, 4, 5), + ); + await expect(promise).rejects.toThrow('write failure'); + }); + + it('uses the direct prisma client when no dedicated worker client exists', async () => { + const plainPrisma = { + workerClient: undefined, + $transaction: vi.fn(async (fn: (tx: unknown) => Promise) => + fn({ auditLog: { create: auditLogCreate } }), + ), + }; + const fallbackWorker = new AuditWorker(plainPrisma as never); + + await fallbackWorker.process( + createJob({ entries: [{ organizationId: 'org-1', action: 'A', entity: 'E1' }] }), + ); + + expect(plainPrisma.$transaction).toHaveBeenCalledOnce(); + }); +}); diff --git a/src/workers/audit.worker.ts b/src/workers/audit.worker.ts new file mode 100644 index 00000000..27e36b7d --- /dev/null +++ b/src/workers/audit.worker.ts @@ -0,0 +1,169 @@ +import { Processor, WorkerHost } from '@nestjs/bullmq'; +import { Inject, Logger, Optional } from '@nestjs/common'; +import { Job } from 'bullmq'; +import { Queues } from '../queues/queues.constants'; +import { AuditJobData } from '../queues/queues.constants'; +import { PrismaService } from '../database/prisma.service'; +import { WorkerMetricsService } from '../modules/metrics/worker-metrics.service'; +import { AuditHashService } from '../modules/audit/audit-hash.service'; + +/** + * How many consecutive audit-entry batches may fail to persist before the + * worker starts logging at `error` level with a `AUDIT_PERSISTENCE_DEGRADED` + * marker. The worker never crashes on persistence failures — audit entries are + * best-effort for the originating request by design (see `AuditListener`) — + * but operators need a loud signal when the audit trail is silently dropping + * writes. + */ +const DEGRADED_FAILURE_THRESHOLD = 5; + +/** + * BullMQ worker that asynchronously persists audit log entries. + * + * High-frequency events (policy checks, authentication attempts, risk + * evaluations) are enqueued by the audit module instead of being written + * synchronously during the request-response cycle. This worker drains the + * queue and writes the entries to PostgreSQL through Prisma: + * + * - Batching: every job carries a batch of entries (`AuditJobData`) which is + * persisted with a single bulk insert inside one transaction, so queue + * volume spikes translate to fewer, larger database round-trips. + * - Hash chaining: entries are chained to the preceding audit hash per + * organization via {@link AuditHashService} to preserve the tamper-evident + * audit history produced by the synchronous `AuditService.record` path. + * - Retries: transient database outages are retried by BullMQ with + * exponential backoff (see the queue registration in `AuditModule`); the + * failure is rethrown so BullMQ tracks attempt counts. + * - Fail-safe: persistence errors are caught and logged — the worker process + * itself never crashes. On the final attempt the error is rethrown so the + * job lands in the dead-letter queue for forensic triage. + */ +@Processor(Queues.Audit) +export class AuditWorker extends WorkerHost { + private readonly logger = new Logger(AuditWorker.name); + + /** Consecutive failed batches — used for the degraded-persistence signal. */ + private consecutiveFailures = 0; + + constructor( + @Optional() @Inject(PrismaService) private readonly prisma?: PrismaService, + @Optional() private readonly auditHashService?: AuditHashService, + @Optional() private readonly workerMetrics?: WorkerMetricsService, + ) { + super(); + } + + async process(job: Job): Promise<{ persisted: number }> { + const jobName = job.name ?? 'audit-persist'; + const execute = async (): Promise<{ persisted: number }> => { + const entries = job.data?.entries ?? []; + if (entries.length === 0) { + return { persisted: 0 }; + } + + try { + const persisted = await this.persistBatch(entries); + this.consecutiveFailures = 0; + this.logger.debug(`Persisted ${persisted} audit entr(ies) from job ${String(job.id)}`); + return { persisted }; + } catch (error) { + this.consecutiveFailures += 1; + const message = (error as Error).message; + + if (this.consecutiveFailures >= DEGRADED_FAILURE_THRESHOLD) { + this.logger.error( + `[AUDIT_PERSISTENCE_DEGRADED] ${this.consecutiveFailures} consecutive failed batches — ` + + `audit entries are being dropped: ${message}`, + ); + } else { + this.logger.warn( + `Failed to persist ${entries.length} audit entr(ies) ` + + `(attempt ${job.attemptsMade + 1}/${job.opts.attempts ?? '?'}): ${message}`, + ); + } + + // Rethrow so BullMQ records the failure, applies backoff and — on the + // final attempt — routes the job to the dead-letter queue. + throw error instanceof Error ? error : new Error(message); + } + }; + + if (this.workerMetrics) { + return this.workerMetrics.instrumentJob(Queues.Audit, jobName, execute); + } + return execute(); + } + + /** + * Persists a batch of audit entries inside a single transaction. + * + * Each entry is hash-chained to the previous entry of its organization so + * the asynchronous trail stays verifiable by `AuditHashService`. Returns the + * number of rows written. + */ + private async persistBatch(entries: AuditJobData['entries']): Promise { + if (!this.prisma) { + this.logger.warn('PrismaService unavailable — dropping audit batch'); + return 0; + } + + // Persist through the dedicated worker client so bulk audit writes are + // never aborted by the API-oriented query timeouts (issue #76). + const client = this.prisma.workerClient ?? this.prisma; + + return client.$transaction(async (tx) => { + let persisted = 0; + // previousHash per organization, chained within the batch too. + const previousHashes = new Map(); + + for (const entry of entries) { + let previousHash = previousHashes.get(entry.organizationId); + if (previousHash === undefined && this.auditHashService) { + previousHash = await this.auditHashService.getLatestHash(entry.organizationId); + } + + let hash: string | null = null; + if (this.auditHashService) { + const createdAt = new Date(); + const result = this.auditHashService.computeEntryHash( + { + organizationId: entry.organizationId, + userId: entry.userId ?? null, + action: entry.action, + entity: entry.entity, + entityId: entry.entityId ?? null, + oldValue: entry.oldValue, + newValue: entry.newValue, + ipAddress: entry.ipAddress ?? null, + device: entry.device ?? null, + createdAt, + }, + previousHash ?? null, + ); + hash = result.hash; + previousHashes.set(entry.organizationId, result.hash); + } + + await tx.auditLog.create({ + data: { + organizationId: entry.organizationId, + userId: entry.userId ?? null, + action: entry.action, + entity: entry.entity, + entityId: entry.entityId ?? null, + oldValue: (entry.oldValue ?? undefined) as never, + newValue: (entry.newValue ?? undefined) as never, + ipAddress: entry.ipAddress ?? null, + device: entry.device ?? null, + requestId: entry.requestId ?? null, + previousHash: previousHash ?? null, + hash, + }, + }); + persisted += 1; + } + + return persisted; + }); + } +} diff --git a/src/workers/index.ts b/src/workers/index.ts index 0ea4e456..e52abf3c 100644 --- a/src/workers/index.ts +++ b/src/workers/index.ts @@ -1,2 +1,3 @@ export * from './workers.module'; export * from './dlq.processor'; +export * from './audit.worker'; diff --git a/src/workers/workers.module.ts b/src/workers/workers.module.ts index 720853ca..8735a9b2 100644 --- a/src/workers/workers.module.ts +++ b/src/workers/workers.module.ts @@ -1,9 +1,13 @@ import { Module } from '@nestjs/common'; +import { BullModule } from '@nestjs/bullmq'; import { BalanceWorker } from './balance.worker'; import { AnalyticsAggregationWorker } from './analytics-aggregation.worker'; import { NotificationDeliveryWorker } from './notification-delivery.worker'; +import { AuditWorker } from './audit.worker'; import { WalletModule } from '../modules/wallets/wallet.module'; import { MetricsModule } from '../modules/metrics/metrics.module'; +import { AuditModule } from '../modules/audit/audit.module'; +import { Queues } from '../queues/queues.constants'; /** * Background job processors. @@ -19,16 +23,34 @@ import { MetricsModule } from '../modules/metrics/metrics.module'; * `worker_job_duration_seconds` or `worker_jobs_total` metrics. */ @Module({ - imports: [WalletModule, MetricsModule], + imports: [ + WalletModule, + MetricsModule, + AuditModule, + // The audit worker consumes the dedicated `audit` queue; register it here + // so BullMQ provisions the queue with bounded retries and exponential + // backoff for transient database outages. + BullModule.registerQueue({ + name: Queues.Audit, + defaultJobOptions: { + attempts: 5, + backoff: { type: 'exponential', delay: 1_000 }, + removeOnComplete: { count: 1_000 }, + removeOnFail: { age: 7 * 24 * 3_600 }, + }, + }), + ], providers: [ NotificationDeliveryWorker, BalanceWorker, AnalyticsAggregationWorker, + AuditWorker, ], exports: [ NotificationDeliveryWorker, BalanceWorker, AnalyticsAggregationWorker, + AuditWorker, ], }) export class WorkersModule {}