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>
152 lines
5.4 KiB
Go
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
|
|
}
|