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
HandleCallback. - Store the raw call in your own table, including rejected ones for debugging.
- Acknowledge right away with
callback.Acknowledgement.Write(w), so the gateway stops retrying. - Process once in a background worker: claim the row, skip what is already fulfilled, 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 table. The sample uses database/sql with ? placeholders and is framework-independent; with PostgreSQL, use your driver's placeholders ($1, $2, …) and the same queries.
Table
CREATE TABLE payment_webhooks (
-- BIGINT AUTO_INCREMENT on MySQL, BIGSERIAL on PostgreSQL
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 payment_webhooks_order
ON payment_webhooks (gateway, order_id, status);Handler
Handle takes the gateway name from the route and the request, and writes the acknowledgement, or a 400 for a rejected call:
package webhooks
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"log"
"net/http"
"time"
myanmarpayments "github.com/laranex/go-myanmar-payments/v4"
)
// Verifier is what every gateway (*kbzpay.Gateway, *wavemoney.Gateway,
// ...) has in common.
type Verifier interface {
HandleCallback(
request *myanmarpayments.CallbackRequest,
) (*myanmarpayments.PaymentCallback, error)
}
type Webhooks struct {
DB *sql.DB
Gateways map[string]Verifier // "kbz-pay" => kbz, "wave-money" => wave, ...
}
// Handle verifies, stores and acknowledges. Processing happens in Work.
// Register it with
// mux.HandleFunc("POST /webhooks/payments/{gateway}", webhooks.Handle).
func (h *Webhooks) Handle(w http.ResponseWriter, r *http.Request) {
name := r.PathValue("gateway")
gateway, ok := h.Gateways[name]
if !ok {
http.NotFound(w, r)
return
}
request, err := myanmarpayments.NewCallbackRequestFromHTTP(r)
if err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
headers, _ := json.Marshal(request.Header)
now := time.Now().Unix()
callback, err := gateway.HandleCallback(request)
if err != nil {
var sigErr *myanmarpayments.SignatureVerificationError
if !errors.As(err, &sigErr) {
http.Error(w, "invalid callback", http.StatusBadRequest)
return
}
_, _ = h.DB.ExecContext(r.Context(), `INSERT INTO payment_webhooks
(gateway, verified, body, headers, last_error,
available_at, created_at)
VALUES (?, 0, ?, ?, ?, ?, ?)`,
name, string(request.Body), string(headers), sigErr.Error(),
now, now)
http.Error(w, "invalid signature", http.StatusBadRequest)
return
}
_, err = h.DB.ExecContext(r.Context(), `INSERT INTO payment_webhooks
(gateway, order_id, status, gateway_status, gateway_reference,
amount, verified, body, headers, available_at, created_at)
VALUES (?, ?, ?, ?, ?, ?, 1, ?, ?, ?, ?)`,
name, callback.OrderID, string(callback.Status),
callback.GatewayStatus, callback.GatewayReference, callback.Amount,
string(request.Body), string(headers), now, now)
if err != nil {
// Not stored: answer 500 so the gateway retries later.
http.Error(w, "try again", http.StatusInternalServerError)
return
}
if err := callback.Acknowledgement.Write(w); err != nil {
log.Printf("write acknowledgement: %v", err)
}
}Worker
Work processes stored webhooks one at a time and is safe to run in several processes:
- Lock: a worker claims a row by setting
locked_until; a crashed worker's claim expires afterlockFor. - Idempotent: the order is only marked paid while
paid_at IS NULL, so a retried or duplicated webhook never fulfills twice. - Retries: a failure keeps
last_errorand waitsbackoffbefore the next attempt, up tomaxAttempts.
const (
maxAttempts = 5
lockFor = 2 * time.Minute
)
// backoff is the wait before each retry.
var backoff = []time.Duration{
10 * time.Second, time.Minute, 5 * time.Minute, 15 * time.Minute,
}
// Work processes stored webhooks until ctx is canceled. It is safe to run
// in several goroutines or processes.
func (h *Webhooks) Work(ctx context.Context) {
for ctx.Err() == nil {
found, err := h.processNext(ctx)
if err != nil {
log.Printf("payment webhooks: %v", err)
}
if !found || err != nil {
select {
case <-ctx.Done():
case <-time.After(2 * time.Second):
}
}
}
}
type webhook struct {
ID int64
OrderID string
Status myanmarpayments.PaymentStatus
GatewayReference string
Amount string
Attempts int
}
func (h *Webhooks) processNext(ctx context.Context) (bool, error) {
now := time.Now().Unix()
var wh webhook
var status string
err := h.DB.QueryRowContext(ctx, `SELECT
id, 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`, maxAttempts, now, now).
Scan(&wh.ID, &wh.OrderID, &status, &wh.GatewayReference,
&wh.Amount, &wh.Attempts)
if errors.Is(err, sql.ErrNoRows) {
return false, nil
}
if err != nil {
return false, err
}
wh.Status = myanmarpayments.PaymentStatus(status)
// Claim the row; a worker that read it at the same time affects 0 rows
// and moves on.
res, err := h.DB.ExecContext(ctx, `UPDATE payment_webhooks
SET locked_until = ?, attempts = attempts + 1
WHERE id = ? AND processed_at IS NULL
AND (locked_until IS NULL OR locked_until < ?)`,
now+int64(lockFor.Seconds()), wh.ID, now)
if err != nil {
return true, err
}
if n, _ := res.RowsAffected(); n != 1 {
return true, nil
}
if err := h.fulfill(ctx, wh); err != nil {
delay := backoff[min(wh.Attempts, len(backoff)-1)]
_, dbErr := h.DB.ExecContext(ctx, `UPDATE payment_webhooks
SET last_error = ?, available_at = ?, locked_until = NULL
WHERE id = ?`,
err.Error(), time.Now().Add(delay).Unix(), wh.ID)
return true, dbErr
}
_, err = h.DB.ExecContext(ctx, `UPDATE payment_webhooks
SET processed_at = ?, last_error = NULL, locked_until = NULL
WHERE id = ?`,
time.Now().Unix(), wh.ID)
return true, err
}
func (h *Webhooks) fulfill(ctx context.Context, wh webhook) error {
if wh.Status != myanmarpayments.StatusSuccessful {
return nil // record failures, cancellations, ... as your app needs
}
var amount string
var paidAt sql.NullInt64
err := h.DB.QueryRowContext(ctx,
`SELECT amount, paid_at FROM orders WHERE number = ?`, wh.OrderID,
).Scan(&amount, &paidAt)
if errors.Is(err, sql.ErrNoRows) {
return fmt.Errorf("order %s not found", wh.OrderID)
}
if err != nil {
return err
}
if paidAt.Valid {
return nil // already fulfilled by an earlier webhook
}
// Yoma MMQR callbacks carry no amount: Yoma fixed it when the order
// was checked out.
expected, err := myanmarpayments.ParseAmount(amount)
if err != nil {
return err
}
if wh.Amount != "" && !expected.Equals(wh.Amount) {
return fmt.Errorf("paid %s, expected %s for order %s",
wh.Amount, amount, wh.OrderID)
}
// The paid_at IS NULL condition keeps this idempotent even if two
// webhooks for one order run at once.
_, err = h.DB.ExecContext(ctx, `UPDATE orders
SET paid_at = ?, gateway_reference = ?
WHERE number = ? AND paid_at IS NULL`,
time.Now().Unix(), wh.GatewayReference, wh.OrderID)
return err
}The sample assumes an orders table with a unique number, the amount as decimal text, and nullable paid_at and gateway_reference. Amount.Equals makes 1000, 1000.0 and 1000.00 compare equal.
Wiring It Up
With net/http, one route serves every gateway and a goroutine runs the worker; any other router builds the CallbackRequest and writes the Acknowledgement the same way (see Framework Integration):
package main
import (
"context"
"database/sql"
"log"
"net/http"
"os"
"os/signal"
"syscall"
"github.com/laranex/go-myanmar-payments/v4/kbzpay"
"github.com/laranex/go-myanmar-payments/v4/wavemoney"
_ "modernc.org/sqlite"
"example.com/shop/webhooks"
)
func main() {
ctx, stop := signal.NotifyContext(
context.Background(), os.Interrupt, syscall.SIGTERM,
)
defer stop()
db, err := sql.Open("sqlite", "payments.sqlite")
if err != nil {
log.Fatal(err)
}
kbz, err := kbzpay.New(kbzpay.ConfigFromEnv(os.Getenv), nil)
if err != nil {
log.Fatal(err)
}
wave, err := wavemoney.New(wavemoney.ConfigFromEnv(os.Getenv), nil)
if err != nil {
log.Fatal(err)
}
hooks := &webhooks.Webhooks{
DB: db,
Gateways: map[string]webhooks.Verifier{
"kbz-pay": kbz,
"wave-money": wave,
},
}
mux := http.NewServeMux()
mux.HandleFunc("POST /webhooks/payments/{gateway}", hooks.Handle)
go hooks.Work(ctx)
log.Fatal(http.ListenAndServe(":8080", mux))
}Use https://shop.test/webhooks/payments/kbz-pay (and so on) as each gateway's callback URL. In production, run the worker as its own process instead: a separate main that calls Work, e.g. as a systemd service.
Replay and Prune
A row that ran out of attempts keeps last_error and processed_at IS NULL. After fixing the cause, make it available again and a worker picks it up:
UPDATE payment_webhooks
SET attempts = 0, available_at = 0, last_error = NULL
WHERE id = 42;To run a stored call through verification again, e.g. after rotating a key, rebuild it from the stored body and headers: decode the headers column into an http.Header with json.Unmarshal, then call myanmarpayments.NewCallbackRequest([]byte(body), header, nil).
Delete old processed and rejected rows once a day; failed rows stay until you replay or delete them:
cutoff := time.Now().AddDate(0, 0, -90).Unix()
_, err := db.ExecContext(ctx, `DELETE FROM payment_webhooks
WHERE created_at < ? AND (processed_at IS NOT NULL OR verified = 0)`,
cutoff)