diff --git a/db/connect.go b/db/connect.go index 01e5ed0..a7092e8 100644 --- a/db/connect.go +++ b/db/connect.go @@ -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 { diff --git a/db/streams.go b/db/streams.go new file mode 100644 index 0000000..40e530e --- /dev/null +++ b/db/streams.go @@ -0,0 +1,139 @@ +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. 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. (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", + }, + "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.; 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 +} diff --git a/middlewares/legacy_identity.go b/middlewares/legacy_identity.go new file mode 100644 index 0000000..b58fd00 --- /dev/null +++ b/middlewares/legacy_identity.go @@ -0,0 +1,76 @@ +package middlewares + +import ( + "strconv" + + "doormile/db" + "doormile/models" + "doormile/utils" + + "github.com/gofiber/fiber/v2" +) + +// LegacyMilerIdentity resolves a rider identity for machine-to-machine ingest +// and puts it in c.Locals("userid"), which is exactly where the normal miler +// handlers read it from. That is the whole point: the ingest routes reuse the +// existing handlers unchanged rather than growing a parallel set with their own +// (inevitably drifting) validation. +// +// It must be mounted *behind* InternalKeyAuth. On its own it is not +// authentication — it names a rider, it does not prove anything about the +// caller. The X-Internal-Key check is what makes that safe, and it is the reason +// this cannot be reached from the public internet. +// +// Two headers, checked in order: +// +// X-Miler-Userid a Doormile userid, used directly +// X-Legacy-Userid a jupiter/Nearle userid, resolved via MilerProfile.Legacyuserid +// +// The identity deliberately does NOT come from the request body. The miler +// telemetry handlers overwrite any body-supplied userid with the token's, which +// is what stops one rider writing another's GPS trail; taking it from the body +// here would reopen that hole from behind the internal key. A header keeps the +// rule intact and works for POST /consignments/logs, whose body is a bare JSON +// array with nowhere to put an id anyway. +func LegacyMilerIdentity(c *fiber.Ctx) error { + if raw := c.Get("X-Miler-Userid"); raw != "" { + userID, err := strconv.Atoi(raw) + if err != nil || userID <= 0 { + return utils.BadRequest(c, "invalid X-Miler-Userid") + } + + var count int64 + if err := db.DB.Model(&models.MilerProfile{}). + Where("userid = ?", userID).Count(&count).Error; err != nil { + return utils.Internal(c, "failed to resolve miler") + } + if count == 0 { + return utils.NotFound(c, "no miler with that userid") + } + + c.Locals("userid", userID) + return c.Next() + } + + raw := c.Get("X-Legacy-Userid") + if raw == "" { + return utils.BadRequest(c, "X-Miler-Userid or X-Legacy-Userid header is required") + } + + legacyID, err := strconv.Atoi(raw) + if err != nil || legacyID <= 0 { + return utils.BadRequest(c, "invalid X-Legacy-Userid") + } + + var profile models.MilerProfile + if err := db.DB.Where("legacyuserid = ?", legacyID).First(&profile).Error; err != nil { + // Unmapped is the expected case for every rider who was never migrated + // from jupiter, so this is a routine 404 rather than an error condition. + // The forwarder treats it as "drop, do not retry" — retrying cannot + // invent a mapping, and NAK-looping on it would wedge the consumer. + return utils.NotFound(c, "no Doormile miler mapped to that legacy userid") + } + + c.Locals("userid", profile.Userid) + return c.Next() +} diff --git a/models/users.go b/models/users.go index 8501c7b..ab0afdf 100644 --- a/models/users.go +++ b/models/users.go @@ -84,14 +84,21 @@ 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 is this rider's userid in the old jupiter/Nearle system, for + // riders migrated from it. It exists so telemetry still arriving over the + // jupiter NATS chain — which identifies a rider by jupiter's userid and has + // no Doormile token — can be resolved to a Doormile rider. Nil for riders + // created natively in Doormile, which is the normal case. Never used for + // authentication: it identifies, the X-Internal-Key authenticates. + 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"` diff --git a/routes/routes.go b/routes/routes.go index f1cee94..3ca6f20 100644 --- a/routes/routes.go +++ b/routes/routes.go @@ -417,6 +417,23 @@ func RegisterRoutes(app *fiber.App, cfg *config.Config) { internal.Get("/agent-decisions/similar", controllers.FindSimilarDecisions) internal.Patch("/agent-decisions/:id/outcome", controllers.UpdateDecisionOutcome) + // Rider telemetry ingest for the jupiter NATS chain. The forwarding worker + // holds no rider JWT — the rider app is still jupiter-shaped and its token is + // jupiter's — so identity arrives as a header and LegacyMilerIdentity turns + // it into c.Locals("userid"). These are the *same handlers* the miler app + // hits under /miler; only the way identity is established differs, so + // validation and storage cannot drift between the two paths. + // + // Only fire-and-forget telemetry is exposed this way. Transactional actions + // (pickup-complete, deliver, payment) are deliberately absent: the rider + // needs a real answer from those, which a queue in front of them cannot give. + ingest := internal.Group("/miler", middlewares.LegacyMilerIdentity) + ingest.Post("/logs", controllers.CreateMilerPeriodicLog) + ingest.Post("/status", controllers.CreateMilerStatus) + ingest.Post("/consignments/logs", controllers.PublishConsignmentLogs) + ingest.Post("/breaks/start", controllers.MilerStartBreak) + ingest.Put("/breaks/end", controllers.MilerEndBreak) + // -------------------- // WEBSOCKET — live miler tracking (no auth, public tracking link) // --------------------