package assignment import ( "math" "os" "strconv" "strings" "time" "doormile/constants" "doormile/db" "doormile/internal/cxstage" "doormile/models" "doormile/utils" "github.com/redis/go-redis/v9" ) // Balanced assignment, Phase B (docs/balanced-assignment-plan.md): riders come // and go during the day. // // - A rider whose app dies right after being given an order never accepts it, // and the order used to sit with them for good. releaseUnaccepted hands it // back to the pool after ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES. // - A rider who starts duty used to wait for the next sweep (up to five // minutes) before getting anything. SweepNear assigns the pending orders // around them straight away. // - An order released by, rejected by or cancelled by a rider is never offered // to that rider again (ridersToSkip). const defaultAcceptTimeoutMinutes = 10 // autoAssignedRemark labels the assignments this package makes. Only those // are ever released: a rider hand-picked by ops or hub staff (whose manual // path records no "assigned by" for the hub console) is never taken away, and // nor is anything assigned before the label existed. const autoAssignedRemark = "Auto-assigned" // acceptTimeout: ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES, default 10; 0 turns the // release off. A typo falls back to the default. func acceptTimeout() int { if v := strings.TrimSpace(os.Getenv("ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES")); v != "" { if n, err := strconv.Atoi(v); err == nil && n >= 0 { return n } utils.Warn("ASSIGNMENT_ACCEPT_TIMEOUT_MINUTES is not a non-negative integer, using the default", "value", v, "default", defaultAcceptTimeoutMinutes) } return defaultAcceptTimeoutMinutes } // releaseUnaccepted returns orders to the pool when the rider they were given // to by auto-assignment has not accepted within the timeout. Only an // auto-assignment (autoAssignedRemark) still "Assigned" // on an order still "Miler_Assigned" qualifies: once the rider accepts // (Accepted / Pickup_Scheduled) or starts the pickup, the order is theirs and // is never taken away. Returns how many orders were released. func releaseUnaccepted() int { timeout := acceptTimeout() if timeout == 0 || db.DB == nil { return 0 } type stale struct { Bookingassignmentid int Bookingid int Mileruserid int } var rows []stale if err := db.DB.Table("bookingassignments AS ba"). Select("ba.bookingassignmentid, ba.bookingid, ba.mileruserid"). Joins("JOIN pickupbookings pb ON pb.bookingid = ba.bookingid"). Where("ba.assignmentstatus = ? AND ba.remarks = ? AND pb.status = ? AND pb.assignedmileruserid = ba.mileruserid", constants.AssignmentAssigned, autoAssignedRemark, constants.BookingMilerAssigned). Where("ba.assignedat < NOW() - make_interval(mins => ?)", timeout). Order("ba.assignedat ASC").Limit(sweepBatch). Scan(&rows).Error; err != nil { utils.Error("ReleaseUnaccepted: could not list stale assignments", "error", err) return 0 } released := 0 for _, r := range rows { if releaseOne(r.Bookingassignmentid, r.Bookingid, r.Mileruserid, timeout) { released++ } } if released > 0 { utils.Info("ReleaseUnaccepted: orders returned to the pool", "released", released, "timeout_minutes", timeout) } return released } // releaseOne returns one order to the pool in a single transaction. Every // write is conditional on the state it was read in, so a rider accepting at // the same moment wins and nothing is released. func releaseOne(assignmentID, bookingID, milerUserID, timeout int) bool { tx := db.DB.Begin() now := time.Now() res := tx.Model(&models.PickupBooking{}). Where("bookingid = ? AND assignedmileruserid = ? AND status = ?", bookingID, milerUserID, constants.BookingMilerAssigned). Updates(map[string]interface{}{ "assignedmileruserid": nil, "status": constants.BookingPendingPickup, "updatedat": now, }) if res.Error != nil || res.RowsAffected == 0 { tx.Rollback() return false } res = tx.Model(&models.BookingAssignment{}). Where("bookingassignmentid = ? AND assignmentstatus = ?", assignmentID, constants.AssignmentAssigned). Updates(map[string]interface{}{ "assignmentstatus": constants.AssignmentReassigned, "remarks": "Not accepted within " + strconv.Itoa(timeout) + " min; returned to the pool", }) if res.Error != nil || res.RowsAffected == 0 { tx.Rollback() return false } // The customer was told "rider assigned"; walk that back so their app does // not show a rider who is not coming. No-op for console bookings. if err := cxstage.Release(tx, bookingID, "Rider did not accept in time", constants.CxActorSystem, nil, "internal/assignment.releaseUnaccepted"); err != nil { tx.Rollback() utils.Error("ReleaseUnaccepted: could not update the customer projection", "booking_id", bookingID, "error", err) return false } // Free the rider if this was the only order they held. var stillOpen int64 tx.Model(&models.BookingAssignment{}). Where("mileruserid = ? AND assignmentstatus IN ?", milerUserID, []string{constants.AssignmentAssigned, constants.AssignmentAccepted}). Count(&stillOpen) if stillOpen == 0 { tx.Model(&models.MilerProfile{}).Where("userid = ? AND availabilitystatus = ?", milerUserID, constants.MilerAssigned). Update("availabilitystatus", constants.MilerAvailable) } if err := tx.Commit().Error; err != nil { utils.Error("ReleaseUnaccepted: commit failed", "booking_id", bookingID, "error", err) return false } utils.Info("ReleaseUnaccepted: order returned to the pool", "booking_id", bookingID, "miler_id", milerUserID) return true } // ridersToSkip lists the riders an order must not go back to: the ones it was // released from, who rejected it, or who cancelled their assignment on it. func ridersToSkip(bookingID int) map[int]bool { skip := map[int]bool{} if db.DB == nil { return skip } var ids []int db.DB.Model(&models.BookingAssignment{}). Where("bookingid = ? AND assignmentstatus IN ?", bookingID, []string{constants.AssignmentReassigned, constants.AssignmentRejected, constants.AssignmentCancelled}). Distinct().Pluck("mileruserid", &ids) for _, id := range ids { skip[id] = true } return skip } // withoutRiders drops the given riders from a GEO search result. func withoutRiders(nearby []redis.GeoLocation, skip map[int]bool) []redis.GeoLocation { if len(skip) == 0 { return nearby } out := nearby[:0:0] for _, loc := range nearby { if id, err := strconv.Atoi(loc.Name); err == nil && skip[id] { continue } out = append(out, loc) } return out } // SweepNear assigns the pending orders around a rider who just started duty, // instead of leaving them for the next periodic sweep. Call it in a goroutine. func SweepNear(lat, lon float64) { defer func() { if r := recover(); r != nil { utils.Error("SweepNear: panic recovered", "error", r) } }() if db.DB == nil || (lat == 0 && lon == 0) { return } pending, err := pendingNear(lat, lon, time.Now()) if err != nil { utils.Error("SweepNear: could not list pending bookings", "error", err) return } assigned := 0 for _, b := range pending { if ok, err := attemptOnce(b.Bookingid, kindFor(b.Bookingsource)); err == nil && ok { assigned++ } } if len(pending) > 0 { utils.Info("SweepNear: rider came online, swept nearby pending bookings", "pending", len(pending), "assigned_or_done", assigned) } } // pendingNear: the sweeper's pending bookings whose pickup lies within the // assignment search radius of (lat, lon). func pendingNear(lat, lon float64, now time.Time) ([]models.PickupBooking, error) { all, err := pendingForSweep(now) if err != nil { return nil, err } var near []models.PickupBooking for _, b := range all { if distanceKm(lat, lon, b.Pickuplatitude, b.Pickuplongitude) <= geoRadiusKm { near = append(near, b) } } return near, nil } // distanceKm is the great-circle distance between two points. func distanceKm(lat1, lon1, lat2, lon2 float64) float64 { const r = 6371.0 toRad := func(d float64) float64 { return d * math.Pi / 180 } dLat, dLon := toRad(lat2-lat1), toRad(lon2-lon1) a := math.Sin(dLat/2)*math.Sin(dLat/2) + math.Cos(toRad(lat1))*math.Cos(toRad(lat2))*math.Sin(dLon/2)*math.Sin(dLon/2) return 2 * r * math.Asin(math.Sqrt(a)) }