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 } } // Orders a rider never accepted go back to the pool first, so this same // sweep can hand them to someone else. releaseUnaccepted() 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", "pickuplatitude", "pickuplongitude"). 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 }