updates on the order bulk fix
This commit is contained in:
@@ -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) {
|
||||
|
||||
266
internal/assignment/balance_pg_test.go
Normal file
266
internal/assignment/balance_pg_test.go
Normal file
@@ -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)
|
||||
}
|
||||
}
|
||||
59
internal/assignment/balance_test.go
Normal file
59
internal/assignment/balance_test.go
Normal file
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
232
internal/assignment/release.go
Normal file
232
internal/assignment/release.go
Normal file
@@ -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))
|
||||
}
|
||||
175
internal/assignment/release_pg_test.go
Normal file
175
internal/assignment/release_pg_test.go
Normal file
@@ -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)
|
||||
}
|
||||
}
|
||||
48
internal/assignment/release_test.go
Normal file
48
internal/assignment/release_test.go
Normal file
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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()))
|
||||
|
||||
Reference in New Issue
Block a user