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. MaxAge 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", Subjects: []string{"mesh.control.*.report", "mesh.control.enrol", "mesh.control.built"}, 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: "BUILDS", Subjects: []string{"mesh.build.request"}, Retention: RetentionWorkQueue, Why: "at least once, one builder at a time; a builder that dies mid-build has its message redelivered", }, { Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}, Retention: RetentionLimits, MaxAge: 7 * 24 * 60 * 60, Why: "a subscriber that was down catches up; tool traffic under the same prefix is excluded by the event token", }, } } // 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. // // Two streams matching one subject is not a warning in NATS; it is accepted, and the message is // stored twice under two retentions. For the mesh that would mean a declaration kept as both // state and a work queue, acknowledged in one and lingering in the other — a divergence nothing // reports and nobody would think to look for. So it is refused here, where the set is written. 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 }