Akar Masalah: Batasan Shared Lock Konvensional

Pada distributed worker queue, race condition kerap terjadi saat worker mengeksekusi job kritis multi-tahap (seperti mutasi ledger keuangan atau provisioning infrastruktur). Pola umum menggunakan distributed lock berbasis Redis (seperti Redlock atau single-key TTL) rentan terhadap kegagalan konsistensi ketika terjadi:

  • Stop-the-World GC Pause / Process Pause: Worker mengklaim lock dengan TTL 10 detik. Pada detik ke-9, worker mengalami garbage collection pause atau disk I/O stall selama 15 detik. TTL kedaluwarsa, queue server menganggap worker mati dan menjadwalkan ulang job ke Worker B. Worker A terbangun dan mengeksekusi commit ke database secara bersamaan dengan Worker B.
  • Network Partition (Split-Brain): Worker terisolasi dari cluster queue tetapi masih memiliki akses ke database target. Worker baru ditugaskan untuk job yang sama oleh orchestrator.
  • Asymmetric Network Delays: Heartbeat perpanjangan lease terlambat tiba di coordinator, memicu pelepasan sewa sepihak.

Mekanisme TTL lock standar tidak menjamin eksklusivitas mutual jika client dapat menahan eksekusi melewati durasi TTL tanpa disadari oleh downstream storage.

Prinsip Dasar: Quorum Lease dan Fencing Token

Pola Quorum Lease memecahkan masalah ini dengan memisahkan otoritas eksekusi menjadi dua komponen: validasi kuorum terdistribusi untuk kepemilikan lease, dan fencing token untuk verifikasi mutasi di sisi database target.

1. Verifikasi Kuorum (Raft/Paxos)

Lease tidak boleh diterbitkan oleh single node tanpa verifikasi mayoritas ($N/2 + 1$). State kepemilikan lease dikelola oleh sistem konsensus seperti etcd atau Consul. Worker hanya diizinkan menjalankan job jika lease tersebut masih divalidasi aktif oleh kuorum independen.

2. Monotonic Fencing Token (Epoch Increment)

Sesuai dengan analisis Martin Kleppmann mengenai distributed locking, lock server wajib menghasilkan token integer yang naik secara monotonik (epoch counter) setiap kali lease baru diberikan. Storage target (misalnya PostgreSQL atau MySQL) harus memvalidasi token ini dalam transaksi sebelum commit:

-- Skema tabel verifikasi epoch di database target
UPDATE critical_jobs 
SET status = 'PROCESSING', 
    last_fencing_token = :current_token, 
    updated_at = NOW()
WHERE id = :job_id 
  AND last_fencing_token < :current_token;

Jika Worker A terbangun dari GC pause dengan token 101, sedangkan Worker B telah mengambil alih dengan token 102, maka query update Worker A akan gagal karena kondisi last_fencing_token < 101 bernilai false.

Arsitektur Eksekusi Multi-Tahap

Job multi-tahap harus dipecah menjadi state machine deterministik di mana setiap transisi state diverifikasi terhadap kuorum lease:

  1. Acquire: Worker mengajukan lease ke etcd dengan TTL pendek (misal 5 detik) dan menerima fencing token generasi terbaru.
  2. Keep-Alive Background Loop: Goroutine/thread terpisah memperpanjang lease secara periodik (misal per 1.5 detik). Jika keep-alive gagal mencapai kuorum, sinyal pembatalan (context cancellation) langsung dikirim ke thread pemrosesan utama.
  3. Pre-Commit Check: Sebelum mutasi eksternal dijalankan, worker memverifikasi status lease lokal dan menyertakan token ke payload database.
  4. State Mutation: Database melakukan CAS (Compare-And-Swap) atomic pada kolom fencing token.

Implementasi Teknis: Worker Lease Berbasis etcd

Contoh implementasi teruji menggunakan Go dan primitif concurrency etcd v3:

package main

import (
	"context"
	"database/sql"
	"errors"
	"fmt"
	"log"
	"time"

	clientv3 "go.etcd.io/etcd/client/v3"
	"go.etcd.io/etcd/client/v3/concurrency"
)

type JobExecutor struct {
	etcdCli *clientv3.Client
	db      *sql.DB
}

func (e *JobExecutor) ExecuteCriticalJob(ctx context.Context, jobID string) error {
	// 1. Buat session dengan TTL 5 detik berbasis kuorum etcd
	sess, err := concurrency.NewSession(e.etcdCli, concurrency.WithTTL(5))
	if err != nil {
		return fmt.Errorf("gagal membuat session kuorum: %w", err)
	}
	defer sess.Close()

	// 2. Akuisisi mutex terdistribusi
	lockKey := fmt.Sprintf("/leases/jobs/%s", jobID)
	mutex := concurrency.NewMutex(sess, lockKey)

	if err := mutex.Lock(ctx); err != nil {
		return fmt.Errorf("gagal mengakuisisi lease: %w", err)
	}
	defer mutex.Unlock(context.Background())

	// Header revision etcd berfungsi sebagai monotonic fencing token
	fencingToken := mutex.Header().Revision

	// 3. Jalankan tahap pekerjaan dengan validasi token di database
	return e.processStageWithFencing(ctx, jobID, fencingToken)
}

func (e *JobExecutor) processStageWithFencing(ctx context.Context, jobID string, token int64) error {
	tx, err := e.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted})
	if err != nil {
		return err
	}
	defer tx.Rollback()

	// Atomic CAS validation
	res, err := tx.ExecContext(ctx, `
		UPDATE job_states 
		SET current_epoch = $1, step = step + 1, updated_at = NOW()
		WHERE id = $2 AND current_epoch < $1`,
		token, jobID,
	)
	if err != nil {
		return fmt.Errorf("error executing query: %w", err)
	}


rowsAffected, err := res.RowsAffected()
	if err != nil {
		return err
	}
	if rowsAffected == 0 {
		return errors.New("fencing token ditolak: lease telah kedaluwarsa atau diambil alih")
	}

	// Jalankan side effect di sini jika CAS berhasil
	return tx.Commit()
}

Mitigasi Node Flapping dan Unstable Network

Partisi jaringan intermiten dapat memicu node flapping, di mana worker berulang kali kehilangan dan merebut kembali lease. Mitigasi wajib diterapkan pada level scheduler:

  • Exponential Backoff dengan Jitter: Saat lease terlepas di luar kendali normal, worker dilarang langsung melakukan re-acquire secara agresif. Terapkan full jitter: sleep = min(cap, base * 2 ** attempt) * rand(0.5, 1.5).
  • Lease Renewal Grace Margin: Tetapkan batas renewal di level $1/3$ durasi TTL. Jika TTL 6 detik dan renewal tidak mendapat ack dalam 2 detik, tandai worker dalam kondisi degraded dan hentikan pengambilan job baru.
  • Poison Pill Abort: Gunakan watchdog timer internal di runtime worker. Jika perpanjangan lease gagal mencapai kuorum sebelum 70% durasi TTL habis, worker mematikan prosesnya sendiri (fast-fail/SIGKILL) sebelum memicu race condition.

Strategi Safe Fallback Saat Kuorum Hilang

Bila koneksi ke cluster kuorum terputus di tengah pemrosesan:

  1. Seketika Batalkan Context: Seluruh I/O keluar (HTTP calls, DB queries) harus di-bind ke context.Context yang otomatis membatalkan operasi saat channel sess.Done() menerima sinyal.
  2. Idempotency Verification: Setiap job wajib mengimplementasikan kunci idempotensi deterministik (misalnya hash input payload). Ketika worker baru mengambil alih job yang ditinggalkan, tahap yang sudah di-commit dengan token lama tidak akan dieksekusi ulang.
  3. Dead Letter Queue (DLQ) Routing: Jika sebuah job gagal diselesaikan setelah $N$ kali pergantian epoch (indikasi payload merusak worker/poison message), pindahkan task ke DLQ untuk audit manual.

Observabilitas dan Metrik Utama

Pantau metrik berikut untuk mendeteksi degradasi konsensus sebelum terjadi kegagalan split-brain:

MetrikTipeBatas Kritis / Anomali
lease_acquisition_latency_msHistogramP99 > 50% durasi heartbeat lease. Mengindikasikan disk lag pada etcd cluster.
lease_drift_secondsGaugeWaktu sejak keep-alive terakhir yang berhasil. Nilai > 50% TTL memicu alert peringatan.
quorum_rejection_totalCounterLonjakan nilai mengindikasikan worker stale yang mencoba menulis ke database (fencing aktif menolak commit).
worker_evictions_totalCounterWorker dipaksa terminate oleh internal watchdog karena kehilangan lease.
Catatan Implementasi: Jangan pernah mengandalkan clock synchronization sistem (seperti NTP) untuk validasi lease di lingkungan terdistribusi. Clock skew dapat memicu drift waktu antar-node. Keamanan sistem harus murni bergantung pada monotonic sequence ordering via fencing token dan validasi kuorum Raft.