fix: assignment retries no longer die with the pod

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>
This commit is contained in:
Suriya
2026-08-10 16:29:32 +05:30
parent 41f0013751
commit 2158031191
5 changed files with 248 additions and 71 deletions

View File

@@ -31,6 +31,10 @@ var streamSubjects = map[string][]string{
"booking.assigned",
"booking.reassigned",
"booking.assignment_failed",
// Drives the assignment retry loop itself, not just a notification that
// one happened — see internal/assignment/queue.go. If this subject is
// missing from the stream, bookings are never assigned at all.
"booking.assignment_requested",
},
"STATUS": {
"booking.status.updated",

View File

@@ -28,9 +28,13 @@ type milerCandidate struct {
activeBookings int64
}
// AssignCRMMiler finds the best available nearby miler for a CRM booking and assigns them.
// It retries up to maxRetries times (retryDelay apart) before logging NO_MILER_AVAILABLE.
// Must be called as a goroutine after tx.Commit() in CreateExpressBooking.
// 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 {
@@ -38,38 +42,7 @@ func AssignCRMMiler(bookingID int) {
}
}()
for attempt := 1; attempt <= maxRetries; attempt++ {
if attempt > 1 {
time.Sleep(retryDelay)
}
utils.Info("CRMAssignment: attempting assignment", "booking_id", bookingID, "attempt", attempt)
done, err := tryAssign(bookingID)
if err != nil {
utils.Error("CRMAssignment: attempt error", "booking_id", bookingID, "attempt", attempt, "error", err)
continue
}
if done {
return
}
utils.Warn("CRMAssignment: no eligible miler found on attempt",
"booking_id", bookingID,
"attempt", attempt,
"remaining", maxRetries-attempt,
)
}
utils.Error("CRMAssignment: NO_MILER_AVAILABLE — all retries exhausted",
"booking_id", bookingID,
"max_retries", maxRetries,
)
// Terminal failure — same handoff as the B2C path. Both entry points must
// publish, or failures arriving via the CRM console stay invisible to the
// DispatchAgent.
publishAssignmentFailed(bookingID, reasonNoMilerAvailable)
enqueue(bookingID, kindExpress)
}
// tryAssign performs a single attempt: queries Redis GEO, scores candidates, commits.

View File

@@ -20,10 +20,18 @@ type providerResult struct {
reliability float64
}
// AssignCustomerMiler is the B2C goroutine entry point.
// It selects both a miler (first-mile pickup) and a provider (delivery routing),
// then commits the assignment. Retries up to maxRetries times with retryDelay in between.
// Must be called as a goroutine after tx.Commit() in CreateCustomerBooking.
// AssignCustomerMiler is the B2C entry point. It selects both a miler
// (first-mile pickup) and a provider (delivery routing), then commits the
// assignment. Call it after tx.Commit() in CreateCustomerBooking.
//
// The attempt and its retries now run on the ASSIGNMENTS stream rather than in
// this process — see queue.go for why. Publishing is fast enough that the
// caller's `go` is no longer strictly needed, but it is harmless and the call
// sites are left as they are.
//
// Terminal failure still hands off to the AI layer's DispatchAgent via
// publishAssignmentFailed, which the worker fires once the last delivery is
// exhausted.
func AssignCustomerMiler(bookingID int) {
defer func() {
if r := recover(); r != nil {
@@ -31,38 +39,7 @@ func AssignCustomerMiler(bookingID int) {
}
}()
for attempt := 1; attempt <= maxRetries; attempt++ {
if attempt > 1 {
time.Sleep(retryDelay)
}
utils.Info("B2CAssignment: attempting assignment", "booking_id", bookingID, "attempt", attempt)
done, err := tryCustomerAssign(bookingID)
if err != nil {
utils.Error("B2CAssignment: attempt error", "booking_id", bookingID, "attempt", attempt, "error", err)
continue
}
if done {
return
}
utils.Warn("B2CAssignment: no eligible miler on attempt",
"booking_id", bookingID,
"attempt", attempt,
"remaining", maxRetries-attempt,
)
}
utils.Error("B2CAssignment: NO_MILER_AVAILABLE — all retries exhausted",
"booking_id", bookingID,
"max_retries", maxRetries,
)
// Terminal failure — hand off to the AI layer's DispatchAgent, which owns
// what happens next (coverage sweep, escalation). Reached only after every
// retry is exhausted, so it fires at most once per booking.
publishAssignmentFailed(bookingID, reasonNoMilerAvailable)
enqueue(bookingID, kindCustomer)
}
// tryCustomerAssign performs one full attempt: GEO miler search → provider selection → commit.

View File

@@ -0,0 +1,218 @@
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)
}

View File

@@ -11,6 +11,7 @@ import (
"doormile/config"
"doormile/controllers"
"doormile/db"
"doormile/internal/assignment"
"doormile/internal/notify"
"doormile/internal/worker"
"doormile/middlewares"
@@ -144,6 +145,10 @@ func main() {
// 7. Start NATS booking worker in background
go worker.StartBookingWorker()
// 8. Start the assignment worker. Every replica runs one; they share a
// durable consumer, so JetStream hands each booking to exactly one of them.
go assignment.StartAssignmentWorker()
// 7. Startup server in a background thread
go func() {
utils.Info("Server starting", "port", cfg.Port)