diff --git a/internal/broker/derived.go b/internal/broker/derived.go new file mode 100644 index 0000000..e1da042 --- /dev/null +++ b/internal/broker/derived.go @@ -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) +} diff --git a/internal/broker/derived_test.go b/internal/broker/derived_test.go new file mode 100644 index 0000000..bafe7d0 --- /dev/null +++ b/internal/broker/derived_test.go @@ -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) + } +} diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go new file mode 100644 index 0000000..59f9caf --- /dev/null +++ b/internal/broker/jetstream.go @@ -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 + } +} diff --git a/internal/broker/jetstream_test.go b/internal/broker/jetstream_test.go new file mode 100644 index 0000000..396cf16 --- /dev/null +++ b/internal/broker/jetstream_test.go @@ -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) + } + }) +} diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 3c83d4c..71c3356 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -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...…`, 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 "" } diff --git a/internal/broker/streams.go b/internal/broker/streams.go index b7be332..90302ff 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -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, diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go index 1fd41cc..bd41654 100644 --- a/internal/broker/streams_test.go +++ b/internal/broker/streams_test.go @@ -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) + } } }