Handling Webhooks (recommended)
The package verifies a webhook and builds the acknowledgement; what you do with it lives in your app. The flow below is the one we recommend for production: it answers the gateway fast, never loses a notification, and fulfills each order exactly once even when the gateway retries.
- Verify the webhook with
this.payments.handleCallback(gateway, request). - Store the raw call in your own table, including rejected ones for debugging.
- Acknowledge right away by returning the callback from an
@AcknowledgeCallback()handler, so the gateway stops retrying. - Process once in a background worker: skip what is already processed, lock the order, retry with backoff, keep the last error.
None of this is part of the package: copy the code into your app and adapt the fulfillment to your own orders. The sample uses the built-in node:sqlite module (DatabaseSync, Node.js 22.13+) so it runs without extra dependencies; with TypeORM, Prisma, MikroORM or another database, keep the same columns and the same conditional UPDATEs. The worker runs in-process; for a queue, run processNext() from a BullMQ processor, and the claim logic stays the same.
Migration
CREATE TABLE IF NOT EXISTS payment_webhooks (
id INTEGER PRIMARY KEY AUTOINCREMENT,
gateway VARCHAR(32) NOT NULL,
order_id VARCHAR(64) NULL,
status VARCHAR(16) NULL, -- PaymentStatus value
gateway_status VARCHAR(64) NULL,
gateway_reference VARCHAR(128) NULL,
amount VARCHAR(32) NULL, -- as the gateway sent it, never a float
verified SMALLINT NOT NULL,
body TEXT NOT NULL, -- the raw request body
headers TEXT NOT NULL, -- JSON
attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT NULL,
available_at INTEGER NOT NULL, -- unix time of the next attempt
locked_until INTEGER NULL,
processed_at INTEGER NULL,
created_at INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS payment_webhooks_order
ON payment_webhooks (gateway, order_id, status);The sample assumes an orders table with a unique number, the amount as a decimal string, and nullable paid_at and gateway_reference.
Model
The sample keeps no entity classes: a small provider opens the database and creates both tables. With an ORM, this is where your entities and repositories go.
// src/database.ts
import { DatabaseSync } from 'node:sqlite';
/** Injection token of the app's SQLite database. */
export const DATABASE = Symbol('DATABASE');
/**
* Opens the app database and creates its tables: the orders the
* checkout routes create and the payment_webhooks table of the recommended
* webhook flow.
*
* @param path A file path, or ':memory:' for tests.
*/
export function openDatabase(path = ':memory:'): DatabaseSync {
const db = new DatabaseSync(path);
db.exec(`
CREATE TABLE IF NOT EXISTS orders (
id INTEGER PRIMARY KEY AUTOINCREMENT,
number VARCHAR(64) NOT NULL UNIQUE,
gateway VARCHAR(32) NOT NULL,
amount VARCHAR(32) NOT NULL, -- decimal string, never a float
gateway_reference VARCHAR(128) NULL,
paid_at INTEGER NULL,
created_at INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS payment_webhooks (
id INTEGER PRIMARY KEY AUTOINCREMENT,
gateway VARCHAR(32) NOT NULL,
order_id VARCHAR(64) NULL,
status VARCHAR(16) NULL, -- PaymentStatus value
gateway_status VARCHAR(64) NULL,
gateway_reference VARCHAR(128) NULL,
amount VARCHAR(32) NULL, -- as sent, never a float
verified SMALLINT NOT NULL,
body TEXT NOT NULL, -- the raw request body
headers TEXT NOT NULL, -- JSON
attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT NULL,
available_at INTEGER NOT NULL, -- unix time of the next attempt
locked_until INTEGER NULL,
processed_at INTEGER NULL,
created_at INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS payment_webhooks_order
ON payment_webhooks (gateway, order_id, status);
`);
return db;
}
/** Unix time in seconds. */
export function now(): number {
return Math.floor(Date.now() / 1000);
}Route and Controller
One route serves every gateway. Gateways post from their own servers, so keep it free of authentication guards and CSRF checks. The app must be created with rawBody: true (see Installation).
// src/webhooks/webhooks.controller.ts
import type {
CallbackRequest,
PaymentCallback,
} from '@laranex/myanmar-payments';
import {
AcknowledgeCallback,
GATEWAY_NAMES,
type GatewayName,
RawCallback,
} from '@laranex/nestjs-myanmar-payments';
import { Controller, NotFoundException, Param, Post } from '@nestjs/common';
import { WebhooksService } from './webhooks.service.js';
/** One route for every gateway: verify, store, acknowledge. */
@Controller('webhooks/payments')
export class WebhooksController {
constructor(private readonly webhooks: WebhooksService) {}
@Post(':gateway')
@AcknowledgeCallback()
async receive(
@Param('gateway') gateway: string,
@RawCallback() request: CallbackRequest,
): Promise<PaymentCallback> {
if (!(GATEWAY_NAMES as readonly string[]).includes(gateway)) {
throw new NotFoundException();
}
// Returning the callback answers with the acknowledgement the gateway
// expects.
return this.webhooks.receive(gateway as GatewayName, request);
}
}Use https://shop.test/webhooks/payments/kbz-pay (and wave-money, aya-pay, yoma-mmqr, cyber-source) as the gateway's callback URL.
Job
The worker is where the order is fulfilled. It is safe to run more than once:
- Lock: each row is claimed with a conditional
UPDATE, so one worker at a time processes it, even across instances. - Idempotent: a row that is already processed is skipped, and
UPDATE orders ... WHERE paid_at IS NULLchecks the order itself again. - Retries:
MAX_ATTEMPTSandBACKOFFretry failures; every attempt is counted and the last error is kept.
// src/webhooks/webhooks.service.ts
import type { DatabaseSync } from 'node:sqlite';
import {
Amount,
type CallbackRequest,
type PaymentCallback,
PaymentStatus,
SignatureVerificationError,
} from '@laranex/myanmar-payments';
import {
type GatewayName,
MyanmarPaymentsService,
} from '@laranex/nestjs-myanmar-payments';
import {
BadRequestException,
Inject,
Injectable,
InternalServerErrorException,
Logger,
type OnApplicationBootstrap,
type OnApplicationShutdown,
} from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { DATABASE, now } from '../database.js';
/** Attempts before a stored webhook is left for a manual replay. */
export const MAX_ATTEMPTS = 5;
/**
* Seconds a worker owns a claimed row; a crashed worker's claim expires after
* this.
*/
export const LOCK_FOR = 120;
/** Seconds to wait before each retry. */
export const BACKOFF = [10, 60, 300, 900];
interface StoredWebhook {
id: number;
gateway: string;
order_id: string;
status: string;
gateway_reference: string | null;
amount: string | null;
attempts: number;
}
/**
* The recommended webhook flow, as app code: verify, store, acknowledge
* immediately, then process each stored call once in a background worker with
* retries. The package never stores webhooks.
*/
@Injectable()
export class WebhooksService
implements OnApplicationBootstrap, OnApplicationShutdown
{
private readonly logger = new Logger(WebhooksService.name);
private timer: NodeJS.Timeout | undefined;
constructor(
@Inject(DATABASE) private readonly db: DatabaseSync,
private readonly payments: MyanmarPaymentsService,
private readonly config: ConfigService,
) {}
/**
* Verifies and stores a callback. Returns the verified callback, whose
* acknowledgement the controller sends; processing happens in the worker.
*/
async receive(
gateway: GatewayName,
request: CallbackRequest,
): Promise<PaymentCallback> {
const headers = JSON.stringify(request.headers);
const time = now();
let callback: PaymentCallback;
try {
callback = await this.payments.handleCallback(gateway, request);
} catch (error) {
if (!(error instanceof SignatureVerificationError)) {
throw error;
}
// Rejected calls are kept for debugging, never processed.
this.db
.prepare(
`INSERT INTO payment_webhooks
(gateway, verified, body, headers, last_error, available_at,
created_at)
VALUES (?, 0, ?, ?, ?, ?, ?)`,
)
.run(gateway, request.body, headers, error.message, time, time);
this.logger.warn(`[${gateway}] callback rejected: ${error.message}`);
throw new BadRequestException('invalid signature');
}
try {
this.db
.prepare(
`INSERT INTO payment_webhooks
(gateway, order_id, status, gateway_status, gateway_reference,
amount, verified, body, headers, available_at, created_at)
VALUES (?, ?, ?, ?, ?, ?, 1, ?, ?, ?, ?)`,
)
.run(
gateway,
callback.orderId,
callback.status,
callback.gatewayStatus,
callback.gatewayReference ?? null,
callback.amount ?? null,
request.body,
headers,
time,
time,
);
} catch (error) {
// Not stored: answer 500 so the gateway retries later.
this.logger.error(
`[${gateway}] could not store callback: ${(error as Error).message}`,
);
throw new InternalServerErrorException('try again');
}
this.logger.log(
`[${gateway}] callback verified order=${callback.orderId}` +
` status=${callback.status} amount=${callback.amount ?? ''}`,
);
return callback;
}
onApplicationBootstrap(): void {
const interval = Number(
this.config.get<string>('WEBHOOK_WORKER_INTERVAL_MS') ?? 2000,
);
if (interval > 0) {
this.start(interval);
}
}
onApplicationShutdown(): void {
this.stop();
this.db.close();
}
/** Processes stored webhooks every `intervalMs` until `stop()`. */
start(intervalMs = 2000): void {
this.stop();
this.timer = setInterval(() => void this.drain(), intervalMs);
this.timer.unref();
}
stop(): void {
if (this.timer !== undefined) {
clearInterval(this.timer);
this.timer = undefined;
}
}
/**
* Processes every row that is due right now. Returns how many were
* attempted.
*/
async drain(): Promise<number> {
let count = 0;
while (await this.processNext()) {
count++;
}
return count;
}
/**
* Claims and processes the next due row. Safe to run in several workers or
* processes: the claim is an UPDATE that only one of them wins.
*/
async processNext(): Promise<boolean> {
const time = now();
const webhook = this.db
.prepare(
`SELECT id, gateway, order_id, status, gateway_reference, amount,
attempts
FROM payment_webhooks
WHERE verified = 1 AND processed_at IS NULL AND attempts < ?
AND available_at <= ?
AND (locked_until IS NULL OR locked_until < ?)
ORDER BY id LIMIT 1`,
)
.get(MAX_ATTEMPTS, time, time) as StoredWebhook | undefined;
if (webhook === undefined) {
return false;
}
const claim = this.db
.prepare(
`UPDATE payment_webhooks SET locked_until = ?, attempts = attempts + 1
WHERE id = ? AND processed_at IS NULL
AND (locked_until IS NULL OR locked_until < ?)`,
)
.run(time + LOCK_FOR, webhook.id, time);
if (claim.changes !== 1) {
return true; // another worker won the claim
}
try {
this.fulfill(webhook);
} catch (error) {
const delay = BACKOFF[Math.min(webhook.attempts, BACKOFF.length - 1)];
this.db
.prepare(
`UPDATE payment_webhooks
SET last_error = ?, available_at = ?, locked_until = NULL
WHERE id = ?`,
)
.run((error as Error).message, now() + delay, webhook.id);
const message = (error as Error).message;
this.logger.warn(
`[${webhook.gateway}] webhook ${webhook.id} failed: ${message}`,
);
return true;
}
this.db
.prepare(
`UPDATE payment_webhooks
SET processed_at = ?, last_error = NULL, locked_until = NULL
WHERE id = ?`,
)
.run(now(), webhook.id);
return true;
}
/**
* Queues a stored webhook again, e.g. after fixing the cause of its last
* error.
*/
replay(id: number): boolean {
const result = this.db
.prepare(
`UPDATE payment_webhooks
SET attempts = 0, available_at = ?, locked_until = NULL,
last_error = NULL
WHERE id = ? AND verified = 1 AND processed_at IS NULL`,
)
.run(now(), id);
return result.changes === 1;
}
/**
* Deletes processed and rejected webhooks older than `days`. Returns how
* many were removed.
*/
prune(days = 30): number {
const before = now() - days * 86400;
const result = this.db
.prepare(
`DELETE FROM payment_webhooks
WHERE created_at < ? AND (processed_at IS NOT NULL OR verified = 0)`,
)
.run(before);
return Number(result.changes);
}
/** Marks the order paid once. Throws to retry later. */
private fulfill(webhook: StoredWebhook): void {
if (webhook.status !== PaymentStatus.Successful) {
return; // record failures, cancellations, ... as your app needs
}
const order = this.db
.prepare('SELECT amount, paid_at FROM orders WHERE number = ?')
.get(webhook.order_id) as
{ amount: string; paid_at: number | null } | undefined;
if (order === undefined) {
throw new Error(`order ${webhook.order_id} not found`);
}
if (order.paid_at !== null) {
return; // already fulfilled by an earlier webhook
}
// Yoma MMQR's callback carries no amount; its amount was fixed at checkout.
if (
webhook.amount !== null &&
!Amount.parse(order.amount).equals(webhook.amount)
) {
throw new Error(
`paid ${webhook.amount}, expected ${order.amount}` +
` for order ${webhook.order_id}`,
);
}
// The paid_at IS NULL condition keeps this idempotent even if two
// webhooks for one order run at once.
this.db
.prepare(
`UPDATE orders SET paid_at = ?, gateway_reference = ?
WHERE number = ? AND paid_at IS NULL`,
)
.run(now(), webhook.gateway_reference, webhook.order_id);
}
}- The worker starts with the app (
onApplicationBootstrap) and runs everyWEBHOOK_WORKER_INTERVAL_MS(2000 by default;0turns it off, for example in tests that calldrain()themselves). - The claim is an
UPDATE ... WHERE locked_until IS NULL OR locked_until < now: when several instances run the worker, only one wins each row, and a crashed worker's claim expires afterLOCK_FORseconds. UPDATE orders ... WHERE paid_at IS NULLkeeps fulfillment idempotent even when two webhooks for one order are processed at once.- A row that fails
MAX_ATTEMPTStimes keepslast_errorand stays unprocessed until you replay it.
Register the controller, the worker and the database in your module, and call app.enableShutdownHooks() in main.ts so the worker stops cleanly:
// src/app.module.ts
import { fileURLToPath } from 'node:url';
import { MyanmarPaymentsModule } from '@laranex/nestjs-myanmar-payments';
import { Module } from '@nestjs/common';
import { ConfigModule, ConfigService } from '@nestjs/config';
import { DATABASE, openDatabase } from './database.js';
import { WebhooksController } from './webhooks/webhooks.controller.js';
import { WebhooksService } from './webhooks/webhooks.service.js';
@Module({
imports: [
ConfigModule.forRoot({ isGlobal: true }),
MyanmarPaymentsModule.forRootAsync({
isGlobal: true,
inject: [ConfigService],
useFactory: (config: ConfigService) => ({ env: config }),
}),
],
controllers: [WebhooksController],
providers: [
{ provide: DATABASE, useFactory: () => openDatabase('payments.sqlite') },
WebhooksService,
],
})
export class AppModule {}Replay and Prune
A failed row keeps last_error and processed_at = NULL. After fixing the cause, process it again; the worker picks it up on its next run:
webhooks.replay(42); // true when a verified, unprocessed row was foundDelete old processed and rejected rows on a schedule, for example with @nestjs/schedule; failed rows stay until you replay or delete them:
@Cron('0 3 * * *')
prune(): void {
this.webhooks.prune(90);
}Testing
Post a signed callback with supertest to check that the webhook is stored and acknowledged, then run the worker yourself:
import request from 'supertest';
import { WebhooksService } from '../src/webhooks/webhooks.service.js';
it('stores, acknowledges and processes a webhook', async () => {
const app = moduleRef.createNestApplication({ rawBody: true });
await app.init();
await request(app.getHttpServer())
.post('/webhooks/payments/wave-money')
.set('Content-Type', 'application/json')
.send(JSON.stringify(signedWavePayload))
.expect(200);
expect(await app.get(WebhooksService).drain()).toBe(1);
});Set WEBHOOK_WORKER_INTERVAL_MS=0 in tests so only drain() processes rows. To send a signed payload, compute the gateway's signature in the test with the secret from your test configuration, as described on each gateway's page and in Testing.