diff --git a/controllers/milerAppController.go b/controllers/milerAppController.go index 4e9b7a0..0641f9f 100644 --- a/controllers/milerAppController.go +++ b/controllers/milerAppController.go @@ -10,6 +10,7 @@ import ( "doormile/constants" "doormile/db" + "doormile/internal/assignment" "doormile/internal/notify" "doormile/models" "doormile/utils" @@ -66,6 +67,12 @@ func MilerStartDuty(c *fiber.Ctx) error { // location ping. Skipped when the app sent no coordinates. indexMilerLocation(milerUserID, req.Lat, req.Lon) + // Give the pending orders around them to this rider now, rather than at + // the next periodic sweep (balanced assignment, Phase B). + if req.Lat != 0 && req.Lon != 0 { + go assignment.SweepNear(req.Lat, req.Lon) + } + return utils.OK(c, fiber.Map{ "dutylogid": dutyLog.Dutylogid, "loginat": dutyLog.Loginat, diff --git a/docs/balanced-assignment-plan.md b/docs/balanced-assignment-plan.md index ba6f67e..43ff696 100644 --- a/docs/balanced-assignment-plan.md +++ b/docs/balanced-assignment-plan.md @@ -1,6 +1,24 @@ # Balanced auto-assignment: implementation plan -**Status:** planned, not started (2026-10-06). +**Status:** Phases A and B built and tested, not committed or deployed (2026-10-06). Phase C is a server setting (`MILER_MAX_ACTIVE_BOOKINGS=20`). + +### Phase B: what was built +- `internal/assignment/release.go` (new): + - **`releaseUnaccepted`** (runs at the start of every sweep): an assignment still `Assigned` on an order still `Miler_Assigned`, older than `ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES` (default 10; 0 = off), goes back to the pool. The order returns to `Pending_Pickup` with no rider, and the assignment becomes `Reassigned` with a note. The customer stage is walked back (`cxstage.Release`), and the rider is set to Available if they hold nothing else. Every write is conditional, so a rider accepting at the same moment wins. + - **Only auto-assignments are released.** They're now labelled `Remarks = "Auto-assigned"` (`autoAssignedRemark`). A rider hand-picked by ops or hub staff is never taken away. (The hub console's manual assign records no "assigned by", so the label is the only reliable marker.) Assignments made before this change have no label and are never released. + - **`ridersToSkip` / `withoutRiders`** (applied in `selectMilerWithAI`): an order is never offered back to a rider it was released from, who rejected it, or who cancelled it. + - **`SweepNear`**: called from `MilerStartDuty` in a goroutine. It assigns the pending orders within 10 km of a rider who just started duty, instead of waiting up to 5 minutes for the next sweep. +- `pendingForSweep` now also loads the pickup coordinates (needed by `pendingNear`). +- Tests: `release_test.go` (settings, skip filter, distance) and `release_pg_test.go` (release after timeout; inside the timeout, accepted and manual orders never released; the rider freed when nothing else is held; a concurrent accept wins; a released or rejected order goes to another rider; the duty-start sweep only picks nearby orders). +- End to end on the real backend: a waiting order was assigned **0.19 s** after the rider started duty. An unaccepted order was released after the timeout and given to the other rider in the same sweep. Accepted and manual orders were untouched after 3 sweeps. Bulk split was still 3/3. + +### Phase A: what was built +- **In hand** = `openStopsToday` (today's open, non-cancelled orders), the same count the ceiling uses. Using all open records ever would let stale, abandoned orders skew the balance. +- `selectMilerWithAI` keeps only the least-loaded riders (`leastLoaded`), then the AI or the fallback chooses among them. It now also returns the full pool. +- `pickBestFromCandidates` → `betterChoice`: in hand ↑, then session ↑ (`sessionStops`, counted since the open `milerdutylogs.loginat`), then distance ↑, then rating ↓. +- `finalizeChoice` (both commit paths): takes `pg_advisory_xact_lock(7270001)`, recounts the pool, keeps the original choice if it's still least loaded, otherwise picks the best of the least loaded. Returns `errNoRiderCapacity` if everyone is at the ceiling (the attempt is retried later). The AI decision id is dropped when the choice changes. +- Tests: `balance_test.go` (ranking) and `balance_pg_test.go` (50 over 5 → 10 each; 23 → 5/5/5/4/4; **50 concurrent → within 1**; newcomers catch up; tie-break; all at the ceiling → waits). +- End to end on the real backend: bulk upload of 9 orders with riders at 0, 2 and 5 km → **3/3/3**. A 4th rider then joins and 4 more orders arrive: the newcomer gets 3, the 4th goes to the nearest (all tied). **Goal:** however many orders and riders there are, split the orders **equally** among the riders, with riders logging in and out all day. --- diff --git a/internal/assignment/ai_layer.go b/internal/assignment/ai_layer.go index 0aef135..b03e348 100644 --- a/internal/assignment/ai_layer.go +++ b/internal/assignment/ai_layer.go @@ -4,6 +4,7 @@ import ( "bytes" "context" "encoding/json" + "errors" "fmt" "net/http" "os" @@ -16,6 +17,7 @@ import ( "doormile/utils" "github.com/redis/go-redis/v9" + "gorm.io/gorm" ) // ─── Request / response types ──────────────────────────────────────────────── @@ -68,18 +70,28 @@ type aiDecisionResponse struct { // audit. Falls back to the original distance/load/rating formula when the AI // layer is unreachable or times out. Returns (nil, nil, false) when there are no // eligible candidates or when the AI layer escalates the booking. -func selectMilerWithAI(booking *models.PickupBooking, nearby []redis.GeoLocation) (*milerCandidate, *uint64, bool) { - candidates, aiCandidates := collectEligibleCandidates(nearby) - if len(candidates) == 0 { - return nil, nil, false +func selectMilerWithAI(booking *models.PickupBooking, nearby []redis.GeoLocation) (*milerCandidate, []*milerCandidate, *uint64, bool) { + // Never offer an order back to a rider it was released from, who + // rejected it or who cancelled it. + nearby = withoutRiders(nearby, ridersToSkip(booking.Bookingid)) + + pool, poolAI := collectEligibleCandidates(nearby) + if len(pool) == 0 { + return nil, nil, nil, false } + // Balance: only the riders holding the fewest orders are considered. The + // AI (or the fallback) chooses among them, so the split stays even however + // many orders and riders there are. The full pool goes back to the caller + // for the final re-check at commit time (finalizeChoice). + candidates, aiCandidates := leastLoaded(pool, poolAI) + decision, err := callDecisionEngine(booking, aiCandidates) if err != nil { utils.Warn("AI_LAYER_FALLBACK: decide-assignment unreachable, using legacy scoring", "booking_id", booking.Bookingid, "error", err) best := pickBestFromCandidates(candidates) - return best, nil, best != nil + return best, pool, nil, best != nil } utils.Info("Assignment: AI layer responded", @@ -95,7 +107,7 @@ func selectMilerWithAI(booking *models.PickupBooking, nearby []redis.GeoLocation "booking_id", booking.Bookingid, "reasoning", decision.Reasoning, ) - return nil, nil, false + return nil, nil, nil, false } var decisionID *uint64 @@ -106,17 +118,42 @@ func selectMilerWithAI(booking *models.PickupBooking, nearby []redis.GeoLocation for _, c := range candidates { if c.profile.Userid == decision.ChosenMilerID { - return c, decisionID, true + return c, pool, decisionID, true } } - // AI returned an ID that is not in our eligibility set — fall back safely. + // AI returned an ID that is not among the least-loaded riders — fall back. utils.Warn("AI_LAYER_FALLBACK: chosen miler not in eligible set, using legacy scoring", "booking_id", booking.Bookingid, "chosen_miler_id", decision.ChosenMilerID, ) best := pickBestFromCandidates(candidates) - return best, nil, best != nil + return best, pool, nil, best != nil +} + +// leastLoaded keeps the candidates holding the fewest orders in hand (and +// the matching AI rows, which are parallel to them; ai may be nil). +func leastLoaded(cands []*milerCandidate, ai []aiCandidate) ([]*milerCandidate, []aiCandidate) { + if len(cands) == 0 { + return cands, ai + } + least := cands[0].activeBookings + for _, c := range cands[1:] { + if c.activeBookings < least { + least = c.activeBookings + } + } + var outC []*milerCandidate + var outAI []aiCandidate + for i, c := range cands { + if c.activeBookings == least { + outC = append(outC, c) + if i < len(ai) { + outAI = append(outAI, ai[i]) + } + } + } + return outC, outAI } // ─── Candidate collection ──────────────────────────────────────────────────── @@ -171,6 +208,7 @@ func collectEligibleCandidates(nearby []redis.GeoLocation) ([]*milerCandidate, [ profile: profile, distanceKm: loc.Dist, activeBookings: activeCount, + sessionStops: sessionStops(db.DB, milerUserID), }) aiCandidates = append(aiCandidates, aiCandidate{ MilerID: milerUserID, @@ -195,8 +233,14 @@ func collectEligibleCandidates(nearby []redis.GeoLocation) ([]*milerCandidate, [ // pile of old or cancelled orders sat at the cap for good and every new order // stayed Pending. func openStopsToday(milerUserID int) int64 { + return openStopsTodayIn(db.DB, milerUserID) +} + +// openStopsTodayIn is openStopsToday on a given handle, so the commit can +// recount inside its own transaction. +func openStopsTodayIn(h *gorm.DB, milerUserID int) int64 { var n int64 - db.DB.Table("bookingassignments AS ba"). + h.Table("bookingassignments AS ba"). Joins("JOIN pickupbookings pb ON pb.bookingid = ba.bookingid"). Where("ba.mileruserid = ? AND ba.assignmentstatus IN ? AND ba.assignedat >= ? AND pb.status <> ?", milerUserID, @@ -221,6 +265,66 @@ func milerHasFreshGPS(milerUserID int) bool { return n > 0 } +// sessionStops counts the orders given to the rider since their current duty +// session started (the latest milerdutylogs row with no logout), excluding +// ones that fell through (cancelled, rejected, reassigned). 0 when the rider +// has no open session. +func sessionStops(h *gorm.DB, milerUserID int) int64 { + var n int64 + h.Table("bookingassignments"). + Where("mileruserid = ? AND assignmentstatus NOT IN ? AND assignedat >= "+ + "(SELECT MAX(loginat) FROM milerdutylogs WHERE userid = ? AND logoutat IS NULL)", + milerUserID, + []string{constants.AssignmentCancelled, constants.AssignmentRejected, constants.AssignmentReassigned}, + milerUserID). + Count(&n) + return n +} + +// errNoRiderCapacity: by the time the commit ran, every candidate had reached +// the per-rider ceiling. The attempt is retried later like "no rider found". +var errNoRiderCapacity = errors.New("every candidate rider is at the order limit") + +// decisionLockKey serialises the final rider choice across attempts and +// replicas (a transaction-scoped Postgres advisory lock). Held only for a +// handful of quick counts and the commit, never across the AI call. +const decisionLockKey int64 = 7_270_001 + +// finalizeChoice re-checks the choice inside the commit transaction. Attempts +// for different orders run in parallel — a bulk upload of 50 starts 50 — and +// each picked "the least-loaded rider" from counts read before the others +// committed, so they would all pile onto the same one. Under the lock the +// loads are recounted: the chosen rider is kept if still among the least +// loaded and under the ceiling, otherwise the best of the least loaded is +// taken instead (same order as betterChoice). +func finalizeChoice(tx *gorm.DB, chosen *milerCandidate, pool []*milerCandidate) (*milerCandidate, error) { + if len(pool) == 0 { + pool = []*milerCandidate{chosen} + } + if err := tx.Exec("SELECT pg_advisory_xact_lock(?)", decisionLockKey).Error; err != nil { + return nil, fmt.Errorf("decision lock: %w", err) + } + ceiling := maxActiveBookings() + var open []*milerCandidate + for _, c := range pool { + c.activeBookings = openStopsTodayIn(tx, c.profile.Userid) + c.sessionStops = sessionStops(tx, c.profile.Userid) + if c.activeBookings < ceiling { + open = append(open, c) + } + } + if len(open) == 0 { + return nil, errNoRiderCapacity + } + least, _ := leastLoaded(open, nil) + for _, c := range least { + if c.profile.Userid == chosen.profile.Userid { + return c, nil // the original choice still holds + } + } + return pickBestFromCandidates(least), nil +} + // startOfISTDay is midnight today in India, the boundary for "today's" // open stops. An absolute instant, so it compares correctly with assignedat // whether the column stores IST wall-clock or a real timestamp. @@ -296,19 +400,31 @@ func fetchHubData(hubID *int) (resolvedHubID int, hubLoad int64, hubCapacity int // eligible slice: score = distance_km*1.0 + active_bookings*2.0 - rating*0.5 func pickBestFromCandidates(candidates []*milerCandidate) *milerCandidate { var best *milerCandidate - bestScore := 1e18 - for _, c := range candidates { - score := c.distanceKm*1.0 + float64(c.activeBookings)*2.0 - c.profile.Rating*0.5 - if score < bestScore { - bestScore = score + if best == nil || betterChoice(c, best) { best = c } } - return best } +// betterChoice is the balancing order: fewest orders in hand, then fewest +// orders this duty session, then nearest to the pickup, then best rated. +// It replaced a weighted score in which distance dominated, so riders near +// the pickups took most of the orders. +func betterChoice(a, b *milerCandidate) bool { + if a.activeBookings != b.activeBookings { + return a.activeBookings < b.activeBookings + } + if a.sessionStops != b.sessionStops { + return a.sessionStops < b.sessionStops + } + if a.distanceKm != b.distanceKm { + return a.distanceKm < b.distanceKm + } + return a.profile.Rating > b.profile.Rating +} + // ─── AI layer HTTP call ────────────────────────────────────────────────────── func callDecisionEngine(booking *models.PickupBooking, candidates []aiCandidate) (aiDecisionResponse, error) { diff --git a/internal/assignment/balance_pg_test.go b/internal/assignment/balance_pg_test.go new file mode 100644 index 0000000..73bce8e --- /dev/null +++ b/internal/assignment/balance_pg_test.go @@ -0,0 +1,266 @@ +package assignment + +import ( + "errors" + "fmt" + "sync" + "testing" + "time" + + "doormile/constants" + "doormile/models" + + "github.com/redis/go-redis/v9" + "gorm.io/gorm" +) + +// Balanced auto-assignment (docs/balanced-assignment-plan.md, Phase A) +// against a real Postgres. Skipped unless REGISTRY_TEST_DSN is set — see +// eligibility_pg_test.go. + +// onDuty adds a rider with live GPS and an open duty session. +func onDuty(t *testing.T, gdb *gorm.DB, id int, sessionStart time.Time) { + t.Helper() + rider(t, gdb, id, constants.MilerAvailable, time.Minute) + mustDo(t, gdb.Create(&models.MilerDutyLog{Userid: id, Loginat: sessionStart}).Error) +} + +// assignOne runs one attempt the way tryAssign does (minus the Redis GEO +// lookup, which `nearby` stands in for) and returns the rider it went to. +func assignOne(t *testing.T, gdb *gorm.DB, bookingID int, nearby []redis.GeoLocation) (int, error) { + t.Helper() + var b models.PickupBooking + if err := gdb.First(&b, bookingID).Error; err != nil { + return 0, err + } + cand, pool, dec, found := selectMilerWithAI(&b, nearby) + if !found { + return 0, nil + } + final, err := commitAssignment(&b, cand, pool, dec) + if err != nil { + return 0, err + } + return final.profile.Userid, nil +} + +// near puts every rider at a different distance from the pickup, so a +// distance-first rule would pile everything onto rider ids[0]. +func near(ids ...int) []redis.GeoLocation { + out := make([]redis.GeoLocation, len(ids)) + for i, id := range ids { + out[i] = redis.GeoLocation{Name: fmt.Sprint(id), Dist: 0.5 + float64(i)*1.5} + } + return out +} + +func perRider(t *testing.T, gdb *gorm.DB) map[int]int64 { + t.Helper() + type row struct { + Mileruserid int + N int64 + } + var rows []row + mustDo(t, gdb.Table("bookingassignments").Select("mileruserid, COUNT(*) AS n"). + Where("assignmentstatus IN ?", []string{constants.AssignmentAssigned, constants.AssignmentAccepted}). + Group("mileruserid").Scan(&rows).Error) + out := map[int]int64{} + for _, r := range rows { + out[r.Mileruserid] = r.N + } + return out +} + +func spread(m map[int]int64, ids ...int) (lo, hi int64) { + lo, hi = 1<<62, -1 + for _, id := range ids { + n := m[id] + if n < lo { + lo = n + } + if n > hi { + hi = n + } + } + return lo, hi +} + +func balanceSetup(t *testing.T) *gorm.DB { + gdb := eligibilityDB(t) + t.Setenv("AI_LAYER_BASE_URL", "http://127.0.0.1:1") // unreachable: the fallback chooses + t.Setenv("MILER_MAX_ACTIVE_BOOKINGS", "20") + return gdb +} + +// 5 riders, 50 orders one after another: 10 each, although rider 1 is the +// nearest to every pickup. +func TestBalancedSplitSequential(t *testing.T) { + gdb := balanceSetup(t) + ids := []int{1, 2, 3, 4, 5} + for _, id := range ids { + onDuty(t, gdb, id, time.Now().Add(-time.Hour)) + } + for i := 0; i < 50; i++ { + b := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{}) + if _, err := assignOne(t, gdb, b, near(ids...)); err != nil { + t.Fatal(err) + } + } + got := perRider(t, gdb) + for _, id := range ids { + if got[id] != 10 { + t.Fatalf("per rider = %v, want 10 each", got) + } + } +} + +// Any count splits within one: 23 orders over 5 riders -> 5,5,5,4,4. +func TestBalancedSplitUneven(t *testing.T) { + gdb := balanceSetup(t) + ids := []int{1, 2, 3, 4, 5} + for _, id := range ids { + onDuty(t, gdb, id, time.Now().Add(-time.Hour)) + } + for i := 0; i < 23; i++ { + b := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{}) + if _, err := assignOne(t, gdb, b, near(ids...)); err != nil { + t.Fatal(err) + } + } + if lo, hi := spread(perRider(t, gdb), ids...); lo != 4 || hi != 5 { + t.Fatalf("spread %d..%d, want 4..5", lo, hi) + } +} + +// A bulk upload: 50 orders attempted at the same moment. Without the final +// re-check under the lock they all saw the same counts and piled onto one +// rider. Still balanced, nobody over the ceiling, one rider per order. +func TestBalancedSplitConcurrentBurst(t *testing.T) { + gdb := balanceSetup(t) + ids := []int{1, 2, 3, 4, 5} + for _, id := range ids { + onDuty(t, gdb, id, time.Now().Add(-time.Hour)) + } + bookings := make([]int, 50) + for i := range bookings { + bookings[i] = order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{}) + } + var wg sync.WaitGroup + errs := make(chan error, len(bookings)) + for _, b := range bookings { + wg.Add(1) + go func(b int) { + defer wg.Done() + if _, err := assignOne(t, gdb, b, near(ids...)); err != nil { + errs <- err + } + }(b) + } + wg.Wait() + close(errs) + for err := range errs { + t.Fatal(err) + } + got := perRider(t, gdb) + if lo, hi := spread(got, ids...); hi-lo > 1 { + t.Fatalf("burst split %v: spread %d..%d, want within 1", got, lo, hi) + } + var total int64 + for _, n := range got { + total += n + } + var perOrder int64 + gdb.Raw("SELECT COALESCE(MAX(n),0) FROM (SELECT COUNT(*) n FROM bookingassignments GROUP BY bookingid) x").Scan(&perOrder) + if total != 50 || perOrder != 1 { + t.Fatalf("assigned %d (want 50), max riders on one order %d (want 1)", total, perOrder) + } +} + +// Riders come and go: 3 riders hold 4 each, 2 log in fresh. The next 8 orders +// go to the newcomers (4 each), then everyone shares. +func TestNewRidersCatchUpThenShare(t *testing.T) { + gdb := balanceSetup(t) + busy := []int{1, 2, 3} + fresh := []int{4, 5} + for _, id := range busy { + onDuty(t, gdb, id, time.Now().Add(-3*time.Hour)) + for i := 0; i < 4; i++ { + order(t, gdb, id, constants.BookingMilerAssigned, constants.AssignmentAccepted, time.Now().Add(-time.Hour)) + } + } + for _, id := range fresh { + onDuty(t, gdb, id, time.Now()) + } + all := append(append([]int{}, busy...), fresh...) + for i := 0; i < 8; i++ { + b := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{}) + id, err := assignOne(t, gdb, b, near(all...)) + if err != nil { + t.Fatal(err) + } + if id != 4 && id != 5 { + t.Fatalf("order %d went to busy rider %d while a newcomer had fewer", i, id) + } + } + if lo, hi := spread(perRider(t, gdb), all...); lo != 4 || hi != 4 { + t.Fatalf("after catch-up spread %d..%d, want all 4", lo, hi) + } + // From here on everyone shares: 5 more orders -> one each. + for i := 0; i < 5; i++ { + b := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{}) + if _, err := assignOne(t, gdb, b, near(all...)); err != nil { + t.Fatal(err) + } + } + if lo, hi := spread(perRider(t, gdb), all...); lo != 5 || hi != 5 { + t.Fatalf("after sharing spread %d..%d, want all 5", lo, hi) + } +} + +// Tie on orders in hand: fewer orders this duty session wins, then distance. +func TestTieBreakSessionThenDistance(t *testing.T) { + gdb := balanceSetup(t) + // Rider 1: nearest, but has had 2 orders this session (both now closed). + onDuty(t, gdb, 1, time.Now().Add(-2*time.Hour)) + for i := 0; i < 2; i++ { + order(t, gdb, 1, constants.BookingConvertedConsignment, constants.AssignmentCompleted, time.Now().Add(-time.Hour)) + } + // Riders 2 and 3: no orders this session; 2 is nearer than 3. + onDuty(t, gdb, 2, time.Now().Add(-2*time.Hour)) + onDuty(t, gdb, 3, time.Now().Add(-2*time.Hour)) + + b := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{}) + id, err := assignOne(t, gdb, b, near(1, 2, 3)) + if err != nil { + t.Fatal(err) + } + if id != 2 { + t.Fatalf("went to %d, want 2 (fewest this session, then nearest)", id) + } +} + +// Everyone at the ceiling: the order waits (no assignment, no error that +// would stop the retries). +func TestAllRidersAtCeilingWaits(t *testing.T) { + gdb := balanceSetup(t) + t.Setenv("MILER_MAX_ACTIVE_BOOKINGS", "2") + for _, id := range []int{1, 2} { + onDuty(t, gdb, id, time.Now().Add(-time.Hour)) + for i := 0; i < 2; i++ { + order(t, gdb, id, constants.BookingMilerAssigned, constants.AssignmentAccepted, time.Now()) + } + } + b := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{}) + id, err := assignOne(t, gdb, b, near(1, 2)) + if err != nil && !errors.Is(err, errNoRiderCapacity) { + t.Fatal(err) + } + if id != 0 { + t.Fatalf("assigned to %d although everyone is at the ceiling", id) + } + var bk models.PickupBooking + mustDo(t, gdb.First(&bk, b).Error) + if bk.Status != constants.BookingPendingPickup || bk.Assignedmileruserid != nil { + t.Fatalf("booking should still be pending: %+v", bk.Status) + } +} diff --git a/internal/assignment/balance_test.go b/internal/assignment/balance_test.go new file mode 100644 index 0000000..e9337d3 --- /dev/null +++ b/internal/assignment/balance_test.go @@ -0,0 +1,59 @@ +package assignment + +import ( + "testing" + + "doormile/models" +) + +func cand(id int, inHand, session int64, km, rating float64) *milerCandidate { + return &milerCandidate{ + profile: models.MilerProfile{Userid: id, Rating: rating}, + distanceKm: km, + activeBookings: inHand, + sessionStops: session, + } +} + +// The balancing order: in hand, then this session, then distance, then rating. +func TestPickBestBalancesBeforeDistance(t *testing.T) { + cases := []struct { + name string + in []*milerCandidate + want int + }{ + {"fewer in hand beats nearer", []*milerCandidate{cand(1, 2, 0, 0.2, 5), cand(2, 1, 9, 9.0, 1)}, 2}, + {"tie in hand: fewer this session", []*milerCandidate{cand(1, 1, 5, 0.2, 5), cand(2, 1, 2, 4.0, 3)}, 2}, + {"tie in hand and session: nearer", []*milerCandidate{cand(1, 1, 2, 3.0, 5), cand(2, 1, 2, 1.0, 3)}, 2}, + {"full tie: better rating", []*milerCandidate{cand(1, 1, 2, 1.0, 4.1), cand(2, 1, 2, 1.0, 4.8)}, 2}, + {"single candidate", []*milerCandidate{cand(7, 3, 3, 3, 3)}, 7}, + } + for _, c := range cases { + if got := pickBestFromCandidates(c.in); got == nil || got.profile.Userid != c.want { + t.Errorf("%s: got %v, want rider %d", c.name, got, c.want) + } + } + if pickBestFromCandidates(nil) != nil { + t.Error("no candidates must give nil") + } +} + +// leastLoaded keeps only the riders with the fewest in hand, with their +// parallel AI rows. +func TestLeastLoaded(t *testing.T) { + cs := []*milerCandidate{cand(1, 2, 0, 1, 4), cand(2, 0, 0, 2, 4), cand(3, 1, 0, 3, 4), cand(4, 0, 5, 4, 4)} + ai := []aiCandidate{{MilerID: 1}, {MilerID: 2}, {MilerID: 3}, {MilerID: 4}} + gotC, gotAI := leastLoaded(cs, ai) + if len(gotC) != 2 || gotC[0].profile.Userid != 2 || gotC[1].profile.Userid != 4 { + t.Fatalf("least loaded = %v", gotC) + } + if len(gotAI) != 2 || gotAI[0].MilerID != 2 || gotAI[1].MilerID != 4 { + t.Fatalf("AI rows must follow: %v", gotAI) + } + if c, a := leastLoaded(nil, nil); len(c) != 0 || len(a) != 0 { + t.Fatal("empty in, empty out") + } + if c, _ := leastLoaded(cs, nil); len(c) != 2 { + t.Fatal("AI rows are optional") + } +} diff --git a/internal/assignment/crm_assignment.go b/internal/assignment/crm_assignment.go index eb4490f..792bd7c 100644 --- a/internal/assignment/crm_assignment.go +++ b/internal/assignment/crm_assignment.go @@ -68,9 +68,16 @@ func maxActiveBookings() int64 { } type milerCandidate struct { - profile models.MilerProfile - distanceKm float64 + profile models.MilerProfile + distanceKm float64 + // activeBookings is the rider's load "in hand": today's open orders that + // are not cancelled (openStopsToday). Balancing gives the next order to + // the rider with the fewest; it is also what the ceiling caps. activeBookings int64 + // sessionStops counts orders given to the rider since they started duty — + // the first tie-break, so a rider who just came online is not passed over + // for one who has been busy all day. + sessionStops int64 } // AssignCRMMiler finds the best available nearby miler for an express booking @@ -144,14 +151,17 @@ func tryAssign(bookingID int) (bool, error) { return false, nil } - candidate, agentDecisionID, found := selectMilerWithAI(&booking, nearby) + candidate, pool, agentDecisionID, found := selectMilerWithAI(&booking, nearby) if !found { return false, nil } - if err := commitAssignment(&booking, candidate, agentDecisionID); err != nil { - if errors.Is(err, errBookingTaken) { + if _, err := commitAssignment(&booking, candidate, pool, agentDecisionID); err != nil { + switch { + case errors.Is(err, errBookingTaken): return true, nil + case errors.Is(err, errNoRiderCapacity): + return false, nil // everyone filled up meanwhile; retry later } return false, fmt.Errorf("commit: %w", err) } @@ -200,7 +210,7 @@ func TryAssignOnce(bookingID int) (AutoAssignResult, error) { return AutoAssignResult{Escalated: true, Reasoning: "no milers within search radius", SearchedRadiusKm: geoRadiusKm}, nil } - candidate, agentDecisionID, found := selectMilerWithAI(&booking, nearby) + candidate, pool, agentDecisionID, found := selectMilerWithAI(&booking, nearby) if !found { return AutoAssignResult{ Escalated: true, @@ -210,12 +220,25 @@ func TryAssignOnce(bookingID int) (AutoAssignResult, error) { }, nil } - if err := commitAssignment(&booking, candidate, agentDecisionID); err != nil { - if errors.Is(err, errBookingTaken) { + chosenID := candidate.profile.Userid + candidate, err = commitAssignment(&booking, candidate, pool, agentDecisionID) + if err != nil { + switch { + case errors.Is(err, errBookingTaken): return AutoAssignResult{Assigned: true}, nil + case errors.Is(err, errNoRiderCapacity): + return AutoAssignResult{ + Escalated: true, + Reasoning: "every nearby rider is at the order limit", + SearchedRadiusKm: geoRadiusKm, + CandidatesFound: len(nearby), + }, nil } return AutoAssignResult{}, fmt.Errorf("commit: %w", err) } + if candidate.profile.Userid != chosenID { + agentDecisionID = nil // balancing overrode the AI's pick; its reasoning no longer applies + } reasoning := "" if agentDecisionID != nil { @@ -252,11 +275,23 @@ func queryNearbyMilers(lat, lon float64) ([]redis.GeoLocation, error) { // commitAssignment writes the BookingAssignment row, updates the booking and the // miler's availability status in a single transaction, then publishes to NATS. -func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate, agentDecisionID *uint64) error { - milerUserID := candidate.profile.Userid - +func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate, pool []*milerCandidate, agentDecisionID *uint64) (*milerCandidate, error) { tx := db.DB.Begin() + // Re-check the choice with fresh loads under the decision lock, so a burst + // of orders spreads across riders instead of all landing on the one that + // looked least loaded when each attempt started. See finalizeChoice. + final, err := finalizeChoice(tx, candidate, pool) + if err != nil { + tx.Rollback() + return nil, err + } + if final != candidate { + agentDecisionID = nil + } + candidate = final + milerUserID := candidate.profile.Userid + // Claim the booking first, and only if it is still unassigned and not // cancelled — see claimBooking. if err := claimBooking(tx, booking.Bookingid, map[string]interface{}{ @@ -265,7 +300,7 @@ func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate, "updatedat": time.Now(), }); err != nil { tx.Rollback() - return err + return nil, err } assignment := models.BookingAssignment{ @@ -274,17 +309,18 @@ func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate, Assignmentstatus: constants.AssignmentAssigned, Assignedat: time.Now(), AgentDecisionID: agentDecisionID, + Remarks: autoAssignedRemark, } if err := tx.Create(&assignment).Error; err != nil { tx.Rollback() - return fmt.Errorf("create BookingAssignment: %w", err) + return nil, fmt.Errorf("create BookingAssignment: %w", err) } if err := tx.Model(&models.MilerProfile{}). Where("userid = ?", milerUserID). Update("availabilitystatus", constants.MilerAssigned).Error; err != nil { tx.Rollback() - return fmt.Errorf("update MilerProfile availability: %w", err) + return nil, fmt.Errorf("update MilerProfile availability: %w", err) } // The customer's "Miler assigned" milestone, in the same transaction as the @@ -298,7 +334,7 @@ func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate, Source: "internal/assignment.commitAssignment", }); err != nil { tx.Rollback() - return fmt.Errorf("record assigned stage: %w", err) + return nil, fmt.Errorf("record assigned stage: %w", err) } tx.Commit() @@ -323,7 +359,7 @@ func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate, // best-effort and off this goroutine's critical path. routing.SequenceMilerStopsAsync(milerUserID) - return nil + return candidate, nil } // publishAssignment sends the booking.assigned event to NATS JetStream. diff --git a/internal/assignment/customer_assignment.go b/internal/assignment/customer_assignment.go index 1290463..c7721c3 100644 --- a/internal/assignment/customer_assignment.go +++ b/internal/assignment/customer_assignment.go @@ -76,7 +76,7 @@ func tryCustomerAssign(bookingID int) (bool, error) { return false, nil } - miler, agentDecisionID, found := selectMilerWithAI(&booking, nearby) + miler, pool, agentDecisionID, found := selectMilerWithAI(&booking, nearby) if !found { return false, nil } @@ -97,13 +97,14 @@ func tryCustomerAssign(bookingID int) (bool, error) { provider = providerResult{company: "Doormile"} } - // Step 3 — ETA based on miler-to-pickup distance. - etaMinutes := calculateETA(miler.distanceKm) - - // Steps 4–6 — Commit to DB and publish to NATS. - if err := commitCustomerAssignment(&booking, miler, provider, etaMinutes, agentDecisionID); err != nil { - if errors.Is(err, errBookingTaken) { + // Steps 3–6 — final balance check, ETA from the chosen rider's distance, + // commit to DB and publish to NATS. + if err := commitCustomerAssignment(&booking, miler, pool, provider, agentDecisionID); err != nil { + switch { + case errors.Is(err, errBookingTaken): return true, nil + case errors.Is(err, errNoRiderCapacity): + return false, nil // everyone filled up meanwhile; retry later } return false, fmt.Errorf("commit: %w", err) } @@ -188,14 +189,25 @@ func calculateETA(distanceKm float64) float64 { func commitCustomerAssignment( booking *models.PickupBooking, miler *milerCandidate, + pool []*milerCandidate, provider providerResult, - etaMinutes float64, agentDecisionID *uint64, ) error { - milerUserID := miler.profile.Userid - tx := db.DB.Begin() + // Final balance check under the decision lock (finalizeChoice). + final, err := finalizeChoice(tx, miler, pool) + if err != nil { + tx.Rollback() + return err + } + if final != miler { + agentDecisionID = nil + } + miler = final + milerUserID := miler.profile.Userid + etaMinutes := calculateETA(miler.distanceKm) + bookingUpdates := map[string]interface{}{ "assignedmileruserid": milerUserID, "status": constants.BookingMilerAssigned, @@ -216,6 +228,7 @@ func commitCustomerAssignment( Assignmentstatus: constants.AssignmentAssigned, Assignedat: time.Now(), AgentDecisionID: agentDecisionID, + Remarks: autoAssignedRemark, } if err := tx.Create(&ba).Error; err != nil { tx.Rollback() diff --git a/internal/assignment/eligibility_pg_test.go b/internal/assignment/eligibility_pg_test.go index b3e0304..39ecf6b 100644 --- a/internal/assignment/eligibility_pg_test.go +++ b/internal/assignment/eligibility_pg_test.go @@ -30,7 +30,7 @@ func eligibilityDB(t *testing.T) *gorm.DB { } gdb := testpg.Open(t, dsn, "assignment_eligibility_test") all := []any{&models.PickupBooking{}, &models.BookingAssignment{}, &models.MilerProfile{}, - &models.BookingStageEvent{}, &models.BookingDestination{}} + &models.BookingStageEvent{}, &models.BookingDestination{}, &models.MilerDutyLog{}} if err := gdb.Migrator().DropTable(all...); err != nil { t.Fatal(err) } @@ -39,7 +39,14 @@ func eligibilityDB(t *testing.T) *gorm.DB { } prev := db.DB db.DB = gdb - t.Cleanup(func() { db.DB = prev }) + // commitAssignment fires customer notifications in a goroutine that can + // outlive the test. Putting back a nil db.DB would make it panic; leaving + // the closed test handle makes it fail quietly instead. + t.Cleanup(func() { + if prev != nil { + db.DB = prev + } + }) t.Setenv("MILER_MAX_ACTIVE_BOOKINGS", "") t.Setenv("ASSIGNMENT_MAX_GPS_AGE_MINUTES", "") return gdb @@ -163,7 +170,7 @@ func TestConcurrentAttemptsAssignOnce(t *testing.T) { go func(i, r int) { defer wg.Done() b := booking - err := commitAssignment(&b, &milerCandidate{profile: models.MilerProfile{Userid: r}}, nil) + _, err := commitAssignment(&b, &milerCandidate{profile: models.MilerProfile{Userid: r}}, nil, nil) mu.Lock() defer mu.Unlock() switch { @@ -198,7 +205,7 @@ func TestCancelledBookingIsNotAssigned(t *testing.T) { mustDo(t, gdb.First(&stale, id).Error) mustDo(t, gdb.Model(&models.PickupBooking{}).Where("bookingid = ?", id).Update("status", constants.BookingCancelled).Error) - err := commitAssignment(&stale, &milerCandidate{profile: models.MilerProfile{Userid: 38}}, nil) + _, err := commitAssignment(&stale, &milerCandidate{profile: models.MilerProfile{Userid: 38}}, nil, nil) if !errors.Is(err, errBookingTaken) { t.Fatalf("want errBookingTaken, got %v", err) } diff --git a/internal/assignment/release.go b/internal/assignment/release.go new file mode 100644 index 0000000..3e4eb5b --- /dev/null +++ b/internal/assignment/release.go @@ -0,0 +1,232 @@ +package assignment + +import ( + "math" + "os" + "strconv" + "strings" + "time" + + "doormile/constants" + "doormile/db" + "doormile/internal/cxstage" + "doormile/models" + "doormile/utils" + + "github.com/redis/go-redis/v9" +) + +// Balanced assignment, Phase B (docs/balanced-assignment-plan.md): riders come +// and go during the day. +// +// - A rider whose app dies right after being given an order never accepts it, +// and the order used to sit with them for good. releaseUnaccepted hands it +// back to the pool after ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES. +// - A rider who starts duty used to wait for the next sweep (up to five +// minutes) before getting anything. SweepNear assigns the pending orders +// around them straight away. +// - An order released by, rejected by or cancelled by a rider is never offered +// to that rider again (ridersToSkip). + +const defaultAcceptTimeoutMinutes = 10 + +// autoAssignedRemark labels the assignments this package makes. Only those +// are ever released: a rider hand-picked by ops or hub staff (whose manual +// path records no "assigned by" for the hub console) is never taken away, and +// nor is anything assigned before the label existed. +const autoAssignedRemark = "Auto-assigned" + +// acceptTimeout: ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES, default 10; 0 turns the +// release off. A typo falls back to the default. +func acceptTimeout() int { + if v := strings.TrimSpace(os.Getenv("ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES")); v != "" { + if n, err := strconv.Atoi(v); err == nil && n >= 0 { + return n + } + utils.Warn("ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES is not a non-negative integer, using the default", + "value", v, "default", defaultAcceptTimeoutMinutes) + } + return defaultAcceptTimeoutMinutes +} + +// releaseUnaccepted returns orders to the pool when the rider they were given +// to by auto-assignment has not accepted within the timeout. Only an +// auto-assignment (autoAssignedRemark) still "Assigned" +// on an order still "Miler_Assigned" qualifies: once the rider accepts +// (Accepted / Pickup_Scheduled) or starts the pickup, the order is theirs and +// is never taken away. Returns how many orders were released. +func releaseUnaccepted() int { + timeout := acceptTimeout() + if timeout == 0 || db.DB == nil { + return 0 + } + type stale struct { + Bookingassignmentid int + Bookingid int + Mileruserid int + } + var rows []stale + if err := db.DB.Table("bookingassignments AS ba"). + Select("ba.bookingassignmentid, ba.bookingid, ba.mileruserid"). + Joins("JOIN pickupbookings pb ON pb.bookingid = ba.bookingid"). + Where("ba.assignmentstatus = ? AND ba.remarks = ? AND pb.status = ? AND pb.assignedmileruserid = ba.mileruserid", + constants.AssignmentAssigned, autoAssignedRemark, constants.BookingMilerAssigned). + Where("ba.assignedat < NOW() - make_interval(mins => ?)", timeout). + Order("ba.assignedat ASC").Limit(sweepBatch). + Scan(&rows).Error; err != nil { + utils.Error("ReleaseUnaccepted: could not list stale assignments", "error", err) + return 0 + } + + released := 0 + for _, r := range rows { + if releaseOne(r.Bookingassignmentid, r.Bookingid, r.Mileruserid, timeout) { + released++ + } + } + if released > 0 { + utils.Info("ReleaseUnaccepted: orders returned to the pool", "released", released, "timeout_minutes", timeout) + } + return released +} + +// releaseOne returns one order to the pool in a single transaction. Every +// write is conditional on the state it was read in, so a rider accepting at +// the same moment wins and nothing is released. +func releaseOne(assignmentID, bookingID, milerUserID, timeout int) bool { + tx := db.DB.Begin() + now := time.Now() + + res := tx.Model(&models.PickupBooking{}). + Where("bookingid = ? AND assignedmileruserid = ? AND status = ?", bookingID, milerUserID, constants.BookingMilerAssigned). + Updates(map[string]interface{}{ + "assignedmileruserid": nil, + "status": constants.BookingPendingPickup, + "updatedat": now, + }) + if res.Error != nil || res.RowsAffected == 0 { + tx.Rollback() + return false + } + res = tx.Model(&models.BookingAssignment{}). + Where("bookingassignmentid = ? AND assignmentstatus = ?", assignmentID, constants.AssignmentAssigned). + Updates(map[string]interface{}{ + "assignmentstatus": constants.AssignmentReassigned, + "remarks": "Not accepted within " + strconv.Itoa(timeout) + " min; returned to the pool", + }) + if res.Error != nil || res.RowsAffected == 0 { + tx.Rollback() + return false + } + // The customer was told "rider assigned"; walk that back so their app does + // not show a rider who is not coming. No-op for console bookings. + if err := cxstage.Release(tx, bookingID, "Rider did not accept in time", + constants.CxActorSystem, nil, "internal/assignment.releaseUnaccepted"); err != nil { + tx.Rollback() + utils.Error("ReleaseUnaccepted: could not update the customer projection", "booking_id", bookingID, "error", err) + return false + } + // Free the rider if this was the only order they held. + var stillOpen int64 + tx.Model(&models.BookingAssignment{}). + Where("mileruserid = ? AND assignmentstatus IN ?", milerUserID, + []string{constants.AssignmentAssigned, constants.AssignmentAccepted}). + Count(&stillOpen) + if stillOpen == 0 { + tx.Model(&models.MilerProfile{}).Where("userid = ? AND availabilitystatus = ?", milerUserID, constants.MilerAssigned). + Update("availabilitystatus", constants.MilerAvailable) + } + if err := tx.Commit().Error; err != nil { + utils.Error("ReleaseUnaccepted: commit failed", "booking_id", bookingID, "error", err) + return false + } + utils.Info("ReleaseUnaccepted: order returned to the pool", "booking_id", bookingID, "miler_id", milerUserID) + return true +} + +// ridersToSkip lists the riders an order must not go back to: the ones it was +// released from, who rejected it, or who cancelled their assignment on it. +func ridersToSkip(bookingID int) map[int]bool { + skip := map[int]bool{} + if db.DB == nil { + return skip + } + var ids []int + db.DB.Model(&models.BookingAssignment{}). + Where("bookingid = ? AND assignmentstatus IN ?", bookingID, + []string{constants.AssignmentReassigned, constants.AssignmentRejected, constants.AssignmentCancelled}). + Distinct().Pluck("mileruserid", &ids) + for _, id := range ids { + skip[id] = true + } + return skip +} + +// withoutRiders drops the given riders from a GEO search result. +func withoutRiders(nearby []redis.GeoLocation, skip map[int]bool) []redis.GeoLocation { + if len(skip) == 0 { + return nearby + } + out := nearby[:0:0] + for _, loc := range nearby { + if id, err := strconv.Atoi(loc.Name); err == nil && skip[id] { + continue + } + out = append(out, loc) + } + return out +} + +// SweepNear assigns the pending orders around a rider who just started duty, +// instead of leaving them for the next periodic sweep. Call it in a goroutine. +func SweepNear(lat, lon float64) { + defer func() { + if r := recover(); r != nil { + utils.Error("SweepNear: panic recovered", "error", r) + } + }() + if db.DB == nil || (lat == 0 && lon == 0) { + return + } + pending, err := pendingNear(lat, lon, time.Now()) + if err != nil { + utils.Error("SweepNear: could not list pending bookings", "error", err) + return + } + assigned := 0 + for _, b := range pending { + if ok, err := attemptOnce(b.Bookingid, kindFor(b.Bookingsource)); err == nil && ok { + assigned++ + } + } + if len(pending) > 0 { + utils.Info("SweepNear: rider came online, swept nearby pending bookings", + "pending", len(pending), "assigned_or_done", assigned) + } +} + +// pendingNear: the sweeper's pending bookings whose pickup lies within the +// assignment search radius of (lat, lon). +func pendingNear(lat, lon float64, now time.Time) ([]models.PickupBooking, error) { + all, err := pendingForSweep(now) + if err != nil { + return nil, err + } + var near []models.PickupBooking + for _, b := range all { + if distanceKm(lat, lon, b.Pickuplatitude, b.Pickuplongitude) <= geoRadiusKm { + near = append(near, b) + } + } + return near, nil +} + +// distanceKm is the great-circle distance between two points. +func distanceKm(lat1, lon1, lat2, lon2 float64) float64 { + const r = 6371.0 + toRad := func(d float64) float64 { return d * math.Pi / 180 } + dLat, dLon := toRad(lat2-lat1), toRad(lon2-lon1) + a := math.Sin(dLat/2)*math.Sin(dLat/2) + + math.Cos(toRad(lat1))*math.Cos(toRad(lat2))*math.Sin(dLon/2)*math.Sin(dLon/2) + return 2 * r * math.Asin(math.Sqrt(a)) +} diff --git a/internal/assignment/release_pg_test.go b/internal/assignment/release_pg_test.go new file mode 100644 index 0000000..a564dbe --- /dev/null +++ b/internal/assignment/release_pg_test.go @@ -0,0 +1,175 @@ +package assignment + +import ( + "testing" + "time" + + "doormile/constants" + "doormile/models" + + "gorm.io/gorm" +) + +// Balanced assignment, Phase B, against a real Postgres. Skipped unless +// REGISTRY_TEST_DSN is set — see eligibility_pg_test.go. + +// autoOrder is order() for an assignment made by auto-assignment (labelled), +// the only kind the release touches. +func autoOrder(t *testing.T, gdb *gorm.DB, rider int, bookingStatus, asgStatus string, at time.Time) int { + t.Helper() + id := order(t, gdb, rider, bookingStatus, asgStatus, at) + mustDo(t, gdb.Model(&models.BookingAssignment{}).Where("bookingid = ?", id).Update("remarks", autoAssignedRemark).Error) + return id +} + +func bookingState(t *testing.T, gdb *gorm.DB, id int) (string, *int) { + t.Helper() + var b models.PickupBooking + mustDo(t, gdb.First(&b, id).Error) + return b.Status, b.Assignedmileruserid +} + +func assignmentState(t *testing.T, gdb *gorm.DB, bookingID int) string { + t.Helper() + var a models.BookingAssignment + mustDo(t, gdb.Where("bookingid = ?", bookingID).Order("bookingassignmentid DESC").First(&a).Error) + return a.Assignmentstatus +} + +// An order the rider never accepted goes back to the pool after the timeout; +// one they accepted, or one still within the timeout, stays with them. +func TestReleaseUnacceptedOrders(t *testing.T) { + gdb := balanceSetup(t) + t.Setenv("ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES", "10") + onDuty(t, gdb, 1, time.Now().Add(-time.Hour)) + mustDo(t, gdb.Model(&models.MilerProfile{}).Where("userid = 1").Update("availabilitystatus", constants.MilerAssigned).Error) + + stale := autoOrder(t, gdb, 1, constants.BookingMilerAssigned, constants.AssignmentAssigned, time.Now().Add(-15*time.Minute)) + fresh := autoOrder(t, gdb, 1, constants.BookingMilerAssigned, constants.AssignmentAssigned, time.Now().Add(-3*time.Minute)) + accepted := autoOrder(t, gdb, 1, constants.BookingPickupScheduled, constants.AssignmentAccepted, time.Now().Add(-40*time.Minute)) + + if n := releaseUnaccepted(); n != 1 { + t.Fatalf("released %d, want 1", n) + } + if st, r := bookingState(t, gdb, stale); st != constants.BookingPendingPickup || r != nil { + t.Fatalf("stale order: status %s rider %v, want Pending_Pickup with no rider", st, r) + } + if s := assignmentState(t, gdb, stale); s != constants.AssignmentReassigned { + t.Fatalf("stale assignment = %s, want Reassigned", s) + } + if st, r := bookingState(t, gdb, fresh); st != constants.BookingMilerAssigned || r == nil { + t.Fatal("an order still inside the timeout must stay with the rider") + } + if st, r := bookingState(t, gdb, accepted); st != constants.BookingPickupScheduled || r == nil { + t.Fatal("an accepted order must never be released") + } + // The rider still holds two orders, so stays Assigned. + var p models.MilerProfile + mustDo(t, gdb.Where("userid = 1").First(&p).Error) + if p.Availabilitystatus != constants.MilerAssigned { + t.Fatalf("rider status %s, want Assigned (still holding orders)", p.Availabilitystatus) + } + + // A rider hand-picked by staff (no auto label) is never taken away. + manual := order(t, gdb, 1, constants.BookingMilerAssigned, constants.AssignmentAssigned, time.Now().Add(-2*time.Hour)) + if n := releaseUnaccepted(); n != 0 { + t.Fatal("a manually assigned order must never be released") + } + if st, r := bookingState(t, gdb, manual); st != constants.BookingMilerAssigned || r == nil { + t.Fatal("manual order changed") + } + + t.Setenv("ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES", "0") + autoOrder(t, gdb, 1, constants.BookingMilerAssigned, constants.AssignmentAssigned, time.Now().Add(-2*time.Hour)) + if n := releaseUnaccepted(); n != 0 { + t.Fatal("timeout 0 must switch the release off") + } +} + +// A rider freed of their only order goes back to Available. +func TestReleaseFreesARiderWithNothingElse(t *testing.T) { + gdb := balanceSetup(t) + onDuty(t, gdb, 1, time.Now().Add(-time.Hour)) + mustDo(t, gdb.Model(&models.MilerProfile{}).Where("userid = 1").Update("availabilitystatus", constants.MilerAssigned).Error) + autoOrder(t, gdb, 1, constants.BookingMilerAssigned, constants.AssignmentAssigned, time.Now().Add(-30*time.Minute)) + if n := releaseUnaccepted(); n != 1 { + t.Fatalf("released %d, want 1", n) + } + var p models.MilerProfile + mustDo(t, gdb.Where("userid = 1").First(&p).Error) + if p.Availabilitystatus != constants.MilerAvailable { + t.Fatalf("rider status %s, want Available", p.Availabilitystatus) + } +} + +// A rider accepting at the same moment wins: nothing is released. +func TestReleaseLosesToAConcurrentAccept(t *testing.T) { + gdb := balanceSetup(t) + onDuty(t, gdb, 1, time.Now().Add(-time.Hour)) + b := order(t, gdb, 1, constants.BookingMilerAssigned, constants.AssignmentAssigned, time.Now().Add(-30*time.Minute)) + var a models.BookingAssignment + mustDo(t, gdb.Where("bookingid = ?", b).First(&a).Error) + // The rider accepts between the release listing and its write. + mustDo(t, gdb.Model(&models.PickupBooking{}).Where("bookingid = ?", b).Update("status", constants.BookingPickupScheduled).Error) + if releaseOne(a.Bookingassignmentid, b, 1, 10) { + t.Fatal("must not release an order the rider has just accepted") + } + if s := assignmentState(t, gdb, b); s != constants.AssignmentAssigned { + t.Fatalf("assignment changed to %s", s) + } +} + +// A released (or rejected) order is never offered back to that rider: it goes +// to someone else even if the first rider is now the least loaded. +func TestReleasedOrderGoesToAnotherRider(t *testing.T) { + gdb := balanceSetup(t) + onDuty(t, gdb, 1, time.Now().Add(-time.Hour)) + onDuty(t, gdb, 2, time.Now().Add(-time.Hour)) + // Rider 2 is busier, so a plain balance would pick rider 1. + order(t, gdb, 2, constants.BookingMilerAssigned, constants.AssignmentAccepted, time.Now()) + b := autoOrder(t, gdb, 1, constants.BookingMilerAssigned, constants.AssignmentAssigned, time.Now().Add(-30*time.Minute)) + if n := releaseUnaccepted(); n != 1 { + t.Fatalf("released %d, want 1", n) + } + id, err := assignOne(t, gdb, b, near(1, 2)) + if err != nil { + t.Fatal(err) + } + if id != 2 { + t.Fatalf("released order went to rider %d, want 2 (never back to rider 1)", id) + } + // Same rule for a rider who rejected it. + r := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{}) + mustDo(t, gdb.Create(&models.BookingAssignment{Bookingid: r, Mileruserid: 2, Assignmentstatus: constants.AssignmentRejected, Assignedat: time.Now()}).Error) + if id, _ := assignOne(t, gdb, r, near(1, 2)); id != 1 { + t.Fatalf("rejected order went to rider %d, want 1", id) + } +} + +// When a rider starts duty, only the pending orders within reach are swept. +func TestPendingNearOnlyPicksNearbyOrders(t *testing.T) { + gdb := balanceSetup(t) + now := time.Now() + mk := func(no string, lat, lon float64) { + nextBookingID++ + mustDo(t, gdb.Create(&models.PickupBooking{Bookingid: nextBookingID, Bookingno: no, Status: constants.BookingPendingPickup, + Bookingsource: constants.BookingSourceExpress, Pickuplatitude: lat, Pickuplongitude: lon, + Createdat: now.Add(-10 * time.Minute)}).Error) + } + mk("DM-NEAR1", 11.0090, 76.9500) // at the rider + mk("DM-NEAR2", 11.0500, 76.9800) // ~5.5 km + mk("DM-FAR", 11.3000, 77.2000) // ~40 km + mk("DM-CHENNAI", 13.0827, 80.2707) // another city + + got, err := pendingNear(11.0090, 76.9500, now) + mustDo(t, err) + names := map[string]bool{} + for _, b := range got { + var full models.PickupBooking + mustDo(t, gdb.First(&full, b.Bookingid).Error) + names[full.Bookingno] = true + } + if len(got) != 2 || !names["DM-NEAR1"] || !names["DM-NEAR2"] { + t.Fatalf("swept %v, want only the two nearby orders", names) + } +} diff --git a/internal/assignment/release_test.go b/internal/assignment/release_test.go new file mode 100644 index 0000000..a5fedbf --- /dev/null +++ b/internal/assignment/release_test.go @@ -0,0 +1,48 @@ +package assignment + +import ( + "math" + "testing" + + "github.com/redis/go-redis/v9" +) + +func TestAcceptTimeout(t *testing.T) { + cases := map[string]int{ + "": defaultAcceptTimeoutMinutes, + "0": 0, // off + "5": 5, + "-1": defaultAcceptTimeoutMinutes, + "soon": defaultAcceptTimeoutMinutes, // a typo must not switch it off + } + for env, want := range cases { + t.Setenv("ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES", env) + if got := acceptTimeout(); got != want { + t.Errorf("ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES=%q: %d, want %d", env, got, want) + } + } +} + +func TestWithoutRiders(t *testing.T) { + nearby := []redis.GeoLocation{{Name: "1"}, {Name: "2"}, {Name: "3"}, {Name: "x"}} + got := withoutRiders(nearby, map[int]bool{2: true}) + if len(got) != 3 || got[0].Name != "1" || got[1].Name != "3" || got[2].Name != "x" { + t.Fatalf("got %v", got) + } + if len(nearby) != 4 || nearby[1].Name != "2" { + t.Fatal("the input slice must not be modified") + } + if got := withoutRiders(nearby, nil); len(got) != 4 { + t.Fatal("nothing to skip, nothing removed") + } +} + +func TestDistanceKm(t *testing.T) { + // RS Puram to Gandhipuram, Coimbatore: about 2.2 km. + if d := distanceKm(11.0090, 76.9500, 11.0182714, 76.9677744); math.Abs(d-2.17) > 0.2 { + t.Fatalf("distance = %.2f km", d) + } + if d := distanceKm(11, 77, 11, 77); d != 0 { + t.Fatalf("same point = %v", d) + } +} diff --git a/internal/assignment/sweeper.go b/internal/assignment/sweeper.go index c8ead22..683e864 100644 --- a/internal/assignment/sweeper.go +++ b/internal/assignment/sweeper.go @@ -120,6 +120,10 @@ func sweepOnce(interval time.Duration) { } } + // Orders a rider never accepted go back to the pool first, so this same + // sweep can hand them to someone else. + releaseUnaccepted() + pending, err := pendingForSweep(time.Now()) if err != nil { utils.Error("PendingSweeper: could not list pending bookings", "error", err) @@ -151,7 +155,7 @@ func sweepOnce(interval time.Duration) { // ExpressDispatchAgent owns them. func pendingForSweep(now time.Time) ([]models.PickupBooking, error) { q := db.DB.Model(&models.PickupBooking{}). - Select("bookingid", "bookingsource"). + Select("bookingid", "bookingsource", "pickuplatitude", "pickuplongitude"). Where("status = ? AND assignedmileruserid IS NULL", constants.BookingPendingPickup). Where("pickuplatitude <> 0 AND pickuplongitude <> 0"). Where("createdat <= ? AND createdat >= ?", now.Add(-sweepMinAge), now.Add(-sweepMaxAge()))