Adds the backend half of the ExpressDispatchAgent flow. Express orders can now be created batch by batch (bulk create only accumulates them, unassigned), then an operator hits one endpoint to hand the whole pending set to the agent for tenant-scoped assignment + road sequencing. The normal B2C flow is untouched. - POST /admin/expressbooking/dispatch: manual trigger. Console-auth, tenant- scoped; gathers the tenant's pending unassigned express orders (or a chosen subset) and publishes express.dispatch_requested. - internal API for the agent: GET /internal/express/riders (tenant's available riders), GET /internal/express/bookings, POST /internal/express/assign (writes the agent's decided assignments with their sequence; re-checks the already- assigned guard so the agent can't double-assign). - booking_assignment_service.go: extracted a behavior-preserving assignMilerTx core; AssignMilerToBooking is unchanged in behavior. assignExpressStops writes a batch, one FCM per rider instead of one per stop. - EXPRESS JetStream stream / express.dispatch_requested subject. - Gated behind EXPRESS_AGENT_ENABLED (default off): deploying this changes nothing until the agent is confirmed running and the flag is flipped. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
160 lines
5.8 KiB
Go
160 lines
5.8 KiB
Go
package db
|
|
|
|
import (
|
|
"errors"
|
|
|
|
"doormile/utils"
|
|
|
|
nats "github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// streamSubjects is the contract between this binary's js.Publish calls and the
|
|
// JetStream server. It lives here, next to the code that publishes, on purpose:
|
|
// the streams used to be created only by an external Python script
|
|
// (Birock/doormile-bookings/setup_jetstream.py) on another machine, and the two
|
|
// drifted.
|
|
//
|
|
// That script has since been reduced to a read-only inspector, because it
|
|
// updated streams with the subject list replaced wholesale. Re-running it would
|
|
// have stripped every subject added here — including booking.assignment_requested,
|
|
// which carries the assignment retry loop itself, so bookings would have stopped
|
|
// reaching riders entirely with nothing logged anywhere. This map is now the only
|
|
// place streams are defined; adding a js.Publish without adding its subject here
|
|
// means the event goes nowhere. Four of the eight subjects this app published were bound to no stream
|
|
// at all — booking.cancelled, booking.outcome, booking.assignment_failed, and
|
|
// chat.room.closed.<id> (CHAT declared the literal "chat.room.closed", which does
|
|
// not match a subject with an extra token). Every publish site is best-effort
|
|
// (`if Js != nil` + warn-log), so those events were failing and being dropped
|
|
// silently rather than surfacing as an error anywhere.
|
|
//
|
|
// Adding a subject here is what makes a new js.Publish actually durable. If you
|
|
// add a publish call and skip this map, the event goes nowhere.
|
|
var streamSubjects = map[string][]string{
|
|
"BOOKINGS": {
|
|
"api.v1.bookings.create",
|
|
"api.v1.bookings.update",
|
|
"api.v1.bookings.cancel",
|
|
},
|
|
"ASSIGNMENTS": {
|
|
"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",
|
|
},
|
|
"EXPRESS": {
|
|
// The console's manual "dispatch" action publishes this; the
|
|
// ExpressDispatchAgent (logistics-ai) consumes it to run tenant-scoped
|
|
// batch assignment + route sequencing over the pending orders. Missing
|
|
// subject → the dispatch is never handed off and bookings sit unassigned
|
|
// until someone assigns them by hand.
|
|
"express.dispatch_requested",
|
|
},
|
|
"STATUS": {
|
|
"booking.status.updated",
|
|
"booking.cancelled",
|
|
"booking.outcome",
|
|
},
|
|
"TRACKING": {
|
|
"miler.location.updated",
|
|
"miler.stalled",
|
|
},
|
|
"NOTIFICATIONS": {
|
|
"notification.send",
|
|
},
|
|
"CHAT": {
|
|
"chat.room.created",
|
|
"chat.room.closed",
|
|
// chat.go publishes chat.room.closed.<bookingid>; the bare subject above
|
|
// does not match it. Both are kept because the bare one may already have
|
|
// messages bound to it from the previous configuration.
|
|
"chat.room.closed.*",
|
|
},
|
|
}
|
|
|
|
// EnsureStreams declares the streams above, idempotently, at startup.
|
|
//
|
|
// It only ever *adds*: a stream that already exists keeps its storage type,
|
|
// retention, limits and every subject it already had, and is updated only when
|
|
// this map names a subject it is missing. Nothing is deleted and no existing
|
|
// subject is dropped, so running this against the streams the Python script
|
|
// created is safe and converges rather than clobbering.
|
|
//
|
|
// Failures are logged and never fatal — NATS is best-effort everywhere in this
|
|
// codebase, and the API must still serve requests when the event bus is down.
|
|
func EnsureStreams() {
|
|
if Js == nil {
|
|
utils.Warn("NATS JetStream not available, skipping stream declaration")
|
|
return
|
|
}
|
|
|
|
for name, want := range streamSubjects {
|
|
info, err := Js.StreamInfo(name)
|
|
|
|
if errors.Is(err, nats.ErrStreamNotFound) {
|
|
// New stream: file storage, so events survive a NATS restart. The
|
|
// externally-created streams are memory-backed; those are left as
|
|
// they are rather than migrated, since storage type cannot be
|
|
// changed on an existing stream.
|
|
if _, addErr := Js.AddStream(&nats.StreamConfig{
|
|
Name: name,
|
|
Subjects: want,
|
|
Storage: nats.FileStorage,
|
|
Replicas: 1,
|
|
}); addErr == nil {
|
|
utils.Info("created JetStream stream", "stream", name, "subjects", want)
|
|
continue
|
|
}
|
|
|
|
// The API runs multiple replicas, all of which call this at boot, so
|
|
// against a cold NATS they race and every loser sees a "name already
|
|
// in use" style failure. Re-reading settles which it was: if the
|
|
// stream is there now, a sibling created it and we simply carry on to
|
|
// reconcile its subjects.
|
|
info, err = Js.StreamInfo(name)
|
|
if err != nil {
|
|
utils.Error("failed to create JetStream stream", "stream", name, "error", err)
|
|
continue
|
|
}
|
|
} else if err != nil {
|
|
utils.Error("failed to inspect JetStream stream", "stream", name, "error", err)
|
|
continue
|
|
}
|
|
|
|
missing := subjectsNotIn(want, info.Config.Subjects)
|
|
if len(missing) == 0 {
|
|
continue
|
|
}
|
|
|
|
cfg := info.Config
|
|
cfg.Subjects = append(append([]string{}, info.Config.Subjects...), missing...)
|
|
if _, err := Js.UpdateStream(&cfg); err != nil {
|
|
utils.Error("failed to add subjects to JetStream stream",
|
|
"stream", name, "missing", missing, "error", err)
|
|
continue
|
|
}
|
|
utils.Info("added subjects to JetStream stream", "stream", name, "added", missing)
|
|
}
|
|
}
|
|
|
|
// subjectsNotIn returns the entries of want that have no exact match in have.
|
|
// Exact matching is deliberate: a stream carrying "chat.room.closed" does not
|
|
// carry "chat.room.closed.*", and treating one as covering the other is the
|
|
// mistake that dropped the chat events in the first place.
|
|
func subjectsNotIn(want, have []string) []string {
|
|
existing := make(map[string]struct{}, len(have))
|
|
for _, s := range have {
|
|
existing[s] = struct{}{}
|
|
}
|
|
|
|
var missing []string
|
|
for _, s := range want {
|
|
if _, ok := existing[s]; !ok {
|
|
missing = append(missing, s)
|
|
}
|
|
}
|
|
return missing
|
|
}
|