Files
doormile_backend/db/streams.go
Suriya 2cbc9e5b13 docs: record that the external stream setup script is now read-only
Birock/doormile-bookings/setup_jetstream.py called update_stream with the
subject list replaced wholesale, from a hardcoded list that had drifted
behind this repo. Re-running it would have stripped booking.cancelled,
booking.outcome, booking.assignment_failed, chat.room.closed.* and --
worst -- booking.assignment_requested, which carries the assignment
retry loop rather than merely reporting on it. Bookings would have
stopped reaching riders with nothing logged, until a restart repaired
the stream via EnsureStreams. A "setup" script that causes an outage
when run is a bad thing to leave lying around.

It is now an inspector: it reports streams, subjects and consumer state
and has no add/update/delete calls at all. Verified it is the only
script anywhere that pointed at doormile-nats (66.116.226.161:4223) --
every other stream-mutating script targets jupiter's NATS.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-11 10:48:53 +05:30

152 lines
5.4 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.
//
// 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
}