Back to Engineering Articles/Order Dispatch & Driver Allocation: Hungarian Algorithm untuk Q-Commerce

Order Dispatch & Driver Allocation: Hungarian Algorithm untuk Q-Commerce

System design order dispatch untuk quick-commerce (q-commerce). Pelajari batching window untuk kumpulkan order, Hungarian/greedy matching driver ke order, Redis GEOADD/GEOSEARCH untuk cari driver terdekat, real-time driver state tracking (idle/picking/delivering), batch delivery multi-order, dan edge cases seperti driver reject, order timeout, serta starvation handling. Dilengkapi implementasi Golang dengan context.Context, error wrapping, struct-based services, channel-based event sourcing, dan concurrent pipeline.

Faisal AffanFaisal Affan
6/20/2026

Order Dispatch & Driver Allocation: Hungarian Algorithm untuk Q-Commerce

"Di quick-commerce, setiap detik berarti. Order yang tidak ter-assign dalam 30 detik berpotensi loss — pelanggan pergi, driver frustrasi, dan SLA hancur."

TL;DR

Sistem dispatch adalah jantung Q-Commerce. Ia harus memutuskan siapa yang mengambil order mana dalam waktu nyaris real-time dengan mempertimbangkan jarak, beban driver, ETA, dan prioritas order. Artikel ini membahas implementasi lengkap dispatch engine di Go: batching window, greedy assignment scoring, driver state machine, Redis geospatial indexing, failover mechanism, dan event sourcing untuk audit trail.


Masalah: Kenapa Dispatch Q-Commerce Sulit?

Bayangkan ini: jam makan siang, 50 order masuk dalam 2 menit. 20 driver online tersebar di berbagai lokasi. Beberapa driver sedang dalam perjalanan mengantar, beberapa baru selesai, yang lain masih menunggu di hub.

Apa yang membuat dispatch Q-Commerce berbeda dari dispatch biasa?

  1. Time-pressure ekstrem: Janji <30 menit dari order ke delivery. Setiap detik di dispatch pipeline adalah biaya.
  2. Multi-order batches: Satu driver bisa bawa 3-5 order sekali jalan — rutenya harus optimal.
  3. Dynamic driver state: Driver bisa reject, cancel, atau tiba-tiba offline. Sistem harus reassign dalam hitungan detik.
  4. Geospatial constraint: Driver hanya efisien dalam radius tertentu dari hub dan rute pengiriman.
  5. Starvation: Driver yang jauh jangan terus-terusan kalah oleh driver yang dekat — mereka perlu "dihangatkan" dengan prioritas.

Core Trade-off

Optimal vs Cepat. Algoritma assignment optimal (Hungarian) punya kompleksitas O(n^3). Untuk 50 order x 20 driver, perlu ~50.000 iterasi — cepat. Tapi untuk 500 order x 200 driver, Hungarian O(100.000^3) tidak feasible. Trade-off antara akurasi matching dan latency dispatch adalah keputusan arsitektural paling kritis di dispatch engine.


Arsitektur Dispatch Engine


System Components

1. Batch Collector: Kumpulin Order, Baru Proses

Tidak semua order bisa langsung di-dispatch satu per satu. Bayangkan order pertama di-assign ke driver A, lalu 3 detik kemudian muncul order yang searah dengan driver A — sayang, sudah terlanjur. Batching mengumpulkan order dalam window waktu tertentu, lalu memproses semuanya sekaligus untuk assignment yang lebih optimal.

package dispatch

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

// BatchConfig controls batch collection behavior.
type BatchConfig struct {
	MaxWindow time.Duration // max wait before flushing (e.g., 5 seconds)
	MinOrders int           // flush immediately if this many orders collected
	MaxOrders int           // hard cap per batch
}

// PendingOrder represents an order waiting to be dispatched.
type PendingOrder struct {
	OrderID   string
	HubID     string
	Lat, Lng  float64
	CreatedAt time.Time
	Priority  int // higher = more urgent
}

// BatchCollector accumulates orders and triggers dispatch.
type BatchCollector struct {
	cfg    BatchConfig
	mu     sync.Mutex
	batch  []PendingOrder
	flush  chan []PendingOrder
	done   chan struct{}
}

// NewBatchCollector creates and starts the collector loop.
func NewBatchCollector(cfg BatchConfig) *BatchCollector {
	bc := &BatchCollector{
		cfg:   cfg,
		batch: make([]PendingOrder, 0, cfg.MaxOrders),
		flush: make(chan []PendingOrder, 1),
		done:  make(chan struct{}),
	}
	go bc.loop()
	return bc
}

// Enqueue adds an order to the current batch.
func (bc *BatchCollector) Enqueue(ctx context.Context, order PendingOrder) error {
	bc.mu.Lock()
	defer bc.mu.Unlock()

	if len(bc.batch) >= bc.cfg.MaxOrders {
		return ErrBatchFull
	}
	bc.batch = append(bc.batch, order)

	if len(bc.batch) >= bc.cfg.MinOrders {
		select {
		case bc.flush <- bc.takeBatch():
		default:
			// flush already queued
		}
	}
	return nil
}

// FlushChan returns a channel that receives batches ready for dispatch.
func (bc *BatchCollector) FlushChan() <-chan []PendingOrder {
	return bc.flush
}

// Stop gracefully shuts down the collector.
func (bc *BatchCollector) Stop() {
	close(bc.done)
}

func (bc *BatchCollector) loop() {
	ticker := time.NewTicker(bc.cfg.MaxWindow)
	defer ticker.Stop()

	for {
		select {
		case <-ticker.C:
			bc.mu.Lock()
			if len(bc.batch) > 0 {
				bc.flush <- bc.takeBatch()
			}
			bc.mu.Unlock()
		case <-bc.done:
			return
		}
	}
}

func (bc *BatchCollector) takeBatch() []PendingOrder {
	batch := bc.batch
	bc.batch = make([]PendingOrder, 0, bc.cfg.MaxOrders)
	return batch
}

var ErrBatchFull = fmt.Errorf("batch collector: queue full")

Mengapa Batching Penting?

Tanpa batching, sistem dispatch riskan terhadap sub-optimal assignment. Contoh: Order A masuk, langsung di-assign ke Driver X yang sedang dekat. 2 detik kemudian Order B masuk — ternyata lokasi Order B sangat dekat dengan Order A, dan Driver Y yang baru selesai delivery bisa ambil keduanya. Dengan batching, kita lihat A dan B bersamaan, assign ke Y, hemat 1 trip.


2. Redis Geospatial: Cari Driver Terdekat

Redis menyediakan perintah GEOADD, GEOSEARCH, dan GEODIST yang sempurna untuk query geospatial real-time. Setiap driver mengirim lokasi setiap beberapa detik (heartbeat), kita update di Redis, lalu dispatch tinggal query.

package dispatch

import (
	"context"
	"fmt"
	"time"

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

const (
	driverLocationKey = "driver:locations"
	driverStateKey    = "driver:%s:state"
	driverLoadKey     = "driver:%s:load"
)

// DriverLocation represents a driver's current position and capacity.
type DriverLocation struct {
	DriverID   string  `json:"driver_id"`
	Lat        float64 `json:"lat"`
	Lng        float64 `json:"lng"`
	State      string  `json:"state"` // idle, to_hub, picking, delivering
	OrderCount int     `json:"order_count"` // current assigned orders
	MaxOrders  int     `json:"max_orders"`  // capacity
}

// GeoRepository wraps Redis geo operations.
type GeoRepository struct {
	rdb *redis.Client
}

// NewGeoRepository creates a GeoRepository.
func NewGeoRepository(rdb *redis.Client) *GeoRepository {
	return &GeoRepository{rdb: rdb}
}

// UpdateDriverLocation stores driver geo + state atomically.
func (r *GeoRepository) UpdateDriverLocation(ctx context.Context, d DriverLocation) error {
	pipe := r.rdb.Pipeline()

	pipe.GeoAdd(ctx, driverLocationKey, &redis.GeoLocation{
		Name:      d.DriverID,
		Latitude:  d.Lat,
		Longitude: d.Lng,
	})

	pipe.Set(ctx, fmt.Sprintf(driverStateKey, d.DriverID), d.State, 30*time.Second)
	pipe.Set(ctx, fmt.Sprintf(driverLoadKey, d.DriverID), d.OrderCount, 30*time.Second)

	_, err := pipe.Exec(ctx)
	return fmt.Errorf("update driver location: %w", err)
}

// FindIdleDriversNearby returns drivers in IDLE state within radius (meters).
func (r *GeoRepository) FindIdleDriversNearby(ctx context.Context, lat, lng float64, radiusMeters float64) ([]DriverLocation, error) {
	res, err := r.rdb.GeoSearch(ctx, driverLocationKey, &redis.GeoSearchQuery{
		Longitude:     lng,
		Latitude:      lat,
		Radius:        radiusMeters,
		RadiusUnit:    "m",
		Sort:          "ASC",
		Count:         50,
		CountAny:      false,
	}).Result()
	if err != nil {
		return nil, fmt.Errorf("geo search: %w", err)
	}

	var drivers []DriverLocation
	for _, name := range res {
		state, err := r.rdb.Get(ctx, fmt.Sprintf(driverStateKey, name)).Result()
		if err != nil || state != "idle" {
			continue
		}
		loadStr, _ := r.rdb.Get(ctx, fmt.Sprintf(driverLoadKey, name)).Result()
		load := 0
		fmt.Sscanf(loadStr, "%d", &load)

		drivers = append(drivers, DriverLocation{
			DriverID:   name,
			OrderCount: load,
			State:      state,
		})
	}
	return drivers, nil
}

// RemoveDriver removes a driver from geo index (e.g., when going offline).
func (r *GeoRepository) RemoveDriver(ctx context.Context, driverID string) error {
	pipe := r.rdb.Pipeline()
	pipe.ZRem(ctx, driverLocationKey, driverID)
	pipe.Del(ctx, fmt.Sprintf(driverStateKey, driverID))
	pipe.Del(ctx, fmt.Sprintf(driverLoadKey, driverID))
	_, err := pipe.Exec(ctx)
	return fmt.Errorf("remove driver: %w", err)
}

Redis GEO Performance

GEOSEARCH menggunakan struktur data sorted set di Redis. Kompleksitas O(log N) untuk setiap anggota ditambah O(N) untuk hasil. Dengan 10.000 driver aktif, query radius 2km biasanya selesai dalam <1ms. Tapi hati-hati: Redis adalah single-threaded — GEOSEARCH dengan radius besar (seluruh kota) bisa block event loop selama puluhan ms.


3. Dispatch Matcher: Greedy Assignment dengan Scoring

Ini adalah inti sistem. Setelah batch terkumpul dan driver idle teridentifikasi, matcher menghitung score untuk setiap pasangan driver-order, lalu membuat assignment optimal.

Formula scoring:

score(d, o) = w1 * normalized_distance + 
              w2 * normalized_eta + 
              w3 * driver_load_penalty + 
              w4 * priority_bonus(o) +
              w5 * anti_starvation(d)
package dispatch

import (
	"context"
	"fmt"
	"math"
	"sort"
	"sync"
)

// MatchingConfig controls scoring weights.
type MatchingConfig struct {
	MaxSearchRadiusMeters float64
	MaxBatchSize          int
	WeightDistance        float64 // w1
	WeightETA             float64 // w2
	WeightLoad            float64 // w3
	WeightPriority        float64 // w4
	WeightStarvation      float64 // w5
	// Anti-starvation: increment every time a driver is skipped
	StarvationIncrement float64
}

// DriverMatchScore holds scoring for a driver-order pair.
type DriverMatchScore struct {
	DriverID string  `json:"driver_id"`
	OrderID  string  `json:"order_id"`
	Distance float64 `json:"distance_m"`
	ETA      float64 `json:"eta_s"`
	Score    float64 `json:"score"` // lower = better match
}

// MatchResult is the output of the matcher.
type MatchResult struct {
	Assignments []Assignment  `json:"assignments"`
	Scores      []DriverMatchScore `json:"scores"`
	Unassigned  []PendingOrder `json:"unassigned"`
}

// Assignment pairs a driver with one or more orders.
type Assignment struct {
	DriverID string   `json:"driver_id"`
	OrderIDs []string `json:"order_ids"`
}

// Matcher performs driver-order matching.
type Matcher struct {
	cfg         MatchingConfig
	mu          sync.Mutex
	starvation  map[string]float64 // driver_id -> accumulated starvation score
}

// NewMatcher creates a Matcher.
func NewMatcher(cfg MatchingConfig) *Matcher {
	return &Matcher{
		cfg:        cfg,
		starvation: make(map[string]float64),
	}
}

// Match assigns orders to drivers using greedy algorithm.
func (m *Matcher) Match(ctx context.Context, orders []PendingOrder, drivers []DriverLocation, geo *GeoRepository) (*MatchResult, error) {
	if len(orders) == 0 || len(drivers) == 0 {
		return &MatchResult{
			Unassigned: orders,
		}, nil
	}

	// Build score matrix: driver -> list of (order, score) sorted ascending
	type orderScore struct {
		order PendingOrder
		score float64
		dist  float64
		eta   float64
	}

	type driverScores struct {
		driver DriverLocation
		scores []orderScore
	}

	var scoredDrivers []driverScores
	for _, d := range drivers {
		// Find distance and ETA via Redis or routing service
		var scores []orderScore
		for _, o := range orders {
			dist, eta, err := m.estimateDistanceETA(ctx, geo, d, o)
			if err != nil {
				continue
			}
			score := m.computeScore(d, o, dist, eta)
			scores = append(scores, orderScore{
				order: o,
				score: score,
				dist:  dist,
				eta:   eta,
			})
		}
		sort.Slice(scores, func(i, j int) bool {
			return scores[i].score < scores[j].score
		})
		scoredDrivers = append(scoredDrivers, driverScores{
			driver: d,
			scores: scores,
		})
	}

	// Greedy assignment: pick lowest score pair, assign, remove from pool
	assignedOrders := make(map[string]bool) // order_id
	assignedDrivers := make(map[string]bool) // driver_id
	var assignments []Assignment
	var matchScores []DriverMatchScore

	for {
		bestScore := math.MaxFloat64
		var bestDriverIdx int
		var bestOrder orderScore
		found := false

		for i, ds := range scoredDrivers {
			if assignedDrivers[ds.driver.DriverID] {
				continue
			}
			for _, s := range ds.scores {
				if assignedOrders[s.order.OrderID] {
					continue
				}
				if s.score < bestScore {
					bestScore = s.score
					bestDriverIdx = i
					bestOrder = s
					found = true
				}
				break // first unassigned is lowest score for this driver
			}
		}

		if !found {
			break
		}

		driver := scoredDrivers[bestDriverIdx].driver
		assignedOrders[bestOrder.order.OrderID] = true
		assignedDrivers[driver.DriverID] = true

		assignments = append(assignments, Assignment{
			DriverID: driver.DriverID,
			OrderIDs: []string{bestOrder.order.OrderID},
		})
		matchScores = append(matchScores, DriverMatchScore{
			DriverID: driver.DriverID,
			OrderID:  bestOrder.order.OrderID,
			Distance: bestOrder.dist,
			ETA:      bestOrder.eta,
			Score:    bestScore,
		})

		// Reset starvation for assigned driver
		m.resetStarvation(driver.DriverID)

		// Check if driver has capacity for more orders (batch delivery)
		if len(assignments) > 0 && driver.OrderCount < driver.MaxOrders {
			// Allow driver to take another order
			delete(assignedDrivers, driver.DriverID)
		}
	}

	// Collect unassigned orders
	var unassigned []PendingOrder
	for _, o := range orders {
		if !assignedOrders[o.OrderID] {
			m.incrementStarvationForUnassigned(scoredDrivers)
			unassigned = append(unassigned, o)
		}
	}

	return &MatchResult{
		Assignments: assignments,
		Scores:      matchScores,
		Unassigned:  unassigned,
	}, nil
}

func (m *Matcher) computeScore(d DriverLocation, o PendingOrder, dist, eta float64) float64 {
	// Normalize distance: 0..1 based on max search radius
	normDist := dist / m.cfg.MaxSearchRadiusMeters
	// Normalize ETA: assume max 1800s (30 min)
	normETA := eta / 1800.0
	// Load penalty: driver capacity usage
	loadPenalty := float64(d.OrderCount) / float64(d.MaxOrders)

	starveBonus := m.getStarvation(d.DriverID)

	score := m.cfg.WeightDistance*normDist +
		m.cfg.WeightETA*normETA +
		m.cfg.WeightLoad*loadPenalty -
		m.cfg.WeightPriority*float64(o.Priority) -
		m.cfg.WeightStarvation*starveBonus

	return score
}

func (m *Matcher) estimateDistanceETA(ctx context.Context, geo *GeoRepository, d DriverLocation, o PendingOrder) (dist, eta float64, err error) {
	// Simplified: use Haversine for distance, estimate speed
	// In production: call OSRM or Google Maps Routing API
	distance := haversine(d.Lat, d.Lng, o.Lat, o.Lng)
	speed := 8.33 // ~30 km/h average for scooters
	eta = distance / speed
	return distance, eta, nil
}

func (m *Matcher) getStarvation(driverID string) float64 {
	m.mu.Lock()
	defer m.mu.Unlock()
	return m.starvation[driverID]
}

func (m *Matcher) resetStarvation(driverID string) {
	m.mu.Lock()
	defer m.mu.Unlock()
	delete(m.starvation, driverID)
}

func (m *Matcher) incrementStarvationForUnassigned(scoredDrivers []driverScores) {
	m.mu.Lock()
	defer m.mu.Unlock()
	for _, ds := range scoredDrivers {
		if !m.isDriverAssigned(ds.driver.DriverID) {
			m.starvation[ds.driver.DriverID] += m.cfg.StarvationIncrement
		}
	}
}

func (m *Matcher) isDriverAssigned(driverID string) bool {
	// In real impl, check against assignments
	return false
}

// haversine calculates distance between two coordinates in meters.
func haversine(lat1, lng1, lat2, lng2 float64) float64 {
	const R = 6371000 // Earth radius in meters
	dLat := (lat2 - lat1) * math.Pi / 180
	dLng := (lng2 - lng1) * math.Pi / 180
	a := math.Sin(dLat/2)*math.Sin(dLat/2) +
		math.Cos(lat1*math.Pi/180)*math.Cos(lat2*math.Pi/180)*
			math.Sin(dLng/2)*math.Sin(dLng/2)
	c := 2 * math.Atan2(math.Sqrt(a), math.Sqrt(1-a))
	return R * c
}

4. Driver State Machine

Setiap driver melewati siklus state yang ketat. State transition ini penting untuk:

  • Visibility: Tim ops tahu persis status setiap driver
  • Correctness: Sistem tahu driver mana yang available untuk dispatch
  • SLA tracking: Waktu di setiap state dihitung untuk analytics
package dispatch

import (
	"context"
	"fmt"
	"time"
)

// DriverState represents the lifecycle of a driver.
type DriverState string

const (
	DriverOffline   DriverState = "OFFLINE"
	DriverIdle      DriverState = "IDLE"
	DriverAssigned  DriverState = "ASSIGNED"
	DriverToHub     DriverState = "TO_HUB"
	DriverPicking   DriverState = "PICKING"
	DriverDelivering DriverState = "DELIVERING"
	DriverCompleted DriverState = "COMPLETED"
)

// ValidTransitions defines allowed state transitions.
var validTransitions = map[DriverState][]DriverState{
	DriverOffline:    {DriverIdle},
	DriverIdle:       {DriverAssigned, DriverOffline},
	DriverAssigned:   {DriverToHub, DriverIdle}, // accept or reject
	DriverToHub:      {DriverPicking, DriverOffline},
	DriverPicking:    {DriverDelivering, DriverOffline},
	DriverDelivering: {DriverDelivering, DriverCompleted, DriverOffline},
	DriverCompleted:  {DriverIdle},
}

// StateEvent represents a driver state change event (event sourcing).
type StateEvent struct {
	DriverID    string      `json:"driver_id"`
	From        DriverState `json:"from"`
	To          DriverState `json:"to"`
	OrderID     string      `json:"order_id,omitempty"`
	Timestamp   time.Time   `json:"timestamp"`
	Reason      string      `json:"reason,omitempty"`
	Coordinates []float64   `json:"coordinates,omitempty"`
}

// StateMachine manages driver lifecycle.
type StateMachine struct {
	events    chan StateEvent
	eventSink EventStore
}

// EventStore persists state events for audit trail.
type EventStore interface {
	Append(ctx context.Context, event StateEvent) error
	GetByDriver(ctx context.Context, driverID string, limit int) ([]StateEvent, error)
}

// NewStateMachine creates a StateMachine.
func NewStateMachine(sink EventStore, bufferSize int) *StateMachine {
	sm := &StateMachine{
		events:    make(chan StateEvent, bufferSize),
		eventSink: sink,
	}
	go sm.processEvents()
	return sm
}

// Transition attempts a state change and emits an event.
func (sm *StateMachine) Transition(ctx context.Context, driverID string, from, to DriverState, opts ...StateOption) error {
	// Validate transition
	allowed, ok := validTransitions[from]
	if !ok {
		return fmt.Errorf("state machine: unknown from state %q", from)
	}

	valid := false
	for _, s := range allowed {
		if s == to {
			valid = true
			break
		}
	}
	if !valid {
		return fmt.Errorf("state machine: invalid transition %s -> %s", from, to)
	}

	event := StateEvent{
		DriverID:  driverID,
		From:      from,
		To:        to,
		Timestamp: time.Now().UTC(),
	}
	for _, opt := range opts {
		opt(&event)
	}

	select {
	case sm.events <- event:
		return nil
	case <-ctx.Done():
		return ctx.Err()
	}
}

func (sm *StateMachine) processEvents() {
	for event := range sm.events {
		ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
		if err := sm.eventSink.Append(ctx, event); err != nil {
			// Log error, retry via DLQ
			// logger.Error("failed to persist state event", "error", err)
		}
		cancel()
	}
}

// StateOption modifies a StateEvent.
type StateOption func(*StateEvent)

// WithOrderID sets the order ID on a state event.
func WithOrderID(orderID string) StateOption {
	return func(e *StateEvent) {
		e.OrderID = orderID
	}
}

// WithReason sets the reason for a state transition.
func WithReason(reason string) StateOption {
	return func(e *StateEvent) {
		e.Reason = reason
	}
}

// WithCoordinates sets driver coordinates at transition time.
func WithCoordinates(lat, lng float64) StateOption {
	return func(e *StateEvent) {
		e.Coordinates = []float64{lat, lng}
	}
}

5. Re-assignment Handler: Driver Reject atau Timeout

Driver tidak selalu menerima assignment. Mereka bisa sibuk, jarak terlalu jauh, atau insentif kurang menarik. Ketika reject atau timeout terjadi, sistem harus reassign ke driver terbaik berikutnya tanpa menyebabkan starvation.

package dispatch

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

const (
	driverAcceptTimeout = 15 * time.Second
	maxReassignAttempts = 3
)

// ReassignmentHandler manages fallback when drivers reject.
type ReassignmentHandler struct {
	matcher  *Matcher
	geo      *GeoRepository
	tracker  *StateMachine
	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)
}

// NewReassignmentHandler creates a handler.
func NewReassignmentHandler(
	matcher *Matcher,
	geo *GeoRepository,
	tracker *StateMachine,
	logger *slog.Logger,
	metrics MetricsRecorder,
) *ReassignmentHandler {
	return &ReassignmentHandler{
		matcher: matcher,
		geo:     geo,
		tracker: tracker,
		logger:  logger,
		metrics: metrics,
	}
}

// HandleReject processes a driver rejection and finds replacement.
func (h *ReassignmentHandler) HandleReject(ctx context.Context, driverID, orderID string) error {
	h.metrics.IncrementCounter("dispatch.reject", map[string]string{"driver": driverID})

	// Record rejection in driver profile (for quality scoring)
	if err := h.recordRejection(ctx, driverID, orderID); err != nil {
		h.logger.Error("record rejection", "error", err)
	}

	// Return driver to IDLE
	if err := h.tracker.Transition(ctx, driverID, DriverAssigned, DriverIdle,
		WithReason("driver_rejected"),
	); err != nil {
		return fmt.Errorf("reassign: transition driver back to idle: %w", err)
	}

	return nil
}

// HandleTimeout processes a driver not responding to assignment.
func (h *ReassignmentHandler) HandleTimeout(ctx context.Context, driverID, orderID string) error {
	h.metrics.IncrementCounter("dispatch.timeout", map[string]string{"driver": driverID})

	// Penalize driver's priority for future assignments
	if err := h.penalizeDriver(ctx, driverID); err != nil {
		h.logger.Error("penalize driver", "error", err)
	}

	// Return driver to IDLE
	if err := h.tracker.Transition(ctx, driverID, DriverAssigned, DriverIdle,
		WithReason("response_timeout"),
	); err != nil {
		return fmt.Errorf("reassign: timeout transition: %w", err)
	}

	return nil
}

// FindReplacement finds the next best driver for an order.
func (h *ReassignmentHandler) FindReplacement(ctx context.Context, order PendingOrder, excludeDriverIDs []string) (*Assignment, error) {
	for attempt := 0; attempt < maxReassignAttempts; attempt++ {
		// Fetch fresh drivers
		drivers, err := h.geo.FindIdleDriversNearby(ctx, order.Lat, order.Lng, 5000)
		if err != nil {
			return nil, fmt.Errorf("reassign find drivers: %w", err)
		}

		// Filter excluded
		excluded := make(map[string]bool, len(excludeDriverIDs))
		for _, id := range excludeDriverIDs {
			excluded[id] = true
		}
		var available []DriverLocation
		for _, d := range drivers {
			if !excluded[d.DriverID] {
				available = append(available, d)
			}
		}

		if len(available) == 0 {
			return nil, ErrNoDriversAvailable
		}

		// Re-match
		result, err := h.matcher.Match(ctx, []PendingOrder{order}, available, h.geo)
		if err != nil {
			return nil, fmt.Errorf("reassign match: %w", err)
		}

		if len(result.Assignments) > 0 {
			h.metrics.IncrementCounter("dispatch.reassign.success", nil)
			return &result.Assignments[0], nil
		}

		// Exponential backoff before retry
		time.Sleep(time.Duration(1<<attempt) * 500 * time.Millisecond)
	}

	h.metrics.IncrementCounter("dispatch.reassign.exhausted", nil)
	return nil, ErrReassignExhausted
}

func (h *ReassignmentHandler) recordRejection(ctx context.Context, driverID, orderID string) error {
	// Implementation: persist to DB for driver quality scoring
	return nil
}

func (h *ReassignmentHandler) penalizeDriver(ctx context.Context, driverID string) error {
	// Implementation: reduce driver priority score
	return nil
}

// Domain errors.
var (
	ErrNoDriversAvailable = fmt.Errorf("dispatch: no drivers available")
	ErrReassignExhausted  = fmt.Errorf("dispatch: reassign attempts exhausted")
)

6. Batch Delivery Multi-Order

Satu driver bisa membawa beberapa order dalam satu trip. Ini meningkatkan efisiensi tapi menambah kompleksitas routing.

package dispatch

import (
	"context"
	"fmt"
	"math"
	"sort"
)

// BatchDelivery manages multi-order trips.
type BatchDelivery struct {
	DriverID     string        `json:"driver_id"`
	OrderIDs     []string      `json:"order_ids"`
	Dropoffs     []Dropoff     `json:"dropoffs"`
	OptimizedETA float64       `json:"optimized_eta_s"`
}

// Dropoff represents a delivery point.
type Dropoff struct {
	OrderID string    `json:"order_id"`
	Lat     float64   `json:"lat"`
	Lng     float64   `json:"lng"`
}

// RouteOptimizer optimizes multi-stop delivery route.
type RouteOptimizer struct {
	geo *GeoRepository
}

// NewRouteOptimizer creates a RouteOptimizer.
func NewRouteOptimizer(geo *GeoRepository) *RouteOptimizer {
	return &RouteOptimizer{geo: geo}
}

// OptimizeDropoffs returns the optimal order of drop-offs using nearest-neighbor.
func (ro *RouteOptimizer) OptimizeDropoffs(ctx context.Context, hubLat, hubLng float64, dropoffs []Dropoff) ([]Dropoff, error) {
	if len(dropoffs) == 0 {
		return nil, fmt.Errorf("route optimizer: empty dropoffs")
	}
	if len(dropoffs) == 1 {
		return dropoffs, nil
	}

	// Nearest-neighbor heuristic from hub
	sorted := make([]Dropoff, len(dropoffs))
	copy(sorted, dropoffs)

	currentLat, currentLng := hubLat, hubLng
	for i := 0; i < len(sorted); i++ {
		nearestIdx := i
		nearestDist := math.MaxFloat64

		for j := i; j < len(sorted); j++ {
			dist := haversine(currentLat, currentLng, sorted[j].Lat, sorted[j].Lng)
			if dist < nearestDist {
				nearestDist = dist
				nearestIdx = j
			}
		}

		// Swap nearest to position i
		sorted[i], sorted[nearestIdx] = sorted[nearestIdx], sorted[i]
		currentLat, currentLng = sorted[i].Lat, sorted[i].Lng
	}

	return sorted, nil
}

// EstimateBatchETA calculates total ETA for a multi-order batch.
func (ro *RouteOptimizer) EstimateBatchETA(ctx context.Context, hubLat, hubLng float64, dropoffs []Dropoff) (float64, error) {
	optimized, err := ro.OptimizeDropoffs(ctx, hubLat, hubLng, dropoffs)
	if err != nil {
		return 0, err
	}

	var totalETA float64
	currentLat, currentLng := hubLat, hubLng

	// Picking time at hub: assume 2 min per order
	totalETA += float64(len(optimized)) * 120

	for _, d := range optimized {
		dist := haversine(currentLat, currentLng, d.Lat, d.Lng)
		speed := 8.33 // m/s (~30 km/h)
		totalETA += dist / speed
		// Drop-off time: 30 seconds per stop
		totalETA += 30
		currentLat, currentLng = d.Lat, d.Lng
	}

	return totalETA, nil
}

Edge Cases dan Penanganannya

1. Driver Reject Loop

Masalah: Seorang driver terus-menerus dapat assignment dan terus reject karena sistem tidak belajar dari pola reject sebelumnya.

Solusi: Rejection tracker dengan sliding window — jika driver reject >3x dalam 30 menit, turunkan priority score atau skip sementara.

func (h *ReassignmentHandler) shouldSkipDriver(driverID string) bool {
	recentRejects, _ := h.rejectionStore.CountRecent(driverID, 30*time.Minute)
	return recentRejects >= 3
}

2. Order Timeout

Masalah: Order sudah di-assign ke driver, tapi driver tidak merespon dalam 15 detik. Order menunggu, pelanggan tidak sabar.

Solusi: Timer goroutine per assignment. Jika timeout, otomatis trigger reassignment dengan event timeout.

func (h *ReassignmentHandler) StartAssignmentTimer(ctx context.Context, order PendingOrder, driverID string) {
	go func() {
		timer := time.NewTimer(driverAcceptTimeout)
		select {
		case <-timer.C:
			h.HandleTimeout(ctx, driverID, order.OrderID)
			h.FindReplacement(ctx, order, []string{driverID})
		case <-ctx.Done():
			timer.Stop()
		}
	}()
}

3. Starvation Driver Jauh

Masalah: Driver yang berada di pinggiran kota jarang dapat order karena selalu kalah skor dengan driver di pusat kota.

Solusi: Anti-starvation score incremental seperti di implementasi Matcher di atas. Setiap kali driver dilewati dalam suatu batch, starvation score naik, membuatnya lebih mungkin terpilih di batch berikutnya.

4. Concurrent Assignment Race

Masalah: Dua goroutine dispatch secara bersamaan meng-assign driver yang sama ke dua order berbeda.

Solusi: Redis locks atau database-level optimistic locking dengan versi state.

func (r *GeoRepository) TryAssignDriver(ctx context.Context, driverID string, version int) (bool, error) {
	// CAS (Check-and-Set) via Redis Lua script
	script := redis.NewScript(`
		local state = redis.call("GET", KEYS[1])
		if state == "idle" then
			redis.call("SET", KEYS[1], "assigned")
			return 1
		end
		return 0
	`)
	result, err := script.Run(ctx, r.rdb, []string{fmt.Sprintf(driverStateKey, driverID)}).Int()
	if err != nil {
		return false, err
	}
	return result == 1, nil
}

5. Mass Driver Drop-off

Masalah: 30 driver tiba-tiba offline (gempa, demonstrasi, traffic massive). Order menumpuk.

Solusi: Dynamic radius expansion — jika idle driver dalam radius normal tidak cukup, expand radius 2x, 3x, hingga 10km. Dashboard alert ke ops team.

func findDriversWithBackoff(ctx context.Context, geo *GeoRepository, lat, lng float64) ([]DriverLocation, error) {
	radii := []float64{2000, 3000, 5000, 10000}
	for _, radius := range radii {
		drivers, err := geo.FindIdleDriversNearby(ctx, lat, lng, radius)
		if err != nil {
			return nil, err
		}
		if len(drivers) > 0 {
			return drivers, nil
		}
	}
	return nil, ErrNoDriversAvailable
}

Key Takeaways

Batch Before You Match

Jangan dispatch order satu per satu. Kumpulkan dalam batch window 3-10 detik untuk assignment yang lebih optimal. Trade-off: latency vs optimality.

Redis GEO is Your Friend

Gunakan GEOSEARCH untuk query driver terdekat dalam O(log N). Tapi jangan lupa — Redis single-threaded, jangan query radius seluruh kota.

State Machine Protects Integrity

Driver state harus strict. Validasi setiap transisi. Gunakan event sourcing untuk audit trail dan debugging.

Anti-Starvation is Design Necessity

Driver pinggiran akan terus kalah tanpa mekanisme starvation. Increment priority setiap batch dilewati. Jangan tunggu komplain driver.

Reassignment is Resilience

Driver reject atau timeout adalah normal flow, bukan error. Siapkan timer per assignment dan reassignment pipeline dengan eksponensial backoff.

Monitor Everything

Track metrics: dispatch latency, match score distribution, reject rate, timeout rate, average batch size. Grafana dashboard adalah ops tool utama.

Penutup

Dispatch engine adalah salah satu komponen paling kritikal di Q-Commerce. Kesalahan di sini langsung terasa: order tidak terkirim, pelanggan komplain, driver frustrasi. Dengan arsitektur yang tepat — batching, geospatial indexing, scoring-based assignment, strict state machine, dan anti-starvation — kamu bisa membangun sistem yang reliable bahkan di peak load 10x normal.

Kode di atas adalah production-oriented template. Untuk production, tambahkan: circuit breaker ke Redis, rate limiter di acceptance endpoint, distributed tracing (OpenTelemetry), dan chaos engineering testing untuk skenario mass driver drop-off.


Related Engineering & Tech Articles

Order Dispatch & Driver Allocation: Hungarian Algorithm u... | Faisal Affan