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 +}