feat: own the JetStream subject contract, and ingest jupiter rider telemetry
Half of this binary's js.Publish calls were bound to no stream at all.
The streams were declared only by an external Python script on another
machine (Birock/doormile-bookings/setup_jetstream.py) and had drifted
from the code: booking.cancelled, booking.outcome and
booking.assignment_failed had no stream, and CHAT declared the literal
"chat.room.closed" while chat.go publishes "chat.room.closed.<id>",
which it does not match. Every publish site is best-effort
(`if db.Js != nil` + warn-log), so those events were failing and being
dropped silently — every cancellation, delivery outcome and assignment
failure since the streams were created.
db.EnsureStreams now declares the streams at startup from a map that
sits next to the code that publishes, so the contract cannot drift
again. It only ever adds: existing streams keep their storage type,
retention, limits and every subject they already have. Nothing is
deleted. Losing the create race against a sibling replica is expected
and reconciles rather than erroring.
Alongside that, /internal/miler/* ingests rider telemetry still arriving
over the jupiter NATS chain. The forwarding worker holds no rider JWT —
the rider app is still jupiter-shaped — so LegacyMilerIdentity resolves
an identity from a header into c.Locals("userid") behind the existing
X-Internal-Key guard. That lets the routes reuse the miler handlers
unchanged instead of growing a parallel set that would drift.
Identity comes from a header, never the body: the telemetry handlers
overwrite a body-supplied userid precisely so one rider cannot write
another's GPS trail, and reading it from the body here would reopen that
from behind the internal key. MilerProfile.Legacyuserid (nullable,
indexed) maps a jupiter userid to a Doormile one.
Only fire-and-forget telemetry is exposed. Transactional actions stay
synchronous — a rider needs a real answer from pickup-complete, which a
queue in front of it cannot give.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -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 {
|
||||
|
||||
139
db/streams.go
Normal file
139
db/streams.go
Normal file
@@ -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.<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",
|
||||
},
|
||||
"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
|
||||
}
|
||||
76
middlewares/legacy_identity.go
Normal file
76
middlewares/legacy_identity.go
Normal file
@@ -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()
|
||||
}
|
||||
@@ -92,6 +92,13 @@ type MilerProfile struct {
|
||||
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"`
|
||||
|
||||
@@ -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)
|
||||
// --------------------
|
||||
|
||||
Reference in New Issue
Block a user