The bus on NATS: both transports behind seams, and the rollout switch #87

Merged
jschoubben merged 40 commits from feat/nats-genesis into main 2026-09-27 17:36:41 +00:00
5 changed files with 408 additions and 12 deletions
Showing only changes of commit 7232d6df4b - Show all commits
+24
View File
@@ -192,6 +192,26 @@ type Manifest struct {
// that every new module would force its predecessors to update. // that every new module would force its predecessors to update.
Claims []Claim `json:"claims,omitempty"` Claims []Claim `json:"claims,omitempty"`
// Seats this module declares of its own, with their protocols (novox/hq ADR 0118). The set
// of seats a mesh has is the mesh's own plus these, derived from what is registered rather
// than written in the controller — closed, and extensible without changing the mesh.
Seats []SeatDeclaration `json:"seats,omitempty"`
// Uses are seats this module sends to. It names the *seat*, never the module holding it, so
// the implementation can be replaced under it and no caller changes. A caller gets publish
// on that seat's inbound subjects and nothing else — not its outbound events, and not a
// subscription to the queue it writes to (design 29 §2).
Uses []string `json:"uses,omitempty"`
// Tools are the tools this module answers — request and reply, awaited.
//
// **New, and not `serves`**, which this manifest already uses for the facts a consumer needs
// in order to reach a provision. Two meanings under one key would be a footgun in the one
// file a module author reads most. Until now a module's tools were known only at runtime,
// from MESH_TOOL_MODULES in its image; declaring them is what lets the mesh check that a
// module claiming a seat answers what that seat's protocol promises (novox/hq ADR 0118).
Tools []string `json:"tools,omitempty"`
// Capabilities the machine must have. A different field from Requires because the remedy // Capabilities the machine must have. A different field from Requires because the remedy
// differs: a missing module can be assigned, and a missing capability means the wrong // differs: a missing module can be assigned, and a missing capability means the wrong
// machine. // machine.
@@ -970,6 +990,10 @@ func ParseManifest(raw []byte) (Manifest, error) {
if wellFormed { if wellFormed {
problems = append(problems, claimProblems(m)...) problems = append(problems, claimProblems(m)...)
} }
// What one manifest can be judged on: a declaration's shape, its scope, and the reserved
// prefix. Whether a seat anybody names exists, and whether a holder answers for it, are
// facts about the catalogue and are checked at registration (CatalogueProblems).
problems = append(problems, declaredSeatProblems(m)...)
if m.Computed != "" && len(m.Resources) > 0 { if m.Computed != "" && len(m.Resources) > 0 {
// One or the other. A module that both ships files and has them computed would leave // One or the other. A module that both ships files and has them computed would leave
// nobody able to say where a given file came from. // nobody able to say where a given file came from.
+5 -3
View File
@@ -100,9 +100,11 @@ func claimProblems(m Manifest) []string {
for _, c := range m.Claims { for _, c := range m.Claims {
seat, known := SeatNamed(c.Name) seat, known := SeatNamed(c.Name)
if !known { if !known {
problems = append(problems, fmt.Sprintf( // Not one of the mesh's own, which no longer means it is not a seat: a module may
"%s claims %q, which is not a seat this mesh defines (novox/hq ADR 0110) — "+ // declare its own (novox/hq ADR 0118), and whether anybody declared *this* one is a
"the seats are: %s", m.Module, c.Name, seatNames())) // fact about the catalogue rather than about this manifest. Deferred to
// CatalogueProblems, which refuses it at registration — the same guarantee ADR 0110
// wanted, at the same moment, from a set nobody maintains by hand.
continue continue
} }
if c.At() != seat.Scope { if c.At() != seat.Scope {
+240
View File
@@ -0,0 +1,240 @@
package catalogue
import (
"fmt"
"sort"
"strings"
)
// Seats a module declares of its own (novox/hq ADR 0118).
//
// The set of seats a mesh has is **derived**: the mesh's own, in seats.go, plus those declared by
// every module it has registered. Still closed — a seat named nowhere is refused — but computed
// from the catalogue rather than written in the controller, which is what ADR 0110 actually
// needed and a hand-maintained table could not keep. Its own evidence: the enumeration done by
// hand while that record was written reported eleven claims where there were thirteen.
//
// **What can be checked from one manifest and what cannot.** A declaration's shape, its scope,
// and the reserved prefix are facts about the manifest in front of you. Whether a seat anybody
// names actually exists, whether two modules declared the same one, and whether a holder
// satisfies the protocol are facts about the *catalogue* — so they are checked at registration,
// by CatalogueProblems, which is the last moment the mesh can still say no.
// meshSeatPrefix is reserved to the mesh. The prefix *is* the reservation rule: no list of
// reserved names to maintain, no way for the mesh's own namespace to be colonised by a manifest,
// and nothing to keep in step when a mesh seat is added.
const meshSeatPrefix = "mesh-"
// A SeatDeclaration is a role a module offers on the bus: what may be sent to it, what it says,
// and what it answers. A caller declares that it uses the *seat*, never the module, so the
// implementation can be replaced under it.
type SeatDeclaration struct {
Name string `json:"name"`
Scope string `json:"scope,omitempty"`
// Accepts are the verbs others may submit work on. Each becomes a work-queue subject, and
// the holder is the only consumer — so exactly one worker does the job, by construction
// rather than by how carefully somebody wrote a subscribe call.
Accepts []string `json:"accepts,omitempty"`
// Emits are the verbs the holder publishes: 1:many, nobody obliged to act.
Emits []string `json:"emits,omitempty"`
// Serves are the verbs the holder answers: request and reply, awaited.
Serves []string `json:"serves,omitempty"`
// RetainSeconds is how long the inbound backlog survives with no holder, zero for the
// mesh's default. Retention belongs to whoever owns the namespace (design 29 §3) — a seat
// owns its own, which is why a seat is also the answer for a module that needs retention
// its events cannot have.
RetainSeconds int `json:"retain-seconds,omitempty"`
}
// At is this declaration's scope, with the default applied. Mesh by default, because a seat
// declared by a module is nearly always "there is one of these in the mesh" — a per-node worker
// is the deliberate case, and says so.
func (s SeatDeclaration) At() string {
if s.Scope == "" {
return ScopeMesh
}
return s.Scope
}
// verbs is everything the protocol names, for the checks that do not care which half.
func (s SeatDeclaration) verbs() []string {
out := append([]string{}, s.Accepts...)
out = append(out, s.Emits...)
return append(out, s.Serves...)
}
// declaredSeatProblems is what one manifest can be judged on alone.
func declaredSeatProblems(m Manifest) []string {
var problems []string
seen := map[string]bool{}
for _, s := range m.Seats {
switch {
case s.Name == "":
problems = append(problems, fmt.Sprintf("%s declares a seat with no name", m.Module))
continue
case !name.MatchString(s.Name):
problems = append(problems, fmt.Sprintf(
"%s declares a seat named %q, which is not a usable name", m.Module, s.Name))
continue
case strings.HasPrefix(s.Name, meshSeatPrefix):
// The mesh's own code dereferences its seats by name — the resolver *is* the thing
// that finds the store — so the prefix is not a convention, it is a namespace.
problems = append(problems, fmt.Sprintf(
"%s declares a seat named %q; %q is reserved to the mesh, which defines its own "+
"seats (novox/hq ADR 0118)", m.Module, s.Name, meshSeatPrefix+"*"))
continue
}
if seen[s.Name] {
problems = append(problems, fmt.Sprintf(
"%s declares the seat %q twice", m.Module, s.Name))
continue
}
seen[s.Name] = true
if _, isMesh := SeatNamed(s.Name); isMesh {
problems = append(problems, fmt.Sprintf(
"%s declares %q, which is a seat the mesh already defines", m.Module, s.Name))
}
switch s.At() {
case ScopeNode, ScopeSite, ScopeMesh:
default:
problems = append(problems, fmt.Sprintf(
"%s declares seat %s at scope %q; a seat is held per node, per site or per mesh",
m.Module, s.Name, s.Scope))
}
if len(s.verbs()) == 0 {
// A seat with no protocol is exclusion with nothing on the other side of it. If a
// module only wants "there is one of me", that is a claim, and saying so keeps the
// word "seat" meaning something callers can depend on.
problems = append(problems, fmt.Sprintf(
"%s declares seat %s with no protocol; a seat says what may be sent to it, what "+
"it emits and what it serves", m.Module, s.Name))
}
for _, v := range s.verbs() {
if !name.MatchString(v) {
problems = append(problems, fmt.Sprintf(
"%s declares %s.%s, which is not a usable verb", m.Module, s.Name, v))
}
}
}
for _, u := range m.Uses {
if !name.MatchString(u) {
problems = append(problems, fmt.Sprintf("%s uses %q, which is not a usable seat name", m.Module, u))
}
}
return problems
}
// A Shelf is every manifest the mesh has registered, by module name.
type Shelf map[string]Manifest
// CatalogueProblems are the rules no single manifest can be judged against.
//
// Run at registration, which is the last moment the mesh can still refuse: after it, a caller is
// bound to a seat and a refusal is an outage rather than a conversation.
func CatalogueProblems(shelf Shelf) []string {
var problems []string
// Who declares what, and who declared it first.
declaredBy := map[string]string{}
declared := map[string]SeatDeclaration{}
for _, module := range shelfOrder(shelf) {
for _, s := range shelf[module].Seats {
if s.Name == "" {
continue
}
if first, taken := declaredBy[s.Name]; taken {
// The second loses. A seat name meaning two different protocols is the failure
// nobody could diagnose afterwards — a caller would bind to whichever happened
// to register first, and the symptom would appear in the other module.
problems = append(problems, fmt.Sprintf(
"%s declares the seat %q, which %s already declares; a seat name means one "+
"protocol", module, s.Name, first))
continue
}
declaredBy[s.Name] = module
declared[s.Name] = s
}
}
exists := func(seat string) bool {
if _, isMesh := SeatNamed(seat); isMesh {
return true
}
_, ok := declaredBy[seat]
return ok
}
for _, module := range shelfOrder(shelf) {
m := shelf[module]
// A `uses` naming nothing is where ADR 0110's guarantee lands under a derived set: the
// same refusal, at the same moment, from a set nobody maintains by hand.
for _, u := range m.Uses {
if !exists(u) {
problems = append(problems, fmt.Sprintf(
"%s uses the seat %q, which no module declares and the mesh does not define",
module, u))
}
}
for _, c := range m.Claims {
if !exists(c.Name) {
problems = append(problems, fmt.Sprintf(
"%s claims the seat %q, which no module declares and the mesh does not define",
module, c.Name))
continue
}
s, isModuleSeat := declared[c.Name]
if !isModuleSeat {
continue // a mesh seat: already judged by claimProblems
}
if c.At() != s.At() {
problems = append(problems, fmt.Sprintf(
"%s claims %s at scope %q, and %s declares it at %s",
module, c.Name, c.At(), declaredBy[c.Name], s.At()))
}
// A holder that does not answer what the seat promises is a caller's timeout, found
// at assignment instead.
if missing := unserved(m, s); len(missing) > 0 {
problems = append(problems, fmt.Sprintf(
"%s claims %s but does not serve %s, which that seat's protocol promises",
module, c.Name, strings.Join(missing, ", ")))
}
}
}
sort.Strings(problems)
return problems
}
// unserved is what a seat's protocol promises and the claimant does not answer. Only the tools
// are checked: `accepts` and `emits` are wired by the runtime from the declaration, while a tool
// is code the module either has or has not written.
func unserved(m Manifest, s SeatDeclaration) []string {
has := map[string]bool{}
for _, t := range m.Tools {
has[t] = true
}
var missing []string
for _, t := range s.Serves {
if !has[t] {
missing = append(missing, t)
}
}
return missing
}
// shelfOrder is the catalogue in a stable order, so two runs report the same problems in the same
// sequence — a refusal that reorders itself is a refusal nobody can diff.
func shelfOrder(shelf Shelf) []string {
out := make([]string, 0, len(shelf))
for k := range shelf {
out = append(out, k)
}
sort.Strings(out)
return out
}
+107
View File
@@ -0,0 +1,107 @@
package catalogue
import (
"strings"
"testing"
)
func problemsFor(t *testing.T, shelf Shelf) string {
t.Helper()
return strings.Join(CatalogueProblems(shelf), "; ")
}
func telegram() Manifest {
return Manifest{Module: "telegram", Tools: []string{"status"}, Seats: []SeatDeclaration{{
Name: "telegram-sender", Scope: ScopeMesh,
Accepts: []string{"send"}, Emits: []string{"delivered", "failed"}, Serves: []string{"status"},
}}, Claims: []Claim{{Name: "telegram-sender", Scope: ScopeMesh}}}
}
// The whole point: a module contributes a capability without the mesh being changed.
func TestAModuleDeclaresItsOwnSeatAndHoldsIt(t *testing.T) {
shop := Manifest{Module: "shop", Uses: []string{"telegram-sender"}}
if got := problemsFor(t, Shelf{"telegram": telegram(), "shop": shop}); got != "" {
t.Fatalf("a declared seat and its caller were refused: %s", got)
}
}
// The prefix is the reservation rule, so there is no list to maintain and none to drift.
func TestAModuleCannotDeclareAMeshSeat(t *testing.T) {
for _, n := range []string{"mesh-broker", "mesh-anything", "mesh-store"} {
m := Manifest{Module: "impostor", Seats: []SeatDeclaration{{Name: n, Accepts: []string{"x"}}}}
got := strings.Join(declaredSeatProblems(m), "; ")
if !strings.Contains(got, "reserved to the mesh") {
t.Fatalf("%q was accepted as a module's seat: %q", n, got)
}
}
}
// A seat name meaning two protocols is the failure nobody could diagnose afterwards.
func TestTwoModulesCannotDeclareTheSameSeat(t *testing.T) {
other := Manifest{Module: "aardvark", Seats: []SeatDeclaration{{
Name: "telegram-sender", Scope: ScopeMesh, Accepts: []string{"something-else"}}}}
got := problemsFor(t, Shelf{"telegram": telegram(), "aardvark": other})
if !strings.Contains(got, "already declares") {
t.Fatalf("both declarations stood: %s", got)
}
// The first declarer keeps it; only the second is refused.
if strings.Count(got, "already declares") != 1 {
t.Fatalf("expected exactly one refusal: %s", got)
}
}
// Where ADR 0110's guarantee lands under a derived set: a typo is refused, not resolved to
// nothing at runtime.
func TestUsingASeatNobodyDeclaresIsRefused(t *testing.T) {
shop := Manifest{Module: "shop", Uses: []string{"telegram-sendr"}}
got := problemsFor(t, Shelf{"telegram": telegram(), "shop": shop})
if !strings.Contains(got, "telegram-sendr") || !strings.Contains(got, "no module declares") {
t.Fatalf("a misspelled seat was accepted: %s", got)
}
}
// A holder that does not answer what the seat promises is a caller's timeout, found here instead.
func TestAHolderMustServeWhatItsSeatPromises(t *testing.T) {
m := telegram()
m.Tools = nil // declares the seat, serves none of it
got := problemsFor(t, Shelf{"telegram": m})
if !strings.Contains(got, "does not serve status") {
t.Fatalf("a holder was accepted that answers nothing its seat promises: %s", got)
}
}
// A seat with no protocol is exclusion with nothing on the other side of it.
func TestASeatWithoutAProtocolIsRefused(t *testing.T) {
m := Manifest{Module: "vague", Seats: []SeatDeclaration{{Name: "something", Scope: ScopeMesh}}}
if got := strings.Join(declaredSeatProblems(m), "; "); !strings.Contains(got, "no protocol") {
t.Fatalf("a seat promising nothing was accepted: %s", got)
}
}
// A claim at the wrong scope is a different seat than the one declared.
func TestAClaimMustMatchTheDeclaredScope(t *testing.T) {
m := telegram()
m.Claims = []Claim{{Name: "telegram-sender", Scope: ScopeNode}}
got := problemsFor(t, Shelf{"telegram": m})
if !strings.Contains(got, "scope") {
t.Fatalf("a claim at the wrong scope was accepted: %s", got)
}
}
// The mesh's own seats still work, and are not shadowed by the derived half.
func TestTheMeshsOwnSeatsAreStillClaimable(t *testing.T) {
m := Manifest{Module: "nats", Claims: []Claim{{Name: "mesh-broker", Scope: ScopeMesh}}}
if got := problemsFor(t, Shelf{"nats": m}); got != "" {
t.Fatalf("a mesh seat was refused by the derived check: %s", got)
}
}
// A refusal that reorders itself between runs is a refusal nobody can diff.
func TestTheProblemsAreStable(t *testing.T) {
shelf := Shelf{"telegram": telegram(), "shop": {Module: "shop", Uses: []string{"nope"}},
"other": {Module: "other", Uses: []string{"also-nope"}}}
first, second := problemsFor(t, shelf), problemsFor(t, shelf)
if first != second {
t.Fatalf("unstable:\n%s\n%s", first, second)
}
}
+32 -9
View File
@@ -54,17 +54,40 @@ func claimed(claims string) []byte {
return []byte(`{"module":"thing","version":"1","provides":[{"name":"npm-package-registry","scope":"mesh"}],"claims":` + claims + `}`) return []byte(`{"module":"thing","version":"1","provides":[{"name":"npm-package-registry","scope":"mesh"}],"claims":` + claims + `}`)
} }
func TestAClaimOnASeatTheMeshDoesNotDefineIsRefused(t *testing.T) { // **The refusal moved, it did not go** (novox/hq ADR 0118, superseding 0110). A module may now
_, err := ParseManifest(claimed(`[{"name":"the-anything","scope":"node"}]`)) // declare its own seats, so whether a claimed seat exists is a fact about the *catalogue* and not
if err == nil { // about the manifest in front of the parser: a claim on a seat another registered module declares
t.Fatal("a module invented a seat by claiming it") // is perfectly good, and the parser cannot tell the two cases apart. So the parser accepts it and
// registration refuses it — the same guarantee, at the same moment work would otherwise start,
// from a set nobody maintains by hand.
func TestAClaimOnASeatNobodyDeclaresIsRefusedAtRegistration(t *testing.T) {
m, err := ParseManifest(claimed(`[{"name":"the-anything","scope":"node"}]`))
if err != nil {
t.Fatalf("the parser judged a claim it cannot judge alone: %v", err)
} }
if !strings.Contains(err.Error(), "the-anything") || !strings.Contains(err.Error(), "not a seat") {
t.Fatalf("the refusal does not say the seat is unknown: %v", err) problems := CatalogueProblems(Shelf{m.Module: m})
if len(problems) == 0 {
t.Fatal("a module invented a seat by claiming it, and registration allowed it")
} }
// And it says what the seats are, because "no" without the list sends somebody reading code. joined := strings.Join(problems, "; ")
if !strings.Contains(err.Error(), "the-packet-filter") { if !strings.Contains(joined, "the-anything") || !strings.Contains(joined, "no module declares") {
t.Fatalf("the refusal does not list the seats: %v", err) t.Fatalf("the refusal does not say the seat is nobody's: %v", problems)
}
}
// And the same claim is fine once something declares that seat, which is the case the parser
// could not have distinguished.
func TestAClaimOnASeatAnotherModuleDeclaresIsAccepted(t *testing.T) {
claimant, err := ParseManifest(claimed(`[{"name":"the-anything","scope":"node"}]`))
if err != nil {
t.Fatal(err)
}
declarer := Manifest{Module: "someone", Seats: []SeatDeclaration{
{Name: "the-anything", Scope: ScopeNode, Accepts: []string{"work"}},
}}
if problems := CatalogueProblems(Shelf{claimant.Module: claimant, "someone": declarer}); len(problems) != 0 {
t.Fatalf("a claim on a declared seat was refused: %v", problems)
} }
} }