Compare commits
9 Commits
90fa4fbb74
...
cfdc99ec9a
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cfdc99ec9a | ||
|
|
ba800aee66 | ||
|
|
cb2660a3da | ||
|
|
0288fb7af8 | ||
|
|
2cbc9e5b13 | ||
|
|
2158031191 | ||
|
|
41f0013751 | ||
|
|
5458d080c6 | ||
|
|
6738d37d7b |
@@ -22,6 +22,11 @@ type Config struct {
|
||||
NatsPassword string
|
||||
AILayerBaseURL string // AI decision-engine service base URL (e.g. http://rider-api:8082)
|
||||
|
||||
// RouteOptimizerURL is the Route Optimization API that orders a rider's
|
||||
// stops (Valhalla-backed road sequencing). Empty disables sequencing: stops
|
||||
// stay unordered rather than assignment failing.
|
||||
RouteOptimizerURL string
|
||||
|
||||
// TrustedProxies is a comma-separated list of reverse-proxy IPs/CIDRs that
|
||||
// are allowed to set X-Forwarded-For. Rate limiting keys on the client IP,
|
||||
// so behind a proxy this MUST be set — otherwise every request appears to
|
||||
@@ -53,12 +58,14 @@ func Load() *Config {
|
||||
NatsUser: getEnv("NATS_USER", "doormile"),
|
||||
NatsPassword: getEnv("NATS_PASSWORD", "Package@321#"),
|
||||
AILayerBaseURL: getEnv("AI_LAYER_BASE_URL", "https://routemate.workolik.com"),
|
||||
TrustedProxies: getEnv("TRUSTED_PROXIES", ""),
|
||||
SMTPHost: getEnv("SMTP_HOST", ""),
|
||||
SMTPPort: getEnv("SMTP_PORT", "465"),
|
||||
SMTPUser: getEnv("SMTP_USER", ""),
|
||||
SMTPPassword: getEnv("SMTP_PASSWORD", ""),
|
||||
SMTPFrom: getEnv("SMTP_FROM", ""),
|
||||
|
||||
RouteOptimizerURL: getEnv("ROUTE_OPTIMIZER_URL", "https://routes.workolik.com"),
|
||||
TrustedProxies: getEnv("TRUSTED_PROXIES", ""),
|
||||
SMTPHost: getEnv("SMTP_HOST", ""),
|
||||
SMTPPort: getEnv("SMTP_PORT", "465"),
|
||||
SMTPUser: getEnv("SMTP_USER", ""),
|
||||
SMTPPassword: getEnv("SMTP_PASSWORD", ""),
|
||||
SMTPFrom: getEnv("SMTP_FROM", ""),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -2,20 +2,27 @@ package constants
|
||||
|
||||
// Miler Availability Statuses
|
||||
const (
|
||||
MilerOffline = "Offline"
|
||||
MilerAvailable = "Available"
|
||||
MilerAssigned = "Assigned"
|
||||
MilerOnPickup = "On_Pickup"
|
||||
MilerAtCustomer = "At_Customer"
|
||||
MilerPickedUp = "Picked_Up"
|
||||
MilerOnDelivery = "On_Delivery"
|
||||
MilerBreak = "Break"
|
||||
MilerBlocked = "Blocked"
|
||||
MilerOffline = "Offline"
|
||||
MilerAvailable = "Available"
|
||||
MilerAssigned = "Assigned"
|
||||
MilerOnPickup = "On_Pickup"
|
||||
MilerAtCustomer = "At_Customer"
|
||||
MilerPickedUp = "Picked_Up"
|
||||
MilerOnDelivery = "On_Delivery"
|
||||
MilerBreak = "Break"
|
||||
MilerBlocked = "Blocked"
|
||||
)
|
||||
|
||||
// App (B2C) Customer Statuses
|
||||
const (
|
||||
CustomerStatusActive = "Active"
|
||||
CustomerStatusBlocked = "Blocked"
|
||||
CustomerStatusDeleted = "Deleted"
|
||||
)
|
||||
|
||||
// Booking Statuses
|
||||
const (
|
||||
BookingPendingPickup = "Pending_Pickup" // customer requested pickup, delivery details not yet known
|
||||
BookingPendingPickup = "Pending_Pickup" // customer requested pickup, delivery details not yet known
|
||||
BookingCreated = "Created"
|
||||
BookingMilerAssigned = "Miler_Assigned"
|
||||
BookingPickupScheduled = "Pickup_Scheduled"
|
||||
@@ -26,16 +33,16 @@ const (
|
||||
|
||||
// Consignment Statuses
|
||||
const (
|
||||
ConsignmentCreated = "Created"
|
||||
ConsignmentInwardedAtHub = "Inwarded_at_Hub"
|
||||
ConsignmentTripsheetLoaded = "Tripsheet_Loaded"
|
||||
ConsignmentInTransit = "In_Transit"
|
||||
ConsignmentOutForDelivery = "Out_for_Delivery"
|
||||
ConsignmentDelivered = "Delivered"
|
||||
ConsignmentRTOInitiated = "RTO_Initiated"
|
||||
ConsignmentReturnedToSender = "Returned_to_Sender"
|
||||
ConsignmentMissing = "Missing"
|
||||
ConsignmentDamaged = "Damaged"
|
||||
ConsignmentCreated = "Created"
|
||||
ConsignmentInwardedAtHub = "Inwarded_at_Hub"
|
||||
ConsignmentTripsheetLoaded = "Tripsheet_Loaded"
|
||||
ConsignmentInTransit = "In_Transit"
|
||||
ConsignmentOutForDelivery = "Out_for_Delivery"
|
||||
ConsignmentDelivered = "Delivered"
|
||||
ConsignmentRTOInitiated = "RTO_Initiated"
|
||||
ConsignmentReturnedToSender = "Returned_to_Sender"
|
||||
ConsignmentMissing = "Missing"
|
||||
ConsignmentDamaged = "Damaged"
|
||||
)
|
||||
|
||||
// Payment Modes
|
||||
@@ -93,8 +100,8 @@ const (
|
||||
|
||||
// Exception Statuses
|
||||
const (
|
||||
ExceptionOpen = "Open"
|
||||
ExceptionOpen = "Open"
|
||||
ExceptionUnderInvestigation = "Under_Investigation"
|
||||
ExceptionResolved = "Resolved"
|
||||
ExceptionClosed = "Closed"
|
||||
ExceptionResolved = "Resolved"
|
||||
ExceptionClosed = "Closed"
|
||||
)
|
||||
|
||||
@@ -1135,6 +1135,10 @@ func GetAdminCustomers(c *fiber.Ctx) error {
|
||||
like := "%" + keyword + "%"
|
||||
query = query.Where("firstname ILIKE ? OR lastname ILIKE ? OR phone ILIKE ?", like, like, like)
|
||||
}
|
||||
// Zone/city scoping. 0 means "all zones".
|
||||
if applocationid := c.QueryInt("applocationid", 0); applocationid != 0 {
|
||||
query = query.Where("applocationid = ?", applocationid)
|
||||
}
|
||||
// A client sees only the customers they have actually delivered to, not
|
||||
// Doormile's whole B2C address book. Doormile staff can ask for one
|
||||
// client's customers with ?tenantid=.
|
||||
@@ -1186,8 +1190,21 @@ func GetAdminCustomers(c *fiber.Ctx) error {
|
||||
data = append(data, fiber.Map{
|
||||
"appcustomerid": cust.Appcustomerid,
|
||||
"name": strings.TrimSpace(cust.Firstname + " " + cust.Lastname),
|
||||
"firstname": cust.Firstname,
|
||||
"lastname": cust.Lastname,
|
||||
"phone": cust.Phone,
|
||||
"email": cust.Email,
|
||||
"doorno": cust.Doorno,
|
||||
"address": cust.Address,
|
||||
"suburb": cust.Suburb,
|
||||
"city": cust.City,
|
||||
"state": cust.State,
|
||||
"postcode": cust.Postcode,
|
||||
"landmark": cust.Landmark,
|
||||
"latitude": cust.Latitude,
|
||||
"longitude": cust.Longitude,
|
||||
"applocationid": cust.Applocationid,
|
||||
"status": cust.Status,
|
||||
"createdat": cust.Createdat,
|
||||
"totalbookings": bookingCounts[cust.Appcustomerid],
|
||||
})
|
||||
@@ -1216,57 +1233,169 @@ func UpdateAdminCustomer(c *fiber.Ctx) error {
|
||||
return utils.NotFound(c, "customer not found")
|
||||
}
|
||||
|
||||
// Pointers so "field omitted" is distinguishable from "field cleared to
|
||||
// empty" — the console sends only what it changed, and a missing address
|
||||
// line must not blank a stored one.
|
||||
req := new(struct {
|
||||
Name string `json:"name"`
|
||||
Phone string `json:"phone"`
|
||||
Email string `json:"email"`
|
||||
Name *string `json:"name"`
|
||||
Firstname *string `json:"firstname"`
|
||||
Lastname *string `json:"lastname"`
|
||||
Phone *string `json:"phone"`
|
||||
Email *string `json:"email"`
|
||||
Doorno *string `json:"doorno"`
|
||||
Address *string `json:"address"`
|
||||
Suburb *string `json:"suburb"`
|
||||
City *string `json:"city"`
|
||||
State *string `json:"state"`
|
||||
Postcode *string `json:"postcode"`
|
||||
Landmark *string `json:"landmark"`
|
||||
Latitude *float64 `json:"latitude"`
|
||||
Longitude *float64 `json:"longitude"`
|
||||
Applocationid *int `json:"applocationid"`
|
||||
})
|
||||
if err := c.BodyParser(req); err != nil {
|
||||
return utils.BadRequest(c, "invalid request body")
|
||||
}
|
||||
|
||||
if req.Phone != "" {
|
||||
if len(req.Phone) != 10 || strings.IndexFunc(req.Phone, func(r rune) bool { return r < '0' || r > '9' }) != -1 {
|
||||
updates := map[string]interface{}{}
|
||||
|
||||
if req.Phone != nil {
|
||||
phone := *req.Phone
|
||||
if len(phone) != 10 || strings.IndexFunc(phone, func(r rune) bool { return r < '0' || r > '9' }) != -1 {
|
||||
return utils.BadRequest(c, "phone must be 10 digits")
|
||||
}
|
||||
customer.Phone = req.Phone
|
||||
updates["phone"] = phone
|
||||
}
|
||||
|
||||
if req.Name != "" {
|
||||
name := strings.TrimSpace(req.Name)
|
||||
// firstname/lastname sent explicitly win; a single name field is split as a
|
||||
// fallback for callers that still send it.
|
||||
if req.Firstname != nil {
|
||||
updates["firstname"] = strings.TrimSpace(*req.Firstname)
|
||||
}
|
||||
if req.Lastname != nil {
|
||||
updates["lastname"] = strings.TrimSpace(*req.Lastname)
|
||||
}
|
||||
if req.Firstname == nil && req.Lastname == nil && req.Name != nil {
|
||||
name := strings.TrimSpace(*req.Name)
|
||||
parts := strings.SplitN(name, " ", 2)
|
||||
customer.Firstname = parts[0]
|
||||
updates["firstname"] = parts[0]
|
||||
if len(parts) > 1 {
|
||||
customer.Lastname = parts[1]
|
||||
updates["lastname"] = parts[1]
|
||||
} else {
|
||||
customer.Lastname = ""
|
||||
updates["lastname"] = ""
|
||||
}
|
||||
}
|
||||
|
||||
if req.Email != "" {
|
||||
customer.Email = req.Email
|
||||
if req.Email != nil {
|
||||
updates["email"] = *req.Email
|
||||
}
|
||||
if req.Doorno != nil {
|
||||
updates["doorno"] = *req.Doorno
|
||||
}
|
||||
if req.Address != nil {
|
||||
updates["address"] = *req.Address
|
||||
}
|
||||
if req.Suburb != nil {
|
||||
updates["suburb"] = *req.Suburb
|
||||
}
|
||||
if req.City != nil {
|
||||
updates["city"] = *req.City
|
||||
}
|
||||
if req.State != nil {
|
||||
updates["state"] = *req.State
|
||||
}
|
||||
if req.Postcode != nil {
|
||||
updates["postcode"] = *req.Postcode
|
||||
}
|
||||
if req.Landmark != nil {
|
||||
updates["landmark"] = *req.Landmark
|
||||
}
|
||||
if req.Latitude != nil {
|
||||
updates["latitude"] = *req.Latitude
|
||||
}
|
||||
if req.Longitude != nil {
|
||||
updates["longitude"] = *req.Longitude
|
||||
}
|
||||
if req.Applocationid != nil {
|
||||
updates["applocationid"] = *req.Applocationid
|
||||
}
|
||||
|
||||
customer.Updatedat = time.Now()
|
||||
if len(updates) == 0 {
|
||||
return utils.BadRequest(c, "no fields to update")
|
||||
}
|
||||
updates["updatedat"] = time.Now()
|
||||
|
||||
if err := db.DB.Model(&models.AppCustomer{}).Where("appcustomerid = ?", customer.Appcustomerid).
|
||||
Updates(map[string]interface{}{
|
||||
"firstname": customer.Firstname,
|
||||
"lastname": customer.Lastname,
|
||||
"phone": customer.Phone,
|
||||
"email": customer.Email,
|
||||
"updatedat": customer.Updatedat,
|
||||
}).Error; err != nil {
|
||||
if err := db.DB.Model(&models.AppCustomer{}).
|
||||
Where("appcustomerid = ?", customer.Appcustomerid).
|
||||
Updates(updates).Error; err != nil {
|
||||
return utils.Internal(c, "failed to update customer")
|
||||
}
|
||||
|
||||
// Return the persisted row so the console reflects exactly what was stored.
|
||||
if err := db.DB.First(&customer, customer.Appcustomerid).Error; err != nil {
|
||||
return utils.Internal(c, "failed to reload customer")
|
||||
}
|
||||
|
||||
return c.JSON(fiber.Map{
|
||||
"success": true,
|
||||
"data": fiber.Map{
|
||||
"appcustomerid": customer.Appcustomerid,
|
||||
"name": strings.TrimSpace(customer.Firstname + " " + customer.Lastname),
|
||||
"firstname": customer.Firstname,
|
||||
"lastname": customer.Lastname,
|
||||
"phone": customer.Phone,
|
||||
"email": customer.Email,
|
||||
"doorno": customer.Doorno,
|
||||
"address": customer.Address,
|
||||
"suburb": customer.Suburb,
|
||||
"city": customer.City,
|
||||
"state": customer.State,
|
||||
"postcode": customer.Postcode,
|
||||
"landmark": customer.Landmark,
|
||||
"latitude": customer.Latitude,
|
||||
"longitude": customer.Longitude,
|
||||
"applocationid": customer.Applocationid,
|
||||
"status": customer.Status,
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
// GetAdminCustomersSummary returns the stat-tile counts the console header
|
||||
// shows, so the client does not have to page through the whole list to derive
|
||||
// them. Scoped the same way as the list: a client sees only their own
|
||||
// customers, Doormile staff can narrow with ?tenantid= / ?applocationid=.
|
||||
func GetAdminCustomersSummary(c *fiber.Ctx) error {
|
||||
base := db.DB.Model(&models.AppCustomer{})
|
||||
|
||||
tenantID, allowed := effectiveTenantID(c)
|
||||
if !allowed {
|
||||
return utils.Forbidden(c, "you can only view your own tenant")
|
||||
}
|
||||
if tenantID != 0 {
|
||||
base = base.Where("appcustomerid IN (?)",
|
||||
db.DB.Model(&models.PickupBooking{}).Select("appcustomerid").
|
||||
Where("tenantid = ?", tenantID))
|
||||
}
|
||||
if applocationid := c.QueryInt("applocationid", 0); applocationid != 0 {
|
||||
base = base.Where("applocationid = ?", applocationid)
|
||||
}
|
||||
|
||||
count := func(status string) int64 {
|
||||
var n int64
|
||||
q := base.Session(&gorm.Session{})
|
||||
if status != "" {
|
||||
q = q.Where("status = ?", status)
|
||||
}
|
||||
q.Count(&n)
|
||||
return n
|
||||
}
|
||||
|
||||
return c.JSON(fiber.Map{
|
||||
"success": true,
|
||||
"data": fiber.Map{
|
||||
"total": count(""),
|
||||
"active": count(constants.CustomerStatusActive),
|
||||
"blocked": count(constants.CustomerStatusBlocked),
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"doormile/db"
|
||||
"doormile/dto"
|
||||
"doormile/internal/assignment"
|
||||
"doormile/internal/routing"
|
||||
"doormile/models"
|
||||
"doormile/utils"
|
||||
|
||||
@@ -2018,13 +2019,54 @@ func HubBatchAssign(c *fiber.Ctx) error {
|
||||
assignedCount++
|
||||
}
|
||||
|
||||
// Batch assignment is exactly the case stop-ordering exists for: a rider
|
||||
// walks out of here with several bookings and, until now, no indication of
|
||||
// what order to run them in. Sequence each rider that actually got work.
|
||||
//
|
||||
// Best-effort and deliberately after the assignments are committed: the
|
||||
// optimizer is a separate service over the network, and it failing must
|
||||
// leave the bookings assigned rather than undoing the batch.
|
||||
sequenced := 0
|
||||
for _, r := range ridersAssigned(results) {
|
||||
if _, err := routing.SequenceMilerStops(r); err != nil {
|
||||
utils.Warn("HubBatchAssign: stop sequencing failed",
|
||||
"miler_userid", r, "error", err)
|
||||
continue
|
||||
}
|
||||
sequenced++
|
||||
}
|
||||
|
||||
return utils.OK(c, fiber.Map{
|
||||
"assigned": assignedCount,
|
||||
"skipped": skippedCount,
|
||||
"results": results,
|
||||
"assigned": assignedCount,
|
||||
"skipped": skippedCount,
|
||||
"riderssequenced": sequenced,
|
||||
"results": results,
|
||||
})
|
||||
}
|
||||
|
||||
// ridersAssigned pulls the distinct miler user ids out of a batch result set,
|
||||
// so each rider is sequenced once rather than once per booking they received.
|
||||
func ridersAssigned(results []fiber.Map) []int {
|
||||
seen := make(map[int]struct{})
|
||||
var ids []int
|
||||
for _, r := range results {
|
||||
ok, _ := r["assigned"].(bool)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
id, isInt := r["mileruserid"].(int)
|
||||
if !isInt {
|
||||
continue
|
||||
}
|
||||
if _, dup := seen[id]; dup {
|
||||
continue
|
||||
}
|
||||
seen[id] = struct{}{}
|
||||
ids = append(ids, id)
|
||||
}
|
||||
return ids
|
||||
}
|
||||
|
||||
// --------------------
|
||||
// HUB REPORT EXPORT
|
||||
// --------------------
|
||||
|
||||
@@ -326,9 +326,14 @@ func GetMilerAssignments(c *fiber.Ctx) error {
|
||||
milerUserID := c.Locals("userid").(int)
|
||||
|
||||
var assignments []models.BookingAssignment
|
||||
// Sequenced stops come first, in the road order internal/routing worked out,
|
||||
// because that is the order the rider should actually ride them. Anything
|
||||
// not yet sequenced (step 0 — a single stop, or the optimizer being
|
||||
// unreachable) falls back to newest-first, which is the old behaviour.
|
||||
if err := db.DB.Where("mileruserid = ? AND assignmentstatus IN ?", milerUserID,
|
||||
[]string{constants.AssignmentAssigned, constants.AssignmentAccepted}).
|
||||
Order("assignedat DESC").Find(&assignments).Error; err != nil {
|
||||
Order("CASE WHEN step > 0 THEN 0 ELSE 1 END ASC, step ASC, assignedat DESC").
|
||||
Find(&assignments).Error; err != nil {
|
||||
return utils.Internal(c, "failed to fetch assignments")
|
||||
}
|
||||
|
||||
|
||||
@@ -163,6 +163,10 @@ func InitNATS(cfg *config.Config) {
|
||||
}
|
||||
|
||||
utils.Info("✅ NATS JetStream connected successfully", "url", cfg.NatsURL)
|
||||
|
||||
// Declare the streams this binary publishes to, so the subject contract is
|
||||
// owned by the code that uses it rather than by an external setup script.
|
||||
EnsureStreams()
|
||||
}
|
||||
|
||||
func getEnv(key, fallback string) string {
|
||||
|
||||
151
db/streams.go
Normal file
151
db/streams.go
Normal file
@@ -0,0 +1,151 @@
|
||||
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",
|
||||
},
|
||||
"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
|
||||
}
|
||||
240
docs/doormile-flow.md
Normal file
240
docs/doormile-flow.md
Normal file
@@ -0,0 +1,240 @@
|
||||
# Doormile — order to delivery
|
||||
|
||||
Every endpoint in the DailyGrubs path, the payload that goes in, and what
|
||||
actually changes when it does. Base URL `https://api.doormile.com/api/v1`.
|
||||
|
||||
Payload shapes are taken from the request structs in this repo, not from
|
||||
documentation — where the two disagree, the code is right. Where a field is
|
||||
optional it says so.
|
||||
|
||||
---
|
||||
|
||||
## Auth
|
||||
|
||||
Two separate logins.
|
||||
|
||||
**Console.** The token carries `tenantid`, which is what scopes a client to
|
||||
their own data. A client passing another tenant's `?tenantid=` gets 403; reading
|
||||
another tenant's resource by id gets 404, so ids aren't probeable.
|
||||
|
||||
```
|
||||
POST /admin/login
|
||||
{ "email": "info@dailygrubs.com", "password": "admin" }
|
||||
|
||||
→ { "token": "eyJ…", "user": { "id": 43, "tenantid": 13 } }
|
||||
```
|
||||
|
||||
**Rider.** Two calls, and `configid` must be **1001** on both — a Doormile login
|
||||
partition with no jupiter equivalent. Omitting it is the most common reason a
|
||||
rider looks like they don't exist.
|
||||
|
||||
```
|
||||
POST /miler/login
|
||||
{ "phone": "9787698259", "configid": 1001 }
|
||||
|
||||
POST /miler/verify-pin
|
||||
{ "phone": "9787698259", "pin": "1234", "configid": 1001, "device_token": "fcm-…" }
|
||||
|
||||
→ { "token": "eyJ…" }
|
||||
```
|
||||
|
||||
`POST /miler/reset-pin` is **admin-only** despite sitting under `/miler`. A phone
|
||||
number is the login identifier, not a secret — left open, reset-then-verify takes
|
||||
over any rider account in two calls. The rider app must not call it.
|
||||
|
||||
---
|
||||
|
||||
## 1. Create a booking
|
||||
|
||||
```
|
||||
POST /admin/expressbooking
|
||||
{
|
||||
"tenantid": 13,
|
||||
"tenantlocationid": 20, // the kitchen. optional — inferred if omitted
|
||||
"customer_phone": "9876500011",
|
||||
"customer_name": "Priya R",
|
||||
"pickupaddress": "DailyGrubs RS Puram Kitchen",
|
||||
"pickuppincode": "641002",
|
||||
"pickuplatitude": 11.004500, "pickuplongitude": 76.961200,
|
||||
"deliveryaddress": "12 Bharathi Rd, Peelamedu",
|
||||
"deliverypincode": "641004",
|
||||
"deliverylatitude": 11.051000, "deliverylongitude": 76.930000,
|
||||
"service_option": "Fast",
|
||||
"finalprice": 65,
|
||||
"parcels": [ { "itemcategory": "Food", "weight": 1.2 } ]
|
||||
}
|
||||
|
||||
→ { "bookingid": 118, "bookingno": "DM-BK-…", "status": "Created" }
|
||||
```
|
||||
|
||||
Writes a `pickupbookings` row plus its parcels and price.
|
||||
|
||||
**On `tenantlocationid`.** This is the client's own site — the kitchen. Omit it
|
||||
and it's resolved from the pickup coordinates: nearest stored site within 150m,
|
||||
falling back to an address match, nil when unsure. It is what per-kitchen
|
||||
reporting groups by.
|
||||
|
||||
`pickuplocationid` is accepted only as a legacy alias and is never stored as
|
||||
given — the column of that name foreign-keys to `appcustomerlocations` (a B2C
|
||||
customer's saved address), so writing a client site id into it fails the insert.
|
||||
|
||||
**CityGate.** The pickup pincode prefix must be an open city: `641` Coimbatore,
|
||||
`600` Chennai, `560` Bengaluru, `500` Hyderabad, `629` Nagercoil. Anything else
|
||||
is refused at creation. jupiter had no such gate.
|
||||
|
||||
## 2. Or create them in bulk
|
||||
|
||||
```
|
||||
POST /admin/expressbooking/bulk
|
||||
{ "bookings": [ { …same shape… }, { … } ] } // max 200
|
||||
|
||||
→ { "results": [
|
||||
{ "index": 0, "success": true, "bookingid": 119, "bookingno": "DM-BK-…" },
|
||||
{ "index": 1, "success": false, "error": "pickup pincode not serviceable" }
|
||||
] }
|
||||
```
|
||||
|
||||
Per-row results, never all-or-nothing — one bad address doesn't lose the other
|
||||
199. Each booking is its own transaction. A client login cannot use bulk to
|
||||
smuggle in another tenant's id; `tenantid` is pinned to the caller's own.
|
||||
|
||||
## 3. A rider is found — automatic
|
||||
|
||||
No call needed. Creation publishes `booking.assignment_requested` to JetStream
|
||||
after the transaction commits. A worker searches riders within 10km via Redis
|
||||
GEO, scores them through the AI layer, and commits the assignment. If nobody is
|
||||
available it retries **5 times, 2 minutes apart**, then publishes
|
||||
`booking.assignment_failed` for the dispatch agent.
|
||||
|
||||
The retry state lives in NATS, not in process memory, so a pod restart no longer
|
||||
loses a booking mid-wait.
|
||||
|
||||
To assign by hand instead:
|
||||
|
||||
```
|
||||
POST /admin/bookings/:id/assign-miler
|
||||
{ "mileruserid": 38 }
|
||||
```
|
||||
|
||||
## 4. Clear a hub queue, and order the stops
|
||||
|
||||
```
|
||||
POST /hub/bookings/batch-assign
|
||||
{ "bookingids": [118, 119, 120], "max_per_rider": 5 }
|
||||
|
||||
→ { "assigned": 3, "skipped": 0, "riderssequenced": 1,
|
||||
"results": [ { "bookingid": 118, "assigned": true,
|
||||
"mileruserid": 38, "distance_km": 1.4 } ] }
|
||||
```
|
||||
|
||||
**This is the only place stops get ordered.** After assigning, each affected
|
||||
rider's whole active set is sent to the Route Optimization API
|
||||
(`POST /api/v1/optimization/doormile/sequence` on `routes.workolik.com`), which
|
||||
returns a road-network sequence via Valhalla — not straight-line distance. The
|
||||
step, per-leg distance, cumulative distance and ETA are written onto
|
||||
`bookingassignments`.
|
||||
|
||||
**Assignment picks _who_; sequencing picks _what order_.** They are separate,
|
||||
and sequencing can only run once a rider is known — which is why bulk *creation*
|
||||
cannot sequence anything, however many bookings you send. jupiter got away with
|
||||
sequencing inside `createdeliveries` because that payload was already one
|
||||
rider's run; Doormile's bulk endpoint is up to 200 bookings across many riders.
|
||||
|
||||
Sequencing is best-effort and runs after the assignments commit: the optimizer
|
||||
is a separate service over the network, and it being down must leave bookings
|
||||
assigned but unordered, never undo the batch.
|
||||
|
||||
## 5. The rider runs the route
|
||||
|
||||
```
|
||||
GET /miler/assignments
|
||||
|
||||
→ [ { "bookingassignmentid": 555, "bookingid": 120,
|
||||
"step": 1, "previouskms": 4.0, "cumulativekms": 4.0,
|
||||
"etaminutes": 14, "cumulativeeta": 14 }, … ]
|
||||
```
|
||||
|
||||
Returned in **step order**. `step: 0` means *not sequenced* — never *first* — and
|
||||
sorts to the end. A rider with one stop is never sequenced, so 0 is common.
|
||||
|
||||
Then, per booking:
|
||||
|
||||
```
|
||||
POST /miler/assignments/:id/accept
|
||||
POST /miler/bookings/:bookingid/reached
|
||||
|
||||
POST /miler/bookings/:bookingid/parcel
|
||||
{ "parcels": [ { "weight": 1.4, "length": 20, "width": 15, "height": 10 } ] }
|
||||
|
||||
POST /miler/bookings/:bookingid/payment
|
||||
{ "amount": 65, "paymentmode": "Cash", "transactionref": "" }
|
||||
|
||||
POST /miler/bookings/:bookingid/pickup-complete
|
||||
```
|
||||
|
||||
`pickup-complete` is the pivot the old system had no concept of. It converts the
|
||||
booking into a **consignment**, recomputes chargeable weight from the dimensions
|
||||
the rider actually measured, carries the kitchen across, and decides routing —
|
||||
matching 3-digit pincode prefixes go straight to `Out_for_Delivery` (hyperlocal),
|
||||
everything else routes via a hub.
|
||||
|
||||
Other rider actions: `vehicle-required` (needs a bigger vehicle), `cancel`
|
||||
(before pickup).
|
||||
|
||||
## 6. Deliver
|
||||
|
||||
```
|
||||
POST /miler/consignments/:id/deliver
|
||||
{ "deliveredtoname": "Priya R",
|
||||
"photourl": "https://…", "receiversignatureurl": "https://…",
|
||||
"lat": 11.051, "lon": 76.93,
|
||||
"otp": "418317" } // only when the tenant requires it
|
||||
```
|
||||
|
||||
Marks the consignment delivered and writes a history row. **Delivery OTP is
|
||||
opt-in per tenant** (`Tenant.Requiredeliveryotp`, default off — off for
|
||||
DailyGrubs). When on it is verified server-side and never serialised outward:
|
||||
returning it would hand the rider the code they are meant to be told.
|
||||
|
||||
Couldn't deliver? `POST /miler/consignments/:id/skip` increments `attemptcount`
|
||||
rather than failing the parcel.
|
||||
|
||||
---
|
||||
|
||||
## Watching it
|
||||
|
||||
| What | Endpoint |
|
||||
|---|---|
|
||||
| One booking end to end | `GET /admin/bookings/:id/track` |
|
||||
| A parcel's GPS trail + history | `GET /admin/consignments/:id/logs` |
|
||||
| Riders today | `GET /admin/milers/summary?from=&to=` |
|
||||
| One rider's logs | `GET /admin/milers/:id/logs` |
|
||||
| Per kitchen | `GET /admin/locations/summary?tenantid=&locationid=` |
|
||||
| Reports | `GET /admin/reports?from=&to=&tenantid=&locationid=&hubid=` |
|
||||
|
||||
Admin miler endpoints key on **`milerprofileid`**, while `assign-miler` takes a
|
||||
**`mileruserid`** in the body — different identity spaces on adjacent endpoints.
|
||||
Worth checking which one you have.
|
||||
|
||||
---
|
||||
|
||||
## State
|
||||
|
||||
| Capability | State |
|
||||
|---|---|
|
||||
| Booking create, bulk, tracking, reports | deployed |
|
||||
| Tenant scoping for client logins | deployed |
|
||||
| Per-kitchen attribution | deployed |
|
||||
| Durable assignment retry (JetStream) | built, **not deployed** |
|
||||
| Stop sequencing (Doormile side) | built, **not deployed** |
|
||||
| `/optimization/doormile/sequence` (routes.workolik.com) | built, **not deployed** |
|
||||
| Multi-stop optimizer service itself | live |
|
||||
|
||||
**Sequencing needs two deploys, not one** — the endpoint in the route-optimizer
|
||||
service (docker-compose on `31.97.228.132`, behind Traefik) and Doormile's
|
||||
client that calls it. Ship one without the other and sequencing fails quietly:
|
||||
bookings stay assigned but unordered, and riders choose their own order.
|
||||
|
||||
Nothing on the client side has moved. The rider app and express console still
|
||||
call `jupiter.nearle.app`. These endpoints exist and are tested; no production
|
||||
traffic uses them yet.
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
218
internal/assignment/queue.go
Normal file
218
internal/assignment/queue.go
Normal 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)
|
||||
}
|
||||
280
internal/routing/optimizer.go
Normal file
280
internal/routing/optimizer.go
Normal file
@@ -0,0 +1,280 @@
|
||||
// Package routing puts a rider's stops in the order they should actually be
|
||||
// run. Assignment decides *who* carries a booking; nothing in Doormile decided
|
||||
// *in what order* a rider with several stops should run them, which is the one
|
||||
// capability jupiter had that Doormile did not.
|
||||
//
|
||||
// It does not solve the routing problem itself. The Route Optimization API
|
||||
// (routes.workolik.com) already does, backed by Valhalla road-network routing
|
||||
// rather than straight-line distance, so this is a client and a writer-back.
|
||||
package routing
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"doormile/constants"
|
||||
"doormile/db"
|
||||
"doormile/models"
|
||||
"doormile/utils"
|
||||
)
|
||||
|
||||
// BaseURL is set from config at startup. Empty disables sequencing entirely,
|
||||
// which is the correct behaviour when the optimizer is not configured: stops
|
||||
// simply stay unsequenced rather than assignment failing.
|
||||
var BaseURL string
|
||||
|
||||
const (
|
||||
// Doormile's own endpoint on the Route Optimization API. It speaks
|
||||
// Doormile's vocabulary (bookingid, pickuplatitude) rather than the
|
||||
// provider one jupiter uses (deliveryid, pickuplat), validates its body,
|
||||
// and returns properly typed numbers. The provider endpoint is left alone:
|
||||
// jupiter is live on it.
|
||||
optimizePath = "/api/v1/optimization/doormile/sequence"
|
||||
|
||||
// Road-network sequencing is not instant — a real call for a handful of
|
||||
// stops took several seconds against Valhalla — but it must not hold a
|
||||
// console request open indefinitely.
|
||||
optimizeTimeout = 30 * time.Second
|
||||
|
||||
// Below this there is nothing to order.
|
||||
minStopsToSequence = 2
|
||||
)
|
||||
|
||||
var httpClient = &http.Client{Timeout: optimizeTimeout}
|
||||
|
||||
// stop is one assignment awaiting sequencing.
|
||||
type stop struct {
|
||||
AssignmentID int
|
||||
BookingID int
|
||||
BookingNo string
|
||||
PickupLat float64
|
||||
PickupLng float64
|
||||
DeliveryLat float64
|
||||
DeliveryLng float64
|
||||
}
|
||||
|
||||
// The Doormile endpoint takes real coordinates rather than the provider
|
||||
// endpoint's strings, and rejects a body it cannot read instead of returning
|
||||
// HTTP 200 with everything silently zeroed — which is what the provider
|
||||
// endpoint does when the field names are wrong, and why this one exists.
|
||||
type optimizeRequestItem struct {
|
||||
Bookingid int `json:"bookingid"`
|
||||
Bookingno string `json:"bookingno,omitempty"`
|
||||
Bookingassignmentid int `json:"bookingassignmentid,omitempty"`
|
||||
Pickuplatitude float64 `json:"pickuplatitude"`
|
||||
Pickuplongitude float64 `json:"pickuplongitude"`
|
||||
Deliverylatitude float64 `json:"deliverylatitude"`
|
||||
Deliverylongitude float64 `json:"deliverylongitude"`
|
||||
}
|
||||
|
||||
type optimizeRequest struct {
|
||||
Tenantid int `json:"tenantid,omitempty"`
|
||||
Mileruserid int `json:"mileruserid,omitempty"`
|
||||
Bookings []optimizeRequestItem `json:"bookings"`
|
||||
}
|
||||
|
||||
type optimizeResponseStop struct {
|
||||
Bookingid int `json:"bookingid"`
|
||||
Bookingassignmentid int `json:"bookingassignmentid"`
|
||||
Step int `json:"step"`
|
||||
Previouskms float64 `json:"previouskms"`
|
||||
Cumulativekms float64 `json:"cumulativekms"`
|
||||
Etaminutes int `json:"etaminutes"`
|
||||
Cumulativeeta int `json:"cumulativeeta"`
|
||||
}
|
||||
|
||||
type optimizeResponse struct {
|
||||
Success bool `json:"success"`
|
||||
Stopcount int `json:"stopcount"`
|
||||
Totalkms float64 `json:"totalkms"`
|
||||
Totaleta int `json:"totaleta"`
|
||||
Stops []optimizeResponseStop `json:"stops"`
|
||||
}
|
||||
|
||||
// Result is one sequenced stop, keyed back to the assignment it came from.
|
||||
type Result struct {
|
||||
AssignmentID int
|
||||
Step int
|
||||
PreviousKM float64
|
||||
CumulativeKM float64
|
||||
ETAMinutes int
|
||||
CumulativeETA int
|
||||
}
|
||||
|
||||
// SequenceMilerStops orders the rider's currently active stops and writes the
|
||||
// result onto their assignments.
|
||||
//
|
||||
// Best-effort by design: every failure path logs and returns an error the
|
||||
// caller is free to ignore. A rider with unsequenced stops is a worse
|
||||
// experience; a rider with no assignment at all is a broken delivery. The
|
||||
// second must never be caused by the first.
|
||||
func SequenceMilerStops(milerUserID int) ([]Result, error) {
|
||||
if BaseURL == "" {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
stops, err := loadActiveStops(milerUserID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(stops) < minStopsToSequence {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
results, err := optimize(stops)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
now := time.Now()
|
||||
for _, r := range results {
|
||||
if err := db.DB.Model(&models.BookingAssignment{}).
|
||||
Where("bookingassignmentid = ?", r.AssignmentID).
|
||||
Updates(map[string]interface{}{
|
||||
"step": r.Step,
|
||||
"previouskms": r.PreviousKM,
|
||||
"cumulativekms": r.CumulativeKM,
|
||||
"etaminutes": r.ETAMinutes,
|
||||
"cumulativeeta": r.CumulativeETA,
|
||||
"sequencedat": now,
|
||||
}).Error; err != nil {
|
||||
utils.Error("routing: failed to persist stop order",
|
||||
"assignment_id", r.AssignmentID, "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
utils.Info("routing: sequenced rider stops",
|
||||
"miler_userid", milerUserID, "stops", len(results))
|
||||
return results, nil
|
||||
}
|
||||
|
||||
// loadActiveStops returns the rider's assignments that still have to be run,
|
||||
// with the coordinates needed to order them.
|
||||
func loadActiveStops(milerUserID int) ([]stop, error) {
|
||||
var rows []struct {
|
||||
Bookingassignmentid int
|
||||
Bookingid int
|
||||
Bookingno string
|
||||
Pickuplatitude float64
|
||||
Pickuplongitude float64
|
||||
Deliverylatitude float64
|
||||
Deliverylongitude float64
|
||||
}
|
||||
|
||||
if err := db.DB.Table("bookingassignments AS ba").
|
||||
Select(`ba.bookingassignmentid, ba.bookingid, b.bookingno,
|
||||
b.pickuplatitude, b.pickuplongitude,
|
||||
b.deliverylatitude, b.deliverylongitude`).
|
||||
Joins("JOIN pickupbookings AS b ON b.bookingid = ba.bookingid").
|
||||
Where("ba.mileruserid = ? AND ba.assignmentstatus IN ?",
|
||||
milerUserID,
|
||||
[]string{constants.AssignmentAssigned, constants.AssignmentAccepted}).
|
||||
Order("ba.assignedat ASC").
|
||||
Scan(&rows).Error; err != nil {
|
||||
return nil, fmt.Errorf("load active stops: %w", err)
|
||||
}
|
||||
|
||||
stops := make([]stop, 0, len(rows))
|
||||
for _, r := range rows {
|
||||
// A stop with no coordinates cannot be ordered, and including it would
|
||||
// let the optimizer treat 0,0 as a real position off the coast of
|
||||
// Africa — which would wreck the ordering for every other stop.
|
||||
if r.Pickuplatitude == 0 || r.Pickuplongitude == 0 ||
|
||||
r.Deliverylatitude == 0 || r.Deliverylongitude == 0 {
|
||||
utils.Warn("routing: skipping stop with missing coordinates",
|
||||
"assignment_id", r.Bookingassignmentid, "booking_id", r.Bookingid)
|
||||
continue
|
||||
}
|
||||
stops = append(stops, stop{
|
||||
AssignmentID: r.Bookingassignmentid,
|
||||
BookingID: r.Bookingid,
|
||||
BookingNo: r.Bookingno,
|
||||
PickupLat: r.Pickuplatitude,
|
||||
PickupLng: r.Pickuplongitude,
|
||||
DeliveryLat: r.Deliverylatitude,
|
||||
DeliveryLng: r.Deliverylongitude,
|
||||
})
|
||||
}
|
||||
return stops, nil
|
||||
}
|
||||
|
||||
// optimize calls the Route Optimization API and maps its answer back onto our
|
||||
// assignment ids.
|
||||
func optimize(stops []stop) ([]Result, error) {
|
||||
items := make([]optimizeRequestItem, 0, len(stops))
|
||||
for _, s := range stops {
|
||||
items = append(items, optimizeRequestItem{
|
||||
Bookingid: s.BookingID,
|
||||
Bookingno: s.BookingNo,
|
||||
Bookingassignmentid: s.AssignmentID,
|
||||
Pickuplatitude: s.PickupLat,
|
||||
Pickuplongitude: s.PickupLng,
|
||||
Deliverylatitude: s.DeliveryLat,
|
||||
Deliverylongitude: s.DeliveryLng,
|
||||
})
|
||||
}
|
||||
|
||||
body, err := json.Marshal(optimizeRequest{Bookings: items})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("marshal stops: %w", err)
|
||||
}
|
||||
|
||||
req, err := http.NewRequest(http.MethodPost, BaseURL+optimizePath, bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("build request: %w", err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
|
||||
resp, err := httpClient.Do(req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("call optimizer: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||
return nil, fmt.Errorf("optimizer returned HTTP %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
var out optimizeResponse
|
||||
if err := json.NewDecoder(resp.Body).Decode(&out); err != nil {
|
||||
return nil, fmt.Errorf("decode optimizer response: %w", err)
|
||||
}
|
||||
if !out.Success || len(out.Stops) == 0 {
|
||||
return nil, fmt.Errorf("optimizer returned no sequence")
|
||||
}
|
||||
|
||||
known := make(map[int]struct{}, len(stops))
|
||||
for _, s := range stops {
|
||||
known[s.AssignmentID] = struct{}{}
|
||||
}
|
||||
|
||||
results := make([]Result, 0, len(out.Stops))
|
||||
for _, d := range out.Stops {
|
||||
if _, ok := known[d.Bookingassignmentid]; !ok {
|
||||
// Never write to an assignment we did not send. Without this an
|
||||
// echoed or stale id could reorder some other rider's work.
|
||||
utils.Warn("routing: optimizer returned unknown assignment",
|
||||
"bookingassignmentid", d.Bookingassignmentid, "bookingid", d.Bookingid)
|
||||
continue
|
||||
}
|
||||
if d.Step <= 0 {
|
||||
continue
|
||||
}
|
||||
results = append(results, Result{
|
||||
AssignmentID: d.Bookingassignmentid,
|
||||
Step: d.Step,
|
||||
PreviousKM: d.Previouskms,
|
||||
CumulativeKM: d.Cumulativekms,
|
||||
ETAMinutes: d.Etaminutes,
|
||||
CumulativeETA: d.Cumulativeeta,
|
||||
})
|
||||
}
|
||||
|
||||
if len(results) == 0 {
|
||||
return nil, fmt.Errorf("optimizer returned no usable steps")
|
||||
}
|
||||
return results, nil
|
||||
}
|
||||
172
internal/routing/optimizer_test.go
Normal file
172
internal/routing/optimizer_test.go
Normal file
@@ -0,0 +1,172 @@
|
||||
package routing
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// stubOptimizer stands in for the Route Optimization API. It captures the
|
||||
// request body so the outbound contract can be asserted, and returns whatever
|
||||
// the test tells it to.
|
||||
func stubOptimizer(t *testing.T, captured *optimizeRequest, respond func(optimizeRequest) optimizeResponse) *httptest.Server {
|
||||
t.Helper()
|
||||
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Path != optimizePath {
|
||||
t.Errorf("posted to %s, want %s", r.URL.Path, optimizePath)
|
||||
}
|
||||
var req optimizeRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
t.Errorf("stub could not decode request: %v", err)
|
||||
}
|
||||
if captured != nil {
|
||||
*captured = req
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
_ = json.NewEncoder(w).Encode(respond(req))
|
||||
}))
|
||||
}
|
||||
|
||||
func testStops() []stop {
|
||||
return []stop{
|
||||
{AssignmentID: 1, BookingID: 101, BookingNo: "A", PickupLat: 11.0045, PickupLng: 76.9612, DeliveryLat: 11.0510, DeliveryLng: 76.9300},
|
||||
{AssignmentID: 2, BookingID: 102, BookingNo: "B", PickupLat: 11.0045, PickupLng: 76.9612, DeliveryLat: 11.0168, DeliveryLng: 76.9558},
|
||||
{AssignmentID: 3, BookingID: 103, BookingNo: "C", PickupLat: 11.0045, PickupLng: 76.9612, DeliveryLat: 10.9938, DeliveryLng: 76.9954},
|
||||
}
|
||||
}
|
||||
|
||||
// The whole point of the Doormile endpoint is that it takes Doormile's field
|
||||
// names. Sending the provider ones (pickuplat/deliverylat) against the provider
|
||||
// endpoint returns HTTP 200 with everything silently zeroed, so this contract is
|
||||
// worth pinning rather than trusting.
|
||||
func TestOptimizeSendsDoormileFieldNames(t *testing.T) {
|
||||
var got optimizeRequest
|
||||
srv := stubOptimizer(t, &got, func(req optimizeRequest) optimizeResponse {
|
||||
return optimizeResponse{Success: true, Stops: []optimizeResponseStop{
|
||||
{Bookingassignmentid: 1, Bookingid: 101, Step: 1},
|
||||
}}
|
||||
})
|
||||
defer srv.Close()
|
||||
BaseURL = srv.URL
|
||||
defer func() { BaseURL = "" }()
|
||||
|
||||
if _, err := optimize(testStops()); err != nil {
|
||||
t.Fatalf("optimize: %v", err)
|
||||
}
|
||||
|
||||
if len(got.Bookings) != 3 {
|
||||
t.Fatalf("sent %d bookings, want 3", len(got.Bookings))
|
||||
}
|
||||
first := got.Bookings[0]
|
||||
if first.Bookingid != 101 || first.Bookingassignmentid != 1 {
|
||||
t.Errorf("identity fields wrong: %+v", first)
|
||||
}
|
||||
// Coordinates must survive as real numbers, not be rounded or stringified
|
||||
// into a different place.
|
||||
if first.Pickuplatitude != 11.0045 || first.Deliverylongitude != 76.93 {
|
||||
t.Errorf("coordinates mangled: %+v", first)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOptimizeMapsResultsBackToAssignments(t *testing.T) {
|
||||
srv := stubOptimizer(t, nil, func(req optimizeRequest) optimizeResponse {
|
||||
// Reverse the order, as a real resequencing would.
|
||||
return optimizeResponse{Success: true, Stops: []optimizeResponseStop{
|
||||
{Bookingassignmentid: 3, Bookingid: 103, Step: 1, Previouskms: 4, Cumulativekms: 4, Etaminutes: 14, Cumulativeeta: 14},
|
||||
{Bookingassignmentid: 2, Bookingid: 102, Step: 2, Previouskms: 5, Cumulativekms: 9, Etaminutes: 8, Cumulativeeta: 22},
|
||||
{Bookingassignmentid: 1, Bookingid: 101, Step: 3, Previouskms: 5, Cumulativekms: 14, Etaminutes: 8, Cumulativeeta: 30},
|
||||
}}
|
||||
})
|
||||
defer srv.Close()
|
||||
BaseURL = srv.URL
|
||||
defer func() { BaseURL = "" }()
|
||||
|
||||
results, err := optimize(testStops())
|
||||
if err != nil {
|
||||
t.Fatalf("optimize: %v", err)
|
||||
}
|
||||
if len(results) != 3 {
|
||||
t.Fatalf("got %d results, want 3", len(results))
|
||||
}
|
||||
if results[0].AssignmentID != 3 || results[0].Step != 1 || results[0].CumulativeETA != 14 {
|
||||
t.Errorf("first result wrong: %+v", results[0])
|
||||
}
|
||||
if results[2].AssignmentID != 1 || results[2].CumulativeKM != 14 {
|
||||
t.Errorf("last result wrong: %+v", results[2])
|
||||
}
|
||||
}
|
||||
|
||||
// These steps get written straight onto assignment rows, so a step for an
|
||||
// assignment we never sent must never make it through — it would reorder some
|
||||
// other rider's work.
|
||||
func TestOptimizeDiscardsUnknownAssignments(t *testing.T) {
|
||||
srv := stubOptimizer(t, nil, func(req optimizeRequest) optimizeResponse {
|
||||
return optimizeResponse{Success: true, Stops: []optimizeResponseStop{
|
||||
{Bookingassignmentid: 999, Bookingid: 999, Step: 1},
|
||||
{Bookingassignmentid: 2, Bookingid: 102, Step: 2, Cumulativekms: 9},
|
||||
}}
|
||||
})
|
||||
defer srv.Close()
|
||||
BaseURL = srv.URL
|
||||
defer func() { BaseURL = "" }()
|
||||
|
||||
results, err := optimize(testStops())
|
||||
if err != nil {
|
||||
t.Fatalf("optimize: %v", err)
|
||||
}
|
||||
if len(results) != 1 || results[0].AssignmentID != 2 {
|
||||
t.Fatalf("unknown assignment leaked through: %+v", results)
|
||||
}
|
||||
}
|
||||
|
||||
// Step 0 means "not sequenced". Persisting it would read as a position.
|
||||
func TestOptimizeDropsZeroSteps(t *testing.T) {
|
||||
srv := stubOptimizer(t, nil, func(req optimizeRequest) optimizeResponse {
|
||||
return optimizeResponse{Success: true, Stops: []optimizeResponseStop{
|
||||
{Bookingassignmentid: 1, Bookingid: 101, Step: 0},
|
||||
{Bookingassignmentid: 2, Bookingid: 102, Step: 1},
|
||||
}}
|
||||
})
|
||||
defer srv.Close()
|
||||
BaseURL = srv.URL
|
||||
defer func() { BaseURL = "" }()
|
||||
|
||||
results, err := optimize(testStops())
|
||||
if err != nil {
|
||||
t.Fatalf("optimize: %v", err)
|
||||
}
|
||||
if len(results) != 1 || results[0].AssignmentID != 2 {
|
||||
t.Fatalf("step 0 was kept: %+v", results)
|
||||
}
|
||||
}
|
||||
|
||||
func TestOptimizeErrorsAreSurfacedNotSilent(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
handler http.HandlerFunc
|
||||
}{
|
||||
{"http 500", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(500) }},
|
||||
{"success false", func(w http.ResponseWriter, r *http.Request) {
|
||||
_ = json.NewEncoder(w).Encode(optimizeResponse{Success: false})
|
||||
}},
|
||||
{"empty stops", func(w http.ResponseWriter, r *http.Request) {
|
||||
_ = json.NewEncoder(w).Encode(optimizeResponse{Success: true, Stops: nil})
|
||||
}},
|
||||
{"garbage body", func(w http.ResponseWriter, r *http.Request) {
|
||||
_, _ = w.Write([]byte("not json"))
|
||||
}},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
srv := httptest.NewServer(tc.handler)
|
||||
defer srv.Close()
|
||||
BaseURL = srv.URL
|
||||
defer func() { BaseURL = "" }()
|
||||
|
||||
if _, err := optimize(testStops()); err == nil {
|
||||
t.Fatal("expected an error, got nil — a failed optimise must not look like success")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
9
main.go
9
main.go
@@ -11,7 +11,9 @@ import (
|
||||
"doormile/config"
|
||||
"doormile/controllers"
|
||||
"doormile/db"
|
||||
"doormile/internal/assignment"
|
||||
"doormile/internal/notify"
|
||||
"doormile/internal/routing"
|
||||
"doormile/internal/worker"
|
||||
"doormile/middlewares"
|
||||
"doormile/migrations"
|
||||
@@ -144,6 +146,13 @@ 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()
|
||||
|
||||
// 9. Point the stop sequencer at the Route Optimization API.
|
||||
routing.BaseURL = cfg.RouteOptimizerURL
|
||||
|
||||
// 7. Startup server in a background thread
|
||||
go func() {
|
||||
utils.Info("Server starting", "port", cfg.Port)
|
||||
|
||||
@@ -125,6 +125,21 @@ type BookingAssignment struct {
|
||||
Riderkms float64 `json:"riderkms" gorm:"column:riderkms;default:0"`
|
||||
Ridercharges float64 `json:"ridercharges" gorm:"column:ridercharges;default:0"`
|
||||
Bonuspoints int `json:"bonuspoints" gorm:"column:bonuspoints;default:0"`
|
||||
|
||||
// Stop sequencing, written by internal/routing from the Route Optimization
|
||||
// API's road-network ordering. Step is 1..N across a rider's currently
|
||||
// active assignments — it says what order to run them in, which nothing in
|
||||
// Doormile decided before: assignment picked *who*, never *in what order*.
|
||||
//
|
||||
// Step 0 means "not sequenced yet", not "first". A rider with a single stop
|
||||
// is never sequenced, and sequencing is best-effort, so 0 is common and must
|
||||
// not be read as a position.
|
||||
Step int `json:"step" gorm:"column:step;default:0"`
|
||||
Previouskms float64 `json:"previouskms" gorm:"column:previouskms;default:0"`
|
||||
Cumulativekms float64 `json:"cumulativekms" gorm:"column:cumulativekms;default:0"`
|
||||
Etaminutes int `json:"etaminutes" gorm:"column:etaminutes;default:0"`
|
||||
Cumulativeeta int `json:"cumulativeeta" gorm:"column:cumulativeeta;default:0"`
|
||||
Sequencedat *time.Time `json:"sequencedat" gorm:"column:sequencedat"`
|
||||
}
|
||||
|
||||
func (BookingAssignment) TableName() string {
|
||||
|
||||
@@ -84,14 +84,20 @@ func (AppUser) TableName() string {
|
||||
}
|
||||
|
||||
type MilerProfile struct {
|
||||
Milerprofileid int `json:"milerprofileid" gorm:"primaryKey;column:milerprofileid"`
|
||||
Userid int `json:"userid" gorm:"column:userid;unique;not null"`
|
||||
Applocationid int `json:"applocationid" gorm:"column:applocationid;default:1"`
|
||||
Displayname string `json:"displayname" gorm:"column:displayname;not null"`
|
||||
Phone string `json:"phone" gorm:"column:phone;not null"`
|
||||
Profilephotourl string `json:"profilephotourl" gorm:"column:profilephotourl"`
|
||||
Vehicleid *int `json:"vehicleid" gorm:"column:vehicleid"`
|
||||
Hubid *int `json:"hubid" gorm:"column:hubid"`
|
||||
Milerprofileid int `json:"milerprofileid" gorm:"primaryKey;column:milerprofileid"`
|
||||
Userid int `json:"userid" gorm:"column:userid;unique;not null"`
|
||||
Applocationid int `json:"applocationid" gorm:"column:applocationid;default:1"`
|
||||
Displayname string `json:"displayname" gorm:"column:displayname;not null"`
|
||||
Phone string `json:"phone" gorm:"column:phone;not null"`
|
||||
Profilephotourl string `json:"profilephotourl" gorm:"column:profilephotourl"`
|
||||
Vehicleid *int `json:"vehicleid" gorm:"column:vehicleid"`
|
||||
Hubid *int `json:"hubid" gorm:"column:hubid"`
|
||||
// Legacyuserid holds this rider's userid in the old jupiter system for the
|
||||
// six riders migrated from it. Nothing reads it any more — the jupiter
|
||||
// telemetry bridge it existed for was removed — but it is kept as a record
|
||||
// of where each migrated rider came from, which is worth having during the
|
||||
// cutover. Nil for riders created natively in Doormile.
|
||||
Legacyuserid *int `json:"legacyuserid,omitempty" gorm:"column:legacyuserid;index"`
|
||||
Defaultvehicletype string `json:"defaultvehicletype" gorm:"column:defaultvehicletype"`
|
||||
Currentlatitude float64 `json:"currentlatitude" gorm:"column:currentlatitude"`
|
||||
Currentlongitude float64 `json:"currentlongitude" gorm:"column:currentlongitude"`
|
||||
@@ -111,21 +117,36 @@ func (MilerProfile) TableName() string {
|
||||
}
|
||||
|
||||
type AppCustomer struct {
|
||||
Appcustomerid int `json:"appcustomerid" gorm:"primaryKey;column:appcustomerid"`
|
||||
Firstname string `json:"firstname" gorm:"column:firstname;not null"`
|
||||
Lastname string `json:"lastname" gorm:"column:lastname"`
|
||||
Phone string `json:"phone" gorm:"column:phone;unique;not null"`
|
||||
Email string `json:"email" gorm:"column:email"`
|
||||
Loginpinhash string `json:"-" gorm:"column:loginpinhash"`
|
||||
Defaultlatitude float64 `json:"defaultlatitude" gorm:"column:defaultlatitude"`
|
||||
Defaultlongitude float64 `json:"defaultlongitude" gorm:"column:defaultlongitude"`
|
||||
Defaultpincode string `json:"defaultpincode" gorm:"column:defaultpincode"`
|
||||
Devicetoken string `json:"device_token,omitempty" gorm:"column:device_token"`
|
||||
Status string `json:"status" gorm:"column:status;default:Active"` // Active, Blocked, Deleted
|
||||
Configid int `json:"configid" gorm:"column:configid;default:1"`
|
||||
Lastloginat *time.Time `json:"lastloginat" gorm:"column:lastloginat"`
|
||||
Createdat time.Time `json:"createdat" gorm:"column:createdat;default:CURRENT_TIMESTAMP"`
|
||||
Updatedat time.Time `json:"updatedat" gorm:"column:updatedat;default:CURRENT_TIMESTAMP"`
|
||||
Appcustomerid int `json:"appcustomerid" gorm:"primaryKey;column:appcustomerid"`
|
||||
Firstname string `json:"firstname" gorm:"column:firstname;not null"`
|
||||
Lastname string `json:"lastname" gorm:"column:lastname"`
|
||||
Phone string `json:"phone" gorm:"column:phone;unique;not null"`
|
||||
Email string `json:"email" gorm:"column:email"`
|
||||
Loginpinhash string `json:"-" gorm:"column:loginpinhash"`
|
||||
Defaultlatitude float64 `json:"defaultlatitude" gorm:"column:defaultlatitude"`
|
||||
Defaultlongitude float64 `json:"defaultlongitude" gorm:"column:defaultlongitude"`
|
||||
Defaultpincode string `json:"defaultpincode" gorm:"column:defaultpincode"`
|
||||
// Flat address on the customer record, kept for parity with the console's
|
||||
// customer page (the reference stored the address this way). Separate from
|
||||
// the normalized appcustomerlocations table, which holds a customer's many
|
||||
// saved delivery addresses; this is the single profile address the console
|
||||
// edits.
|
||||
Doorno string `json:"doorno" gorm:"column:doorno"`
|
||||
Address string `json:"address" gorm:"column:address"`
|
||||
Suburb string `json:"suburb" gorm:"column:suburb"`
|
||||
City string `json:"city" gorm:"column:city"`
|
||||
State string `json:"state" gorm:"column:state"`
|
||||
Postcode string `json:"postcode" gorm:"column:postcode"`
|
||||
Landmark string `json:"landmark" gorm:"column:landmark"`
|
||||
Latitude float64 `json:"latitude" gorm:"column:latitude"`
|
||||
Longitude float64 `json:"longitude" gorm:"column:longitude"`
|
||||
Applocationid int `json:"applocationid" gorm:"column:applocationid"` // city/zone, for console scoping
|
||||
Devicetoken string `json:"device_token,omitempty" gorm:"column:device_token"`
|
||||
Status string `json:"status" gorm:"column:status;default:Active"` // Active, Blocked, Deleted
|
||||
Configid int `json:"configid" gorm:"column:configid;default:1"`
|
||||
Lastloginat *time.Time `json:"lastloginat" gorm:"column:lastloginat"`
|
||||
Createdat time.Time `json:"createdat" gorm:"column:createdat;default:CURRENT_TIMESTAMP"`
|
||||
Updatedat time.Time `json:"updatedat" gorm:"column:updatedat;default:CURRENT_TIMESTAMP"`
|
||||
}
|
||||
|
||||
func (AppCustomer) TableName() string {
|
||||
|
||||
@@ -235,6 +235,7 @@ func RegisterRoutes(app *fiber.App, cfg *config.Config) {
|
||||
|
||||
// B2C App customers
|
||||
adminAuth.Get("/customers", controllers.GetAdminCustomers)
|
||||
adminAuth.Get("/customers/summary", controllers.GetAdminCustomersSummary)
|
||||
adminAuth.Patch("/customers/:id", controllers.UpdateAdminCustomer)
|
||||
|
||||
// Hubs
|
||||
|
||||
Reference in New Issue
Block a user