The pipeline was observable from a merge to an artifact and went dark where it touched a machine: a node's report is control traffic only the control plane reads, so nothing said which version a machine runs, or that it refused to (novox/hq ADR 0134). The control plane now states both under the seat it holds — a role's events belong to the role and keep their address when the holder is replaced — and only when the report is news, because a machine reconciles every minute and a fact per report would be a fact per minute per machine. Whether a report is news is the store's answer: it holds the previous one, so the listener returns it and the server states the fact. That also gives the catch-up replay a subject the controller may publish: it was published as a module's event from a module called "control-plane", which does not exist, so the controller's own account refused it and every catalogue that asked what it missed was answered with nothing.
253 lines
12 KiB
Go
253 lines
12 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"
|
|
|
|
// ControllerSeat is the role the control plane holds, and ControllerStates are the facts it states
|
|
// under it (novox/hq ADR 0134).
|
|
//
|
|
// **Written here as well as in `link`, and a test keeps them agreeing.** `link` imports `broker`, so
|
|
// `broker` cannot import `link`; a grant naming a subject the controller never publishes is authority
|
|
// nobody uses, and a controller publishing one the grant omits is refused at the moment it has
|
|
// something to say.
|
|
const ControllerSeat = "mesh-controller"
|
|
|
|
var ControllerStates = []string{"applied", "refused", "built-before"}
|
|
|
|
// 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"),
|
|
// A build's outcome, which is the build-machine role's own event now (ADR 0121) rather than a
|
|
// message on the control branch. Same three audiences, one publish: whoever asked, this, and the
|
|
// catalogue.
|
|
seatEventSubject("mesh-build-machine", "built"),
|
|
// The forge's merges: what moved a source, so the mesh builds what that source produces
|
|
// without anybody telling it (novox/hq 04-ISSUES/131). Appended, because the index is a name.
|
|
moduleEventSubject("gitea", "pull.merged"),
|
|
}
|
|
|
|
// 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
|
|
}
|
|
|
|
// seatEventSubject is where a role's own event lands, derived the same way a holder's permission is.
|
|
func seatEventSubject(seat, verb string) string {
|
|
return "mesh.seat." + seat + ".event." + verb
|
|
}
|
|
|
|
// 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
|
|
}
|