Checkout & Payment Q-Commerce: Saga Pattern untuk Konsistensi Terdistribusi
System design checkout flow untuk quick-commerce (q-commerce) dengan Saga orchestration pattern untuk distributed transaction consistency. Pelajari reserve stock, create order PENDING, charge payment via payment gateway, confirm order, publish event via transactional outbox, compensating transaction untuk rollback (release stock, cancel order), idempotency key untuk payment callback, webhook signature verification dengan HMAC, outbox relay worker CDC-based, dan edge cases seperti double charge, partial refund, payment timeout, dan network partition. Dilengkapi implementasi Golang dengan context.Context, interface-based services, step-based saga executor, dan database transaction patterns.
- Checkout & Payment Q-Commerce: Saga Pattern untuk Konsistensi Terdistribusi
- Masalah: Kenapa Checkout Q-Commerce Butuh Saga?
- Saga Pattern vs 2PC
- Arsitektur Checkout Saga
- Flow Sukses (Happy Path)
- Flow Failure with Compensation
- System Components
- 1. Saga Orchestrator: Koordinator Utama
- 2. Checkout Service: Mendefinisikan Steps
- 3. Idempotency Handler: Anti Double Charge
- 4. Payment Webhook Handler: HMAC Signature + Idempotent Processing
- 5. Transactional Outbox: Reliable Event Publishing
- 6. Service Clients (gRPC)
- Architecture Overview
- Edge Cases dan Penanganannya
- 1. Double Charge (Payment Gateway Retry)
- 2. Partial Refund vs Full Refund
- 3. Outbox Relay Duplikasi
- 4. Payment Timeout / Gateway Down
- 5. Concurrent Checkout (Stock Race Condition)
- Key Takeaways
Checkout & Payment Q-Commerce: Saga Pattern untuk Konsistensi Terdistribusi
"Di Q-Commerce, checkout adalah titik paling kritis dalam seluruh sistem. Jika gagal di tengah, kamu tidak bisa rollback database — karena stock sudah di-reserve, payment mungkin sudah ter-charge, dan driver sudah dalam perjalanan."
TL;DR
Checkout flow di Q-Commerce melibatkan minimal 4 service berbeda: Inventory, Order, Payment, dan Notification. Tidak ada database tunggal yang mengkoordinasikan semuanya. Saga Pattern — khususnya Orchestration Saga — adalah solusi untuk menjaga konsistensi di distributed transaction tanpa two-phase commit (2PC) yang berat. Artikel ini membahas implementasi lengkap: saga orchestrator dengan step/compensation, idempotency key untuk payment idempotency, transactional outbox untuk reliable event publishing, webhook signature verification, dan edge cases yang sering terjadi di production.
Masalah: Kenapa Checkout Q-Commerce Butuh Saga?
Flow checkout sederhana secara konseptual:
- Reserve stock di Inventory Service
- Create order dengan status PENDING di Order Service
- Charge payment via Payment Gateway (PG)
- Confirm order (PENDING → CONFIRMED)
- Publish event (dispatch, notification, analytics)
Masalahnya: Setiap langkah adalah service yang berbeda dengan database sendiri. Jika payment berhasil tapi order gagal confirm — kamu punya paid order yang tidak terkirim ('em'). Jika stock reserve berhasil tapi payment gagal — stock terblokir tanpa digunakan.
Saga Pattern vs 2PC
2PC (Two-Phase Commit)
Strong consistency. Tapi blocking — semua resource terkunci sampai coordinator memutuskan. Tidak scalable untuk microservices. Tidak feasible untuk external payment gateway.
Saga (Orchestration)
Eventually consistent. Setiap langkah punya **compensating action** untuk rollback. Non-blocking. Cocok untuk long-running transaction lintas service. Tapi butuh careful error handling.
Keputusan Arsitektural
Di Q-Commerce, Saga Orchestration lebih cocok daripada Choreography. Kenapa? Orchestration terpusat: flow checkout rumit (banyak branching berdasarkan payment method, promo, voucher), dan kamu perlu visibilitas penuh terhadap status setiap transaksi. Choreography bagus untuk flow sederhana, tapi checkout Q-Commerce bukan flow sederhana.
Arsitektur Checkout Saga
Flow Sukses (Happy Path)
Flow Failure with Compensation
System Components
1. Saga Orchestrator: Koordinator Utama
Saga orchestrator adalah heart dari checkout system. Ia menjalankan step secara berurutan, dan jika ada yang gagal, menjalankan compensating actions secara reverse order.
package checkout
import (
"context"
"fmt"
"log/slog"
"time"
)
// SagaStep represents one step in the saga.
type SagaStep struct {
Name string
Execute func(ctx context.Context, ctx SagaContext) error
Compensate func(ctx context.Context, ctx SagaContext) error
}
// SagaContext carries data between saga steps.
type SagaContext struct {
OrderID string
UserID string
Items []OrderItem
Amount int64 // in smallest currency unit (cents)
PaymentMethod string
IdempotencyKey string
PaymentTxID string
ReservationID string
}
// OrderItem represents a product in the order.
type OrderItem struct {
ProductID string `json:"product_id"`
Quantity int `json:"quantity"`
PriceCents int64 `json:"price_cents"`
}
// SagaOrchestrator executes sagas with compensation.
type SagaOrchestrator struct {
steps []SagaStep
logger *slog.Logger
metrics MetricsRecorder
}
// MetricsRecorder for observability.
type MetricsRecorder interface {
IncrementCounter(name string, labels map[string]string)
ObserveHistogram(name string, value float64, labels map[string]string)
}
// NewSagaOrchestrator creates an orchestrator.
func NewSagaOrchestrator(steps []SagaStep, logger *slog.Logger, metrics MetricsRecorder) *SagaOrchestrator {
return &SagaOrchestrator{
steps: steps,
logger: logger,
metrics: metrics,
}
}
// SagaResult contains the outcome of saga execution.
type SagaResult struct {
Success bool `json:"success"`
OrderID string `json:"order_id"`
CompletedAt string `json:"completed_at,omitempty"`
FailedAtStep string `json:"failed_at_step,omitempty"`
Compensated []string `json:"compensated_steps,omitempty"`
}
// Execute runs the saga: steps forward, compensates on failure.
func (so *SagaOrchestrator) Execute(ctx context.Context, sagaCtx SagaContext) *SagaResult {
completedSteps := make([]string, 0, len(so.steps))
for _, step := range so.steps {
select {
case <-ctx.Done():
// Context cancelled (timeout, client disconnect)
so.compensate(ctx, sagaCtx, completedSteps)
return &SagaResult{
Success: false,
OrderID: sagaCtx.OrderID,
FailedAtStep: step.Name,
Compensated: completedSteps,
}
default:
}
so.logger.Info("saga step executing", "step", step.Name, "order_id", sagaCtx.OrderID)
start := time.Now()
if err := step.Execute(ctx, sagaCtx); err != nil {
so.logger.Error("saga step failed",
"step", step.Name,
"order_id", sagaCtx.OrderID,
"error", err,
)
so.metrics.IncrementCounter("saga.step.failure", map[string]string{
"step": step.Name,
})
// Compensate completed steps in reverse order
compensated := so.compensate(ctx, sagaCtx, completedSteps)
so.metrics.IncrementCounter("saga.failure", map[string]string{
"failed_step": step.Name,
})
return &SagaResult{
Success: false,
OrderID: sagaCtx.OrderID,
FailedAtStep: step.Name,
Compensated: compensated,
}
}
so.metrics.ObserveHistogram("saga.step.duration_ms",
float64(time.Since(start).Milliseconds()),
map[string]string{"step": step.Name},
)
completedSteps = append(completedSteps, step.Name)
}
so.metrics.IncrementCounter("saga.success", nil)
return &SagaResult{
Success: true,
OrderID: sagaCtx.OrderID,
CompletedAt: time.Now().UTC().Format(time.RFC3339),
}
}
// compensate runs compensating actions in reverse order.
func (so *SagaOrchestrator) compensate(ctx context.Context, sagaCtx SagaContext, completedSteps []string) []string {
compensated := make([]string, 0, len(completedSteps))
// Walk completed steps in reverse
for i := len(completedSteps) - 1; i >= 0; i-- {
stepName := completedSteps[i]
// Find the step
for _, step := range so.steps {
if step.Name != stepName {
continue
}
if step.Compensate == nil {
so.logger.Warn("step has no compensation", "step", stepName)
continue
}
so.logger.Info("compensating step", "step", stepName, "order_id", sagaCtx.OrderID)
if err := step.Compensate(ctx, sagaCtx); err != nil {
so.logger.Error("compensation failed",
"step", stepName,
"order_id", sagaCtx.OrderID,
"error", err,
)
// Continue compensating other steps — don't stop on failure
}
compensated = append(compensated, stepName)
break
}
}
return compensated
}2. Checkout Service: Mendefinisikan Steps
Checkout service mendefinisikan saga steps dan menghubungkan ke service lain.
package checkout
import (
"context"
"fmt"
"time"
)
// CheckoutService orchestrates the full checkout flow.
type CheckoutService struct {
inventoryClient *InventoryClient
orderClient *OrderClient
paymentClient *PaymentClient
outboxWriter *OutboxWriter
logger *slog.Logger
metrics MetricsRecorder
}
// NewCheckoutService creates the checkout service.
func NewCheckoutService(
inventory *InventoryClient,
order *OrderClient,
payment *PaymentClient,
outbox *OutboxWriter,
logger *slog.Logger,
metrics MetricsRecorder,
) *CheckoutService {
return &CheckoutService{
inventoryClient: inventory,
orderClient: order,
paymentClient: payment,
outboxWriter: outbox,
logger: logger,
metrics: metrics,
}
}
// Checkout initiates the checkout saga.
func (cs *CheckoutService) Checkout(ctx context.Context, req CheckoutRequest) (*SagaResult, error) {
if err := req.Validate(); err != nil {
return nil, fmt.Errorf("checkout: invalid request: %w", err)
}
// Generate order ID and idempotency key
orderID := generateOrderID()
idempotencyKey := fmt.Sprintf("checkout:%s:%d", orderID, time.Now().Unix())
sagaCtx := SagaContext{
OrderID: orderID,
UserID: req.UserID,
Items: req.Items,
Amount: calculateTotal(req.Items),
PaymentMethod: req.PaymentMethod,
IdempotencyKey: idempotencyKey,
}
steps := []SagaStep{
{
Name: "reserve_stock",
Execute: func(ctx context.Context, sc SagaContext) error {
resp, err := cs.inventoryClient.ReserveStock(ctx, sc.OrderID, sc.Items)
if err != nil {
return fmt.Errorf("reserve stock: %w", err)
}
sc.ReservationID = resp.ReservationID
return nil
},
Compensate: func(ctx context.Context, sc SagaContext) error {
return cs.inventoryClient.ReleaseStock(ctx, sc.OrderID, sc.Items)
},
},
{
Name: "create_order_pending",
Execute: func(ctx context.Context, sc SagaContext) error {
return cs.orderClient.CreateOrder(ctx, sc.OrderID, sc.UserID, sc.Items, OrderStatusPending)
},
Compensate: func(ctx context.Context, sc SagaContext) error {
return cs.orderClient.CancelOrder(ctx, sc.OrderID, "payment_failed")
},
},
{
Name: "charge_payment",
Execute: func(ctx context.Context, sc SagaContext) error {
resp, err := cs.paymentClient.Charge(ctx, PaymentRequest{
IdempotencyKey: sc.IdempotencyKey,
Amount: sc.Amount,
Currency: "IDR",
PaymentMethod: sc.PaymentMethod,
OrderID: sc.OrderID,
})
if err != nil {
return fmt.Errorf("charge payment: %w", err)
}
sc.PaymentTxID = resp.TransactionID
return nil
},
Compensate: func(ctx context.Context, sc SagaContext) error {
if sc.PaymentTxID != "" {
return cs.paymentClient.Refund(ctx, sc.PaymentTxID, sc.Amount)
}
return nil
},
},
{
Name: "confirm_order",
Execute: func(ctx context.Context, sc SagaContext) error {
return cs.orderClient.ConfirmOrder(ctx, sc.OrderID, sc.PaymentTxID)
},
Compensate: func(ctx context.Context, sc SagaContext) error {
return cs.orderClient.CancelOrder(ctx, sc.OrderID, "compensation_after_confirm")
},
},
{
Name: "publish_events",
Execute: func(ctx context.Context, sc SagaContext) error {
events := []OutboxEvent{
{
AggregateType: "order",
AggregateID: sc.OrderID,
EventType: "order.confirmed",
Payload: marshalJSON(OrderConfirmedEvent{OrderID: sc.OrderID, Amount: sc.Amount}),
},
{
AggregateType: "payment",
AggregateID: sc.OrderID,
EventType: "payment.success",
Payload: marshalJSON(PaymentSuccessEvent{OrderID: sc.OrderID, TxID: sc.PaymentTxID}),
},
{
AggregateType: "dispatch",
AggregateID: sc.OrderID,
EventType: "dispatch.requested",
Payload: marshalJSON(DispatchRequestedEvent{OrderID: sc.OrderID, Items: sc.Items}),
},
}
return cs.outboxWriter.WriteBatch(ctx, events)
},
Compensate: nil, // Events cannot be unpublished; downstream handles dedup
},
}
orchestrator := NewSagaOrchestrator(steps, cs.logger, cs.metrics)
result := orchestrator.Execute(ctx, sagaCtx)
return result, nil
}
// CheckoutRequest is the API input.
type CheckoutRequest struct {
UserID string `json:"user_id"`
Items []OrderItem `json:"items"`
PaymentMethod string `json:"payment_method"`
}
func (r CheckoutRequest) Validate() error {
if r.UserID == "" {
return fmt.Errorf("user_id required")
}
if len(r.Items) == 0 {
return fmt.Errorf("at least one item required")
}
if r.PaymentMethod == "" {
return fmt.Errorf("payment_method required")
}
return nil
}
func calculateTotal(items []OrderItem) int64 {
var total int64
for _, item := range items {
total += item.PriceCents * int64(item.Quantity)
}
return total
}
func generateOrderID() string {
return fmt.Sprintf("ORD-%d", time.Now().UnixNano())
}
// Order status constants.
type OrderStatus string
const (
OrderStatusPending OrderStatus = "PENDING"
OrderStatusConfirmed OrderStatus = "CONFIRMED"
OrderStatusCancelled OrderStatus = "CANCELLED"
)Compensation Strategy
Perhatikan bahwa publish_events tidak memiliki compensate. Ini intentional — sekali event ter-publish ke message broker, tidak bisa ditarik. Downstream services harus handle idempotency dan eventual consistency. Contoh: Dispatch service mungkin sudah meng-assign driver — ketika menerima order.cancelled event, dispatch harus cancel assignment dan mencari order baru.
3. Idempotency Handler: Anti Double Charge
Idempotency key adalah mekanisme paling penting di payment flow. Tanpa ini, satu klik "Bayar" bisa menghasilkan 2 charge jika user refresh atau network timeout.
package checkout
import (
"bytes"
"context"
"encoding/json"
"fmt"
"net/http"
"time"
"github.com/redis/go-redis/v9"
)
// IdempotencyMiddleware provides idempotency for HTTP handlers.
type IdempotencyMiddleware struct {
rdb *redis.Client
ttl time.Duration
}
// NewIdempotencyMiddleware creates the middleware.
func NewIdempotencyMiddleware(rdb *redis.Client, ttl time.Duration) *IdempotencyMiddleware {
return &IdempotencyMiddleware{
rdb: rdb,
ttl: ttl,
}
}
// IdempotencyResult is stored in Redis.
type IdempotencyResult struct {
StatusCode int `json:"status_code"`
Headers http.Header `json:"headers"`
Body json.RawMessage `json:"body"`
}
// WrapHandler wraps an HTTP handler with idempotency check.
func (im *IdempotencyMiddleware) WrapHandler(next http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
key := r.Header.Get("Idempotency-Key")
if key == "" {
http.Error(w, "Idempotency-Key header required", http.StatusBadRequest)
return
}
// Check if we already processed this key
existing, err := im.getResult(r.Context(), key)
if err == nil && existing != nil {
// Return cached response
for k, vals := range existing.Headers {
for _, v := range vals {
w.Header().Add(k, v)
}
}
w.WriteHeader(existing.StatusCode)
w.Write(existing.Body)
return
}
// Lock the idempotency key (prevent concurrent requests)
locked, err := im.lockKey(r.Context(), key)
if err != nil || !locked {
http.Error(w, "Concurrent request detected", http.StatusConflict)
return
}
defer im.unlockKey(r.Context(), key)
// Capture response
recorder := &responseRecorder{
ResponseWriter: w,
statusCode: http.StatusOK,
}
next(recorder, r)
// Store result
im.storeResult(r.Context(), key, &IdempotencyResult{
StatusCode: recorder.statusCode,
Headers: recorder.Header(),
Body: recorder.body.Bytes(),
})
}
}
func (im *IdempotencyMiddleware) lockKey(ctx context.Context, key string) (bool, error) {
ok, err := im.rdb.SetNX(ctx, fmt.Sprintf("idem:lock:%s", key), "1", 5*time.Second).Result()
if err != nil {
return false, fmt.Errorf("idempotency lock: %w", err)
}
return ok, nil
}
func (im *IdempotencyMiddleware) unlockKey(ctx context.Context, key string) {
im.rdb.Del(ctx, fmt.Sprintf("idem:lock:%s", key))
}
func (im *IdempotencyMiddleware) getResult(ctx context.Context, key string) (*IdempotencyResult, error) {
data, err := im.rdb.Get(ctx, fmt.Sprintf("idem:result:%s", key)).Bytes()
if err == redis.Nil {
return nil, nil
}
if err != nil {
return nil, err
}
var res IdempotencyResult
if err := json.Unmarshal(data, &res); err != nil {
return nil, err
}
return &res, nil
}
func (im *IdempotencyMiddleware) storeResult(ctx context.Context, key string, result *IdempotencyResult) {
data, _ := json.Marshal(result)
im.rdb.Set(ctx, fmt.Sprintf("idem:result:%s", key), data, im.ttl)
}
// responseRecorder captures HTTP response for idempotency.
type responseRecorder struct {
http.ResponseWriter
statusCode int
body bytes.Buffer
}
func (rr *responseRecorder) WriteHeader(code int) {
rr.statusCode = code
rr.ResponseWriter.WriteHeader(code)
}
func (rr *responseRecorder) Write(b []byte) (int, error) {
rr.body.Write(b)
return rr.ResponseWriter.Write(b)
}4. Payment Webhook Handler: HMAC Signature + Idempotent Processing
Payment gateway mengirim callback (webhook) ketika status transaksi berubah. Webhook harus:
- Verifikasi signature (HMAC-SHA256)
- Cek idempotency (jika webhook duplikat)
- Update order status
- Handle berbagai status: success, failed, expired
package checkout
import (
"crypto/hmac"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"log/slog"
"net/http"
"time"
)
// PaymentWebhookHandler processes payment gateway callbacks.
type PaymentWebhookHandler struct {
webhookSecret string
orderClient *OrderClient
paymentClient *PaymentClient
idempotency *IdempotencyMiddleware
logger *slog.Logger
}
// NewPaymentWebhookHandler creates the handler.
func NewPaymentWebhookHandler(
secret string,
orderClient *OrderClient,
paymentClient *PaymentClient,
idempotency *IdempotencyMiddleware,
logger *slog.Logger,
) *PaymentWebhookHandler {
return &PaymentWebhookHandler{
webhookSecret: secret,
orderClient: orderClient,
paymentClient: paymentClient,
idempotency: idempotency,
logger: logger,
}
}
// WebhookPayload from payment gateway.
type WebhookPayload struct {
Event string `json:"event"` // payment.success, payment.failed
OrderID string `json:"order_id"`
TxID string `json:"transaction_id"`
Amount int64 `json:"amount"`
Status string `json:"status"`
Signature string `json:"signature"`
Timestamp string `json:"timestamp"`
}
// HandleWebhook processes an incoming webhook.
func (h *PaymentWebhookHandler) HandleWebhook(w http.ResponseWriter, r *http.Request) {
// Read body
body, err := io.ReadAll(r.Body)
if err != nil {
http.Error(w, "cannot read body", http.StatusBadRequest)
return
}
defer r.Body.Close()
// Parse payload
var payload WebhookPayload
if err := json.Unmarshal(body, &payload); err != nil {
http.Error(w, "invalid payload", http.StatusBadRequest)
return
}
// Verify HMAC signature
if !h.verifySignature(body, payload.Signature) {
h.logger.Warn("webhook invalid signature", "order_id", payload.OrderID)
http.Error(w, "invalid signature", http.StatusUnauthorized)
return
}
h.logger.Info("webhook received",
"event", payload.Event,
"order_id", payload.OrderID,
"tx_id", payload.TxID,
)
// Process based on event type
switch payload.Event {
case "payment.success":
err = h.handlePaymentSuccess(r.Context(), payload)
case "payment.failed":
err = h.handlePaymentFailed(r.Context(), payload)
case "payment.expired":
err = h.handlePaymentExpired(r.Context(), payload)
default:
h.logger.Warn("unknown webhook event", "event", payload.Event)
http.Error(w, "unknown event", http.StatusBadRequest)
return
}
if err != nil {
h.logger.Error("webhook processing failed",
"event", payload.Event,
"order_id", payload.OrderID,
"error", err,
)
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusOK)
json.NewEncoder(w).Encode(map[string]string{
"status": "processed",
"tx_id": payload.TxID,
})
}
func (h *PaymentWebhookHandler) verifySignature(body []byte, signature string) bool {
mac := hmac.New(sha256.New, []byte(h.webhookSecret))
mac.Write(body)
expected := hex.EncodeToString(mac.Sum(nil))
return hmac.Equal([]byte(expected), []byte(signature))
}
func (h *PaymentWebhookHandler) handlePaymentSuccess(ctx context.Context, payload WebhookPayload) error {
// Idempotency: check if already processed
key := fmt.Sprintf("webhook:payment.success:%s", payload.TxID)
processed, _ := h.idempotency.getResult(ctx, key)
if processed != nil {
return nil // Already processed
}
// Record payment success
if err := h.paymentClient.RecordSuccess(ctx, payload.TxID, payload.OrderID); err != nil {
return fmt.Errorf("record payment success: %w", err)
}
// Confirm order (saga continuing from webhook)
if err := h.orderClient.ConfirmOrder(ctx, payload.OrderID, payload.TxID); err != nil {
return fmt.Errorf("confirm order from webhook: %w", err)
}
return nil
}
func (h *PaymentWebhookHandler) handlePaymentFailed(ctx context.Context, payload WebhookPayload) error {
key := fmt.Sprintf("webhook:payment.failed:%s", payload.TxID)
processed, _ := h.idempotency.getResult(ctx, key)
if processed != nil {
return nil
}
// Cancel order — triggers saga compensation
if err := h.orderClient.CancelOrder(ctx, payload.OrderID, "payment_failed"); err != nil {
return fmt.Errorf("cancel order payment failed: %w", err)
}
return nil
}
func (h *PaymentWebhookHandler) handlePaymentExpired(ctx context.Context, payload WebhookPayload) error {
// Similar to failed — release stock and cancel
return h.handlePaymentFailed(ctx, payload)
}5. Transactional Outbox: Reliable Event Publishing
Masalah klasik di microservices: dual-write. Kamu update database (order confirmed) dan publish event ke Kafka/RabbitMQ. Jika publish gagal setelah DB commit — event hilang. Jika publish sukses tapi DB rollback — event phantom.
Solusi: Transactional Outbox Pattern. INSERT event ke outbox table di dalam transaksi database yang sama dengan update order. Kemudian relay worker membaca outbox table dan publish ke message broker.
package checkout
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"log/slog"
"time"
)
// OutboxEvent represents an event to be published.
type OutboxEvent struct {
ID string `json:"id"`
AggregateType string `json:"aggregate_type"`
AggregateID string `json:"aggregate_id"`
EventType string `json:"event_type"`
Payload []byte `json:"payload"`
CreatedAt time.Time `json:"created_at"`
PublishedAt *time.Time `json:"published_at,omitempty"`
}
// OutboxWriter writes events into the outbox table within a DB transaction.
type OutboxWriter struct {
db *sql.DB
logger *slog.Logger
}
// NewOutboxWriter creates the writer.
func NewOutboxWriter(db *sql.DB, logger *slog.Logger) *OutboxWriter {
return &OutboxWriter{
db: db,
logger: logger,
}
}
// WriteBatch inserts multiple outbox events in a single transaction.
func (ow *OutboxWriter) WriteBatch(ctx context.Context, events []OutboxEvent) error {
tx, err := ow.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback() // no-op if committed
stmt, err := tx.PrepareContext(ctx, `
INSERT INTO outbox (id, aggregate_type, aggregate_id, event_type, payload, created_at)
VALUES ($1, $2, $3, $4, $5, $6)
`)
if err != nil {
return fmt.Errorf("prepare: %w", err)
}
defer stmt.Close()
now := time.Now().UTC()
for _, event := range events {
event.CreatedAt = now
if event.ID == "" {
event.ID = generateEventID()
}
_, err := stmt.ExecContext(ctx,
event.ID, event.AggregateType, event.AggregateID,
event.EventType, event.Payload, event.CreatedAt,
)
if err != nil {
return fmt.Errorf("insert outbox event %s: %w", event.EventType, err)
}
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit outbox tx: %w", err)
}
return nil
}
// OutboxRelay reads from outbox table and publishes to message broker.
type OutboxRelay struct {
db *sql.DB
broker MessageBroker
logger *slog.Logger
pollInterval time.Duration
batchSize int
done chan struct{}
}
// MessageBroker interface abstracts the message broker.
type MessageBroker interface {
Publish(ctx context.Context, topic string, key string, payload []byte) error
}
// NewOutboxRelay creates a relay worker.
func NewOutboxRelay(
db *sql.DB,
broker MessageBroker,
logger *slog.Logger,
pollInterval time.Duration,
batchSize int,
) *OutboxRelay {
return &OutboxRelay{
db: db,
broker: broker,
logger: logger,
pollInterval: pollInterval,
batchSize: batchSize,
done: make(chan struct{}),
}
}
// Start begins the polling loop.
func (or *OutboxRelay) Start() {
go func() {
ticker := time.NewTicker(or.pollInterval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
or.processBatch()
case <-or.done:
return
}
}
}()
or.logger.Info("outbox relay started", "poll_interval", or.pollInterval)
}
// Stop stops the relay gracefully.
func (or *OutboxRelay) Stop() {
close(or.done)
}
func (or *OutboxRelay) processBatch() {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
// Read unpublished events (using FOR UPDATE SKIP LOCKED for concurrent safety)
rows, err := or.db.QueryContext(ctx, `
SELECT id, aggregate_type, aggregate_id, event_type, payload, created_at
FROM outbox
WHERE published_at IS NULL
ORDER BY created_at ASC
LIMIT $1
FOR UPDATE SKIP LOCKED
`, or.batchSize)
if err != nil {
or.logger.Error("outbox relay query", "error", err)
return
}
defer rows.Close()
var events []OutboxEvent
for rows.Next() {
var e OutboxEvent
if err := rows.Scan(&e.ID, &e.AggregateType, &e.AggregateID, &e.EventType, &e.Payload, &e.CreatedAt); err != nil {
or.logger.Error("scan outbox row", "error", err)
continue
}
events = append(events, e)
}
for _, event := range events {
topic := mapTopic(event.AggregateType, event.EventType)
if err := or.broker.Publish(ctx, topic, event.AggregateID, event.Payload); err != nil {
or.logger.Error("publish outbox event",
"event_id", event.ID,
"topic", topic,
"error", err,
)
continue
}
// Mark as published
now := time.Now()
if _, err := or.db.ExecContext(ctx,
`UPDATE outbox SET published_at = $1 WHERE id = $2`,
now, event.ID,
); err != nil {
or.logger.Error("mark outbox published", "event_id", event.ID, "error", err)
}
}
}
func mapTopic(aggregateType, eventType string) string {
return fmt.Sprintf("%s.%s", aggregateType, eventType)
}
func generateEventID() string {
return fmt.Sprintf("evt-%d", time.Now().UnixNano())
}FOR UPDATE SKIP LOCKED
FOR UPDATE SKIP LOCKED adalah PostgreSQL feature yang memungkinkan multiple relay worker instances mengambil batch berbeda secara concurrent tanpa saling blocking. Worker A mengambil 10 event pertama, Worker B mengambil 10 berikutnya. Ini penting untuk high-throughput checkout system. Tanpa SKIP LOCKED, worker akan saling menunggu.
6. Service Clients (gRPC)
Berikut adalah contoh client untuk memanggil Inventory dan Order service via gRPC.
package checkout
import (
"context"
"fmt"
"time"
pb "checkout/gen/go/inventory/v1"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
// ReservationResponse from inventory service.
type ReservationResponse struct {
ReservationID string
ExpiresAt time.Time
}
// InventoryClient calls the inventory service.
type InventoryClient struct {
conn *grpc.ClientConn
client pb.InventoryServiceClient
}
// NewInventoryClient creates a gRPC client.
func NewInventoryClient(addr string) (*InventoryClient, error) {
conn, err := grpc.Dial(addr,
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithDefaultCallOptions(grpc.MaxCallRecvMsgSize(4*1024*1024)),
)
if err != nil {
return nil, fmt.Errorf("inventory grpc dial: %w", err)
}
return &InventoryClient{
conn: conn,
client: pb.NewInventoryServiceClient(conn),
}, nil
}
// ReserveStock calls inventory service to reserve items.
func (ic *InventoryClient) ReserveStock(ctx context.Context, orderID string, items []OrderItem) (*ReservationResponse, error) {
req := &pb.ReserveStockRequest{
OrderId: orderID,
Items: convertItems(items),
}
ctx, cancel := context.WithTimeout(ctx, 3*time.Second)
defer cancel()
resp, err := ic.client.ReserveStock(ctx, req)
if err != nil {
return nil, fmt.Errorf("reserve stock: %w", err)
}
return &ReservationResponse{
ReservationID: resp.ReservationId,
ExpiresAt: resp.ExpiresAt.AsTime(),
}, nil
}
// ReleaseStock releases previously reserved stock.
func (ic *InventoryClient) ReleaseStock(ctx context.Context, orderID string, items []OrderItem) error {
req := &pb.ReleaseStockRequest{
OrderId: orderID,
Items: convertItems(items),
}
ctx, cancel := context.WithTimeout(ctx, 3*time.Second)
defer cancel()
_, err := ic.client.ReleaseStock(ctx, req)
return fmt.Errorf("release stock: %w", err)
}
func convertItems(items []OrderItem) []*pb.OrderItem {
result := make([]*pb.OrderItem, len(items))
for i, item := range items {
result[i] = &pb.OrderItem{
ProductId: item.ProductID,
Quantity: int32(item.Quantity),
}
}
return result
}
// OrderClient calls the order service.
type OrderClient struct {
conn *grpc.ClientConn
client pb.OrderServiceClient
}
func convertItems2(items []OrderItem) []*pb.OrderItem {
// same as convertItems
return convertItems(items)
}Architecture Overview
Edge Cases dan Penanganannya
1. Double Charge (Payment Gateway Retry)
Masalah: Payment gateway mengirim webhook payment.success dua kali karena network glitch. Jika tidak ditangani, order di-confirm dua kali, event dipublish dua kali, dan charge terjadi dua kali.
Solusi: Idempotency key pada level transaksi. Di webhook handler, cek transaction_id di Redis sebelum memproses. Juga, ConfirmOrder harus idempoten — jika order sudah CONFIRMED, return sukses tanpa effect.
func (oc *OrderClient) ConfirmOrder(ctx context.Context, orderID, txID string) error {
// Check current status first
order, err := oc.GetOrder(ctx, orderID)
if err != nil {
return err
}
if order.Status == string(OrderStatusConfirmed) {
return nil // Already confirmed, idempotent
}
if order.Status == string(OrderStatusCancelled) {
return fmt.Errorf("cannot confirm cancelled order: %s", orderID)
}
// Proceed with confirmation
// ...
}2. Partial Refund vs Full Refund
Masalah: Checkout gagal di step confirm_order setelah payment berhasil. Compensating action melakukan refund penuh. Tapi beberapa payment method (QRIS, BNPL) tidak support partial refund — refund harus full amount.
Solusi: Saga compensation untuk payment step harus full refund. Simpan payment_status di order service sehingga refund hanya dilakukan sekali (idempotent).
3. Outbox Relay Duplikasi
Masalah: Outbox relay worker crash setelah publish ke Kafka tapi sebelum update published_at. Retry → publish duplikat.
Solusi: Idempotent consumer di sisi penerima. Setiap consumer harus track event ID yang sudah diproses di dedup cache (Redis set dengan TTL 24 jam).
// EventConsumer with dedup.
type EventConsumer struct {
rdb *redis.Client
handler func(ctx context.Context, event OutboxEvent) error
}
func (ec *EventConsumer) Process(ctx context.Context, event OutboxEvent) error {
// Dedup check
processed, err := ec.rdb.SIsMember(ctx, "events:processed", event.ID).Result()
if err == nil && processed {
return nil // Already processed
}
if err := ec.handler(ctx, event); err != nil {
return err
}
// Mark processed (TTL 24h for cleanup)
ec.rdb.SAdd(ctx, "events:processed", event.ID)
ec.rdb.Expire(ctx, "events:processed", 24*time.Hour)
return nil
}4. Payment Timeout / Gateway Down
Masalah: Payment gateway lambat atau down. Orchestrator menunggu response hingga timeout. Jika timeout terjadi, kompensasi harus segera dijalankan — tapi bagaimana kalau payment sebenarnya sukses (gateway slow)?
Solusi: Async reconciliation. Jika charge timeout:
- Kompensasi: release stock, cancel order
- Tandai payment sebagai
pending_reconciliation - Webhook
payment.successyang datang belakangan → deteksi order sudah cancelled, trigger refund
// ReconciliationHandler resolves ambiguous payment states.
type ReconciliationHandler struct {
paymentClient *PaymentClient
orderClient *OrderClient
}
func (rh *ReconciliationHandler) HandleLateWebhook(ctx context.Context, txID, orderID string) error {
order, err := rh.orderClient.GetOrder(ctx, orderID)
if err != nil {
return err
}
if order.Status == string(OrderStatusCancelled) {
// Payment succeeded but order was cancelled due to timeout
// Need to refund the payment
return rh.paymentClient.Refund(ctx, txID, order.TotalAmount)
}
return nil
}5. Concurrent Checkout (Stock Race Condition)
Masalah: Dua user checkout item yang sama secara bersamaan. Inventory service harus handle concurrent reservation correctly.
Solusi: Optimistic locking di inventory database. Jika dua transaksi mencoba reserve stock yang sama, yang kedua akan gagal (conflict) dan harus retry.
-- Inventory: optimistic locking via version column
UPDATE inventory_items
SET reserved_quantity = reserved_quantity + $1,
version = version + 1
WHERE product_id = $2
AND available_quantity - reserved_quantity >= $1
AND version = $3; -- optimistic lock checkfunc (ic *InventoryClient) ReserveStockWithRetry(ctx context.Context, req ReserveRequest, maxRetries int) error {
for attempt := 0; attempt < maxRetries; attempt++ {
_, err := ic.ReserveStock(ctx, req.OrderID, req.Items)
if err == nil {
return nil
}
if !isConflictError(err) {
return err // Non-retryable error
}
// Exponential backoff
time.Sleep(time.Duration(100*(1<<attempt)) * time.Millisecond)
}
return fmt.Errorf("reserve stock: max retries exhausted")
}Key Takeaways
Saga Orchestration > Choreography
Untuk checkout flow yang kompleks dengan banyak branching, orchestration memberikan visibilitas penuh dan kontrol error yang lebih baik. Choreography terlalu implicit untuk bisnis-critical flow.
Idempotency is Non-Negotiable
Setiap endpoint yang memproses payment harus idempotent. Idempotency-Key header di request, caching hasil di Redis, dan cek duplikat di database adalah minimum requirement.
Transactional Outbox Saves You
Dual-write (DB + message broker) adalah source of truth inconsistency. Outbox pattern menjamin atomicity: either both DB update and event publish happen, or neither.
Compensating Actions Must Be Idempotent
Release stock bisa dipanggil dua kali. Cancel order bisa dipanggil dua kali. Pastikan compensating actions safe untuk dijalankan multiple times.
Webhook Signature Verification
Jangan pernah trust webhook tanpa verifikasi. HMAC-SHA256 adalah minimum. Simpan webhook secret di environment variable, bukan di codebase.
Async Reconciliation for Ambiguous States
Payment timeout menciptakan state ambiguous: apakah payment sukses atau tidak? Async reconciliation dengan jadwal periodik (setiap 5 menit) membersihkan state ini.
Penutup
Distributed transaction di Q-Commerce tidak bisa dihindari — setiap checkout melibatkan multiple service dengan database sendiri-sendiri. Saga Pattern bukan pengganti transaksi database, tapi strategi untuk menjaga eventual consistency di distributed system. Kuncinya: step yang bisa di-compensate, idempotency di semua level, dan outbox untuk reliable event publishing.
Kode di atas adalah production-grade template. Untuk production, tambahkan: dead letter queue (DLQ) untuk outbox events yang gagal publish, distributed tracing (OpenTelemetry) untuk debugging saga flow, dan automated chaos testing untuk memvalidasi saga resilience terhadap network partition dan service failure.