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>
318 lines
11 KiB
Go
318 lines
11 KiB
Go
package controllers
|
|
|
|
import (
|
|
"encoding/json"
|
|
"os"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"doormile/constants"
|
|
"doormile/db"
|
|
"doormile/models"
|
|
"doormile/utils"
|
|
|
|
"github.com/gofiber/fiber/v2"
|
|
)
|
|
|
|
// expressAgentEnabled gates the whole express-batch handoff. Default OFF, so
|
|
// deploying this code changes nothing: a bulk create keeps assigning each
|
|
// booking inline exactly as before, and no batch event is emitted. Flip
|
|
// EXPRESS_AGENT_ENABLED=true ONLY once the ExpressDispatchAgent is confirmed
|
|
// running and consuming express.batch_created — otherwise batches would suppress
|
|
// inline assignment with nothing to pick them up, and bulk bookings would sit
|
|
// unassigned. Read at request time so it can be toggled without a redeploy.
|
|
func expressAgentEnabled() bool {
|
|
return strings.EqualFold(os.Getenv("EXPRESS_AGENT_ENABLED"), "true")
|
|
}
|
|
|
|
// expressDispatchSubject is the JetStream subject the ExpressDispatchAgent binds
|
|
// to. It must be present on the EXPRESS stream (db/streams.go) before anything
|
|
// publishes here — JetStream silently drops messages on uncovered subjects.
|
|
const expressDispatchSubject = "express.dispatch_requested"
|
|
|
|
// publishExpressBatch emits one express.dispatch_requested event per tenant when
|
|
// an operator dispatches their pending orders. Best-effort like every publish in
|
|
// this codebase: a NATS outage degrades to hand-assignment from the console, it
|
|
// never fails the request.
|
|
func publishExpressBatch(byTenant map[int][]int) {
|
|
if db.Js == nil {
|
|
if len(byTenant) > 0 {
|
|
utils.Warn("Express: JetStream unavailable, dispatch not handed to agent",
|
|
"tenants", len(byTenant))
|
|
}
|
|
return
|
|
}
|
|
for tenantID, bookingIDs := range byTenant {
|
|
if len(bookingIDs) == 0 {
|
|
continue
|
|
}
|
|
payload := map[string]interface{}{
|
|
"tenantid": tenantID,
|
|
"booking_ids": bookingIDs,
|
|
"created_at": time.Now().UnixMilli(),
|
|
}
|
|
data, err := json.Marshal(payload)
|
|
if err != nil {
|
|
utils.Warn("Express: failed to marshal dispatch event", "tenantid", tenantID, "error", err)
|
|
continue
|
|
}
|
|
if _, err := db.Js.Publish(expressDispatchSubject, data); err != nil {
|
|
utils.Warn("Express: failed to publish dispatch event", "tenantid", tenantID, "error", err)
|
|
continue
|
|
}
|
|
utils.Info("Express: published "+expressDispatchSubject,
|
|
"tenantid", tenantID, "bookings", len(bookingIDs))
|
|
}
|
|
}
|
|
|
|
// This file is the internal API surface the ExpressDispatchAgent (logistics-ai)
|
|
// uses to run the express-batch flow: read the tenant's available riders, read
|
|
// the batch's bookings, and write back the assignments it decided after calling
|
|
// the Route Optimization API. All three sit under /internal (InternalKeyAuth),
|
|
// never exposed to a console or app token.
|
|
//
|
|
// The division of labour: the agent decides *who* and *what order* (rider pool +
|
|
// routes.workolik). Go stays the single writer of assignment state — the agent
|
|
// never writes the DB directly, it posts its decision here and this reuses the
|
|
// same transactional assignment path the consoles use.
|
|
|
|
// DispatchExpressBatch is the manual trigger the console operator hits once a
|
|
// batch of express orders has piled up. Orders are created batch by batch (bulk
|
|
// create only accumulates them, unassigned); this hands the whole pending set
|
|
// for the tenant to the ExpressDispatchAgent in one go — the "take over from
|
|
// here" button.
|
|
//
|
|
// Console auth, tenant-scoped: a client login dispatches only its own pending
|
|
// orders; Doormile staff pass ?tenantid= to dispatch for one client. An optional
|
|
// body {"booking_ids":[...]} dispatches a chosen subset instead of everything
|
|
// pending — always still pinned to the resolved tenant so no cross-tenant id can
|
|
// be smuggled in.
|
|
func DispatchExpressBatch(c *fiber.Ctx) error {
|
|
tenantID, allowed := effectiveTenantID(c)
|
|
if !allowed {
|
|
return utils.Forbidden(c, "you can only dispatch your own tenant")
|
|
}
|
|
if tenantID == 0 {
|
|
return utils.BadRequest(c, "tenantid is required (staff must pass ?tenantid=)")
|
|
}
|
|
|
|
var body struct {
|
|
BookingIDs []int `json:"booking_ids"`
|
|
}
|
|
_ = c.BodyParser(&body) // body is optional
|
|
|
|
// Only ever the tenant's own, unassigned, still-pending express orders.
|
|
q := db.DB.Model(&models.PickupBooking{}).
|
|
Where("tenantid = ? AND bookingsource = ? AND assignedmileruserid IS NULL AND status = ?",
|
|
tenantID, constants.BookingSourceExpress, constants.BookingPendingPickup)
|
|
if len(body.BookingIDs) > 0 {
|
|
q = q.Where("bookingid IN ?", body.BookingIDs)
|
|
}
|
|
|
|
var bookingIDs []int
|
|
if err := q.Pluck("bookingid", &bookingIDs).Error; err != nil {
|
|
return utils.Internal(c, "failed to gather pending bookings")
|
|
}
|
|
if len(bookingIDs) == 0 {
|
|
return utils.OK(c, fiber.Map{
|
|
"queued": 0,
|
|
"booking_ids": []int{},
|
|
"message": "no pending express bookings to dispatch",
|
|
})
|
|
}
|
|
|
|
publishExpressBatch(map[int][]int{tenantID: bookingIDs})
|
|
|
|
resp := fiber.Map{
|
|
"queued": len(bookingIDs),
|
|
"booking_ids": bookingIDs,
|
|
}
|
|
if !expressAgentEnabled() {
|
|
// The batch was published, but with the flag off the bulk path may have
|
|
// already assigned these inline and the agent may not be consuming — make
|
|
// that visible rather than implying work was dispatched.
|
|
resp["warning"] = "EXPRESS_AGENT_ENABLED is off; the agent may not be consuming this event"
|
|
}
|
|
return utils.OK(c, resp)
|
|
}
|
|
|
|
// GetExpressRiders returns a tenant's riders that are free to take work, with the
|
|
// location and hub the agent needs to distribute stops. Optional ?city= narrows
|
|
// to one operating zone (applocationid) so a Coimbatore batch is never handed to
|
|
// a Nagercoil rider.
|
|
func GetExpressRiders(c *fiber.Ctx) error {
|
|
tenantID := c.QueryInt("tenantid", 0)
|
|
if tenantID == 0 {
|
|
return utils.BadRequest(c, "tenantid is required")
|
|
}
|
|
|
|
type riderRow struct {
|
|
Userid int `json:"miler_user_id"`
|
|
Displayname string `json:"displayname"`
|
|
Phone string `json:"phone"`
|
|
Currentlatitude float64 `json:"latitude"`
|
|
Currentlongitude float64 `json:"longitude"`
|
|
Hubid *int `json:"hubid"`
|
|
Applocationid int `json:"applocationid"`
|
|
Availabilitystatus string `json:"availabilitystatus"`
|
|
Devicetoken string `json:"has_device_token"`
|
|
}
|
|
|
|
q := db.DB.Table("milerprofiles AS mp").
|
|
Select(`mp.userid, mp.displayname, mp.phone, mp.currentlatitude,
|
|
mp.currentlongitude, mp.hubid, mp.applocationid,
|
|
mp.availabilitystatus, mp.device_token`).
|
|
Joins("JOIN appusers AS u ON u.userid = mp.userid").
|
|
Where("u.tenantid = ? AND u.roleid = ? AND mp.availabilitystatus = ?",
|
|
tenantID, 5, constants.MilerAvailable)
|
|
|
|
if city := c.QueryInt("city", 0); city != 0 {
|
|
q = q.Where("mp.applocationid = ?", city)
|
|
}
|
|
|
|
var rows []riderRow
|
|
if err := q.Scan(&rows).Error; err != nil {
|
|
return utils.Internal(c, "failed to load riders")
|
|
}
|
|
|
|
// Expose only whether a device token exists, never the token itself.
|
|
out := make([]fiber.Map, 0, len(rows))
|
|
for _, r := range rows {
|
|
out = append(out, fiber.Map{
|
|
"miler_user_id": r.Userid,
|
|
"displayname": r.Displayname,
|
|
"phone": r.Phone,
|
|
"latitude": r.Currentlatitude,
|
|
"longitude": r.Currentlongitude,
|
|
"hubid": r.Hubid,
|
|
"applocationid": r.Applocationid,
|
|
"availabilitystatus": r.Availabilitystatus,
|
|
"has_device_token": r.Devicetoken != "",
|
|
})
|
|
}
|
|
return c.JSON(fiber.Map{"success": true, "riders": out, "total": len(out)})
|
|
}
|
|
|
|
// GetExpressBookings returns the coordinates and kitchen for a set of booking
|
|
// ids — everything the agent needs to feed the optimizer, and nothing it does
|
|
// not. ?ids=1,2,3.
|
|
func GetExpressBookings(c *fiber.Ctx) error {
|
|
idsParam := c.Query("ids")
|
|
if idsParam == "" {
|
|
return utils.BadRequest(c, "ids is required, e.g. ?ids=1,2,3")
|
|
}
|
|
|
|
ids := make([]int, 0)
|
|
for _, part := range strings.Split(idsParam, ",") {
|
|
part = strings.TrimSpace(part)
|
|
if part == "" {
|
|
continue
|
|
}
|
|
n, err := strconv.Atoi(part)
|
|
if err != nil {
|
|
return utils.BadRequest(c, "ids must be a comma-separated list of integers")
|
|
}
|
|
ids = append(ids, n)
|
|
}
|
|
if len(ids) == 0 {
|
|
return utils.BadRequest(c, "ids is required")
|
|
}
|
|
|
|
type bookingRow struct {
|
|
Bookingid int `json:"booking_id"`
|
|
Bookingno string `json:"booking_no"`
|
|
Tenantid *int `json:"tenantid"`
|
|
Tenantlocationid *int `json:"tenantlocationid"`
|
|
Pickuplatitude float64 `json:"pickuplatitude"`
|
|
Pickuplongitude float64 `json:"pickuplongitude"`
|
|
Deliverylatitude float64 `json:"deliverylatitude"`
|
|
Deliverylongitude float64 `json:"deliverylongitude"`
|
|
Pickuppincode string `json:"pickuppincode"`
|
|
Deliverypincode string `json:"deliverypincode"`
|
|
Status string `json:"status"`
|
|
Assignedmileruserid *int `json:"assignedmileruserid"`
|
|
}
|
|
|
|
var rows []bookingRow
|
|
if err := db.DB.Model(&models.PickupBooking{}).
|
|
Select(`bookingid, bookingno, tenantid, tenantlocationid,
|
|
pickuplatitude, pickuplongitude, deliverylatitude, deliverylongitude,
|
|
pickuppincode, deliverypincode, status, assignedmileruserid`).
|
|
Where("bookingid IN ?", ids).
|
|
Scan(&rows).Error; err != nil {
|
|
return utils.Internal(c, "failed to load bookings")
|
|
}
|
|
|
|
return c.JSON(fiber.Map{"success": true, "bookings": rows, "total": len(rows)})
|
|
}
|
|
|
|
// AssignExpressBatch writes the agent's decided assignments. Body:
|
|
//
|
|
// { "assignments": [ {booking_id, miler_user_id, step, previouskms,
|
|
// cumulativekms, etaminutes, cumulativeeta}, ... ] }
|
|
//
|
|
// Each row is assigned in its own transaction with its sequence already set, and
|
|
// each miler is notified once for the whole batch. Per-row results mirror the
|
|
// bulk-create shape so a single bad booking id never fails the batch.
|
|
func AssignExpressBatch(c *fiber.Ctx) error {
|
|
var req struct {
|
|
Assignments []ExpressStop `json:"assignments"`
|
|
}
|
|
if err := c.BodyParser(&req); err != nil {
|
|
return utils.BadRequest(c, "invalid request body")
|
|
}
|
|
if len(req.Assignments) == 0 {
|
|
return utils.BadRequest(c, "assignments is required and must not be empty")
|
|
}
|
|
|
|
// Guard against writing to a booking that is not actually the tenant's or is
|
|
// already assigned — the agent is trusted but the writeback must still be the
|
|
// place assignment invariants are enforced, not the agent.
|
|
bookingIDs := make([]int, 0, len(req.Assignments))
|
|
for _, a := range req.Assignments {
|
|
bookingIDs = append(bookingIDs, a.BookingID)
|
|
}
|
|
assignable := map[int]bool{}
|
|
var existing []struct {
|
|
Bookingid int
|
|
Assignedmileruserid *int
|
|
Status string
|
|
}
|
|
db.DB.Model(&models.PickupBooking{}).
|
|
Select("bookingid, assignedmileruserid, status").
|
|
Where("bookingid IN ?", bookingIDs).
|
|
Scan(&existing)
|
|
for _, b := range existing {
|
|
assignable[b.Bookingid] = b.Assignedmileruserid == nil &&
|
|
b.Status != constants.BookingCancelled
|
|
}
|
|
|
|
toAssign := make([]ExpressStop, 0, len(req.Assignments))
|
|
results := make([]ExpressAssignResult, 0, len(req.Assignments))
|
|
for _, a := range req.Assignments {
|
|
if !assignable[a.BookingID] {
|
|
results = append(results, ExpressAssignResult{
|
|
BookingID: a.BookingID, MilerUserID: a.MilerUserID,
|
|
Success: false, Error: "booking already assigned or not eligible"})
|
|
continue
|
|
}
|
|
toAssign = append(toAssign, a)
|
|
}
|
|
|
|
results = append(results, assignExpressStops(toAssign)...)
|
|
|
|
assigned := 0
|
|
for _, r := range results {
|
|
if r.Success {
|
|
assigned++
|
|
}
|
|
}
|
|
return c.JSON(fiber.Map{
|
|
"success": true,
|
|
"assigned": assigned,
|
|
"total": len(req.Assignments),
|
|
"results": results,
|
|
})
|
|
}
|