diff --git a/CLAUDE.md b/CLAUDE.md index c06000dac..d67e6c3d5 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -12,6 +12,7 @@ Monorepo for the Ethio Djibouti Railway (EDR) digital platform. Contains the Fre | `edr-freight-web/portal` | `@edr/freight-portal` | React frontend for freight customer/portal users | 5173 | | `edr-freight-web/backoffice` | `@edr/freight-backoffice` | React frontend for freight backoffice employees | 5183 | | `edr-passenger-api` | `@edr/passenger-api` | NestJS API for passenger management | 3002 | +| `edr-payment-api` | `@edr/payment-api` | NestJS payment microservice (intents, webhooks) | 3003 | | `edr-passenger-web/portal` | `@edr/passenger-portal` | React frontend for passenger customer/portal users | 5174 | | `edr-passenger-web/backoffice` | `@edr/passenger-backoffice` | React frontend for passenger backoffice employees | 5184 | @@ -73,6 +74,7 @@ The `@CurrentUser`, `@Roles`, and `@Public` decorators in `@edr/api-common` are - `edr-freight-web/portal`: 5173 - `edr-freight-web/backoffice`: 5183 - `edr-passenger-api`: 3002 +- `edr-payment-api`: 3003 - `edr-passenger-web/portal`: 5174 - `edr-passenger-web/backoffice`: 5184 @@ -80,6 +82,7 @@ The `@CurrentUser`, `@Roles`, and `@Public` decorators in `@edr/api-common` are - `postgres-freight` (port 5433): database `edr_freight` — freight API only. - `postgres-passenger` (port 5434): database `edr_passenger` — passenger API only. +- `edr_payment` schema — lives in the same Postgres database as the domain system (whatever the passenger `DATABASE_URL` points at) but is owned exclusively by `apps/edr-payment-api`. Dedicated DB user, no cross-schema FKs, domain apps have no grants on it (see `docs/payment-service/`). - Each app owns its own DB. No cross-database joins; cross-domain data flows through API calls or message queues. ## Adding a new module to a NestJS app diff --git a/apps/edr-payment-api/nest-cli.json b/apps/edr-payment-api/nest-cli.json new file mode 100644 index 000000000..89d7d6c57 --- /dev/null +++ b/apps/edr-payment-api/nest-cli.json @@ -0,0 +1,8 @@ +{ + "$schema": "https://json.schemastore.org/nest-cli", + "collection": "@nestjs/schematics", + "sourceRoot": "src", + "compilerOptions": { + "deleteOutDir": false + } +} diff --git a/apps/edr-payment-api/package.json b/apps/edr-payment-api/package.json new file mode 100644 index 000000000..3ce7804a9 --- /dev/null +++ b/apps/edr-payment-api/package.json @@ -0,0 +1,74 @@ +{ + "name": "@edr/payment-api", + "version": "0.0.0", + "private": true, + "description": "EDR Payment Microservice — owns provider integration, payment intents, webhooks, and outbox notifications for the whole platform", + "scripts": { + "clean": "node -e \"const fs=require('fs'); fs.rmSync('dist',{recursive:true,force:true}); fs.rmSync('.tsbuildinfo',{force:true});\"", + "predev": "pnpm run clean", + "dev": "nest start --watch", + "prebuild": "pnpm run clean", + "build": "nest build", + "start": "node dist/main.js", + "lint": "eslint src", + "test": "jest", + "type-check": "tsc --noEmit", + "migration:run": "ts-node src/scripts/migrate.ts", + "migration:revert": "ts-node src/scripts/migrate-revert.ts" + }, + "dependencies": { + "@edr/api-common": "workspace:*", + "@edr/payment-providers": "workspace:*", + "@edr/types": "workspace:*", + "@nestjs/axios": "^4.0.1", + "@nestjs/common": "^11.0.0", + "@nestjs/config": "^4.0.0", + "@nestjs/core": "^11.0.0", + "@nestjs/platform-express": "^11.0.0", + "@nestjs/schedule": "^6.0.0", + "@nestjs/swagger": "^11.4.2", + "@nestjs/typeorm": "^11.0.1", + "axios": "^1.16.1", + "class-transformer": "^0.5.1", + "class-validator": "^0.14.1", + "dotenv": "^17.4.2", + "pg": "^8.13.0", + "reflect-metadata": "^0.2.2", + "rxjs": "^7.8.1", + "typeorm": "0.3.30" + }, + "devDependencies": { + "@edr/eslint-config": "workspace:*", + "@edr/tsconfig": "workspace:*", + "@nestjs/cli": "^11.0.0", + "@nestjs/schematics": "^11.0.0", + "@nestjs/testing": "^11.0.0", + "@types/express": "^5.0.0", + "@types/jest": "^29.5.13", + "@types/node": "^20.14.0", + "@types/pg": "^8.6.7", + "jest": "^29.7.0", + "ts-jest": "^29.2.5", + "ts-loader": "^9.5.1", + "ts-node": "^10.9.2", + "tsconfig-paths": "^4.2.0", + "typescript": "^5.5.4" + }, + "jest": { + "moduleFileExtensions": [ + "js", + "json", + "ts" + ], + "rootDir": "src", + "testRegex": ".*\\.spec\\.ts$", + "transform": { + "^.+\\.(t|j)s$": "ts-jest" + }, + "collectCoverageFrom": [ + "**/*.(t|j)s" + ], + "coverageDirectory": "../coverage", + "testEnvironment": "node" + } +} diff --git a/apps/edr-payment-api/src/app.module.ts b/apps/edr-payment-api/src/app.module.ts new file mode 100644 index 000000000..4a99472d0 --- /dev/null +++ b/apps/edr-payment-api/src/app.module.ts @@ -0,0 +1,51 @@ +import { Module } from "@nestjs/common"; +import { ConfigModule, ConfigService } from "@nestjs/config"; +import { ScheduleModule } from "@nestjs/schedule"; +import { TypeOrmModule, TypeOrmModuleOptions } from "@nestjs/typeorm"; +import appConfig from "./config/app.config"; +import databaseConfig from "./config/database.config"; +import notifierConfig from "./config/notifier.config"; +import telebirrConfig from "./config/telebirr.config"; +import waafiConfig from "./config/waafi.config"; +import cbeConfig from "./config/cbe.config"; +import ebirrConfig from "./config/ebirr.config"; +import cardConfig from "./config/card.config"; +import dmoneyConfig from "./config/dmoney.config"; +import { HealthModule } from "./modules/health/health.module"; +import { IntentsModule } from "./modules/intents/intents.module"; +import { OutboxModule } from "./modules/outbox/outbox.module"; +import { ProvidersModule } from "./modules/providers/providers.module"; +import { ReconciliationModule } from "./modules/reconciliation/reconciliation.module"; +import { WebhooksModule } from "./modules/webhooks/webhooks.module"; + +@Module({ + imports: [ + ConfigModule.forRoot({ + isGlobal: true, + load: [ + appConfig, + databaseConfig, + notifierConfig, + telebirrConfig, + waafiConfig, + cbeConfig, + ebirrConfig, + cardConfig, + dmoneyConfig, + ], + }), + TypeOrmModule.forRootAsync({ + inject: [ConfigService], + useFactory: (config: ConfigService) => + config.get("database") as TypeOrmModuleOptions, + }), + ScheduleModule.forRoot(), + HealthModule, + ProvidersModule, + IntentsModule, + WebhooksModule, + OutboxModule, + ReconciliationModule, + ], +}) +export class AppModule {} diff --git a/apps/edr-payment-api/src/common/guards/service-auth.guard.ts b/apps/edr-payment-api/src/common/guards/service-auth.guard.ts new file mode 100644 index 000000000..5b5fb39c5 --- /dev/null +++ b/apps/edr-payment-api/src/common/guards/service-auth.guard.ts @@ -0,0 +1,55 @@ +import { + CanActivate, + ExecutionContext, + Injectable, + Logger, + UnauthorizedException, +} from "@nestjs/common"; +import { ConfigService } from "@nestjs/config"; +import { timingSafeEqual } from "node:crypto"; +import { Request } from "express"; + +/** + * Shared-secret service-to-service auth for the internal surface (/payments/*). + * Callers send `x-service-token: ` (or `Authorization: Bearer …`). + * Webhook endpoints are intentionally NOT behind this guard — they are provider-facing and + * authenticate via signature verification instead. + */ +@Injectable() +export class ServiceAuthGuard implements CanActivate { + private readonly logger = new Logger(ServiceAuthGuard.name); + private readonly token: string; + private warned = false; + + constructor(config: ConfigService) { + this.token = config.get("app.serviceAuthToken") ?? ""; + if (!this.token && process.env.NODE_ENV === "production") { + throw new Error("SERVICE_AUTH_TOKEN must be set in production"); + } + } + + canActivate(context: ExecutionContext): boolean { + if (!this.token) { + if (!this.warned) { + this.logger.warn( + "SERVICE_AUTH_TOKEN unset — internal endpoints are UNGUARDED (dev only)", + ); + this.warned = true; + } + return true; + } + + const request = context.switchToHttp().getRequest(); + const header = request.headers["x-service-token"]; + const bearer = request.headers.authorization?.replace(/^Bearer\s+/i, ""); + const presented = + (Array.isArray(header) ? header[0] : header) ?? bearer ?? ""; + + const expected = Buffer.from(this.token); + const actual = Buffer.from(presented); + const valid = + expected.length === actual.length && timingSafeEqual(expected, actual); + if (!valid) throw new UnauthorizedException("Invalid service token"); + return true; + } +} diff --git a/apps/edr-payment-api/src/config/app.config.ts b/apps/edr-payment-api/src/config/app.config.ts new file mode 100644 index 000000000..c1128a872 --- /dev/null +++ b/apps/edr-payment-api/src/config/app.config.ts @@ -0,0 +1,21 @@ +import { registerAs } from "@nestjs/config"; + +export default registerAs("app", () => ({ + port: parseInt(process.env.PORT ?? "3003", 10), + /** + * Shared secret for service-to-service auth (apps -> /payments/*, payment -> mark-paid). + * Required in production; in development an empty value disables the guard with a warning. + * TODO: integrate @tria-plc IAM / mTLS as the long-term mechanism (docs/payment-service §14). + */ + serviceAuthToken: process.env.SERVICE_AUTH_TOKEN ?? "", + reconciliation: { + /** How often the stale-intent sweep runs. */ + sweepIntervalMs: parseInt( + process.env.RECONCILE_SWEEP_INTERVAL_MS ?? "60000", + 10, + ), + /** An intent is "stale" when non-terminal and untouched for this long. */ + staleAfterMs: parseInt(process.env.RECONCILE_STALE_AFTER_MS ?? "60000", 10), + batchSize: parseInt(process.env.RECONCILE_BATCH_SIZE ?? "20", 10), + }, +})); diff --git a/apps/edr-payment-api/src/config/card.config.ts b/apps/edr-payment-api/src/config/card.config.ts new file mode 100644 index 000000000..fa9d4b30c --- /dev/null +++ b/apps/edr-payment-api/src/config/card.config.ts @@ -0,0 +1,9 @@ +import { registerAs } from "@nestjs/config"; + +export default registerAs("card", () => ({ + baseUrl: process.env.CARD_BASE_URL || "", + apiKey: process.env.CARD_API_KEY || "", + webhookSecret: process.env.CARD_WEBHOOK_SECRET || "", + webhookUrl: process.env.CARD_WEBHOOK_URL || "", + returnUrl: process.env.CARD_RETURN_URL || "", +})); diff --git a/apps/edr-payment-api/src/config/cbe.config.ts b/apps/edr-payment-api/src/config/cbe.config.ts new file mode 100644 index 000000000..eb2cd689f --- /dev/null +++ b/apps/edr-payment-api/src/config/cbe.config.ts @@ -0,0 +1,9 @@ +import { registerAs } from "@nestjs/config"; + +export default registerAs("cbe", () => ({ + baseUrl: process.env.CBE_BASE_URL || "", + merchantId: process.env.CBE_MERCHANT_ID || "", + secretKey: process.env.CBE_SECRET_KEY || "", + notifyUrl: process.env.CBE_NOTIFY_URL || "", + returnUrl: process.env.CBE_RETURN_URL || "", +})); diff --git a/apps/edr-payment-api/src/config/database.config.ts b/apps/edr-payment-api/src/config/database.config.ts new file mode 100644 index 000000000..43d8ca564 --- /dev/null +++ b/apps/edr-payment-api/src/config/database.config.ts @@ -0,0 +1,36 @@ +import { registerAs } from "@nestjs/config"; +import { TypeOrmModuleOptions } from "@nestjs/typeorm"; +import { DataSourceOptions } from "typeorm"; + +/** + * Shared connection options for the Nest TypeORM module and the standalone DataSource + * (migration CLI). Payment tables live in the SAME Postgres database as the domain system + * (edr_database by default) but in the dedicated `edr_payment` schema; logical ownership is + * enforced with a dedicated DB user in non-dev environments (grants only on this schema). + */ +export function buildDataSourceOptions(): DataSourceOptions { + return { + type: "postgres", + host: process.env.DB_HOST ?? "localhost", + port: parseInt(process.env.DB_PORT ?? "5432", 10), + username: process.env.DB_USER ?? "edr", + password: process.env.DB_PASSWORD ?? "", + database: process.env.DB_NAME ?? "edr_database", + schema: process.env.DB_SCHEMA ?? "edr_payment", + entities: [__dirname + "/../**/*.entity.{ts,js}"], + migrations: [__dirname + "/../migrations/*.{ts,js}"], + // Schema changes go through migrations only — never synchronize (house rule). + synchronize: false, + logging: process.env.NODE_ENV === "development", + }; +} + +export default registerAs( + "database", + (): TypeOrmModuleOptions => ({ + ...buildDataSourceOptions(), + autoLoadEntities: true, + // Run pending migrations on boot (main.ts ensures the database/schema exist first). + migrationsRun: true, + }), +); diff --git a/apps/edr-payment-api/src/config/dmoney.config.ts b/apps/edr-payment-api/src/config/dmoney.config.ts new file mode 100644 index 000000000..78751f8d5 --- /dev/null +++ b/apps/edr-payment-api/src/config/dmoney.config.ts @@ -0,0 +1,10 @@ +import { registerAs } from "@nestjs/config"; + +export default registerAs("dmoney", () => ({ + baseUrl: process.env.DMONEY_BASE_URL ?? "", + appId: process.env.DMONEY_APP_ID ?? "", + appSecret: process.env.DMONEY_APP_SECRET ?? "", + publicKey: process.env.DMONEY_PUBLIC_KEY ?? "", + privateKey: process.env.DMONEY_PRIVATE_KEY ?? "", + notifyUrl: process.env.DMONEY_NOTIFY_URL ?? "", +})); diff --git a/apps/edr-payment-api/src/config/ebirr.config.ts b/apps/edr-payment-api/src/config/ebirr.config.ts new file mode 100644 index 000000000..4c5dd68e6 --- /dev/null +++ b/apps/edr-payment-api/src/config/ebirr.config.ts @@ -0,0 +1,9 @@ +import { registerAs } from "@nestjs/config"; + +export default registerAs("ebirr", () => ({ + baseUrl: process.env.EBIRR_BASE_URL || "", + merchantCode: process.env.EBIRR_MERCHANT_CODE || "", + secretKey: process.env.EBIRR_SECRET_KEY || "", + notifyUrl: process.env.EBIRR_NOTIFY_URL || "", + returnUrl: process.env.EBIRR_RETURN_URL || "", +})); diff --git a/apps/edr-payment-api/src/config/ensure-schema.ts b/apps/edr-payment-api/src/config/ensure-schema.ts new file mode 100644 index 000000000..9f7a5bb52 --- /dev/null +++ b/apps/edr-payment-api/src/config/ensure-schema.ts @@ -0,0 +1,52 @@ +import { Client } from "pg"; + +const IDENTIFIER = /^[a-z_][a-z0-9_]*$/; + +function connectionEnv() { + return { + host: process.env.DB_HOST ?? "localhost", + port: parseInt(process.env.DB_PORT ?? "5432", 10), + user: process.env.DB_USER ?? "edr", + password: process.env.DB_PASSWORD ?? "", + }; +} + +/** + * Dev/bootstrap convenience: make sure the `edr_payment` schema exists in the shared + * database before TypeORM initializes (the migrations table itself lives in the schema, so + * migrations cannot create it). In production the schema/grants are provisioned out-of-band + * by ops; this is then a no-op. + */ +export async function ensurePaymentSchema(): Promise { + const database = process.env.DB_NAME ?? "edr_database"; + const schema = process.env.DB_SCHEMA ?? "edr_payment"; + if (!IDENTIFIER.test(database) || !IDENTIFIER.test(schema)) { + throw new Error( + `Invalid DB_NAME/DB_SCHEMA identifier: ${database}/${schema}`, + ); + } + + let client = new Client({ ...connectionEnv(), database }); + try { + await client.connect(); + } catch (err) { + // 3D000 = database does not exist — create it from the maintenance DB, then reconnect. + if ((err as { code?: string }).code !== "3D000") throw err; + await client.end().catch(() => undefined); + const admin = new Client({ ...connectionEnv(), database: "postgres" }); + await admin.connect(); + try { + await admin.query(`CREATE DATABASE "${database}"`); + } finally { + await admin.end(); + } + client = new Client({ ...connectionEnv(), database }); + await client.connect(); + } + + try { + await client.query(`CREATE SCHEMA IF NOT EXISTS "${schema}"`); + } finally { + await client.end(); + } +} diff --git a/apps/edr-payment-api/src/config/notifier.config.ts b/apps/edr-payment-api/src/config/notifier.config.ts new file mode 100644 index 000000000..cddc66761 --- /dev/null +++ b/apps/edr-payment-api/src/config/notifier.config.ts @@ -0,0 +1,15 @@ +import { registerAs } from "@nestjs/config"; + +export default registerAs("notifier", () => ({ + /** mark-paid callback URL per owning service (PaymentService discriminator routes here). */ + passengerUrl: + process.env.PAYMENT_NOTIFY_PASSENGER_URL ?? + "http://localhost:3002/internal/payments/mark-paid", + freightUrl: + process.env.PAYMENT_NOTIFY_FREIGHT_URL ?? + "http://localhost:3001/internal/payments/mark-paid", + relayIntervalMs: parseInt(process.env.OUTBOX_RELAY_INTERVAL_MS ?? "5000", 10), + maxAttempts: parseInt(process.env.OUTBOX_MAX_ATTEMPTS ?? "10", 10), + httpTimeoutMs: parseInt(process.env.NOTIFY_HTTP_TIMEOUT_MS ?? "10000", 10), + relayBatchSize: parseInt(process.env.OUTBOX_RELAY_BATCH_SIZE ?? "20", 10), +})); diff --git a/apps/edr-payment-api/src/config/telebirr.config.ts b/apps/edr-payment-api/src/config/telebirr.config.ts new file mode 100644 index 000000000..8e5d1712a --- /dev/null +++ b/apps/edr-payment-api/src/config/telebirr.config.ts @@ -0,0 +1,16 @@ +import { registerAs } from "@nestjs/config"; + +export default registerAs("telebirr", () => ({ + baseUrl: process.env.TELEBIRR_BASE_URL ?? "", + webBaseUrl: process.env.TELEBIRR_WEB_BASE_URL ?? "", + fabricAppId: process.env.TELEBIRR_FABRIC_APP_ID ?? "", + appSecret: process.env.TELEBIRR_APP_SECRET ?? "", + merchantAppId: process.env.TELEBIRR_MERCHANT_APP_ID ?? "", + merchantCode: process.env.TELEBIRR_MERCHANT_CODE ?? "", + notifyUrl: process.env.TELEBIRR_NOTIFY_URL ?? "", + returnUrl: process.env.TELEBIRR_RETURN_URL ?? "", + timeoutExpress: process.env.TELEBIRR_TIMEOUT_EXPRESS ?? "15m", + privateKey: process.env.TELEBIRR_PRIVATE_KEY ?? "", + publicKey: process.env.TELEBIRR_PUBLIC_KEY ?? "", + insecureTls: process.env.TELEBIRR_INSECURE_TLS === "true", +})); diff --git a/apps/edr-payment-api/src/config/waafi.config.ts b/apps/edr-payment-api/src/config/waafi.config.ts new file mode 100644 index 000000000..05624922f --- /dev/null +++ b/apps/edr-payment-api/src/config/waafi.config.ts @@ -0,0 +1,27 @@ +import { registerAs } from "@nestjs/config"; + +export default registerAs("waafi", () => ({ + // `/asm` is appended in the provider; use sandbox by default, switch to + // https://api.waafipay.net in production. + baseUrl: process.env.WAAFI_BASE_URL ?? "https://sandbox.waafipay.net", + // HPP credentials (Hosted Payment Page family). + merchantUid: process.env.WAAFI_MERCHANT_UID ?? "", + storeId: process.env.WAAFI_STORE_ID ?? "", + hppKey: process.env.WAAFI_HPP_KEY ?? "", + // HMAC secret returned once by WEBHOOK_REGISTER; verifies inbound webhooks. + webhookSecret: process.env.WAAFI_WEBHOOK_SECRET ?? "", + // Wallet payment method (EVC/ZAAD/Sahal) — MWALLET_ACCOUNT requires the payer phone up front. + paymentMethod: process.env.WAAFI_PAYMENT_METHOD ?? "MWALLET_ACCOUNT", + // Waafi has no ETB; when set this overrides the asserted currency (USD/DJF/SLSH). + currency: process.env.WAAFI_CURRENCY ?? "DJF", + // Browser redirect targets after the hosted page completes/fails (UX only; webhook is source of truth). + successUrl: process.env.WAAFI_HPP_SUCCESS_URL ?? "", + failureUrl: process.env.WAAFI_HPP_FAILURE_URL ?? "", + // Callback data format: 1 = POST, 2 = GET, 4 = Result Token. + respDataFormat: Number(process.env.WAAFI_HPP_RESP_FORMAT ?? "1"), + // Registered webhook URL (reference only; registration is performed out-of-band). + notifyUrl: process.env.WAAFI_NOTIFY_URL ?? "", + // DEV ONLY: disable TLS cert verification. The Waafi sandbox serves a *.waafi.com cert that + // does not match sandbox.waafipay.net (ERR_TLS_CERT_ALTNAME_INVALID). Never enable in prod. + insecureTls: process.env.WAAFI_INSECURE_TLS === "true", +})); diff --git a/apps/edr-payment-api/src/data-source.ts b/apps/edr-payment-api/src/data-source.ts new file mode 100644 index 000000000..2aa708300 --- /dev/null +++ b/apps/edr-payment-api/src/data-source.ts @@ -0,0 +1,8 @@ +import "dotenv/config"; +import { DataSource } from "typeorm"; +import { buildDataSourceOptions } from "./config/database.config"; + +/** Standalone DataSource for the TypeORM CLI and the migrate script. */ +export const AppDataSource = new DataSource(buildDataSourceOptions()); + +export default AppDataSource; diff --git a/apps/edr-payment-api/src/main.ts b/apps/edr-payment-api/src/main.ts new file mode 100644 index 000000000..c9ad69494 --- /dev/null +++ b/apps/edr-payment-api/src/main.ts @@ -0,0 +1,54 @@ +import "reflect-metadata"; +import "dotenv/config"; +import { NestFactory } from "@nestjs/core"; +import { ValidationPipe } from "@nestjs/common"; +import { DocumentBuilder, SwaggerModule } from "@nestjs/swagger"; +import { AppModule } from "./app.module"; +import { ensurePaymentSchema } from "./config/ensure-schema"; + +async function bootstrap() { + // The edr_payment database/schema must exist before TypeORM boots (migrationsRun: true). + await ensurePaymentSchema(); + + // rawBody: true buffers the unparsed request body onto req.rawBody so webhook handlers + // (e.g. Waafi HMAC verification) can sign over the exact bytes the provider signed. + const app = await NestFactory.create(AppModule, { rawBody: true }); + + app.useGlobalPipes( + new ValidationPipe({ + whitelist: true, + transform: true, + forbidUnknownValues: false, + }), + ); + + const config = new DocumentBuilder() + .setTitle("EDR Payment API") + .setDescription( + "Platform payment microservice: payment intents, provider integration, the single " + + "registered webhook per provider, and reliable (outbox) notification of the owning app. " + + "Internal endpoints (/payments/*) require the x-service-token header; /webhooks/* is the " + + "only public surface. See docs/payment-service/.", + ) + .setVersion("1.0.0") + .addApiKey( + { type: "apiKey", name: "x-service-token", in: "header" }, + "service-token", + ) + .build(); + SwaggerModule.setup( + "api-docs", + app, + SwaggerModule.createDocument(app, config), + { + customSiteTitle: "EDR Payment API", + swaggerOptions: { persistAuthorization: true }, + }, + ); + + const port = process.env.PORT ?? 3003; + await app.listen(port); + console.log(`🚀 EDR Payment API running on port ${port}`); + console.log(`📚 Swagger: http://localhost:${port}/api-docs`); +} +bootstrap(); diff --git a/apps/edr-payment-api/src/migrations/1781136000000-InitPaymentSchema.ts b/apps/edr-payment-api/src/migrations/1781136000000-InitPaymentSchema.ts new file mode 100644 index 000000000..bbff2e7f1 --- /dev/null +++ b/apps/edr-payment-api/src/migrations/1781136000000-InitPaymentSchema.ts @@ -0,0 +1,124 @@ +import { MigrationInterface, QueryRunner } from "typeorm"; + +/** + * Initial edr_payment schema: payment_intent, payment_webhook_event, notification_outbox. + * + * Enum-valued columns are varchar on purpose (values mirror the @edr/types enums) so new + * providers/statuses never need an ALTER TYPE. uuid defaults use gen_random_uuid() (built into + * Postgres 13+; no extension required). + */ +export class InitPaymentSchema1781136000000 implements MigrationInterface { + name = "InitPaymentSchema1781136000000"; + + public async up(queryRunner: QueryRunner): Promise { + // Defensive — bootstrap (ensure-schema) normally creates this before migrations run. + await queryRunner.query(`CREATE SCHEMA IF NOT EXISTS "edr_payment"`); + + await queryRunner.query(` + CREATE TABLE "edr_payment"."payment_intent" ( + "id" uuid NOT NULL DEFAULT gen_random_uuid(), + "created_at" timestamptz NOT NULL DEFAULT now(), + "updated_at" timestamptz NOT NULL DEFAULT now(), + "deleted_at" timestamptz, + "service" varchar(16) NOT NULL, + "reference_type" varchar(16) NOT NULL, + "reference_id" varchar(64) NOT NULL, + "merchant_order_id" varchar(64) NOT NULL, + "provider" varchar(16) NOT NULL, + "provider_order_id" varchar(128), + "provider_txn_id" varchar(128), + "amount_minor" integer NOT NULL, + "confirmed_amount_minor" integer, + "currency" varchar(8) NOT NULL, + "status" varchar(24) NOT NULL DEFAULT 'REQUIRES_ACTION', + "client_action" jsonb, + "failure_code" varchar(64), + "failure_message" text, + "idempotency_key" varchar(128), + "expires_at" timestamptz, + "paid_at" timestamptz, + "raw_initiation" jsonb, + CONSTRAINT "pk_payment_intent" PRIMARY KEY ("id"), + CONSTRAINT "uq_payment_intent_merchant_order_id" UNIQUE ("merchant_order_id") + ) + `); + // One ACTIVE intent per domain order; terminal-failed attempts remain as audit rows. + await queryRunner.query(` + CREATE UNIQUE INDEX "uq_payment_intent_active_reference" + ON "edr_payment"."payment_intent" ("service", "reference_type", "reference_id") + WHERE status NOT IN ('FAILED','CANCELLED') AND deleted_at IS NULL + `); + await queryRunner.query(` + CREATE INDEX "idx_payment_intent_provider_txn" + ON "edr_payment"."payment_intent" ("provider_txn_id") + `); + await queryRunner.query(` + CREATE INDEX "idx_payment_intent_sweep" + ON "edr_payment"."payment_intent" ("status", "updated_at") + `); + await queryRunner.query(` + CREATE INDEX "idx_payment_intent_idempotency" + ON "edr_payment"."payment_intent" ("service", "idempotency_key") + `); + + await queryRunner.query(` + CREATE TABLE "edr_payment"."payment_webhook_event" ( + "id" uuid NOT NULL DEFAULT gen_random_uuid(), + "created_at" timestamptz NOT NULL DEFAULT now(), + "updated_at" timestamptz NOT NULL DEFAULT now(), + "deleted_at" timestamptz, + "provider" varchar(16) NOT NULL, + "external_event_id" varchar(191) NOT NULL, + "merchant_order_id" varchar(64), + "provider_txn_id" varchar(128), + "signature_valid" boolean NOT NULL DEFAULT false, + "status" varchar(64), + "payload" jsonb NOT NULL, + "received_at" timestamptz NOT NULL DEFAULT now(), + "processed_at" timestamptz, + "processing_error" text, + CONSTRAINT "pk_payment_webhook_event" PRIMARY KEY ("id") + ) + `); + // The webhook dedupe key: duplicate provider deliveries hit this and short-circuit. + await queryRunner.query(` + CREATE UNIQUE INDEX "uq_payment_webhook_event_external" + ON "edr_payment"."payment_webhook_event" ("provider", "external_event_id") + `); + + await queryRunner.query(` + CREATE TABLE "edr_payment"."notification_outbox" ( + "id" uuid NOT NULL DEFAULT gen_random_uuid(), + "created_at" timestamptz NOT NULL DEFAULT now(), + "updated_at" timestamptz NOT NULL DEFAULT now(), + "deleted_at" timestamptz, + "event_type" varchar(32) NOT NULL, + "service" varchar(16) NOT NULL, + "intent_id" uuid NOT NULL, + "reference_type" varchar(16) NOT NULL, + "reference_id" varchar(64) NOT NULL, + "payload" jsonb NOT NULL, + "status" varchar(16) NOT NULL DEFAULT 'PENDING', + "attempts" integer NOT NULL DEFAULT 0, + "next_retry_at" timestamptz, + "last_error" text, + "sent_at" timestamptz, + CONSTRAINT "pk_notification_outbox" PRIMARY KEY ("id") + ) + `); + await queryRunner.query(` + CREATE INDEX "idx_notification_outbox_relay" + ON "edr_payment"."notification_outbox" ("status", "next_retry_at") + `); + await queryRunner.query(` + CREATE INDEX "idx_notification_outbox_intent" + ON "edr_payment"."notification_outbox" ("intent_id") + `); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP TABLE "edr_payment"."notification_outbox"`); + await queryRunner.query(`DROP TABLE "edr_payment"."payment_webhook_event"`); + await queryRunner.query(`DROP TABLE "edr_payment"."payment_intent"`); + } +} diff --git a/apps/edr-payment-api/src/modules/health/health.controller.ts b/apps/edr-payment-api/src/modules/health/health.controller.ts new file mode 100644 index 000000000..86eee1cc0 --- /dev/null +++ b/apps/edr-payment-api/src/modules/health/health.controller.ts @@ -0,0 +1,16 @@ +import { Controller, Get } from "@nestjs/common"; +import { ApiOperation, ApiTags } from "@nestjs/swagger"; + +@ApiTags("Health") +@Controller("health") +export class HealthController { + @Get() + @ApiOperation({ summary: "Liveness probe" }) + check() { + return { + status: "ok", + service: "edr-payment-api", + timestamp: new Date().toISOString(), + }; + } +} diff --git a/apps/edr-payment-api/src/modules/health/health.module.ts b/apps/edr-payment-api/src/modules/health/health.module.ts new file mode 100644 index 000000000..40b7bdfae --- /dev/null +++ b/apps/edr-payment-api/src/modules/health/health.module.ts @@ -0,0 +1,7 @@ +import { Module } from "@nestjs/common"; +import { HealthController } from "./health.controller"; + +@Module({ + controllers: [HealthController], +}) +export class HealthModule {} diff --git a/apps/edr-payment-api/src/modules/intents/dto/initiate-payment.dto.ts b/apps/edr-payment-api/src/modules/intents/dto/initiate-payment.dto.ts new file mode 100644 index 000000000..20465fc0e --- /dev/null +++ b/apps/edr-payment-api/src/modules/intents/dto/initiate-payment.dto.ts @@ -0,0 +1,100 @@ +import { + IsEnum, + IsIn, + IsInt, + IsOptional, + IsPositive, + IsString, + Length, + MaxLength, +} from "class-validator"; +import { ApiProperty, ApiPropertyOptional } from "@nestjs/swagger"; +import { + InitiatePaymentRequest, + PaymentPlatform, + PaymentReferenceType, + PaymentService, + ProviderMethod, +} from "@edr/types"; + +/** Wire shape is the shared `InitiatePaymentRequest` contract from @edr/types. */ +export class InitiatePaymentRequestDto implements InitiatePaymentRequest { + @ApiProperty({ enum: PaymentService }) + @IsEnum(PaymentService) + service!: PaymentService; + + @ApiProperty({ enum: PaymentReferenceType }) + @IsEnum(PaymentReferenceType) + referenceType!: PaymentReferenceType; + + @ApiProperty({ + description: + "Domain order id (booking/shipment id) — already validated by the calling app", + }) + @IsString() + @Length(1, 64) + referenceId!: string; + + @ApiPropertyOptional({ + description: + "Human-readable order ref shown on provider pages; defaults to referenceId", + }) + @IsOptional() + @IsString() + @MaxLength(64) + orderRef?: string; + + @ApiProperty({ + description: + "Authoritative amount in minor units, computed server-side by the app", + }) + @IsInt() + @IsPositive() + amountMinor!: number; + + @ApiProperty({ example: "ETB" }) + @IsString() + @Length(3, 8) + currency!: string; + + @ApiProperty({ enum: ProviderMethod }) + @IsEnum(ProviderMethod) + provider!: ProviderMethod; + + @ApiPropertyOptional({ enum: ["web", "mobile"] }) + @IsOptional() + @IsIn(["web", "mobile"]) + platform?: PaymentPlatform; + + @ApiPropertyOptional({ + description: + "Payer wallet MSISDN for providers that pre-fill it (e.g. Waafi MWALLET_ACCOUNT)", + }) + @IsOptional() + @IsString() + @MaxLength(32) + payerAccount?: string; + + @ApiPropertyOptional({ + description: "Caller key to dedupe retried initiations", + }) + @IsOptional() + @IsString() + @MaxLength(128) + idempotencyKey?: string; +} + +export class IntentReferenceQueryDto { + @ApiProperty({ enum: PaymentService }) + @IsEnum(PaymentService) + service!: PaymentService; + + @ApiProperty({ enum: PaymentReferenceType }) + @IsEnum(PaymentReferenceType) + referenceType!: PaymentReferenceType; + + @ApiProperty() + @IsString() + @Length(1, 64) + referenceId!: string; +} diff --git a/apps/edr-payment-api/src/modules/intents/entities/payment-intent.entity.ts b/apps/edr-payment-api/src/modules/intents/entities/payment-intent.entity.ts new file mode 100644 index 000000000..68c519604 --- /dev/null +++ b/apps/edr-payment-api/src/modules/intents/entities/payment-intent.entity.ts @@ -0,0 +1,135 @@ +import { Column, Entity, Index } from "typeorm"; +import { BaseEntity } from "@edr/api-common"; +import { + ClientAction, + PaymentReferenceType, + PaymentService, + ProviderMethod, + ProviderPaymentStatus, +} from "@edr/types"; + +/** + * One payment attempt for one domain order — the platform-wide source of truth for payment + * state. `reference_id` is a soft reference into the owning app's schema (never a FK; see + * docs/payment-service/architecture.md §5). + * + * Enum-valued columns are stored as varchar (values mirror the shared @edr/types enums) so + * adding a provider/status never needs an ALTER TYPE migration. + */ +@Entity({ name: "payment_intent" }) +// One ACTIVE intent per domain order; FAILED/CANCELLED attempts may accumulate as audit rows. +@Index( + "uq_payment_intent_active_reference", + ["service", "referenceType", "referenceId"], + { + unique: true, + where: `status NOT IN ('FAILED','CANCELLED') AND deleted_at IS NULL`, + }, +) +@Index("idx_payment_intent_sweep", ["status", "updatedAt"]) +@Index("idx_payment_intent_idempotency", ["service", "idempotencyKey"]) +export class PaymentIntent extends BaseEntity { + /** Owning domain app — routing discriminator for notifications. */ + @Column({ name: "service", type: "varchar", length: 16 }) + service!: PaymentService; + + @Column({ name: "reference_type", type: "varchar", length: 16 }) + referenceType!: PaymentReferenceType; + + /** Domain order id (booking/shipment). Soft reference — no cross-schema FK. */ + @Column({ name: "reference_id", type: "varchar", length: 64 }) + referenceId!: string; + + /** Provider-facing reference, prefixed PSG-/FRT- so webhooks route before a DB lookup. */ + @Column({ + name: "merchant_order_id", + type: "varchar", + length: 64, + unique: true, + }) + merchantOrderId!: string; + + @Column({ name: "provider", type: "varchar", length: 16 }) + provider!: ProviderMethod; + + /** Provider-side order/session id (prepay id, HPP orderId, …). */ + @Column({ + name: "provider_order_id", + type: "varchar", + length: 128, + nullable: true, + }) + providerOrderId?: string | null; + + /** Final provider transaction id, set on terminal success. */ + @Index("idx_payment_intent_provider_txn") + @Column({ + name: "provider_txn_id", + type: "varchar", + length: 128, + nullable: true, + }) + providerTxnId?: string | null; + + /** App-asserted authoritative amount in minor units. */ + @Column({ name: "amount_minor", type: "integer" }) + amountMinor!: number; + + /** Provider-reported amount; reconciled against amount_minor (e.g. Waafi truncates decimals). */ + @Column({ name: "confirmed_amount_minor", type: "integer", nullable: true }) + confirmedAmountMinor?: number | null; + + @Column({ name: "currency", type: "varchar", length: 8 }) + currency!: string; + + /** State machine: REQUIRES_ACTION → PROCESSING → SUCCEEDED | FAILED | CANCELLED (absorbing). */ + @Column({ + name: "status", + type: "varchar", + length: 24, + default: ProviderPaymentStatus.REQUIRES_ACTION, + }) + status!: ProviderPaymentStatus; + + /** Redirect/launch payload returned to the app for the user to complete payment. */ + @Column({ name: "client_action", type: "jsonb", nullable: true }) + clientAction?: ClientAction | null; + + @Column({ name: "failure_code", type: "varchar", length: 64, nullable: true }) + failureCode?: string | null; + + @Column({ name: "failure_message", type: "text", nullable: true }) + failureMessage?: string | null; + + /** Caller-supplied initiate dedupe key (in addition to the per-reference upsert). */ + @Column({ + name: "idempotency_key", + type: "varchar", + length: 128, + nullable: true, + }) + idempotencyKey?: string | null; + + @Column({ name: "expires_at", type: "timestamptz", nullable: true }) + expiresAt?: Date | null; + + @Column({ name: "paid_at", type: "timestamptz", nullable: true }) + paidAt?: Date | null; + + /** Audit copy of the provider initiation request/response (secrets redacted upstream). */ + @Column({ name: "raw_initiation", type: "jsonb", nullable: true }) + rawInitiation?: Record | null; +} + +/** Statuses that keep the per-reference unique index "active" (block a new intent). */ +export const ACTIVE_INTENT_STATUSES = [ + ProviderPaymentStatus.REQUIRES_ACTION, + ProviderPaymentStatus.PROCESSING, + ProviderPaymentStatus.SUCCEEDED, +] as const; + +export const TERMINAL_INTENT_STATUSES = [ + ProviderPaymentStatus.SUCCEEDED, + ProviderPaymentStatus.FAILED, + ProviderPaymentStatus.CANCELLED, +] as const; diff --git a/apps/edr-payment-api/src/modules/intents/intents.controller.ts b/apps/edr-payment-api/src/modules/intents/intents.controller.ts new file mode 100644 index 000000000..858ad3bae --- /dev/null +++ b/apps/edr-payment-api/src/modules/intents/intents.controller.ts @@ -0,0 +1,69 @@ +import { + Body, + Controller, + Get, + Param, + ParseUUIDPipe, + Post, + Query, + UseGuards, +} from "@nestjs/common"; +import { ApiOperation, ApiTags } from "@nestjs/swagger"; +import { PaymentIntentSnapshot } from "@edr/types"; +import { ServiceAuthGuard } from "../../common/guards/service-auth.guard"; +import { + InitiatePaymentRequestDto, + IntentReferenceQueryDto, +} from "./dto/initiate-payment.dto"; +import { IntentsService } from "./intents.service"; + +/** + * Internal surface — called only by the domain apps (service-authenticated), never by + * browsers. Domain validation ("is this booking payable", authoritative amount) has already + * happened in the calling app. + */ +@ApiTags("Payments (internal)") +@UseGuards(ServiceAuthGuard) +@Controller("payments") +export class IntentsController { + constructor(private readonly intentsService: IntentsService) {} + + @Post("initiate") + @ApiOperation({ + summary: + "Create (or idempotently reuse) a payment intent and open a provider session", + description: + "One active intent per (service, referenceType, referenceId). Re-initiating a non-terminal intent returns the existing clientAction.", + }) + async initiate( + @Body() dto: InitiatePaymentRequestDto, + ): Promise { + return this.intentsService.initiate(dto); + } + + @Get("intents/:id") + @ApiOperation({ + summary: "Intent status by id (pull/reconcile)", + description: + "Stale non-terminal intents trigger a provider status query before returning.", + }) + async getIntent( + @Param("id", ParseUUIDPipe) id: string, + ): Promise { + return this.intentsService.getIntent(id); + } + + @Get("intents") + @ApiOperation({ + summary: "Active intent status by domain reference (pull/reconcile)", + }) + async getIntentByReference( + @Query() query: IntentReferenceQueryDto, + ): Promise { + return this.intentsService.getIntentByReference( + query.service, + query.referenceType, + query.referenceId, + ); + } +} diff --git a/apps/edr-payment-api/src/modules/intents/intents.module.ts b/apps/edr-payment-api/src/modules/intents/intents.module.ts new file mode 100644 index 000000000..9e9774ad9 --- /dev/null +++ b/apps/edr-payment-api/src/modules/intents/intents.module.ts @@ -0,0 +1,21 @@ +import { Module } from "@nestjs/common"; +import { TypeOrmModule } from "@nestjs/typeorm"; +import { ProvidersModule } from "../providers/providers.module"; +import { NotificationOutbox } from "../outbox/entities/notification-outbox.entity"; +import { PaymentIntent } from "./entities/payment-intent.entity"; +import { IntentsController } from "./intents.controller"; +import { IntentsRepository } from "./intents.repository"; +import { IntentsService } from "./intents.service"; + +@Module({ + // NotificationOutbox is registered here because terminal transitions insert outbox rows + // inside the intent-finalizing transaction (transactional outbox). + imports: [ + TypeOrmModule.forFeature([PaymentIntent, NotificationOutbox]), + ProvidersModule, + ], + controllers: [IntentsController], + providers: [IntentsService, IntentsRepository], + exports: [IntentsService, IntentsRepository], +}) +export class IntentsModule {} diff --git a/apps/edr-payment-api/src/modules/intents/intents.repository.ts b/apps/edr-payment-api/src/modules/intents/intents.repository.ts new file mode 100644 index 000000000..a17407069 --- /dev/null +++ b/apps/edr-payment-api/src/modules/intents/intents.repository.ts @@ -0,0 +1,73 @@ +import { Injectable } from "@nestjs/common"; +import { InjectRepository } from "@nestjs/typeorm"; +import { In, LessThan, Not, Repository } from "typeorm"; +import { BaseRepository } from "@edr/api-common"; +import { + PaymentReferenceType, + PaymentService, + ProviderPaymentStatus, +} from "@edr/types"; +import { PaymentIntent } from "./entities/payment-intent.entity"; + +@Injectable() +export class IntentsRepository extends BaseRepository { + constructor( + @InjectRepository(PaymentIntent) + repository: Repository, + ) { + super(repository); + } + + /** The single non-FAILED/CANCELLED intent for a domain order (matches the partial unique index). */ + async findActiveByReference( + service: PaymentService, + referenceType: PaymentReferenceType, + referenceId: string, + ): Promise { + return this.repository.findOne({ + where: { + service, + referenceType, + referenceId, + status: Not( + In([ProviderPaymentStatus.FAILED, ProviderPaymentStatus.CANCELLED]), + ), + }, + order: { createdAt: "DESC" }, + }); + } + + async findByMerchantOrderId( + merchantOrderId: string, + ): Promise { + return this.repository.findOne({ where: { merchantOrderId } }); + } + + async findByIdempotencyKey( + service: PaymentService, + idempotencyKey: string, + ): Promise { + return this.repository.findOne({ + where: { service, idempotencyKey }, + order: { createdAt: "DESC" }, + }); + } + + /** Non-terminal intents untouched since `updatedBefore` — input for the reconciliation sweep. */ + async findStale( + updatedBefore: Date, + limit: number, + ): Promise { + return this.repository.find({ + where: { + status: In([ + ProviderPaymentStatus.REQUIRES_ACTION, + ProviderPaymentStatus.PROCESSING, + ]), + updatedAt: LessThan(updatedBefore), + }, + order: { updatedAt: "ASC" }, + take: limit, + }); + } +} diff --git a/apps/edr-payment-api/src/modules/intents/intents.service.ts b/apps/edr-payment-api/src/modules/intents/intents.service.ts new file mode 100644 index 000000000..faec489f9 --- /dev/null +++ b/apps/edr-payment-api/src/modules/intents/intents.service.ts @@ -0,0 +1,342 @@ +import { + BadRequestException, + Inject, + Injectable, + Logger, + NotFoundException, +} from "@nestjs/common"; +import { DataSource, QueryFailedError } from "typeorm"; +import { createMerchantOrderId } from "@edr/payment-providers"; +import { + InitiatePaymentRequest, + MERCHANT_ORDER_PREFIX, + PaymentIntentSnapshot, + PaymentReferenceType, + PaymentService, + ProviderPaymentStatus, + ProviderStatus, +} from "@edr/types"; +import { + PAYMENT_PROVIDER_MAP, + PaymentProviderMap, +} from "../providers/providers.module"; +import { NotificationOutbox } from "../outbox/entities/notification-outbox.entity"; +import { buildOutboxRow } from "../outbox/payment-event.factory"; +import { + PaymentIntent, + TERMINAL_INTENT_STATUSES, +} from "./entities/payment-intent.entity"; +import { IntentsRepository } from "./intents.repository"; + +const PG_UNIQUE_VIOLATION = "23505"; +/** Don't hit the provider again if the intent was refreshed this recently. */ +const REFRESH_MIN_AGE_MS = 5_000; + +/** Result of a provider signal (webhook or status query) applied to the state machine. */ +export interface ProviderResultInput { + status: ProviderPaymentStatus; + providerTxnId?: string; + paidAt?: Date; + confirmedAmountMinor?: number; + failureCode?: string; + failureMessage?: string; +} + +@Injectable() +export class IntentsService { + private readonly logger = new Logger(IntentsService.name); + + constructor( + private readonly intentsRepository: IntentsRepository, + // DataSource is used only for the finalize transaction (intent update + outbox insert + // must commit atomically); routine access still goes through the custom repository. + private readonly dataSource: DataSource, + @Inject(PAYMENT_PROVIDER_MAP) + private readonly providers: PaymentProviderMap, + ) {} + + /* ------------------------------------------------------------------ initiate */ + + async initiate( + request: InitiatePaymentRequest, + ): Promise { + if (request.idempotencyKey) { + const byKey = await this.intentsRepository.findByIdempotencyKey( + request.service, + request.idempotencyKey, + ); + if (byKey) return this.toSnapshot(byKey); + } + + const existing = await this.intentsRepository.findActiveByReference( + request.service, + request.referenceType, + request.referenceId, + ); + if (existing) { + const reusable = await this.reuseOrRetire(existing); + if (reusable) return this.toSnapshot(reusable); + } + + const provider = this.providers.get(request.provider); + if (!provider) { + throw new BadRequestException( + `Unsupported payment provider: ${request.provider}`, + ); + } + + const merchantOrderId = `${MERCHANT_ORDER_PREFIX[request.service]}${createMerchantOrderId()}`; + const result = await provider.initiate({ + merchantOrderId, + orderRef: request.orderRef ?? request.referenceId, + amountMinor: request.amountMinor, + currency: request.currency, + platform: request.platform, + payerAccount: request.payerAccount, + }); + + try { + const intent = await this.intentsRepository.create({ + service: request.service, + referenceType: request.referenceType, + referenceId: request.referenceId, + merchantOrderId, + provider: request.provider, + providerOrderId: result.providerOrderId, + amountMinor: request.amountMinor, + currency: request.currency, + status: ProviderPaymentStatus.REQUIRES_ACTION, + clientAction: result.clientAction, + idempotencyKey: request.idempotencyKey ?? null, + expiresAt: result.expiresAt, + rawInitiation: result.rawInitiation, + }); + this.logger.log( + `intent ${intent.id} created: ${request.service}/${request.referenceType}/${request.referenceId} via ${request.provider} (${merchantOrderId})`, + ); + return this.toSnapshot(intent); + } catch (err) { + // Concurrent initiate for the same reference lost the partial-unique race — return the + // winner's intent. The provider session we just opened is simply abandoned. + if ( + err instanceof QueryFailedError && + (err.driverError as { code?: string })?.code === PG_UNIQUE_VIOLATION + ) { + const winner = await this.intentsRepository.findActiveByReference( + request.service, + request.referenceType, + request.referenceId, + ); + if (winner) return this.toSnapshot(winner); + } + throw err; + } + } + + /** + * Decide whether an existing active intent can be returned as-is. An expired + * REQUIRES_ACTION intent is retired (CANCELLED, no notification — nothing was paid) + * so a fresh provider session can be opened. + */ + private async reuseOrRetire( + intent: PaymentIntent, + ): Promise { + const expired = + intent.status === ProviderPaymentStatus.REQUIRES_ACTION && + intent.expiresAt != null && + intent.expiresAt.getTime() < Date.now(); + if (!expired) return intent; + + await this.intentsRepository.update(intent.id, { + status: ProviderPaymentStatus.CANCELLED, + failureCode: "EXPIRED", + failureMessage: "Provider session expired before the payer acted", + }); + return null; + } + + /* ------------------------------------------------------------------ lookups */ + + async getIntent(id: string): Promise { + const intent = await this.intentsRepository.findById(id); + if (!intent) throw new NotFoundException("PaymentIntent not found"); + return this.toSnapshot(await this.refreshIfStale(intent)); + } + + async getIntentByReference( + service: PaymentService, + referenceType: PaymentReferenceType, + referenceId: string, + ): Promise { + const intent = await this.intentsRepository.findActiveByReference( + service, + referenceType, + referenceId, + ); + if (!intent) throw new NotFoundException("PaymentIntent not found"); + return this.toSnapshot(await this.refreshIfStale(intent)); + } + + /** + * Pull-side reconciliation: when a polled intent is non-terminal and stale, ask the + * provider for the truth and run the answer through the state machine. The browser + * redirect never confirms payment — this query (or a webhook) does. + */ + private async refreshIfStale(intent: PaymentIntent): Promise { + const refreshable = + intent.status === ProviderPaymentStatus.REQUIRES_ACTION || + intent.status === ProviderPaymentStatus.PROCESSING; + const stale = intent.updatedAt.getTime() < Date.now() - REFRESH_MIN_AGE_MS; + const provider = this.providers.get(intent.provider); + if (!refreshable || !stale || !provider) return intent; + + try { + const status = await provider.queryStatus(intent.merchantOrderId); + await this.applyProviderResult( + intent.id, + this.fromProviderStatus(status), + ); + return (await this.intentsRepository.findById(intent.id)) ?? intent; + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + this.logger.warn( + `queryStatus failed for intent ${intent.id}: ${message}; returning cached`, + ); + return intent; + } + } + + fromProviderStatus(status: ProviderStatus): ProviderResultInput { + return { + status: status.status, + providerTxnId: status.providerTxnId, + failureCode: status.failureCode, + failureMessage: status.failureMessage, + }; + } + + /* ------------------------------------------------------------------ state machine */ + + /** + * Advance the intent state machine with a verified provider signal. Terminal states are + * absorbing; a terminal transition writes the notification_outbox row IN THE SAME + * TRANSACTION as the intent update (transactional outbox — architecture.md §8). + */ + async applyProviderResult( + intentId: string, + result: ProviderResultInput, + ): Promise<{ alreadyTerminal: boolean }> { + return this.dataSource.transaction(async (manager) => { + const intent = await manager + .getRepository(PaymentIntent) + .createQueryBuilder("intent") + .setLock("pessimistic_write") + .where("intent.id = :intentId", { intentId }) + .getOne(); + if (!intent) throw new NotFoundException("PaymentIntent not found"); + + if ( + (TERMINAL_INTENT_STATUSES as readonly ProviderPaymentStatus[]).includes( + intent.status, + ) + ) { + return { alreadyTerminal: true }; + } + + if (result.status === ProviderPaymentStatus.SUCCEEDED) { + const paidAt = result.paidAt ?? new Date(); + intent.status = ProviderPaymentStatus.SUCCEEDED; + intent.providerTxnId = result.providerTxnId ?? intent.providerTxnId; + intent.paidAt = paidAt; + intent.confirmedAmountMinor = + result.confirmedAmountMinor ?? intent.confirmedAmountMinor; + intent.failureCode = null; + intent.failureMessage = null; + await manager.save(intent); + await manager.getRepository(NotificationOutbox).save( + buildOutboxRow(intent, { + eventType: "payment.succeeded", + providerTxnId: intent.providerTxnId ?? undefined, + paidAt, + }), + ); + if ( + result.confirmedAmountMinor != null && + result.confirmedAmountMinor !== intent.amountMinor + ) { + this.logger.error( + `intent ${intent.id} amount mismatch: asserted=${intent.amountMinor} confirmed=${result.confirmedAmountMinor}`, + ); + } + this.logger.log( + `intent ${intent.id} SUCCEEDED (txn=${intent.providerTxnId ?? "n/a"})`, + ); + return { alreadyTerminal: false }; + } + + if ( + result.status === ProviderPaymentStatus.FAILED || + result.status === ProviderPaymentStatus.CANCELLED + ) { + intent.status = result.status; + intent.providerTxnId = result.providerTxnId ?? intent.providerTxnId; + intent.failureCode = result.failureCode ?? null; + intent.failureMessage = result.failureMessage ?? null; + await manager.save(intent); + await manager.getRepository(NotificationOutbox).save( + buildOutboxRow(intent, { + eventType: "payment.failed", + failureCode: result.failureCode, + failureMessage: result.failureMessage, + }), + ); + this.logger.log( + `intent ${intent.id} ${result.status} (${result.failureCode ?? "n/a"})`, + ); + return { alreadyTerminal: false }; + } + + // Non-terminal: REQUIRES_ACTION may move to PROCESSING; never the reverse. + if ( + result.status === ProviderPaymentStatus.PROCESSING && + intent.status === ProviderPaymentStatus.REQUIRES_ACTION + ) { + intent.status = ProviderPaymentStatus.PROCESSING; + } + intent.providerTxnId = result.providerTxnId ?? intent.providerTxnId; + await manager.save(intent); + return { alreadyTerminal: false }; + }); + } + + /** Expire an abandoned intent (reconciliation sweep) — CANCELLED + payment.failed event. */ + async expireIntent(intentId: string): Promise { + await this.applyProviderResult(intentId, { + status: ProviderPaymentStatus.CANCELLED, + failureCode: "EXPIRED", + failureMessage: "Payment session expired before completion", + }); + } + + /* ------------------------------------------------------------------ mapping */ + + toSnapshot(intent: PaymentIntent): PaymentIntentSnapshot { + return { + intentId: intent.id, + service: intent.service, + referenceType: intent.referenceType, + referenceId: intent.referenceId, + merchantOrderId: intent.merchantOrderId, + provider: intent.provider, + status: intent.status, + amountMinor: intent.amountMinor, + currency: intent.currency, + clientAction: intent.clientAction ?? undefined, + providerTxnId: intent.providerTxnId ?? undefined, + paidAt: intent.paidAt?.toISOString(), + failureCode: intent.failureCode ?? undefined, + failureMessage: intent.failureMessage ?? undefined, + expiresAt: intent.expiresAt?.toISOString(), + }; + } +} diff --git a/apps/edr-payment-api/src/modules/outbox/entities/notification-outbox.entity.ts b/apps/edr-payment-api/src/modules/outbox/entities/notification-outbox.entity.ts new file mode 100644 index 000000000..55b34e968 --- /dev/null +++ b/apps/edr-payment-api/src/modules/outbox/entities/notification-outbox.entity.ts @@ -0,0 +1,55 @@ +import { Column, Entity, Index } from "typeorm"; +import { BaseEntity } from "@edr/api-common"; +import { + PaymentEvent, + PaymentEventType, + PaymentReferenceType, + PaymentService, +} from "@edr/types"; + +export type OutboxStatus = "PENDING" | "SENT" | "FAILED"; + +/** + * Transactional outbox: a row is inserted in the SAME transaction that finalizes an intent, + * so "payment succeeded" and "a notification is owed" commit or roll back together. The relay + * drains PENDING rows and retries until acked (at-least-once delivery; consumers are idempotent). + */ +@Entity({ name: "notification_outbox" }) +@Index("idx_notification_outbox_relay", ["status", "nextRetryAt"]) +export class NotificationOutbox extends BaseEntity { + @Column({ name: "event_type", type: "varchar", length: 32 }) + eventType!: PaymentEventType; + + /** Routing discriminator — which app's mark-paid endpoint the relay delivers to. */ + @Column({ name: "service", type: "varchar", length: 16 }) + service!: PaymentService; + + @Index("idx_notification_outbox_intent") + @Column({ name: "intent_id", type: "uuid" }) + intentId!: string; + + @Column({ name: "reference_type", type: "varchar", length: 16 }) + referenceType!: PaymentReferenceType; + + @Column({ name: "reference_id", type: "varchar", length: 64 }) + referenceId!: string; + + /** The full versioned event envelope delivered verbatim to the consumer. */ + @Column({ name: "payload", type: "jsonb" }) + payload!: PaymentEvent; + + @Column({ name: "status", type: "varchar", length: 16, default: "PENDING" }) + status!: OutboxStatus; + + @Column({ name: "attempts", type: "integer", default: 0 }) + attempts!: number; + + @Column({ name: "next_retry_at", type: "timestamptz", nullable: true }) + nextRetryAt?: Date | null; + + @Column({ name: "last_error", type: "text", nullable: true }) + lastError?: string | null; + + @Column({ name: "sent_at", type: "timestamptz", nullable: true }) + sentAt?: Date | null; +} diff --git a/apps/edr-payment-api/src/modules/outbox/outbox-relay.service.ts b/apps/edr-payment-api/src/modules/outbox/outbox-relay.service.ts new file mode 100644 index 000000000..f1d8076a5 --- /dev/null +++ b/apps/edr-payment-api/src/modules/outbox/outbox-relay.service.ts @@ -0,0 +1,119 @@ +import { + Inject, + Injectable, + Logger, + OnModuleDestroy, + OnModuleInit, +} from "@nestjs/common"; +import { ConfigService } from "@nestjs/config"; +import { SchedulerRegistry } from "@nestjs/schedule"; +import { NotificationOutbox } from "./entities/notification-outbox.entity"; +import { OutboxRepository } from "./outbox.repository"; +import { + PAYMENT_EVENT_PUBLISHER, + PaymentEventPublisher, +} from "./publisher/payment-event-publisher"; + +const RELAY_INTERVAL_NAME = "outbox-relay"; +/** Retry backoff: base doubles per attempt, capped. */ +const BACKOFF_BASE_MS = 10_000; +const BACKOFF_CAP_MS = 10 * 60_000; + +/** + * Drains the transactional outbox: PENDING rows are published (HTTP now, RabbitMQ later), + * marked SENT on ack, retried with exponential backoff on failure, and flagged FAILED after + * OUTBOX_MAX_ATTEMPTS (an alertable condition — delivery is at-least-once, never dropped + * silently). A crash between commit and publish only delays delivery. + */ +@Injectable() +export class OutboxRelayService implements OnModuleInit, OnModuleDestroy { + private readonly logger = new Logger(OutboxRelayService.name); + private readonly intervalMs: number; + private readonly maxAttempts: number; + private readonly batchSize: number; + private draining = false; + + constructor( + config: ConfigService, + private readonly outboxRepository: OutboxRepository, + private readonly schedulerRegistry: SchedulerRegistry, + @Inject(PAYMENT_EVENT_PUBLISHER) + private readonly publisher: PaymentEventPublisher, + ) { + this.intervalMs = config.get("notifier.relayIntervalMs") ?? 5_000; + this.maxAttempts = config.get("notifier.maxAttempts") ?? 10; + this.batchSize = config.get("notifier.relayBatchSize") ?? 20; + } + + onModuleInit(): void { + const interval = setInterval(() => void this.drain(), this.intervalMs); + this.schedulerRegistry.addInterval(RELAY_INTERVAL_NAME, interval); + } + + onModuleDestroy(): void { + if (this.schedulerRegistry.doesExist("interval", RELAY_INTERVAL_NAME)) { + this.schedulerRegistry.deleteInterval(RELAY_INTERVAL_NAME); + } + } + + /** One relay pass; re-entrant ticks are skipped so slow deliveries don't overlap. */ + async drain(): Promise { + if (this.draining) return; + this.draining = true; + try { + const due = await this.outboxRepository.findDue(this.batchSize); + for (const row of due) { + await this.deliver(row); + } + } catch (err) { + this.logger.error( + `relay pass failed: ${err instanceof Error ? err.message : String(err)}`, + ); + } finally { + this.draining = false; + } + } + + private async deliver(row: NotificationOutbox): Promise { + try { + await this.publisher.publish(row.payload); + await this.outboxRepository.markSent(row.id); + } catch (err) { + const message = this.describeError(err); + const attempts = row.attempts + 1; + const exhausted = attempts >= this.maxAttempts; + const backoffMs = Math.min( + BACKOFF_BASE_MS * 2 ** row.attempts, + BACKOFF_CAP_MS, + ); + await this.outboxRepository.markAttemptFailed( + row, + message, + exhausted ? null : new Date(Date.now() + backoffMs), + exhausted, + ); + if (exhausted) { + // ALERT: a paid order may not be confirmed in the owning app — needs operator action. + this.logger.error( + `outbox ${row.id} (${row.eventType} intent=${row.intentId}) FAILED after ${attempts} attempts: ${message}`, + ); + } else { + this.logger.warn( + `outbox ${row.id} delivery attempt ${attempts} failed (retry in ${backoffMs}ms): ${message}`, + ); + } + } + } + + /** Connection failures surface as AggregateError with an empty message — dig out the code. */ + private describeError(err: unknown): string { + if (err instanceof Error) { + if (err.message) return err.message; + const code = (err as { code?: string }).code; + if (code) return code; + const inner = (err as { errors?: unknown[] }).errors?.[0]; + if (inner instanceof Error && inner.message) return inner.message; + } + return String(err); + } +} diff --git a/apps/edr-payment-api/src/modules/outbox/outbox.module.ts b/apps/edr-payment-api/src/modules/outbox/outbox.module.ts new file mode 100644 index 000000000..85f464afe --- /dev/null +++ b/apps/edr-payment-api/src/modules/outbox/outbox.module.ts @@ -0,0 +1,20 @@ +import { Module } from "@nestjs/common"; +import { HttpModule } from "@nestjs/axios"; +import { TypeOrmModule } from "@nestjs/typeorm"; +import { NotificationOutbox } from "./entities/notification-outbox.entity"; +import { OutboxRelayService } from "./outbox-relay.service"; +import { OutboxRepository } from "./outbox.repository"; +import { HttpPaymentEventPublisher } from "./publisher/http-payment-event-publisher"; +import { PAYMENT_EVENT_PUBLISHER } from "./publisher/payment-event-publisher"; + +@Module({ + imports: [TypeOrmModule.forFeature([NotificationOutbox]), HttpModule], + providers: [ + OutboxRepository, + OutboxRelayService, + // Swap to RabbitPaymentEventPublisher here when the broker lands — nothing else changes. + { provide: PAYMENT_EVENT_PUBLISHER, useClass: HttpPaymentEventPublisher }, + ], + exports: [OutboxRepository], +}) +export class OutboxModule {} diff --git a/apps/edr-payment-api/src/modules/outbox/outbox.repository.ts b/apps/edr-payment-api/src/modules/outbox/outbox.repository.ts new file mode 100644 index 000000000..20000d9c6 --- /dev/null +++ b/apps/edr-payment-api/src/modules/outbox/outbox.repository.ts @@ -0,0 +1,58 @@ +import { Injectable } from "@nestjs/common"; +import { InjectRepository } from "@nestjs/typeorm"; +import { Repository } from "typeorm"; +import { BaseRepository } from "@edr/api-common"; +import { NotificationOutbox } from "./entities/notification-outbox.entity"; + +@Injectable() +export class OutboxRepository extends BaseRepository { + constructor( + @InjectRepository(NotificationOutbox) + repository: Repository, + ) { + super(repository); + } + + /** + * PENDING rows whose retry time has come, oldest first. The relay runs as a single + * non-overlapping loop per instance; with multiple service instances this should move to a + * SELECT … FOR UPDATE SKIP LOCKED claim. + */ + async findDue(limit: number): Promise { + return this.repository + .createQueryBuilder("outbox") + .where(`outbox.status = 'PENDING'`) + .andWhere( + "(outbox.next_retry_at IS NULL OR outbox.next_retry_at <= now())", + ) + .orderBy("outbox.created_at", "ASC") + .take(limit) + .getMany(); + } + + async markSent(id: string): Promise { + await this.update(id, { + status: "SENT", + sentAt: new Date(), + lastError: null, + }); + } + + async markAttemptFailed( + row: NotificationOutbox, + error: string, + nextRetryAt: Date | null, + exhausted: boolean, + ): Promise { + await this.update(row.id, { + attempts: row.attempts + 1, + lastError: error, + nextRetryAt, + status: exhausted ? "FAILED" : "PENDING", + }); + } + + async countBacklog(): Promise { + return this.repository.count({ where: { status: "PENDING" } }); + } +} diff --git a/apps/edr-payment-api/src/modules/outbox/payment-event.factory.ts b/apps/edr-payment-api/src/modules/outbox/payment-event.factory.ts new file mode 100644 index 000000000..85b67c444 --- /dev/null +++ b/apps/edr-payment-api/src/modules/outbox/payment-event.factory.ts @@ -0,0 +1,66 @@ +import { randomUUID } from "node:crypto"; +import { + PaymentEvent, + PaymentFailedEvent, + PaymentSucceededEvent, +} from "@edr/types"; +import { PaymentIntent } from "../intents/entities/payment-intent.entity"; +import { NotificationOutbox } from "./entities/notification-outbox.entity"; + +/** + * Build a ready-to-insert outbox row for a terminal intent. Pure (no DI) so the intents + * state machine can insert it inside its own DB transaction without a module cycle. + * The row id is generated here because the event envelope embeds it as `eventId`. + */ +export function buildOutboxRow( + intent: PaymentIntent, + terminal: + | { eventType: "payment.succeeded"; providerTxnId?: string; paidAt: Date } + | { + eventType: "payment.failed"; + failureCode?: string; + failureMessage?: string; + }, +): Partial { + const id = randomUUID(); + const base = { + version: 1 as const, + eventId: id, + occurredAt: new Date().toISOString(), + service: intent.service, + intentId: intent.id, + referenceType: intent.referenceType, + referenceId: intent.referenceId, + merchantOrderId: intent.merchantOrderId, + provider: intent.provider, + amountMinor: intent.amountMinor, + currency: intent.currency, + }; + + const event: PaymentEvent = + terminal.eventType === "payment.succeeded" + ? ({ + ...base, + eventType: "payment.succeeded", + providerTxnId: terminal.providerTxnId, + paidAt: terminal.paidAt.toISOString(), + } satisfies PaymentSucceededEvent) + : ({ + ...base, + eventType: "payment.failed", + failureCode: terminal.failureCode, + failureMessage: terminal.failureMessage, + } satisfies PaymentFailedEvent); + + return { + id, + eventType: event.eventType, + service: intent.service, + intentId: intent.id, + referenceType: intent.referenceType, + referenceId: intent.referenceId, + payload: event, + status: "PENDING", + attempts: 0, + }; +} diff --git a/apps/edr-payment-api/src/modules/outbox/publisher/http-payment-event-publisher.ts b/apps/edr-payment-api/src/modules/outbox/publisher/http-payment-event-publisher.ts new file mode 100644 index 000000000..867119d09 --- /dev/null +++ b/apps/edr-payment-api/src/modules/outbox/publisher/http-payment-event-publisher.ts @@ -0,0 +1,53 @@ +import { Injectable, Logger } from "@nestjs/common"; +import { ConfigService } from "@nestjs/config"; +import { HttpService } from "@nestjs/axios"; +import { firstValueFrom } from "rxjs"; +import { PaymentEvent, PaymentService } from "@edr/types"; +import { PaymentEventPublisher } from "./payment-event-publisher"; + +/** + * Delivers events by POSTing to the owning app's idempotent mark-paid endpoint, routed by + * the `service` discriminator. Authenticated with the shared service token (the same secret + * the apps use to call /payments/initiate). + */ +@Injectable() +export class HttpPaymentEventPublisher implements PaymentEventPublisher { + private readonly logger = new Logger(HttpPaymentEventPublisher.name); + private readonly routes: Record; + private readonly timeoutMs: number; + private readonly serviceToken: string; + + constructor( + config: ConfigService, + private readonly http: HttpService, + ) { + this.routes = { + [PaymentService.PASSENGER]: + config.get("notifier.passengerUrl") ?? "", + [PaymentService.FREIGHT]: config.get("notifier.freightUrl") ?? "", + }; + this.timeoutMs = config.get("notifier.httpTimeoutMs") ?? 10_000; + this.serviceToken = config.get("app.serviceAuthToken") ?? ""; + } + + async publish(event: PaymentEvent): Promise { + const url = this.routes[event.service]; + if (!url) { + throw new Error( + `No mark-paid URL configured for service ${event.service}`, + ); + } + + const response = await firstValueFrom( + this.http.post(url, event, { + timeout: this.timeoutMs, + headers: this.serviceToken + ? { "x-service-token": this.serviceToken } + : {}, + }), + ); + this.logger.log( + `delivered ${event.eventType} (${event.eventId}) to ${event.service} — HTTP ${response.status}`, + ); + } +} diff --git a/apps/edr-payment-api/src/modules/outbox/publisher/payment-event-publisher.ts b/apps/edr-payment-api/src/modules/outbox/publisher/payment-event-publisher.ts new file mode 100644 index 000000000..7586c1509 --- /dev/null +++ b/apps/edr-payment-api/src/modules/outbox/publisher/payment-event-publisher.ts @@ -0,0 +1,13 @@ +import { PaymentEvent } from "@edr/types"; + +/** + * Publisher port (architecture.md §12): how a payment event leaves this service. + * HTTP implementation now; a RabbitMQ implementation later is a DI swap only — the outbox + * and relay stay exactly as they are. + */ +export interface PaymentEventPublisher { + /** Deliver one event; throw on failure so the relay can retry with backoff. */ + publish(event: PaymentEvent): Promise; +} + +export const PAYMENT_EVENT_PUBLISHER = Symbol("PAYMENT_EVENT_PUBLISHER"); diff --git a/apps/edr-payment-api/src/modules/providers/providers.module.ts b/apps/edr-payment-api/src/modules/providers/providers.module.ts new file mode 100644 index 000000000..af2983ee4 --- /dev/null +++ b/apps/edr-payment-api/src/modules/providers/providers.module.ts @@ -0,0 +1,46 @@ +import { Module } from "@nestjs/common"; +import { HttpModule } from "@nestjs/axios"; +import { + CardProvider, + CbeBirrProvider, + DMoneyProvider, + EBirrProvider, + PaymentProvider, + TelebirrProvider, + WaafiProvider, +} from "@edr/payment-providers"; +import { ProviderMethod } from "@edr/types"; + +/** Injection token for the Map used to select a gateway. */ +export const PAYMENT_PROVIDER_MAP = Symbol("PAYMENT_PROVIDER_MAP"); + +export type PaymentProviderMap = Map; + +const providerClasses = [ + TelebirrProvider, + CbeBirrProvider, + EBirrProvider, + CardProvider, + WaafiProvider, + DMoneyProvider, +]; + +/** + * Thin DI wiring around @edr/payment-providers — the exact provider set the passenger app + * used to construct, relocated here. After cutover this service is the only consumer of the + * provider SDK and of the provider secrets (config/{waafi,telebirr,…}.config.ts). + */ +@Module({ + imports: [HttpModule.register({ timeout: 10_000 })], + providers: [ + ...providerClasses, + { + provide: PAYMENT_PROVIDER_MAP, + useFactory: (...providers: PaymentProvider[]): PaymentProviderMap => + new Map(providers.map((provider) => [provider.method, provider])), + inject: providerClasses, + }, + ], + exports: [PAYMENT_PROVIDER_MAP, ...providerClasses], +}) +export class ProvidersModule {} diff --git a/apps/edr-payment-api/src/modules/reconciliation/reconciliation.module.ts b/apps/edr-payment-api/src/modules/reconciliation/reconciliation.module.ts new file mode 100644 index 000000000..b424dcfec --- /dev/null +++ b/apps/edr-payment-api/src/modules/reconciliation/reconciliation.module.ts @@ -0,0 +1,10 @@ +import { Module } from "@nestjs/common"; +import { IntentsModule } from "../intents/intents.module"; +import { ProvidersModule } from "../providers/providers.module"; +import { ReconciliationService } from "./reconciliation.service"; + +@Module({ + imports: [IntentsModule, ProvidersModule], + providers: [ReconciliationService], +}) +export class ReconciliationModule {} diff --git a/apps/edr-payment-api/src/modules/reconciliation/reconciliation.service.ts b/apps/edr-payment-api/src/modules/reconciliation/reconciliation.service.ts new file mode 100644 index 000000000..5216798f2 --- /dev/null +++ b/apps/edr-payment-api/src/modules/reconciliation/reconciliation.service.ts @@ -0,0 +1,114 @@ +import { + Inject, + Injectable, + Logger, + OnModuleDestroy, + OnModuleInit, +} from "@nestjs/common"; +import { ConfigService } from "@nestjs/config"; +import { SchedulerRegistry } from "@nestjs/schedule"; +import { ProviderPaymentStatus } from "@edr/types"; +import { + PAYMENT_PROVIDER_MAP, + PaymentProviderMap, +} from "../providers/providers.module"; +import { PaymentIntent } from "../intents/entities/payment-intent.entity"; +import { IntentsRepository } from "../intents/intents.repository"; +import { IntentsService } from "../intents/intents.service"; + +const SWEEP_INTERVAL_NAME = "reconciliation-sweep"; + +/** + * Safety net (architecture.md §7.4): webhooks get lost, users abandon hosted pages. The sweep + * queries the provider for stale non-terminal intents and feeds the answer through the same + * state machine the webhooks use; intents whose provider session expired are CANCELLED. + */ +@Injectable() +export class ReconciliationService implements OnModuleInit, OnModuleDestroy { + private readonly logger = new Logger(ReconciliationService.name); + private readonly intervalMs: number; + private readonly staleAfterMs: number; + private readonly batchSize: number; + private sweeping = false; + + constructor( + config: ConfigService, + private readonly intentsRepository: IntentsRepository, + private readonly intentsService: IntentsService, + private readonly schedulerRegistry: SchedulerRegistry, + @Inject(PAYMENT_PROVIDER_MAP) + private readonly providers: PaymentProviderMap, + ) { + this.intervalMs = + config.get("app.reconciliation.sweepIntervalMs") ?? 60_000; + this.staleAfterMs = + config.get("app.reconciliation.staleAfterMs") ?? 60_000; + this.batchSize = config.get("app.reconciliation.batchSize") ?? 20; + } + + onModuleInit(): void { + const interval = setInterval(() => void this.sweep(), this.intervalMs); + this.schedulerRegistry.addInterval(SWEEP_INTERVAL_NAME, interval); + } + + onModuleDestroy(): void { + if (this.schedulerRegistry.doesExist("interval", SWEEP_INTERVAL_NAME)) { + this.schedulerRegistry.deleteInterval(SWEEP_INTERVAL_NAME); + } + } + + async sweep(): Promise { + if (this.sweeping) return; + this.sweeping = true; + try { + const cutoff = new Date(Date.now() - this.staleAfterMs); + const stale = await this.intentsRepository.findStale( + cutoff, + this.batchSize, + ); + for (const intent of stale) { + await this.reconcileIntent(intent); + } + } catch (err) { + this.logger.error( + `sweep failed: ${err instanceof Error ? err.message : String(err)}`, + ); + } finally { + this.sweeping = false; + } + } + + private async reconcileIntent(intent: PaymentIntent): Promise { + try { + const provider = this.providers.get(intent.provider); + if (provider) { + const status = await provider.queryStatus(intent.merchantOrderId); + const result = this.intentsService.fromProviderStatus(status); + if (result.status !== intent.status || result.providerTxnId) { + await this.intentsService.applyProviderResult(intent.id, result); + } + if ( + result.status === ProviderPaymentStatus.SUCCEEDED || + result.status === ProviderPaymentStatus.FAILED || + result.status === ProviderPaymentStatus.CANCELLED + ) { + this.logger.log(`reconciled intent ${intent.id} → ${result.status}`); + return; + } + } + + // Provider still says pending (or is unknown): expire only once the session is dead. + if (intent.expiresAt && intent.expiresAt.getTime() < Date.now()) { + await this.intentsService.expireIntent(intent.id); + this.logger.log( + `expired abandoned intent ${intent.id} (${intent.merchantOrderId})`, + ); + } + } catch (err) { + // Per-intent failures must not stall the sweep; the row stays stale and is retried. + this.logger.warn( + `reconcile failed for intent ${intent.id}: ${err instanceof Error ? err.message : String(err)}`, + ); + } + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/entities/payment-webhook-event.entity.ts b/apps/edr-payment-api/src/modules/webhooks/entities/payment-webhook-event.entity.ts new file mode 100644 index 000000000..2e104bf17 --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/entities/payment-webhook-event.entity.ts @@ -0,0 +1,57 @@ +import { Column, Entity, Index } from "typeorm"; +import { BaseEntity } from "@edr/api-common"; +import { ProviderMethod } from "@edr/types"; + +/** + * Idempotency + audit record for every inbound provider webhook. The unique + * (provider, external_event_id) pair is the dedupe key: a duplicate insert hits the unique + * violation and the handler short-circuits with a 200 ack. + */ +@Entity({ name: "payment_webhook_event" }) +@Index("uq_payment_webhook_event_external", ["provider", "externalEventId"], { + unique: true, +}) +export class PaymentWebhookEvent extends BaseEntity { + @Column({ name: "provider", type: "varchar", length: 16 }) + provider!: ProviderMethod; + + /** Provider event id when given (e.g. Waafi X-Webhook-Event-Id), else derived from the payload. */ + @Column({ name: "external_event_id", type: "varchar", length: 191 }) + externalEventId!: string; + + @Column({ + name: "merchant_order_id", + type: "varchar", + length: 64, + nullable: true, + }) + merchantOrderId?: string | null; + + @Column({ + name: "provider_txn_id", + type: "varchar", + length: 128, + nullable: true, + }) + providerTxnId?: string | null; + + @Column({ name: "signature_valid", type: "boolean", default: false }) + signatureValid!: boolean; + + /** Raw provider status string as sent (pre-mapping). */ + @Column({ name: "status", type: "varchar", length: 64, nullable: true }) + status?: string | null; + + /** Full webhook body — hostile input, stored verbatim for audit/replay analysis. */ + @Column({ name: "payload", type: "jsonb" }) + payload!: Record; + + @Column({ name: "received_at", type: "timestamptz", default: () => "now()" }) + receivedAt!: Date; + + @Column({ name: "processed_at", type: "timestamptz", nullable: true }) + processedAt?: Date | null; + + @Column({ name: "processing_error", type: "text", nullable: true }) + processingError?: string | null; +} diff --git a/apps/edr-payment-api/src/modules/webhooks/handlers/card-webhook.service.ts b/apps/edr-payment-api/src/modules/webhooks/handlers/card-webhook.service.ts new file mode 100644 index 000000000..8f9c9a9ab --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/handlers/card-webhook.service.ts @@ -0,0 +1,35 @@ +import { Injectable } from "@nestjs/common"; +import { CardProvider, CardWebhookPayload } from "@edr/payment-providers"; +import { WebhookProcessorService } from "../webhook-processor.service"; + +@Injectable() +export class CardWebhookService { + constructor( + private readonly provider: CardProvider, + private readonly processor: WebhookProcessorService, + ) {} + + async handle(payload: CardWebhookPayload, signature: string): Promise { + const signatureValid = this.provider.verifyWebhookSignature( + payload as unknown as Record, + signature, + ); + const object = payload.data.object; + const mapped = this.provider.mapWebhookStatus(object.status); + + await this.processor.process({ + provider: this.provider.method, + externalEventId: `${payload.id}_${payload.type}`, + merchantOrderId: object.metadata.merchantOrderId, + providerTxnId: object.transaction_id, + signatureValid, + rawStatus: object.status, + payload: payload as unknown as Record, + result: { + status: mapped, + providerTxnId: object.transaction_id, + failureCode: object.status, + }, + }); + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/handlers/cbe-birr-webhook.service.ts b/apps/edr-payment-api/src/modules/webhooks/handlers/cbe-birr-webhook.service.ts new file mode 100644 index 000000000..6621aa3b2 --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/handlers/cbe-birr-webhook.service.ts @@ -0,0 +1,30 @@ +import { Injectable } from "@nestjs/common"; +import { CbeBirrProvider, CbeBirrWebhookPayload } from "@edr/payment-providers"; +import { WebhookProcessorService } from "../webhook-processor.service"; + +@Injectable() +export class CbeBirrWebhookService { + constructor( + private readonly provider: CbeBirrProvider, + private readonly processor: WebhookProcessorService, + ) {} + + async handle(payload: CbeBirrWebhookPayload): Promise { + const signatureValid = this.provider.verifyWebhookSignature( + payload as unknown as Record, + ); + const mapped = this.provider.mapWebhookStatus(payload.status); + const providerTxnId = payload.transactionId ?? payload.orderId; + + await this.processor.process({ + provider: this.provider.method, + externalEventId: `${payload.orderId}_${payload.status}`, + merchantOrderId: payload.merchantOrderId, + providerTxnId, + signatureValid, + rawStatus: payload.status, + payload: payload as unknown as Record, + result: { status: mapped, providerTxnId, failureCode: payload.status }, + }); + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/handlers/dmoney-webhook.service.ts b/apps/edr-payment-api/src/modules/webhooks/handlers/dmoney-webhook.service.ts new file mode 100644 index 000000000..51937dee4 --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/handlers/dmoney-webhook.service.ts @@ -0,0 +1,34 @@ +import { Injectable } from "@nestjs/common"; +import { DMoneyProvider, DMoneyWebhookPayload } from "@edr/payment-providers"; +import { WebhookProcessorService } from "../webhook-processor.service"; + +@Injectable() +export class DMoneyWebhookService { + constructor( + private readonly provider: DMoneyProvider, + private readonly processor: WebhookProcessorService, + ) {} + + async handle(payload: DMoneyWebhookPayload): Promise { + const signatureValid = this.provider.verifyWebhookSignature( + payload as unknown as Record, + ); + const mapped = this.provider.mapWebhookStatus(payload.status); + + await this.processor.process({ + provider: this.provider.method, + externalEventId: `${payload.orderId}_${payload.status}`, + merchantOrderId: payload.merchantOrderId, + providerTxnId: payload.transactionId, + signatureValid, + rawStatus: payload.status, + payload: payload as unknown as Record, + result: { + status: mapped, + providerTxnId: payload.transactionId, + paidAt: payload.paidAt ? new Date(payload.paidAt) : undefined, + failureCode: payload.status, + }, + }); + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/handlers/ebirr-webhook.service.ts b/apps/edr-payment-api/src/modules/webhooks/handlers/ebirr-webhook.service.ts new file mode 100644 index 000000000..3f8b923cd --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/handlers/ebirr-webhook.service.ts @@ -0,0 +1,33 @@ +import { Injectable } from "@nestjs/common"; +import { EBirrProvider, EBirrWebhookPayload } from "@edr/payment-providers"; +import { WebhookProcessorService } from "../webhook-processor.service"; + +@Injectable() +export class EBirrWebhookService { + constructor( + private readonly provider: EBirrProvider, + private readonly processor: WebhookProcessorService, + ) {} + + async handle(payload: EBirrWebhookPayload): Promise { + const signatureValid = this.provider.verifyWebhookSignature( + payload as unknown as Record, + ); + const mapped = this.provider.mapWebhookStatus(payload.tradeStatus); + + await this.processor.process({ + provider: this.provider.method, + externalEventId: `${payload.orderNo}_${payload.tradeStatus}_${payload.timestamp}`, + merchantOrderId: payload.orderNo, + providerTxnId: payload.tradeNo, + signatureValid, + rawStatus: payload.tradeStatus, + payload: payload as unknown as Record, + result: { + status: mapped, + providerTxnId: payload.tradeNo, + failureCode: payload.tradeStatus, + }, + }); + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/handlers/telebirr-webhook.service.ts b/apps/edr-payment-api/src/modules/webhooks/handlers/telebirr-webhook.service.ts new file mode 100644 index 000000000..e207e1d21 --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/handlers/telebirr-webhook.service.ts @@ -0,0 +1,46 @@ +import { Injectable } from "@nestjs/common"; +import { + TelebirrProvider, + TelebirrWebhookPayload, +} from "@edr/payment-providers"; +import { WebhookProcessorService } from "../webhook-processor.service"; + +@Injectable() +export class TelebirrWebhookService { + constructor( + private readonly provider: TelebirrProvider, + private readonly processor: WebhookProcessorService, + ) {} + + async handle(payload: TelebirrWebhookPayload): Promise { + // TODO: re-enable Telebirr public-key signature verification — skipped for now + // (carried over from the passenger handler; see telebirr.provider verifyWebhookSignature). + const signatureValid = true; + + const mapped = this.provider.mapWebhookTradeStatus(payload.trade_status); + const providerTxnId = payload.trans_id ?? payload.payment_order_id; + + await this.processor.process({ + provider: this.provider.method, + externalEventId: `${payload.payment_order_id}_${payload.trade_status}`, + merchantOrderId: payload.merch_order_id, + providerTxnId, + signatureValid, + rawStatus: payload.trade_status, + payload: payload as unknown as Record, + result: { + status: mapped, + providerTxnId, + paidAt: this.parseEpochSeconds(payload.trans_end_time), + failureCode: payload.trade_status, + }, + }); + } + + private parseEpochSeconds(raw: string | undefined): Date | undefined { + if (!raw) return undefined; + const n = parseInt(raw, 10); + if (Number.isNaN(n)) return undefined; + return new Date(n * 1000); + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/handlers/waafi-webhook.service.ts b/apps/edr-payment-api/src/modules/webhooks/handlers/waafi-webhook.service.ts new file mode 100644 index 000000000..232957f79 --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/handlers/waafi-webhook.service.ts @@ -0,0 +1,82 @@ +import { Injectable, Logger } from "@nestjs/common"; +import { + WaafiProvider, + WaafiWebhookHeaders, + WaafiWebhookPayload, +} from "@edr/payment-providers"; +import { WebhookProcessorService } from "../webhook-processor.service"; + +/** Reject webhooks whose timestamp is older than this (replay protection). */ +const WAAFI_REPLAY_WINDOW_SECONDS = 300; + +@Injectable() +export class WaafiWebhookService { + private readonly logger = new Logger(WaafiWebhookService.name); + + constructor( + private readonly provider: WaafiProvider, + private readonly processor: WebhookProcessorService, + ) {} + + async handle( + payload: WaafiWebhookPayload, + rawBody: string, + headers: WaafiWebhookHeaders, + ): Promise { + // Unsigned validation ping sent on registration — acknowledge without verifying or persisting. + if (payload.event === "webhook.test") { + this.logger.log("Waafi webhook.test ping received"); + return; + } + + const { payment } = payload; + const eventId = headers["x-webhook-event-id"]; + const timestamp = headers["x-webhook-timestamp"]; + const signature = headers["x-webhook-signature"]; + + const signatureValid = + this.isFresh(timestamp) && + this.provider.verifyWebhookSignature( + rawBody, + signature, + timestamp, + eventId, + ); + + const mapped = this.provider.mapWebhookStatus(payment.status); + + await this.processor.process({ + provider: this.provider.method, + // X-Webhook-Event-Id is unique per event; fall back to a derived id if absent. + externalEventId: eventId ?? `${payment.transaction_id}_${payment.status}`, + merchantOrderId: payment.reference_id, + providerTxnId: payment.transaction_id, + signatureValid, + rawStatus: payment.status, + payload: payload as unknown as Record, + result: { + status: mapped, + providerTxnId: payment.transaction_id, + paidAt: this.parseDate(payment.date), + failureCode: payment.status, + failureMessage: payment.description, + }, + }); + } + + /** True when the webhook timestamp (unix seconds) is within the replay window. */ + private isFresh(timestamp: string | undefined): boolean { + if (!timestamp) return false; + const ts = parseInt(timestamp, 10); + if (Number.isNaN(ts)) return false; + const now = Math.floor(Date.now() / 1000); + return Math.abs(now - ts) <= WAAFI_REPLAY_WINDOW_SECONDS; + } + + /** Parse Waafi's "YYYY-MM-DD HH:mm:ss" payment date; undefined when unparseable. */ + private parseDate(raw: string | undefined): Date | undefined { + if (!raw) return undefined; + const d = new Date(raw); + return Number.isNaN(d.getTime()) ? undefined : d; + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/webhook-events.repository.ts b/apps/edr-payment-api/src/modules/webhooks/webhook-events.repository.ts new file mode 100644 index 000000000..7d4aeff95 --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/webhook-events.repository.ts @@ -0,0 +1,44 @@ +import { Injectable } from "@nestjs/common"; +import { InjectRepository } from "@nestjs/typeorm"; +import { QueryFailedError, Repository } from "typeorm"; +import { BaseRepository } from "@edr/api-common"; +import { PaymentWebhookEvent } from "./entities/payment-webhook-event.entity"; + +const PG_UNIQUE_VIOLATION = "23505"; + +@Injectable() +export class WebhookEventsRepository extends BaseRepository { + constructor( + @InjectRepository(PaymentWebhookEvent) + repository: Repository, + ) { + super(repository); + } + + /** + * Insert the event, relying on the unique (provider, external_event_id) index for dedupe. + * Returns null when the event was already recorded (duplicate delivery / provider replay). + */ + async createDeduped( + data: Partial, + ): Promise { + try { + return await this.create(data); + } catch (err) { + if ( + err instanceof QueryFailedError && + (err.driverError as { code?: string })?.code === PG_UNIQUE_VIOLATION + ) { + return null; + } + throw err; + } + } + + async markProcessed(id: string, processingError?: string): Promise { + await this.update(id, { + processedAt: new Date(), + processingError: processingError ?? null, + }); + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/webhook-processor.service.ts b/apps/edr-payment-api/src/modules/webhooks/webhook-processor.service.ts new file mode 100644 index 000000000..6be66b705 --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/webhook-processor.service.ts @@ -0,0 +1,112 @@ +import { Injectable, Logger } from "@nestjs/common"; +import { + MERCHANT_ORDER_PREFIX, + PaymentService, + ProviderMethod, +} from "@edr/types"; +import { IntentsRepository } from "../intents/intents.repository"; +import { + IntentsService, + ProviderResultInput, +} from "../intents/intents.service"; +import { WebhookEventsRepository } from "./webhook-events.repository"; + +/** A provider webhook reduced to the fields the shared pipeline needs. */ +export interface NormalizedWebhook { + provider: ProviderMethod; + /** Provider event id (or a deterministic derivation) — the dedupe key. */ + externalEventId: string; + merchantOrderId: string; + providerTxnId?: string; + signatureValid: boolean; + /** Raw provider status string, stored for audit. */ + rawStatus: string; + payload: Record; + /** Mapped outcome to feed the intent state machine. */ + result: ProviderResultInput; +} + +/** + * The shared webhook pipeline every provider handler funnels into: + * persist+dedupe → signature gate → intent lookup → prefix/service cross-check → + * state machine → mark processed. Always returns (never throws) so controllers can + * ack 200 fast — providers like Waafi time out at 5s and do not retry. + */ +@Injectable() +export class WebhookProcessorService { + private readonly logger = new Logger(WebhookProcessorService.name); + + constructor( + private readonly webhookEvents: WebhookEventsRepository, + private readonly intentsRepository: IntentsRepository, + private readonly intentsService: IntentsService, + ) {} + + async process(webhook: NormalizedWebhook): Promise { + const { provider, merchantOrderId } = webhook; + + const eventRow = await this.webhookEvents.createDeduped({ + provider, + externalEventId: webhook.externalEventId, + merchantOrderId, + providerTxnId: webhook.providerTxnId ?? null, + signatureValid: webhook.signatureValid, + status: webhook.rawStatus, + payload: webhook.payload, + }); + if (!eventRow) { + this.logger.log( + `${provider} webhook duplicate: ${webhook.externalEventId} — short-circuit OK`, + ); + return; + } + + if (!webhook.signatureValid) { + this.logger.warn( + `${provider} webhook signature invalid/stale for ref=${merchantOrderId}`, + ); + await this.webhookEvents.markProcessed(eventRow.id, "signature-invalid"); + return; + } + + const intent = + await this.intentsRepository.findByMerchantOrderId(merchantOrderId); + if (!intent) { + // Tolerated: webhook may have raced the intent commit, or the reference is foreign. + // The provider gets a 200; retry/poll/reconciliation converges later. + this.logger.warn( + `${provider} webhook: no PaymentIntent for ref=${merchantOrderId}`, + ); + await this.webhookEvents.markProcessed(eventRow.id, "intent-not-found"); + return; + } + + // Integrity guard (§10): the stateless prefix and the stored discriminator must agree. + const expectedPrefix = + MERCHANT_ORDER_PREFIX[intent.service as PaymentService]; + if (expectedPrefix && !merchantOrderId.startsWith(expectedPrefix)) { + this.logger.error( + `${provider} webhook: merchantOrderId ${merchantOrderId} prefix does not match stored service ${intent.service} — refusing to process`, + ); + await this.webhookEvents.markProcessed( + eventRow.id, + "service-prefix-mismatch", + ); + return; + } + + try { + await this.intentsService.applyProviderResult(intent.id, webhook.result); + await this.webhookEvents.markProcessed(eventRow.id); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + this.logger.error( + `${provider} webhook processing failed for ${merchantOrderId}: ${message}`, + ); + await this.webhookEvents.markProcessed( + eventRow.id, + `processing-error: ${message}`, + ); + } + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/webhooks.controller.ts b/apps/edr-payment-api/src/modules/webhooks/webhooks.controller.ts new file mode 100644 index 000000000..95ca51adc --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/webhooks.controller.ts @@ -0,0 +1,143 @@ +import { + All, + Body, + Controller, + Headers, + HttpCode, + HttpStatus, + Logger, + Post, + Req, +} from "@nestjs/common"; +import { ApiOperation, ApiTags } from "@nestjs/swagger"; +import { + CardWebhookPayload, + CbeBirrWebhookPayload, + DMoneyWebhookPayload, + EBirrWebhookPayload, + TelebirrWebhookPayload, + WaafiWebhookHeaders, + WaafiWebhookPayload, +} from "@edr/payment-providers"; +import { TelebirrWebhookService } from "./handlers/telebirr-webhook.service"; +import { CbeBirrWebhookService } from "./handlers/cbe-birr-webhook.service"; +import { EBirrWebhookService } from "./handlers/ebirr-webhook.service"; +import { CardWebhookService } from "./handlers/card-webhook.service"; +import { WaafiWebhookService } from "./handlers/waafi-webhook.service"; +import { DMoneyWebhookService } from "./handlers/dmoney-webhook.service"; + +/** + * The ONLY public surface of the payment service — the single registered webhook URL per + * provider for the whole platform. No service auth here (provider-facing); trust comes from + * signature verification inside each handler. Every route acks 2xx fast and never rethrows: + * Waafi times out at 5s and does NOT retry. + */ +@ApiTags("Provider Webhooks") +@Controller("webhooks") +export class WebhooksController { + private readonly logger = new Logger(WebhooksController.name); + + constructor( + private readonly telebirr: TelebirrWebhookService, + private readonly cbeBirr: CbeBirrWebhookService, + private readonly eBirr: EBirrWebhookService, + private readonly card: CardWebhookService, + private readonly waafi: WaafiWebhookService, + private readonly dMoney: DMoneyWebhookService, + ) {} + + @All("telebirr") + @HttpCode(HttpStatus.OK) + @ApiOperation({ + summary: "Telebirr payment notification callback (Ethiopia)", + }) + async receiveTelebirr(@Body() payload: TelebirrWebhookPayload) { + this.logger.log("Telebirr webhook called"); + try { + await this.telebirr.handle(payload); + } catch (err) { + this.logger.error(`Telebirr webhook handler threw: ${this.message(err)}`); + } + return { code: "0", message: "OK" }; + } + + @Post("cbe-birr") + @HttpCode(HttpStatus.OK) + @ApiOperation({ + summary: "CBE Birr payment notification callback (Ethiopia)", + }) + async receiveCbeBirr(@Body() payload: CbeBirrWebhookPayload) { + try { + await this.cbeBirr.handle(payload); + } catch (err) { + this.logger.error(`CBE Birr webhook handler threw: ${this.message(err)}`); + } + return { success: true }; + } + + @Post("ebirr") + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: "eBirr payment notification callback (Ethiopia)" }) + async receiveEBirr(@Body() payload: EBirrWebhookPayload) { + try { + await this.eBirr.handle(payload); + } catch (err) { + this.logger.error(`eBirr webhook handler threw: ${this.message(err)}`); + } + return { code: "0000", message: "success" }; + } + + @Post("card") + @HttpCode(HttpStatus.OK) + @ApiOperation({ + summary: "Card payment notification callback (International)", + }) + async receiveCard( + @Body() payload: CardWebhookPayload, + @Headers("stripe-signature") signature: string, + ) { + try { + await this.card.handle(payload, signature); + } catch (err) { + this.logger.error(`Card webhook handler threw: ${this.message(err)}`); + } + return { received: true }; + } + + @Post("waafi") + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: "Waafi payment notification callback (Djibouti)" }) + async receiveWaafi( + @Body() payload: WaafiWebhookPayload, + @Headers() headers: WaafiWebhookHeaders, + @Req() req: { rawBody?: Buffer }, + ) { + this.logger.log( + `Waafi webhook hit: event=${payload?.event ?? "unknown"} eventId=${headers["x-webhook-event-id"] ?? "n/a"}`, + ); + try { + // HMAC verification must sign over the exact raw bytes Waafi sent, not re-serialized JSON. + const rawBody = req.rawBody?.toString("utf8") ?? ""; + await this.waafi.handle(payload, rawBody, headers); + } catch (err) { + this.logger.error(`Waafi webhook handler threw: ${this.message(err)}`); + } + return { responseCode: "2001", responseMsg: "Success" }; + } + + @Post("dmoney") + @HttpCode(HttpStatus.OK) + @ApiOperation({ summary: "D-Money payment notification callback (Djibouti)" }) + async receiveDMoney(@Body() payload: DMoneyWebhookPayload) { + try { + await this.dMoney.handle(payload); + } catch (err) { + this.logger.error(`D-Money webhook handler threw: ${this.message(err)}`); + } + return { success: true }; + } + + private message(err: unknown): string { + return err instanceof Error ? err.message : String(err); + } +} diff --git a/apps/edr-payment-api/src/modules/webhooks/webhooks.module.ts b/apps/edr-payment-api/src/modules/webhooks/webhooks.module.ts new file mode 100644 index 000000000..86b94032b --- /dev/null +++ b/apps/edr-payment-api/src/modules/webhooks/webhooks.module.ts @@ -0,0 +1,34 @@ +import { Module } from "@nestjs/common"; +import { TypeOrmModule } from "@nestjs/typeorm"; +import { IntentsModule } from "../intents/intents.module"; +import { ProvidersModule } from "../providers/providers.module"; +import { PaymentWebhookEvent } from "./entities/payment-webhook-event.entity"; +import { WebhookEventsRepository } from "./webhook-events.repository"; +import { WebhookProcessorService } from "./webhook-processor.service"; +import { WebhooksController } from "./webhooks.controller"; +import { TelebirrWebhookService } from "./handlers/telebirr-webhook.service"; +import { CbeBirrWebhookService } from "./handlers/cbe-birr-webhook.service"; +import { EBirrWebhookService } from "./handlers/ebirr-webhook.service"; +import { CardWebhookService } from "./handlers/card-webhook.service"; +import { WaafiWebhookService } from "./handlers/waafi-webhook.service"; +import { DMoneyWebhookService } from "./handlers/dmoney-webhook.service"; + +@Module({ + imports: [ + TypeOrmModule.forFeature([PaymentWebhookEvent]), + IntentsModule, + ProvidersModule, + ], + controllers: [WebhooksController], + providers: [ + WebhookEventsRepository, + WebhookProcessorService, + TelebirrWebhookService, + CbeBirrWebhookService, + EBirrWebhookService, + CardWebhookService, + WaafiWebhookService, + DMoneyWebhookService, + ], +}) +export class WebhooksModule {} diff --git a/apps/edr-payment-api/src/scripts/migrate-revert.ts b/apps/edr-payment-api/src/scripts/migrate-revert.ts new file mode 100644 index 000000000..88b1500e5 --- /dev/null +++ b/apps/edr-payment-api/src/scripts/migrate-revert.ts @@ -0,0 +1,18 @@ +import "dotenv/config"; +import { AppDataSource } from "../data-source"; + +/** `pnpm --filter @edr/payment-api migration:revert` — undo the most recent migration. */ +async function main(): Promise { + await AppDataSource.initialize(); + try { + await AppDataSource.undoLastMigration(); + console.log("reverted last migration"); + } finally { + await AppDataSource.destroy(); + } +} + +main().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/apps/edr-payment-api/src/scripts/migrate.ts b/apps/edr-payment-api/src/scripts/migrate.ts new file mode 100644 index 000000000..05f59f1b9 --- /dev/null +++ b/apps/edr-payment-api/src/scripts/migrate.ts @@ -0,0 +1,23 @@ +import "dotenv/config"; +import { AppDataSource } from "../data-source"; +import { ensurePaymentSchema } from "../config/ensure-schema"; + +/** `pnpm --filter @edr/payment-api migration:run` — ensure schema, then run pending migrations. */ +async function main(): Promise { + await ensurePaymentSchema(); + await AppDataSource.initialize(); + try { + const applied = await AppDataSource.runMigrations(); + for (const migration of applied) { + console.log(`applied: ${migration.name}`); + } + if (applied.length === 0) console.log("no pending migrations"); + } finally { + await AppDataSource.destroy(); + } +} + +main().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/apps/edr-payment-api/tsconfig.build.json b/apps/edr-payment-api/tsconfig.build.json new file mode 100644 index 000000000..64f86c6bd --- /dev/null +++ b/apps/edr-payment-api/tsconfig.build.json @@ -0,0 +1,4 @@ +{ + "extends": "./tsconfig.json", + "exclude": ["node_modules", "test", "dist", "**/*spec.ts"] +} diff --git a/apps/edr-payment-api/tsconfig.json b/apps/edr-payment-api/tsconfig.json new file mode 100644 index 000000000..467c474ee --- /dev/null +++ b/apps/edr-payment-api/tsconfig.json @@ -0,0 +1,14 @@ +{ + "extends": "@edr/tsconfig/nestjs.json", + "compilerOptions": { + "baseUrl": "./", + "outDir": "./dist", + "rootDir": "./src", + "noEmit": false, + "incremental": true, + "tsBuildInfoFile": "./.tsbuildinfo", + "module": "node16", + "moduleResolution": "node16" + }, + "include": ["src"] +} diff --git a/package.json b/package.json index 399ded0bf..575de8c41 100644 --- a/package.json +++ b/package.json @@ -6,6 +6,7 @@ "dev": "turbo run dev", "dev:freight": "turbo run dev --filter=@edr/freight-api... --filter=@edr/freight-portal... --filter=@edr/freight-backoffice... --filter=@edr/ui-common...", "dev:passenger": "turbo run dev --filter=@edr/passenger-api... --filter=@edr/passenger-portal... --filter=@edr/passenger-backoffice...", + "dev:payment": "turbo run dev --filter=@edr/payment-api...", "build": "turbo run build", "build:freight": "turbo run build --filter=@edr/freight-api... --filter=@edr/freight-portal... --filter=@edr/freight-backoffice...", "build:passenger": "turbo run build --filter=@edr/passenger-api... --filter=@edr/passenger-portal... --filter=@edr/passenger-backoffice...", diff --git a/packages/types/src/common/payments.ts b/packages/types/src/common/payments.ts index 18fbfc8b1..0cf2174c3 100644 --- a/packages/types/src/common/payments.ts +++ b/packages/types/src/common/payments.ts @@ -74,3 +74,98 @@ export interface PaymentProvider { initiate(input: ProviderInitiationInput): Promise; queryStatus(merchantOrderId: string): Promise; } + +/* ------------------------------------------------------------------------------------------------ + * Payment microservice contracts (docs/payment-service) + * + * Shared shapes exchanged between the payment microservice (apps/edr-payment-api) and the + * domain apps (passenger/freight). Both sides import these so the wire format cannot drift. + * ---------------------------------------------------------------------------------------------- */ + +/** Which domain app owns the order being paid for. Routing discriminator on every intent. */ +export enum PaymentService { + PASSENGER = "PASSENGER", + FREIGHT = "FREIGHT", +} + +/** What kind of domain order the intent references (soft reference — never a cross-schema FK). */ +export enum PaymentReferenceType { + BOOKING = "BOOKING", + SHIPMENT = "SHIPMENT", +} + +/** `merchant_order_id` prefix per owning service — lets a webhook be routed before a DB lookup. */ +export const MERCHANT_ORDER_PREFIX: Record = { + [PaymentService.PASSENGER]: "PSG-", + [PaymentService.FREIGHT]: "FRT-", +}; + +/** Body of `POST /payments/initiate` on the payment service (internal, service-authenticated). */ +export interface InitiatePaymentRequest { + service: PaymentService; + referenceType: PaymentReferenceType; + /** Domain order id (booking/shipment id). Soft reference; the app has already validated it. */ + referenceId: string; + /** Human-readable order ref (e.g. booking ref) shown on provider pages. Defaults to referenceId. */ + orderRef?: string; + /** App-asserted authoritative amount in minor units (computed server-side by the domain app). */ + amountMinor: number; + currency: string; + provider: ProviderMethod; + platform?: PaymentPlatform; + payerAccount?: string; + /** Optional caller key to dedupe retried initiations beyond the per-reference upsert. */ + idempotencyKey?: string; +} + +/** Response of `POST /payments/initiate` and shape of intent lookups. */ +export interface PaymentIntentSnapshot { + intentId: string; + service: PaymentService; + referenceType: PaymentReferenceType; + referenceId: string; + merchantOrderId: string; + provider: ProviderMethod; + status: ProviderPaymentStatus; + amountMinor: number; + currency: string; + clientAction?: ClientAction; + providerTxnId?: string; + paidAt?: string; + failureCode?: string; + failureMessage?: string; + expiresAt?: string; +} + +export type PaymentEventType = "payment.succeeded" | "payment.failed"; + +/** Versioned envelope delivered (at-least-once) to the owning app's mark-paid consumer. */ +interface PaymentEventBase { + version: 1; + /** Outbox row id — stable across redeliveries; consumers may use it as a dedupe key. */ + eventId: string; + eventType: PaymentEventType; + occurredAt: string; + service: PaymentService; + intentId: string; + referenceType: PaymentReferenceType; + referenceId: string; + merchantOrderId: string; + provider: ProviderMethod; + amountMinor: number; + currency: string; +} + +export interface PaymentSucceededEvent extends PaymentEventBase { + eventType: "payment.succeeded"; + providerTxnId?: string; + paidAt: string; +} + +export interface PaymentFailedEvent extends PaymentEventBase { + eventType: "payment.failed"; + failureCode?: string; + failureMessage?: string; +} + +export type PaymentEvent = PaymentSucceededEvent | PaymentFailedEvent; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 9b91c2013..3d4aa9a32 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -672,6 +672,112 @@ importers: specifier: ^5.5.4 version: 5.9.3 + apps/edr-payment-api: + dependencies: + '@edr/api-common': + specifier: workspace:* + version: link:../../packages/api-common + '@edr/payment-providers': + specifier: workspace:* + version: link:../../packages/payment-providers + '@edr/types': + specifier: workspace:* + version: link:../../packages/types + '@nestjs/axios': + specifier: ^4.0.1 + version: 4.0.1(@nestjs/common@11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(axios@1.17.0)(rxjs@7.8.2) + '@nestjs/common': + specifier: ^11.0.0 + version: 11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2) + '@nestjs/config': + specifier: ^4.0.0 + version: 4.0.4(@nestjs/common@11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(rxjs@7.8.2) + '@nestjs/core': + specifier: ^11.0.0 + version: 11.1.24(@nestjs/common@11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/microservices@11.1.24)(@nestjs/platform-express@11.1.24)(reflect-metadata@0.2.2)(rxjs@7.8.2) + '@nestjs/platform-express': + specifier: ^11.0.0 + version: 11.1.24(@nestjs/common@11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@11.1.24) + '@nestjs/schedule': + specifier: ^6.0.0 + version: 6.1.3(@nestjs/common@11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@11.1.24) + '@nestjs/swagger': + specifier: ^11.4.2 + version: 11.4.4(@nestjs/common@11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@11.1.24)(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2) + '@nestjs/typeorm': + specifier: ^11.0.1 + version: 11.0.1(@nestjs/common@11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@11.1.24)(reflect-metadata@0.2.2)(rxjs@7.8.2)(typeorm@0.3.30(babel-plugin-macros@3.1.0)(pg@8.21.0)(ts-node@10.9.2(@types/node@20.19.42)(typescript@5.9.3))) + axios: + specifier: ^1.16.1 + version: 1.17.0 + class-transformer: + specifier: ^0.5.1 + version: 0.5.1 + class-validator: + specifier: ^0.14.1 + version: 0.14.4 + dotenv: + specifier: ^17.4.2 + version: 17.4.2 + pg: + specifier: ^8.13.0 + version: 8.21.0 + reflect-metadata: + specifier: ^0.2.2 + version: 0.2.2 + rxjs: + specifier: ^7.8.1 + version: 7.8.2 + typeorm: + specifier: 0.3.30 + version: 0.3.30(babel-plugin-macros@3.1.0)(pg@8.21.0)(ts-node@10.9.2(@types/node@20.19.42)(typescript@5.9.3)) + devDependencies: + '@edr/eslint-config': + specifier: workspace:* + version: link:../../packages/config/eslint-config + '@edr/tsconfig': + specifier: workspace:* + version: link:../../packages/config/tsconfig + '@nestjs/cli': + specifier: ^11.0.0 + version: 11.0.21(@types/node@20.19.42)(prettier@3.8.3) + '@nestjs/schematics': + specifier: ^11.0.0 + version: 11.1.0(chokidar@4.0.3)(prettier@3.8.3)(typescript@5.9.3) + '@nestjs/testing': + specifier: ^11.0.0 + version: 11.1.24(@nestjs/common@11.1.24(class-transformer@0.5.1)(class-validator@0.14.4)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@11.1.24)(@nestjs/microservices@11.1.24)(@nestjs/platform-express@11.1.24) + '@types/express': + specifier: ^5.0.0 + version: 5.0.6 + '@types/jest': + specifier: ^29.5.13 + version: 29.5.14 + '@types/node': + specifier: ^20.14.0 + version: 20.19.42 + '@types/pg': + specifier: ^8.6.7 + version: 8.20.0 + jest: + specifier: ^29.7.0 + version: 29.7.0(@types/node@20.19.42)(babel-plugin-macros@3.1.0)(ts-node@10.9.2(@types/node@20.19.42)(typescript@5.9.3)) + ts-jest: + specifier: ^29.2.5 + version: 29.4.11(@babel/core@7.29.7)(@jest/transform@29.7.0)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.29.7))(jest-util@29.7.0)(jest@29.7.0(@types/node@20.19.42)(babel-plugin-macros@3.1.0)(ts-node@10.9.2(@types/node@20.19.42)(typescript@5.9.3)))(typescript@5.9.3) + ts-loader: + specifier: ^9.5.1 + version: 9.6.0(loader-utils@1.4.2)(typescript@5.9.3)(webpack@5.106.0) + ts-node: + specifier: ^10.9.2 + version: 10.9.2(@types/node@20.19.42)(typescript@5.9.3) + tsconfig-paths: + specifier: ^4.2.0 + version: 4.2.0 + typescript: + specifier: ^5.5.4 + version: 5.9.3 + packages/api-common: dependencies: '@edr/types':