Back to Engineering Articles/Real-time Inventory Dark Store: Cegah Oversell di Q-Commerce

Real-time Inventory Dark Store: Cegah Oversell di Q-Commerce

Desain inventory per-hub untuk Astro (instant grocery). Reservation pattern dengan Redis DECRBY atomic + TTL, async sync ke Postgres, stock partitioning untuk hot SKU, reaper job untuk reservation leak, dan cycle counting untuk stok fisik vs sistem. Implementasi Golang lengkap.

Faisal AffanFaisal Affan
6/20/2026
Real-time Inventory Dark Store: Cegah Oversell di Q-Commerce — image 1 of 4
1 / 4

Real-time Inventory Dark Store: Cegah Oversell di Q-Commerce

TL;DR

Oversell di q-commerce beda konsekuensinya dengan e-commerce. Kalau Tokopedia oversell, refund aja. Kalau Astro oversell Indomie, customer nunggu 30 menit — dan gak dapet. Ini bencana operational + reputasi. Artikel ini membahas reservation pattern dengan Redis atomic DECRBY + TTL, async sync ke Postgres, stock partitioning untuk hot SKU, reaper job untuk reservation leak, dan cycle counting. Kode Golang lengkap.


Masalah: Kenapa Inventory Q-Commerce Lebih Sulit dari E-Commerce?

Stok Fisik per Hub

Tiap dark store punya stok sendiri-sendiri. 100 unit Indomie di hub A, 0 di hub B. Sistem harus tahu.

Reservation Window

Dari add-to-cart sampai checkout ada jeda 5-15 menit. Stok harus di-reserve, bukan di-hold forever.

Real-time Required

Stok harus akurat dalam milidetik. Postgres eventual consistency gak cukup.

Hot SKU Contention

Indomie, Telur, Beras — 20% SKU menghasilkan 80% order. Redis key contention tinggi.

Production Horror Story

Sebuah q-commerce startup di Indonesia oversell 400% saat promo "Indomie Goreng 2 pack gratis". Penyebabnya: semua request baca stok dari Postgres replica yang lag 3 detik. 4.000 order diterima, stok cuma 1.000. 3.000 order dibatalkan. Customer trust hancur dalam 10 menit.


Arsitektur Inventory: Redis (Operational) + Postgres (System of Record)


Reservation Lifecycle State Machine


1. Inventory Service: Reservation Pattern dengan Redis Lua

Kenapa Redis bukan Postgres?

The Read/Write Ratio Problem

Di q-commerce, stok dicek 50-100x lebih sering dari pada di-update. Setiap halaman produk, setiap search, setiap cart view — semuanya baca stok. Redis bisa handle 100.000+ reads/s dengan latency < 1ms. Postgres? Di angka yang sama, CPU langsung 100% karena lock contention.

Redis Lua Script untuk Atomic Reserve/Release

package inventory

import (
	"context"
	"fmt"
	"math/rand"
	"sync"
	"time"

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

// Reservation represents a stock reservation with TTL.
type Reservation struct {
	SKU         string    `json:"sku"`
	HubID       string    `json:"hub_id"`
	CartID      string    `json:"cart_id"`
	UserID      string    `json:"user_id"`
	Quantity    int       `json:"quantity"`
	Status      string    `json:"status"` // RESERVED, CONFIRMED, RELEASED, EXPIRED
	CreatedAt   time.Time `json:"created_at"`
	ExpiresAt   time.Time `json:"expires_at"`
}

// InventoryConfig holds configuration for the inventory service.
type InventoryConfig struct {
	DefaultReservationTTL time.Duration
	MaxReservationQty     int
	StockKeyPrefix        string
	ReservationKeyPrefix  string
	LockKeyPrefix         string
	NumStockBuckets       int // For hot key partitioning
}

// DefaultInventoryConfig returns sensible defaults.
func DefaultInventoryConfig() *InventoryConfig {
	return &InventoryConfig{
		DefaultReservationTTL: 15 * time.Minute,
		MaxReservationQty:     50,
		StockKeyPrefix:        "inv:stock",
		ReservationKeyPrefix:  "inv:reservation",
		LockKeyPrefix:         "inv:lock",
		NumStockBuckets:       10,
	}
}

// InventoryReservationService handles stock reservations with Redis atomic operations.
type InventoryReservationService struct {
	config   *InventoryConfig
	rdb      *redis.Client
	luaHash  string
	mu       sync.RWMutex
}

// NewInventoryReservationService creates a new InventoryReservationService.
func NewInventoryReservationService(ctx context.Context, rdb *redis.Client, config *InventoryConfig) (*InventoryReservationService, error) {
	s := &InventoryReservationService{
		config: config,
		rdb:    rdb,
	}

	// Reserve Lua script: atomic check-and-decrement with TTL.
	reserveScript := `
-- KEYS[1]: stock key (e.g., inv:stock:hub:sku)
-- KEYS[2]: reservation set key
-- ARGV[1]: quantity
-- ARGV[2]: reservation ID
-- ARGV[3]: TTL in seconds
-- ARGV[4]: current timestamp
-- ARGV[5]: user ID
-- ARGV[6]: cart ID

local stock = redis.call("GET", KEYS[1])
if not stock then
    return {0, "SKU_NOT_FOUND"}
end

local available = tonumber(stock)
if available < tonumber(ARGV[1]) then
    return {0, "INSUFFICIENT_STOCK", available}
end

-- Atomic decrement
local remaining = redis.call("DECRBY", KEYS[1], tonumber(ARGV[1]))

-- Store reservation with TTL
redis.call("HSET", KEYS[2],
    "reservation_id", ARGV[2],
    "sku", KEYS[1],
    "quantity", ARGV[1],
    "user_id", ARGV[5],
    "cart_id", ARGV[6],
    "status", "RESERVED",
    "created_at", ARGV[4],
    "expires_at", tonumber(ARGV[4]) + tonumber(ARGV[3])
)
redis.call("EXPIRE", KEYS[2], tonumber(ARGV[3]))

return {1, "RESERVED", remaining}
`

	hash, err := rdb.ScriptLoad(ctx, reserveScript).Result()
	if err != nil {
		return nil, fmt.Errorf("inventory: load reserve lua: %w", err)
	}
	s.luaHash = hash
	return s, nil
}

// ReserveStock atomically reserves stock for a SKU.
func (s *InventoryReservationService) ReserveStock(ctx context.Context, hubID, sku string, quantity int, cartID, userID string) (*Reservation, error) {
	if quantity <= 0 || quantity > s.config.MaxReservationQty {
		return nil, fmt.Errorf("inventory: invalid quantity %d (max %d)", quantity, s.config.MaxReservationQty)
	}

	stockKey := s.stockKey(hubID, sku)
	reservationID := s.generateReservationID(cartID, sku)
	reservationKey := fmt.Sprintf("%s:%s", s.config.ReservationKeyPrefix, reservationID)
	ttl := int(s.config.DefaultReservationTTL.Seconds())
	now := time.Now().Unix()

	result, err := s.rdb.EvalSha(ctx, s.luaHash, []string{stockKey, reservationKey},
		quantity, reservationID, ttl, now, userID, cartID,
	).Result()
	if err != nil {
		return nil, fmt.Errorf("inventory: evalsha reserve: %w", err)
	}

	vals, ok := result.([]interface{})
	if !ok || len(vals) < 2 {
		return nil, fmt.Errorf("inventory: unexpected redis response: %v", result)
	}

	success, _ := vals[0].(int64)
	if success == 0 {
		msg, _ := vals[1].(string)
		if len(vals) >= 3 {
			available, _ := vals[2].(int64)
			return nil, &InsufficientStockError{
				SKU:          sku,
				Requested:    quantity,
				Available:    int(available),
			}
		}
		return nil, fmt.Errorf("inventory: reserve failed: %s", msg)
	}

	return &Reservation{
		SKU:       sku,
		HubID:     hubID,
		CartID:    cartID,
		UserID:    userID,
		Quantity:  quantity,
		Status:    "RESERVED",
		CreatedAt: time.Now(),
		ExpiresAt: time.Now().Add(s.config.DefaultReservationTTL),
	}, nil
}

// ReleaseStock releases a reservation (cart abandoned, item removed).
func (s *InventoryReservationService) ReleaseStock(ctx context.Context, reservationID string) error {
	releaseScript := `
-- KEYS[1]: reservation key
-- KEYS[2]: stock key
-- ARGV[1]: current timestamp

local res = redis.call("HGETALL", KEYS[1])
if #res == 0 then
    return {0, "RESERVATION_NOT_FOUND"}
end

-- Extract fields
local quantity = 0
local status = ""
for i = 1, #res, 2 do
    if res[i] == "quantity" then quantity = tonumber(res[i+1]) end
    if res[i] == "status" then status = res[i+1] end
end

if status ~= "RESERVED" then
    return {0, "WRONG_STATUS", status}
end

-- Update status
redis.call("HSET", KEYS[1], "status", "RELEASED", "released_at", ARGV[1])

-- Return stock
redis.call("INCRBY", KEYS[2], quantity)

return {1, "RELEASED", quantity}
`

	result, err := s.rdb.Eval(ctx, releaseScript, []string{
		fmt.Sprintf("%s:%s", s.config.ReservationKeyPrefix, reservationID),
		s.stockKeyFromReservationID(reservationID),
	}, time.Now().Unix()).Result()
	if err != nil {
		return fmt.Errorf("inventory: release stock: %w", err)
	}

	vals, ok := result.([]interface{})
	if !ok || len(vals) < 2 {
		return fmt.Errorf("inventory: unexpected release response: %v", result)
	}

	success, _ := vals[0].(int64)
	if success == 0 {
		msg, _ := vals[1].(string)
		return fmt.Errorf("inventory: release failed: %s", msg)
	}

	return nil
}

// GetStock returns the current available stock for a SKU.
func (s *InventoryReservationService) GetStock(ctx context.Context, hubID, sku string) (int, error) {
	key := s.stockKey(hubID, sku)
	val, err := s.rdb.Get(ctx, key).Int()
	if err != nil {
		if err == redis.Nil {
			return 0, nil
		}
		return 0, fmt.Errorf("inventory: get stock: %w", err)
	}
	return val, nil
}

func (s *InventoryReservationService) stockKey(hubID, sku string) string {
	return fmt.Sprintf("%s:%s:%s", s.config.StockKeyPrefix, hubID, sku)
}

func (s *InventoryReservationService) generateReservationID(cartID, sku string) string {
	return fmt.Sprintf("%s:%s:%d", cartID, sku, time.Now().UnixNano())
}

func (s *InventoryReservationService) stockKeyFromReservationID(reservationID string) string {
	// In production, parse from the reservation hash.
	// For now, derive from the reservation ID format.
	return fmt.Sprintf("%s:%s", s.config.StockKeyPrefix, reservationID)
}

// InsufficientStockError is returned when there is not enough stock.
type InsufficientStockError struct {
	SKU       string
	Requested int
	Available int
}

func (e *InsufficientStockError) Error() string {
	return fmt.Sprintf("inventory: insufficient stock for %s: requested %d, available %d",
		e.SKU, e.Requested, e.Available)
}

// IsInsufficientStock reports whether err is an InsufficientStockError.
func IsInsufficientStock(err error) bool {
	_, ok := err.(*InsufficientStockError)
	return ok
}

Edge Cases

Edge Case: Double Reserve

User add-to-cart item yang sama dua kali. Tanpa idempotency, stok ter-reserve 2x. Solusi: cart-level dedup. Cek apakah SKU sudah di-reserve untuk cart yang sama sebelum reserve.

Edge Case: Race Condition Confirm

User checkout tepat saat TTL reservation expired. Dua hal terjadi: reaper release stok, dan confirm reservation. Solusi: confirm harus cek status reservation di Redis (atomic) sebelum proceed.


2. Stock Sync Worker: Async Redis to Postgres dengan Outbox Pattern

Redis adalah operational store yang cepat. Postgres adalah system of record yang reliable. Keduanya harus sinkron. Tapi sinkron di setiap write? Too slow. Solusinya: outbox pattern.

package inventory

import (
	"context"
	"encoding/json"
	"fmt"
	"log/slog"
	"time"

	"github.com/jackc/pgx/v5/pgxpool"
	"github.com/redis/go-redis/v9"
	"github.com/segmentio/kafka-go"
)

// StockEvent represents a stock change event.
type StockEvent struct {
	EventType     string `json:"event_type"`     // RESERVED, CONFIRMED, RELEASED, EXPIRED
	HubID         string `json:"hub_id"`
	SKU           string `json:"sku"`
	Quantity      int    `json:"quantity"`
	ReservationID string `json:"reservation_id"`
	UserID        string `json:"user_id"`
	CartID        string `json:"cart_id"`
	Timestamp     int64  `json:"timestamp"`
}

// StockSyncWorker asynchronously synchronizes Redis stock state to Postgres.
type StockSyncWorker struct {
	rdb        *redis.Client
	pool       *pgxpool.Pool
	kafkaReader *kafka.Reader
	logger     *slog.Logger
	batchSize  int
}

// NewStockSyncWorker creates a new StockSyncWorker.
func NewStockSyncWorker(
	rdb *redis.Client,
	pool *pgxpool.Pool,
	kafkaReader *kafka.Reader,
	logger *slog.Logger,
) *StockSyncWorker {
	return &StockSyncWorker{
		rdb:        rdb,
		pool:       pool,
		kafkaReader: kafkaReader,
		logger:     logger,
		batchSize:  100,
	}
}

// Run starts the sync worker loop. Blocks until context is cancelled.
func (w *StockSyncWorker) Run(ctx context.Context) error {
	w.logger.InfoContext(ctx, "stock sync worker started")

	for {
		select {
		case <-ctx.Done():
			return ctx.Err()
		default:
			if err := w.processBatch(ctx); err != nil {
				w.logger.ErrorContext(ctx, "sync worker batch failed",
					"error", err,
					"retry_in_seconds", 5,
				)
				time.Sleep(5 * time.Second)
			}
		}
	}
}

// processBatch reads a batch of events from Kafka and syncs to Postgres.
func (w *StockSyncWorker) processBatch(ctx context.Context) error {
	ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
	defer cancel()

	messages := make([]kafka.Message, 0, w.batchSize)

	for i := 0; i < w.batchSize; i++ {
		msg, err := w.kafkaReader.ReadMessage(ctx)
		if err != nil {
			if i > 0 {
				break // Process what we have.
			}
			return fmt.Errorf("sync worker: read message: %w", err)
		}
		messages = append(messages, msg)
	}

	if len(messages) == 0 {
		return nil
	}

	tx, err := w.pool.Begin(ctx)
	if err != nil {
		return fmt.Errorf("sync worker: begin tx: %w", err)
	}
	defer tx.Rollback(ctx)

	for _, msg := range messages {
		var event StockEvent
		if err := json.Unmarshal(msg.Value, &event); err != nil {
			w.logger.ErrorContext(ctx, "sync worker: unmarshal event",
				"error", err,
				"key", string(msg.Key),
			)
			continue
		}

		if err := w.syncEvent(ctx, tx, &event); err != nil {
			return fmt.Errorf("sync worker: sync event %s: %w", event.EventType, err)
		}
	}

	if err := tx.Commit(ctx); err != nil {
		return fmt.Errorf("sync worker: commit tx: %w", err)
	}

	w.logger.InfoContext(ctx, "stock sync batch completed",
		"events_processed", len(messages),
	)

	return nil
}

// syncEvent applies a stock event to Postgres.
func (w *StockSyncWorker) syncEvent(ctx context.Context, tx pgx.Tx, event *StockEvent) error {
	switch event.EventType {
	case "RESERVED", "CONFIRMED":
		return w.applyReservation(ctx, tx, event)
	case "RELEASED", "EXPIRED":
		return w.applyRelease(ctx, tx, event)
	default:
		return fmt.Errorf("unknown event type: %s", event.EventType)
	}
}

func (w *StockSyncWorker) applyReservation(ctx context.Context, tx pgx.Tx, event *StockEvent) error {
	query := `
		INSERT INTO stock_reservations (
			reservation_id, hub_id, sku, user_id, cart_id,
			quantity, status, created_at, expires_at
		) VALUES ($1, $2, $3, $4, $5, $6, $7, NOW(), NOW() + INTERVAL '15 minutes')
		ON CONFLICT (reservation_id) DO NOTHING
	`

	_, err := tx.Exec(ctx, query,
		event.ReservationID,
		event.HubID,
		event.SKU,
		event.UserID,
		event.CartID,
		event.Quantity,
		event.EventType,
	)
	if err != nil {
		return fmt.Errorf("apply reservation: %w", err)
	}

	// Update the inventory table.
	updateQuery := `
		INSERT INTO hub_inventory (hub_id, sku, reserved_quantity, updated_at)
		VALUES ($1, $2, $3, NOW())
		ON CONFLICT (hub_id, sku)
		DO UPDATE SET reserved_quantity = hub_inventory.reserved_quantity + $3,
		              updated_at = NOW()
	`
	_, err = tx.Exec(ctx, updateQuery, event.HubID, event.SKU, event.Quantity)
	if err != nil {
		return fmt.Errorf("apply inventory update: %w", err)
	}

	return nil
}

func (w *StockSyncWorker) applyRelease(ctx context.Context, tx pgx.Tx, event *StockEvent) error {
	query := `
		UPDATE stock_reservations
		SET status = $2, released_at = NOW()
		WHERE reservation_id = $1 AND status IN ('RESERVED', 'CONFIRMED')
	`
	result, err := tx.Exec(ctx, query, event.ReservationID, event.EventType)
	if err != nil {
		return fmt.Errorf("apply release: %w", err)
	}

	if result.RowsAffected() == 0 {
		w.logger.WarnContext(ctx, "reservation not found for release",
			"reservation_id", event.ReservationID,
		)
	}

	// Decrement reserved quantity.
	updateQuery := `
		UPDATE hub_inventory
		SET reserved_quantity = GREATEST(0, reserved_quantity - $3),
		    updated_at = NOW()
		WHERE hub_id = $1 AND sku = $2
	`
	_, err = tx.Exec(ctx, updateQuery, event.HubID, event.SKU, event.Quantity)
	if err != nil {
		return fmt.Errorf("apply inventory release: %w", err)
	}

	return nil
}

// PostgresSchema returns the DDL for the inventory tables.
// Run this during migration.
func PostgresSchema() string {
	return `
	CREATE TABLE IF NOT EXISTS hub_inventory (
		hub_id VARCHAR(50) NOT NULL,
		sku VARCHAR(100) NOT NULL,
		available_quantity INTEGER NOT NULL DEFAULT 0,
		reserved_quantity INTEGER NOT NULL DEFAULT 0,
		pending_quantity INTEGER NOT NULL DEFAULT 0,
		damaged_quantity INTEGER NOT NULL DEFAULT 0,
		updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
		PRIMARY KEY (hub_id, sku)
	);

	CREATE TABLE IF NOT EXISTS stock_reservations (
		reservation_id VARCHAR(200) PRIMARY KEY,
		hub_id VARCHAR(50) NOT NULL,
		sku VARCHAR(100) NOT NULL,
		user_id VARCHAR(100) NOT NULL,
		cart_id VARCHAR(100) NOT NULL,
		quantity INTEGER NOT NULL,
		status VARCHAR(20) NOT NULL DEFAULT 'RESERVED',
		created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
		expires_at TIMESTAMPTZ,
		confirmed_at TIMESTAMPTZ,
		released_at TIMESTAMPTZ,
		cancelled_at TIMESTAMPTZ
	);

	CREATE INDEX idx_reservations_status ON stock_reservations(status);
	CREATE INDEX idx_reservations_user ON stock_reservations(user_id);
	CREATE INDEX idx_reservations_cart ON stock_reservations(cart_id);
	CREATE INDEX idx_hub_inventory_updated ON hub_inventory(updated_at);
	`
}

3. Reaper Job: Cleanup Reservation Leak

Problem

Setiap flash sale atau promo besar, ribuan reservation dibuat. Tapi banyak yang abandoned — user add-to-cart tapi gak checkout. Tanpa cleanup, stok terkunci selamanya.

Solution

Reaper job jalan setiap 30 detik, scan reservation yang expired, release otomatis.

Reaper query Redis keys dengan TTL expired

Release stok dan update status dalam 1 atomic operation

Biar StockSyncWorker sync ke Postgres

Track release rate, leak percentage, average hold time

package inventory

import (
	"context"
	"fmt"
	"log/slog"
	"sync"
	"time"

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

// ReaperConfig holds configuration for the reservation reaper.
type ReaperConfig struct {
	// Interval between reaper runs.
	Interval time.Duration
	// BatchSize is the max number of reservations to process per run.
	BatchSize int
	// GracePeriod is how long after TTL expiry before reaping.
	GracePeriod time.Duration
}

// DefaultReaperConfig returns sensible defaults.
func DefaultReaperConfig() *ReaperConfig {
	return &ReaperConfig{
		Interval:    30 * time.Second,
		BatchSize:   1000,
		GracePeriod: 30 * time.Second,
	}
}

// ReaperMetrics tracks reaper performance.
type ReaperMetrics struct {
	mu              sync.Mutex
	TotalExpired    int64
	TotalReleased   int64
	TotalErrors     int64
	LastRunDuration time.Duration
	LastRunTime     time.Time
}

// Reaper periodically cleans up expired reservations.
type Reaper struct {
	config      *ReaperConfig
	rdb         *redis.Client
	kafkaWriter *kafka.Writer
	logger      *slog.Logger
	metrics     *ReaperMetrics
}

// NewReaper creates a new Reaper.
func NewReaper(config *ReaperConfig, rdb *redis.Client, kafkaWriter *kafka.Writer, logger *slog.Logger) *Reaper {
	return &Reaper{
		config:      config,
		rdb:         rdb,
		kafkaWriter: kafkaWriter,
		logger:      logger,
		metrics:     &ReaperMetrics{},
	}
}

// Run starts the reaper loop. Blocks until context is cancelled.
func (r *Reaper) Run(ctx context.Context) error {
	r.logger.InfoContext(ctx, "reaper started",
		"interval", r.config.Interval,
		"batch_size", r.config.BatchSize,
	)

	ticker := time.NewTicker(r.config.Interval)
	defer ticker.Stop()

	// Run immediately on start.
	if err := r.reap(ctx); err != nil {
		r.logger.ErrorContext(ctx, "initial reap failed", "error", err)
	}

	for {
		select {
		case <-ctx.Done():
			return ctx.Err()
		case <-ticker.C:
			start := time.Now()
			if err := r.reap(ctx); err != nil {
				r.metrics.mu.Lock()
				r.metrics.TotalErrors++
				r.metrics.mu.Unlock()
				r.logger.ErrorContext(ctx, "reap cycle failed", "error", err)
			}
			r.metrics.mu.Lock()
			r.metrics.LastRunDuration = time.Since(start)
			r.metrics.LastRunTime = time.Now()
			r.metrics.mu.Unlock()
		}
	}
}

// reap performs one cycle of expired reservation cleanup.
func (r *Reaper) reap(ctx context.Context) error {
	ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
	defer cancel()

	reapScript := `
-- KEYS[1]: reservation key or pattern
-- ARGV[1]: current timestamp
-- ARGV[2]: grace period (seconds)

-- Find expired reservations
local cursor = "0"
local expired = {}
local pattern = "inv:reservation:*"

repeat
    local scan = redis.call("SCAN", cursor, "MATCH", pattern, "COUNT", 100)
    cursor = scan[1]
    local keys = scan[2]
    
    for _, key in ipairs(keys) do
        local expires_at = redis.call("HGET", key, "expires_at")
        local status = redis.call("HGET", key, "status")
        
        if expires_at and status == "RESERVED" then
            local expiry = tonumber(expires_at)
            if tonumber(ARGV[1]) >= (expiry + tonumber(ARGV[2])) then
                -- This reservation is expired
                local quantity = redis.call("HGET", key, "quantity")
                local sku = redis.call("HGET", key, "sku")
                
                -- Release stock
                redis.call("INCRBY", sku, quantity)
                
                -- Update status
                redis.call("HSET", key, "status", "EXPIRED", "released_at", ARGV[1])
                redis.call("EXPIRE", key, 86400) -- Keep for 24h audit
                
                table.insert(expired, key)
                table.insert(expired, quantity)
            end
        end
    end
until cursor == "0"

return expired
`

	now := time.Now().Unix()
	graceSec := int(r.config.GracePeriod.Seconds())

	result, err := r.rdb.Eval(ctx, reapScript, []string{}, now, graceSec).Result()
	if err != nil {
		return fmt.Errorf("reaper: eval reap script: %w", err)
	}

	expiredList, ok := result.([]interface{})
	if !ok {
		return nil
	}

	releasedCount := len(expiredList) / 2

	if releasedCount == 0 {
		return nil
	}

	// Publish release events to Kafka for Postgres sync.
	for i := 0; i < len(expiredList); i += 2 {
		reservationKey, _ := expiredList[i].(string)
		quantity, _ := expiredList[i+1].(int64)

		event := StockEvent{
			EventType:     "EXPIRED",
			ReservationID: extractReservationID(reservationKey),
			Quantity:      int(quantity),
			Timestamp:     now,
		}

		eventBytes, _ := json.Marshal(event)
		err := r.kafkaWriter.WriteMessages(ctx, kafka.Message{
			Key:   []byte(event.ReservationID),
			Value: eventBytes,
		})
		if err != nil {
			r.logger.ErrorContext(ctx, "reaper: publish event failed",
				"error", err,
				"reservation", event.ReservationID,
			)
		}
	}

	r.metrics.mu.Lock()
	r.metrics.TotalExpired += int64(len(expiredList))
	r.metrics.TotalReleased += int64(releasedCount)
	r.metrics.mu.Unlock()

	r.logger.InfoContext(ctx, "reaper cycle completed",
		"reservations_released", releasedCount,
	)

	return nil
}

// Metrics returns a copy of the current reaper metrics.
func (r *Reaper) Metrics() ReaperMetrics {
	r.metrics.mu.Lock()
	defer r.metrics.mu.Unlock()
	return ReaperMetrics{
		TotalExpired:    r.metrics.TotalExpired,
		TotalReleased:   r.metrics.TotalReleased,
		TotalErrors:     r.metrics.TotalErrors,
		LastRunDuration: r.metrics.LastRunDuration,
		LastRunTime:     r.metrics.LastRunTime,
	}
}

func extractReservationID(key string) string {
	// Format: inv:reservation:<cart_id>:<sku>:<timestamp>
	prefix := "inv:reservation:"
	if len(key) > len(prefix) {
		return key[len(prefix):]
	}
	return key
}

4. Hot Key Mitigation: Stock Partitioning

Problem

SKU populer ("Indomie Goreng", "Telur 1kg", "Beras 5kg") adalah hot key. 20% SKU menghasilkan 80% traffic. Semua request tulis ke Redis key yang sama menyebabkan contention.

Solution: Random Bucket Selection

package inventory

import (
	"context"
	"fmt"
	"math/rand"
	"time"
)

// StockPartitioner distributes hot SKU stock across multiple Redis keys.
type StockPartitioner struct {
	numBuckets int
	prefix     string
}

// NewStockPartitioner creates a new StockPartitioner.
func NewStockPartitioner(numBuckets int) *StockPartitioner {
	if numBuckets <= 0 {
		numBuckets = 10
	}
	return &StockPartitioner{
		numBuckets: numBuckets,
		prefix:     "inv:partition",
	}
}

// PartitionKey returns a specific bucket key for a hub/SKU combination.
func (p *StockPartitioner) PartitionKey(hubID, sku string, bucketIdx int) string {
	return fmt.Sprintf("%s:%s:%s:%d", p.prefix, hubID, sku, bucketIdx)
}

// SelectBucket deterministically selects a bucket for a hub/SKU, with random
// fallback to spread load across buckets.
func (p *StockPartitioner) SelectBucket(hubID, sku string, useRandom bool) int {
	if useRandom {
		return rand.Intn(p.numBuckets)
	}

	// Deterministic hash for consistent bucket selection per SKU.
	// This reduces cache misses since same SKU tends to hit same bucket.
	hash := hashString(hubID + ":" + sku)
	return hash % p.numBuckets
}

// PartitionedReserve reserves stock across partitioned buckets with fallback.
func (p *StockPartitioner) PartitionedReserve(
	svc *InventoryReservationService,
	ctx context.Context,
	hubID, sku string,
	quantity int,
	cartID, userID string,
) (*Reservation, error) {
	// Try primary bucket first, then fallback to random.
	for attempt := 0; attempt < p.numBuckets; attempt++ {
		bucketIdx := p.SelectBucket(sku, hubID, attempt > 0)
		partitionedSKU := fmt.Sprintf("%s:%d", sku, bucketIdx)

		reservation, err := svc.ReserveStock(ctx, hubID, partitionedSKU, quantity, cartID, userID)
		if err == nil {
			return reservation, nil
		}

		if !IsInsufficientStock(err) {
			return nil, err // Real error, not stock issue.
		}

		// Try next bucket.
		continue
	}

	return nil, &InsufficientStockError{
		SKU:       sku,
		Requested: quantity,
		Available: 0,
	}
}

// InitializePartitionedStock initializes partitioned stock for a hub/SKU.
func (p *StockPartitioner) InitializePartitionedStock(
	ctx context.Context,
	rdb *redis.Client,
	hubID, sku string,
	totalStock int,
	ttl time.Duration,
) error {
	stockPerBucket := totalStock / p.numBuckets
	remainder := totalStock % p.numBuckets

	pipe := rdb.Pipeline()

	for i := 0; i < p.numBuckets; i++ {
		bucketStock := stockPerBucket
		if i < remainder {
			bucketStock++
		}
		key := p.PartitionKey(hubID, sku, i)
		pipe.Set(ctx, key, bucketStock, ttl)
	}

	_, err := pipe.Exec(ctx)
	if err != nil {
		return fmt.Errorf("stock partitioner: init stock: %w", err)
	}
	return nil
}

func hashString(s string) int {
	h := 0
	for _, c := range s {
		h = h*31 + int(c)
	}
	if h < 0 {
		h = -h
	}
	return h
}

5. Cycle Counting: Stok Fisik vs Sistem

Problem

Stok di sistem dan stok di dark store sering berbeda. Picker ambil barang dari rak → stok sistem gak update. Barang rusak → gak dicatat. Stok hilang (shrinkage) → sistem tahu-nya masih ada.

Solution: Cycle Counting dengan Delta Detection

package inventory

import (
	"context"
	"fmt"
	"log/slog"
	"time"

	"github.com/jackc/pgx/v5/pgxpool"
)

// CycleCountService handles physical vs system stock reconciliation.
type CycleCountService struct {
	pool   *pgxpool.Pool
	logger *slog.Logger
}

// NewCycleCountService creates a new CycleCountService.
func NewCycleCountService(pool *pgxpool.Pool, logger *slog.Logger) *CycleCountService {
	return &CycleCountService{
		pool:   pool,
		logger: logger,
	}
}

// CountResult represents the result of a cycle count for a single SKU.
type CountResult struct {
	HubID            string
	SKU              string
	PhysicalStock    int
	SystemStock      int
	ReservedStock    int
	AvailableStock   int // system=available - reserved
	Delta            int
	DeltaPercent     float64
	CountedAt        time.Time
	NeedsAdjustment  bool
}

// RecordCycleCount records a physical cycle count and detects discrepancies.
func (s *CycleCountService) RecordCycleCount(ctx context.Context, hubID, sku string, physicalCount int, countedBy string) (*CountResult, error) {
	// Get current system stock from Postgres.
	var systemAvailable, systemReserved int
	err := s.pool.QueryRow(ctx, `
		SELECT available_quantity, reserved_quantity
		FROM hub_inventory
		WHERE hub_id = $1 AND sku = $2
	`, hubID, sku).Scan(&systemAvailable, &systemReserved)
	if err != nil {
		return nil, fmt.Errorf("cycle count: get system stock: %w", err)
	}

	effectiveSystem := systemAvailable - systemReserved
	delta := physicalCount - effectiveSystem
	deltaPct := 0.0
	if effectiveSystem > 0 {
		deltaPct = float64(delta) / float64(effectiveSystem) * 100
	}

	result := &CountResult{
		HubID:           hubID,
		SKU:             sku,
		PhysicalStock:   physicalCount,
		SystemStock:     systemAvailable,
		ReservedStock:   systemReserved,
		AvailableStock:  effectiveSystem,
		Delta:           delta,
		DeltaPercent:    deltaPct,
		CountedAt:       time.Now(),
		NeedsAdjustment: delta != 0,
	}

	// Record the count in the audit table.
	_, err = s.pool.Exec(ctx, `
		INSERT INTO cycle_counts (
			hub_id, sku, physical_quantity, system_quantity,
			delta, counted_by, counted_at
		) VALUES ($1, $2, $3, $4, $5, $6, NOW())
	`, hubID, sku, physicalCount, effectiveSystem, delta, countedBy)
	if err != nil {
		return nil, fmt.Errorf("cycle count: insert record: %w", err)
	}

	// If delta exceeds threshold, flag for adjustment.
	if deltaPct > 5.0 || deltaPct < -5.0 {
		_, err = s.pool.Exec(ctx, `
			INSERT INTO inventory_adjustments (
				hub_id, sku, previous_quantity, adjusted_quantity,
				delta, reason, status, created_at
			) VALUES ($1, $2, $3, $4, $5, $6, 'PENDING', NOW())
		`, hubID, sku, effectiveSystem, physicalCount, delta, "cycle_count_discrepancy")
		if err != nil {
			return nil, fmt.Errorf("cycle count: create adjustment: %w", err)
		}

		s.logger.InfoContext(ctx, "cycle count discrepancy flagged",
			"hub_id", hubID,
			"sku", sku,
			"delta", delta,
			"delta_pct", deltaPct,
		)
	}

	return result, nil
}

// GetDiscrepancyReport returns all SKUs with significant discrepancies.
func (s *CycleCountService) GetDiscrepancyReport(ctx context.Context, hubID string, thresholdPct float64) ([]*CountResult, error) {
	rows, err := s.pool.Query(ctx, `
		SELECT 
			cc.hub_id, cc.sku, cc.physical_quantity, 
			hi.available_quantity, hi.reserved_quantity,
			cc.delta, cc.counted_at
		FROM cycle_counts cc
		JOIN hub_inventory hi ON hi.hub_id = cc.hub_id AND hi.sku = cc.sku
		WHERE cc.hub_id = $1
		  AND ABS(cc.delta::float / NULLIF(hi.available_quantity - hi.reserved_quantity, 0)) * 100 > $2
		  AND cc.counted_at > NOW() - INTERVAL '24 hours'
		ORDER BY ABS(cc.delta) DESC
	`, hubID, thresholdPct)
	if err != nil {
		return nil, fmt.Errorf("cycle count: get discrepancy report: %w", err)
	}
	defer rows.Close()

	var results []*CountResult
	for rows.Next() {
		var r CountResult
		err := rows.Scan(&r.HubID, &r.SKU, &r.PhysicalStock, &r.SystemStock,
			&r.ReservedStock, &r.Delta, &r.CountedAt)
		if err != nil {
			return nil, fmt.Errorf("cycle count: scan result: %w", err)
		}
		r.AvailableStock = r.SystemStock - r.ReservedStock
		if r.AvailableStock > 0 {
			r.DeltaPercent = float64(r.Delta) / float64(r.AvailableStock) * 100
		}
		r.NeedsAdjustment = true
		results = append(results, &r)
	}

	return results, nil
}

// CycleCountSchema returns DDL for cycle counting tables.
func CycleCountSchema() string {
	return `
	CREATE TABLE IF NOT EXISTS cycle_counts (
		id BIGSERIAL PRIMARY KEY,
		hub_id VARCHAR(50) NOT NULL,
		sku VARCHAR(100) NOT NULL,
		physical_quantity INTEGER NOT NULL,
		system_quantity INTEGER NOT NULL,
		delta INTEGER NOT NULL,
		counted_by VARCHAR(100) NOT NULL,
		counted_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
		notes TEXT
	);

	CREATE INDEX idx_cycle_counts_hub ON cycle_counts(hub_id, counted_at);

	CREATE TABLE IF NOT EXISTS inventory_adjustments (
		id BIGSERIAL PRIMARY KEY,
		hub_id VARCHAR(50) NOT NULL,
		sku VARCHAR(100) NOT NULL,
		previous_quantity INTEGER NOT NULL,
		adjusted_quantity INTEGER NOT NULL,
		delta INTEGER NOT NULL,
		reason VARCHAR(200) NOT NULL,
		status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
		approved_by VARCHAR(100),
		created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
		approved_at TIMESTAMPTZ
	);

	CREATE INDEX idx_adjustments_status ON inventory_adjustments(status);
	`
}

Orchestrasi: Inventory Service

package inventory

import (
	"context"
	"encoding/json"
	"fmt"
	"log/slog"
	"net/http"
	"time"

	"github.com/jackc/pgx/v5/pgxpool"
	"github.com/redis/go-redis/v9"
	"github.com/segmentio/kafka-go"
)

// InventoryService orchestrates all inventory operations.
type InventoryService struct {
	reservation *InventoryReservationService
	partitioner *StockPartitioner
	reaper      *Reaper
	syncWorker  *StockSyncWorker
	cycleCount  *CycleCountService
	rdb         *redis.Client
	pool        *pgxpool.Pool
	kafkaWriter *kafka.Writer
	logger      *slog.Logger
}

// NewInventoryService creates a fully wired inventory service.
func NewInventoryService(
	ctx context.Context,
	rdb *redis.Client,
	pool *pgxpool.Pool,
	logger *slog.Logger,
) (*InventoryService, error) {
	config := DefaultInventoryConfig()

	reservation, err := NewInventoryReservationService(ctx, rdb, config)
	if err != nil {
		return nil, fmt.Errorf("inventory: init reservation: %w", err)
	}

	kafkaWriter := &kafka.Writer{
		Addr:     kafka.TCP("localhost:9092"),
		Topic:    "inventory.events",
		Balancer: &kafka.Hash{},
		Async:    true,
	}

	kafkaReader := kafka.NewReader(kafka.ReaderConfig{
		Brokers:   []string{"localhost:9092"},
		Topic:     "inventory.events",
		GroupID:   "inventory-sync-worker",
		MinBytes:  10,
		MaxBytes:  10e6,
		MaxWait:   1 * time.Second,
	})

	reaper := NewReaper(DefaultReaperConfig(), rdb, kafkaWriter, logger)
	syncWorker := NewStockSyncWorker(rdb, pool, kafkaReader, logger)
	cycleCount := NewCycleCountService(pool, logger)

	return &InventoryService{
		reservation: reservation,
		partitioner: NewStockPartitioner(10),
		reaper:      reaper,
		syncWorker:  syncWorker,
		cycleCount:  cycleCount,
		rdb:         rdb,
		pool:        pool,
		kafkaWriter: kafkaWriter,
		logger:      logger,
	}, nil
}

// Start starts all background workers.
func (s *InventoryService) Start(ctx context.Context) {
	go func() {
		if err := s.syncWorker.Run(ctx); err != nil && err != context.Canceled {
			s.logger.ErrorContext(ctx, "sync worker exited", "error", err)
		}
	}()

	go func() {
		if err := s.reaper.Run(ctx); err != nil && err != context.Canceled {
			s.logger.ErrorContext(ctx, "reaper exited", "error", err)
		}
	}()
}

// ReserveForCart handles cart-level reservation with hot key mitigation.
func (s *InventoryService) ReserveForCart(ctx context.Context, hubID string, items []CartItem, cartID, userID string) ([]*Reservation, error) {
	reservations := make([]*Reservation, 0, len(items))

	for _, item := range items {
		reservation, err := s.partitioner.PartitionedReserve(
			s.reservation, ctx, hubID, item.SKU, item.Quantity, cartID, userID,
		)
		if err != nil {
			// Rollback previous reservations on failure.
			for _, r := range reservations {
				if releaseErr := s.reservation.ReleaseStock(ctx, fmt.Sprintf("%s:%s:%d", cartID, r.SKU, r.CreatedAt.UnixNano())); releaseErr != nil {
					s.logger.ErrorContext(ctx, "rollback reservation failed",
						"error", releaseErr,
						"sku", r.SKU,
					)
				}
			}
			return nil, fmt.Errorf("inventory: reserve cart item %s: %w", item.SKU, err)
		}
		reservations = append(reservations, reservation)
	}

	return reservations, nil
}

// CartItem represents an item in a shopping cart.
type CartItem struct {
	SKU      string `json:"sku"`
	Quantity int    `json:"quantity"`
}

// HTTPHandler returns HTTP handlers for the inventory service.
func (s *InventoryService) HTTPHandler() http.Handler {
	mux := http.NewServeMux()

	mux.HandleFunc("GET /inventory/stock/{hub_id}/{sku}", func(w http.ResponseWriter, r *http.Request) {
		hubID := r.PathValue("hub_id")
		sku := r.PathValue("sku")

		stock, err := s.reservation.GetStock(r.Context(), hubID, sku)
		if err != nil {
			http.Error(w, `{"error":"internal_error"}`, http.StatusInternalServerError)
			return
		}

		json.NewEncoder(w).Encode(map[string]interface{}{
			"hub_id": hubID,
			"sku":    sku,
			"stock":  stock,
		})
	})

	mux.HandleFunc("POST /inventory/reserve", func(w http.ResponseWriter, r *http.Request) {
		var req struct {
			HubID    string     `json:"hub_id"`
			CartID   string     `json:"cart_id"`
			UserID   string     `json:"user_id"`
			Items    []CartItem `json:"items"`
		}
		if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
			http.Error(w, `{"error":"invalid_request"}`, http.StatusBadRequest)
			return
		}

		reservations, err := s.ReserveForCart(r.Context(), req.HubID, req.Items, req.CartID, req.UserID)
		if err != nil {
			http.Error(w, `{"error":"`+err.Error()+`"}`, http.StatusConflict)
			return
		}

		json.NewEncoder(w).Encode(map[string]interface{}{
			"status":       "reserved",
			"reservations": reservations,
		})
	})

	mux.HandleFunc("POST /inventory/release/{reservation_id}", func(w http.ResponseWriter, r *http.Request) {
		reservationID := r.PathValue("reservation_id")
		if err := s.reservation.ReleaseStock(r.Context(), reservationID); err != nil {
			http.Error(w, `{"error":"`+err.Error()+`"}`, http.StatusBadRequest)
			return
		}
		json.NewEncoder(w).Encode(map[string]string{"status": "released"})
	})

	mux.HandleFunc("POST /inventory/cycle-count", func(w http.ResponseWriter, r *http.Request) {
		var req struct {
			HubID       string `json:"hub_id"`
			SKU         string `json:"sku"`
			PhysicalQty int    `json:"physical_quantity"`
			CountedBy   string `json:"counted_by"`
		}
		if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
			http.Error(w, `{"error":"invalid_request"}`, http.StatusBadRequest)
			return
		}

		result, err := s.cycleCount.RecordCycleCount(r.Context(), req.HubID, req.SKU, req.PhysicalQty, req.CountedBy)
		if err != nil {
			http.Error(w, `{"error":"`+err.Error()+`"}`, http.StatusInternalServerError)
			return
		}

		json.NewEncoder(w).Encode(result)
	})

	return mux
}

Key Takeaways

Redis Atomic Operations

Gunakan Redis Lua script untuk check-and-decrement atomic. Single key contention dihindari dengan stock partitioning.

TTL-based Auto Release

Setiap reservation punya TTL 15 menit. Reaper job otomatis cleanup expired reservation tiap 30 detik.

Async Sync ke Postgres

Outbox pattern via Kafka. Redis untuk operational (ms-latency), Postgres untuk system of record (reliable).

Stock Partitioning

Hot SKU di-partition ke 10 bucket random. Kurangi contention di key yang sama.

Cycle Counting

Stok fisik vs sistem di-reconcile via cycle counting. Delta >5% auto-flag buat adjustment.

Idempotency & Rollback

Cart-level idempotency cegah double reserve. Partial failure trigger rollback semua reservation.

Bottom Line

Inventory q-commerce bukan soal "stok berapa" — tapi soal "stok ini sudah di-reserve siapa, berapa lama, dan apa yang terjadi kalau gak jadi dibeli?" Kombinasi Redis (fast + atomic) + Postgres (reliable + queryable) + Reaper (auto-cleanup) + Cycle Count (detect discrepancy) adalah fondasi yang sudah terbukti di Astro, Instacart, Gorillas, dan Getir.


Related Engineering & Tech Articles

Real-time Inventory Dark Store: Cegah Oversell di Q-Commerce | Faisal Affan