diff --git a/db/streams.go b/db/streams.go index 40e530e..5b689f9 100644 --- a/db/streams.go +++ b/db/streams.go @@ -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", diff --git a/internal/assignment/crm_assignment.go b/internal/assignment/crm_assignment.go index 3a3252d..c66f972 100644 --- a/internal/assignment/crm_assignment.go +++ b/internal/assignment/crm_assignment.go @@ -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. diff --git a/internal/assignment/customer_assignment.go b/internal/assignment/customer_assignment.go index f8030ee..cb2431b 100644 --- a/internal/assignment/customer_assignment.go +++ b/internal/assignment/customer_assignment.go @@ -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. diff --git a/internal/assignment/queue.go b/internal/assignment/queue.go new file mode 100644 index 0000000..973522b --- /dev/null +++ b/internal/assignment/queue.go @@ -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) +} diff --git a/main.go b/main.go index 2d49cc0..0a22e79 100644 --- a/main.go +++ b/main.go @@ -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)