updates on the sweeper and eligibility things in the order
This commit is contained in:
@@ -149,17 +149,16 @@ func collectEligibleCandidates(nearby []redis.GeoLocation) ([]*milerCandidate, [
|
||||
continue
|
||||
}
|
||||
|
||||
// Only today's open stops count towards the cap. An order assigned
|
||||
// yesterday or last month and never closed (rider app not updated,
|
||||
// order abandoned) used to hold the rider at the cap forever, so the
|
||||
// one rider actually working was skipped for every new order.
|
||||
var activeCount int64
|
||||
db.DB.Model(&models.BookingAssignment{}).
|
||||
Where("mileruserid = ? AND assignmentstatus IN ? AND assignedat >= ?", milerUserID, []string{
|
||||
constants.AssignmentAssigned,
|
||||
constants.AssignmentAccepted,
|
||||
}, startOfISTDay(time.Now())).
|
||||
Count(&activeCount)
|
||||
// Only a rider whose app is actually reporting can take an order. A
|
||||
// status left at "Assigned" by an app that stopped weeks ago used to
|
||||
// make that rider look like the least busy one, so orders went to
|
||||
// nobody. Compared in SQL against the database clock, which is right
|
||||
// whatever type the column has.
|
||||
if !milerHasFreshGPS(milerUserID) {
|
||||
continue
|
||||
}
|
||||
|
||||
activeCount := openStopsToday(milerUserID)
|
||||
|
||||
if activeCount >= maxActiveBookings() {
|
||||
continue
|
||||
@@ -189,6 +188,39 @@ func collectEligibleCandidates(nearby []redis.GeoLocation) ([]*milerCandidate, [
|
||||
return candidates, aiCandidates
|
||||
}
|
||||
|
||||
// openStopsToday is what counts towards the per-rider cap: the rider's open
|
||||
// assignments (Assigned / Accepted) made since midnight India time, on orders
|
||||
// that are not cancelled. Before, every open record of any age counted, and
|
||||
// nothing closed a record when ops cancelled its order — so a rider with a
|
||||
// pile of old or cancelled orders sat at the cap for good and every new order
|
||||
// stayed Pending.
|
||||
func openStopsToday(milerUserID int) int64 {
|
||||
var n int64
|
||||
db.DB.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,
|
||||
[]string{constants.AssignmentAssigned, constants.AssignmentAccepted},
|
||||
startOfISTDay(time.Now()),
|
||||
constants.BookingCancelled).
|
||||
Count(&n)
|
||||
return n
|
||||
}
|
||||
|
||||
// milerHasFreshGPS reports whether the rider's last position is recent enough
|
||||
// (ASSIGNMENT_MAX_GPS_AGE_MINUTES, default 15; 0 turns the check off).
|
||||
func milerHasFreshGPS(milerUserID int) bool {
|
||||
maxAge := maxGPSAgeMinutes()
|
||||
if maxAge == 0 {
|
||||
return true
|
||||
}
|
||||
var n int64
|
||||
db.DB.Model(&models.MilerProfile{}).
|
||||
Where("userid = ? AND lastlocationupdatedat >= NOW() - make_interval(mins => ?)", milerUserID, maxAge).
|
||||
Count(&n)
|
||||
return n > 0
|
||||
}
|
||||
|
||||
// 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.
|
||||
|
||||
@@ -3,6 +3,7 @@ package assignment
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"strconv"
|
||||
@@ -18,6 +19,7 @@ import (
|
||||
"doormile/utils"
|
||||
|
||||
"github.com/redis/go-redis/v9"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -33,6 +35,24 @@ const (
|
||||
defaultMaxActive = 3
|
||||
)
|
||||
|
||||
// defaultMaxGPSAgeMinutes: a rider whose app has not reported a position for
|
||||
// longer than this is not offered orders. Override with
|
||||
// ASSIGNMENT_MAX_GPS_AGE_MINUTES; 0 turns the check off.
|
||||
const defaultMaxGPSAgeMinutes = 15
|
||||
|
||||
// maxGPSAgeMinutes is read per call, like maxActiveBookings. A typo falls back
|
||||
// to the default rather than switching the check off.
|
||||
func maxGPSAgeMinutes() int {
|
||||
if v := strings.TrimSpace(os.Getenv("ASSIGNMENT_MAX_GPS_AGE_MINUTES")); v != "" {
|
||||
if n, err := strconv.Atoi(v); err == nil && n >= 0 {
|
||||
return n
|
||||
}
|
||||
utils.Warn("ASSIGNMENT_MAX_GPS_AGE_MINUTES is not a non-negative integer, using the default",
|
||||
"value", v, "default", defaultMaxGPSAgeMinutes)
|
||||
}
|
||||
return defaultMaxGPSAgeMinutes
|
||||
}
|
||||
|
||||
// maxActiveBookings reads the per-miler concurrent-stop cap, read per call so
|
||||
// it can be changed without a redeploy. A non-numeric or non-positive value
|
||||
// falls back to the default rather than uncapping the fleet by typo.
|
||||
@@ -70,6 +90,28 @@ func AssignCRMMiler(bookingID int) {
|
||||
enqueue(bookingID, kindExpress)
|
||||
}
|
||||
|
||||
// errBookingTaken: the booking got a rider (or was cancelled) between this
|
||||
// attempt reading it and committing. Not a failure — the work is done.
|
||||
var errBookingTaken = errors.New("booking already assigned or cancelled")
|
||||
|
||||
// claimBooking sets the booking's rider only if it still has none and is not
|
||||
// cancelled, in the caller's transaction. Two attempts can run for one booking
|
||||
// at once — the retry queue, the pending-order sweeper, a hub "auto-assign"
|
||||
// tap — and without this both would commit and the booking would end up with
|
||||
// two riders.
|
||||
func claimBooking(tx *gorm.DB, bookingID int, updates map[string]interface{}) error {
|
||||
res := tx.Model(&models.PickupBooking{}).
|
||||
Where("bookingid = ? AND assignedmileruserid IS NULL AND status <> ?", bookingID, constants.BookingCancelled).
|
||||
Updates(updates)
|
||||
if res.Error != nil {
|
||||
return fmt.Errorf("update PickupBooking: %w", res.Error)
|
||||
}
|
||||
if res.RowsAffected == 0 {
|
||||
return errBookingTaken
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// tryAssign performs a single attempt: queries Redis GEO, scores candidates, commits.
|
||||
// Returns (true, nil) on success or when the booking no longer needs assignment.
|
||||
// Returns (false, nil) when no eligible miler was found (retry warranted).
|
||||
@@ -108,6 +150,9 @@ func tryAssign(bookingID int) (bool, error) {
|
||||
}
|
||||
|
||||
if err := commitAssignment(&booking, candidate, agentDecisionID); err != nil {
|
||||
if errors.Is(err, errBookingTaken) {
|
||||
return true, nil
|
||||
}
|
||||
return false, fmt.Errorf("commit: %w", err)
|
||||
}
|
||||
|
||||
@@ -166,6 +211,9 @@ func TryAssignOnce(bookingID int) (AutoAssignResult, error) {
|
||||
}
|
||||
|
||||
if err := commitAssignment(&booking, candidate, agentDecisionID); err != nil {
|
||||
if errors.Is(err, errBookingTaken) {
|
||||
return AutoAssignResult{Assigned: true}, nil
|
||||
}
|
||||
return AutoAssignResult{}, fmt.Errorf("commit: %w", err)
|
||||
}
|
||||
|
||||
@@ -209,6 +257,17 @@ func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate,
|
||||
|
||||
tx := db.DB.Begin()
|
||||
|
||||
// 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{}{
|
||||
"assignedmileruserid": milerUserID,
|
||||
"status": constants.BookingMilerAssigned,
|
||||
"updatedat": time.Now(),
|
||||
}); err != nil {
|
||||
tx.Rollback()
|
||||
return err
|
||||
}
|
||||
|
||||
assignment := models.BookingAssignment{
|
||||
Bookingid: booking.Bookingid,
|
||||
Mileruserid: milerUserID,
|
||||
@@ -221,16 +280,6 @@ func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate,
|
||||
return fmt.Errorf("create BookingAssignment: %w", err)
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
if err := tx.Model(booking).Updates(map[string]interface{}{
|
||||
"assignedmileruserid": milerUserID,
|
||||
"status": constants.BookingMilerAssigned,
|
||||
"updatedat": now,
|
||||
}).Error; err != nil {
|
||||
tx.Rollback()
|
||||
return fmt.Errorf("update PickupBooking: %w", err)
|
||||
}
|
||||
|
||||
if err := tx.Model(&models.MilerProfile{}).
|
||||
Where("userid = ?", milerUserID).
|
||||
Update("availabilitystatus", constants.MilerAssigned).Error; err != nil {
|
||||
|
||||
@@ -2,6 +2,7 @@ package assignment
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"math"
|
||||
"strconv"
|
||||
@@ -101,6 +102,9 @@ func tryCustomerAssign(bookingID int) (bool, error) {
|
||||
|
||||
// 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) {
|
||||
return true, nil
|
||||
}
|
||||
return false, fmt.Errorf("commit: %w", err)
|
||||
}
|
||||
|
||||
@@ -192,6 +196,20 @@ func commitCustomerAssignment(
|
||||
|
||||
tx := db.DB.Begin()
|
||||
|
||||
bookingUpdates := map[string]interface{}{
|
||||
"assignedmileruserid": milerUserID,
|
||||
"status": constants.BookingMilerAssigned,
|
||||
"updatedat": time.Now(),
|
||||
}
|
||||
if provider.company != "" {
|
||||
bookingUpdates["providercompany"] = provider.company
|
||||
}
|
||||
// Claim first, only if still unassigned and not cancelled (claimBooking).
|
||||
if err := claimBooking(tx, booking.Bookingid, bookingUpdates); err != nil {
|
||||
tx.Rollback()
|
||||
return err
|
||||
}
|
||||
|
||||
ba := models.BookingAssignment{
|
||||
Bookingid: booking.Bookingid,
|
||||
Mileruserid: milerUserID,
|
||||
@@ -204,20 +222,6 @@ func commitCustomerAssignment(
|
||||
return fmt.Errorf("create BookingAssignment: %w", err)
|
||||
}
|
||||
|
||||
bookingUpdates := map[string]interface{}{
|
||||
"assignedmileruserid": milerUserID,
|
||||
"status": constants.BookingMilerAssigned,
|
||||
"updatedat": time.Now(),
|
||||
}
|
||||
if provider.company != "" {
|
||||
bookingUpdates["providercompany"] = provider.company
|
||||
}
|
||||
|
||||
if err := tx.Model(booking).Updates(bookingUpdates).Error; err != nil {
|
||||
tx.Rollback()
|
||||
return fmt.Errorf("update PickupBooking: %w", err)
|
||||
}
|
||||
|
||||
if err := tx.Model(&models.MilerProfile{}).
|
||||
Where("userid = ?", milerUserID).
|
||||
Update("availabilitystatus", constants.MilerAssigned).Error; err != nil {
|
||||
|
||||
246
internal/assignment/eligibility_pg_test.go
Normal file
246
internal/assignment/eligibility_pg_test.go
Normal file
@@ -0,0 +1,246 @@
|
||||
package assignment
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"doormile/constants"
|
||||
"doormile/db"
|
||||
"doormile/internal/testpg"
|
||||
"doormile/models"
|
||||
|
||||
"github.com/redis/go-redis/v9"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// Who auto-assignment may offer an order to, against a real Postgres. Skipped
|
||||
// unless REGISTRY_TEST_DSN is set; the DSN must be a THROWAWAY database — the
|
||||
// tables below are dropped and recreated in their own schema. See
|
||||
// internal/ai/registry/store_integration_test.go for how to start one.
|
||||
|
||||
func eligibilityDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
dsn := os.Getenv("REGISTRY_TEST_DSN")
|
||||
if dsn == "" {
|
||||
t.Skip("REGISTRY_TEST_DSN not set; skipping Postgres assignment test")
|
||||
}
|
||||
gdb := testpg.Open(t, dsn, "assignment_eligibility_test")
|
||||
all := []any{&models.PickupBooking{}, &models.BookingAssignment{}, &models.MilerProfile{},
|
||||
&models.BookingStageEvent{}, &models.BookingDestination{}}
|
||||
if err := gdb.Migrator().DropTable(all...); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := gdb.AutoMigrate(all...); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
prev := db.DB
|
||||
db.DB = gdb
|
||||
t.Cleanup(func() { db.DB = prev })
|
||||
t.Setenv("MILER_MAX_ACTIVE_BOOKINGS", "")
|
||||
t.Setenv("ASSIGNMENT_MAX_GPS_AGE_MINUTES", "")
|
||||
return gdb
|
||||
}
|
||||
|
||||
func mustDo(t *testing.T, err error) {
|
||||
t.Helper()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
var nextBookingID = 1000
|
||||
|
||||
// order creates a booking and, when rider != 0, an assignment to that rider.
|
||||
func order(t *testing.T, gdb *gorm.DB, rider int, bookingStatus, asgStatus string, assignedAt time.Time) int {
|
||||
t.Helper()
|
||||
nextBookingID++
|
||||
id := nextBookingID
|
||||
b := models.PickupBooking{Bookingid: id, Bookingno: fmt.Sprintf("DM-T%06d", id), Status: bookingStatus,
|
||||
Bookingsource: constants.BookingSourceExpress, Pickuplatitude: 11.0053, Pickuplongitude: 76.9511}
|
||||
if rider != 0 {
|
||||
r := rider
|
||||
b.Assignedmileruserid = &r
|
||||
}
|
||||
mustDo(t, gdb.Create(&b).Error)
|
||||
if rider != 0 {
|
||||
mustDo(t, gdb.Create(&models.BookingAssignment{Bookingid: id, Mileruserid: rider,
|
||||
Assignmentstatus: asgStatus, Assignedat: assignedAt}).Error)
|
||||
}
|
||||
return id
|
||||
}
|
||||
|
||||
func rider(t *testing.T, gdb *gorm.DB, id int, status string, gpsAge time.Duration) {
|
||||
t.Helper()
|
||||
seen := time.Now().Add(-gpsAge)
|
||||
mustDo(t, gdb.Create(&models.MilerProfile{Userid: id, Displayname: fmt.Sprintf("rider %d", id),
|
||||
Phone: fmt.Sprintf("90000%05d", id), Availabilitystatus: status, Lastlocationupdatedat: &seen}).Error)
|
||||
}
|
||||
|
||||
// The question that started this: Rajan has 10 orders from earlier days still
|
||||
// open, and a cancelled one from today. A new order must still be offered to
|
||||
// him — only today's real, open work counts.
|
||||
func TestOldAndCancelledOrdersDoNotBlockARider(t *testing.T) {
|
||||
gdb := eligibilityDB(t)
|
||||
const rajan = 38
|
||||
rider(t, gdb, rajan, constants.MilerAssigned, time.Minute)
|
||||
for i := 0; i < 10; i++ {
|
||||
order(t, gdb, rajan, constants.BookingMilerAssigned, constants.AssignmentAccepted, time.Now().AddDate(0, 0, -(i+1)))
|
||||
}
|
||||
order(t, gdb, rajan, constants.BookingCancelled, constants.AssignmentAssigned, time.Now())
|
||||
order(t, gdb, rajan, constants.BookingMilerAssigned, constants.AssignmentAssigned, time.Now())
|
||||
|
||||
if n := openStopsToday(rajan); n != 1 {
|
||||
t.Fatalf("open stops today = %d, want 1 (10 old + 1 cancelled must not count)", n)
|
||||
}
|
||||
got, _ := collectEligibleCandidates([]redis.GeoLocation{{Name: fmt.Sprint(rajan), Dist: 0.03}})
|
||||
if len(got) != 1 {
|
||||
t.Fatal("Rajan must be offered the new order")
|
||||
}
|
||||
}
|
||||
|
||||
// Today's real work still counts towards the cap.
|
||||
func TestTodaysOpenOrdersStillCountTowardsTheCap(t *testing.T) {
|
||||
gdb := eligibilityDB(t)
|
||||
rider(t, gdb, 6, constants.MilerAssigned, time.Minute)
|
||||
for i := 0; i < 3; i++ {
|
||||
order(t, gdb, 6, constants.BookingMilerAssigned, constants.AssignmentAccepted, time.Now())
|
||||
}
|
||||
if got, _ := collectEligibleCandidates([]redis.GeoLocation{{Name: "6"}}); len(got) != 0 {
|
||||
t.Fatal("a rider with 3 open orders today is at the cap")
|
||||
}
|
||||
t.Setenv("MILER_MAX_ACTIVE_BOOKINGS", "5")
|
||||
if got, _ := collectEligibleCandidates([]redis.GeoLocation{{Name: "6"}}); len(got) != 1 {
|
||||
t.Fatal("raising the cap must let them take more")
|
||||
}
|
||||
}
|
||||
|
||||
// A rider whose app stopped reporting is not offered orders, whatever their
|
||||
// status says — but only while the check is on.
|
||||
func TestStaleGPSRiderIsSkipped(t *testing.T) {
|
||||
gdb := eligibilityDB(t)
|
||||
rider(t, gdb, 23, constants.MilerOnDelivery, 6*24*time.Hour) // last GPS 6 days ago
|
||||
rider(t, gdb, 38, constants.MilerAssigned, 2*time.Minute)
|
||||
nearby := []redis.GeoLocation{{Name: "23", Dist: 0.5}, {Name: "38", Dist: 0.9}}
|
||||
|
||||
got, _ := collectEligibleCandidates(nearby)
|
||||
if len(got) != 1 || got[0].profile.Userid != 38 {
|
||||
t.Fatalf("only the rider with live GPS may be offered the order, got %d candidates", len(got))
|
||||
}
|
||||
t.Setenv("ASSIGNMENT_MAX_GPS_AGE_MINUTES", "0")
|
||||
if got, _ := collectEligibleCandidates(nearby); len(got) != 2 {
|
||||
t.Fatal("with the check off both riders are candidates")
|
||||
}
|
||||
}
|
||||
|
||||
// An offline rider stays out, fresh GPS or not.
|
||||
func TestOfflineRiderIsSkipped(t *testing.T) {
|
||||
gdb := eligibilityDB(t)
|
||||
rider(t, gdb, 46, constants.MilerOffline, time.Minute)
|
||||
if got, _ := collectEligibleCandidates([]redis.GeoLocation{{Name: "46"}}); len(got) != 0 {
|
||||
t.Fatal("offline rider must not be offered orders")
|
||||
}
|
||||
}
|
||||
|
||||
// Two attempts racing on one booking — queue retry, sweeper, hub button —
|
||||
// must end with exactly one rider.
|
||||
func TestConcurrentAttemptsAssignOnce(t *testing.T) {
|
||||
gdb := eligibilityDB(t)
|
||||
rider(t, gdb, 38, constants.MilerAvailable, time.Minute)
|
||||
rider(t, gdb, 21, constants.MilerAvailable, time.Minute)
|
||||
id := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{})
|
||||
|
||||
var booking models.PickupBooking
|
||||
mustDo(t, gdb.First(&booking, id).Error)
|
||||
var wg sync.WaitGroup
|
||||
var mu sync.Mutex
|
||||
won, taken, other := 0, 0, 0
|
||||
for i, r := range []int{38, 21, 38, 21, 38, 21} {
|
||||
wg.Add(1)
|
||||
go func(i, r int) {
|
||||
defer wg.Done()
|
||||
b := booking
|
||||
err := commitAssignment(&b, &milerCandidate{profile: models.MilerProfile{Userid: r}}, nil)
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
switch {
|
||||
case err == nil:
|
||||
won++
|
||||
case errors.Is(err, errBookingTaken):
|
||||
taken++
|
||||
default:
|
||||
other++
|
||||
t.Errorf("attempt %d: %v", i, err)
|
||||
}
|
||||
}(i, r)
|
||||
}
|
||||
wg.Wait()
|
||||
if won != 1 || taken != 5 || other != 0 {
|
||||
t.Fatalf("won=%d taken=%d other=%d, want exactly one winner", won, taken, other)
|
||||
}
|
||||
var n int64
|
||||
gdb.Model(&models.BookingAssignment{}).Where("bookingid = ?", id).Count(&n)
|
||||
if n != 1 {
|
||||
t.Fatalf("assignment rows = %d, want 1", n)
|
||||
}
|
||||
}
|
||||
|
||||
// A cancelled booking is never assigned, even by an attempt that read it
|
||||
// before the cancel.
|
||||
func TestCancelledBookingIsNotAssigned(t *testing.T) {
|
||||
gdb := eligibilityDB(t)
|
||||
rider(t, gdb, 38, constants.MilerAvailable, time.Minute)
|
||||
id := order(t, gdb, 0, constants.BookingPendingPickup, "", time.Time{})
|
||||
var stale models.PickupBooking
|
||||
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)
|
||||
if !errors.Is(err, errBookingTaken) {
|
||||
t.Fatalf("want errBookingTaken, got %v", err)
|
||||
}
|
||||
var n int64
|
||||
gdb.Model(&models.BookingAssignment{}).Where("bookingid = ?", id).Count(&n)
|
||||
if n != 0 {
|
||||
t.Fatal("a cancelled booking must get no assignment")
|
||||
}
|
||||
}
|
||||
|
||||
// The sweeper retries every unassigned pending booking in its window — no
|
||||
// matter how many — and nothing else.
|
||||
func TestSweeperPicksEveryUnassignedPendingBooking(t *testing.T) {
|
||||
gdb := eligibilityDB(t)
|
||||
t.Setenv("EXPRESS_AGENT_ENABLED", "")
|
||||
t.Setenv("ASSIGNMENT_SWEEP_MAX_AGE_HOURS", "")
|
||||
now := time.Now()
|
||||
mk := func(no string, status string, rider *int, lat float64, age time.Duration, source string) {
|
||||
nextBookingID++
|
||||
mustDo(t, gdb.Create(&models.PickupBooking{Bookingid: nextBookingID, Bookingno: no, Status: status,
|
||||
Assignedmileruserid: rider, Bookingsource: source, Pickuplatitude: lat, Pickuplongitude: 76.95,
|
||||
Createdat: now.Add(-age)}).Error)
|
||||
}
|
||||
r := 38
|
||||
for i := 0; i < 25; i++ { // a pile of pending orders: all of them are retried
|
||||
mk(fmt.Sprintf("DM-P%02d", i), constants.BookingPendingPickup, nil, 11.0, time.Duration(10+i)*time.Minute, constants.BookingSourceExpress)
|
||||
}
|
||||
mk("DM-CX", constants.BookingPendingPickup, nil, 11.0, 30*time.Minute, constants.BookingSourceCustomerApp)
|
||||
mk("DM-NEW", constants.BookingPendingPickup, nil, 11.0, 30*time.Second, constants.BookingSourceExpress) // its own first attempt runs
|
||||
mk("DM-OLD", constants.BookingPendingPickup, nil, 11.0, 5*24*time.Hour, constants.BookingSourceExpress) // abandoned
|
||||
mk("DM-ASG", constants.BookingMilerAssigned, &r, 11.0, time.Hour, constants.BookingSourceExpress) // has a rider
|
||||
mk("DM-CAN", constants.BookingCancelled, nil, 11.0, time.Hour, constants.BookingSourceExpress) // cancelled
|
||||
mk("DM-NOLOC", constants.BookingPendingPickup, nil, 0, time.Hour, constants.BookingSourceExpress) // no pickup point
|
||||
|
||||
got, err := pendingForSweep(now)
|
||||
mustDo(t, err)
|
||||
if len(got) != 26 {
|
||||
t.Fatalf("swept %d bookings, want 26 (25 console + 1 customer)", len(got))
|
||||
}
|
||||
t.Setenv("EXPRESS_AGENT_ENABLED", "true")
|
||||
got, _ = pendingForSweep(now)
|
||||
if len(got) != 1 || got[0].Bookingsource != constants.BookingSourceCustomerApp {
|
||||
t.Fatalf("with the express agent on, only customer bookings are swept; got %d", len(got))
|
||||
}
|
||||
}
|
||||
164
internal/assignment/sweeper.go
Normal file
164
internal/assignment/sweeper.go
Normal file
@@ -0,0 +1,164 @@
|
||||
package assignment
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"doormile/constants"
|
||||
"doormile/db"
|
||||
"doormile/models"
|
||||
"doormile/utils"
|
||||
)
|
||||
|
||||
// The pending-order sweeper.
|
||||
//
|
||||
// A booking gets assignment attempts when it is created, and retries for a
|
||||
// limited window (ASSIGNMENT_RETRY_WINDOW_MINUTES on the queue; about ten
|
||||
// minutes on the in-process fallback when NATS is down). A booking whose window
|
||||
// closed while every rider was busy, off duty or blocked then sat in Pending
|
||||
// for good — nothing looked at it again, even once riders were free.
|
||||
//
|
||||
// The sweeper closes that gap: every ASSIGNMENT_SWEEP_SECONDS (default 300) it
|
||||
// makes one assignment attempt for every unassigned pending booking, however
|
||||
// many there are. One attempt per booking per sweep, run in-process, so a
|
||||
// sweep never multiplies queue messages. claimBooking makes a booking that is
|
||||
// assigned meanwhile — by the queue, by hand, by another replica — a no-op.
|
||||
|
||||
const (
|
||||
defaultSweepSeconds = 300
|
||||
defaultSweepMaxAgeHours = 72
|
||||
// sweepBatch caps one sweep's work; the rest are next sweep's.
|
||||
sweepBatch = 200
|
||||
// sweepMinAge leaves a just-created booking to its own first attempt.
|
||||
sweepMinAge = 2 * time.Minute
|
||||
sweepLockKey = "assignment:pending-sweep:lock"
|
||||
)
|
||||
|
||||
// sweepInterval: ASSIGNMENT_SWEEP_SECONDS, default 300; 0 turns the sweeper
|
||||
// off. Read once at start.
|
||||
func sweepInterval() time.Duration {
|
||||
if v := strings.TrimSpace(os.Getenv("ASSIGNMENT_SWEEP_SECONDS")); v != "" {
|
||||
if n, err := strconv.Atoi(v); err == nil && n >= 0 {
|
||||
if n > 0 && n < 30 {
|
||||
n = 30 // a sweep makes one attempt per pending booking; keep it sane
|
||||
}
|
||||
return time.Duration(n) * time.Second
|
||||
}
|
||||
utils.Warn("ASSIGNMENT_SWEEP_SECONDS is not a non-negative integer, using the default",
|
||||
"value", v, "default_seconds", defaultSweepSeconds)
|
||||
}
|
||||
return defaultSweepSeconds * time.Second
|
||||
}
|
||||
|
||||
// sweepMaxAge: ASSIGNMENT_SWEEP_MAX_AGE_HOURS, default 72. Older pending
|
||||
// bookings are left alone — they are almost certainly abandoned, and offering
|
||||
// them to a rider now would send someone to a pickup nobody is waiting at.
|
||||
func sweepMaxAge() time.Duration {
|
||||
if v := strings.TrimSpace(os.Getenv("ASSIGNMENT_SWEEP_MAX_AGE_HOURS")); v != "" {
|
||||
if n, err := strconv.Atoi(v); err == nil && n > 0 {
|
||||
return time.Duration(n) * time.Hour
|
||||
}
|
||||
utils.Warn("ASSIGNMENT_SWEEP_MAX_AGE_HOURS is not a positive integer, using the default",
|
||||
"value", v, "default_hours", defaultSweepMaxAgeHours)
|
||||
}
|
||||
return defaultSweepMaxAgeHours * time.Hour
|
||||
}
|
||||
|
||||
// kindFor picks the assignment path for a booking source: customer-app
|
||||
// bookings go through the B2C path, everything else through the express one —
|
||||
// the same split the create handlers make.
|
||||
func kindFor(source string) string {
|
||||
if source == constants.BookingSourceCustomerApp {
|
||||
return kindCustomer
|
||||
}
|
||||
return kindExpress
|
||||
}
|
||||
|
||||
// sweepSkipsExpress: with EXPRESS_AGENT_ENABLED=true, bulk console bookings
|
||||
// are deliberately left for the ExpressDispatchAgent to batch, so the sweeper
|
||||
// must not assign console bookings behind its back.
|
||||
func sweepSkipsExpress() bool {
|
||||
return strings.EqualFold(os.Getenv("EXPRESS_AGENT_ENABLED"), "true")
|
||||
}
|
||||
|
||||
// StartPendingSweeper runs the sweep loop. Call once at boot, in a goroutine.
|
||||
func StartPendingSweeper() {
|
||||
interval := sweepInterval()
|
||||
if interval == 0 {
|
||||
utils.Info("PendingSweeper: disabled (ASSIGNMENT_SWEEP_SECONDS=0)")
|
||||
return
|
||||
}
|
||||
utils.Info("PendingSweeper: started", "interval", interval.String(), "max_age", sweepMaxAge().String())
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
for range ticker.C {
|
||||
sweepOnce(interval)
|
||||
}
|
||||
}
|
||||
|
||||
// sweepOnce makes one attempt for each eligible pending booking. Only one
|
||||
// replica sweeps at a time (a Redis lock that expires before the next tick);
|
||||
// without Redis every replica sweeps, which claimBooking keeps correct.
|
||||
func sweepOnce(interval time.Duration) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
utils.Error("PendingSweeper: panic recovered", "error", r)
|
||||
}
|
||||
}()
|
||||
if db.DB == nil {
|
||||
return
|
||||
}
|
||||
if db.Rdb != nil {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
got, err := db.Rdb.SetNX(ctx, sweepLockKey, "1", interval-10*time.Second).Result()
|
||||
cancel()
|
||||
if err == nil && !got {
|
||||
return // another replica has this sweep
|
||||
}
|
||||
}
|
||||
|
||||
pending, err := pendingForSweep(time.Now())
|
||||
if err != nil {
|
||||
utils.Error("PendingSweeper: could not list pending bookings", "error", err)
|
||||
return
|
||||
}
|
||||
if len(pending) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
assigned, failed := 0, 0
|
||||
for _, b := range pending {
|
||||
ok, err := attemptOnce(b.Bookingid, kindFor(b.Bookingsource))
|
||||
switch {
|
||||
case err != nil:
|
||||
failed++
|
||||
utils.Warn("PendingSweeper: attempt failed", "booking_id", b.Bookingid, "error", err)
|
||||
case ok:
|
||||
assigned++
|
||||
}
|
||||
}
|
||||
utils.Info("PendingSweeper: swept pending bookings",
|
||||
"pending", len(pending), "assigned_or_done", assigned, "errors", failed,
|
||||
"still_waiting", len(pending)-assigned-failed)
|
||||
}
|
||||
|
||||
// pendingForSweep lists the bookings a sweep retries: pending, no rider, with
|
||||
// a pickup location, created between sweepMaxAge ago and sweepMinAge ago,
|
||||
// oldest first, at most sweepBatch. Console bookings are left out while the
|
||||
// ExpressDispatchAgent owns them.
|
||||
func pendingForSweep(now time.Time) ([]models.PickupBooking, error) {
|
||||
q := db.DB.Model(&models.PickupBooking{}).
|
||||
Select("bookingid", "bookingsource").
|
||||
Where("status = ? AND assignedmileruserid IS NULL", constants.BookingPendingPickup).
|
||||
Where("pickuplatitude <> 0 AND pickuplongitude <> 0").
|
||||
Where("createdat <= ? AND createdat >= ?", now.Add(-sweepMinAge), now.Add(-sweepMaxAge()))
|
||||
if sweepSkipsExpress() {
|
||||
q = q.Where("bookingsource = ?", constants.BookingSourceCustomerApp)
|
||||
}
|
||||
var pending []models.PickupBooking
|
||||
err := q.Order("createdat ASC").Limit(sweepBatch).Find(&pending).Error
|
||||
return pending, err
|
||||
}
|
||||
81
internal/assignment/sweeper_test.go
Normal file
81
internal/assignment/sweeper_test.go
Normal file
@@ -0,0 +1,81 @@
|
||||
package assignment
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"doormile/constants"
|
||||
)
|
||||
|
||||
func TestSweepInterval(t *testing.T) {
|
||||
cases := []struct {
|
||||
env string
|
||||
want time.Duration
|
||||
}{
|
||||
{"", defaultSweepSeconds * time.Second},
|
||||
{"0", 0}, // off
|
||||
{"120", 120 * time.Second},
|
||||
{"5", 30 * time.Second}, // floored: one attempt per pending booking per sweep
|
||||
{"soon", defaultSweepSeconds * time.Second},
|
||||
{"-1", defaultSweepSeconds * time.Second},
|
||||
}
|
||||
for _, c := range cases {
|
||||
t.Setenv("ASSIGNMENT_SWEEP_SECONDS", c.env)
|
||||
if got := sweepInterval(); got != c.want {
|
||||
t.Errorf("ASSIGNMENT_SWEEP_SECONDS=%q: %v, want %v", c.env, got, c.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestSweepMaxAge(t *testing.T) {
|
||||
cases := map[string]time.Duration{
|
||||
"": defaultSweepMaxAgeHours * time.Hour,
|
||||
"24": 24 * time.Hour,
|
||||
"0": defaultSweepMaxAgeHours * time.Hour, // 0 would sweep nothing
|
||||
"two": defaultSweepMaxAgeHours * time.Hour,
|
||||
}
|
||||
for env, want := range cases {
|
||||
t.Setenv("ASSIGNMENT_SWEEP_MAX_AGE_HOURS", env)
|
||||
if got := sweepMaxAge(); got != want {
|
||||
t.Errorf("ASSIGNMENT_SWEEP_MAX_AGE_HOURS=%q: %v, want %v", env, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestKindFor(t *testing.T) {
|
||||
if kindFor(constants.BookingSourceCustomerApp) != kindCustomer {
|
||||
t.Error("customer-app bookings take the customer path")
|
||||
}
|
||||
for _, s := range []string{constants.BookingSourceExpress, "", "anything"} {
|
||||
if kindFor(s) != kindExpress {
|
||||
t.Errorf("%q should take the express path", s)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestSweepSkipsExpressFollowsTheAgentFlag(t *testing.T) {
|
||||
t.Setenv("EXPRESS_AGENT_ENABLED", "")
|
||||
if sweepSkipsExpress() {
|
||||
t.Error("agent off: console bookings are swept")
|
||||
}
|
||||
t.Setenv("EXPRESS_AGENT_ENABLED", "true")
|
||||
if !sweepSkipsExpress() {
|
||||
t.Error("agent on: console bookings are left to the agent")
|
||||
}
|
||||
}
|
||||
|
||||
func TestMaxGPSAgeMinutes(t *testing.T) {
|
||||
cases := map[string]int{
|
||||
"": defaultMaxGPSAgeMinutes,
|
||||
"0": 0, // check off
|
||||
"30": 30,
|
||||
"-5": defaultMaxGPSAgeMinutes,
|
||||
"half": defaultMaxGPSAgeMinutes, // a typo must not switch the check off
|
||||
}
|
||||
for env, want := range cases {
|
||||
t.Setenv("ASSIGNMENT_MAX_GPS_AGE_MINUTES", env)
|
||||
if got := maxGPSAgeMinutes(); got != want {
|
||||
t.Errorf("ASSIGNMENT_MAX_GPS_AGE_MINUTES=%q: %d, want %d", env, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user