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>
219 lines
6.7 KiB
Go
219 lines
6.7 KiB
Go
package assignment
|
|
|
|
import (
|
|
"encoding/json"
|
|
"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"`
|
|
}
|
|
|
|
// 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})
|
|
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
|
|
}
|
|
|
|
// 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: JetStream will not redeliver past MaxDeliver, so the
|
|
// terminal failure has to be published here or it never fires at all. Ack
|
|
// rather than Nak so the message is not left to expire silently.
|
|
if attempt >= maxRetries {
|
|
utils.Error("AssignmentWorker: NO_MILER_AVAILABLE — all attempts exhausted",
|
|
"booking_id", req.BookingID, "attempts", attempt)
|
|
publishAssignmentFailed(req.BookingID, reasonNoMilerAvailable)
|
|
_ = msg.Ack()
|
|
return
|
|
}
|
|
|
|
utils.Warn("AssignmentWorker: no eligible miler, will retry",
|
|
"booking_id", req.BookingID,
|
|
"attempt", attempt,
|
|
"remaining", maxRetries-attempt,
|
|
"retry_in", retryDelay.String(),
|
|
)
|
|
_ = msg.NakWithDelay(retryDelay)
|
|
}
|
|
|
|
// 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) {
|
|
for attempt := 1; attempt <= maxRetries; attempt++ {
|
|
if attempt > 1 {
|
|
time.Sleep(retryDelay)
|
|
}
|
|
|
|
utils.Info("Assignment(inline): attempting", "booking_id", bookingID, "attempt", attempt)
|
|
|
|
done, err := attemptOnce(bookingID, kind)
|
|
if err != nil {
|
|
utils.Error("Assignment(inline): attempt error",
|
|
"booking_id", bookingID, "attempt", attempt, "error", err)
|
|
continue
|
|
}
|
|
if done {
|
|
return
|
|
}
|
|
|
|
utils.Warn("Assignment(inline): no eligible miler",
|
|
"booking_id", bookingID, "attempt", attempt, "remaining", maxRetries-attempt)
|
|
}
|
|
|
|
utils.Error("Assignment(inline): NO_MILER_AVAILABLE — all retries exhausted",
|
|
"booking_id", bookingID, "max_retries", maxRetries)
|
|
publishAssignmentFailed(bookingID, reasonNoMilerAvailable)
|
|
}
|