From 44ba33eda2597484695046263505208960610077 Mon Sep 17 00:00:00 2001 From: dharaneesh-r Date: Fri, 25 Sep 2026 16:30:38 +0530 Subject: [PATCH] updates on the admincontroller and the hubcity fix and queue orders --- controllers/adminController.go | 55 ++++++- controllers/hubCity_test.go | 84 +++++++++++ controllers/milerAppController.go | 5 + controllers/milerController.go | 41 +++-- internal/assignment/queue.go | 193 ++++++++++++++++++++++-- internal/assignment/retrywindow_test.go | 78 ++++++++++ models/users.go | 16 ++ 7 files changed, 441 insertions(+), 31 deletions(-) create mode 100644 controllers/hubCity_test.go create mode 100644 internal/assignment/retrywindow_test.go diff --git a/controllers/adminController.go b/controllers/adminController.go index 4988454..ab8f26c 100644 --- a/controllers/adminController.go +++ b/controllers/adminController.go @@ -1623,9 +1623,57 @@ func GetHubs(c *fiber.Ctx) error { if err := query.Find(&hubs).Error; err != nil { return utils.Internal(c, "failed to fetch hubs") } + attachHubCities(hubs) return utils.List(c, hubs, int64(len(hubs))) } +// attachHubCities fills Hub.City from applocations, in ONE query for the whole +// page rather than one per hub — this list is read on every console page load +// through ZoneContext, so a per-row lookup would be twenty round trips for a +// field that comes from a five-row table. +// +// A hub whose applocationid matches nothing keeps an empty City. That is the +// honest answer, and callers already treat "" as "unknown": ZoneContext skips +// its city comparison rather than matching everything. +func attachHubCities(hubs []models.Hub) { + if len(hubs) == 0 { + return + } + ids := make([]int, 0, len(hubs)) + seen := map[int]bool{} + for _, h := range hubs { + if h.Applocationid != 0 && !seen[h.Applocationid] { + seen[h.Applocationid] = true + ids = append(ids, h.Applocationid) + } + } + if len(ids) == 0 { + return + } + + var locs []models.AppLocation + if err := db.DB.Where("applocationid IN ?", ids).Find(&locs).Error; err != nil { + // A failed lookup leaves every City empty, which is the same state the + // response had before this existed. Refusing the whole hub list because + // one derived label could not be resolved would be worse. + return + } + byID := make(map[int]string, len(locs)) + for _, l := range locs { + byID[l.Applocationid] = l.Applocationname + } + applyHubCities(hubs, byID) +} + +// applyHubCities writes the resolved city onto each hub. Split from the query so +// the mapping — including what happens to a hub whose applocation is missing — +// can be tested without a database. +func applyHubCities(hubs []models.Hub, byID map[int]string) { + for i := range hubs { + hubs[i].City = byID[hubs[i].Applocationid] + } +} + func CreateHub(c *fiber.Ctx) error { req := new(dto.HubCreateRequest) if err := c.BodyParser(req); err != nil { @@ -1659,7 +1707,12 @@ func GetHubDetails(c *fiber.Ctx) error { if err := db.DB.Where("hubid = ? AND deletedat IS NULL", id).First(&hub).Error; err != nil { return utils.NotFound(c, "hub not found") } - return utils.OK(c, hub) + // A one-element slice, because attachHubCities writes THROUGH the slice — + // handing it `[]models.Hub{hub}` would fill a copy and return the original + // with City still empty. + one := []models.Hub{hub} + attachHubCities(one) + return utils.OK(c, one[0]) } func UpdateHub(c *fiber.Ctx) error { diff --git a/controllers/hubCity_test.go b/controllers/hubCity_test.go new file mode 100644 index 0000000..bddc991 --- /dev/null +++ b/controllers/hubCity_test.go @@ -0,0 +1,84 @@ +package controllers + +import ( + "encoding/json" + "testing" + + "doormile/models" +) + +// A hub's city is derived, not stored. `hubs` carries only `applocationid`, and +// for as long as the response carried that and nothing else, every console +// screen that asked a hub what city it was in got undefined: +// +// - ZoneContext.matchesZone compares an order's address text against the +// hub's city. An empty city makes that comparison unreachable, so the only +// matcher left was a 35km radius — and an order stored without coordinates +// then belonged to no zone at all. +// - fetchAppLocations names each city `hub.city || hub.hubname`, so the +// Pricing and report pickers offered "Coimbatore Jupiter Hub" where they +// meant "Coimbatore". +func TestApplyHubCities(t *testing.T) { + locations := map[int]string{1: "Coimbatore", 3: "Bangalore"} + + t.Run("resolves each hub against its applocation", func(t *testing.T) { + hubs := []models.Hub{ + {Hubid: 1, Hubname: "Coimbatore Jupiter Hub", Applocationid: 1}, + {Hubid: 5, Hubname: "Bangalore Earth Hub", Applocationid: 3}, + } + applyHubCities(hubs, locations) + if hubs[0].City != "Coimbatore" || hubs[1].City != "Bangalore" { + t.Fatalf("got %q and %q", hubs[0].City, hubs[1].City) + } + }) + + t.Run("an unknown applocation leaves the city empty, not guessed", func(t *testing.T) { + // Empty is the honest answer and callers already read it as "unknown": + // ZoneContext skips its city comparison rather than matching everything. + hubs := []models.Hub{{Hubid: 99, Hubname: "Somewhere Hub", Applocationid: 42}} + applyHubCities(hubs, locations) + if hubs[0].City != "" { + t.Fatalf("expected empty city, got %q", hubs[0].City) + } + }) + + t.Run("a hub with no applocation at all is left alone", func(t *testing.T) { + hubs := []models.Hub{{Hubid: 98, Hubname: "Orphan Hub"}} + applyHubCities(hubs, locations) + if hubs[0].City != "" { + t.Fatalf("expected empty city, got %q", hubs[0].City) + } + }) + + t.Run("the city follows the applocation, never the name", func(t *testing.T) { + // Hub 22 is a live example of why this is worth asserting: it is named + // "Chennai Comet Hub" but sits on applocationid 1, because CreateCityHub + // copies the creating staff member's location. + hubs := []models.Hub{{Hubid: 22, Hubname: "Chennai Comet Hub", Applocationid: 1}} + applyHubCities(hubs, locations) + if hubs[0].City != "Coimbatore" { + t.Fatalf("city must come from applocationid, got %q", hubs[0].City) + } + }) + + t.Run("an empty list is not an error", func(t *testing.T) { + applyHubCities(nil, locations) + applyHubCities([]models.Hub{}, locations) + }) +} + +// The field has to survive JSON encoding, since the console reads `hub.city`. +// `gorm:"-"` keeps it out of the SQL; it must not also keep it out of the body. +func TestHubCityIsSerialised(t *testing.T) { + b, err := json.Marshal(models.Hub{Hubid: 5, Hubname: "Bangalore Earth Hub", City: "Bangalore"}) + if err != nil { + t.Fatal(err) + } + var back map[string]any + if err := json.Unmarshal(b, &back); err != nil { + t.Fatal(err) + } + if back["city"] != "Bangalore" { + t.Fatalf("city missing from the hub response: %s", b) + } +} diff --git a/controllers/milerAppController.go b/controllers/milerAppController.go index 94a4cb0..3cab450 100644 --- a/controllers/milerAppController.go +++ b/controllers/milerAppController.go @@ -61,6 +61,11 @@ func MilerStartDuty(c *fiber.Ctx) error { }) db.DB.Model(&models.AppUser{}).Where("userid = ?", milerUserID).Update("onduty", 1) + // Make the rider findable by auto-assignment straight away, from the + // position they started duty at, instead of only after the app's first + // location ping. Skipped when the app sent no coordinates. + indexMilerLocation(milerUserID, req.Lat, req.Lon) + return utils.OK(c, fiber.Map{ "dutylogid": dutyLog.Dutylogid, "loginat": dutyLog.Loginat, diff --git a/controllers/milerController.go b/controllers/milerController.go index d01c829..f08fb7a 100644 --- a/controllers/milerController.go +++ b/controllers/milerController.go @@ -352,6 +352,32 @@ func UpdateMilerProfile(c *fiber.Ctx) error { return utils.OK(c, profile) } +// indexMilerLocation records a rider's position where auto-assignment looks +// for riders: the `miler:gps:{userid}` key and the `milers:locations` GEO set +// that internal/assignment GEOSEARCHes around a pickup. +// +// Used by UpdateMilerLocation and MilerStartDuty. Before, only the location +// ping wrote here, so a rider who had tapped "Start duty" was Available in the +// database but invisible to assignment until their app sent its first GPS +// ping — orders created in that gap were never offered to them. +func indexMilerLocation(milerUserID int, lat, lon float64) { + if db.Rdb == nil || lat == 0 || lon == 0 { + return + } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + + redisKey := fmt.Sprintf("miler:gps:%d", milerUserID) + val := fmt.Sprintf("%f,%f", lat, lon) + db.Rdb.Set(ctx, redisKey, val, 30*time.Minute) + + db.Rdb.GeoAdd(ctx, "milers:locations", &redis.GeoLocation{ + Name: strconv.Itoa(milerUserID), + Latitude: lat, + Longitude: lon, + }) +} + func UpdateMilerLocation(c *fiber.Ctx) error { milerUserID := c.Locals("userid").(int) @@ -380,20 +406,7 @@ func UpdateMilerLocation(c *fiber.Ctx) error { return utils.Internal(c, "failed to update location") } - if db.Rdb != nil { - ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) - defer cancel() - - redisKey := fmt.Sprintf("miler:gps:%d", milerUserID) - val := fmt.Sprintf("%f,%f", req.Latitude, req.Longitude) - db.Rdb.Set(ctx, redisKey, val, 30*time.Minute) - - db.Rdb.GeoAdd(ctx, "milers:locations", &redis.GeoLocation{ - Name: strconv.Itoa(milerUserID), - Latitude: req.Latitude, - Longitude: req.Longitude, - }) - } + indexMilerLocation(milerUserID, req.Latitude, req.Longitude) return utils.OK(c, fiber.Map{ "latitude": req.Latitude, diff --git a/internal/assignment/queue.go b/internal/assignment/queue.go index 973522b..1f945e5 100644 --- a/internal/assignment/queue.go +++ b/internal/assignment/queue.go @@ -2,6 +2,9 @@ package assignment import ( "encoding/json" + "os" + "strconv" + "strings" "time" "doormile/db" @@ -43,6 +46,102 @@ const ( type assignmentRequest struct { BookingID int `json:"booking_id"` Kind string `json:"kind"` + + // Extended retry. One "round" is one JetStream message with up to + // maxRetries deliveries (~8 minutes). A round that ends with no miler + // found re-queues the booking as a new round, until the retry window + // runs out. All three are omitempty so messages already on the stream + // (published before this field existed) still decode: they are treated + // as round 1 and their window starts when first seen. + FirstQueuedAt int64 `json:"first_queued_at,omitempty"` // unix seconds, first enqueue + Round int `json:"round,omitempty"` // 1-based + NotBefore int64 `json:"not_before,omitempty"` // unix seconds; wait until then before attempting +} + +// ---- Extended retry -------------------------------------------------------- +// +// Auto-assignment used to give up for good after one round — 5 attempts over +// ~8 minutes. A booking created while every nearby rider was full, on a +// break, or not yet broadcasting GPS then sat in pending_pickup forever, even +// after riders freed up minutes later: nothing ever looked at it again, and +// nothing on the booking said so. That is how DM-664517 got stuck. +// +// Now a failed round re-queues the booking and keeps trying every +// extendedRetryDelay until the retry window closes (default 2h, env +// ASSIGNMENT_RETRY_WINDOW_MINUTES). Each attempt re-reads the booking and +// stops as soon as it is cancelled or has a rider — including one assigned by +// hand — so a manual assignment ends the loop. +// +// maxRetries is deliberately NOT raised to get this. It is also the durable +// consumer's MaxDeliver, which is stored on the NATS server; subscribing with +// a different value fails ("subscribe failed") and the worker would then stop +// assigning anything at all. Re-queuing new rounds keeps the consumer config +// identical. +// +// booking.assignment_failed still fires once, at the end of the FIRST round, +// exactly when it did before — the DispatchAgent and ops alerting keep their +// timing and aren't sent a duplicate for every later round. +const ( + extendedRetryDelay = 5 * time.Minute + defaultRetryWindowMinutes = 120 +) + +// retryWindow reads ASSIGNMENT_RETRY_WINDOW_MINUTES per call, like +// maxActiveBookings, so it can be tuned without a redeploy. 0 restores the old +// single-round behaviour; a non-numeric or negative value falls back to the +// default rather than disabling retries by typo. +func retryWindow() time.Duration { + if v := strings.TrimSpace(os.Getenv("ASSIGNMENT_RETRY_WINDOW_MINUTES")); v != "" { + if n, err := strconv.Atoi(v); err == nil && n >= 0 { + return time.Duration(n) * time.Minute + } + utils.Warn("ASSIGNMENT_RETRY_WINDOW_MINUTES is not a non-negative integer, using the default", + "value", v, "default_minutes", defaultRetryWindowMinutes) + } + return defaultRetryWindowMinutes * time.Minute +} + +// withinRetryWindow reports whether another round may start now. +func withinRetryWindow(firstQueuedAt int64, now time.Time) bool { + if firstQueuedAt == 0 { + return true + } + return now.Sub(time.Unix(firstQueuedAt, 0)) < retryWindow() +} + +// requeueNextRound publishes the booking as a new round that waits +// extendedRetryDelay before its first attempt. Returns false when it could not +// be queued (JetStream down or publish error) — the caller then gives up as +// the old code did, rather than dropping it silently. +func requeueNextRound(req assignmentRequest, now time.Time) bool { + if db.Js == nil { + return false + } + next := req + if next.FirstQueuedAt == 0 { + next.FirstQueuedAt = now.Unix() + } + if next.Round < 1 { + next.Round = 1 + } + next.Round++ + next.NotBefore = now.Add(extendedRetryDelay).Unix() + + data, err := json.Marshal(next) + if err != nil { + utils.Error("Assignment: marshal next round failed", "booking_id", req.BookingID, "error", err) + return false + } + if _, err := db.Js.Publish(subjectAssignmentRequested, data); err != nil { + utils.Error("Assignment: publish next round failed", "booking_id", req.BookingID, "error", err) + return false + } + utils.Warn("Assignment: no miler this round, retrying later", + "booking_id", req.BookingID, + "next_round", next.Round, + "retry_in", extendedRetryDelay.String(), + ) + return true } // enqueue publishes an assignment request, falling back to the old in-process @@ -60,7 +159,12 @@ func enqueue(bookingID int, kind string) { return } - data, err := json.Marshal(assignmentRequest{BookingID: bookingID, Kind: kind}) + data, err := json.Marshal(assignmentRequest{ + BookingID: bookingID, + Kind: kind, + FirstQueuedAt: time.Now().Unix(), + Round: 1, + }) if err != nil { utils.Error("Assignment: marshal failed, retrying in-process", "booking_id", bookingID, "error", err) @@ -149,6 +253,16 @@ func handleAssignmentMessage(msg *nats.Msg) { return } + // A later round waits before its first attempt. JetStream has no delayed + // publish, so the wait is a NAK — it uses one of the round's deliveries, + // leaving maxRetries-1 attempts, which is fine at this cadence. + if req.NotBefore > 0 { + if wait := time.Until(time.Unix(req.NotBefore, 0)); wait > time.Second { + _ = msg.NakWithDelay(wait) + return + } + } + // NumDelivered counts this delivery, so it runs 1..maxRetries. attempt := 1 if md, err := msg.Metadata(); err == nil { @@ -167,52 +281,99 @@ func handleAssignmentMessage(msg *nats.Msg) { return } - // Last delivery: JetStream will not redeliver past MaxDeliver, so the - // terminal failure has to be published here or it never fires at all. Ack - // rather than Nak so the message is not left to expire silently. + // Last delivery of this round: JetStream will not redeliver past + // MaxDeliver, so what happens next is decided here. Ack rather than Nak so + // the message is not left to expire silently. if attempt >= maxRetries { - utils.Error("AssignmentWorker: NO_MILER_AVAILABLE — all attempts exhausted", - "booking_id", req.BookingID, "attempts", attempt) - publishAssignmentFailed(req.BookingID, reasonNoMilerAvailable) + round := req.Round + if round < 1 { + round = 1 + } + if round == 1 { + utils.Error("AssignmentWorker: NO_MILER_AVAILABLE — first round exhausted", + "booking_id", req.BookingID, "attempts", attempt) + publishAssignmentFailed(req.BookingID, reasonNoMilerAvailable) + } + // Only "no miler found" earns another round. A hard error (booking + // gone, no pickup coordinates) won't fix itself by waiting. + now := time.Now() + if err == nil && withinRetryWindow(req.FirstQueuedAt, now) && requeueNextRound(req, now) { + _ = msg.Ack() + return + } + utils.Error("AssignmentWorker: NO_MILER_AVAILABLE — giving up", + "booking_id", req.BookingID, + "rounds", round, + "retry_window", retryWindow().String(), + "last_error", err, + ) _ = msg.Ack() return } + delay := retryDelay + if req.Round > 1 { + delay = extendedRetryDelay + } utils.Warn("AssignmentWorker: no eligible miler, will retry", "booking_id", req.BookingID, + "round", req.Round, "attempt", attempt, "remaining", maxRetries-attempt, - "retry_in", retryDelay.String(), + "retry_in", delay.String(), ) - _ = msg.NakWithDelay(retryDelay) + _ = msg.NakWithDelay(delay) } // runInline is the pre-JetStream behaviour, kept only as the fallback path when // the event bus is down. It holds its retries in memory and does not survive a // restart — which is exactly the weakness the queue exists to fix. func runInline(bookingID int, kind string) { - for attempt := 1; attempt <= maxRetries; attempt++ { + firstQueuedAt := time.Now().Unix() + failurePublished := false + + for attempt := 1; ; attempt++ { if attempt > 1 { - time.Sleep(retryDelay) + if attempt <= maxRetries { + time.Sleep(retryDelay) + } else { + time.Sleep(extendedRetryDelay) + } } utils.Info("Assignment(inline): attempting", "booking_id", bookingID, "attempt", attempt) done, err := attemptOnce(bookingID, kind) if err != nil { + // Same rule as the queue: a hard error doesn't earn the extended + // window, but it still gets the original first-round attempts. utils.Error("Assignment(inline): attempt error", "booking_id", bookingID, "attempt", attempt, "error", err) + if attempt >= maxRetries { + break + } continue } if done { return } - utils.Warn("Assignment(inline): no eligible miler", - "booking_id", bookingID, "attempt", attempt, "remaining", maxRetries-attempt) + if attempt == maxRetries && !failurePublished { + utils.Error("Assignment(inline): NO_MILER_AVAILABLE — first round exhausted", + "booking_id", bookingID, "attempts", attempt) + publishAssignmentFailed(bookingID, reasonNoMilerAvailable) + failurePublished = true + } + if attempt >= maxRetries && !withinRetryWindow(firstQueuedAt, time.Now()) { + break + } + + utils.Warn("Assignment(inline): no eligible miler", "booking_id", bookingID, "attempt", attempt) } - utils.Error("Assignment(inline): NO_MILER_AVAILABLE — all retries exhausted", - "booking_id", bookingID, "max_retries", maxRetries) - publishAssignmentFailed(bookingID, reasonNoMilerAvailable) + utils.Error("Assignment(inline): NO_MILER_AVAILABLE — giving up", + "booking_id", bookingID, "retry_window", retryWindow().String()) + if !failurePublished { + publishAssignmentFailed(bookingID, reasonNoMilerAvailable) + } } diff --git a/internal/assignment/retrywindow_test.go b/internal/assignment/retrywindow_test.go new file mode 100644 index 0000000..6e17c3f --- /dev/null +++ b/internal/assignment/retrywindow_test.go @@ -0,0 +1,78 @@ +package assignment + +import ( + "encoding/json" + "testing" + "time" +) + +func TestRetryWindowDefaultAndOverride(t *testing.T) { + t.Setenv("ASSIGNMENT_RETRY_WINDOW_MINUTES", "") + if got := retryWindow(); got != 120*time.Minute { + t.Fatalf("default window = %v, want 2h", got) + } + + t.Setenv("ASSIGNMENT_RETRY_WINDOW_MINUTES", "45") + if got := retryWindow(); got != 45*time.Minute { + t.Fatalf("override window = %v, want 45m", got) + } + + // 0 restores the old single-round behaviour. + t.Setenv("ASSIGNMENT_RETRY_WINDOW_MINUTES", "0") + if got := retryWindow(); got != 0 { + t.Fatalf("zero window = %v, want 0", got) + } + + // A typo must not disable retries. + for _, bad := range []string{"abc", "-5", "1.5"} { + t.Setenv("ASSIGNMENT_RETRY_WINDOW_MINUTES", bad) + if got := retryWindow(); got != 120*time.Minute { + t.Fatalf("bad value %q gave %v, want the 2h default", bad, got) + } + } +} + +func TestWithinRetryWindow(t *testing.T) { + t.Setenv("ASSIGNMENT_RETRY_WINDOW_MINUTES", "120") + now := time.Now() + + if !withinRetryWindow(now.Add(-10*time.Minute).Unix(), now) { + t.Fatal("10 minutes in should still retry") + } + if withinRetryWindow(now.Add(-121*time.Minute).Unix(), now) { + t.Fatal("121 minutes in should give up") + } + // Messages queued before this change carry no timestamp; they get a + // window starting now rather than being dropped. + if !withinRetryWindow(0, now) { + t.Fatal("a message without first_queued_at should retry") + } + + t.Setenv("ASSIGNMENT_RETRY_WINDOW_MINUTES", "0") + if withinRetryWindow(now.Add(-time.Second).Unix(), now) { + t.Fatal("window 0 should never start another round") + } +} + +// Messages already on the ASSIGNMENTS stream when this ships were published +// with only booking_id and kind. They must still decode and act as round 1. +func TestAssignmentRequestDecodesOldPayload(t *testing.T) { + var req assignmentRequest + if err := json.Unmarshal([]byte(`{"booking_id":664517,"kind":"express"}`), &req); err != nil { + t.Fatalf("old payload failed to decode: %v", err) + } + if req.BookingID != 664517 || req.Kind != kindExpress { + t.Fatalf("decoded %+v", req) + } + if req.Round != 0 || req.FirstQueuedAt != 0 || req.NotBefore != 0 { + t.Fatalf("new fields should be zero on an old payload, got %+v", req) + } +} + +// requeueNextRound must refuse (return false) rather than panic when +// JetStream is not connected, so the caller falls back to giving up cleanly. +func TestRequeueWithoutJetStream(t *testing.T) { + if requeueNextRound(assignmentRequest{BookingID: 1, Kind: kindExpress, Round: 1}, time.Now()) { + t.Fatal("requeue should report false with no JetStream connection") + } +} diff --git a/models/users.go b/models/users.go index 07930c4..e363ebc 100644 --- a/models/users.go +++ b/models/users.go @@ -35,6 +35,22 @@ type Hub struct { Createdby int `json:"createdby" gorm:"column:createdby"` Updatedby int `json:"updatedby" gorm:"column:updatedby"` Deletedat *time.Time `json:"deletedat,omitempty" gorm:"column:deletedat"` + + // City is the applocations row this hub sits in, resolved on read and never + // stored (`gorm:"-"` keeps it out of both the SELECT and the INSERT). + // + // It exists because the console asks a hub what city it is in, and until now + // nothing answered: the hub response carried `applocationid` and no name, so + // `hub.city` was undefined on every screen that reached for it. That is not a + // cosmetic gap: + // + // - ZoneContext.matchesZone falls back to comparing an order's address text + // against the hub's city. An empty city makes that branch unreachable, + // leaving a 35km radius as the only thing still matching. + // - fetchAppLocations names each city `hub.city || hub.hubname`, so the + // Pricing and report pickers offered "Coimbatore Jupiter Hub" where they + // meant "Coimbatore". + City string `json:"city" gorm:"-"` } func (Hub) TableName() string {