AssignCustomerMiler and AssignCRMMiler retried five times, two minutes apart, using time.Sleep inside a bare goroutine — about ten minutes of state held only in one pod's memory. Any restart dropped every retry in flight, and nothing recorded it: the booking just stayed unassigned forever with no failure event, because publishAssignmentFailed only runs at the end of a loop that no longer existed. Deploying during a quiet patch was enough to lose bookings this way, and it happened during today's rollout. Retries now run on the ASSIGNMENTS stream. Each entry point publishes one booking.assignment_requested message; a durable consumer performs a single attempt per delivery and NAKs with retryDelay when no miler is available, so JetStream owns both the waiting and the delivery count. A pod dying mid-wait costs nothing — the message is still on the server and another replica takes it. Behaviour is deliberately unchanged from the caller's side: same five attempts, same two-minute spacing, same publishAssignmentFailed handoff to the DispatchAgent. The failure event is fired explicitly on the last delivery, since JetStream stops redelivering at MaxDeliver and would otherwise let the booking fail silently again. runInline keeps the old loop as a fallback for when JetStream is down. Assignment is how a booking reaches a rider, so it must not become dependent on the event bus: an outage should cost durability, which is what we had before, not stop bookings being assigned at all. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
272 lines
8.5 KiB
Go
272 lines
8.5 KiB
Go
package assignment
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"time"
|
|
|
|
"doormile/constants"
|
|
"doormile/db"
|
|
"doormile/models"
|
|
"doormile/utils"
|
|
|
|
"github.com/redis/go-redis/v9"
|
|
)
|
|
|
|
const (
|
|
maxRetries = 5
|
|
retryDelay = 2 * time.Minute
|
|
geoRadiusKm = 10.0
|
|
geoMaxCount = 10
|
|
maxActive = 3
|
|
)
|
|
|
|
type milerCandidate struct {
|
|
profile models.MilerProfile
|
|
distanceKm float64
|
|
activeBookings int64
|
|
}
|
|
|
|
// AssignCRMMiler finds the best available nearby miler for an express booking
|
|
// and assigns them. Call it after tx.Commit() in CreateExpressBooking.
|
|
//
|
|
// Like the B2C path, the attempt and its retries run on the ASSIGNMENTS stream
|
|
// rather than in this process — see queue.go. Both entry points still reach
|
|
// publishAssignmentFailed on terminal failure, or failures arriving via the
|
|
// console would stay invisible to the DispatchAgent.
|
|
func AssignCRMMiler(bookingID int) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
utils.Error("CRMAssignment: panic recovered", "booking_id", bookingID, "error", r)
|
|
}
|
|
}()
|
|
|
|
enqueue(bookingID, kindExpress)
|
|
}
|
|
|
|
// 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).
|
|
// Returns (false, err) on hard errors (booking missing, DB failure).
|
|
func tryAssign(bookingID int) (bool, error) {
|
|
var booking models.PickupBooking
|
|
if err := db.DB.Preload("Parcels").Preload("ServiceOptions").First(&booking, bookingID).Error; err != nil {
|
|
return false, fmt.Errorf("load booking: %w", err)
|
|
}
|
|
|
|
// If the booking was cancelled or already assigned between retries, stop.
|
|
if booking.Status == constants.BookingCancelled || booking.Assignedmileruserid != nil {
|
|
utils.Info("CRMAssignment: booking no longer needs assignment",
|
|
"booking_id", bookingID,
|
|
"status", booking.Status,
|
|
)
|
|
return true, nil
|
|
}
|
|
|
|
if booking.Pickuplatitude == 0 || booking.Pickuplongitude == 0 {
|
|
return false, fmt.Errorf("booking %d has no pickup coordinates", bookingID)
|
|
}
|
|
|
|
nearby, err := queryNearbyMilers(booking.Pickuplatitude, booking.Pickuplongitude)
|
|
if err != nil {
|
|
utils.Warn("CRMAssignment: GEO query failed", "booking_id", bookingID, "error", err)
|
|
return false, nil
|
|
}
|
|
if len(nearby) == 0 {
|
|
return false, nil
|
|
}
|
|
|
|
candidate, agentDecisionID, found := selectMilerWithAI(&booking, nearby)
|
|
if !found {
|
|
return false, nil
|
|
}
|
|
|
|
if err := commitAssignment(&booking, candidate, agentDecisionID); err != nil {
|
|
return false, fmt.Errorf("commit: %w", err)
|
|
}
|
|
|
|
return true, nil
|
|
}
|
|
|
|
// AutoAssignResult carries enough detail for a synchronous caller (the hub
|
|
// console's manual "auto-assign" trigger) to report what happened, since that
|
|
// caller can't rely on the fire-and-forget logging AssignCRMMiler normally uses.
|
|
type AutoAssignResult struct {
|
|
Assigned bool
|
|
Escalated bool
|
|
MilerUserID int
|
|
MilerName string
|
|
DistanceKm float64
|
|
Reasoning string
|
|
SearchedRadiusKm float64
|
|
CandidatesFound int
|
|
}
|
|
|
|
// TryAssignOnce performs a single assignment attempt (GEOSEARCH → AI decision
|
|
// → commit) with no retry loop, for a booking that wasn't auto-assigned at
|
|
// creation time. AssignCRMMiler/AssignCustomerMiler retry over ~10 minutes,
|
|
// which is too slow for a hub staff member waiting on a synchronous response;
|
|
// this wraps the same single-attempt core (tryAssign) used internally by both.
|
|
func TryAssignOnce(bookingID int) (AutoAssignResult, error) {
|
|
var booking models.PickupBooking
|
|
if err := db.DB.Preload("Parcels").Preload("ServiceOptions").First(&booking, bookingID).Error; err != nil {
|
|
return AutoAssignResult{}, fmt.Errorf("load booking: %w", err)
|
|
}
|
|
|
|
if booking.Status == constants.BookingCancelled || booking.Assignedmileruserid != nil {
|
|
return AutoAssignResult{Assigned: true}, nil
|
|
}
|
|
|
|
if booking.Pickuplatitude == 0 || booking.Pickuplongitude == 0 {
|
|
return AutoAssignResult{}, fmt.Errorf("booking %d has no pickup coordinates", bookingID)
|
|
}
|
|
|
|
nearby, err := queryNearbyMilers(booking.Pickuplatitude, booking.Pickuplongitude)
|
|
if err != nil {
|
|
return AutoAssignResult{Escalated: true, Reasoning: "miler location search unavailable", SearchedRadiusKm: geoRadiusKm}, nil
|
|
}
|
|
if len(nearby) == 0 {
|
|
return AutoAssignResult{Escalated: true, Reasoning: "no milers within search radius", SearchedRadiusKm: geoRadiusKm}, nil
|
|
}
|
|
|
|
candidate, agentDecisionID, found := selectMilerWithAI(&booking, nearby)
|
|
if !found {
|
|
return AutoAssignResult{
|
|
Escalated: true,
|
|
Reasoning: "no eligible miler after evaluation",
|
|
SearchedRadiusKm: geoRadiusKm,
|
|
CandidatesFound: len(nearby),
|
|
}, nil
|
|
}
|
|
|
|
if err := commitAssignment(&booking, candidate, agentDecisionID); err != nil {
|
|
return AutoAssignResult{}, fmt.Errorf("commit: %w", err)
|
|
}
|
|
|
|
reasoning := ""
|
|
if agentDecisionID != nil {
|
|
var ad models.AgentDecision
|
|
if db.DB.Where("id = ?", *agentDecisionID).First(&ad).Error == nil {
|
|
reasoning = ad.Reasoning
|
|
}
|
|
}
|
|
|
|
return AutoAssignResult{
|
|
Assigned: true,
|
|
MilerUserID: candidate.profile.Userid,
|
|
MilerName: candidate.profile.Displayname,
|
|
DistanceKm: candidate.distanceKm,
|
|
Reasoning: reasoning,
|
|
CandidatesFound: len(nearby),
|
|
}, nil
|
|
}
|
|
|
|
// queryNearbyMilers runs GEOSEARCH on milers:locations and returns up to geoMaxCount
|
|
// milers within geoRadiusKm km, sorted nearest-first, with distances populated.
|
|
func queryNearbyMilers(lat, lon float64) ([]redis.GeoLocation, error) {
|
|
if db.Rdb == nil {
|
|
return nil, fmt.Errorf("Redis not available")
|
|
}
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
|
|
defer cancel()
|
|
|
|
locs, err := db.Rdb.GeoSearchLocation(ctx, "milers:locations", &redis.GeoSearchLocationQuery{
|
|
GeoSearchQuery: redis.GeoSearchQuery{
|
|
Longitude: lon,
|
|
Latitude: lat,
|
|
Radius: geoRadiusKm,
|
|
RadiusUnit: "km",
|
|
Sort: "ASC",
|
|
Count: geoMaxCount,
|
|
},
|
|
WithDist: true,
|
|
}).Result()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return locs, nil
|
|
}
|
|
|
|
// commitAssignment writes the BookingAssignment row, updates the booking and the
|
|
// miler's availability status in a single transaction, then publishes to NATS.
|
|
func commitAssignment(booking *models.PickupBooking, candidate *milerCandidate, agentDecisionID *uint64) error {
|
|
milerUserID := candidate.profile.Userid
|
|
|
|
tx := db.DB.Begin()
|
|
|
|
assignment := models.BookingAssignment{
|
|
Bookingid: booking.Bookingid,
|
|
Mileruserid: milerUserID,
|
|
Assignmentstatus: constants.AssignmentAssigned,
|
|
Assignedat: time.Now(),
|
|
AgentDecisionID: agentDecisionID,
|
|
}
|
|
if err := tx.Create(&assignment).Error; err != nil {
|
|
tx.Rollback()
|
|
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 {
|
|
tx.Rollback()
|
|
return fmt.Errorf("update MilerProfile availability: %w", err)
|
|
}
|
|
|
|
tx.Commit()
|
|
|
|
utils.Info("CRMAssignment: assigned",
|
|
"booking_id", booking.Bookingid,
|
|
"miler_id", milerUserID,
|
|
"distance_km", candidate.distanceKm,
|
|
"active_bookings", candidate.activeBookings,
|
|
"agent_decision_id", agentDecisionID,
|
|
)
|
|
|
|
publishAssignment(booking, milerUserID)
|
|
notifyMilerNewAssignment(candidate.profile, booking.Bookingid)
|
|
notifyCustomerMilerAssigned(booking, candidate.profile.Displayname)
|
|
|
|
return nil
|
|
}
|
|
|
|
// publishAssignment sends the booking.assigned event to NATS JetStream.
|
|
// Non-fatal: logs a warning and returns if NATS is unavailable or publish fails.
|
|
func publishAssignment(booking *models.PickupBooking, milerUserID int) {
|
|
if db.Js == nil {
|
|
return
|
|
}
|
|
|
|
payload := map[string]interface{}{
|
|
"booking_id": booking.Bookingid,
|
|
"booking_no": booking.Bookingno,
|
|
"miler_id": milerUserID,
|
|
"provider_company": booking.Providercompany,
|
|
"provider_hub": booking.Providerlocation,
|
|
"assigned_at": time.Now().UnixMilli(),
|
|
}
|
|
|
|
data, err := json.Marshal(payload)
|
|
if err != nil {
|
|
utils.Warn("CRMAssignment: failed to marshal NATS payload", "booking_id", booking.Bookingid, "error", err)
|
|
return
|
|
}
|
|
|
|
if _, err := db.Js.Publish("booking.assigned", data); err != nil {
|
|
utils.Warn("CRMAssignment: NATS publish failed", "booking_id", booking.Bookingid, "error", err)
|
|
}
|
|
}
|