The bus on NATS: both transports behind seams, and the rollout switch #87
@@ -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
|
||||
}
|
||||
|
||||
@@ -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"
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user