Back to Engineering Articles/Checkout & Payment Q-Commerce: Saga Pattern untuk Konsistensi Terdistribusi

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.

Faisal AffanFaisal Affan
6/20/2026

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:

  1. Reserve stock di Inventory Service
  2. Create order dengan status PENDING di Order Service
  3. Charge payment via Payment Gateway (PG)
  4. Confirm order (PENDING → CONFIRMED)
  5. 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:

  1. Verifikasi signature (HMAC-SHA256)
  2. Cek idempotency (jika webhook duplikat)
  3. Update order status
  4. 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:

  1. Kompensasi: release stock, cancel order
  2. Tandai payment sebagai pending_reconciliation
  3. Webhook payment.success yang 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 check
func (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.


Related Engineering & Tech Articles

Checkout & Payment Q-Commerce: Saga Pattern untuk Konsist... | Faisal Affan