Files
mesh-controller/internal/broker/streams.go
T
jochen 751e39186c Heal what is known, under a brake, and say every repair (hq to-be 45 Phase 3)
Research 031 counted the repairs people made by hand: a push to unstick a
plan waiting on a report, a controller restarted to make an object again, a
plan closed, a consumer re-made from now. Each was the ordinary path taken
again by someone who noticed. The healer registry makes each a registered
response to one condition kind, with a budget, a settle and its event:

- H1 sent-not-reported: ask the machine's node-engine to report again
  (mesh.node.<n>.ask.report); if it does not report what it was sent, send
  it again, never moving a build a policy or a plan holds back
- H2 stalled: close a plan whose wait is superseded or finished
- H3 holder-silent / consumer-lost: the send's own assertion of the bus's
  objects (issue 208's note)
- H4 consumer-behind: consumer-reset, only for a consumer the stream table
  marks resettable (the controller's own events consumer)
- H5 is the identity provider's own repair (ADR 0224 §5), registered only

Success is the observation clearing the condition, never the healer; a spent
budget hands the condition to the operator, urgent, with what was tried, and
no healer touches it again. Every act is begun in the store before it is made
(migration 0070), kept in the condition's tried as "healer Hn" and said as
the seat event healer-acted; a heal is never a hand act. More than twelve acts
in an hour stop every healer until an hour after the last, said urgently.
Only the lease holder heals.

S15 is live: a cause repaired by hand twice in a fortnight raises
healer-wanted, naming the healer that was not enough where one exists. D6's
far-behind finding has its own kind, consumer-behind. Nodes are granted the
question; the controller's grant gains healer-acted (genesis lock in
mesh-host). `healers` lists the registry, the acts and the brake; status
counts the week's heals.
2026-10-06 14:26:26 +02:00

342 lines
17 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
// Direct lets a client read a subject's last message without a consumer, which is how a
// runtime reads its own membership with no JetStream API beyond one request (ADR 0160).
Direct bool
}
// AssignmentsStream holds every assignment's membership, the newest per subject.
const AssignmentsStream = "ASSIGNMENTS"
// 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: AssignmentsStream,
Subjects: []string{"mesh.assignment.*.*"},
Retention: RetentionLastPerSubject,
Direct: true,
Why: "one membership per assignment, always the newest: what the mesh issued this module " +
"on this machine to serve and to reach (ADR 0160); read directly by the runtime it is for",
},
{
Name: EventsStream,
// 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",
},
}
}
// DeliverSubjectFor is where a push consumer's messages land.
//
// **Per consumer, which means per stream as well as per name** (novox/hq 04-ISSUES/146). A push
// consumer delivers onto an ordinary subject, and everything subscribed to that subject gets a
// copy. The controller holds a consumer called `controller` on CONTROL and another called
// `controller` on EVENTS; while both were given `_DELIVER.controller`, the one process holding
// both subscriptions acted on every message twice — a joining machine was enrolled twice from one
// request, and the second enrolment minted a credential that replaced the one the machine had just
// been handed. Every report and every followed event doubled the same way, silently: nothing is
// redelivered, no count is wrong, the work simply happens twice.
//
// The stream belongs in it because the pair is what identifies a consumer — the server scopes a
// durable's name to its stream, and this subject was the one place that scoping was dropped. It
// stays inside what a controller may already subscribe (`_DELIVER.controller.>`).
func DeliverSubjectFor(c Consumer) string {
return "_DELIVER." + c.Name + "." + c.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",
// What is wrong, said as it changes (novox/hq to-be 45 §2): a condition raised, changed in
// severity, resolver or silence, and cleared. The operator-channel's holder and any other surface
// consume them; the controller tells nobody itself.
"condition-raised", "condition-changed", "condition-cleared",
// And the self-check's heartbeat, at the end of every run (to-be 45 §4, S10): watched from a
// machine that is not the control node, so the controller going quiet is itself said.
"doctor-heartbeat",
// And a value given by hand, replaced after its module's first good start (novox/hq ADR 0228).
"secret-replaced",
// And every act a healer takes on a condition (novox/hq to-be 45 §7, Phase 3): a repair the mesh
// made by itself is said like one a person made, never quietly.
"healer-acted"}
// BusAdvisories are what the bus server says about the mesh's own account that the controller
// reads (novox/hq to-be 45 §3, S9): a durable consumer that handed a message over as often as it
// may and gave up on it, and one that was deleted. Read-only: an advisory is the server's to
// publish, and the controller's subscription changes nothing on the bus.
var BusAdvisories = []string{
"$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>",
"$JS.EVENT.ADVISORY.CONSUMER.DELETED.>",
}
// 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("node-build-agent", "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"),
// The retired build role's outcome too, while the handover runs (novox/hq ADR 0190): the one
// build machine keeps answering on its seat until build-agent replaces it, and the outcome that
// registers build-agent itself comes from there. Appended, for the same reason as above; goes
// with the retired seat row.
seatEventSubject("mesh-build-machine", "built"),
// **Every provider's standing** (novox/hq ADR 0224): a consumer it has failed for minutes, and
// that consumer recovered. The one pattern on this list, and a narrow one — two named events,
// from whichever module provides — because the rule is about every provider, and a list of
// providers here would be a list somebody forgets to extend. On 2026-10-05 the identity provider
// failed every consumer for a day and only its journal said so (issue 179). Appended, because
// the index is a name.
moduleEventSubject("*", ProvisionerFailing),
moduleEventSubject("*", ProvisionerRecovered),
}
// The provider standing events, by their local names. Written here as well as in the catalogue
// (catalogue.ProvisionerEvents), which this package cannot import; a test keeps them agreeing.
const (
ProvisionerFailing = "provisioner.failing"
ProvisionerRecovered = "provisioner.recovered"
)
// 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
}
// EventsStream holds every module's and every role's events, a build's log among them.
const EventsStream = "EVENTS"
// 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,
// One at a time (novox/hq issue 175): acting on a merge builds for minutes, and an
// announcement handed over behind it must wait on the server, not time out on the
// client and come back to be acted on again.
MaxAckPending: 1,
FromNow: true,
// **The one consumer the mesh resets by itself** (healer H4, issue 248): it fell a week
// behind once and held every merge after it. What a reset drops is caught up elsewhere —
// a merge by the catch-up pass that reads the forge (issue 266), a build's outcome from the
// build records a plan settles from (issue 214), a provider's failing word by the provider
// saying it again every quarter of an hour (ADR 0224).
Resettable: "what it drops is caught up: merges by the catch-up pass (issue 266), build outcomes " +
"from the build records (issue 214), a provider's failing word said again (ADR 0224)",
Why: "the events the mesh's own controller reacts to, one at a time; 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
}