From c88255c09144b1d02ab1686ea974decc15cd833f Mon Sep 17 00:00:00 2001 From: Fatima Aminu <81179612+phertyameen@users.noreply.github.com> Date: Thu, 24 Sep 2026 21:36:14 +0100 Subject: [PATCH] feat(backend): add provider circuit breaker, structured logging, and API hardening Wraps the outbound payment-rail initiation in a per-rail opossum circuit breaker configured from the environment, so a degraded provider is short-circuited instead of being hammered on every request. Business rejections do not count toward opening the circuit, a volume threshold keeps a low-traffic rail from tripping on one error, and an open breaker surfaces as a clear 503. The breaker is created lazily so the existing service constructor and tests are unchanged. Adds pino-backed structured JSON logging for production (with a LOG_JSON override), promoting the Nest logger context and request id to real fields, plus one structured access-log line per completed response. Completes the admin OpenAPI surface: reconciliation metrics, ledger integrity, payment sweep and settlement run now have typed response DTOs, and previously undocumented query/path parameters are declared, with a reflection-based spec that fails if any admin handler loses its operation or response metadata. Configures an explicit, configurable request body size limit on the JSON and urlencoded parsers while preserving the raw-body capture the payment webhook HMAC verification depends on. Closes #1797 Closes #1796 Closes #1795 Closes #1794 --- backend/.env.example | 18 +++ backend/package.json | 2 + .../src/common/structured-logger.service.ts | 113 ++++++++++++++++++ .../structured-request-logger.middleware.ts | 52 ++++++++ .../src/credits/credits-admin.controller.ts | 51 +++++++- .../dto/ledger-integrity-response.dto.ts | 65 ++++++++++ .../credits/dto/payment-sweep-response.dto.ts | 18 +++ .../dto/settlement-run-response.dto.ts | 19 +++ backend/src/main.ts | 56 ++++++++- backend/src/payments/admin-swagger.spec.ts | 45 +++++++ .../reconciliation-metrics-response.dto.ts | 22 ++++ .../src/payments/payments-admin.controller.ts | 13 +- backend/src/payments/payments.service.spec.ts | 33 ++++- backend/src/payments/payments.service.ts | 36 +++++- backend/src/payments/utils/circuit-breaker.ts | 85 +++++++++++++ 15 files changed, 612 insertions(+), 16 deletions(-) create mode 100644 backend/src/common/structured-logger.service.ts create mode 100644 backend/src/common/structured-request-logger.middleware.ts create mode 100644 backend/src/credits/dto/ledger-integrity-response.dto.ts create mode 100644 backend/src/credits/dto/payment-sweep-response.dto.ts create mode 100644 backend/src/credits/dto/settlement-run-response.dto.ts create mode 100644 backend/src/payments/admin-swagger.spec.ts create mode 100644 backend/src/payments/dto/reconciliation-metrics-response.dto.ts create mode 100644 backend/src/payments/utils/circuit-breaker.ts diff --git a/backend/.env.example b/backend/.env.example index f655b283..2d1249ac 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -1,5 +1,15 @@ PORT=6000 NODE_ENV=development + +# Maximum request body size accepted by the JSON and urlencoded parsers. +# The webhook parser still captures the exact raw bytes for HMAC checks. +REQUEST_BODY_LIMIT=1mb + +# Structured JSON logging +# Production automatically uses pino; set LOG_JSON=true to enable it in +# another environment. LOG_JSON=false does not disable the production default. +LOG_JSON= + DATABASE_NAME=your_database_name DATABASE_PASSWORD=your_database_password DATABASE_USERNAME=your_database_username @@ -84,6 +94,14 @@ PAYMENT_WEBHOOK_SECRET=change-me-in-every-environment # payment CONFIRMED on its own. PAYMENT_VERIFY_TIMEOUT_MS=3000 +# Payment provider circuit breaker (issue #1797) +# Bounds one outbound provider initiation, opens after enough provider +# failures, and waits this long before trying a half-open request again. +# The breaker also requires ten requests in its rolling window before it can open. +PAYMENT_PROVIDER_BREAKER_TIMEOUT_MS=10000 +PAYMENT_PROVIDER_BREAKER_ERROR_THRESHOLD_PERCENT=50 +PAYMENT_PROVIDER_BREAKER_RESET_TIMEOUT_MS=30000 + # Payments reconciliation engine (issue #1572) # A payment isn't eligible for its first reconciliation pass until it's # been AWAITING_CONFIRMATION for at least this long — gives the webhook a diff --git a/backend/package.json b/backend/package.json index cd1b50cf..0ea1da7b 100644 --- a/backend/package.json +++ b/backend/package.json @@ -71,12 +71,14 @@ "nestjs-command": "^3.1.5", "nodemailer": "^7.0.12", "nodemailer-mjml": "^1.6.0", + "opossum": "^8.2.0", "otplib": "^13.3.0", "passport": "^0.7.0", "passport-jwt": "^4.0.1", "passport-local": "^1.0.0", "pdfkit": "^0.17.2", "pg": "^8.16.3", + "pino": "^9.5.0", "qrcode": "^1.5.4", "reflect-metadata": "^0.2.0", "rxjs": "^7.8.1", diff --git a/backend/src/common/structured-logger.service.ts b/backend/src/common/structured-logger.service.ts new file mode 100644 index 00000000..5aff9151 --- /dev/null +++ b/backend/src/common/structured-logger.service.ts @@ -0,0 +1,113 @@ +import type { LoggerService, LogLevel } from '@nestjs/common'; +import pino from 'pino'; +import { currentRequestId } from './request-context'; + +type PinoLevel = 'fatal' | 'error' | 'warn' | 'info' | 'debug'; + +/** + * Nest-compatible logger backed by pino. Nest supplies the logger context as + * the first optional parameter, so it is promoted to a structured field + * rather than being flattened into the message text. + */ +export class StructuredLoggerService implements LoggerService { + private readonly logger = pino(); + + log(message: any, ...optionalParams: any[]): void { + this.write('info', message, optionalParams); + } + + error(message: any, ...optionalParams: any[]): void { + this.write('error', message, optionalParams); + } + + warn(message: any, ...optionalParams: any[]): void { + this.write('warn', message, optionalParams); + } + + debug(message: any, ...optionalParams: any[]): void { + this.write('debug', message, optionalParams); + } + + verbose(message: any, ...optionalParams: any[]): void { + this.write('debug', message, optionalParams); + } + + fatal(message: any, ...optionalParams: any[]): void { + this.write('fatal', message, optionalParams); + } + + setLogLevels(levels: LogLevel[]): void { + if (levels.length === 0) { + this.logger.level = 'silent'; + } else if (levels.includes('verbose') || levels.includes('debug')) { + this.logger.level = 'debug'; + } else if (levels.includes('log')) { + this.logger.level = 'info'; + } else if (levels.includes('warn')) { + this.logger.level = 'warn'; + } else if (levels.includes('error')) { + this.logger.level = 'error'; + } else { + this.logger.level = 'fatal'; + } + } + + private write( + level: PinoLevel, + message: any, + optionalParams: any[], + ): void { + const fields: Record = {}; + const context = + typeof optionalParams[0] === 'string' ? optionalParams[0] : undefined; + const requestId = currentRequestId(); + + if (context) { + fields.context = context; + } + if (requestId) { + fields.requestId = requestId; + } + if (message instanceof Error) { + fields.err = message; + } + + const text = this.stringifyMessage(message); + switch (level) { + case 'fatal': + this.logger.fatal(fields, text); + break; + case 'error': + this.logger.error(fields, text); + break; + case 'warn': + this.logger.warn(fields, text); + break; + case 'info': + this.logger.info(fields, text); + break; + case 'debug': + this.logger.debug(fields, text); + break; + } + } + + private stringifyMessage(message: any): string { + if (typeof message === 'string') { + return message; + } + if (message instanceof Error) { + return message.stack ?? message.message; + } + try { + const serialized = JSON.stringify(message); + return serialized === undefined ? String(message) : serialized; + } catch { + try { + return String(message); + } catch { + return '[unserializable log message]'; + } + } + } +} diff --git a/backend/src/common/structured-request-logger.middleware.ts b/backend/src/common/structured-request-logger.middleware.ts new file mode 100644 index 00000000..d749d3cb --- /dev/null +++ b/backend/src/common/structured-request-logger.middleware.ts @@ -0,0 +1,52 @@ +import { Injectable, NestMiddleware } from '@nestjs/common'; +import pino from 'pino'; +import { NextFunction, Request, Response } from 'express'; +import { currentRequestId } from './request-context'; + +/** + * Emits one structured record after a response finishes. It only attaches a + * listener and never writes headers or buffers, so the existing security and + * request-context middleware remain in control of the response. + */ +@Injectable() +export class StructuredRequestLoggerMiddleware implements NestMiddleware { + private readonly logger = pino(); + + use(req: Request, res: Response, next: NextFunction): void { + const startedAt = Date.now(); + let logged = false; + + const logCompletedResponse = (): void => { + if (logged) { + return; + } + logged = true; + + const fields: Record = { + method: req.method, + url: req.originalUrl, + path: req.path, + statusCode: res.statusCode, + durationMs: Date.now() - startedAt, + }; + const responseRequestId = + typeof res.getHeader === 'function' + ? res.getHeader('x-request-id') + : undefined; + const requestId = + currentRequestId() ?? + (typeof req.header === 'function' + ? req.header('x-request-id') + : undefined) ?? + (typeof responseRequestId === 'string' ? responseRequestId : undefined); + if (requestId) { + fields.requestId = requestId; + } + + this.logger.info(fields, 'HTTP request'); + }; + + res.on('finish', logCompletedResponse); + next(); + } +} diff --git a/backend/src/credits/credits-admin.controller.ts b/backend/src/credits/credits-admin.controller.ts index 23061e45..0154ace0 100644 --- a/backend/src/credits/credits-admin.controller.ts +++ b/backend/src/credits/credits-admin.controller.ts @@ -14,6 +14,9 @@ import { import { ApiBearerAuth, ApiOperation, + ApiParam, + ApiProduces, + ApiQuery, ApiResponse, ApiTags, } from '@nestjs/swagger'; @@ -61,6 +64,9 @@ import { SettlementBatchBreakdownResponseDto, SettlementBatchResponseDto, } from './dto/settlement-response.dto'; +import { LedgerIntegrityResponseDto } from './dto/ledger-integrity-response.dto'; +import { PaymentSweepResponseDto } from './dto/payment-sweep-response.dto'; +import { SettlementRunResponseDto } from './dto/settlement-run-response.dto'; /** * Admin surface for the credit ledger (issue #1575): account policy, @@ -89,6 +95,7 @@ export class CreditsAdminController { @Get('accounts') @ApiOperation({ summary: 'List ledger accounts' }) + @ApiQuery({ name: 'currency', required: false, type: String }) @ApiResponse({ status: 200, type: [LedgerAccountResponseDto] }) async listAccounts( @Query('currency') currency?: string, @@ -101,6 +108,9 @@ export class CreditsAdminController { @Get('accounts/export') @ApiOperation({ summary: 'Export ledger accounts as xlsx' }) + @ApiQuery({ name: 'currency', required: false, type: String }) + @ApiProduces('application/vnd.openxmlformats-officedocument.spreadsheetml.sheet') + @ApiResponse({ status: 200, description: 'xlsx workbook stream' }) async exportAccounts( @Query('currency') currency: string | undefined, @Res({ passthrough: true }) res: Response, @@ -165,6 +175,7 @@ export class CreditsAdminController { } @Patch('accounts/:id') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Update an account’s policy (overdraft, payout address, freeze)', description: @@ -190,6 +201,8 @@ export class CreditsAdminController { @Get('balances/:userId') @ApiOperation({ summary: 'A member’s credit balance' }) + @ApiParam({ name: 'userId', type: String, format: 'uuid' }) + @ApiQuery({ name: 'currency', required: false, type: String }) @ApiResponse({ status: 200, type: CreditBalanceResponseDto }) async getBalance( @Param('userId', ParseUUIDPipe) userId: string, @@ -227,8 +240,14 @@ export class CreditsAdminController { 'reports drift, plus any transaction whose debits and credits do not ' + 'cancel. Both lists empty is the healthy state.', }) - checkIntegrity(@Query('currency') currency?: string) { - return this.ledger.checkIntegrity(currency); + @ApiQuery({ name: 'currency', required: false, type: String }) + @ApiResponse({ status: 200, type: LedgerIntegrityResponseDto }) + async checkIntegrity( + @Query('currency') currency?: string, + ): Promise { + return LedgerIntegrityResponseDto.fromView( + await this.ledger.checkIntegrity(currency), + ); } // ── revenue splits ───────────────────────────────────────────────────── @@ -264,6 +283,7 @@ export class CreditsAdminController { } @Get('splits/:id') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Get one revenue split config' }) @ApiResponse({ status: 200, type: RevenueSplitConfigResponseDto }) async getSplit( @@ -274,6 +294,7 @@ export class CreditsAdminController { } @Put('splits/:id/recipients') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Replace a config’s recipients', description: @@ -290,6 +311,7 @@ export class CreditsAdminController { } @Post('splits/:id/active') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Activate or deactivate a config' }) @ApiResponse({ status: 200, type: RevenueSplitConfigResponseDto }) async setActive( @@ -329,6 +351,7 @@ export class CreditsAdminController { // ── payment integration ──────────────────────────────────────────────── @Post('payments/:paymentId/split-config') + @ApiParam({ name: 'paymentId', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Attach a revenue split config to a payment', description: @@ -350,6 +373,7 @@ export class CreditsAdminController { } @Post('payments/:paymentId/top-up') + @ApiParam({ name: 'paymentId', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Mark a payment as funding the payer’s credit balance', description: @@ -368,8 +392,11 @@ export class CreditsAdminController { @ApiOperation({ summary: 'Run the confirmed-payment credit sweep immediately', }) - sweepPayments() { - return this.paymentCredits.sweepConfirmedPayments(); + @ApiResponse({ status: 200, type: PaymentSweepResponseDto }) + async sweepPayments(): Promise { + return PaymentSweepResponseDto.fromView( + await this.paymentCredits.sweepConfirmedPayments(), + ); } // ── settlement ───────────────────────────────────────────────────────── @@ -382,8 +409,11 @@ export class CreditsAdminController { 'Safe to call at any time — the same guarantees the scheduled job ' + 'relies on.', }) - runSettlement() { - return this.settlement.runSettlement(); + @ApiResponse({ status: 200, type: SettlementRunResponseDto }) + async runSettlement(): Promise { + return SettlementRunResponseDto.fromView( + await this.settlement.runSettlement(), + ); } @Post('settlement/batches') @@ -413,6 +443,11 @@ export class CreditsAdminController { @Get('settlement/batches') @ApiOperation({ summary: 'List settlement batches, newest first' }) + @ApiQuery({ + name: 'status', + required: false, + enum: SettlementBatchStatus, + }) @ApiResponse({ status: 200, type: [SettlementBatchResponseDto] }) async listBatches( @Query('status') status?: SettlementBatchStatus, @@ -422,6 +457,7 @@ export class CreditsAdminController { } @Get('settlement/batches/:id') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Full breakdown of one batch', description: @@ -437,6 +473,7 @@ export class CreditsAdminController { } @Post('settlement/batches/:id/execute') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Advance one batch by a step (submit pending, poll submitted)', }) @@ -457,6 +494,7 @@ export class CreditsAdminController { } @Post('settlement/batches/:id/retry') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Re-queue a batch’s failed payouts', description: @@ -480,6 +518,7 @@ export class CreditsAdminController { } @Post('settlement/batches/:id/abandon') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Give up on a batch and release its unsettled claims', description: diff --git a/backend/src/credits/dto/ledger-integrity-response.dto.ts b/backend/src/credits/dto/ledger-integrity-response.dto.ts new file mode 100644 index 00000000..18ad4e0a --- /dev/null +++ b/backend/src/credits/dto/ledger-integrity-response.dto.ts @@ -0,0 +1,65 @@ +import { ApiProperty } from '@nestjs/swagger'; +import type { LedgerIntegrityReport } from '../ledger.service'; + +/** One account whose cached balance disagrees with its ledger entries. */ +export class LedgerBalanceDriftResponseDto { + @ApiProperty() accountId: string; + @ApiProperty({ description: 'Balance stored on the account row' }) + materialized: number; + @ApiProperty({ description: 'Balance derived from the append-only entries' }) + derived: number; + + static fromView( + row: LedgerIntegrityReport['balanceDrift'][number], + ): LedgerBalanceDriftResponseDto { + const dto = new LedgerBalanceDriftResponseDto(); + dto.accountId = row.accountId; + dto.materialized = row.materialized; + dto.derived = row.derived; + return dto; + } +} + +/** One transaction whose debit and credit legs do not cancel. */ +export class LedgerUnbalancedTransactionResponseDto { + @ApiProperty() transactionId: string; + @ApiProperty() debits: number; + @ApiProperty() credits: number; + + static fromView( + row: LedgerIntegrityReport['unbalancedTransactions'][number], + ): LedgerUnbalancedTransactionResponseDto { + const dto = new LedgerUnbalancedTransactionResponseDto(); + dto.transactionId = row.transactionId; + dto.debits = row.debits; + dto.credits = row.credits; + return dto; + } +} + +/** + * A full re-derivation of the ledger cache. Empty drift and imbalance lists + * are the healthy state; both are retained in the response so an operator + * can diagnose the exact account or transaction involved. + */ +export class LedgerIntegrityResponseDto { + @ApiProperty() accountsChecked: number; + @ApiProperty({ type: [LedgerBalanceDriftResponseDto] }) + balanceDrift: LedgerBalanceDriftResponseDto[]; + @ApiProperty({ type: [LedgerUnbalancedTransactionResponseDto] }) + unbalancedTransactions: LedgerUnbalancedTransactionResponseDto[]; + + static fromView( + report: LedgerIntegrityReport, + ): LedgerIntegrityResponseDto { + const dto = new LedgerIntegrityResponseDto(); + dto.accountsChecked = report.accountsChecked; + dto.balanceDrift = report.balanceDrift.map((row) => + LedgerBalanceDriftResponseDto.fromView(row), + ); + dto.unbalancedTransactions = report.unbalancedTransactions.map((row) => + LedgerUnbalancedTransactionResponseDto.fromView(row), + ); + return dto; + } +} diff --git a/backend/src/credits/dto/payment-sweep-response.dto.ts b/backend/src/credits/dto/payment-sweep-response.dto.ts new file mode 100644 index 00000000..2dfc24c4 --- /dev/null +++ b/backend/src/credits/dto/payment-sweep-response.dto.ts @@ -0,0 +1,18 @@ +import { ApiProperty } from '@nestjs/swagger'; +import type { PaymentSweepSummary } from '../payment-credits.service'; + +/** Counts produced by one confirmed-payment credit sweep pass. */ +export class PaymentSweepResponseDto { + @ApiProperty({ description: 'Confirmed payments considered by the pass' }) + candidates: number; + @ApiProperty({ description: 'Payments whose credit effect was applied' }) + applied: number; + @ApiProperty({ description: 'Payments skipped because they were already handled' }) + skipped: number; + @ApiProperty({ description: 'Payments whose application failed during the pass' }) + failed: number; + + static fromView(summary: PaymentSweepSummary): PaymentSweepResponseDto { + return Object.assign(new PaymentSweepResponseDto(), summary); + } +} diff --git a/backend/src/credits/dto/settlement-run-response.dto.ts b/backend/src/credits/dto/settlement-run-response.dto.ts new file mode 100644 index 00000000..1573782b --- /dev/null +++ b/backend/src/credits/dto/settlement-run-response.dto.ts @@ -0,0 +1,19 @@ +import { ApiProperty } from '@nestjs/swagger'; +import type { SettlementRunSummary } from '../settlement.service'; + +/** Aggregate progress and outcomes from one settlement pass. */ +export class SettlementRunResponseDto { + @ApiProperty() batchesCreated: number; + @ApiProperty() batchesExecuted: number; + @ApiProperty() payoutsSubmitted: number; + @ApiProperty() payoutsConfirmed: number; + @ApiProperty() payoutsFailed: number; + @ApiProperty() payoutsAwaitingRail: number; + @ApiProperty() entriesSettled: number; + @ApiProperty({ type: String, isArray: true }) + notes: string[]; + + static fromView(summary: SettlementRunSummary): SettlementRunResponseDto { + return Object.assign(new SettlementRunResponseDto(), summary); + } +} diff --git a/backend/src/main.ts b/backend/src/main.ts index 2e8d4826..60ab7fcf 100644 --- a/backend/src/main.ts +++ b/backend/src/main.ts @@ -1,16 +1,62 @@ import { NestFactory } from '@nestjs/core'; import { ValidationPipe } from '@nestjs/common'; +import { NestExpressApplication } from '@nestjs/platform-express'; import { SwaggerModule, DocumentBuilder } from '@nestjs/swagger'; import { AppModule } from './app.module'; import { AllExceptionsFilter } from './common/http-exception.filter'; import { randomUUID } from 'crypto'; -import { NextFunction, Request, Response } from 'express'; +import { NextFunction, Request, Response, json, urlencoded } from 'express'; +import { StructuredLoggerService } from './common/structured-logger.service'; +import { StructuredRequestLoggerMiddleware } from './common/structured-request-logger.middleware'; + +interface RawBodyRequest extends Request { + rawBody?: Buffer; +} async function bootstrap() { - // rawBody is needed by the payment webhook controller to verify HMAC - // signatures against the exact bytes the provider signed, not a - // re-serialized copy of the parsed JSON body. - const app = await NestFactory.create(AppModule, { rawBody: true }); + /** + * Production deployments use newline-delimited JSON logs so they can be + * indexed directly by a log aggregator. LOG_JSON=true is an explicit + * override for other environments; development keeps Nest's readable logger. + */ + const structuredLoggingEnabled = + process.env.NODE_ENV === 'production' || process.env.LOG_JSON === 'true'; + const structuredLogger = structuredLoggingEnabled + ? new StructuredLoggerService() + : undefined; + const app = await NestFactory.create(AppModule, { + rawBody: true, + bodyParser: false, + ...(structuredLogger ? { logger: structuredLogger } : {}), + }); + + if (structuredLoggingEnabled) { + const requestLogger = new StructuredRequestLoggerMiddleware(); + app.use((req, res, next) => requestLogger.use(req, res, next)); + } + + const requestBodyLimit = process.env.REQUEST_BODY_LIMIT ?? '1mb'; + const captureRawBody = ( + req: Request, + _res: Response, + buf: Buffer, + ): void => { + (req as RawBodyRequest).rawBody = Buffer.from(buf); + }; + + /** + * Nest's implicit body parsers have no configured request-size bound, so + * register both parsers explicitly. The verify callback preserves the exact + * bytes for webhook HMAC verification while REQUEST_BODY_LIMIT caps memory. + */ + app.use(json({ limit: requestBodyLimit, verify: captureRawBody })); + app.use( + urlencoded({ + extended: true, + limit: requestBodyLimit, + verify: captureRawBody, + }), + ); app.use((req: Request, res: Response, next: NextFunction) => { res.setHeader('x-dns-prefetch-control', 'off'); diff --git a/backend/src/payments/admin-swagger.spec.ts b/backend/src/payments/admin-swagger.spec.ts new file mode 100644 index 00000000..70f285ac --- /dev/null +++ b/backend/src/payments/admin-swagger.spec.ts @@ -0,0 +1,45 @@ +import 'reflect-metadata'; +import { CreditsAdminController } from '../credits/credits-admin.controller'; +import { PaymentsAdminController } from './payments-admin.controller'; + +type Handler = (...args: any[]) => unknown; +type SwaggerResponse = { + description?: string; + type?: unknown; +}; + +describe('admin Swagger metadata', () => { + it('documents every admin handler with a response', () => { + for (const Controller of [ + PaymentsAdminController, + CreditsAdminController, + ]) { + const prototype = Controller.prototype as unknown as Record< + string, + Handler + >; + const methodNames = Object.getOwnPropertyNames(prototype).filter( + (name) => name !== 'constructor', + ); + + expect(methodNames.length).toBeGreaterThan(0); + for (const methodName of methodNames) { + const handler = prototype[methodName]; + expect(Reflect.getMetadata('swagger/apiOperation', handler)).toBeDefined(); + + const responses = Reflect.getMetadata( + 'swagger/apiResponse', + handler, + ) as Record | undefined; + expect(responses).toBeDefined(); + expect(Object.keys(responses ?? {}).length).toBeGreaterThan(0); + expect( + Object.values(responses ?? {}).some( + (response) => + response.type !== undefined || Boolean(response.description), + ), + ).toBe(true); + } + } + }); +}); diff --git a/backend/src/payments/dto/reconciliation-metrics-response.dto.ts b/backend/src/payments/dto/reconciliation-metrics-response.dto.ts new file mode 100644 index 00000000..7316d061 --- /dev/null +++ b/backend/src/payments/dto/reconciliation-metrics-response.dto.ts @@ -0,0 +1,22 @@ +import { ApiProperty } from '@nestjs/swagger'; +import type { ReconciliationMetrics } from '../reconciliation.service'; + +/** + * The operational signal exposed to administrators when the reconciliation + * queue needs attention. The service owns the calculation; this DTO keeps the + * wire shape explicit and stable for OpenAPI consumers. + */ +export class ReconciliationMetricsResponseDto { + @ApiProperty({ description: 'Payments currently awaiting manual review' }) + manualReviewQueueDepth: number; + + @ApiProperty({ description: 'Queue depth at which alerting is enabled' }) + alertThreshold: number; + + @ApiProperty({ description: 'True when the queue exceeds the alert threshold' }) + alerting: boolean; + + static fromView(view: ReconciliationMetrics): ReconciliationMetricsResponseDto { + return Object.assign(new ReconciliationMetricsResponseDto(), view); + } +} diff --git a/backend/src/payments/payments-admin.controller.ts b/backend/src/payments/payments-admin.controller.ts index 2d0a2c54..4d76c9f3 100644 --- a/backend/src/payments/payments-admin.controller.ts +++ b/backend/src/payments/payments-admin.controller.ts @@ -10,6 +10,7 @@ import { import { ApiBearerAuth, ApiOperation, + ApiParam, ApiResponse, ApiTags, } from '@nestjs/swagger'; @@ -24,6 +25,7 @@ import { RefundsService } from './refunds.service'; import { AdminActionLogService } from '../admin-audit/admin-action-log.service'; import { AdminActionType } from '../admin-audit/admin-action-type.enum'; import { PaymentResponseDto } from './dto/payment-response.dto'; +import { ReconciliationMetricsResponseDto } from './dto/reconciliation-metrics-response.dto'; import { ResolvePaymentManuallyDto } from './dto/resolve-payment-manually.dto'; import { VoidPaymentDto } from './dto/void-payment.dto'; import { CreateRefundDto } from './dto/create-refund.dto'; @@ -65,11 +67,15 @@ export class PaymentsAdminController { summary: 'Reconciliation metrics — manual-review queue depth and alert status', }) - getMetrics() { - return this.reconciliationService.getMetrics(); + @ApiResponse({ status: 200, type: ReconciliationMetricsResponseDto }) + async getMetrics(): Promise { + return ReconciliationMetricsResponseDto.fromView( + await this.reconciliationService.getMetrics(), + ); } @Post(':id/force-reconcile') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Immediately re-verify one payment against the provider, bypassing the due-schedule', @@ -83,6 +89,7 @@ export class PaymentsAdminController { } @Post(':id/resolve-manually') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Resolve a MANUAL_REVIEW payment by hand (reason required, audited)', @@ -109,6 +116,7 @@ export class PaymentsAdminController { } @Post(':id/void') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Void a MANUAL_REVIEW payment without resolving it (reason required, audited)', @@ -131,6 +139,7 @@ export class PaymentsAdminController { } @Post(':id/refunds') + @ApiParam({ name: 'id', type: String, format: 'uuid' }) @ApiOperation({ summary: 'Issue a (partial) refund against a CONFIRMED/PARTIALLY_REFUNDED payment', diff --git a/backend/src/payments/payments.service.spec.ts b/backend/src/payments/payments.service.spec.ts index d42e7f12..473e113d 100644 --- a/backend/src/payments/payments.service.spec.ts +++ b/backend/src/payments/payments.service.spec.ts @@ -3,6 +3,7 @@ import { ConflictException, ForbiddenException, NotFoundException, + ServiceUnavailableException, } from '@nestjs/common'; import { PaymentsService } from './payments.service'; import { Payment } from './entities/payment.entity'; @@ -66,7 +67,9 @@ describe('PaymentsService', () => { initiate: jest.fn().mockResolvedValue({ providerReference: 'ref-1' }), }; railRegistry = { get: jest.fn().mockReturnValue(railAdapter) }; - config = { get: jest.fn().mockReturnValue(30) }; + config = { + get: jest.fn((_key: string, defaultValue: unknown) => defaultValue), + }; metrics = { recordPaymentTransition: jest.fn() }; service = new PaymentsService( repository as any, @@ -84,7 +87,7 @@ describe('PaymentsService', () => { expect(repository.save).not.toHaveBeenCalled(); }); - it('creates a new INITIATED payment then progresses it to AWAITING_CONFIRMATION', async () => { + it('routes a successful provider initiation through the per-rail breaker and progresses it to AWAITING_CONFIRMATION', async () => { repository.findOne.mockResolvedValueOnce(null); // no existing idempotency-key row repository.findOne.mockResolvedValueOnce(null); // booking is free @@ -96,6 +99,32 @@ describe('PaymentsService', () => { expect(repository.save).toHaveBeenCalledTimes(2); }); + it('opens the rail breaker and returns 503 without retrying the provider', async () => { + repository.findOne.mockResolvedValue(null); + railAdapter.initiate.mockRejectedValue(new Error('provider unavailable')); + + for (let attempt = 0; attempt < 10; attempt += 1) { + await expect( + service.initiate( + `user-${attempt}`, + `key-${attempt}`, + makeDto({ bookingId: `booking-${attempt}` }), + ), + ).rejects.toThrow('provider unavailable'); + } + + expect(railAdapter.initiate).toHaveBeenCalledTimes(10); + + await expect( + service.initiate( + 'user-10', + 'key-10', + makeDto({ bookingId: 'booking-10' }), + ), + ).rejects.toThrow(ServiceUnavailableException); + expect(railAdapter.initiate).toHaveBeenCalledTimes(10); + }); + it('replays the same Idempotency-Key and returns the original payment without creating a duplicate', async () => { const existing = { id: 'p-1', diff --git a/backend/src/payments/payments.service.ts b/backend/src/payments/payments.service.ts index b2e6ccc8..b227fd13 100644 --- a/backend/src/payments/payments.service.ts +++ b/backend/src/payments/payments.service.ts @@ -19,6 +19,10 @@ import { } from './enums/payment-status.enum'; import { assertValidTransition } from './payment-state-machine'; import { PaymentRailRegistry } from './payment-rail-registry'; +import { PaymentRail } from './enums/payment-rail.enum'; +import type { PaymentInitiationResult } from './interfaces/payment-rail-adapter.interface'; +import { createPaymentProviderCircuitBreaker } from './utils/circuit-breaker'; +import type { PaymentProviderCircuitBreaker } from './utils/circuit-breaker'; const USER_IDEMPOTENCY_KEY_CONSTRAINT = 'uq_payments_user_id_idempotency_key'; const BOOKING_NON_TERMINAL_CONSTRAINT = 'uq_payments_booking_id_non_terminal'; @@ -31,6 +35,11 @@ const BLOCKING_STATUSES_FOR_NEW_PAYMENT = [ @Injectable() export class PaymentsService { + private readonly breakers = new Map< + string, + PaymentProviderCircuitBreaker<[Payment], PaymentInitiationResult> + >(); + constructor( @InjectRepository(Payment) private readonly paymentRepository: Repository, @@ -159,10 +168,35 @@ export class PaymentsService { }); } + /** + * Breakers are created lazily and kept per rail so one provider's outage + * cannot suppress traffic to a healthy rail. Opossum closes a breaker after + * a successful half-open action, so no manual reset is needed. + */ + private getBreaker( + rail: PaymentRail, + ): PaymentProviderCircuitBreaker<[Payment], PaymentInitiationResult> { + let breaker = this.breakers.get(rail); + if (breaker) { + return breaker; + } + + const railAdapter = this.railRegistry.get(rail); + const action = (paymentToInitiate: Payment) => + railAdapter.initiate(paymentToInitiate); + breaker = createPaymentProviderCircuitBreaker( + action, + this.config, + rail.toLowerCase(), + ); + this.breakers.set(rail, breaker); + return breaker; + } + private async progressToAwaitingConfirmation( payment: Payment, ): Promise { - const result = await this.railRegistry.get(payment.rail).initiate(payment); + const result = await this.getBreaker(payment.rail).fire(payment); payment.providerReference = result.providerReference; this.transitionStatus(payment, PaymentStatus.AWAITING_CONFIRMATION); return this.paymentRepository.save(payment); diff --git a/backend/src/payments/utils/circuit-breaker.ts b/backend/src/payments/utils/circuit-breaker.ts new file mode 100644 index 00000000..83f9f804 --- /dev/null +++ b/backend/src/payments/utils/circuit-breaker.ts @@ -0,0 +1,85 @@ +import { + ConflictException, + ServiceUnavailableException, +} from '@nestjs/common'; +import { ConfigService } from '@nestjs/config'; +import CircuitBreakerModule = require('opossum'); + +export interface PaymentProviderCircuitBreaker< + TArgs extends unknown[], + TResult, +> { + fire(...args: TArgs): Promise; +} + +interface PaymentProviderCircuitBreakerInstance< + TArgs extends unknown[], + TResult, +> extends PaymentProviderCircuitBreaker { + fallback( + handler: (...args: [...TArgs, unknown]) => TResult | Promise, + ): PaymentProviderCircuitBreakerInstance; +} + +interface CircuitBreakerOptions { + timeout: number; + errorThresholdPercentage: number; + resetTimeout: number; + volumeThreshold: number; + name: string; + errorFilter: (error: unknown) => boolean; +} + +interface CircuitBreakerConstructor { + isOurError(error: Error): boolean; + + new ( + action: (...args: TArgs) => Promise, + options: CircuitBreakerOptions, + ): PaymentProviderCircuitBreakerInstance; +} + +const CircuitBreaker = + CircuitBreakerModule as unknown as CircuitBreakerConstructor; + +/** + * Creates a breaker for one external provider action. A ConflictException is + * a domain/business rejection and is filtered out; transport, timeout, and + * other provider failures count toward opening the circuit. The volume + * threshold keeps a low-traffic rail from opening after a single transient + * error, while the fallback turns an open circuit into a clear HTTP 503 + * response without hiding the original error while the rail is still closed. + */ +export function createPaymentProviderCircuitBreaker< + TArgs extends unknown[], + TResult, +>( + action: (...args: TArgs) => Promise, + config: ConfigService, + name: string, +): PaymentProviderCircuitBreaker { + const breaker = new CircuitBreaker(action, { + timeout: config.get('PAYMENT_PROVIDER_BREAKER_TIMEOUT_MS', 10000), + errorThresholdPercentage: config.get( + 'PAYMENT_PROVIDER_BREAKER_ERROR_THRESHOLD_PERCENT', + 50, + ), + resetTimeout: config.get( + 'PAYMENT_PROVIDER_BREAKER_RESET_TIMEOUT_MS', + 30000, + ), + volumeThreshold: 10, + name: `payment-provider-${name}`, + errorFilter: (error: unknown) => error instanceof ConflictException, + }); + + return breaker.fallback((...args) => { + const error = args[args.length - 1]; + if (error instanceof Error && CircuitBreaker.isOurError(error)) { + throw new ServiceUnavailableException( + 'Payment provider temporarily unavailable, please retry', + ); + } + throw error; + }); +}