From 88dc2fc1da7171c174f773ced8073fd099a98aae Mon Sep 17 00:00:00 2001 From: Bogunrot <317332203+Bogunrot@users.noreply.github.com> Date: Mon, 28 Sep 2026 10:34:14 +0000 Subject: [PATCH] feat: queue reliability hardening and migration verification checks Closes #214, closes #217, closes #219, closes #222. Migration verification (#214): - scripts/verify-migrations.sh now runs static checks that need no database: missing/empty migration.sql files, unbalanced quotes/parentheses, no executable SQL, duplicate migration directories, and the Prisma _ naming convention. Database apply/status and shadow-database drift checks run only when DATABASE_URL/ SHADOW_DATABASE_URL are set. - CI gains a migrations (static) job running npm run db:verify:static on every pull request alongside the existing database-backed job. - CONTRIBUTING.md documents the verification workflow and naming convention. Audit log persistence worker (#217): - New AuditWorker (src/workers/audit.worker.ts) drains the dedicated `audit` BullMQ queue and persists batched entries in one transaction per job using the dedicated worker Prisma pool, preserving hash chaining via AuditHashService. Transient database failures are retried with exponential backoff (5 attempts), logged at warn/error without crashing the worker, and terminal failures land in the dead-letter queue. AuditService.queueRecord groups entries per organization and falls back to synchronous persistence if the queue is unavailable. Webhook circuit breaker (#219): - New WebhookCircuitBreakerService tracks consecutive delivery failures per endpoint host in Redis (in-memory fallback) and trips OPEN after 5 consecutive failures, half-opens after a cooldown window, and closes again after consecutive successful trials. The webhook processor fail-fasts deliveries to OPEN domains and records every outcome, while the existing 5-attempt exponential backoff with jitter stays intact. Queue health endpoint (#222): - New BullMQHealthIndicator probes every registered queue with a hard 2s timeout, reports waiting/active/completed/failed/delayed/paused counts and Redis connectivity, and exposes GET /health/queues with Swagger docs and a 503 when Redis is unreachable or all queues fail their probes. --- .github/workflows/ci.yml | 21 +- CONTRIBUTING.md | 34 +- package.json | 3 +- scripts/verify-migrations.sh | 178 +++++++++- src/modules/audit/audit.module.ts | 13 + src/modules/audit/audit.service.ts | 65 +++- src/modules/health/health.controller.ts | 40 +++ src/modules/health/health.module.ts | 5 +- src/modules/health/index.ts | 1 + .../health/indicators/bullmq.health.spec.ts | 159 +++++++++ .../health/indicators/bullmq.health.ts | 199 +++++++++++ src/modules/webhooks/index.ts | 1 + .../webhook-circuit-breaker.service.spec.ts | 155 ++++++++ .../webhook-circuit-breaker.service.ts | 334 ++++++++++++++++++ src/modules/webhooks/webhook.module.ts | 4 +- .../webhooks/webhooks.processor.spec.ts | 90 +++++ src/modules/webhooks/webhooks.processor.ts | 43 +++ src/queues/queues.constants.ts | 22 ++ src/workers/audit.worker.spec.ts | 180 ++++++++++ src/workers/audit.worker.ts | 169 +++++++++ src/workers/index.ts | 1 + src/workers/workers.module.ts | 24 +- 22 files changed, 1717 insertions(+), 24 deletions(-) create mode 100644 src/modules/health/indicators/bullmq.health.spec.ts create mode 100644 src/modules/health/indicators/bullmq.health.ts create mode 100644 src/modules/webhooks/services/webhook-circuit-breaker.service.spec.ts create mode 100644 src/modules/webhooks/services/webhook-circuit-breaker.service.ts create mode 100644 src/workers/audit.worker.spec.ts create mode 100644 src/workers/audit.worker.ts diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index fe34481a..575bf5e1 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -58,8 +58,27 @@ jobs: - run: npm ci --include=dev - run: npm run typecheck + # 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 + needs: build + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-node@v4 + with: + node-version: 20 + cache: npm + - run: npm ci --include=dev + - name: Verify migration files (no database required) + run: npm run db:verify:static + migrations: - name: migrations + name: migrations (database) runs-on: ubuntu-latest needs: build services: diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 277c4be0..b5f2e781 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -34,8 +34,38 @@ its own microservice. 1. `npm run typecheck && npm run lint && npm run test` all pass. 2. Database schema changes include a Prisma migration. -3. New endpoints are documented with OpenAPI/Swagger decorators. -4. Cross-repo contracts (response envelope, entity/enum names) still match `astroid-web` and `astroid-sdk`. +3. `npm run db:verify:static` passes — see [Database migration verification](#database-migration-verification). +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`. + +## 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 two jobs +on every pull request: + +| 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. ## Branch strategy diff --git a/package.json b/package.json index c1a52524..25b26cdd 100644 --- a/package.json +++ b/package.json @@ -25,7 +25,8 @@ "prisma:deploy": "prisma migrate deploy", "prisma:seed": "ts-node prisma/seed.ts", "db:seed": "ts-node prisma/seed.ts", - "db:verify": "scripts/verify-migrations.sh" + "db:verify": "scripts/verify-migrations.sh", + "db:verify:static": "DATABASE_URL= SHADOW_DATABASE_URL= scripts/verify-migrations.sh" }, "prisma": { "seed": "ts-node prisma/seed.ts" diff --git a/scripts/verify-migrations.sh b/scripts/verify-migrations.sh index 3ab4d482..58b2086b 100755 --- a/scripts/verify-migrations.sh +++ b/scripts/verify-migrations.sh @@ -3,34 +3,182 @@ # Verifies the Prisma database migrations for the Astroid API. # # What it does, in order: -# 1. Generates the Prisma client. -# 2. Applies every pending migration to the target database (idempotent). -# 3. Drift check: rebuilds the schema purely from the committed migrations -# (in an ephemeral shadow database) and fails if it does not match -# prisma/schema.prisma. This catches schema edits that were never -# captured in a migration. +# 1. Static verification (no database required): +# - Every migration directory contains a migration.sql file. +# - Every migration.sql is non-empty and contains executable SQL +# statements (not just comments/whitespace). +# - Migration directory names follow the Prisma convention +# `_` (plus the reserved `0_init`-style +# migration and `migration_lock.toml`). +# - Unbalanced parentheses / unclosed quote heuristics catch truncated +# or hand-broken SQL before it reaches a database. +# 2. Generates the Prisma client. +# 3. (DATABASE_URL set) Applies every pending migration (idempotent) and +# checks `prisma migrate status`. +# 4. (SHADOW_DATABASE_URL set) Drift check: rebuilds the schema purely from +# the committed migrations in an ephemeral shadow database and fails if +# it does not match prisma/schema.prisma. This catches schema edits that +# were never captured in a migration. # # Env: -# DATABASE_URL (required) Target PostgreSQL the migrations are -# applied to — e.g. a fresh ephemeral CI database. -# SHADOW_DATABASE_URL (optional but recommended) An empty scratch database -# used for the drift check. When unset the drift check -# is skipped and only apply + status are verified. +# DATABASE_URL (optional) Target PostgreSQL the migrations are +# applied to. When unset the script runs in static +# mode only — suitable for containerized CI runners +# without an active database connection. +# SHADOW_DATABASE_URL (optional) An empty scratch database used for the +# drift check. Requires DATABASE_URL to be set. # -# Exit code 0 when migrations apply cleanly and stay in sync with the schema. +# Exit codes: +# 0 verification passed +# 1 a migration file is missing/empty, invalid SQL, a name violates the +# naming convention, migrations fail to apply, or schema drift exists. set -euo pipefail -: "${DATABASE_URL:?DATABASE_URL is required}" +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +MIGRATIONS_DIR="${REPO_ROOT}/prisma/migrations" +SCHEMA_FILE="${REPO_ROOT}/prisma/schema.prisma" -echo "==> Generating Prisma client" +# The Prisma naming convention is a 14-digit UTC timestamp followed by an +# underscore-separated name (e.g. 20260830174000_sync_schema). Prisma also +# permits a trailing suffix when a name is regenerated (-, +, <, >). The +# baseline `0_init` migration (any single-0-prefixed name) is allowed. +readonly MIGRATION_NAME_RE='^(0|[0-9]{14})_[a-z0-9_]+([-+<>][a-zA-Z0-9_]+)?$' + +errors=0 + +fail() { + echo "!! ${1}" >&2 + errors=$((errors + 1)) +} + +# -------------------------------------------------------------------------- +# 1. Static verification — pure filesystem checks, no database required. +# -------------------------------------------------------------------------- +echo "==> Verifying migrations directory structure (${MIGRATIONS_DIR})" + +if [[ ! -f "${SCHEMA_FILE}" ]]; then + fail "Prisma schema not found at ${SCHEMA_FILE}" +fi + +if [[ ! -d "${MIGRATIONS_DIR}" ]]; then + fail "Migrations directory not found at ${MIGRATIONS_DIR}" + # Nothing else can be verified without the directory. + exit 1 +fi + +# --------------------------------------------------------------------------- +# SQL structural sanity check for a single migration file. +# A migration is considered structurally sound when it contains at least one +# executable statement (anything that is not a comment/whitespace) and has +# balanced quotes. Parentheses balance is checked only when the file does not +# use dollar-quoted strings or `$$` bodies (functions/triggers), where naive +# counting is unreliable. +# --------------------------------------------------------------------------- +check_sql_integrity() { + local file="$1" + local name + name="$(basename "$(dirname "${file}")")" + + if [[ ! -s "${file}" ]]; then + fail "${name}: migration.sql is empty" + return + fi + + # Statements = non-empty lines that are not pure comments or whitespace. + if ! grep -Ev '^[[:space:]]*(--.*)?$' "${file}" > /dev/null; then + fail "${name}: migration.sql contains no executable SQL (comments/blank lines only)" + return + fi + + local content + content="$(cat "${file}")" + + # Unclosed single-quote detection: strip line comments, then count the + # single-quote characters left over. An ODD total means a quote literal was + # left open — common in truncated or hand-broken SQL files. + local quotes + quotes="$(sed 's/--.*$//' "${file}" | tr -cd "'" | wc -c)" + if [[ $((quotes % 2)) -ne 0 ]]; then + fail "${name}: migration.sql appears to contain an unclosed single-quoted string" + fi + + # Skip parenthesis balance when the file declares functions, triggers or + # dollar-quoted bodies — a naive count would false-positive there. + if ! grep -Eq '\$\$|CREATE[[:space:]]+(OR[[:space:]]+REPLACE[[:space:]]+)?(FUNCTION|TRIGGER|PROCEDURE)' "${file}"; then + local open close + open="$(sed 's/--.*$//' "${file}" | tr -cd '(' | wc -c)" + close="$(sed 's/--.*$//' "${file}" | tr -cd ')' | wc -c)" + if [[ "${open}" -ne "${close}" ]]; then + fail "${name}: migration.sql has unbalanced parentheses (${open} open vs ${close} close)" + fi + fi +} + +found_migrations=0 +for dir in "${MIGRATIONS_DIR}"/*; do + [[ -d "${dir}" ]] || continue + found_migrations=$((found_migrations + 1)) + name="$(basename "${dir}")" + + if [[ "${name}" == "migration_lock.toml" || "${name}" == *.toml ]]; then + continue + fi + + if [[ ! "${name}" =~ ${MIGRATION_NAME_RE} ]]; then + fail "'${name}' violates the migration naming convention '_'" + fi + + sql_file="${dir}/migration.sql" + if [[ ! -f "${sql_file}" ]]; then + fail "'${name}' is missing its migration.sql file" + continue + fi + + check_sql_integrity "${sql_file}" +done + +if [[ "${found_migrations}" -eq 0 ]]; then + fail "No migration directories found under prisma/migrations" +fi + +echo " Checked ${found_migrations} migration director(ies)" + +# Duplicated migration names (case variants) would silently shadow history. +duplicate="$(find "${MIGRATIONS_DIR}" -maxdepth 1 -type d \( -iname '[0-9]*' -o -iname '0_*' \) -printf '%f\n' 2>/dev/null | tr '[:upper:]' '[:lower:]' | sort | uniq -d || true)" +if [[ -n "${duplicate}" ]]; then + fail "Duplicate migration directory names detected (case-insensitive): ${duplicate//$'\n'/, }" +fi + +if [[ "${errors}" -gt 0 ]]; then + echo "!! Static migration verification failed with ${errors} error(s)" >&2 + exit 1 +fi +echo "==> Static migration verification passed" + +# -------------------------------------------------------------------------- +# 2. Prisma client generation (also validates schema.prisma syntax). +# -------------------------------------------------------------------------- +echo "==> Generating Prisma client (validates schema.prisma syntax)" npx prisma generate +# -------------------------------------------------------------------------- +# 3. Apply migrations when a database is available (skipped in static mode). +# -------------------------------------------------------------------------- +if [[ -z "${DATABASE_URL:-}" ]]; then + echo "==> DATABASE_URL unset — skipping database apply/status/drift checks" + echo "==> Migration verification passed (static mode)" + exit 0 +fi + echo "==> Applying migrations to ${DATABASE_URL}" npx prisma migrate deploy echo "==> Checking migration status" npx prisma migrate status +# -------------------------------------------------------------------------- +# 4. Drift check (requires an empty shadow database). +# -------------------------------------------------------------------------- if [[ -n "${SHADOW_DATABASE_URL:-}" ]]; then echo "==> Drift check: rebuilding schema from migrations only" echo " shadow database: ${SHADOW_DATABASE_URL}" @@ -56,4 +204,4 @@ else echo "!! SHADOW_DATABASE_URL unset — skipping drift check" >&2 fi -echo "==> Migration verification passed" \ No newline at end of file +echo "==> Migration verification passed" 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/audit/audit.service.ts b/src/modules/audit/audit.service.ts index d2abe8bc..724c833d 100644 --- a/src/modules/audit/audit.service.ts +++ b/src/modules/audit/audit.service.ts @@ -1,5 +1,7 @@ -import { Injectable } from '@nestjs/common'; +import { Injectable, Logger, Optional } from '@nestjs/common'; +import { InjectQueue } from '@nestjs/bullmq'; import { Prisma } from '@prisma/client'; +import { Queue } from 'bullmq'; import { AuditRepository, CreateAuditLogData } from './audit.repository'; import { AuditHashService } from './audit-hash.service'; import { @@ -8,6 +10,7 @@ import { toPrismaPagination, } from '../../common/helpers/pagination'; import { Paginated } from '../../common/interfaces/api-response.interface'; +import { Queues, AuditLogJobEntry } from '../../queues/queues.constants'; const SORTABLE = ['createdAt', 'action', 'entity']; @@ -20,14 +23,74 @@ type ExportedAuditLog = Prisma.AuditLogGetPayload<{ * Writes and queries the immutable audit trail. Records Who / When / Where / * Why / Old / New for every important action. Never updates or deletes. * Integrates cryptographic hash chaining for tamper-evident audit history. + * + * High-frequency events may instead be persisted asynchronously via + * {@link queueRecord}, which enqueues entries onto the BullMQ `audit` queue; + * the `AuditWorker` (src/workers/audit.worker.ts) batch-writes them without + * blocking the request-response cycle. */ @Injectable() export class AuditService { + private readonly logger = new Logger(AuditService.name); + constructor( private readonly repository: AuditRepository, private readonly hashService: AuditHashService, + @Optional() + @InjectQueue(Queues.Audit) + private readonly auditQueue?: Queue, ) {} + /** + * Persists audit entries asynchronously through the BullMQ `audit` queue. + * Entries are grouped by organization (one job per organization) so the + * worker can batch-insert and hash-chain them efficiently. Enqueue failures + * are logged and never propagate — the audit trail must not break callers. + * + * @returns the number of jobs enqueued (one per organization represented). + */ + async queueRecord(entries: AuditLogJobEntry[]): Promise { + if (entries.length === 0) return 0; + if (!this.auditQueue) { + this.logger.warn('Audit queue unavailable — falling back to synchronous persistence'); + for (const entry of entries) { + await this.record(entry as CreateAuditLogData); + } + return 0; + } + + const byOrganization = new Map(); + for (const entry of entries) { + const bucket = byOrganization.get(entry.organizationId); + if (bucket) { + bucket.push(entry); + } else { + byOrganization.set(entry.organizationId, [entry]); + } + } + + let enqueued = 0; + for (const [organizationId, orgEntries] of byOrganization) { + try { + await this.auditQueue.add( + 'audit-persist', + { entries: orgEntries }, + { attempts: 5, backoff: { type: 'exponential', delay: 1_000 } }, + ); + enqueued += 1; + } catch (error) { + this.logger.error( + `Failed to enqueue ${orgEntries.length} audit entr(ies) for organization ${organizationId}: ` + + `${(error as Error).message} — persisting synchronously instead`, + ); + for (const entry of orgEntries) { + await this.record(entry as CreateAuditLogData); + } + } + } + return enqueued; + } + async record(data: CreateAuditLogData) { const previousHash = await this.hashService.getLatestHash(data.organizationId); const createdAt = new Date(); diff --git a/src/modules/health/health.controller.ts b/src/modules/health/health.controller.ts index a08efaff..bcbc4e7a 100644 --- a/src/modules/health/health.controller.ts +++ b/src/modules/health/health.controller.ts @@ -10,6 +10,10 @@ import { DatabaseMigrationHealthIndicator, MigrationHealthReport, } from './indicators/database-migration.health'; +import { + BullMQHealthIndicator, + QueuesHealthReport, +} from './indicators/bullmq.health'; import { JwtAuthGuard } from '../../common/guards/jwt-auth.guard'; import { RolesGuard } from '../../common/guards/roles.guard'; import { Roles } from '../../common/decorators/roles.decorator'; @@ -21,6 +25,7 @@ export class HealthController { constructor( private readonly stellarHealthIndicator: StellarHealthIndicator, private readonly databaseMigrationIndicator: DatabaseMigrationHealthIndicator, + private readonly bullmqHealthIndicator: BullMQHealthIndicator, ) {} @Get('stellar') @@ -75,4 +80,39 @@ export class HealthController { return report; } + + @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 c0b143e8..fe914309 100644 --- a/src/modules/health/health.module.ts +++ b/src/modules/health/health.module.ts @@ -1,11 +1,12 @@ import { Module } from '@nestjs/common'; import { StellarHealthIndicator } from './indicators/stellar.health'; import { DatabaseMigrationHealthIndicator } from './indicators/database-migration.health'; +import { BullMQHealthIndicator } from './indicators/bullmq.health'; import { HealthController } from './health.controller'; @Module({ controllers: [HealthController], - providers: [StellarHealthIndicator, DatabaseMigrationHealthIndicator], - exports: [StellarHealthIndicator, DatabaseMigrationHealthIndicator], + providers: [StellarHealthIndicator, DatabaseMigrationHealthIndicator, BullMQHealthIndicator], + exports: [StellarHealthIndicator, DatabaseMigrationHealthIndicator, BullMQHealthIndicator], }) export class HealthModule {} diff --git a/src/modules/health/index.ts b/src/modules/health/index.ts index b4f17c67..4b75aec6 100644 --- a/src/modules/health/index.ts +++ b/src/modules/health/index.ts @@ -2,3 +2,4 @@ export * from './health.module'; export * from './health.controller'; 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 f8fee92f..db2ab6c7 100644 --- a/src/modules/webhooks/webhook.module.ts +++ b/src/modules/webhooks/webhook.module.ts @@ -5,6 +5,7 @@ import { WebhookService } from './webhook.service'; import { WebhookRepository } from './webhook.repository'; import { WebhookDispatcher } from './webhook.dispatcher'; import { WebhookDeliveryService } from './services/webhook-delivery.service'; +import { WebhookCircuitBreakerService } from './services/webhook-circuit-breaker.service'; import { WebhookWorker } from './workers/webhook.worker'; import { WebhooksProcessor } from './webhooks.processor'; import { Queues } from '../../queues/queues.constants'; @@ -57,9 +58,10 @@ import type { RegisterQueueOptions } from '@nestjs/bullmq'; WebhookRepository, WebhookDispatcher, WebhookDeliveryService, + WebhookCircuitBreakerService, WebhookWorker, WebhooksProcessor, ], - exports: [WebhookService], + exports: [WebhookService, WebhookCircuitBreakerService], }) export class WebhookModule {} diff --git a/src/modules/webhooks/webhooks.processor.spec.ts b/src/modules/webhooks/webhooks.processor.spec.ts index d757e657..a56cae65 100644 --- a/src/modules/webhooks/webhooks.processor.spec.ts +++ b/src/modules/webhooks/webhooks.processor.spec.ts @@ -3,6 +3,9 @@ import { Job, UnrecoverableError } from 'bullmq'; import { WebhooksProcessor } from './webhooks.processor'; import { WebhookJobData } from './types/webhook-job.types'; import { createHmac } from 'crypto'; +import { + WebhookCircuitBreakerService, +} from './services/webhook-circuit-breaker.service'; describe('WebhooksProcessor', () => { let processor: WebhooksProcessor; @@ -235,4 +238,91 @@ describe('WebhooksProcessor', () => { await expect(processor.process(job)).rejects.toThrow('HTTP 503'); }); }); + + 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('records failures and trips the circuit after consecutive HTTP 500s', async () => { + 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'); + + 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 89a285c8..55461d0c 100644 --- a/src/modules/webhooks/webhooks.processor.ts +++ b/src/modules/webhooks/webhooks.processor.ts @@ -7,6 +7,7 @@ import { WebhookJobData, WebhookJobResult } from './types/webhook-job.types'; import { signWebhookPayload } from './utils/signing'; import { PrismaService } from '../../database/prisma.service'; import { WorkerMetricsService } from '../../modules/metrics/worker-metrics.service'; +import { WebhookCircuitBreakerService } from './services/webhook-circuit-breaker.service'; /** * BullMQ job processor for webhook event delivery with exponential backoff + jitter. @@ -18,6 +19,11 @@ import { WorkerMetricsService } from '../../modules/metrics/worker-metrics.servi * - 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. @@ -37,6 +43,7 @@ export class WebhooksProcessor extends WorkerHost { @Optional() @Inject(PrismaService) private readonly prisma?: PrismaService, @Optional() private readonly configService?: ConfigService, @Optional() private readonly workerMetrics?: WorkerMetricsService, + @Optional() private readonly circuitBreaker?: WebhookCircuitBreakerService, ) { super(); } @@ -56,6 +63,24 @@ export class WebhooksProcessor extends WorkerHost { const execute = async (): Promise => { 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; @@ -125,9 +150,27 @@ export class WebhooksProcessor extends WorkerHost { if (isLastAttempt) { 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 c2bc761c..d82eea9a 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]; @@ -51,3 +53,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 20763842..c97b4f78 100644 --- a/src/workers/workers.module.ts +++ b/src/workers/workers.module.ts @@ -1,10 +1,14 @@ import { Module } from '@nestjs/common'; +import { BullModule } from '@nestjs/bullmq'; import { BalanceWorker } from './balance.worker'; import { WebhookDeliveryWorker } from './webhook-delivery.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. @@ -20,18 +24,36 @@ 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, WebhookDeliveryWorker, BalanceWorker, AnalyticsAggregationWorker, + AuditWorker, ], exports: [ NotificationDeliveryWorker, WebhookDeliveryWorker, BalanceWorker, AnalyticsAggregationWorker, + AuditWorker, ], }) export class WorkersModule {}