Issue 127 stood because nothing compared the two halves. Every manifest was well-formed on its own and every derivation correct on its own, and no cross-module subscription in the mesh matched anything — a subscription that matches nothing is not an error, it is silence. Two checks, because the mistake is possible at two scales. Per manifest: an event is a local name, and `module.` is refused with the name to write instead. A module emitting under what reads as another module's name is refused too, pointing at the seat, where a name outlives whoever holds it. Across the catalogue: where a consumed event's emitter is present, it must emit that event. It cannot demand a live emitter for everything — a module lives in its own repository and may be installed long before the one whose events it wants — so the rule is narrower and still catches this. It found two real dangling subscriptions the moment it ran. Wildcards were undecided and two manifests needed them: `*` is one name and `**` is the rest, spelled the mesh's way and derived to `>` here and `#` on the old bus. A manifest naming either would stop being true when the wire changed, which is the whole reason names are local. And the field documentation taught the old form, examples included — which is why the drift was uniform across 37 manifests rather than scattered. Nobody was guessing; everybody followed the comment.
229 lines
9.5 KiB
Go
229 lines
9.5 KiB
Go
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
|
|
// 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 04-ISSUES/127).
|
|
Emits []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
|
|
}
|
|
|
|
// 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.<node>.>` 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)
|
|
}
|