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 }