From 4de10e32e370830e9e4e30b8b7243c2088e80e8a Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 27 Sep 2026 02:59:05 +0200 Subject: [PATCH] The bus's objects are raised on every start, and one switch says which bus MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two of 1.7's three remaining pieces. **Raised on every start, not created once at genesis.** A stream somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was replaced all have records and no objects — and a node whose consumer is missing hears nothing while everything else about it looks correct. The order is not a preference: a consumer on a stream that does not exist is refused *naming the stream*, so somebody reading that refusal goes looking for a deletion instead of a reversed pair of lines. Pinned by a test, along with the one thing about seats that reads like an omission and is not — a seat's work queue is asserted whether or not anybody holds it, because work queues until a holder appears, so installing the module a week later flushes the backlog instead of having lost it. Against a real server: every object accepted, asserting twice changes nothing (a start that failed the second time is a controller that cannot restart), a machine joining an already-raised bus is accepted, each node's consumer is bound to its own declaration subject and no other's, and CONTROL does not dead-letter — because the store window's bound belongs to the controller and a server that gave up first would discard the push the stream exists to protect. **Which bus this mesh is on is one fact, read in one place.** Every seam the change went behind ships both implementations; this is what the rollout flips. Being told about both is refused at start rather than warned about: a mesh half on each is one where a declaration goes out on one bus and the report comes back on the other, and every component logs success while it happens — ADR 0074's failure arriving through configuration instead of through code. The refusal names both variables and says which to unset, because whoever reads it has to choose and the wrong choice is a rollout half done. --- cmd/mesh-controller/push.go | 52 +++++++++++++++ internal/broker/onnats.go | 58 ++++++++++++++++ internal/broker/onnats_test.go | 40 +++++++++++ internal/broker/raise.go | 73 ++++++++++++++++++++ internal/broker/raise_live_test.go | 104 +++++++++++++++++++++++++++++ internal/broker/streams_test.go | 90 +++++++++++++++++++++++++ 6 files changed, 417 insertions(+) create mode 100644 internal/broker/onnats.go create mode 100644 internal/broker/onnats_test.go create mode 100644 internal/broker/raise.go create mode 100644 internal/broker/raise_live_test.go diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 2ee6875..d02e40c 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -77,12 +77,33 @@ func serve(ctx context.Context) error { "reconnect. Set %s and %s.\n", broker.AddressVar, broker.CertificateVar) } + // **Which bus this mesh is on, read once** (novox/hq ADR 0116 step 5). Both clients ship; both + // being live is refused, because a mesh half on each is one where a declaration goes out on one + // and the report comes back on the other, and every component logs success while it happens. + busAddress, onNATS, err := broker.OnNATS() + if err != nil { + return err + } + if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil { + return err + } + work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known} server, err := link.Connect(work, work) if err != nil { return err } defer server.Close() + + // The bus's own objects, asserted on every start. **Not created once at genesis**: a stream + // somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was + // replaced all have records and no objects — and a node whose consumer is missing hears nothing + // while everything else about it looks correct. + if onNATS { + if err := raiseTheBus(ctx, inv, busAddress); err != nil { + return err + } + } // And build results nobody was waiting for. A build triggered any other way than `build` // would otherwise be reported into the void, which is the same as not reporting it. server.Records(builds{inv}) @@ -664,3 +685,34 @@ func wouldSend(ctx context.Context, open *stores, } return out, nil } + +// raiseTheBus asserts the streams and consumers the mesh's own traffic needs. +// +// **Every start, and it says what it did.** The objects are the mesh's, created by nothing else — +// the controller is their only writer (design 25 §3) — so a mesh that came up without them is one +// where nodes connect, authenticate, and hear nothing. Said rather than silent for the reason the +// first line of `serve` is said: a log that is quiet on success and loud on failure reads as broken +// when it is working. +func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string) error { + js, err := broker.Dial(address) + if err != nil { + return fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w", + address, err) + } + defer js.Close() + + nodes, err := inv.Nodes(ctx) + if err != nil { + return err + } + names := make([]string, 0, len(nodes)) + for _, n := range nodes { + names = append(names, n.Name) + } + if err := broker.Raise(js, names); err != nil { + return err + } + fmt.Printf("the bus at %s has its streams, and %d machine(s) can hear a declaration\n", + address, len(names)) + return nil +} diff --git a/internal/broker/onnats.go b/internal/broker/onnats.go new file mode 100644 index 0000000..f6efa54 --- /dev/null +++ b/internal/broker/onnats.go @@ -0,0 +1,58 @@ +package broker + +import ( + "fmt" + "strings" + + "github.com/novox/mesh-controller/internal/envfile" +) + +// Whether this mesh's own traffic is on the bus being built. +// +// **One switch, read in one place** (novox/hq ADR 0116 step 5). Every seam the bus change went +// behind ships both implementations, and until the rollout every one of them chooses the bus the +// mesh runs on today. This is what the rollout flips, and it is deliberately a single fact rather +// than a fact per component: a controller whose outbound is on one bus and whose inbound is on the +// other is a mesh that hears nothing, and no test of either half would catch it. + +// NATSVar is where the controller finds the bus being built. Unset is the ordinary case and means +// the mesh runs on the bus it has always run on. +const NATSVar = "MESH_BUS_NATS" + +// OnNATS is the address of the bus being built, and whether the mesh is on it. +// +// Read from the node's own settings rather than baked in, for the reason the broker's address is +// (novox/hq 04-ISSUES/102): an address recorded once does not follow a node's ports. +func OnNATS() (address string, on bool, err error) { + address, err = envfile.Placed(NATSVar) + if err != nil { + return "", false, err + } + address = strings.TrimSpace(address) + if address == "" { + return "", false, nil + } + return address, true, nil +} + +// MustBeOneBus refuses a configuration that names both buses for the mesh's own traffic. +// +// **Both clients ship and that is the point; both being live is not.** The rollout moves every node +// at once (ADR 0116 step 5): a mesh half on each is one where a declaration goes out on one bus and +// the report comes back on the other, and nothing anywhere says so — every component would log +// success. Refused at start, where it can be said in one sentence. +func MustBeOneBus(amqp, nats string) error { + if strings.TrimSpace(amqp) != "" && strings.TrimSpace(nats) != "" { + return fmt.Errorf( + "this control plane is told about both buses (%s and %s) and can only be on one. A mesh "+ + "half on each is one where a declaration goes out on one and the report comes back "+ + "on the other, and every component reports success while it happens. The rollout "+ + "moves every node at once: unset %s to stay, or unset %s to move", + AMQPVarName, NATSVar, NATSVar, AMQPVarName) + } + return nil +} + +// AMQPVarName is the variable naming the bus the mesh runs on today. Named here rather than +// imported from the link package, for the one direction of dependency. +const AMQPVarName = "MESH_BROKER_AMQP" diff --git a/internal/broker/onnats_test.go b/internal/broker/onnats_test.go new file mode 100644 index 0000000..19e78e9 --- /dev/null +++ b/internal/broker/onnats_test.go @@ -0,0 +1,40 @@ +package broker + +import ( + "strings" + "testing" +) + +// Which bus the mesh is on is one fact, and being told about both is refused. +// +// **Not a warning.** A mesh half on each bus is one where a declaration goes out on one and the +// report comes back on the other, and every component reports success while it happens — which is +// the exact failure ADR 0074 exists to catch, arriving through configuration instead of through code. +func TestBeingToldAboutBothBusesIsRefused(t *testing.T) { + err := MustBeOneBus("amqps://broker:5671/", "nats://bus:4222") + if err == nil { + t.Fatal("a control plane told about both buses was allowed to start") + } + // The remedy is in the words, because whoever reads this has to choose one and the wrong choice + // is a rollout half done. + for _, want := range []string{AMQPVarName, NATSVar, "unset"} { + if !strings.Contains(err.Error(), want) { + t.Errorf("the refusal does not mention %s: %v", want, err) + } + } +} + +// One bus, or none, is ordinary. None is a control plane that publishes nothing and holds records, +// which several of its own commands are. +func TestOneBusOrNeitherIsAllowed(t *testing.T) { + for _, c := range []struct{ what, amqp, nats string }{ + {"the bus the mesh runs on today", "amqps://broker:5671/", ""}, + {"the bus being built", "", "nats://bus:4222"}, + {"neither", "", ""}, + {"neither, with whitespace for an address", " ", "\t"}, + } { + if err := MustBeOneBus(c.amqp, c.nats); err != nil { + t.Errorf("%s was refused: %v", c.what, err) + } + } +} diff --git a/internal/broker/raise.go b/internal/broker/raise.go new file mode 100644 index 0000000..6f3d15a --- /dev/null +++ b/internal/broker/raise.go @@ -0,0 +1,73 @@ +package broker + +import "fmt" + +// Bringing the bus's own objects into being, in the one order that works. +// +// **Asserted on every start rather than created once at genesis.** A stream somebody deleted, a mesh +// raised from a restored backup, or a bus whose data directory was replaced all have records and no +// objects — and a node whose consumer is missing hears nothing while everything else about it looks +// correct. Idempotence is the whole requirement, and the parts are already idempotent; this is the +// order they have to be asked in. + +// Raiser is everything asserting the bus's objects needs of a connection to it. +type Raiser interface { + Asserter + Ensurer +} + +// Raise asserts the mesh's streams, the controller's own consumers, and one consumer per node. +// +// **The order is not a preference.** A consumer on a stream that does not exist is refused, and the +// refusal names the stream rather than the order — so somebody reading it goes looking for a deleted +// stream instead of a reversed pair of lines. Nodes last, because the one a node reads lives on a +// stream the mesh's own set defines. +func Raise(r Raiser, nodes []string) error { + if err := AssertMeshStreams(r); err != nil { + return err + } + if err := AssertMeshConsumers(r); err != nil { + return err + } + if err := AssertNodeConsumers(r, nodes); err != nil { + return err + } + return nil +} + +// RaiseSeats asserts one work queue per declared seat, and the worker of whoever holds it. +// +// Separate from Raise because it is answered by a different question: the mesh's own objects exist +// because the mesh does, and a seat's exist because a module declaring one was registered. Kept +// beside it so the order is visible — a holder's worker needs the seat's stream, and a seat's stream +// needs nothing. +func RaiseSeats(r Raiser, seats []DeclaredSeat, holders map[string]Holder) error { + for _, s := range SeatStreams(seats) { + if err := r.EnsureStream(s); err != nil { + return fmt.Errorf("asserting the work queue for %s: %w", s.Name, err) + } + } + for _, s := range seats { + h, held := holders[s.Name] + if !held { + // **The stream exists and the consumer does not, on purpose.** Work queues until a + // holder appears, so installing the module a week after something started sending to it + // flushes the backlog instead of having lost it. + continue + } + c, needed := HolderConsumerFor(h.Node, h.Module, s) + if !needed { + continue + } + if err := r.EnsureConsumer(c); err != nil { + return fmt.Errorf("asserting how %s on %s works %s: %w", h.Module, h.Node, s.Name, err) + } + } + return nil +} + +// Holder is which module on which machine holds a seat. +type Holder struct { + Node string + Module string +} diff --git a/internal/broker/raise_live_test.go b/internal/broker/raise_live_test.go new file mode 100644 index 0000000..bf468a0 --- /dev/null +++ b/internal/broker/raise_live_test.go @@ -0,0 +1,104 @@ +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) + // Deleted before, so what this test asserts is what it finds — and after, so the next test does + // not inherit it. Deleting a stream takes its consumers with it, which is why this is enough. + clear := func() { + for _, s := range MeshStreams() { + _ = js.Context().DeleteStream(s.Name) + } + } + clear() + t.Cleanup(clear) + 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) + } +} diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go index bd41654..dc70bb8 100644 --- a/internal/broker/streams_test.go +++ b/internal/broker/streams_test.go @@ -144,3 +144,93 @@ func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) { } } } + +// The order the bus's objects are asserted in, because getting it wrong is a refusal that names the +// wrong thing: a consumer on a stream that does not exist is refused naming the *stream*, so +// somebody reading it goes looking for a deletion instead of a reversed pair of lines. +func TestTheBusesObjectsAreAssertedStreamsBeforeConsumers(t *testing.T) { + r := &recording{} + if err := Raise(r, []string{"anchor", "laptop"}); err != nil { + t.Fatal(err) + } + + // Every stream before every consumer. + firstConsumer := -1 + for i, step := range r.steps { + if strings.HasPrefix(step, "consumer ") && firstConsumer < 0 { + firstConsumer = i + } + if strings.HasPrefix(step, "stream ") && firstConsumer >= 0 { + t.Fatalf("a stream was asserted after a consumer: %v", r.steps) + } + } + if firstConsumer < 0 { + t.Fatalf("no consumer was asserted: %v", r.steps) + } + + // And every node got one, named after it — without which that node hears nothing while + // everything else about it looks correct. + for _, node := range []string{"anchor", "laptop"} { + if !containsStep(r.steps, "consumer NODES/"+node) { + t.Errorf("%s was given no way to hear its declaration: %v", node, r.steps) + } + } + // And the controller its own, on both streams it reads. + for _, want := range []string{"consumer CONTROL/controller", "consumer EVENTS/controller"} { + if !containsStep(r.steps, want) { + t.Errorf("the controller is missing %s: %v", want, r.steps) + } + } +} + +// A seat's work queue is asserted whether or not anybody holds it; the holder's worker only when +// somebody does. **The stream without the consumer is the point**: work queues until a holder +// appears, so installing the module later flushes the backlog instead of having lost it. +func TestASeatsQueueExistsBeforeItsHolderDoes(t *testing.T) { + seats := []DeclaredSeat{{Name: "telegram-sender", Accepts: []string{"send"}}} + + unheld := &recording{} + if err := RaiseSeats(unheld, seats, nil); err != nil { + t.Fatal(err) + } + if !containsStep(unheld.steps, "stream SEAT_TELEGRAM_SENDER") { + t.Fatalf("a declared seat got no work queue: %v", unheld.steps) + } + for _, step := range unheld.steps { + if strings.HasPrefix(step, "consumer ") { + t.Fatalf("a seat nobody holds got a worker: %v", unheld.steps) + } + } + + held := &recording{} + if err := RaiseSeats(held, seats, map[string]Holder{ + "telegram-sender": {Node: "anchor", Module: "telegram"}, + }); err != nil { + t.Fatal(err) + } + if !containsStep(held.steps, "consumer SEAT_TELEGRAM_SENDER/SEAT_TELEGRAM_SENDER_worker") { + t.Fatalf("the seat's holder got no worker: %v", held.steps) + } +} + +// recording is a connection to the bus that writes down what it was asked for. +type recording struct{ steps []string } + +func (r *recording) EnsureStream(s Stream) error { + r.steps = append(r.steps, "stream "+s.Name) + return nil +} + +func (r *recording) EnsureConsumer(c Consumer) error { + r.steps = append(r.steps, "consumer "+c.Stream+"/"+c.Name) + return nil +} + +func containsStep(steps []string, want string) bool { + for _, s := range steps { + if s == want { + return true + } + } + return false +}