The decoder named the forge's merge subject as the fourth thing followed and the list was three long: every message that fell through to that switch panicked the control plane (2026-09-28). The entry was written and lost between two attempts at the same edit. A test now walks the list; the composed grants and the genesis template carry the subject.
242 lines
11 KiB
Go
242 lines
11 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"),
|
|
// 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
|
|
}
|