package broker import ( "fmt" "sort" "strings" ) // 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 // Push asks the server to deliver to a subject rather than wait to be pulled. // // For the mesh's own consumer, where the controller wants every message to arrive in the one // loop it already runs: pulling would mean a second goroutine fetching batches and handing // them over, and a loop that acts on one message at a time is the property the store window // depends on. A queue group implies this, because a group has nothing to pull from. Push bool // 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 // Emits are the verbs the seat's holder publishes under the seat's own name. An event about a // role belongs here rather than in the holder's namespace, because the name then outlives // whoever fills it (novox/hq ADR 0121, 04-ISSUES/127). Emits []string // Serves are the verbs the holder answers, request and reply. Serves []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) { // A module that reacts to anything — a module's events or a role's (novox/hq ADR 0121). Watching // a role was missing here, so the one module that does it got no consumer at all: it started, // connected, and its graph stayed empty with nothing anywhere reporting why. if p.Kind != KindModule || (len(p.Consumes) == 0 && len(p.Watches) == 0) { return Consumer{}, false } perms, err := PermissionsFor(p) if err != nil { return Consumer{}, false } // Events, wherever they live: a module's own namespace, and the namespace of any role it watches // (novox/hq ADR 0121). Tool subjects and inboxes are subscribed directly and are not a consumer's // business, which is why this is a filter and not the whole list. var filters []string for _, s := range perms.Subscribe { if strings.Contains(s, ".event.") { 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 } // NodeConsumer is the durable consumer a node reads its own declaration through. // // **Derived from a node existing, and created by the controller, because a host cannot create it.** // A host's account may subscribe its own declaration subject and publish its own ack subject, and // reaches no part of the JetStream API — which is correct (the controller is the only writer of // consumer definitions, design 25 §3) and means the consumer must be waiting before the host binds // to it. Named after the node, because the node's ack grant is `$JS.ACK.NODES..>` and a // consumer named anything else is one the host cannot acknowledge a delivery from. // // **No max-deliver, and a long ack wait.** A declaration is settled only after the node has applied // it and reported, which is minutes on a machine pulling images; and a declaration the mesh cannot // get a node to accept is not one to dead-letter, because the stream keeps only the newest per node // anyway — so there is exactly one message per node to redeliver, for as long as that node is away. func NodeConsumer(node string) Consumer { return Consumer{ Name: node, Stream: "NODES", Filters: []string{"mesh.node." + node + ".declare"}, Push: true, AckWaitSeconds: 300, Why: "how " + node + " hears what it should be; last-per-subject, so a node that was away " + "gets exactly the current declaration and nothing older", } } // AssertNodeConsumers brings every known node's declaration consumer into being. // // Asserted on start as well as created at enrolment, for the reason the streams are: a mesh raised // from a restored backup, or one whose bus was recreated, has node records and no consumers, and a // node whose consumer is missing hears nothing while everything else about it looks correct. func AssertNodeConsumers(e Ensurer, nodes []string) error { for _, n := range nodes { if err := e.EnsureConsumer(NodeConsumer(n)); err != nil { return fmt.Errorf("asserting how %s hears its declaration: %w", n, err) } } return nil } // 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) }