From cb77f35a279029a2866ccb0c7c9b9e99b1111703 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 27 Sep 2026 17:22:30 +0200 Subject: [PATCH] A module may watch a role's events, and the catch-up turns out to be unnecessary MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Moving the build outcome onto its role broke the one module that consumes it, and my own agreement check passed anyway. The catalogue's subscription derived `mesh.mod.mesh-build-machine.event.built` — a module namespace for a role's event, which no such module owns — so it started, connected, and its graph stayed empty. The check compared names, and the names agreed: the build machine does emit `built`. Only the subjects disagreed, and a subscription that matches nothing is silence. A consumed name is a module's event unless it names a role, and this package cannot tell by looking — so whoever resolved the declaration says which, the way it already does for a seat held or used. A module that watches a role gets the role's event subject and a consumer filtered on it; watching grants subscribe and nothing else, because hearing what a role announced is not taking part in it. The check now compares the two halves that actually have to match — the subject a consumer subscribes against the subject an emitter publishes — with a case pinning that it catches this exact confusion. Comparing names was checking the easy half. **And that answered the open question about catch-up: there is nothing to build.** The mechanism exists because a queue on the old bus receives only what is published after it is bound, so everything built before the catalogue existed was announced to nobody. A stream is a log and a consumer is a position in it: a consumer created afterwards starts at the beginning, so the builds are simply there. Asked of a real server, since the whole decision rested on it — three builds published with nothing listening, then a consumer created, and all three waiting for it. --- internal/broker/agreement_catalogue_test.go | 112 ++++++++++++++++++++ internal/broker/derived.go | 11 +- internal/broker/nats.go | 17 +++ internal/broker/raise_live_test.go | 51 +++++++++ internal/broker/users.go | 4 +- internal/inventory/busrecords.go | 20 +++- 6 files changed, 211 insertions(+), 4 deletions(-) diff --git a/internal/broker/agreement_catalogue_test.go b/internal/broker/agreement_catalogue_test.go index 7b9e0ec..9a61b3d 100644 --- a/internal/broker/agreement_catalogue_test.go +++ b/internal/broker/agreement_catalogue_test.go @@ -4,6 +4,7 @@ import ( "encoding/json" "os" "path/filepath" + "sort" "strings" "testing" @@ -122,3 +123,114 @@ func theCataloguesEvents(t *testing.T) ([]AnEmitter, []AConsumer, []DeclaredSeat } return emitters, consumers, seats } + +// **Do the derived subjects meet, not just the names?** +// +// The check above compares what a consumer asks for against what an emitter says it emits, by name. It +// passed while the catalogue's subscription pointed at `mesh.mod.mesh-build-machine.event.built` — a +// module namespace for a role's event, which no emitter owns. The names agreed; the subjects did not, +// and the graph stayed empty. +// +// So this compares the thing that actually has to match: the subject a consumer subscribes against the +// subject an emitter publishes. It is the last place the two halves can be held together, because +// after this the server is the only thing that knows and it says nothing — a subscription that matches +// nothing is silence. +func TestTheCataloguesDerivedSubjectsMeet(t *testing.T) { + emitters, consumers, seats := theCataloguesEvents(t) + + // Every subject something publishes: a module's own events, and the events of every role. + published := map[string]bool{} + for _, e := range emitters { + for _, name := range e.Emits { + published["mesh.mod."+e.Module+".event."+name] = true + } + } + for _, s := range seats { + for _, name := range s.Emits { + published["mesh.seat."+s.Name+".event."+name] = true + } + } + + byName := map[string]DeclaredSeat{} + for _, s := range seats { + byName[s.Name] = s + } + + var lonely []string + for _, c := range consumers { + principal := Principal{Kind: KindModule, Node: "one", Module: c.Module, PasswordHash: "x"} + for _, want := range c.Consumes { + emitter, event, named := strings.Cut(want, ".") + if named { + if s, isASeat := byName[emitter]; isASeat { + principal.Watches = append(principal.Watches, + Seat{Name: s.Name, Emits: []string{event}}) + continue + } + } + principal.Consumes = append(principal.Consumes, want) + } + perms, err := PermissionsFor(principal) + if err != nil { + t.Fatalf("%s: %v", c.Module, err) + } + for _, subject := range perms.Subscribe { + if !strings.Contains(subject, ".event.") { + continue + } + if reaches(subject, published) { + continue + } + // A wildcard over emitters reaches whatever arrives later, and an emitter that is not + // installed is ordinary — both are already excused by the check above, so only a subject + // that can never match anything gets here. + if strings.Contains(subject, "*") || strings.Contains(subject, ">") { + continue + } + lonely = append(lonely, c.Module+" subscribes "+subject+", which nothing publishes") + } + } + if len(lonely) > 0 { + sort.Strings(lonely) + t.Fatalf("%d subscription(s) derive to a subject no emitter owns:\n %s", + len(lonely), strings.Join(lonely, "\n ")) + } +} + +// And it catches the thing it exists for: a role's event read as a module's. +func TestTheDerivedSubjectCheckCatchesARolesEventReadAsAModules(t *testing.T) { + published := map[string]bool{"mesh.seat.mesh-build-machine.event.built": true} + // What the derivation produced before a consumed seat name was resolved as one. + if reaches("mesh.mod.mesh-build-machine.event.built", published) { + t.Fatal("a module namespace was treated as reaching a role's event, which is the bug") + } + // And the corrected one does reach it. + if !reaches("mesh.seat.mesh-build-machine.event.built", published) { + t.Fatal("the role's own subject does not reach the role's event") + } +} + +// reaches says whether a subscribed subject admits any published one. +func reaches(subject string, published map[string]bool) bool { + for p := range published { + if admitsSubject(strings.Split(subject, "."), strings.Split(p, ".")) { + return true + } + } + return false +} + +func admitsSubject(pattern, subject []string) bool { + for i, token := range pattern { + if token == ">" { + return i < len(subject) + } + if i >= len(subject) { + return false + } + if token != "*" && token != subject[i] { + return false + } + } + return len(pattern) == len(subject) +} diff --git a/internal/broker/derived.go b/internal/broker/derived.go index b4b0cfa..f5f514a 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -3,6 +3,7 @@ package broker import ( "fmt" "sort" + "strings" ) // Streams and consumers derived from what modules declare. @@ -103,16 +104,22 @@ type DeclaredSeat struct { // from its name: a module with a consumer per event would need an ack permission per consumer, // and the permission list would stop being derivable from the declaration. func ConsumerFor(p Principal) (Consumer, bool) { - if p.Kind != KindModule || len(p.Consumes) == 0 { + // A module that reacts to anything — a module's events or a role's (novox/hq ADR 0121). Watching + // a role was missing here, so the one module that does it got no consumer at all: it started, + // connected, and its graph stayed empty with nothing anywhere reporting why. + if p.Kind != KindModule || (len(p.Consumes) == 0 && len(p.Watches) == 0) { return Consumer{}, false } perms, err := PermissionsFor(p) if err != nil { return Consumer{}, false } + // Events, wherever they live: a module's own namespace, and the namespace of any role it watches + // (novox/hq ADR 0121). Tool subjects and inboxes are subscribed directly and are not a consumer's + // business, which is why this is a filter and not the whole list. var filters []string for _, s := range perms.Subscribe { - if len(s) > 9 && s[:9] == "mesh.mod." { + if strings.Contains(s, ".event.") { filters = append(filters, s) } } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 295f2e4..3c18343 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -60,6 +60,15 @@ type Principal struct { Holds []Seat Uses []Seat + // Watches are seats whose events this principal consumes. Separate from Consumes because a + // role's event lives under the seat's namespace and not a module's, and this package cannot tell + // a seat's name from a module's by looking at it — whoever resolved the declaration can, and + // does (novox/hq ADR 0121). + // + // **Found by a consumer reading nothing.** The catalogue consumes the build machine's outcome; + // with that name read as a module's, its subscription pointed at `mesh.mod.mesh-build-machine.…`, + // a namespace no such module owns. Every service started and the graph stayed empty. + Watches []Seat // Invokes are the tools a person may call, as `.`; a single `*` is every tool, // for an administrator. Only meaningful for KindPerson. @@ -260,6 +269,14 @@ func PermissionsFor(p Principal) (Permissions, error) { sub = append(sub, subject) } + // 2b. Events of a role it watches, under the seat's own namespace. Subscribe only: watching a + // role is hearing what it announced, not taking part in it. + for _, w := range p.Watches { + for _, e := range w.Emits { + sub = append(sub, seatSubject(w, "event", e)) + } + } + // 3. Seats it holds: full participation. for _, s := range p.Holds { for _, a := range s.Accepts { diff --git a/internal/broker/raise_live_test.go b/internal/broker/raise_live_test.go index 6281679..202babc 100644 --- a/internal/broker/raise_live_test.go +++ b/internal/broker/raise_live_test.go @@ -143,3 +143,54 @@ func TestRaisingAMeshRolesWorkQueue(t *testing.T) { t.Fatalf("the holder got no worker on the role's queue: %v", err) } } + +// **A consumer created after the fact still sees what came before it**, which is why the mesh needs no +// catch-up at all on this bus (novox/hq 04-ISSUES/050). +// +// On the bus the mesh runs on today a queue receives only what is published after it is bound, so +// everything built before the catalogue existed was announced to nobody — and on a fresh mesh that is +// always the foundation, because those are the things the catalogue needed in order to exist. A whole +// mechanism was built for it: the catalogue asks, the controller re-publishes. +// +// A stream is a log and a consumer is a position in it. A consumer created later starts at the +// beginning by default, so the builds are simply there. Asked of a real server rather than assumed, +// because the whole decision about whether to keep that mechanism rests on it. +func TestAConsumerCreatedAfterwardsStillSeesWhatCameBefore(t *testing.T) { + js := aLiveBus(t) + if err := AssertMeshStreams(js); err != nil { + t.Fatal(err) + } + if err := js.Context().PurgeStream("EVENTS"); err != nil { + t.Fatal(err) + } + + // Genesis: things are built before anything is listening. + built := []string{"base", "store", "mesh-catalog"} + for _, m := range built { + if _, err := js.Context().Publish("mesh.seat.mesh-build-machine.event.built", + []byte(`{"module":"`+m+`"}`)); err != nil { + t.Fatal(err) + } + } + + // Now the catalogue is installed and the controller creates its consumer. + c, ok := ConsumerFor(Principal{Kind: KindModule, Node: "one", Module: "mesh-catalog", + Watches: []Seat{{Name: "mesh-build-machine", Emits: []string{"built"}}}, PasswordHash: "x"}) + if !ok { + t.Fatal("a module that watches a role got no consumer") + } + t.Cleanup(func() { _ = js.Context().DeleteConsumer(c.Stream, c.Name) }) + if err := js.EnsureConsumer(c); err != nil { + t.Fatal(err) + } + + info, err := js.Context().ConsumerInfo(c.Stream, c.Name) + if err != nil { + t.Fatal(err) + } + if info.NumPending != uint64(len(built)) { + t.Fatalf("a consumer created after %d builds has %d waiting for it — if this is 0 the mesh "+ + "does need a catch-up after all, and the reasoning for deleting it is wrong", + len(built), info.NumPending) + } +} diff --git a/internal/broker/users.go b/internal/broker/users.go index f6fc289..17d0ce4 100644 --- a/internal/broker/users.go +++ b/internal/broker/users.go @@ -29,6 +29,8 @@ type Declared struct { Holds []Seat // Uses are the seats this module sends to. Uses []Seat + // Watches are the seats whose events it consumes. + Watches []Seat } // Records is what composing a user list needs to know about the mesh, and nothing more. @@ -59,7 +61,7 @@ func Users(r Records) ([]Principal, error) { out = append(out, Principal{ Kind: KindModule, Node: node, Module: d.Module, Emits: d.Emits, Consumes: d.Consumes, Serves: d.Serves, - Holds: d.Holds, Uses: d.Uses, + Holds: d.Holds, Uses: d.Uses, Watches: d.Watches, }) } } diff --git a/internal/inventory/busrecords.go b/internal/inventory/busrecords.go index 4b29bef..a30e624 100644 --- a/internal/inventory/busrecords.go +++ b/internal/inventory/busrecords.go @@ -3,6 +3,7 @@ package inventory import ( "context" "fmt" + "strings" "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/catalogue" @@ -93,10 +94,27 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) { // declaredFor is one module's manifest as the composer needs it: what it says about itself, and the // protocol of every seat it holds or uses. func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaration) broker.Declared { + // A consumed name is a module's event unless it names a seat, and only somebody holding the seat + // set can tell (novox/hq ADR 0121). Split here, because the composer cannot look at a name and + // know — and a role's event read as a module's is a subscription to a namespace nobody owns. + var fromModules []string + var watches []broker.Seat + for _, c := range m.Consumes { + emitter, event, named := strings.Cut(c, ".") + if named { + if s, isASeat := seats[emitter]; isASeat { + watches = append(watches, broker.Seat{Name: s.Name, Emits: []string{event}}) + continue + } + } + fromModules = append(fromModules, c) + } + d := broker.Declared{ Module: m.Module, Emits: m.Emits, - Consumes: m.Consumes, + Consumes: fromModules, + Watches: watches, // The tools it answers, which is `tools` and not `serves`: the manifest's `serves` is the // facts a consumer needs to reach a provision, a different meaning under a similar word. Serves: m.Tools,