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, }) }