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.
230 lines
10 KiB
Go
230 lines
10 KiB
Go
package broker
|
|
|
|
import (
|
|
"fmt"
|
|
"sort"
|
|
)
|
|
|
|
// The mesh's own streams.
|
|
//
|
|
// **These four and no more** (novox/hq ADR 0116 task 1.4, as revised by ADR 0118). An earlier
|
|
// reading had the controller create *every* stream at genesis, from a fixed set. That is only the
|
|
// mesh's own half: a seat's streams are created when the module declaring it is registered, and a
|
|
// module's durable consumers when it is assigned — neither of which has happened at genesis. What
|
|
// is here is the foundation, which exists before any module does.
|
|
//
|
|
// The controller is the only writer of stream definitions (design 25 §3). A module declares
|
|
// nothing about them and cannot reach the JetStream API to make one.
|
|
|
|
// Retention is how a stream decides what to keep, which is the whole of what distinguishes the
|
|
// mesh's four relationships on the wire (design 29 §4).
|
|
type Retention string
|
|
|
|
const (
|
|
// RetentionWorkQueue: a message is removed once a consumer acknowledges it. Exactly one
|
|
// worker does the work, and a worker that dies has its message redelivered.
|
|
RetentionWorkQueue Retention = "workqueue"
|
|
// RetentionLastPerSubject: only the newest message on each subject survives. This is the
|
|
// state shape — a node that was away gets exactly the current declaration and nothing older.
|
|
RetentionLastPerSubject Retention = "last_per_subject"
|
|
// RetentionLimits: kept until it ages or the stream fills. Events, where a subscriber that
|
|
// was down catches up and nobody is obliged to act.
|
|
RetentionLimits Retention = "limits"
|
|
)
|
|
|
|
// A Stream is one of the mesh's own, as the controller asserts it.
|
|
type Stream struct {
|
|
Name string
|
|
Subjects []string
|
|
Retention Retention
|
|
// MaxAge in seconds, zero for unbounded. Per stream — JetStream has no per-subject age,
|
|
// which is why differing retention between modules would mean a stream each.
|
|
MaxAge int
|
|
// MaxMsgsPerSubject caps each subject independently, so one noisy emitter cannot push
|
|
// another's events out of a shared stream. Verified: with a cap of 3, ten messages on one
|
|
// subject and one on another leave four in the stream, not three.
|
|
MaxMsgsPerSubject int
|
|
// Why is carried into the assertion so an operator reading the server's own state finds the
|
|
// reason there, rather than only in a repository they may not have.
|
|
Why string
|
|
}
|
|
|
|
// MeshStreams is the foundation set, in the order a person reads it.
|
|
//
|
|
// **CONTROL names its subjects rather than taking `mesh.control.>`**, because heartbeats live
|
|
// under that prefix and must not be persisted: a lost heartbeat is the next heartbeat, and a
|
|
// stream of them is a stream of the least valuable messages the mesh sends, competing for the
|
|
// same retention as the ones that matter.
|
|
//
|
|
// **EVENTS filters on the `event` token**, which is the reason that token exists. A module's
|
|
// namespace carries both its events and its tool calls; a filter of `mesh.mod.*.>` would persist
|
|
// every tool invocation in the mesh, and a tool call must never be persisted (design 25 §3 keeps
|
|
// tools on core NATS, where a lost call is a timeout the caller already handles).
|
|
func MeshStreams() []Stream {
|
|
return []Stream{
|
|
{
|
|
Name: "CONTROL",
|
|
// A build's outcome is no longer here: it is the build-machine seat's own event, so one
|
|
// publish reaches whoever asked, the controller and the catalogue (novox/hq ADR 0121).
|
|
Subjects: []string{"mesh.control.*.report", "mesh.control.enrol"},
|
|
Retention: RetentionWorkQueue,
|
|
Why: "the store-window guarantee (ADR 0083): the controller naks with a delay while its " +
|
|
"store is away and the message is redelivered; nothing is dropped",
|
|
},
|
|
{
|
|
Name: "NODES",
|
|
Subjects: []string{"mesh.node.*.declare"},
|
|
Retention: RetentionLastPerSubject,
|
|
Why: "one declaration per node, always the newest; a node that sees sequence n refuses " +
|
|
"n-1 by construction (issue 107)",
|
|
},
|
|
{
|
|
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,
|
|
Why: "a subscriber that was down catches up; tool traffic under the same prefix is " +
|
|
"excluded by the event token; per-subject caps keep a noisy emitter from " +
|
|
"evicting a quiet one without splitting the stream",
|
|
},
|
|
}
|
|
}
|
|
|
|
// An Asserter is the part of a JetStream connection stream assertion needs. Narrow on purpose: it
|
|
// keeps this testable without a server, and keeps the client library out of everything that only
|
|
// wants to know what the streams are.
|
|
type Asserter interface {
|
|
// EnsureStream creates the stream if absent and updates it to match if present. It must be
|
|
// idempotent: the controller asserts on every start, not only at genesis.
|
|
EnsureStream(s Stream) error
|
|
}
|
|
|
|
// AssertMeshStreams brings the foundation set into being, in order, and says which one failed
|
|
// rather than that something did.
|
|
//
|
|
// Asserted on every start rather than created once at genesis, because a stream that was deleted,
|
|
// or a mesh raised from a restored backup, must converge rather than run without the guarantee
|
|
// its messages assume. Idempotence is the whole requirement.
|
|
func AssertMeshStreams(a Asserter) error {
|
|
for _, s := range MeshStreams() {
|
|
if err := a.EnsureStream(s); err != nil {
|
|
return fmt.Errorf("asserting stream %s: %w", s.Name, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Overlaps reports subject filters claimed by more than one stream.
|
|
//
|
|
// **Corrected against the server**: an earlier version of this comment said NATS accepts
|
|
// overlapping streams and stores the message twice. It does not — it refuses the second stream
|
|
// with "subjects overlap with an existing stream" (verified against nats-server 2.10). The check
|
|
// still earns its place, for a different reason: the server's refusal arrives when the controller
|
|
// is applying, naming one stream, at a moment when the mesh is half-configured. This one arrives
|
|
// where the set is written, names both, and cannot reach a running mesh.
|
|
//
|
|
// It also decides a design question. Because overlap is refused rather than merged, a shared
|
|
// EVENTS stream and a per-module stream cannot coexist — the module's would be refused — so
|
|
// "one stream for most, its own for a module that wants different retention" is not an option
|
|
// the server allows. It is all of one or all of the other.
|
|
func Overlaps() []string {
|
|
seen := map[string]string{}
|
|
var clashes []string
|
|
for _, s := range MeshStreams() {
|
|
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
|
|
}
|
|
|
|
// The mesh's own consumers.
|
|
//
|
|
// A seat's streams and a module's consumers are derived from declarations (derived.go). These two
|
|
// are not: **the controller is not a module and files no manifest**, so its authority and its
|
|
// subscriptions cannot come from a declaration that does not exist. They are named here, where the
|
|
// mesh's own streams are named, and narrowly — a controller subscribing `mesh.mod.*.event.>` would
|
|
// hear every event in the mesh, which it has no business doing and which would make its permission
|
|
// list stop explaining anything.
|
|
|
|
// ControllerName is the controller's durable consumer on each stream it reads, and the name its
|
|
// ack subject is derived from (nats.go: `$JS.ACK.<stream>.controller.>`).
|
|
const ControllerName = "controller"
|
|
|
|
// ControllerFollows are the events the controller reacts to: the catalogue saying a module's
|
|
// current version moved, and a catalogue that has just started saying it may have missed builds.
|
|
//
|
|
// **Derived the same way a module's subscription is**, from the emitter and the bare local event
|
|
// name, rather than written out. They were written out while the catalogue still spelled its events
|
|
// as the old bus's routing keys, and the moment those were converted (novox/hq 04-ISSUES/127) a
|
|
// hard-coded pair became a controller listening to a subject nothing publishes — the same fault, from
|
|
// the other side. Deriving them means the conversion could not leave these behind.
|
|
var ControllerFollows = []string{
|
|
moduleEventSubject("mesh-catalog", "upgraded"),
|
|
moduleEventSubject("mesh-catalog", "catching-up"),
|
|
}
|
|
|
|
// moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so
|
|
// what the controller subscribes and what the emitter is permitted to publish cannot drift apart.
|
|
func moduleEventSubject(module, event string) string {
|
|
return "mesh.mod." + module + ".event." + event
|
|
}
|
|
|
|
// MeshConsumers is what the controller consumes, in the order a person reads it.
|
|
//
|
|
// **Unlimited redelivery on CONTROL, deliberately.** The store window's bound is the controller's,
|
|
// not the server's (window.go): a message is held with a nak-and-delay until the controller either
|
|
// takes it or gives up and says so. A max-deliver here would dead-letter a push that was being
|
|
// held through a store restart — the exact message the stream exists to protect — some minutes
|
|
// before the controller had finished deciding about it.
|
|
func MeshConsumers() []Consumer {
|
|
return []Consumer{
|
|
{
|
|
Name: ControllerName,
|
|
Stream: "CONTROL",
|
|
Push: true,
|
|
AckWaitSeconds: 30,
|
|
Why: "the controller is the single consumer of what nodes say; explicit ack and no " +
|
|
"max-deliver, because the store window's bound is the controller's own",
|
|
},
|
|
{
|
|
Name: ControllerName,
|
|
Stream: "EVENTS",
|
|
Filters: ControllerFollows,
|
|
Push: true,
|
|
AckWaitSeconds: 30,
|
|
MaxDeliver: 5,
|
|
Why: "the two events the mesh's own controller reacts to; after max-deliver it " +
|
|
"dead-letters, because an announcement it cannot act on will not become actionable",
|
|
},
|
|
}
|
|
}
|
|
|
|
// Ensurer is the part of a JetStream connection consumer assertion needs, narrow for the reason
|
|
// Asserter is.
|
|
type Ensurer interface {
|
|
EnsureConsumer(c Consumer) error
|
|
}
|
|
|
|
// AssertMeshConsumers brings the controller's own consumers into being, and says which one failed.
|
|
//
|
|
// After the streams, necessarily: a consumer on a stream that does not exist is refused, and the
|
|
// refusal names the stream rather than the order.
|
|
func AssertMeshConsumers(e Ensurer) error {
|
|
for _, c := range MeshConsumers() {
|
|
if err := e.EnsureConsumer(c); err != nil {
|
|
return fmt.Errorf("asserting consumer %s on %s: %w", c.Name, c.Stream, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|