diff --git a/backend/.env.example b/backend/.env.example index 69a43f84..2a77a906 100644 --- a/backend/.env.example +++ b/backend/.env.example @@ -10,6 +10,16 @@ # PORT and NODE_ENV are read by src/main.ts and the auth cookie settings. 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= + # Exact, comma-separated browser origins. Whitespace is ignored. In development # and test, an empty value falls back to localhost:3000 and localhost:3001; in # production an unset or empty value is an empty allowlist and denies all @@ -135,6 +145,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 cc00152e..bb4a9cd9 100644 --- a/backend/package.json +++ b/backend/package.json @@ -77,12 +77,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 70423c5e..1ce4e20b 100644 --- a/backend/src/main.ts +++ b/backend/src/main.ts @@ -1,5 +1,6 @@ import { NestFactory } from '@nestjs/core'; import { Logger, 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'; @@ -10,16 +11,61 @@ import { resolveGracefulShutdownTimeout, } from './common/graceful-shutdown'; import { HttpLoggingInterceptor } from './common/http-logging.interceptor'; -import { randomUUID } from 'crypto'; -import { NextFunction, Request, Response } from 'express'; +import { NextFunction, Request, Response, json, urlencoded } from 'express'; import { initTracing } from './common/tracing'; +import { StructuredLoggerService } from './common/structured-logger.service'; +import { StructuredRequestLoggerMiddleware } from './common/structured-request-logger.middleware'; + +interface RawBodyRequest extends Request { + rawBody?: Buffer; +} async function bootstrap() { initTracing(); - // 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 9e67e2ed..a9967253 100644 --- a/backend/src/payments/payments-admin.controller.ts +++ b/backend/src/payments/payments-admin.controller.ts @@ -11,6 +11,7 @@ import { import { ApiBearerAuth, ApiOperation, + ApiParam, ApiQuery, ApiResponse, ApiTags, @@ -26,6 +27,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 { ReconciliationRunResponseDto } from './dto/reconciliation-run-response.dto'; import { ResolvePaymentManuallyDto } from './dto/resolve-payment-manually.dto'; import { VoidPaymentDto } from './dto/void-payment.dto'; @@ -68,8 +70,11 @@ 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(), + ); } @Get('reconciliation-runs') @@ -93,6 +98,7 @@ export class PaymentsAdminController { } @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', @@ -106,6 +112,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)', @@ -132,6 +139,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)', @@ -154,6 +162,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 3ffa9ebb..2e39b836 100644 --- a/backend/src/payments/payments.service.ts +++ b/backend/src/payments/payments.service.ts @@ -20,6 +20,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'; @@ -32,6 +36,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, @@ -172,12 +181,37 @@ 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 withSpan( 'payments.rail.initiate', - () => this.railRegistry.get(payment.rail).initiate(payment), + () => this.getBreaker(payment.rail).fire(payment), { 'payment.id': payment.id, 'payment.rail': payment.rail, 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; + }); +}