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) } }