diff --git a/controllers/adminController.go b/controllers/adminController.go index 327b5fa..894d5d8 100644 --- a/controllers/adminController.go +++ b/controllers/adminController.go @@ -2967,6 +2967,8 @@ func AdminCancelBooking(c *fiber.Ctx) error { "booking_id", booking.Bookingid, "error", err) } + closeOpenAssignments(booking.Bookingid, "Order cancelled by Doormile operations") + if booking.Assignedmileruserid != nil { db.DB.Model(&models.MilerProfile{}). Where("userid = ?", *booking.Assignedmileruserid). @@ -3055,6 +3057,8 @@ func AdminBulkCancelBookings(c *fiber.Ctx) error { "booking_id", booking.Bookingid, "error", err) } + closeOpenAssignments(booking.Bookingid, "Order cancelled by Doormile operations (bulk)") + if booking.Assignedmileruserid != nil { db.DB.Model(&models.MilerProfile{}). Where("userid = ?", *booking.Assignedmileruserid). @@ -4118,3 +4122,20 @@ func opsActorID(c *fiber.Ctx) *int { } return nil } + +// closeOpenAssignments closes the rider's open assignment on a booking ops +// cancelled. Without it the record stayed Assigned/Accepted for good, and +// auto-assignment counted it against the rider's cap — enough cancelled +// orders and the rider was never offered another. Best effort: the +// cancellation itself has already been saved. +func closeOpenAssignments(bookingID int, remark string) { + if err := db.DB.Model(&models.BookingAssignment{}). + Where("bookingid = ? AND assignmentstatus IN ?", bookingID, + []string{constants.AssignmentAssigned, constants.AssignmentAccepted}). + Updates(map[string]interface{}{ + "assignmentstatus": constants.AssignmentCancelled, + "remarks": remark, + }).Error; err != nil { + utils.Error("cancel: could not close the rider's assignment", "booking_id", bookingID, "error", err) + } +} diff --git a/controllers/cancel_frees_rider_pg_test.go b/controllers/cancel_frees_rider_pg_test.go new file mode 100644 index 0000000..c494ef5 --- /dev/null +++ b/controllers/cancel_frees_rider_pg_test.go @@ -0,0 +1,35 @@ +package controllers + +import ( + "testing" + + "doormile/constants" + "doormile/models" +) + +// Cancelling an order from the console must close the rider's assignment on +// it. It used to stay Assigned/Accepted for good, and auto-assignment counted +// it against the rider's cap — enough cancelled orders and the rider was never +// offered another. Uses rtoTestDB (consignmentReturn_pg_test.go): skipped +// unless REGISTRY_TEST_DSN points at a throwaway database. +func TestCancelClosesTheRidersAssignment(t *testing.T) { + gdb := rtoTestDB(t) + rider := 38 + must(t, gdb.Create(&models.PickupBooking{Bookingid: 71, Bookingno: "DM-T71", Status: constants.BookingMilerAssigned, + Assignedmileruserid: &rider}).Error) + must(t, gdb.Create(&models.BookingAssignment{Bookingid: 71, Mileruserid: rider, Assignmentstatus: constants.AssignmentAccepted}).Error) + // A closed assignment on the same booking must be left as it is. + must(t, gdb.Create(&models.BookingAssignment{Bookingid: 71, Mileruserid: 21, Assignmentstatus: constants.AssignmentRejected}).Error) + + closeOpenAssignments(71, "Order cancelled by Doormile operations") + + var rows []models.BookingAssignment + must(t, gdb.Where("bookingid = ?", 71).Order("mileruserid").Find(&rows).Error) + got := map[int]string{} + for _, r := range rows { + got[r.Mileruserid] = r.Assignmentstatus + } + if got[38] != constants.AssignmentCancelled || got[21] != constants.AssignmentRejected { + t.Fatalf("assignments after cancel = %v", got) + } +} diff --git a/docs/balanced-assignment-plan.md b/docs/balanced-assignment-plan.md new file mode 100644 index 0000000..ba6f67e --- /dev/null +++ b/docs/balanced-assignment-plan.md @@ -0,0 +1,133 @@ +# Balanced auto-assignment: implementation plan + +**Status:** planned, not started (2026-10-06). +**Goal:** however many orders and riders there are, split the orders **equally** among the riders, with riders logging in and out all day. + +--- + +## 0. Where things stand (context for whoever picks this up) + +### Already built (2026-10-06) + +| Fix | Where | +|---|---| +| Rider search falls back to `GEORADIUS` on Redis < 6.2; `/ready` shows `redis_geo` | `internal/milergeo/` (committed `9b94be1`) | +| Only **today's**, non-cancelled open orders count towards the per-rider limit | `internal/assignment/ai_layer.go` → `openStopsToday` | +| Only riders with **live GPS** (default 15 min) are offered orders | `ai_layer.go` → `milerHasFreshGPS`, setting `ASSIGNMENT_MAX_GPS_AGE_MINUTES` | +| Console cancel (single and bulk) **frees the rider** (closes the assignment) | `controllers/adminController.go` → `closeOpenAssignments` | +| **Pending-order sweeper:** retries every unassigned pending order every 5 min (last 72 h) | `internal/assignment/sweeper.go`, started in `main.go` | +| **No double assignment:** the order is claimed only if it still has no rider and isn't cancelled | `crm_assignment.go` → `claimBooking` (used by both commit paths) | + +Tests: `internal/assignment/eligibility_pg_test.go`, `sweeper_test.go`, `controllers/cancel_frees_rider_pg_test.go` (real-Postgres tests skip unless `REGISTRY_TEST_DSN` is set). + +### How a rider is chosen today (what this plan changes) + +1. Find riders within **10 km** of the pickup (Redis GEO), whose status allows work, with live GPS, holding fewer than `MILER_MAX_ACTIVE_BOOKINGS` (default **3**) of today's open orders. +2. Ask the AI service (`routemate …/decide-assignment`). That endpoint currently returns **404**, so the backend falls back to: + ``` + score = distance_km + 2 × open_orders_today − 0.5 × rating (lowest wins) + ``` + +**The problem:** distance dominates, so riders near the pickups get more orders and the split is not equal. `MILER_MAX_ACTIVE_BOOKINGS` is only a ceiling and can't make it equal. **A code change is required.** + +### Relevant facts found in the code + +- A rider **can't end duty** while holding Assigned/Accepted orders (`MilerEndDuty`). +- If a rider's app dies, their orders **stay with them**: nothing releases an order a rider never accepted. (`InternalReassign` exists in `adminController.go`, but nothing calls it automatically.) +- Nothing reacts when a rider **starts duty**: they only get orders from the next retry or sweep. +- Each duty session is recorded in `milerdutylogs` (`loginat`, `logoutat`). +- The rider app sends GPS about **every 30 s** in the background (`PUT /miler/location`). + +--- + +## 1. The balancing rule + +Riders come and go, so "same number of orders **today**" isn't fair: a rider logging in at 3 PM would get *every* new order until they catch up. + +| Rule | With riders joining and leaving | +|---|---| +| Equal orders today | ❌ late joiners get flooded | +| **Equal orders in hand right now** ✅ | fair at every moment; a new rider gets a fair share of *new* orders; a rider who leaves just drops out | + +**Rule:** each new order goes to the eligible rider holding the **fewest unfinished orders right now**. Ties are broken by: +1. fewest orders **this duty session** (since `milerdutylogs.loginat`); +2. nearest to the pickup; +3. highest rating. + +Balancing happens **among riders near each pickup** (10 km), so each area balances its own riders. + +| Situation | Expected | +|---|---| +| 5 riders, 50 orders | 10 each, at most 1 apart (if the ceiling allows) | +| 23 orders, 5 riders | 5, 5, 5, 4, 4 | +| 3 riders hold 4 each; 2 riders log in | the next 8 orders go to the 2 new riders, then everyone shares | +| A rider goes offline (no GPS) | gets nothing new; the others share | + +**"In hand"** = assignments `Assigned`/`Accepted` on orders that aren't `Cancelled`, `Delivered` or otherwise finished. Note: for hyperlocal parcels the assignment stays open until delivery, which is correct, because the rider is still carrying it. + +--- + +## 2. Code changes (`doormile_backend`) + +### Phase A: equal split (core) · ~½ day + +| File / function | Change | +|---|---| +| `internal/assignment/ai_layer.go` → `collectEligibleCandidates` | For each eligible rider, compute **in-hand count** (replaces the today-only count for ranking; the today-only count can stay as the ceiling check) and **session count** (assignments since the latest open `milerdutylogs.loginat`). Store both on `milerCandidate` / `aiCandidate`. | +| `ai_layer.go` → `selectMilerWithAI` | Before calling the AI, **keep only candidates with the minimum in-hand count**. The AI or fallback then chooses among them, so balance holds even when the AI endpoint returns. | +| `ai_layer.go` → `pickBestFromCandidates` | Replace the weighted formula with ordering by in-hand ↑, session count ↑, distance ↑, rating ↓. | +| `internal/assignment/crm_assignment.go` → `commitAssignment`, `customer_assignment.go` → `commitCustomerAssignment` | **Bulk safety:** 50 orders arriving at once run in parallel and would all pick the same "least-loaded" rider. Inside the commit transaction, lock the rider's `milerprofiles` row (`SELECT … FOR UPDATE`), **recount** their in-hand orders, and refuse (try the next candidate) if they're at the ceiling or no longer the least loaded. Keep `claimBooking` as is. | + +### Phase B: riders joining and leaving · ~½ day + +| File / function | Change | +|---|---| +| `controllers/milerAppController.go` → `MilerStartDuty` | **Rider comes online:** after duty starts (and the GPS is indexed), trigger an immediate sweep of pending orders near them (new helper in `sweeper.go`, e.g. `SweepNear(lat, lon)`). New riders get orders within seconds, not up to 5 min. | +| `internal/assignment/sweeper.go` | **Rider stops responding:** release orders that are still **Assigned** (never Accepted) after `ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES` (new, e.g. 10): close that assignment as `Reassigned`, clear `pickupbookings.assignedmileruserid`, set the status back to `Pending_Pickup`, and let the sweeper reassign. **Never release an Accepted order.** Reuse the logic of `InternalReassign` where possible. | + +### Phase C: settings (server, no code) + +| Setting | Recommended | Purpose | +|---|---|---| +| `MILER_MAX_ACTIVE_BOOKINGS` | **20** on the server (code default stays 3) | safety ceiling only; balancing decides the split | +| `ASSIGNMENT_MAX_GPS_AGE_MINUTES` | 15 (exists) | only riders whose app is running | +| `ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES` | 10 (new, Phase B) | release orders never accepted | +| `ASSIGNMENT_SWEEP_SECONDS` | 300 (exists) | pending-order retry interval | +| `ASSIGNMENT_SWEEP_MAX_AGE_HOURS` | 72 (exists) | older pending orders are left alone | + +--- + +## 3. Tests (real Postgres, same pattern as `eligibility_pg_test.go`) + +1. 5 riders, 50 orders assigned one by one → 10 each, max − min ≤ 1. +2. **50 orders created concurrently** → still balanced, no rider over the ceiling, every order exactly 1 rider. +3. 3 busy riders + 2 who just started duty → new orders go to the newcomers until level. +4. A rider with stale GPS gets nothing new. +5. An order still Assigned after the timeout is released and goes to the least-loaded rider; an **Accepted** order is never released. +6. Riders outside 10 km are never used. +7. A tie on in-hand count is broken by session count, then distance, then rating. +8. Starting duty triggers assignment of nearby pending orders (Phase B). + +Plus a unit test for the new ranking (`pickBestFromCandidates`) and the new setting parser. + +--- + +## 4. Rollout + +1. Deploy **Phase A** with `MILER_MAX_ACTIVE_BOOKINGS=20`. For one day, compare orders per rider (should be within 1 of each other among riders in the same area). +2. Deploy **Phase B**. +3. If riders are sent too far, add a cap such as "prefer balance, but not more than X km further than the nearest eligible rider" (`ASSIGNMENT_BALANCE_MAX_EXTRA_KM`). + +--- + +## 5. Unchanged + +The 10 km radius, the live-GPS rule, cancel freeing the rider, the pending-order sweeper, `claimBooking`, manual assignment (`AssignMilerToBooking`), hub batch assign (`HubBatchAssign`, its own cap of 5), the express dispatch agent, and the customer-app flow. + +--- + +## 6. Open points (decide before or during the work) + +- **Distance vs balance:** do you want the extra-km cap from rollout step 3 from day one? +- **The AI endpoint** (`/api/v1/doormile/decide-assignment` on `routemate.workolik.com`) is missing. Once restored, it will choose only among the least-loaded riders (Phase A guarantees this). +- **Hub batch assign** has its own greedy logic and cap (5). Should it use the same balancing later? diff --git a/internal/assignment/ai_layer.go b/internal/assignment/ai_layer.go index 7314511..0aef135 100644 --- a/internal/assignment/ai_layer.go +++ b/internal/assignment/ai_layer.go @@ -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. diff --git a/internal/assignment/crm_assignment.go b/internal/assignment/crm_assignment.go index 947aba8..eb4490f 100644 --- a/internal/assignment/crm_assignment.go +++ b/internal/assignment/crm_assignment.go @@ -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 { diff --git a/internal/assignment/customer_assignment.go b/internal/assignment/customer_assignment.go index d8cbc65..1290463 100644 --- a/internal/assignment/customer_assignment.go +++ b/internal/assignment/customer_assignment.go @@ -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 { diff --git a/internal/assignment/eligibility_pg_test.go b/internal/assignment/eligibility_pg_test.go new file mode 100644 index 0000000..b3e0304 --- /dev/null +++ b/internal/assignment/eligibility_pg_test.go @@ -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)) + } +} diff --git a/internal/assignment/sweeper.go b/internal/assignment/sweeper.go new file mode 100644 index 0000000..c8ead22 --- /dev/null +++ b/internal/assignment/sweeper.go @@ -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 +} diff --git a/internal/assignment/sweeper_test.go b/internal/assignment/sweeper_test.go new file mode 100644 index 0000000..d534e32 --- /dev/null +++ b/internal/assignment/sweeper_test.go @@ -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) + } + } +} diff --git a/main.go b/main.go index 5af2200..f186ac1 100644 --- a/main.go +++ b/main.go @@ -239,6 +239,9 @@ func main() { // 8. Start the assignment worker. Every replica runs one; they share a // durable consumer, so JetStream hands each booking to exactly one of them. go assignment.StartAssignmentWorker() + // Retries every unassigned pending booking on a timer, so a booking whose + // retry window closed while no rider was free is still picked up later. + go assignment.StartPendingSweeper() // 9. Point the stop sequencer at the Route Optimization API. routing.BaseURL = cfg.RouteOptimizerURL