The bus on NATS: both transports behind seams, and the rollout switch #87

Merged
jschoubben merged 40 commits from feat/nats-genesis into main 2026-09-27 17:36:41 +00:00
7 changed files with 569 additions and 19 deletions
Showing only changes of commit aa74bd86ca - Show all commits
+178
View File
@@ -0,0 +1,178 @@
package broker
import (
"fmt"
"sort"
)
// Streams and consumers derived from what modules declare.
//
// The mesh's own four exist before any module does (streams.go). Everything here is the other
// half: a seat's stream comes into being when the module declaring it is **registered**, and a
// consumer when a module is **assigned** — which is why ADR 0116's task 1.4 had to be narrowed to
// the foundation set. Neither has happened at genesis.
//
// All of it is a pure function of declarations. The controller is still the only writer; this is
// only what it writes.
// A Consumer is a durable subscription the controller creates on a module's behalf. A module
// declares what it reacts to, never how delivery works, so it does not name these and cannot
// misconfigure them.
type Consumer struct {
Name string
Stream string
// Filters are the subjects this consumer receives. One consumer per module with several
// filters, rather than one per consumed event: its ack subject is derived from its name, and
// a module with five consumers would need five ack permissions to ack its own deliveries.
Filters []string
// Queue is the queue group, set for a seat's worker so that "exactly one holder" survives a
// seat later being relaxed to several. Authority and delivery are kept separate on purpose.
Queue string
// AckWaitSeconds before an unacknowledged delivery is redelivered.
AckWaitSeconds int
// MaxDeliver before the message is dead-lettered; zero for the mesh's default.
MaxDeliver int
Why string
}
// seatStreamName is the stream holding a seat's inbound work. Named after the seat rather than
// the module holding it, because the holder can change and the queued work must not care — which
// is the whole reason a caller addresses a seat instead of a module.
func seatStreamName(seat string) string { return "SEAT_" + upperSnake(seat) }
// SeatStreams is one work queue per declared seat, created when the declaring module is
// registered rather than when it is assigned.
//
// **The stream exists before anyone holds the seat, and that is the point.** Work queues until a
// holder appears, so installing the telegram module a week after something started sending to it
// flushes the backlog instead of having lost it. A stream created at assignment would make "the
// holder is not here yet" mean "your messages are gone".
func SeatStreams(seats []DeclaredSeat) []Stream {
sorted := append([]DeclaredSeat(nil), seats...)
sort.Slice(sorted, func(i, j int) bool { return sorted[i].Name < sorted[j].Name })
var out []Stream
for _, s := range sorted {
if len(s.Accepts) == 0 {
// A seat that only emits and serves needs no stream: its events ride EVENTS and its
// tools are core request/reply, which is never persisted.
continue
}
retain := s.RetainSeconds
if retain == 0 {
retain = 7 * 24 * 60 * 60
}
out = append(out, Stream{
Name: seatStreamName(s.Name),
Subjects: []string{"mesh.seat." + s.Name + ".accept.>"},
Retention: RetentionWorkQueue,
MaxAge: retain,
Why: fmt.Sprintf("work submitted to the %s seat; one holder consumes it, and it "+
"queues while nobody does", s.Name),
})
}
return out
}
// A DeclaredSeat is a seat as the catalogue knows it. Mirrored here rather than imported so this
// package stays free of the catalogue's own types — the same reason the host mirrors the
// contracts instead of importing the sdk.
type DeclaredSeat struct {
Name string
Accepts []string
RetainSeconds int
}
// ConsumerFor is the durable consumer a module's declarations imply, or false when it subscribes
// to nothing and needs none.
//
// One per module, with every consumed subject as a filter, because its ack permission is derived
// from its name: a module with a consumer per event would need an ack permission per consumer,
// and the permission list would stop being derivable from the declaration.
func ConsumerFor(p Principal) (Consumer, bool) {
if p.Kind != KindModule || len(p.Consumes) == 0 {
return Consumer{}, false
}
perms, err := PermissionsFor(p)
if err != nil {
return Consumer{}, false
}
var filters []string
for _, s := range perms.Subscribe {
if len(s) > 9 && s[:9] == "mesh.mod." {
filters = append(filters, s)
}
}
if len(filters) == 0 {
return Consumer{}, false
}
sort.Strings(filters)
return Consumer{
Name: consumerDurable(p),
Stream: consumerStream(p),
Filters: filters,
AckWaitSeconds: 30,
MaxDeliver: 5,
Why: "what " + p.Module + " declared it consumes; after max-deliver it dead-letters",
}, true
}
// HolderConsumerFor is the worker a seat's holder gets on that seat's work queue.
//
// **A queue group even though the seat guarantees one holder.** The seat is *authority* — who may
// be the telegram sender — and the queue group is *delivery*. Tie delivery to the seat and the
// day somebody allows two holders for throughput, every message is processed twice with nothing
// reporting it. Kept separate, relaxing one changes nothing about the other.
func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) {
if len(seat.Accepts) == 0 {
return Consumer{}, false
}
return Consumer{
Name: "SEAT_" + upperSnake(seat.Name) + "_worker",
Stream: seatStreamName(seat.Name),
Filters: []string{"mesh.seat." + seat.Name + ".accept.>"},
Queue: "holders",
AckWaitSeconds: 60,
MaxDeliver: 5,
Why: fmt.Sprintf("%s on %s holds %s; it acknowledges after the work is done, so a "+
"crash mid-work redelivers rather than loses", module, node, seat.Name),
}, true
}
// AllOverlaps reports subject filters claimed by more than one stream, across the mesh's own and
// every derived one.
//
// NATS refuses an overlapping stream rather than merging it (verified against nats-server 2.10:
// "subjects overlap with an existing stream"), so this is not a subtle divergence — it is a
// registration that fails. Catching it here names both streams, before a half-applied mesh does.
func AllOverlaps(seats []DeclaredSeat) []string {
all := append(MeshStreams(), SeatStreams(seats)...)
seen := map[string]string{}
var clashes []string
for _, s := range all {
for _, subject := range s.Subjects {
if first, ok := seen[subject]; ok {
clashes = append(clashes, fmt.Sprintf("%s and %s both claim %s", first, s.Name, subject))
continue
}
seen[subject] = s.Name
}
}
sort.Strings(clashes)
return clashes
}
// upperSnake makes a stream name from a seat name. NATS stream names may not contain a dot,
// a space or a wildcard, and a hyphen is legal but reads badly beside the mesh's own.
func upperSnake(s string) string {
out := []rune(s)
for i, r := range out {
switch {
case r >= 'a' && r <= 'z':
out[i] = r - 32
case r == '-' || r == '.':
out[i] = '_'
}
}
return string(out)
}
+119
View File
@@ -0,0 +1,119 @@
package broker
import (
"strings"
"testing"
)
func telegramSeat() DeclaredSeat {
return DeclaredSeat{Name: "telegram-sender", Accepts: []string{"send"}}
}
// The stream exists from registration, not assignment: work queues until a holder appears, so
// installing the module a week later flushes the backlog rather than having lost it.
func TestASeatGetsAWorkQueueOfItsOwn(t *testing.T) {
got := SeatStreams([]DeclaredSeat{telegramSeat()})
if len(got) != 1 {
t.Fatalf("expected one stream, got %d", len(got))
}
s := got[0]
if s.Retention != RetentionWorkQueue {
t.Fatalf("a seat's inbound queue retains as %q; one holder must take each message once", s.Retention)
}
if s.Subjects[0] != "mesh.seat.telegram-sender.accept.>" {
t.Fatalf("filters on %v", s.Subjects)
}
}
// A seat that only emits and serves needs no stream: its events ride EVENTS and its tools are
// core request/reply, which is never persisted.
func TestASeatThatAcceptsNothingGetsNoStream(t *testing.T) {
if got := SeatStreams([]DeclaredSeat{{Name: "announcer"}}); len(got) != 0 {
t.Fatalf("a seat with no inbound work got %d stream(s)", len(got))
}
}
// Retention belongs to whoever owns the namespace, and a seat owns its own.
func TestASeatsRetentionIsItsOwn(t *testing.T) {
s := SeatStreams([]DeclaredSeat{{Name: "slow", Accepts: []string{"work"}, RetainSeconds: 30 * 24 * 60 * 60}})
if s[0].MaxAge != 30*24*60*60 {
t.Fatalf("the seat's declared retention was not used: %d", s[0].MaxAge)
}
d := SeatStreams([]DeclaredSeat{telegramSeat()})
if d[0].MaxAge == 0 {
t.Fatal("a seat that declares no retention got an unbounded queue")
}
}
// NATS refuses an overlapping stream outright, so a clash here is a registration that fails.
func TestNoDerivedStreamOverlapsTheMeshsOwn(t *testing.T) {
seats := []DeclaredSeat{telegramSeat(), {Name: "licensing-master", Accepts: []string{"report"}}}
if c := AllOverlaps(seats); len(c) != 0 {
t.Fatalf("overlapping filters: %v", c)
}
}
// One consumer per module, with every consumed subject as a filter — because its ack permission
// is derived from its name, and a consumer per event would need an ack permission per consumer.
func TestAModuleGetsOneConsumerCarryingEveryFilter(t *testing.T) {
c, ok := ConsumerFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
Consumes: []string{"shop.order.placed", "billing.invoice.sent"}, PasswordHash: "x"})
if !ok {
t.Fatal("a module that consumes got no consumer")
}
if len(c.Filters) != 2 {
t.Fatalf("expected both subjects as filters, got %v", c.Filters)
}
perms, _ := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
Consumes: []string{"shop.order.placed"}, PasswordHash: "x"})
ack := "$JS.ACK." + c.Stream + "." + c.Name + ".>"
found := false
for _, p := range perms.Publish {
if p == ack {
found = true
}
}
if !found {
t.Fatalf("the consumer is named %q but the ack permission is %v; a module could not ack "+
"its own deliveries", c.Name, perms.Publish)
}
}
// A module that subscribes to nothing needs no consumer, and creating one would leave an object
// nothing reads and everything has to maintain.
func TestAModuleThatConsumesNothingGetsNoConsumer(t *testing.T) {
if _, ok := ConsumerFor(Principal{Kind: KindModule, Node: "one", Module: "shop",
Emits: []string{"order.placed"}, PasswordHash: "x"}); ok {
t.Fatal("a pure emitter got a consumer")
}
}
// The seat is authority and the queue group is delivery. Tie them together and the day somebody
// allows two holders, every message is processed twice with nothing reporting it.
func TestAHoldersWorkerUsesAQueueGroupAnyway(t *testing.T) {
c, ok := HolderConsumerFor("one", "telegram", telegramSeat())
if !ok {
t.Fatal("the holder of a seat with inbound work got no worker")
}
if c.Queue == "" {
t.Fatal("the worker is not in a queue group, so a second holder would double-process")
}
if c.Stream != "SEAT_TELEGRAM_SENDER" {
t.Fatalf("the worker reads %q, not the seat's own stream", c.Stream)
}
if c.MaxDeliver == 0 {
t.Fatal("a failing worker would redeliver forever rather than dead-letter")
}
}
// The stream is named after the seat, not its holder: the holder can change and the queued work
// must not care.
func TestASeatsStreamIsNamedAfterTheSeat(t *testing.T) {
name := seatStreamName("telegram-sender")
if strings.Contains(name, "telegram-sender") {
t.Fatalf("%q keeps characters a stream name may not hold", name)
}
if name != "SEAT_TELEGRAM_SENDER" {
t.Fatalf("unexpected stream name %q", name)
}
}
+142
View File
@@ -0,0 +1,142 @@
package broker
import (
"errors"
"fmt"
"time"
"github.com/nats-io/nats.go"
)
// The JetStream side of the controller: the one place the mesh's streams and consumers are
// actually created.
//
// Everything that decides *what* they are is pure and lives beside this (streams.go, derived.go).
// This is only the part that talks to a server, kept small on purpose: a bug in a subject filter
// should be findable in a unit test, and only a bug in "did the server accept it" should need one
// running.
// A JetStream is a connection to the bus, as the controller uses it.
type JetStream struct {
conn *nats.Conn
js nats.JetStreamContext
}
// Dial connects and returns the controller's JetStream handle.
func Dial(url string, opts ...nats.Option) (*JetStream, error) {
// A name, because a connection nobody can identify in the server's own monitoring is one
// nobody can attribute a problem to.
opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second))
conn, err := nats.Connect(url, opts...)
if err != nil {
return nil, fmt.Errorf("connecting to the bus at %s: %w", url, err)
}
js, err := conn.JetStream()
if err != nil {
conn.Close()
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", url, err)
}
return &JetStream{conn: conn, js: js}, nil
}
func (j *JetStream) Close() {
if j.conn != nil {
j.conn.Close()
}
}
// EnsureStream creates the stream if it is absent and brings it to match if it is present.
//
// **Idempotent, because the controller asserts on every start** rather than creating once at
// genesis: a stream somebody deleted, or a mesh raised from a restored backup, has to converge
// rather than run without the guarantee its messages assume.
//
// An update, not a delete and recreate. Recreating would discard every message the stream holds
// and every consumer's position in it — which for CONTROL means the pushes being held through a
// store restart, exactly the guarantee the stream exists for.
func (j *JetStream) EnsureStream(s Stream) error {
want := &nats.StreamConfig{
Name: s.Name,
Subjects: s.Subjects,
Retention: retentionOf(s.Retention),
MaxAge: time.Duration(s.MaxAge) * time.Second,
MaxMsgsPerSubject: int64(s.MaxMsgsPerSubject),
Description: s.Why,
}
if s.Retention == RetentionLastPerSubject {
// Last-per-subject is a limits stream with one message kept per subject, not a
// retention policy of its own — the state shape, spelled the way the server spells it.
want.Retention = nats.LimitsPolicy
want.MaxMsgsPerSubject = 1
want.MaxAge = 0
}
switch _, err := j.js.StreamInfo(s.Name); {
case err == nil:
if _, err := j.js.UpdateStream(want); err != nil {
return fmt.Errorf("bringing stream %s to match: %w", s.Name, err)
}
return nil
case errors.Is(err, nats.ErrStreamNotFound):
if _, err := j.js.AddStream(want); err != nil {
return fmt.Errorf("creating stream %s: %w", s.Name, err)
}
return nil
default:
return fmt.Errorf("asking about stream %s: %w", s.Name, err)
}
}
// EnsureConsumer creates or updates one durable consumer.
//
// Explicit acknowledgement throughout: a consumer that acknowledges on delivery cannot redeliver
// work its holder died in the middle of, which is the whole difference between a queue and a
// firehose.
func (j *JetStream) EnsureConsumer(c Consumer) error {
want := &nats.ConsumerConfig{
Durable: c.Name,
AckPolicy: nats.AckExplicitPolicy,
AckWait: time.Duration(c.AckWaitSeconds) * time.Second,
MaxDeliver: c.MaxDeliver,
DeliverGroup: c.Queue,
DeliverSubject: "",
Description: c.Why,
}
switch len(c.Filters) {
case 0:
case 1:
want.FilterSubject = c.Filters[0]
default:
want.FilterSubjects = c.Filters
}
// A queue group needs a delivery subject: a pull consumer has no group, and declaring one
// without the other is refused by the server with a message that does not say which half is
// missing.
if c.Queue != "" {
want.DeliverSubject = "_DELIVER." + c.Name
}
switch _, err := j.js.ConsumerInfo(c.Stream, c.Name); {
case err == nil:
if _, err := j.js.UpdateConsumer(c.Stream, want); err != nil {
return fmt.Errorf("bringing consumer %s on %s to match: %w", c.Name, c.Stream, err)
}
return nil
case errors.Is(err, nats.ErrConsumerNotFound):
if _, err := j.js.AddConsumer(c.Stream, want); err != nil {
return fmt.Errorf("creating consumer %s on %s: %w", c.Name, c.Stream, err)
}
return nil
default:
return fmt.Errorf("asking about consumer %s on %s: %w", c.Name, c.Stream, err)
}
}
func retentionOf(r Retention) nats.RetentionPolicy {
switch r {
case RetentionWorkQueue:
return nats.WorkQueuePolicy
default:
return nats.LimitsPolicy
}
}
+71
View File
@@ -0,0 +1,71 @@
package broker
import (
"os"
"testing"
)
// Against a real server, because the questions here are all "does the server accept this" —
// which a mock would answer by agreeing with whatever this file already believes.
//
// Skipped unless MESH_TEST_NATS names one, so the ordinary suite stays fast and offline:
//
// docker run -d --rm --name t -p 14222:4222 nats:2.10-alpine -js
// MESH_TEST_NATS=nats://127.0.0.1:14222 go test ./internal/broker/ -run TestAgainstARealServer
func TestAgainstARealServer(t *testing.T) {
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
js, err := Dial(url)
if err != nil {
t.Fatal(err)
}
defer js.Close()
t.Run("the mesh's own streams are accepted", func(t *testing.T) {
if err := AssertMeshStreams(js); err != nil {
t.Fatal(err)
}
})
t.Run("asserting again changes nothing and fails nothing", func(t *testing.T) {
if err := AssertMeshStreams(js); err != nil {
t.Fatalf("the second assertion failed, so the controller cannot restart: %v", err)
}
})
t.Run("a seat's work queue is accepted beside them", func(t *testing.T) {
seats := []DeclaredSeat{{Name: "telegram-sender", Accepts: []string{"send"}}}
for _, s := range SeatStreams(seats) {
if err := js.EnsureStream(s); err != nil {
t.Fatal(err)
}
}
if c := AllOverlaps(seats); len(c) != 0 {
t.Fatalf("overlaps the server would refuse: %v", c)
}
})
t.Run("a module's consumer is accepted and is idempotent", func(t *testing.T) {
c, ok := ConsumerFor(Principal{Kind: KindModule, Node: "one", Module: "audit",
Consumes: []string{"shop.order.placed", "billing.invoice.sent"}, PasswordHash: "x"})
if !ok {
t.Fatal("no consumer derived")
}
if err := js.EnsureConsumer(c); err != nil {
t.Fatal(err)
}
if err := js.EnsureConsumer(c); err != nil {
t.Fatalf("the second assertion failed: %v", err)
}
})
t.Run("a holder's worker is accepted with its queue group", func(t *testing.T) {
c, _ := HolderConsumerFor("one", "telegram",
DeclaredSeat{Name: "telegram-sender", Accepts: []string{"send"}})
if err := js.EnsureConsumer(c); err != nil {
t.Fatal(err)
}
})
}
+29 -8
View File
@@ -203,7 +203,7 @@ func PermissionsFor(p Principal) (Permissions, error) {
// module received would be redelivered forever, refused by the permission list it already
// has (design 25 §4). Scoped to this principal's own consumer name, so it can ack its own
// deliveries and no other's.
pub = append(pub, "$JS.ACK."+consumerName(p)+".>")
pub = append(pub, "$JS.ACK."+consumerStream(p)+"."+consumerDurable(p)+".>")
}
sort.Strings(pub)
@@ -232,17 +232,38 @@ func seatSubject(s Seat, kind, verb string) string {
return "mesh.seat." + s.Name + "." + kind + "." + verb
}
// consumerName is the durable consumer the controller derives for this principal. It is here
// rather than in the caller because the permission and the consumer must agree by construction —
// two places deriving the same name is how a module ends up unable to ack its own deliveries.
func consumerName(p Principal) string {
// consumerStream and consumerDurable are the two halves of a consumer's identity, and they are
// two functions because conflating them was a real bug.
//
// **A durable name may not contain a dot; an ack subject is built from two names that do.** The
// server acknowledges on `$JS.ACK.<stream>.<consumer>.…`, so a single string "EVENTS.one_audit"
// reads correctly inside the permission and is rejected as a consumer name — *nats: invalid
// consumer name*. Caught against a running server, and worth the comment because the shape of
// the failure if it had not been is the one design 25 §4 warns about: a consumer that cannot ack
// has every message redelivered forever, and its permission list looks right while it happens.
//
// They are derived here, beside the permission that must match them, because two places deriving
// the same name is how a module ends up unable to ack its own deliveries.
func consumerStream(p Principal) string {
switch p.Kind {
case KindModule:
return "EVENTS." + p.Node + "_" + p.Module
return "EVENTS"
case KindNode:
return "NODES." + p.Node
return "NODES"
case KindController:
return "CONTROL.controller"
return "CONTROL"
}
return ""
}
func consumerDurable(p Principal) string {
switch p.Kind {
case KindModule:
return p.Node + "_" + p.Module
case KindNode:
return p.Node
case KindController:
return "controller"
}
return ""
}
+5 -2
View File
@@ -83,8 +83,11 @@ func MeshStreams() []Stream {
Why: "at least once, one builder at a time; a builder that dies mid-build has its message redelivered",
},
{
Name: "EVENTS",
Subjects: []string{"mesh.mod.*.event.>"},
Name: "EVENTS",
// A seat's own events ride here too: they are 1:many like any event, and the
// `event` token keeps them clear of both the seat's work queue (`accept`) and its
// tools (`tool`), which must not be persisted.
Subjects: []string{"mesh.mod.*.event.>", "mesh.seat.*.event.>"},
Retention: RetentionLimits,
MaxAge: 7 * 24 * 60 * 60,
MaxMsgsPerSubject: 10000,
+25 -9
View File
@@ -73,16 +73,32 @@ func TestTheEventsStreamDoesNotCaptureToolCalls(t *testing.T) {
events = s
}
}
if len(events.Subjects) != 1 || events.Subjects[0] != "mesh.mod.*.event.>" {
t.Fatalf("EVENTS filters on %v", events.Subjects)
// Nothing a tool call rides may match any of the filters — a module's or a seat's.
for _, tool := range []string{
"mesh.mod.billing.tool.status",
"mesh.seat.telegram-sender.tool.status",
"mesh.seat.telegram-sender.accept.send", // work, not an event: its own stream
} {
for _, f := range events.Subjects {
if subjectMatches(f, tool) {
t.Fatalf("%q matches the events filter %q, so it would be persisted here", tool, f)
}
}
}
// A tool subject the composer would actually produce must not match that filter.
tool := "mesh.mod.billing.tool.status"
if subjectMatches(events.Subjects[0], tool) {
t.Fatalf("%q matches the events filter, so every tool call would be persisted", tool)
}
if !subjectMatches(events.Subjects[0], "mesh.mod.billing.event.order.placed") {
t.Fatal("an event does not match the events filter")
// And both kinds of event do match.
for _, event := range []string{
"mesh.mod.billing.event.order.placed",
"mesh.seat.telegram-sender.event.delivered",
} {
matched := false
for _, f := range events.Subjects {
if subjectMatches(f, event) {
matched = true
}
}
if !matched {
t.Fatalf("%q matches no events filter, so nothing would keep it", event)
}
}
}