380 lines
12 KiB
Go
380 lines
12 KiB
Go
package assignment
|
|
|
|
import (
|
|
"encoding/json"
|
|
"os"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"doormile/db"
|
|
"doormile/utils"
|
|
|
|
nats "github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// Assignment retries used to live in a bare goroutine: five attempts, two
|
|
// minutes apart, held together by time.Sleep. That is roughly ten minutes of
|
|
// state kept only in one pod's memory. Any restart — a deploy, an OOM kill, a
|
|
// node drain — silently dropped every retry in flight, and nothing recorded
|
|
// that it had happened. The booking simply stayed unassigned forever.
|
|
//
|
|
// The retry now lives in JetStream. One message per booking is published to
|
|
// ASSIGNMENTS; a durable consumer performs a single attempt per delivery and
|
|
// NAKs with a delay when no miler is available, so JetStream owns the waiting
|
|
// and the redelivery count. A pod dying mid-wait costs nothing: the message is
|
|
// still on the server and another replica picks it up.
|
|
const (
|
|
subjectAssignmentRequested = "booking.assignment_requested"
|
|
assignmentDurable = "assignment-worker"
|
|
assignmentStream = "ASSIGNMENTS"
|
|
|
|
// Must exceed the time a single attempt takes. An attempt does a Redis
|
|
// GEOSEARCH plus the AI call to routemate, which is capped at 5s, so 60s is
|
|
// generous. This is unrelated to retryDelay, which is how long JetStream
|
|
// waits after a NAK before redelivering.
|
|
assignmentAckWait = 60 * time.Second
|
|
)
|
|
|
|
// assignmentKind distinguishes the two entry points, which differ only in which
|
|
// single-attempt function they call.
|
|
const (
|
|
kindCustomer = "customer"
|
|
kindExpress = "express"
|
|
)
|
|
|
|
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
|
|
// retry loop when JetStream is unavailable.
|
|
//
|
|
// The fallback matters: assignment is how a booking reaches a rider, so it must
|
|
// not become dependent on the event bus being up. A NATS outage should make
|
|
// retries non-durable again, which is what we had before, not stop bookings
|
|
// being assigned at all.
|
|
func enqueue(bookingID int, kind string) {
|
|
if db.Js == nil {
|
|
utils.Warn("Assignment: JetStream unavailable, retrying in-process",
|
|
"booking_id", bookingID, "kind", kind)
|
|
runInline(bookingID, kind)
|
|
return
|
|
}
|
|
|
|
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)
|
|
runInline(bookingID, kind)
|
|
return
|
|
}
|
|
|
|
if _, err := db.Js.Publish(subjectAssignmentRequested, data); err != nil {
|
|
utils.Error("Assignment: publish failed, retrying in-process",
|
|
"booking_id", bookingID, "error", err)
|
|
runInline(bookingID, kind)
|
|
return
|
|
}
|
|
|
|
utils.Info("Assignment: queued", "booking_id", bookingID, "kind", kind)
|
|
}
|
|
|
|
// attemptOnce runs exactly one assignment attempt for either entry point.
|
|
func attemptOnce(bookingID int, kind string) (bool, error) {
|
|
if kind == kindCustomer {
|
|
return tryCustomerAssign(bookingID)
|
|
}
|
|
return tryAssign(bookingID)
|
|
}
|
|
|
|
// StartAssignmentWorker consumes queued assignment requests. Safe to run on
|
|
// every replica: they share one durable consumer, so JetStream hands each
|
|
// message to exactly one of them.
|
|
func StartAssignmentWorker() {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
utils.Error("AssignmentWorker: panic recovered, restarting", "error", r)
|
|
time.Sleep(5 * time.Second)
|
|
go StartAssignmentWorker()
|
|
}
|
|
}()
|
|
|
|
if db.Js == nil {
|
|
utils.Warn("AssignmentWorker: JetStream unavailable, worker not started")
|
|
return
|
|
}
|
|
|
|
sub, err := db.Js.PullSubscribe(
|
|
subjectAssignmentRequested,
|
|
assignmentDurable,
|
|
nats.BindStream(assignmentStream),
|
|
nats.AckWait(assignmentAckWait),
|
|
nats.MaxDeliver(maxRetries),
|
|
)
|
|
if err != nil {
|
|
utils.Error("AssignmentWorker: subscribe failed", "error", err)
|
|
return
|
|
}
|
|
|
|
utils.Info("AssignmentWorker: running",
|
|
"subject", subjectAssignmentRequested, "max_attempts", maxRetries)
|
|
|
|
for {
|
|
msgs, err := sub.Fetch(5, nats.MaxWait(5*time.Second))
|
|
if err != nil {
|
|
if err == nats.ErrTimeout {
|
|
continue
|
|
}
|
|
utils.Warn("AssignmentWorker: fetch error", "error", err)
|
|
time.Sleep(time.Second)
|
|
continue
|
|
}
|
|
for _, msg := range msgs {
|
|
handleAssignmentMessage(msg)
|
|
}
|
|
}
|
|
}
|
|
|
|
func handleAssignmentMessage(msg *nats.Msg) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
utils.Error("AssignmentWorker: panic in handler", "error", r)
|
|
_ = msg.NakWithDelay(retryDelay)
|
|
}
|
|
}()
|
|
|
|
var req assignmentRequest
|
|
if err := json.Unmarshal(msg.Data, &req); err != nil {
|
|
utils.Error("AssignmentWorker: invalid payload, dropping", "error", err)
|
|
_ = msg.Term()
|
|
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 {
|
|
attempt = int(md.NumDelivered)
|
|
}
|
|
|
|
utils.Info("AssignmentWorker: attempting assignment",
|
|
"booking_id", req.BookingID, "kind", req.Kind, "attempt", attempt)
|
|
|
|
done, err := attemptOnce(req.BookingID, req.Kind)
|
|
if err != nil {
|
|
utils.Error("AssignmentWorker: attempt error",
|
|
"booking_id", req.BookingID, "attempt", attempt, "error", err)
|
|
} else if done {
|
|
_ = msg.Ack()
|
|
return
|
|
}
|
|
|
|
// 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 {
|
|
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", delay.String(),
|
|
)
|
|
_ = 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) {
|
|
firstQueuedAt := time.Now().Unix()
|
|
failurePublished := false
|
|
|
|
for attempt := 1; ; attempt++ {
|
|
if attempt > 1 {
|
|
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
|
|
}
|
|
|
|
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 — giving up",
|
|
"booking_id", bookingID, "retry_window", retryWindow().String())
|
|
if !failurePublished {
|
|
publishAssignmentFailed(bookingID, reasonNoMilerAvailable)
|
|
}
|
|
}
|