233 lines
8.2 KiB
Go
233 lines
8.2 KiB
Go
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))
|
|
}
|