Pada arsitektur berbasis event, jaminan urutan data (ordering guarantee) absolut jarang ditemukan di boundary eksternal. Layanan pihak ketiga seperti payment gateway atau logistik sering mengirim webhook secara paralel melalui beberapa jalur jaringan. Masalah muncul saat event anak (misalnya payment.captured) tiba sebelum event induk (order.created) selesai ditulis ke database lokal. Relasi foreign key gagal, menghasilkan fenomena event yatim (orphaned event).

Akar Masalah: Tiga Anti-Pattern Penanganan Event Yatim

Menolak atau mengabaikan event anak yang kehilangan referensi induknya memicu masalah operasional serius:

  • Fail-Fast dengan HTTP 5xx: Mengembalikan status error memaksa vendor memicu exponential backoff retry. Jika volume trafik tinggi, teknik ini menciptakan retry storm yang membebani gateway tanpa mempercepat ketersediaan data induk.
  • Drop Silang (HTTP 200 + Discard): Menelan event secara senyap demi menjaga kestabilan sistem menyebabkan silent data loss. State transaksi menjadi inkonsisten secara permanen.
  • Thread Sleep / Polling Blocking: Menahan thread worker HTTP untuk menunggu parent selesai ditulis akan menghabiskan connection pool worker dalam hitungan detik.

Arsitektur Pola lost+found

Pola lost+found meminjam konsep utilitas fsck pada file system Unix: saat file system menemukan blok inode teralokasi tanpa entri direktori induk, inode tersebut diamankan ke direktori darurat /lost+found alih-alih dihapus.

Pada API ingest, payload anak yang belum memiliki induk dialihkan ke holding buffer (quarantine store) berbasis kuncian ID induk (parent_id) dengan Time-to-Live (TTL). Begitu event induk tiba dan berhasil di-commit ke database, sistem mengeksekusi mekanisme rekonsiliasi untuk menarik, memproses, dan membersihkan seluruh event anak yang tertahan di buffer.

Implementasi Rekonsiliasi Otomatis (Go + Redis)

Contoh berikut mendemonstrasikan penanganan event order.payment_received yang tiba mendahului order.created menggunakan Redis List dan Set sebagai quarantine buffer.

package main

import (
	"context"
	"database/sql"
	"encoding/json"
	"errors"
	"fmt"
	"net/http"
	"time"

	"github.com/redis/go-redis/v9"
)

type WebhookEvent struct {
	EventID   string          `json:"event_id"`
	EventType string          `json:"event_type"`
	ParentID  string          `json:"parent_id,omitempty"`
	Payload   json.RawMessage `json:"payload"`
}

type IngestService struct {
	db  *sql.DB
	rdb *redis.Client
}

const quarantineTTL = 24 * time.Hour

func (s *IngestService) HandleWebhook(w http.ResponseWriter, r *http.Request) {
	var ev WebhookEvent
	if err := json.NewDecoder(r.Body).Decode(&ev); err != nil {
		http.Error(w, "invalid payload", http.StatusBadRequest)
		return
	}

	ctx := r.Context()

	switch ev.EventType {
	case "order.created":
		if err := s.processParent(ctx, ev); err != nil {
			http.Error(w, "failed to process parent", http.StatusInternalServerError)
			return
		}
	case "order.payment_received":
		if err := s.processChild(ctx, ev); err != nil {
			http.Error(w, "failed to process child", http.StatusInternalServerError)
			return
		}
	}

	w.WriteHeader(http.StatusAccepted)
}

func (s *IngestService) processParent(ctx context.Context, ev WebhookEvent) error {
	// 1. Simpan entitas induk ke DB
	_, err := s.db.ExecContext(ctx, "INSERT INTO orders (id, data) VALUES ($1, $2) ON CONFLICT DO NOTHING", ev.EventID, string(ev.Payload))
	if err != nil {
		return err
	}

	// 2. Kuras payload anak dari buffer lost+found
	bufferKey := fmt.Sprintf("lost_found:%s", ev.EventID)
	for {
		rawChild, err := s.rdb.RPop(ctx, bufferKey).Bytes()
		if errors.Is(err, redis.Nil) {
			break // Buffer kosong
		}
		if err != nil {
			return err
		}

		var childEvent WebhookEvent
		if err := json.Unmarshal(rawChild, &childEvent); err == nil {
			// Eksekusi logic bisnis event anak
			s.applyChildLogic(ctx, childEvent)
		}
	}
	return nil
}

func (s *IngestService) processChild(ctx context.Context, ev WebhookEvent) error {
	// 1. Cek keberadaan parent di DB
	var exists bool
	err := s.db.QueryRowContext(ctx, "SELECT EXISTS(SELECT 1 FROM orders WHERE id = $1)", ev.ParentID).Scan(&exists)
	if err != nil {
		return err
	}

	// 2. Jika parent sudah ada, langsung eksekusi
	if exists {
		return s.applyChildLogic(ctx, ev)
	}

	// 3. Parent belum ada: karantina payload ke lost+found buffer
	bufferKey := fmt.Sprintf("lost_found:%s", ev.ParentID)
	raw, err := json.Marshal(ev)
	if err != nil {
		return err
	}

	pipe := s.rdb.TxPipeline()
	pipe.LPush(ctx, bufferKey, raw)
	pipe.Expire(ctx, bufferKey, quarantineTTL)
	_, err = pipe.Exec(ctx)
	return err
}

func (s *IngestService) applyChildLogic(ctx context.Context, ev WebhookEvent) error {
	_, err := s.db.ExecContext(ctx, "INSERT INTO payments (id, order_id, data) VALUES ($1, $2, $3) ON CONFLICT DO NOTHING", ev.EventID, ev.ParentID, string(ev.Payload))
	return err
}
ponytail: implementasi di atas mengasumsikan drain buffer satu arah; tambahkan distributed lock (misal: Redlock) jika worker parent beroperasi secara paralel pada partisi data yang sama.

Validasi Idempotensi dan Skema

Buffer karantina rentan terhadap payload duplikat akibat retry transmisi dari vendor. Terapkan dua filter validasi wajib:

  1. Idempotency Layer: Gunakan Redis Set atau kolom DB unik (dedup:{event_id}) sebelum memasukkan item ke holding list. Jika event_id sudah terdaftar di buffer, abaikan payload tambahan.
  2. Contract Validation: Validasi JSON schema sebelum karantina. Jangan menampung payload yang tidak valid secara sintaksis ke dalam lost+found; tolak langsung dengan status HTTP 422 Unprocessable Entity.

Observabilitas: Menangani Event Yatim Kedaluwarsa

Event di buffer karantina tidak boleh dibiarkan hilang saat TTL habis tanpa notifikasi. Ketika parent tidak kunjung tiba dalam batas SLA (misal: 24 jam), event yatim harus dialihkan ke Dead-Letter Queue (DLQ) atau tabel arsip khusus untuk investigasi manual.

Metrik minimal yang wajib dipantau:

  • quarantine_events_stored_total: Counter event anak yang dialihkan ke buffer. Spike pada metrik ini menandakan latensi tinggi pada pipeline event induk.
  • quarantine_events_reconciled_total: Counter event yang berhasil dikuras saat event induk tiba.
  • quarantine_events_expired_total: Counter event yatim yang hangus karena parent tidak pernah tiba. Picu alert P2 jika nilai metrik ini melebihi ambang batas toleransi.

Gunakan Redis Keypspace Notifications (__keyevent@*__:expired) atau background worker terjadwal untuk memindahkan key lost_found:* yang kedaluwarsa ke database relasional jangka panjang dengan status PERMANENTLY_ORPHANED.