From aecac5bda20a704610aab77af70bb912b765eb25 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 28 Sep 2026 03:36:16 +0200 Subject: [PATCH] One bus: the AMQP transport is gone from the controller The mesh runs on the seat's bus alone (novox/hq ADR 0131, design 28 task 5.5). The old transport's consume loop, build request, tool ask, management API and account scoping are deleted, and the bus switch with them; the controller connects to the broker seat and to nothing else. The store-window tests keep their assertions on a bus-less fake, and the tests that only made sense for the old transport's in-memory holding go with it. --- cmd/mesh-builder/main.go | 96 +------ cmd/mesh-builder/where_test.go | 6 +- cmd/mesh-controller/ask.go | 2 +- cmd/mesh-controller/build.go | 120 +++----- cmd/mesh-controller/modules.go | 72 +---- cmd/mesh-controller/nodes.go | 10 - cmd/mesh-controller/push.go | 24 +- cmd/mesh-controller/upgrades.go | 18 +- go.mod | 3 +- go.sum | 4 - internal/broker/agreement_test.go | 67 ----- internal/broker/broker_test.go | 13 - internal/broker/management.go | 366 ------------------------- internal/broker/module_account_test.go | 96 ------- internal/broker/onnats.go | 30 +- internal/broker/onnats_test.go | 32 --- internal/link/ask.go | 82 ++---- internal/link/build.go | 94 ------- internal/link/builds_current.go | 202 -------------- internal/link/bus.go | 46 ---- internal/link/bus_nats.go | 52 +--- internal/link/enrolment.go | 21 +- internal/link/receive.go | 13 - internal/link/receive_current.go | 321 ---------------------- internal/link/receive_current_test.go | 71 ----- internal/link/receive_fake_test.go | 94 +++++++ internal/link/receive_nats.go | 12 +- internal/link/receive_nats_test.go | 4 + internal/link/report_retry_test.go | 18 -- internal/link/serve.go | 88 +----- internal/link/store_window_test.go | 52 +--- 31 files changed, 214 insertions(+), 1915 deletions(-) delete mode 100644 internal/broker/agreement_test.go delete mode 100644 internal/broker/management.go delete mode 100644 internal/broker/module_account_test.go delete mode 100644 internal/link/builds_current.go delete mode 100644 internal/link/receive_current.go delete mode 100644 internal/link/receive_current_test.go create mode 100644 internal/link/receive_fake_test.go diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index 6aded14..5f9ddb8 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -17,12 +17,7 @@ package main import ( "context" - "crypto/sha256" - "crypto/tls" - "crypto/x509" - "encoding/hex" "encoding/json" - "errors" "fmt" "net/url" "os" @@ -30,8 +25,6 @@ import ( "strings" "syscall" - amqp "github.com/rabbitmq/amqp091-go" - "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/builder" "github.com/novox/mesh-controller/internal/link" @@ -52,7 +45,6 @@ const usage = `mesh-builder — builds modules for the mesh It consumes build requests and answers with what it made. Nothing is listened on and nothing is dialled except the broker. - MESH_BROKER_AMQP where the broker is, with this builder's own credential MESH_BROKER_FILE a file the mesh sealed to this machine holding the same MESH_REGISTRY host:port to publish artifacts to, when the mesh has not said MESH_BINDING a file the mesh wrote saying where the artifact store is @@ -135,42 +127,17 @@ func run() error { // machine told about both would take work from one and answer on the other, and every log line would // say it was fine. func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) { - // **The credential decides, before any variable does.** A machine moved to the new bus was - // handed a credential for it and nothing else changed in its environment; that credential - // names the bus by scheme, so it is enough to know which bus to take work from. - if credential.onTheNewBus() { - js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint) - if err != nil { - return nil, err - } - return link.MachineOverNATS(js, on), nil + // **The credential names the bus, and there is one** (novox/hq ADR 0131, design 28 task 5.5). + // A credential for the mesh's bus carries user, password and fingerprint beside the address, + // and that is enough to dial it, pinned. + if !credential.onTheNewBus() { + return nil, fmt.Errorf("the credential at hand names %q, which is not the mesh's bus", credential.URL) } - address, onNATS, err := broker.OnNATS() + js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint) if err != nil { return nil, err } - if err := broker.MustBeOneBus(credential.URL, address); err != nil { - return nil, err - } - if onNATS { - js, err := broker.Dial(address) - if err != nil { - return nil, fmt.Errorf("cannot reach the bus at %s: %w", address, err) - } - return link.MachineOverNATS(js, on), nil - } - - conn, err := dial(credential) - if err != nil { - // Not quoted back: the URL carries this builder's broker password. - return nil, fmt.Errorf("cannot reach the broker: %w", err) - } - channel, err := conn.Channel() - if err != nil { - conn.Close() - return nil, err - } - return link.MachineOverCurrent(conn, channel, on), nil + return link.MachineOverNATS(js, on), nil } // answer does one build and says what happened, whichever way it went. @@ -452,13 +419,8 @@ func brokerFrom() (Credential, error) { // broker is then verified against whatever this machine already trusts. return Credential{URL: said}, nil } - url := strings.TrimSpace(os.Getenv("MESH_BROKER_AMQP")) - if url == "" { - return Credential{}, fmt.Errorf( - "neither MESH_BROKER_FILE nor MESH_BROKER_AMQP: a builder with no broker has " + - "nothing to build") - } - return Credential{URL: url}, nil + return Credential{}, fmt.Errorf( + "no MESH_BROKER_FILE: a build machine with no credential for the bus has nothing to build") } // Credential is what a build machine is given so it can reach the broker. @@ -491,43 +453,3 @@ func (c Credential) natsURL() string { } return "nats://" + c.User + ":" + c.Password + "@" + rest } - -// dial opens the connection, pinning the broker's certificate when there is one to pin. -func dial(held Credential) (*amqp.Connection, error) { - if held.Fingerprint == "" { - return amqp.Dial(held.URL) - } - return amqp.DialTLS(held.URL, pinning(held.Fingerprint)) -} - -// pinning is a TLS configuration that trusts exactly one certificate. -// -// InsecureSkipVerify with a VerifyPeerCertificate is **pinning, not skipping**: the standard chain -// check is replaced, not removed, and what replaces it is stricter — one certificate is accepted -// rather than every certificate a public authority would sign. -// -// Its own function so a test can drive it against a real handshake. A pin check that is only ever -// exercised through a broker is a pin check nothing tests. -func pinning(fingerprint string) *tls.Config { - return &tls.Config{ - InsecureSkipVerify: true, - VerifyPeerCertificate: func(raw [][]byte, _ [][]*x509.Certificate) error { - if len(raw) == 0 { - return errors.New("the broker presented no certificate") - } - // The leaf, and in the same spelling the mesh writes it — `sha256:` and 64 hex - // characters. Comparing a bare digest against a written fingerprint never matches, - // and the failure is indistinguishable from being pointed at the wrong broker. - sum := sha256.Sum256(raw[0]) - got := "sha256:" + hex.EncodeToString(sum[:]) - if got != fingerprint { - return fmt.Errorf( - "this is not the broker this builder was told about\n expected %s\n "+ - "got %s\nEither this mesh's broker was replaced, or this builder is "+ - "being pointed at something else. Retrying will not help", - fingerprint, got) - } - return nil - }, - } -} diff --git a/cmd/mesh-builder/where_test.go b/cmd/mesh-builder/where_test.go index 5571493..5a55bdb 100644 --- a/cmd/mesh-builder/where_test.go +++ b/cmd/mesh-builder/where_test.go @@ -15,6 +15,8 @@ import ( "strings" "testing" "time" + + "github.com/novox/mesh-controller/internal/broker" ) // Where a builder publishes. @@ -193,9 +195,9 @@ func TestThePinIsComparedInTheSpellingTheMeshWritesIt(t *testing.T) { } } -// handshakeWith runs the builder's own pin check against an address. +// handshakeWith runs the pin check the builder dials with against an address. func handshakeWith(address, pin string) error { - conn, err := tls.Dial("tcp", address, pinning(pin)) + conn, err := tls.Dial("tcp", address, broker.PinnedToFingerprint(pin)) if err != nil { return err } diff --git a/cmd/mesh-controller/ask.go b/cmd/mesh-controller/ask.go index 42423d5..9006b79 100644 --- a/cmd/mesh-controller/ask.go +++ b/cmd/mesh-controller/ask.go @@ -43,7 +43,7 @@ func askCommand(ctx context.Context, args []string) error { } defer server.Close() - answer, err := link.Ask(ctx, server.Channel(), module, tool, arguments, *wait) + answer, err := link.Ask(ctx, server.Bus(), module, tool, arguments, *wait) if err != nil { return err } diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index d1189a6..df62f6d 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -2,8 +2,6 @@ package main import ( "context" - "crypto/rand" - "encoding/base64" "encoding/json" "errors" "flag" @@ -50,7 +48,7 @@ func buildOn(ctx context.Context, base string, wait time.Duration) error { continue } for _, b := range e.Manifest.Build.On { - if b.Module == base { + if standsOnModule(b, base) { on = append(on, e) break } @@ -261,85 +259,45 @@ func builderCommand(ctx context.Context, args []string) error { } name := positionals[1] - management, err := broker.ManagementFromEnvironment() + // **The build machine's credential is a module's credential** (novox/hq ADR 0131, design 28 + // task 5.5): minted into the mesh's records and sealed to the machine as the builder module's + // broker secret, usable at the next push — the same act `module issue` performs, and the same + // account the composed user list carries. Nothing is created on a server; the bus reads the list. + open, err := openStores(ctx) if err != nil { return err } - - // The same shape of secret a token carries: enough entropy that guessing is not a strategy, - // and safe to put in a URL because that is where it goes. - raw := make([]byte, 32) - if _, err := rand.Read(raw); err != nil { + defer open.Close() + inv := open.inventory + shelf, err := inv.Catalogue(ctx) + if err != nil { return err } - password := base64.RawURLEncoding.EncodeToString(raw) - if err := management.CreateBuilderAccount(ctx, name, password); err != nil { + m, known := shelf[*module] + if !known { + return fmt.Errorf("%s is not in the catalogue; `module add` it first", *module) + } + fmt.Printf("build machine %s: ", name) + node := *forNode + if node == "" { + entries, err := inv.Catalogued(ctx) + if err != nil { + return err + } + for _, e := range entries { + if e.Manifest.Module == *module && len(e.On) > 0 { + node = e.On[0] + } + } + } + if node == "" { + return fmt.Errorf("%s is assigned nowhere; `assign %s` first, or say --node", *module, *module) + } + address, err := broker.BusAddress() + if err != nil { return err } - - fmt.Printf("broker account %s created, scoped to the %s queue and the %s exchange\n\n", - name, link.BuildQueue, link.Exchange) - - if *forNode != "" { - known, err := broker.FromEnvironment() - if err != nil { - return fmt.Errorf("cannot deliver a credential without knowing where the broker is: %w", err) - } - inv, err := openInventory(ctx) - if err != nil { - return err - } - defer inv.Close() - - brokerAddr, err := brokerReachableAt(ctx, inv, known, *forNode) - if err != nil { - return err - } - - // The URL and what verifies the broker, together. A mesh's broker presents a certificate - // of the mesh's own, which is in no public trust store — so a URL on its own reaches only - // a broker somebody else vouches for, and the connection fails at TLS with an error about - // an unknown authority rather than about a missing pin. - // - // **The same two facts a node's token carries** (novox/hq ADR 0004), delivered the same - // way: out of band relative to the broker, so what is trusted does not come from the thing - // being trusted. - held, err := json.Marshal(struct { - URL string `json:"url"` - Fingerprint string `json:"fingerprint,omitempty"` - }{ - URL: fmt.Sprintf("amqps://%s:%s@%s/", name, password, brokerAddr), - Fingerprint: known.Fingerprint, - }) - if err != nil { - return err - } - if err := inv.AcceptSecretForModule(ctx, *forNode, *module, "broker", string(held)); err != nil { - return err - } - // Not printed. It is sealed to that machine and the mesh cannot read it back, which is - // the whole point — printing it here would put the one copy that matters on a terminal. - fmt.Printf(" sealed to %s, for the %s module. It arrives with the next push.\n", - *forNode, *module) - fmt.Printf(" run `push %s` to send it\n", *forNode) - return nil - } - - // The whole line only when the address is known. A URL with a placeholder where the host - // should be is a URL somebody pastes and then debugs, and the placeholder is the last thing - // they look at. - if known, err := broker.FromEnvironment(); err == nil { - fmt.Printf(" MESH_BROKER_AMQP=amqps://%s:%s@%s/\n\n", name, password, known.Address) - } else { - fmt.Printf(" the password is %s\n\n", password) - fmt.Printf(" This control plane has no %s, so it cannot say where the broker is.\n"+ - " Put the password in MESH_BROKER_AMQP on the build machine.\n\n", - broker.AddressVar) - } - // Shown once, like a token, and for the same reason: what is stored is the broker's own hash - // of it, and a control plane that could show it back would be a control plane that holds it. - fmt.Println("This is the only time it is shown.") - return nil + return issueOnTheNewBus(ctx, inv, m, node, address) } // buildBehind builds every module the mesh holds older than its source has. @@ -634,16 +592,10 @@ func heldBy(ctx context.Context) map[string]string { // **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5). On the bus // the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus // being built it dials, because a build request is a one-shot and holds nothing else. -func askOver(server *link.Server) (link.Builders, error) { - address, onNATS, err := broker.OnNATS() +func askOver(_ *link.Server) (link.Builders, error) { + address, err := broker.BusAddress() if err != nil { return nil, err } - if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), address); err != nil { - return nil, err - } - if onNATS { - return link.BuildsOverNATS(address) - } - return link.BuildsOverCurrent(server.Channel()), nil + return link.BuildsOverNATS(address) } diff --git a/cmd/mesh-controller/modules.go b/cmd/mesh-controller/modules.go index 6cfae4d..1d64c80 100644 --- a/cmd/mesh-controller/modules.go +++ b/cmd/mesh-controller/modules.go @@ -2,8 +2,6 @@ package main import ( "context" - "crypto/rand" - "encoding/base64" "encoding/json" "errors" "flag" @@ -261,77 +259,11 @@ func moduleCommand(ctx context.Context, args []string) error { // made in entirely different ways: on the bus the mesh runs on today an account is a // management call, and on the bus being built it is a row the next composition writes into // the server's user list (novox/hq design 25 §4). - busAddress, onNATS, err := broker.OnNATS() + busAddress, err := broker.BusAddress() if err != nil { return err } - if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil { - return err - } - if onNATS { - return issueOnTheNewBus(ctx, inv, m, *forNode, busAddress) - } - - management, err := broker.ManagementFromEnvironment() - if err != nil { - return err - } - // The foundation owns the bus; make sure it exists before a module binds onto it. - if err := management.EnsureEventExchanges(ctx); err != nil { - return err - } - - secret := make([]byte, 32) - if _, err := rand.Read(secret); err != nil { - return err - } - password := base64.RawURLEncoding.EncodeToString(secret) - account, err := management.CreateModuleAccount(ctx, *forNode, module, password, m.Emits, m.Consumes) - if err != nil { - return err - } - // A consumer's queue, with its dead-letter, is the foundation's to declare — its own account - // may not (ADR 0043). Made now, so it exists before the module binds onto it. - if len(m.Consumes) > 0 { - if err := management.EnsureModuleQueue(ctx, *forNode, module); err != nil { - return err - } - } - - known, err := broker.FromEnvironment() - if err != nil { - return fmt.Errorf("cannot deliver a credential without knowing where the broker is: %w", err) - } - brokerAddr, err := brokerReachableAt(ctx, inv, known, *forNode) - if err != nil { - return err - } - // The URL and what verifies the broker, together — a mesh's broker presents its own - // certificate, in no public trust store, so a URL alone fails at TLS (as `builder issue`). - held, err := json.Marshal(struct { - URL string `json:"url"` - Fingerprint string `json:"fingerprint,omitempty"` - Node string `json:"node"` - Module string `json:"module"` - }{ - URL: fmt.Sprintf("amqps://%s:%s@%s/", account, password, brokerAddr), - Fingerprint: known.Fingerprint, - // The node and module the account is for, so the runtime names its queue as the mesh - // scoped it (..events) without a manifest having to interpolate a node. - Node: *forNode, - Module: module, - }) - if err != nil { - return err - } - if err := inv.AcceptSecretForModule(ctx, *forNode, module, "broker", string(held)); err != nil { - return err - } - fmt.Printf("broker account %s created for %s, scoped to what it emits and consumes\n", - account, module) - fmt.Printf(" sealed to %s. It arrives with the next push — `push %s` to send it\n", - *forNode, *forNode) - return nil + return issueOnTheNewBus(ctx, inv, m, *forNode, busAddress) default: return fmt.Errorf("module has no %q; it has add, list, moved, forget and issue", args[0]) diff --git a/cmd/mesh-controller/nodes.go b/cmd/mesh-controller/nodes.go index d75ee40..6f54179 100644 --- a/cmd/mesh-controller/nodes.go +++ b/cmd/mesh-controller/nodes.go @@ -10,7 +10,6 @@ import ( "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/inventory" - "github.com/novox/mesh-controller/internal/link" "github.com/novox/mesh-controller/internal/token" ) @@ -258,15 +257,6 @@ func tokenCommand(ctx context.Context, args []string) error { // chicken-and-egg entirely: the mesh runs the broker, so a joining node's credentials can // exist before it does. The one-time secret IS the password, so a node's first connection is // already authenticated and enrolment is what happens over it. - if management, err := broker.ManagementFromEnvironment(); err == nil { - if err := management.CreateNodeAccount(ctx, issued.Node.Name, issued.Secret); err != nil { - return err - } - fmt.Printf("broker account %s created, scoped to %s and the %s exchange\n\n", - issued.Node.Name, link.QueueFor(issued.Node.Name), link.Exchange) - } else if !errors.Is(err, broker.ErrNotConfigured) { - return err - } made := token.Token{Node: issued.Node.Name, Signer: key.Public, Secret: issued.Secret, Adopted: issued.Node.Adopted} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index d1f0335..d4afdf4 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -42,16 +42,10 @@ func reportUnhostable(node string, plan catalogue.Resolution) { // variable moves it). The streams and this controller's consumers are raised first on the new bus, // so nothing served here finds them missing. func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) { - busAddress, onNATS, err := broker.OnNATS() + busAddress, err := broker.BusAddress() if err != nil { return nil, err } - if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil { - return nil, err - } - if !onNATS { - return link.Connect(enroller, listener) - } if inv != nil { if err := raiseTheBus(ctx, inv, busAddress); err != nil { return nil, err @@ -88,11 +82,6 @@ func serve(ctx context.Context) error { } fmt.Printf("signing as %s\n", key.Fingerprint()[:16]) - management, err := broker.ManagementFromEnvironment() - if err != nil && !errors.Is(err, broker.ErrNotConfigured) { - return err - } - // Where the broker is and what to expect there, so a node can be told how to come back // without a person and a new token. known, err := broker.FromEnvironment() @@ -107,16 +96,9 @@ func serve(ctx context.Context) error { // **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, - OnNATS: onNATS} + work := link.Enrolment{Inventory: inv, Identity: ident, Broker: known, + OnNATS: true} server, err := connectLink(ctx, inv, work, work) if err != nil { return err diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index 647b8f8..51f7f05 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -8,6 +8,7 @@ import ( "strings" "time" + "github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/link" ) @@ -310,11 +311,8 @@ func orderByBases(entries []inventory.Entry) []inventory.Entry { seen[name] = true if e.Manifest.Build != nil { for _, on := range e.Manifest.Build.On { - if on.Module == "" || on.Module == name || !inSet[on.Module] { - continue - } for _, base := range entries { - if base.Manifest.Module == on.Module { + if base.Manifest.Module != name && inSet[base.Manifest.Module] && standsOnModule(on, base.Manifest.Module) { place(base, seen) } } @@ -336,10 +334,20 @@ func standsOn(entries []inventory.Entry, module string) bool { continue } for _, on := range e.Manifest.Build.On { - if on.Module == module { + if standsOnModule(on, module) { return true } } } return false } + +// standsOnModule is whether a base names the module: as written in a manifest (`module`), or as +// recorded after a build, when the mesh has replaced it with the artifact it resolved to +// (`artifact-store:///@…`). A recorded manifest is what the catalogue holds. +func standsOnModule(on catalogue.BuildsOn, module string) bool { + if on.Module == module { + return true + } + return strings.HasPrefix(on.Image, "artifact-store://"+module+"/") +} diff --git a/go.mod b/go.mod index 47c7bce..ac7fec7 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,7 @@ go 1.26.0 require ( github.com/jackc/pgx/v5 v5.10.0 - github.com/rabbitmq/amqp091-go v1.14.0 + github.com/nats-io/nats.go v1.54.0 golang.org/x/crypto v0.57.0 ) @@ -13,7 +13,6 @@ require ( github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/klauspost/compress v1.20.0 // indirect - github.com/nats-io/nats.go v1.54.0 // indirect github.com/nats-io/nkeys v0.4.16 // indirect github.com/nats-io/nuid v1.0.1 // indirect golang.org/x/net v0.58.0 // indirect diff --git a/go.sum b/go.sum index 5bfecc5..7edc48e 100644 --- a/go.sum +++ b/go.sum @@ -19,15 +19,11 @@ github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek= -github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= -go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= -go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M= golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA= golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= diff --git a/internal/broker/agreement_test.go b/internal/broker/agreement_test.go deleted file mode 100644 index 00f692a..0000000 --- a/internal/broker/agreement_test.go +++ /dev/null @@ -1,67 +0,0 @@ -package broker_test - -import ( - "regexp" - "testing" - - "github.com/novox/mesh-controller/internal/broker" - "github.com/novox/mesh-controller/internal/link" -) - -// The names are written twice, so a test keeps them agreeing. -// -// `link` imports `broker`, so `broker` cannot import `link` — the queue and exchange names -// therefore exist in both. A scoped account naming a queue nothing publishes to produces a -// builder that takes no work and says nothing about why, which is the worst kind of silence. -// -// An external test package, because it may import both without either importing the other. -func TestTheNamesTheBrokerScopesAreTheNamesTheLinkUses(t *testing.T) { - for _, agreed := range []struct { - what string - scoped string - actually string - }{ - {"the build queue", broker.BuildQueueName, link.BuildQueue}, - {"the exchange", broker.ExchangeName, link.Exchange}, - {"a node's queue", broker.QueueFor("somewhere"), link.QueueFor("somewhere")}, - } { - if agreed.scoped != agreed.actually { - t.Errorf("%s: the broker scopes %q and the link uses %q", - agreed.what, agreed.scoped, agreed.actually) - } - } -} - -func TestABuilderMayWriteToAReplyQueueAndReadNoNodesDeclarations(t *testing.T) { - // The scoping, checked as patterns rather than by connecting: what it may write must include - // the queue an asker actually waits on, and what it may read must not include any node's. - // - // This exists because the first version scoped writes to `amq.gen-*` — the name one broker - // happens to generate — and the builder built, could not answer, and the connection simply - // closed saying only "not allowed to publish to exchange \'\'". Answering through the - // exchange is what removed the need for any of that. - write := regexp.MustCompile("^" + regexp.QuoteMeta(broker.ExchangeName) + "$") - if !write.MatchString(broker.ExchangeName) { - t.Error("a builder may not write to the exchange, so it can take work and never answer") - } - // And not the default exchange, where permission is per exchange rather than per queue — a - // builder allowed to use it could publish into any node's queue. - // - // Confirmed against a real broker as well, and worth recording how that nearly went wrong: - // an unconfirmed publish is asynchronous, so a refusal arrives as a channel close afterwards - // and a naive check reports success. With publisher confirms the broker's refusal is - // immediate. **A negative security assertion made against an asynchronous call is not an - // assertion.** - if write.MatchString("") { - t.Error("a builder may publish to the default exchange, and so into any node's queue") - } - - read := regexp.MustCompile("^" + regexp.QuoteMeta(broker.BuildQueueName) + "$") - if !read.MatchString(broker.BuildQueueName) { - t.Error("a builder may not read the build queue") - } - if read.MatchString(link.QueueFor("someone-else")) { - // A build machine is not a node, and a node's queue carries its declarations. - t.Error("a builder may read another machine's declarations") - } -} diff --git a/internal/broker/broker_test.go b/internal/broker/broker_test.go index 8baf100..f9abd06 100644 --- a/internal/broker/broker_test.go +++ b/internal/broker/broker_test.go @@ -194,16 +194,3 @@ func TestTheAddressPortFollowsThePortTwin(t *testing.T) { t.Fatalf("the address is %q; the node put the bus on 5679", b.Address) } } - -func TestTheManagementPortFollowsThePortTwin(t *testing.T) { - t.Setenv(ManagementVar, "http://guest:guest@127.0.0.1:15672") - t.Setenv(ManagementVar+"_FILE", "") - t.Setenv(ManagementVar+"_PORT", "15673") - m, err := ManagementFromEnvironment() - if err != nil { - t.Fatal(err) - } - if m.base.Host != "127.0.0.1:15673" { - t.Fatalf("the management API is at %q; the node put it on 15673", m.base.Host) - } -} diff --git a/internal/broker/management.go b/internal/broker/management.go deleted file mode 100644 index 2693b69..0000000 --- a/internal/broker/management.go +++ /dev/null @@ -1,366 +0,0 @@ -package broker - -import ( - "bytes" - "context" - "encoding/json" - "fmt" - "github.com/novox/mesh-controller/internal/envfile" - "io" - "net/http" - "net/url" - "regexp" - "strings" - "time" -) - -// The mesh runs the broker, so there is no chicken-and-egg in a node needing an account before it -// can connect: the account is created when the token is issued, and the one-time secret in that -// token IS the password. A node's first connection is already authenticated, and enrolment is -// what happens over it. -// -// novox/hq ADR 0004's *a node holds its own identity and nothing else* is why the account is per -// node rather than shared. A shared enrolment account would let any node consume another's queue, -// which is the shared-credential fault that record exists to remove, reappearing at the transport. - -// ManagementVar holds the broker's management API, credentials included. -const ManagementVar = "MESH_BROKER_MANAGEMENT" - -// safeName is what a node may be called at the broker. -// -// The name goes into a URL path and into permission patterns, which are regular expressions. A -// name carrying a `.` or a `*` would silently widen what that node may reach — so it is -// constrained here rather than escaped later, because an escape that is forgotten once is a node -// reading everybody's queues. -var safeName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,62}$`) - -// Management is the broker's administrative interface. -type Management struct { - base *url.URL - client *http.Client -} - -// ManagementFromEnvironment reads where the management API is, if it is configured. -// -// On the port MESH_BROKER_MANAGEMENT_PORT names when the node moved it (novox/hq 04-ISSUES/102); -// the URL's own port otherwise. -func ManagementFromEnvironment() (*Management, error) { - raw, err := envfile.Placed(ManagementVar) - if err != nil { - return nil, err - } - if raw == "" { - return nil, ErrNotConfigured - } - base, err := url.Parse(raw) - if err != nil || base.Host == "" { - // The value carries a password, so it is not quoted back. - return nil, fmt.Errorf("%s is not a usable URL", ManagementVar) - } - return &Management{base: base, client: &http.Client{Timeout: 15 * time.Second}}, nil -} - -// QueueFor is the queue a node consumes from. One per node, named after it. -func QueueFor(node string) string { return "node." + node } - -// ExchangeName is where nodes publish what they have to say. One exchange, and the control plane -// is the only consumer behind it (novox/hq ADR 0006 — one consumer, so two cannot silently split -// the traffic between them). -const ExchangeName = "mesh" - -// BuildQueueName is where build work waits. Duplicated from `link` rather than imported, for the -// same reason QueueFor above is: this package must not depend on the one that uses it, and a -// constant that differed would be caught by the test that asserts they agree. -const BuildQueueName = "builds" - -// The events bus (novox/hq ADR 0042): one topic exchange every event rides, a second for tool -// RPC kept apart, and a dead-letter home for a poison event. The foundation owns these — a module's -// account cannot declare them, only bind its own queue to the events one. -const ( - EventsExchangeName = "mesh.events" - RPCExchangeName = "mesh.rpc" - DeadExchangeName = "mesh.events.dead" -) - -// ModuleQueueFor is the durable queue a module consumes its events from — one per module per node -// (novox/hq ADR 0042), named so the account that may read it is exactly this module's. -func ModuleQueueFor(node, module string) string { return node + "." + module + ".events" } - -// modulePermissions is a module's authority on the bus, derived from its manifest (novox/hq -// ADR 0043): what it consumes and what it emits, and nothing else. Pure, so the scope is tested as -// patterns without a broker — the way a builder's is. -// -// A note on the limit: the broker's write permission is per exchange, not per routing key (LavinMQ -// has no topic permissions), so an emitting module is granted the events exchange whole. ADR 0042's -// origin reservation — a module publishes only under `module..*` — is stamped by the sdk, not -// enforced here; that gap is the broker's, and is recorded rather than hidden. A pure consumer like -// the audit logger is unaffected: it is granted no write to the exchange at all. -func modulePermissions(node, module string, emits, consumes []string) (configure, write, read string) { - queue := regexp.QuoteMeta(ModuleQueueFor(node, module)) - events := regexp.QuoteMeta(EventsExchangeName) - rpc := regexp.QuoteMeta(RPCExchangeName) - // A module serves each of its tools on its own queue, namespaced by the module (novox/hq - // ADR 0047) — serve.. — so the account may declare, bind and read exactly its own, - // and no other module's. - serve := "serve\\." + regexp.QuoteMeta(module) + "\\..*" - - // Declare its own events queue and its own tool serve queues. - configure = "^(" + queue + "|" + serve + ")$" - - // Write to bind its queue and serve queues (binding is a write on the queue), and to the RPC - // exchange to publish replies (ADR 0047: replies ride mesh.rpc, never the default exchange, which - // would let it publish into any queue). To the events exchange only if it emits. - writes := []string{queue, serve, rpc} - if len(emits) > 0 { - writes = append(writes, events) - } - write = "^(" + strings.Join(writes, "|") + ")$" - - // Read its own queue and serve queues to consume them, and the RPC exchange to bind its serve - // queues onto. The events exchange to bind onto only if it consumes. - reads := []string{queue, serve, rpc} - if len(consumes) > 0 { - reads = append(reads, events) - } - read = "^(" + strings.Join(reads, "|") + ")$" - return configure, write, read -} - -// CreateModuleAccount gives an assigned module its own broker account, scoped by what it emits and -// consumes (novox/hq ADR 0043). The account name carries the node so the same module on two machines -// holds two accounts, each sealed to its own; the permissions carry the module so one module cannot -// read another's queue. Generic — the builder is one instance of this rule, not a separate kind. -func (m *Management) CreateModuleAccount(ctx context.Context, node, module, password string, emits, consumes []string) (string, error) { - if !safeName.MatchString(node) { - return "", fmt.Errorf("%q cannot be part of a broker account: it is a permission pattern", node) - } - if !safeName.MatchString(module) { - return "", fmt.Errorf("%q cannot be part of a broker account: it is a permission pattern", module) - } - account := node + "-" + module - if !safeName.MatchString(account) { - return "", fmt.Errorf("%q is not a usable broker account name", account) - } - - if err := m.put(ctx, "/api/users/"+url.PathEscape(account), - map[string]string{"password": password, "tags": ""}); err != nil { - return "", fmt.Errorf("cannot create the broker account for %s on %s: %w", module, node, err) - } - - configure, write, read := modulePermissions(node, module, emits, consumes) - if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(account), map[string]string{ - "configure": configure, "write": write, "read": read, - }); err != nil { - return "", fmt.Errorf("cannot scope the broker account for %s on %s: %w", module, node, err) - } - return account, nil -} - -// EnsureModuleQueue declares a consuming module's queue with its dead-letter exchange, idempotently. -// The foundation declares it because a scoped module account may not: the broker refuses a queue with -// a dead-letter exchange to a non-administrator (novox/hq ADR 0043), so a consumer passively checks -// the queue the mesh made rather than declaring its own. -func (m *Management) EnsureModuleQueue(ctx context.Context, node, module string) error { - queue := ModuleQueueFor(node, module) - if err := m.put(ctx, "/api/queues/%2f/"+url.PathEscape(queue), map[string]any{ - "durable": true, - "arguments": map[string]any{"x-dead-letter-exchange": DeadExchangeName}, - }); err != nil { - return fmt.Errorf("cannot declare the queue for %s on %s: %w", module, node, err) - } - return nil -} - -// EnsureEventExchanges declares the bus's exchanges and the dead-letter home, idempotently. The -// foundation owns them (a module's account may not declare an exchange), and a dead-letter exchange -// with no queue behind it drops what it receives — so a durable queue bound to `#` retains a poison -// event for inspection, which is the whole reason the trail exists. -func (m *Management) EnsureEventExchanges(ctx context.Context) error { - for _, exchange := range []string{EventsExchangeName, RPCExchangeName, DeadExchangeName} { - if err := m.put(ctx, "/api/exchanges/%2f/"+url.PathEscape(exchange), - map[string]any{"type": "topic", "durable": true}); err != nil { - return fmt.Errorf("cannot declare the %s exchange: %w", exchange, err) - } - } - if err := m.put(ctx, "/api/queues/%2f/"+url.PathEscape(DeadExchangeName), - map[string]any{"durable": true}); err != nil { - return fmt.Errorf("cannot declare the dead-letter queue: %w", err) - } - if err := m.post(ctx, "/api/bindings/%2f/e/"+url.PathEscape(DeadExchangeName)+ - "/q/"+url.PathEscape(DeadExchangeName), map[string]string{"routing_key": "#"}); err != nil { - return fmt.Errorf("cannot bind the dead-letter queue: %w", err) - } - return nil -} - -// CreateNodeAccount gives a node its own broker account, with the token's secret as the password. -// -// Scoped so a node can reach its own queue and the one exchange, and nothing else. The patterns -// are anchored: a node called `laptop` must not be able to read `laptop-of-somebody-else`. -func (m *Management) CreateNodeAccount(ctx context.Context, node, password string) error { - if !safeName.MatchString(node) { - return fmt.Errorf( - "%q cannot be a broker account name: it becomes part of a permission pattern, so it "+ - "is lower-case letters, digits and dashes", node) - } - - if err := m.put(ctx, "/api/users/"+url.PathEscape(node), - map[string]string{"password": password, "tags": ""}); err != nil { - return fmt.Errorf("cannot create the broker account for %s: %w", node, err) - } - - queue := regexp.QuoteMeta(QueueFor(node)) - if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(node), map[string]string{ - "configure": "^" + queue + "$", - "write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + queue + ")$", - "read": "^" + queue + "$", - }); err != nil { - return fmt.Errorf("cannot scope the broker account for %s: %w", node, err) - } - return nil -} - -// CreateBuilderAccount scopes an account to taking build work and answering it. -// -// **A build machine is not a node**, and giving it a node's account would let it read another -// machine's declarations. What it needs is narrower and different: read the build queue, and -// write to the exchange and to whatever temporary queue an asker is waiting on. -// -// The reply queues are the reason `write` is not simply the exchange. `RequestBuild` declares an -// exclusive queue with a generated name and waits on it, so a builder that could not write to it -// could take work and never answer — which is the failure that looks like a builder that is not -// running. -func (m *Management) CreateBuilderAccount(ctx context.Context, name, password string) error { - if !safeName.MatchString(name) { - return fmt.Errorf( - "%q cannot be a broker account name: it becomes part of a permission pattern, so it "+ - "is lower-case letters, digits and dashes", name) - } - - if err := m.put(ctx, "/api/users/"+url.PathEscape(name), - map[string]string{"password": password, "tags": ""}); err != nil { - return fmt.Errorf("cannot create the broker account for %s: %w", name, err) - } - - builds := regexp.QuoteMeta(BuildQueueName) - if err := m.put(ctx, "/api/permissions/%2f/"+url.PathEscape(name), map[string]string{ - // It declares the build queue, because whichever builder starts first must be able to — - // and a queue nobody may declare is a queue that exists only if the control plane has - // already run, which makes the order they start in matter. - "configure": "^" + builds + "$", - // Two exchanges, and nothing else. **Not the default exchange**: permission there is - // granted per exchange rather than per queue, so a builder allowed to use it could - // publish into any node's queue — the privilege a build machine most obviously should - // not have. Answers go through the node exchange; announcing what was built goes through - // the events exchange, which is a different act with a different audience (novox/hq - // ADR 0072). A builder that could answer and not announce would leave the module graph - // knowing less than the registry does. - "write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + regexp.QuoteMeta(EventsExchangeName) + ")$", - // The build queue and nothing else. Not another machine's declarations. - "read": "^" + builds + "$", - }); err != nil { - return fmt.Errorf("cannot scope the broker account for %s: %w", name, err) - } - return nil -} - -// RemoveNodeAccount withdraws a node's access. -func (m *Management) RemoveNodeAccount(ctx context.Context, node string) error { - if !safeName.MatchString(node) { - return fmt.Errorf("%q is not a broker account name", node) - } - return m.do(ctx, http.MethodDelete, "/api/users/"+url.PathEscape(node), nil) -} - -// Accounts lists the broker's users, so a picture can be read from the system rather than assumed -// (novox/hq ADR 0018). -func (m *Management) Accounts(ctx context.Context) ([]string, error) { - body, err := m.get(ctx, "/api/users") - if err != nil { - return nil, err - } - var users []struct { - Name string `json:"name"` - } - if err := json.Unmarshal(body, &users); err != nil { - return nil, err - } - names := make([]string, 0, len(users)) - for _, u := range users { - names = append(names, u.Name) - } - return names, nil -} - -func (m *Management) put(ctx context.Context, path string, body any) error { - return m.do(ctx, http.MethodPut, path, body) -} - -func (m *Management) post(ctx context.Context, path string, body any) error { - return m.do(ctx, http.MethodPost, path, body) -} - -func (m *Management) get(ctx context.Context, path string) ([]byte, error) { - request, err := m.request(ctx, http.MethodGet, path, nil) - if err != nil { - return nil, err - } - response, err := m.client.Do(request) - if err != nil { - return nil, err - } - defer response.Body.Close() - if response.StatusCode >= 300 { - return nil, fmt.Errorf("the broker's management API answered %s to GET %s", - response.Status, path) - } - return io.ReadAll(io.LimitReader(response.Body, 1<<20)) -} - -func (m *Management) do(ctx context.Context, method, path string, body any) error { - request, err := m.request(ctx, method, path, body) - if err != nil { - return err - } - response, err := m.client.Do(request) - if err != nil { - return err - } - defer response.Body.Close() - if response.StatusCode >= 300 { - detail, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) - return fmt.Errorf("the broker's management API answered %s to %s %s: %s", - response.Status, method, path, strings.TrimSpace(string(detail))) - } - return nil -} - -func (m *Management) request(ctx context.Context, method, path string, body any) (*http.Request, error) { - var payload io.Reader - if body != nil { - raw, err := json.Marshal(body) - if err != nil { - return nil, err - } - payload = bytes.NewReader(raw) - } - - // Path joined by hand rather than through url.Parse: %2f is the default vhost and must reach - // the broker still encoded. Parsing would decode it to a slash and address a different route. - target := strings.TrimSuffix(m.base.String(), "/") - if user := m.base.User; user != nil { - target = strings.TrimSuffix(m.base.Scheme+"://"+m.base.Host, "/") - } - request, err := http.NewRequestWithContext(ctx, method, target+path, payload) - if err != nil { - return nil, err - } - if user := m.base.User; user != nil { - password, _ := user.Password() - request.SetBasicAuth(user.Username(), password) - } - if body != nil { - request.Header.Set("Content-Type", "application/json") - } - return request, nil -} diff --git a/internal/broker/module_account_test.go b/internal/broker/module_account_test.go deleted file mode 100644 index 2f401fc..0000000 --- a/internal/broker/module_account_test.go +++ /dev/null @@ -1,96 +0,0 @@ -package broker_test - -import ( - "context" - "encoding/json" - "net/http" - "net/http/httptest" - "regexp" - "strings" - "testing" - - "github.com/novox/mesh-controller/internal/broker" -) - -// The audit logger consumes everything and emits nothing. Its account must let it declare and read -// its own queue and read the events exchange to bind onto — and must not reach another module's -// queue, nor grant any write to the events exchange (novox/hq ADR 0043). -func TestAConsumerReadsItsOwnQueueAndTheEventsExchangeAndNoOthers(t *testing.T) { - // Rebuilt through CreateModuleAccount's own path by asking for the same scope it would apply. - // The queue this module reads: - mine := broker.ModuleQueueFor("anchor", "audit-logger") - other := broker.ModuleQueueFor("anchor", "plex") - - read := scope(t, "read", "anchor", "audit-logger", nil, []string{"#"}) - if !read.MatchString(mine) { - t.Error("the audit logger may not read its own queue, so it consumes nothing") - } - if !read.MatchString(broker.EventsExchangeName) { - t.Error("the audit logger may not read the events exchange, so it cannot bind onto it") - } - if read.MatchString(other) { - t.Error("the audit logger may read another module's queue") - } - - // It emits nothing, so it is granted no write to the events exchange — only its own queue, to bind. - write := scope(t, "write", "anchor", "audit-logger", nil, []string{"#"}) - if write.MatchString(broker.EventsExchangeName) { - t.Error("a pure consumer was granted write to the events exchange") - } - if !write.MatchString(mine) { - t.Error("the audit logger may not write to its own queue, so it cannot bind it") - } -} - -// An emitter is granted the events exchange to write; a consumer is not. -func TestAnEmitterMayWriteTheEventsExchangeAndAConsumerMayNot(t *testing.T) { - emitter := scope(t, "write", "anchor", "umami", []string{"module.umami.site.created"}, nil) - if !emitter.MatchString(broker.EventsExchangeName) { - t.Error("an emitting module may not write the events exchange, so it cannot emit") - } -} - -// scope reconstructs one of the three permission patterns CreateModuleAccount would apply, by -// reading it back from a captured request against a stub management API. -func scope(t *testing.T, which, node, module string, emits, consumes []string) *regexp.Regexp { - t.Helper() - pat := capturePermission(t, which, node, module, emits, consumes) - re, err := regexp.Compile(pat) - if err != nil { - t.Fatalf("the %s pattern does not compile: %v", which, err) - } - return re -} - -// capturePermission runs CreateModuleAccount against a stub management API and returns the pattern -// it set for `which` ("configure"/"write"/"read"). The scope is tested where it is applied, not -// reconstructed by the test — so a change to the mapping cannot pass a test that hard-codes the old -// one. -func capturePermission(t *testing.T, which, node, module string, emits, consumes []string) string { - t.Helper() - var captured map[string]string - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - if strings.HasPrefix(r.URL.Path, "/api/permissions/") { - _ = json.NewDecoder(r.Body).Decode(&captured) - } - w.WriteHeader(http.StatusNoContent) - })) - defer server.Close() - - t.Setenv(broker.ManagementVar, server.URL) - m, err := broker.ManagementFromEnvironment() - if err != nil { - t.Fatalf("stub management not usable: %v", err) - } - if _, err := m.CreateModuleAccount(context.Background(), node, module, "pw", emits, consumes); err != nil { - t.Fatalf("CreateModuleAccount: %v", err) - } - if captured == nil { - t.Fatal("no permissions were set") - } - pattern, ok := captured[which] - if !ok { - t.Fatalf("no %s permission was set; got %v", which, captured) - } - return pattern -} diff --git a/internal/broker/onnats.go b/internal/broker/onnats.go index 245f429..1314977 100644 --- a/internal/broker/onnats.go +++ b/internal/broker/onnats.go @@ -68,24 +68,16 @@ func BareAddress(address string) string { return "nats://" + bare } -// 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) +// BusAddress is where the mesh's bus is, with this process's credential, and refuses to be empty: +// there is one bus, and a control plane without it can hold records and answer nothing (novox/hq +// ADR 0131, design 28 task 5.5). +func BusAddress() (string, error) { + address, _, err := OnNATS() + if err != nil { + return "", err } - return nil + if address == "" { + return "", fmt.Errorf("this control plane has no %s, so it cannot reach the mesh's bus", NATSVar) + } + return address, 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 index f21bc87..099f1e7 100644 --- a/internal/broker/onnats_test.go +++ b/internal/broker/onnats_test.go @@ -1,43 +1,11 @@ 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) - } - } -} // The controller's own credential arrives in its address, and has to be readable out of it — its user // is created by the installer at a bootstrap password, before the controller exists to mint one. diff --git a/internal/link/ask.go b/internal/link/ask.go index 6d3a19c..41adda6 100644 --- a/internal/link/ask.go +++ b/internal/link/ask.go @@ -3,16 +3,13 @@ package link import ( "context" "encoding/json" + "errors" "fmt" "time" - amqp "github.com/rabbitmq/amqp091-go" + "github.com/nats-io/nats.go" ) -// RPCExchange is where a module's tools are asked over the broker, keyed `.`, and -// where the answer comes back, keyed by the asker's reply queue (novox/hq ADR 0047). -const RPCExchange = "mesh.rpc" - // Answer is what a module's tool replies: one of the two, never both. type Answer struct { Result json.RawMessage `json:"result,omitempty"` @@ -27,70 +24,35 @@ type Answer struct { // grants, and should not. The control plane already holds a connection that may, so a person or // an agent asks through it, and every question passes one process where an audit belongs. // -// The reply queue is the caller's own, server-named and exclusive, bound to the RPC exchange under -// its own name: a serving module answers through that exchange and never the default one, whose -// permission is per exchange rather than per queue. The correlation is checked rather than -// assumed, as every RPC here is. -func Ask(ctx context.Context, channel *amqp.Channel, module, tool string, args json.RawMessage, +// The answer comes back on the asker's own inbox, which only the asker may read; the serving +// module answers there and nowhere else. The bus refuses a request nothing serves at once, so a +// module that is down or a tool that does not exist is said now rather than after the whole wait. +func Ask(ctx context.Context, bus Bus, module, tool string, args json.RawMessage, timeout time.Duration) (Answer, error) { if len(args) == 0 { args = json.RawMessage(`{}`) } - replies, err := channel.QueueDeclare("", false, true, true, false, nil) - if err != nil { - return Answer{}, err - } - if err := channel.QueueBind(replies.Name, replies.Name, RPCExchange, false, nil); err != nil { - return Answer{}, fmt.Errorf("cannot bind a reply queue to %s: %w", RPCExchange, err) - } - answers, err := channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil) - if err != nil { - return Answer{}, err - } - - // Mandatory, so a request nothing consumes comes straight back: a module that is down, or a - // tool that does not exist, is said at once rather than after the whole wait. - returned := channel.NotifyReturn(make(chan amqp.Return, 1)) - id := fmt.Sprintf("ask-%d", time.Now().UnixNano()) key := module + "." + tool - if err := channel.PublishWithContext(ctx, RPCExchange, key, true, false, amqp.Publishing{ - ContentType: "application/json", - CorrelationId: id, - ReplyTo: replies.Name, - Body: args, - }); err != nil { - return Answer{}, fmt.Errorf("cannot ask %s: %w", key, err) - } - - waiting, cancel := context.WithTimeout(ctx, timeout) - defer cancel() - for { - select { - case back := <-returned: - if back.CorrelationId == id { - return Answer{}, fmt.Errorf( - "nothing serves %s: no runtime has bound %q on the broker. The module is not "+ - "assigned, its runtime is not up, or it serves no such tool — `status` "+ - "says whether the machine carrying it has applied", module, key) - } - case <-waiting.Done(): + reply, err := bus.AskTool(ctx, module, tool, args, timeout) + if err != nil { + if errors.Is(err, nats.ErrNoResponders) { + return Answer{}, fmt.Errorf( + "nothing serves %s: no runtime has bound %q on the bus. The module is not "+ + "assigned, its runtime is not up, or it serves no such tool — `status` says "+ + "whether the machine carrying it has applied", module, key) + } + if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, nats.ErrTimeout) { return Answer{}, fmt.Errorf( "%s did not answer within %s. Its runtime serves %q when it is up and has bound "+ - "the broker — `status` says whether the machine carrying it has applied", + "the bus — `status` says whether the machine carrying it has applied", module, timeout, key) - case delivery, ok := <-answers: - if !ok { - return Answer{}, fmt.Errorf("the connection closed while waiting for %s", key) - } - if delivery.CorrelationId != id { - continue - } - var answer Answer - if err := json.Unmarshal(delivery.Body, &answer); err != nil { - return Answer{}, fmt.Errorf("%s answered with something unreadable: %w", key, err) - } - return answer, nil } + return Answer{}, fmt.Errorf("cannot ask %s: %w", key, err) } + var answer Answer + if err := json.Unmarshal(reply, &answer); err != nil { + return Answer{}, fmt.Errorf("%s answered with something unreadable: %w", key, err) + } + return answer, nil } diff --git a/internal/link/build.go b/internal/link/build.go index f097233..45dad1f 100644 --- a/internal/link/build.go +++ b/internal/link/build.go @@ -1,12 +1,7 @@ package link import ( - "context" "encoding/json" - "fmt" - "time" - - amqp "github.com/rabbitmq/amqp091-go" ) // Asking a machine to build a module, and hearing what came out. @@ -22,21 +17,6 @@ import ( // holds no opinion about what they contain, and a host that also built things would be a host // with a container runtime requirement and a git dependency (novox/hq ADR 0005). -// BuildQueue is where build requests wait. One queue, so several build machines can share the -// work and each request is done exactly once — which is what a queue is for and what a -// per-machine routing key would not give. -const BuildQueue = "builds" - -// KeyBuilt is what a builder publishes when it has finished, successfully or not. -const KeyBuilt = "built" - -// ReplyQueue is where the answer to one request goes. -// -// **Named here rather than left to the broker**, so it can be scoped and reasoned about. A broker -// generates its own name for an unnamed queue, and a builder permitted to write to whatever that -// convention happens to produce works on one broker and silently cannot answer on another. -func ReplyQueue(id string) string { return BuildQueue + ".reply." + id } - // BuildRequest is one module to build. type BuildRequest struct { // ID correlates the answer with the asking. Not the module name: two builds of one module can @@ -118,77 +98,3 @@ type MadeArtifact struct { Kind string `json:"kind"` Reference string `json:"reference"` } - -// RequestBuild asks for a module to be built and waits for the answer. -// -// Waiting rather than returning immediately, because the thing a person wants after asking for a -// build is to know whether it worked. A build that is dispatched and forgotten needs somewhere to -// look afterwards, and there is nowhere yet. -func RequestBuild(ctx context.Context, channel *amqp.Channel, request BuildRequest, - timeout time.Duration) (BuildResult, error) { - - // Its own queue for the answer, declared before the ask. Consuming from the shared exchange - // would mean competing with the control plane's own consumer for a message meant for this - // caller — which is the fault this package's own doc comment records having had. - replies, err := channel.QueueDeclare(ReplyQueue(request.ID), false, true, true, false, nil) - if err != nil { - return BuildResult{}, err - } - // Bound to the exchange, and the answer comes back through it. - // - // **A builder never publishes to the default exchange**, because permission there is per - // exchange and not per queue — a builder allowed to use it could publish into any node's - // queue, which is the privilege a build machine most obviously should not have. Found by - // running it: the builder built, could not answer, and the connection closed saying only - // "not allowed to publish to exchange ''". - // - // The cost is that every asker sees every result, which is why the correlation is checked - // below rather than assumed. - if err := channel.QueueBind(replies.Name, KeyBuilt, Exchange, false, nil); err != nil { - return BuildResult{}, err - } - answers, err := channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil) - if err != nil { - return BuildResult{}, err - } - - body, err := json.Marshal(request) - if err != nil { - return BuildResult{}, err - } - if err := channel.PublishWithContext(ctx, "", BuildQueue, false, false, amqp.Publishing{ - ContentType: "application/json", - DeliveryMode: amqp.Persistent, - CorrelationId: request.ID, - ReplyTo: replies.Name, - Body: body, - }); err != nil { - return BuildResult{}, err - } - - waiting, cancel := context.WithTimeout(ctx, timeout) - defer cancel() - for { - select { - case <-waiting.Done(): - return BuildResult{}, fmt.Errorf( - "no builder answered within %s. Something must be consuming %q, and nothing is "+ - "— or it is building something that takes longer than this", - timeout, BuildQueue) - case delivery, ok := <-answers: - if !ok { - return BuildResult{}, fmt.Errorf("the connection closed while waiting for a build") - } - var result BuildResult - if err := json.Unmarshal(delivery.Body, &result); err != nil { - return BuildResult{}, fmt.Errorf("a builder answered with something unreadable: %w", err) - } - if result.ID != request.ID { - // Somebody else's answer on this queue. Ignored rather than returned, because - // returning it would attribute one build's outcome to another's. - continue - } - return result, nil - } - } -} diff --git a/internal/link/builds_current.go b/internal/link/builds_current.go deleted file mode 100644 index 57165d2..0000000 --- a/internal/link/builds_current.go +++ /dev/null @@ -1,202 +0,0 @@ -package link - -import ( - "context" - "encoding/json" - "errors" - "fmt" - "time" - - amqp "github.com/rabbitmq/amqp091-go" -) - -// The build flow on the bus the mesh runs on today. -// -// Moved behind the seam rather than changed. The queue, the reply binding and the correlation are -// what they were, because the mesh is running on this. - -// currentBuilds asks for builds over a channel. -type currentBuilds struct{ channel *amqp.Channel } - -// BuildsOverCurrent is the asking side on the bus the mesh has. -func BuildsOverCurrent(channel *amqp.Channel) Builders { return currentBuilds{channel: channel} } - -func (b currentBuilds) Close() {} - -func (b currentBuilds) Submit(ctx context.Context, request BuildRequest, - wait time.Duration) (BuildResult, error) { - - // Its own queue for the answer, declared before the ask. Consuming from the shared exchange - // would mean competing with the controller's own consumer for a message meant for this caller. - replies, err := b.channel.QueueDeclare(ReplyQueue(request.ID), false, true, true, false, nil) - if err != nil { - return BuildResult{}, err - } - // **A builder never publishes to the default exchange**, because permission there is per - // exchange and not per queue — a builder allowed to use it could publish into any node's queue, - // which is the privilege a build machine most obviously should not have. The cost is that every - // asker sees every result, which is why the correlation is checked below rather than assumed. - if err := b.channel.QueueBind(replies.Name, KeyBuilt, Exchange, false, nil); err != nil { - return BuildResult{}, err - } - answers, err := b.channel.ConsumeWithContext(ctx, replies.Name, "", true, true, false, false, nil) - if err != nil { - return BuildResult{}, err - } - - body, err := json.Marshal(request) - if err != nil { - return BuildResult{}, err - } - if err := b.channel.PublishWithContext(ctx, "", BuildQueue, false, false, amqp.Publishing{ - ContentType: "application/json", - DeliveryMode: amqp.Persistent, - CorrelationId: request.ID, - ReplyTo: replies.Name, - Body: body, - }); err != nil { - return BuildResult{}, err - } - - waiting, cancel := context.WithTimeout(ctx, wait) - defer cancel() - for { - select { - case <-waiting.Done(): - return BuildResult{}, waitingFor(wait) - case delivery, ok := <-answers: - if !ok { - return BuildResult{}, errors.New("the connection closed while waiting for a build") - } - result, mine, err := theOutcomeOf(delivery.Body, request.ID) - if err != nil { - return BuildResult{}, err - } - if mine { - return result, nil - } - } - } -} - -// --- the machine's side --------------------------------------------------------------------- - -type currentMachine struct { - conn *amqp.Connection - channel *amqp.Channel - on string -} - -// MachineOverCurrent takes build work over a channel. -func MachineOverCurrent(conn *amqp.Connection, channel *amqp.Channel, on string) BuildMachine { - return ¤tMachine{conn: conn, channel: channel, on: on} -} - -func (m *currentMachine) Close() {} - -func (m *currentMachine) Take(ctx context.Context, do func(context.Context, Build)) error { - if _, err := m.channel.QueueDeclare(BuildQueue, true, false, false, false, nil); err != nil { - return err - } - // One at a time. A machine that took five requests at once would run five container builds - // against one runtime and finish all of them slower than it would have finished the first — and - // the queue is what shares work between machines, so nothing is lost by it. - if err := m.channel.Qos(1, 0, false); err != nil { - return err - } - // Not auto-acknowledged: a request acknowledged on arrival is a build that vanishes if this - // process dies mid-way, with nobody waiting on it ever hearing why. - requests, err := m.channel.ConsumeWithContext(ctx, BuildQueue, "mesh-builder", - false, false, false, false, nil) - if err != nil { - return err - } - for { - select { - case <-ctx.Done(): - return nil - case delivery, ok := <-requests: - if !ok { - return errors.New("the broker closed the connection") - } - var request BuildRequest - if err := json.Unmarshal(delivery.Body, &request); err != nil { - // Unreadable: rejected rather than retried, because the next attempt reads the same - // bytes. Nobody waiting hears an answer, which is correct — there was no request. - _ = delivery.Reject(false) - continue - } - do(ctx, ¤tBuild{request: request, delivery: delivery, on: m.on, channel: m.channel}) - } - } -} - -type currentBuild struct { - request BuildRequest - delivery amqp.Delivery - on string - channel *amqp.Channel -} - -func (b *currentBuild) Request() BuildRequest { return b.request } - -// Announce answers and announces, which on this bus are two publishes to two exchanges. -// -// The reply goes to whoever asked, correlated to their request; the announcement says to the whole -// mesh that a module now exists at a commit (novox/hq ADR 0072). Only a successful build is -// announced: a failed one produced no module version, and announcing one would put something in the -// graph that was never made. -func (b *currentBuild) Announce(ctx context.Context, result BuildResult) error { - body, err := json.Marshal(result) - if err != nil { - return err - } - if err := b.channel.PublishWithContext(ctx, Exchange, KeyBuilt, false, false, amqp.Publishing{ - ContentType: "application/json", - CorrelationId: result.ID, - Body: body, - }); err != nil { - return fmt.Errorf("cannot answer a build request: %w", err) - } - if result.Failed != "" || result.Commit == "" { - return nil - } - - // **Announced under both names on this bus, for exactly as long as this bus lives.** - // - // A build's outcome belongs to the role now (novox/hq ADR 0121), so a catalogue built from the - // current manifests listens for the role's name. A catalogue that is *already running* listens for - // the module's, because that is what it was told when it was installed. A rename on a live bus - // needs the publisher and the subscriber to change together, and a merge cannot promise that: one - // of them is deployed first, and in that window the graph silently stops being updated — which is - // the failure this whole change was cleaning up after. - // - // So both, and the order stops mattering. The module's own name goes with the bus, in step 5's - // retirement list; nothing has ever run on the bus being built, so there is no legacy name there - // and this doubling has no counterpart. - announced := announcementOf(result) - if err := EmitEvent(ctx, OverCurrent{Channel: b.channel}, KeyModuleBuilt, "builder", b.on, - announced); err != nil { - return err - } - return EmitEvent(ctx, OverCurrent{Channel: b.channel}, KeyRoleBuilt, TheBuildMachine, b.on, - announced) -} - -func (b *currentBuild) Done() error { return b.delivery.Ack(false) } - -func (b *currentBuild) Hold(time.Duration) error { - // No delayed redelivery on this bus: handed back at once, which is what it has always done. - return b.delivery.Nack(false, true) -} - -// announcementOf is what the mesh is told about a finished build. One function, so the two -// transports cannot describe the same build differently. -func announcementOf(result BuildResult) map[string]any { - return map[string]any{ - "module": ModuleOf(result.Manifest), "commit": result.Commit, - "repository": result.Repository, "path": result.Path, "ref": result.Ref, - "manifest": json.RawMessage(result.Manifest), "against": result.Against, - "made": result.Made, - } -} diff --git a/internal/link/bus.go b/internal/link/bus.go index b96d3cf..dbe4f46 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -7,7 +7,6 @@ import ( "time" "github.com/nats-io/nats.go" - amqp "github.com/rabbitmq/amqp091-go" ) // Bus is what the controller needs of the mesh's bus, **in the mesh's own words rather than a @@ -39,51 +38,6 @@ type Bus interface { // --- The bus the mesh runs on today ----------------------------------------------------- -// OverCurrent is the bus the mesh runs on today, until the rollout. -type OverCurrent struct{ Channel *amqp.Channel } - -func (b OverCurrent) PublishEvent(ctx context.Context, key, source, node string, body []byte) error { - id, err := eventID() - if err != nil { - return err - } - return b.Channel.PublishWithContext(ctx, EventsExchange, key, false, false, amqp.Publishing{ - ContentType: "application/json", - DeliveryMode: amqp.Persistent, - MessageId: id, - Timestamp: time.Now().UTC(), - Body: body, - Headers: amqp.Table{ - "x-event-id": id, - "x-source": source, - "x-node": node, - "x-time": time.Now().UTC().Format(time.RFC3339), - "content-type": "application/json", - }, - }) -} - -// AskTool is implemented over the existing reply-queue machinery in ask.go; this seam does not -// change how it works today. -func (b OverCurrent) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) { - answer, err := Ask(ctx, b.Channel, module, tool, args, timeout) - if err != nil { - return nil, err - } - return answer.Result, nil -} - -func (b OverCurrent) PublishDeclaration(ctx context.Context, node string, body []byte) error { - // To the queue directly rather than through an exchange: a declaration is for one node, and - // routing it by name through a shared exchange would mean a binding per node that nothing - // removes when a node is retired. - return b.Channel.PublishWithContext(ctx, "", QueueFor(node), false, false, amqp.Publishing{ - ContentType: "application/json", - DeliveryMode: amqp.Persistent, - Body: body, - }) -} - // --- NATS, the bus being built ---------------------------------------------------------------- // OverNATS is the bus as a JetStream context. diff --git a/internal/link/bus_nats.go b/internal/link/bus_nats.go index f4e4511..b902da4 100644 --- a/internal/link/bus_nats.go +++ b/internal/link/bus_nats.go @@ -1,66 +1,16 @@ package link import ( - "context" - "fmt" - "time" - - "github.com/nats-io/nats.go" - "github.com/novox/mesh-controller/internal/broker" ) -// OverNats is the controller's outbound on the bus being built: the same three acts the other -// transport has, on the subjects the permissions were derived for (design 25). A declaration is a -// JetStream publish into NODES, where the node's own consumer waits for it; an event is announced on -// the subject its name derives to; a tool is asked by request and reply on the module's tool subject. -type OverNats struct{ JS *broker.JetStream } - -// declareSubject is where one node's declaration lands — the NODES stream's subject for it, and the -// only subject that node's consumer delivers. The host subscribes exactly this. -func declareSubject(node string) string { return "mesh.node." + node + ".declare" } - -func (b OverNats) PublishDeclaration(ctx context.Context, node string, body []byte) error { - publish, cancel := context.WithTimeout(ctx, 15*time.Second) - defer cancel() - if _, err := b.JS.Context().Publish(declareSubject(node), body, nats.Context(publish)); err != nil { - return fmt.Errorf("declaring to %s: %w", node, err) - } - return nil -} - -func (b OverNats) PublishEvent(ctx context.Context, key, source, node string, body []byte) error { - // The key is the subject: the controller's own events are named in full, and what a module - // emits is derived before it reaches here. Headers carry the envelope the other transport put - // in message properties (ADR 0042), so a consumer reads who and when without the payload. - msg := nats.NewMsg(key) - msg.Data = body - msg.Header.Set("x-source", source) - msg.Header.Set("x-node", node) - msg.Header.Set("x-time", time.Now().UTC().Format(time.RFC3339Nano)) - if err := b.JS.Conn().PublishMsg(msg); err != nil { - return fmt.Errorf("announcing %s: %w", key, err) - } - return nil -} - -func (b OverNats) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) { - ask, cancel := context.WithTimeout(ctx, timeout) - defer cancel() - reply, err := b.JS.Conn().RequestWithContext(ask, "mesh.mod."+module+".tool."+tool, args) - if err != nil { - return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err) - } - return reply.Data, nil -} - // ConnectNats is Connect for the bus being built: the controller's inbound and outbound over one // JetStream connection the caller has already raised the streams on. Nothing is declared here — // the streams and the controller's consumers are asserted by Raise, before anything is served. func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server { return &Server{ inbound: Nats(js), - bus: OverNats{JS: js}, + bus: OverNATS{JS: js.Context(), Conn: js.Conn()}, js: js, enroller: enroller, listener: listener, diff --git a/internal/link/enrolment.go b/internal/link/enrolment.go index 750ecab..615a60e 100644 --- a/internal/link/enrolment.go +++ b/internal/link/enrolment.go @@ -24,10 +24,9 @@ import ( // — it holds both grants and asks each for its part, which is what the process running them is // for. type Enrolment struct { - Inventory *inventory.Inventory - Identity *identity.Identity - Management *broker.Management - Broker broker.Broker + Inventory *inventory.Inventory + Identity *identity.Identity + Broker broker.Broker // OnNATS says the mesh's own traffic is on the bus being built, so a node's credential is // minted into the mesh's records and composed into the bus's user list rather than pushed @@ -194,17 +193,6 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply Enrol reply.Password = password } - case e.Management != nil: - password, err := freshPassword() - if err != nil { - return EnrolReply{}, err - } - if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil { - log.Printf("%s is enrolled and its broker password could not be replaced, so it keeps "+ - "the token's secret as its password: %v", node.Name, err) - } else { - reply.Password = password - } } if profile != nil { @@ -261,9 +249,6 @@ func freshPassword() (string, error) { var _ Enroller = Enrolment{} -// ErrNoBrokerManagement is returned when an account cannot be made because nothing was configured. -var ErrNoBrokerManagement = errors.New("no broker management configured") - // Heard records what a node reported about itself. // // A node states; the owning context writes (novox/hq ADR 0006). What a node says it applied is diff --git a/internal/link/receive.go b/internal/link/receive.go index d96276e..962e9d4 100644 --- a/internal/link/receive.go +++ b/internal/link/receive.go @@ -78,19 +78,6 @@ type Control interface { // Took settles the message: acted on, or understood and needing no action. Took() error - // About names what this message is about — a node's report, one module's move, one build's - // outcome — and is said before the store is asked. - // - // A transport that holds messages **in memory** uses it to set aside anything older it is - // holding about the same thing: the older is the past, and letting it come back after the - // newer was acted on would undo the newer. - // - // **This is the one thing holding-in-memory can do that holding-in-the-server cannot**, and - // naming it here rather than hiding it is deliberate. On the bus being built the message - // belongs to the server and comes back whatever happened meanwhile, so this is ignored and the - // digest a report carries answers the same question instead (window.go, design 25 §3). - About(what string) - // Hold keeps the message and asks for it again after the delay — the store window. Hold(after time.Duration) error diff --git a/internal/link/receive_current.go b/internal/link/receive_current.go deleted file mode 100644 index 7272805..0000000 --- a/internal/link/receive_current.go +++ /dev/null @@ -1,321 +0,0 @@ -package link - -import ( - "context" - "errors" - "fmt" - "time" - - amqp "github.com/rabbitmq/amqp091-go" -) - -// The consume side on the bus the mesh runs on today. -// -// Everything here was the serving loop's until the seam went in: the queues, the binds, the -// prefetch, and the list of messages the store could not take yet. It moved rather than changed — -// the behaviour this transport has is the behaviour it had, because the mesh is running on it and -// a bus nothing speaks yet is no reason to alter the one every node is on (ADR 0116). - -// Prefetch is how many messages the bus hands the controller before it has settled them. -// -// More than one because a message the store could not take is held, unsettled, while the loop -// goes on answering others — an enrolment above all, which a host is waiting on (novox/hq issue -// 083). Bounded, because what is held is also what the bus has not kept on its own disk as -// pending. -const Prefetch = 64 - -// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to -// answer — an enrolment above all — while others wait for the store. -const PrefetchHeadroom = 8 - -// TryAgainAfter is how often the held are looked at. A store comes back in seconds, and a report a -// few seconds late is still current. -const TryAgainAfter = 2 * time.Second - -// currentInbound consumes what nodes say over the bus the mesh has. -type currentInbound struct { - conn *amqp.Connection - channel *amqp.Channel - // upgrades and catchups are bound only when something is listening (Also). - upgrades bool - catchups bool - // held is every message the store could not take, by the bus's own delivery tag. Kept here - // rather than in the serving loop because holding a delivery unacknowledged is this - // transport's way of keeping it, and the other's is to hand it back to the server. - held map[uint64]*holding - // again is how often the held are looked at; zero means TryAgainAfter. Set by tests. - again time.Duration -} - -// holding is one message kept for the store, and when to try it again. -type holding struct { - message *currentControl - due time.Time - about string -} - -// Current is the consume side of the bus the mesh runs on today. -func Current(conn *amqp.Connection, channel *amqp.Channel) Inbound { - return ¤tInbound{conn: conn, channel: channel, held: map[uint64]*holding{}} -} - -// Also binds the queue one more kind arrives on. -// -// The kinds nodes publish all share one queue and are bound at Connect, because a node may -// publish any of them and binding one while forgetting another is a message the bus accepts, finds -// no queue for, and drops — the publisher sees success and the consumer sees nothing. The two that -// are events get their own queue each, and only when something is listening. -func (c *currentInbound) Also(kind string) error { - switch kind { - case KindSourceMoved: - // Not followed on the bus the mesh is leaving: the forge's merges are announced on the - // new one, and this transport goes with the move (design 28, task 5.5). - return nil - case KindModuleMoved: - if err := c.bindEvent(UpgradeQueue, KeyModuleUpgraded); err != nil { - return err - } - c.upgrades = true - case KindCatchUp: - if err := c.bindEvent(CatchUpQueue, KeyCatchingUp); err != nil { - return err - } - c.catchups = true - default: - return fmt.Errorf("nothing binds a queue for %s on this bus", kind) - } - return nil -} - -func (c *currentInbound) bindEvent(queue, key string) error { - if _, err := c.channel.QueueDeclare(queue, true, false, false, false, nil); err != nil { - return fmt.Errorf("cannot declare the %s queue: %w", queue, err) - } - if err := c.channel.QueueBind(queue, key, EventsExchange, false, nil); err != nil { - return fmt.Errorf("cannot bind %s to %s/%s: %w", queue, EventsExchange, key, err) - } - return nil -} - -func (c *currentInbound) Close() {} - -// Receive consumes until the context ends. -// -// One consumer per queue, deliberately: with two on one queue the bus would round-robin between -// them and each would receive half of what it expects — a fault this project has already had, -// between a module's daemon and its capability server. -func (c *currentInbound) Receive(ctx context.Context, act func(context.Context, Control)) error { - // A bounded prefetch rather than one. The loop still takes messages one at a time; what the - // prefetch buys is that a message the store could not take can be held while the loop goes on - // to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything - // held goes back to the bus if the controller stops, because nothing held is acknowledged. - if err := c.channel.Qos(Prefetch, 0, false); err != nil { - return err - } - - deliveries, err := c.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane", - false, false, false, false, nil) - if err != nil { - return err - } - - // Its own queue and its own consumer for each event, for the reason above: two consumers on - // one queue split its messages between them, and an upgrade or a catch-up request going to - // whichever half was not listening is a gap that looks like a working mesh. - var upgrades, catchups <-chan amqp.Delivery - if c.upgrades { - upgrades, err = c.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades", - false, false, false, false, nil) - if err != nil { - return err - } - } - if c.catchups { - catchups, err = c.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup", - false, false, false, false, nil) - if err != nil { - return err - } - } - - closed := c.conn.NotifyClose(make(chan *amqp.Error, 1)) - - again := c.again - if again == 0 { - again = TryAgainAfter - } - ticker := time.NewTicker(again) - defer ticker.Stop() - - for { - select { - case <-ctx.Done(): - return nil - case <-ticker.C: - if ctx.Err() != nil { - return nil - } - c.retryHeld(ctx, act) - case delivery, ok := <-catchups: - if !ok { - if catchups != nil { - return errors.New("the bus stopped delivering catch-up requests") - } - continue - } - act(ctx, c.wrap(KindCatchUp, delivery)) - case delivery, ok := <-upgrades: - // A nil channel blocks for ever, so this case simply never fires when nothing is - // listening for upgrades. Closed is different, and means the bus stopped. - if !ok { - if upgrades != nil { - return errors.New("the bus stopped delivering upgrades") - } - continue - } - act(ctx, c.wrap(KindModuleMoved, delivery)) - case reason := <-closed: - // Said rather than returned quietly. A controller whose bus connection dropped is a - // mesh where nothing can be told anything, and the reason is the first thing anybody - // will want. - return fmt.Errorf("the bus connection closed: %v", reason) - case delivery, ok := <-deliveries: - if !ok { - return errors.New("the bus stopped delivering") - } - kind, known := kindOfKey[delivery.RoutingKey] - if !known { - // Rejected without requeue: a message nothing understands will not be understood - // on the next attempt either, and requeuing it would spin. - _ = delivery.Reject(false) - continue - } - act(ctx, c.wrap(kind, delivery)) - } - } -} - -// kindOfKey is how this transport's addressing becomes what the mesh calls a message. -var kindOfKey = map[string]string{ - KeyEnrol: KindEnrolment, - KeyReport: KindReport, - KeyAlive: KindHeartbeat, - KeyBuilt: KindBuilt, - KeyModuleUpgraded: KindModuleMoved, - KeyCatchingUp: KindCatchUp, -} - -func (c *currentInbound) wrap(kind string, delivery amqp.Delivery) *currentControl { - return ¤tControl{kind: kind, delivery: delivery, on: c} -} - -// retryHeld hands every message whose delay has passed back to the loop. Each handler holds it -// again, settles it, or lets it go past the bound. -func (c *currentInbound) retryHeld(ctx context.Context, act func(context.Context, Control)) { - now := time.Now() - due := make([]*currentControl, 0, len(c.held)) - for _, h := range c.held { - if !h.due.After(now) { - due = append(due, h.message) - } - } - for _, m := range due { - if ctx.Err() != nil { - return - } - act(ctx, m) - } -} - -// currentControl is one delivery from the bus the mesh has, as the controller reads it. -type currentControl struct { - kind string - delivery amqp.Delivery - on *currentInbound - // about is what this message is about, as the handler named it; empty until it does. - about string - // first is when this message was first held for the store; zero while it has not been. - first time.Time -} - -func (m *currentControl) Kind() string { return m.kind } -func (m *currentControl) Body() []byte { return m.delivery.Body } -func (m *currentControl) Redelivered() bool { return m.delivery.Redelivered } - -func (m *currentControl) HeldFor() time.Duration { - if m.first.IsZero() { - return 0 - } - return time.Since(m.first) -} - -// Answer publishes to the reply queue the request named. -func (m *currentControl) Answer(ctx context.Context, body []byte) error { - if m.delivery.ReplyTo == "" { - return errors.New("that request named no reply queue, so nothing can be told the answer") - } - return m.on.channel.PublishWithContext(ctx, "", m.delivery.ReplyTo, false, false, - amqp.Publishing{ - ContentType: "application/json", - CorrelationId: m.delivery.CorrelationId, - Body: body, - }) -} - -func (m *currentControl) Took() error { - m.forget() - return m.delivery.Ack(false) -} - -// Drop rejects without requeue: on this bus that is what "understood, and not worth another -// attempt" is spelled as, and it is what feeds a dead-letter queue where one is configured. -func (m *currentControl) Drop() error { - m.forget() - return m.delivery.Reject(false) -} - -// About names what this message is about, and lets go of whatever is held about the same thing: -// the held one is the past, and acting on it after this one would undo this one. Acknowledged -// rather than left to come back, because a held message nothing will act on is a place in the -// prefetch nothing gets back. -func (m *currentControl) About(what string) { - m.about = what - if what == "" { - return - } - for tag, h := range m.on.held { - if h.about != what || tag == m.delivery.DeliveryTag { - continue - } - delete(m.on.held, tag) - _ = h.message.delivery.Ack(false) - } -} - -// Hold keeps the message unacknowledged and sets it aside to be handed back after the delay. -// -// Held no further than the prefetch leaves room: past that the bus would hand the loop nothing -// new — enrolments included — until something held was let go. A message that cannot be held says -// so, and the handler settles it its own way. -func (m *currentControl) Hold(after time.Duration) error { - if m.on.held == nil { - m.on.held = map[uint64]*holding{} - } - if _, already := m.on.held[m.delivery.DeliveryTag]; !already { - if len(m.on.held) >= Prefetch-PrefetchHeadroom { - return fmt.Errorf("%d messages are already held for the store, and holding more "+ - "would stop the queue", len(m.on.held)) - } - m.first = time.Now() - } - m.on.held[m.delivery.DeliveryTag] = &holding{ - message: m, due: time.Now().Add(after), about: m.about, - } - return nil -} - -func (m *currentControl) forget() { - if m.on != nil { - delete(m.on.held, m.delivery.DeliveryTag) - } -} diff --git a/internal/link/receive_current_test.go b/internal/link/receive_current_test.go deleted file mode 100644 index 462e3f3..0000000 --- a/internal/link/receive_current_test.go +++ /dev/null @@ -1,71 +0,0 @@ -package link - -import ( - "context" - "encoding/json" - "io" - "log" - "testing" - "time" - - amqp "github.com/rabbitmq/amqp091-go" -) - -// The harness for the consume side on the bus the mesh runs on today. -// -// Messages arrive through the seam, so what these tests exercise is the controller's decision -// about a message and this transport's way of keeping one — which is what the seam separated. A -// fake acknowledger stands in for the bus, because what is asserted is how a message was settled -// and that needs no server. - -// settled is how the bus was told to settle one message. -type settled struct{ acked, nacked, requeued, rejected bool } - -func (a *settled) Ack(uint64, bool) error { a.acked = true; return nil } -func (a *settled) Nack(_ uint64, _ bool, requeue bool) error { - a.nacked, a.requeued = true, requeue - return nil -} -func (a *settled) Reject(uint64, bool) error { a.rejected = true; return nil } - -// unsettled is a message the controller has neither taken nor let go: it is held, and the bus will -// hand it to whatever consumes next if the controller stops. -func (a *settled) unsettled() bool { return !a.acked && !a.nacked && !a.rejected } - -var tag uint64 - -func quiet() *log.Logger { return log.New(io.Discard, "", 0) } - -// serving is a controller with nothing but a way of receiving, ready for a listener, a recorder, -// an upgrader or a replayer to be set on it. -func serving() (*Server, *currentInbound) { - in := ¤tInbound{held: map[uint64]*holding{}} - return &Server{inbound: in, bus: OverCurrent{}, log: quiet()}, in -} - -// sends is one message arriving over this transport, as the controller reads it. -func (c *currentInbound) sends(t *testing.T, to *settled, kind string, v any) Control { - t.Helper() - body, err := json.Marshal(v) - if err != nil { - t.Fatal(err) - } - tag++ - return ¤tControl{kind: kind, on: c, delivery: amqp.Delivery{ - Acknowledger: to, Body: body, DeliveryTag: tag, - }} -} - -// dueNow brings every held message forward, so a test need not wait out the backoff a real store -// restart would be given (RedeliverAfter). -func (c *currentInbound) dueNow() { - for _, h := range c.held { - h.due = time.Now().Add(-time.Second) - } -} - -// retries hands every held message back to the controller, the way the ticker does. -func (c *currentInbound) retries(ctx context.Context, s *Server) { - c.dueNow() - c.retryHeld(ctx, s.act) -} diff --git a/internal/link/receive_fake_test.go b/internal/link/receive_fake_test.go new file mode 100644 index 0000000..ad71890 --- /dev/null +++ b/internal/link/receive_fake_test.go @@ -0,0 +1,94 @@ +package link + +import ( + "context" + "encoding/json" + "testing" + "time" +) + +// A harness for the controller's decisions about a message, with no bus behind it. +// +// What these tests exercise is the server's verdict — taken, dropped, held for the store, let go +// once the store has been away too long — and the transport's part in that is only to keep a held +// message and hand it back. A fake that does exactly that stands in for the bus, so what is +// asserted is how a message was settled and the decision needs no server. + +// settled is how the bus was told to settle one message. +type settled struct{ acked, nacked, requeued, rejected bool } + +// unsettled is a message the controller has neither taken nor let go: it is held, and the bus will +// hand it to whatever consumes next if the controller stops. +func (a *settled) unsettled() bool { return !a.acked && !a.nacked && !a.rejected } + +// fakeInbound keeps the messages the controller held, the way a stream would. +type fakeInbound struct{ held map[uint64]*fakeControl } + +var tag uint64 + +// serving is a controller with nothing but a way of receiving, ready for a listener, a recorder, +// an upgrader or a replayer to be set on it. +func serving() (*Server, *fakeInbound) { + in := &fakeInbound{held: map[uint64]*fakeControl{}} + return &Server{inbound: in, log: quiet()}, in +} + +func (c *fakeInbound) Also(string) error { return nil } +func (c *fakeInbound) Close() {} +func (c *fakeInbound) Receive(context.Context, func(context.Context, Control)) error { + return nil +} + +// sends is one message arriving, as the controller reads it. +func (c *fakeInbound) sends(t *testing.T, to *settled, kind string, v any) Control { + t.Helper() + body, err := json.Marshal(v) + if err != nil { + t.Fatal(err) + } + tag++ + return &fakeControl{kind: kind, body: body, tag: tag, to: to, on: c} +} + +// retries hands every held message back to the controller, the way a stream redelivers. +func (c *fakeInbound) retries(ctx context.Context, s *Server) { + for tag, m := range c.held { + delete(c.held, tag) + m.redelivered = true + s.act(ctx, m) + } +} + +type fakeControl struct { + kind string + body []byte + tag uint64 + to *settled + on *fakeInbound + redelivered bool + since time.Time +} + +func (m *fakeControl) Kind() string { return m.kind } +func (m *fakeControl) Body() []byte { return m.body } +func (m *fakeControl) Redelivered() bool { return m.redelivered } +func (m *fakeControl) About(string) {} +func (m *fakeControl) Answer(context.Context, []byte) error { return nil } +func (m *fakeControl) Took() error { m.to.acked = true; return nil } +func (m *fakeControl) Drop() error { m.to.rejected = true; return nil } + +// HeldFor is how long the store has been waited on for this message — zero on a first delivery. +func (m *fakeControl) HeldFor() time.Duration { + if m.since.IsZero() { + return 0 + } + return time.Since(m.since) +} + +func (m *fakeControl) Hold(time.Duration) error { + if m.since.IsZero() { + m.since = time.Now() + } + m.on.held[m.tag] = m + return nil +} diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index 967d19f..a4b21af 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -23,6 +23,10 @@ import ( // it first could not take it, and a controller that restarts mid-window has nothing to lose. // natsInbound consumes what nodes and modules say over NATS. +// Prefetch is how many controls the controller holds unacknowledged at once: the store window +// (ADR 0083) is the server's, and this is what it may hand this process ahead of its acting. +const Prefetch = 64 + type natsInbound struct { js *broker.JetStream // follows is the kinds asked for beyond what nodes say (Also). The events those are are the @@ -222,14 +226,6 @@ func (m *natsControl) HeldFor() time.Duration { return time.Since(first) } -// About is nothing here, and that is the point. -// -// Setting a held message aside when a newer one about the same thing arrives is what a controller -// holding deliveries in memory can do. A naked message belongs to the server and comes back -// whatever happened meanwhile, so the question "is this the past?" is answered by what the message -// says instead — the digest of the declaration a report is about (window.go, design 25 §3). -func (m *natsControl) About(string) {} - // Answer publishes to the reply subject the request carries **in its payload**. // // Not `Respond`, and not the message's reply field: a message a JetStream consumer delivers has had diff --git a/internal/link/receive_nats_test.go b/internal/link/receive_nats_test.go index e7b9a59..a10ed34 100644 --- a/internal/link/receive_nats_test.go +++ b/internal/link/receive_nats_test.go @@ -4,6 +4,8 @@ import ( "context" "encoding/json" "errors" + "io" + "log" "os" "sync" "testing" @@ -120,6 +122,8 @@ func eventually(t *testing.T, what string, is func() bool) { // A report published by a node reaches the controller, is recorded, and is acknowledged — so the // stream does not hold it. A work queue is the check: what is acknowledged leaves it. +func quiet() *log.Logger { return log.New(io.Discard, "", 0) } + func TestNatsAReportIsHeardAndLeavesTheStream(t *testing.T) { js := aBus(t) store := &counted{} diff --git a/internal/link/report_retry_test.go b/internal/link/report_retry_test.go index 7549437..6b46a33 100644 --- a/internal/link/report_retry_test.go +++ b/internal/link/report_retry_test.go @@ -46,24 +46,6 @@ func TestAReportTheStoreCouldNotTakeIsHeldAndOneItRefusedIsNot(t *testing.T) { } } -// A newer report from the same node supersedes one of its reports still held: recorded after the -// newer, the older would overwrite what the node is doing now. -func TestANewerReportSupersedesAHeldOneFromTheSameNode(t *testing.T) { - s, in := serving() - s.listener = heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))} - older, newer, other := &settled{}, &settled{}, &settled{} - s.act(context.Background(), in.sends(t, older, KindReport, aReport("anchor", "d1"))) - s.act(context.Background(), in.sends(t, other, KindReport, aReport("laptop", "d7"))) - s.act(context.Background(), in.sends(t, newer, KindReport, aReport("anchor", "d2"))) - if !older.acked { - t.Fatalf("the older report was not set aside by the newer: %+v", older) - } - if !newer.unsettled() || !other.unsettled() || len(in.held) != 2 { - t.Fatalf("the newer report and another node's were not both held: newer %+v other %+v, %d held", - newer, other, len(in.held)) - } -} - // A store that has not come back within the bound is not restarting: the report is let go, loudly, // rather than held for ever. func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) { diff --git a/internal/link/serve.go b/internal/link/serve.go index 563e23e..63f4e7c 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -11,10 +11,6 @@ import ( "log" "os" "time" - - amqp "github.com/rabbitmq/amqp091-go" - - "github.com/novox/mesh-controller/internal/envfile" ) // AMQPVar is the controller's own connection to the bus the mesh runs on today. @@ -71,8 +67,6 @@ type Upgrader interface { type Server struct { inbound Inbound bus Bus - conn *amqp.Connection - channel *amqp.Channel js *broker.JetStream enroller Enroller @@ -127,85 +121,18 @@ func (s *Server) Answers(r Replayer) error { // On the port MESH_BROKER_AMQP_PORT names when the node's settings moved the broker (novox/hq // 04-ISSUES/102) — the URL is genesis's, sealed, and its port is the one thing in it the node may // have moved since. -func Connect(enroller Enroller, listener Listener) (*Server, error) { - url, err := envfile.Placed(AMQPVar) - if err != nil { - return nil, err - } - if url == "" { - return nil, fmt.Errorf( - "this control plane has no %s, so it cannot reach its broker. Nodes talk to it over "+ - "the broker and nowhere else, so without this it can hold records and answer "+ - "nothing", AMQPVar) - } - - conn, err := amqp.Dial(url) - if err != nil { - // Not quoted back: the URL carries the controller's own bus password. - return nil, fmt.Errorf("cannot reach the broker named in %s: %w", AMQPVar, err) - } - channel, err := conn.Channel() - if err != nil { - conn.Close() - return nil, err - } - - // Declared here rather than assumed. The controller is the only thing that may create them — a - // node's account can write to this exchange and read its own queue, and configure nothing - // else, so a node arriving before the controller has ever run finds nothing and says so, - // rather than quietly creating a topology nobody designed. - if err := channel.ExchangeDeclare(Exchange, "direct", true, false, false, false, nil); err != nil { - conn.Close() - return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err) - } - // The events exchange too. The controller is not the only publisher on it — modules announce - // onto it with their own accounts — but it is the only thing permitted to create it, for the - // same reason it is the only thing permitted to create the direct one. - if err := channel.ExchangeDeclare(EventsExchange, "topic", true, false, false, false, nil); err != nil { - conn.Close() - return nil, fmt.Errorf("cannot declare the %s exchange: %w", EventsExchange, err) - } - if _, err := channel.QueueDeclare(ControlQueue, true, false, false, false, nil); err != nil { - conn.Close() - return nil, fmt.Errorf("cannot declare the %s queue: %w", ControlQueue, err) - } - // Every key a node may publish. Binding one and forgetting another is a message the bus - // accepts, finds no queue for, and drops — the publisher sees success and the consumer sees - // nothing. That is exactly what happened to reports: `report` was left unbound while `enrol` - // worked, so nodes announced what they had applied into a void for an afternoon. - for _, key := range []string{KeyEnrol, KeyReport, KeyAlive, KeyBuilt} { - if err := channel.QueueBind(ControlQueue, key, Exchange, false, nil); err != nil { - conn.Close() - return nil, fmt.Errorf("cannot bind %s to %s/%s: %w", ControlQueue, Exchange, key, err) - } - } - - return &Server{ - inbound: Current(conn, channel), - bus: OverCurrent{Channel: channel}, - conn: conn, - channel: channel, - enroller: enroller, - listener: listener, - log: newLog(), - }, nil +// Connect is ConnectNats: the mesh has one bus (novox/hq ADR 0131, design 28 task 5.5). Kept as +// the name callers know; the transport it opened before is gone with the bus it spoke to. +func Connect(js *broker.JetStream, enroller Enroller, listener Listener) *Server { + return ConnectNats(js, enroller, listener) } func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) } -// Channel is the controller's channel, for the command line's own publishing. -func (s *Server) Channel() *amqp.Channel { return s.channel } - func (s *Server) Close() { if s.inbound != nil { s.inbound.Close() } - if s.channel != nil { - _ = s.channel.Close() - } - if s.conn != nil { - _ = s.conn.Close() - } if s.js != nil { s.js.Close() } @@ -349,7 +276,6 @@ func (s *Server) reported(ctx context.Context, m Control) { _ = m.Drop() return } - m.About("report " + report.Node) what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared) if s.listener != nil { @@ -482,9 +408,6 @@ func (s *Server) wasBuilt(ctx context.Context, m Control) { _ = m.Drop() return } - // Each build result its own subject: none supersedes another, and recording one twice is - // harmless — the build is kept by its id. - m.About("build " + digest(m.Body())) err := s.recorder.Built(ctx, result) switch s.decide(ctx, m, fmt.Sprintf("a build result from %s", result.On), "", "", err) { case Hold: @@ -516,7 +439,6 @@ func (s *Server) wasBuilt(ctx context.Context, m Control) { // the catalogue next restarts (issue 083). One request stands for all: a newer one supersedes one // still held. func (s *Server) catchingUp(ctx context.Context, m Control) { - m.About("catch-up") if s.replayer == nil { s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay") _ = m.Took() @@ -569,7 +491,6 @@ func (s *Server) moved(ctx context.Context, m Control) { _ = m.Took() return } - m.About("upgrade " + u.Module) if u.Module == "" { s.log.Printf("an upgrade announcement named no module; ignored") _ = m.Took() @@ -620,7 +541,6 @@ func (s *Server) sourceMoved(ctx context.Context, m Control) { _ = m.Took() return } - m.About("merge " + moved.Owner + "/" + moved.Repo + " into " + moved.Base) if moved.Commit == "" || moved.Repo == "" { s.log.Printf("a merge announcement named no repository or no commit; ignored") _ = m.Took() diff --git a/internal/link/store_window_test.go b/internal/link/store_window_test.go index dbe994f..1dbf60b 100644 --- a/internal/link/store_window_test.go +++ b/internal/link/store_window_test.go @@ -3,7 +3,6 @@ package link import ( "context" "errors" - "fmt" "testing" "github.com/jackc/pgx/v5/pgconn" @@ -18,7 +17,8 @@ func (r recordsWith) Built(context.Context, BuildResult) error { return r.err } type upgradesWith struct{ err error } -func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err } +func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err } +func (u upgradesWith) SourceMoved(context.Context, SourceMoved) error { return u.err } type replaysWith struct{ err error } @@ -106,51 +106,3 @@ func TestAMessageHandledDuringShutdownIsLeftForTheBus(t *testing.T) { t.Fatalf("a build result handled during shutdown was settled, and so lost: %+v", *to) } } - -// Two identical build results: the newer sets the older aside rather than leaving it unsettled -// for ever, holding a place in the prefetch. -func TestAnIdenticalBuildResultSetsTheHeldOneAside(t *testing.T) { - s, in := serving() - s.recorder = recordsWith{err: restarting} - built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"} - first, second := &settled{}, &settled{} - s.act(context.Background(), in.sends(t, first, KindBuilt, built)) - s.act(context.Background(), in.sends(t, second, KindBuilt, built)) - if !first.acked || !second.unsettled() || len(in.held) != 1 { - t.Fatalf("an identical build result did not set the held one aside: first %+v second %+v, %d held", - *first, *second, len(in.held)) - } -} - -// What is held stops short of the prefetch, so the loop always has room to answer an enrolment. -func TestWhatIsHeldLeavesRoomInThePrefetch(t *testing.T) { - s, in := serving() - s.recorder = recordsWith{err: restarting} - var last *settled - for i := 0; i < Prefetch; i++ { - last = &settled{} - s.act(context.Background(), in.sends(t, last, KindBuilt, BuildResult{On: "anchor", Commit: fmt.Sprint(i)})) - } - if len(in.held) != Prefetch-PrefetchHeadroom { - t.Fatalf("%d messages were held; the ceiling is %d", len(in.held), Prefetch-PrefetchHeadroom) - } - if last.unsettled() { - t.Fatalf("a message past the ceiling was held: %+v", *last) - } -} - -// An upgrade handled during shutdown is left for the bus too — the upgrader's error is the -// cancelled context, which is no answer about the announcement. -func TestAnUpgradeHandledDuringShutdownIsLeftForTheBus(t *testing.T) { - ctx, cancel := context.WithCancel(context.Background()) - cancel() - s, in := serving() - s.upgrader = upgradesWith{err: context.Canceled} - to := &settled{} - s.act(ctx, in.sends(t, to, KindModuleMoved, Upgraded{Module: "gitea", Commit: "abcdef0123"})) - if !to.unsettled() { - t.Fatalf("an upgrade was settled during shutdown, and so lost: %+v", *to) - } -} - -func (u upgradesWith) SourceMoved(context.Context, SourceMoved) error { return nil }