ADR 0121, first half. The `mesh-*` seats said who does a job and nothing about what may be said to them or by them, so the mesh had roles it could not describe. They take the same three fields a module's seat has now, and the machinery that already derives a work queue, a holder's worker and a permission set from a declared seat does it for these too. The build-machine role accepts a build and emits an outcome, so `mesh.build.request`, `mesh.control.built` and the BUILDS stream are gone. A work queue shared by several build machines is what a seat's `accepts` already is, and keeping a second mechanism for it was two places a permission could be wrong. The controller's own side of a seat is a named list rather than something derived: it is not a module and declares no `uses`, so which roles the mesh itself submits work to has to be stated — and stating it makes that question answerable. Two things this caught: **The followed event subjects were hard-coded and had just gone stale.** They were written out while the catalogue still spelled its events as the old bus's routing keys, so converting those (issue 127) turned the pair into a controller listening to a subject nothing publishes — the same fault as the issue, from the other side. They derive from the emitter and the event name now, through the same function the permission uses, so the two cannot drift apart. **A role's queue exists before its holder**, checked against a real server, and asserting twice changes nothing. Work queues until somebody arrives to do it, so assigning a build machine later flushes the backlog instead of having lost it.
231 lines
9.6 KiB
Go
231 lines
9.6 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 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) {
|
|
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)
|
|
}
|