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>
140 lines
4.7 KiB
Go
140 lines
4.7 KiB
Go
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
|
|
}
|