Files
mesh-controller/internal/broker/raise_live_test.go
jschoubben cb77f35a27 A module may watch a role's events, and the catch-up turns out to be unnecessary
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.
2026-09-27 17:22:30 +02:00

197 lines
8.0 KiB
Go

package broker
import (
"os"
"testing"
"github.com/nats-io/nats.go"
)
// Raising the bus's objects against a real server.
//
// The pure tests above say what is asked for and in what order. Only a server can say whether it
// accepts them — and two of these are claims about the server's own behaviour that nothing else
// could answer: that asserting twice changes nothing, and that a consumer really is bound to the one
// subject its node is allowed to read.
//
// docker run -d --rm --name t -p 14227:4222 nats:2.10-alpine -js
// MESH_TEST_NATS=nats://127.0.0.1:14227 go test ./internal/broker/ -run TestRaising
func aLiveBus(t *testing.T) *JetStream {
t.Helper()
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
js, err := Dial(url)
if err != nil {
t.Fatal(err)
}
t.Cleanup(js.Close)
// **Nothing is deleted here, deliberately.** These objects are the mesh's own and every live
// test in every package shares one server: a test that deleted a stream to get a clean slate
// took it out from under whatever was running beside it, and the failure landed in the other
// test as "stream not found" — which reads as a bug in the code under test. Raise is idempotent
// by requirement, so asserting against whatever is already there is both safe and the realistic
// case.
return js
}
// Every object the mesh's own traffic needs, accepted by a real server, and asserting again changes
// nothing — which is the whole requirement, because this runs on every start.
func TestRaisingTheBusIsAcceptedAndIdempotent(t *testing.T) {
js := aLiveBus(t)
if err := Raise(js, []string{"anchor", "laptop"}); err != nil {
t.Fatalf("a real server refused the mesh's own objects: %v", err)
}
// Twice, with nothing in between. A start that failed the second time is a controller that
// cannot restart.
if err := Raise(js, []string{"anchor", "laptop"}); err != nil {
t.Fatalf("asserting the bus's objects a second time failed, so a restart would: %v", err)
}
// And again with a machine that was not there before, which is what enrolling one is.
if err := Raise(js, []string{"anchor", "laptop", "workstation"}); err != nil {
t.Fatalf("a machine joining an already-raised bus was refused: %v", err)
}
for _, s := range MeshStreams() {
if _, err := js.Context().StreamInfo(s.Name); err != nil {
t.Errorf("stream %s is not there: %v", s.Name, err)
}
}
for _, c := range MeshConsumers() {
if _, err := js.Context().ConsumerInfo(c.Stream, c.Name); err != nil {
t.Errorf("the controller's consumer on %s is not there: %v", c.Stream, err)
}
}
for _, node := range []string{"anchor", "laptop", "workstation"} {
info, err := js.Context().ConsumerInfo("NODES", node)
if err != nil {
t.Errorf("%s has no way to hear its declaration: %v", node, err)
continue
}
// **Its own subject and no other node's.** A consumer filtered on anything wider is a node
// reading another machine's declaration, and its own ack grant would not cover it either.
if info.Config.FilterSubject != "mesh.node."+node+".declare" {
t.Errorf("%s's consumer reads %q", node, info.Config.FilterSubject)
}
if info.Config.AckPolicy != nats.AckExplicitPolicy {
t.Errorf("%s's consumer acknowledges on delivery, so a declaration it died applying is "+
"never sent again", node)
}
}
}
// The store window needs unlimited redelivery on CONTROL: the bound belongs to the controller, and a
// server that dead-lettered first would discard the push the stream exists to protect.
func TestTheControlConsumerDoesNotDeadLetterBeforeTheControllerGivesUp(t *testing.T) {
js := aLiveBus(t)
if err := Raise(js, nil); err != nil {
t.Fatal(err)
}
info, err := js.Context().ConsumerInfo("CONTROL", ControllerName)
if err != nil {
t.Fatal(err)
}
if info.Config.MaxDeliver > 0 {
t.Fatalf("max-deliver is %d: a push held through a store restart would be dead-lettered "+
"before the controller finished deciding about it", info.Config.MaxDeliver)
}
}
// A role's work queue exists before anybody holds it, against a real server.
//
// **The queue before the holder is the point** (novox/hq ADR 0121): work queues until somebody arrives
// to do it, so assigning a build machine a week after something started asking for builds flushes the
// backlog instead of having lost it. A stream created at assignment would make "the holder is not here
// yet" mean "your requests are gone".
func TestRaisingAMeshRolesWorkQueue(t *testing.T) {
js := aLiveBus(t)
seats := []DeclaredSeat{{Name: "mesh-build-machine", Accepts: []string{"build"},
Emits: []string{"built"}}}
t.Cleanup(func() { _ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") })
if err := RaiseSeats(js, seats, nil); err != nil {
t.Fatalf("a real server refused a role's work queue: %v", err)
}
info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE")
if err != nil {
t.Fatalf("the role has no work queue: %v", err)
}
if info.Config.Retention != nats.WorkQueuePolicy {
t.Errorf("the queue retains as %v: work a holder took must leave it, or the next holder does "+
"it again", info.Config.Retention)
}
if len(info.Config.Subjects) != 1 || info.Config.Subjects[0] != "mesh.seat.mesh-build-machine.accept.>" {
t.Errorf("it carries %v rather than the role's own inbound subjects", info.Config.Subjects)
}
// Nobody holds it, so there is no worker — and asserting again changes nothing, because this runs
// on every start.
if err := RaiseSeats(js, seats, nil); err != nil {
t.Fatalf("asserting a role's queue a second time failed, so a restart would: %v", err)
}
// And once somebody holds it, the worker appears on that same queue.
if err := RaiseSeats(js, seats, map[string]Holder{
"mesh-build-machine": {Node: "anchor", Module: "builder"},
}); err != nil {
t.Fatal(err)
}
if _, err := js.Context().ConsumerInfo("SEAT_MESH_BUILD_MACHINE",
"SEAT_MESH_BUILD_MACHINE_worker"); err != nil {
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)
}
}