From e5c2eb20f239d914e93384321e5d46c07b86ce42 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 29 Sep 2026 17:36:50 +0200 Subject: [PATCH 1/4] A token is an account on the bus, and genesis can place the list MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit novox/hq 04-ISSUES/146. The composed user list names an enrolment user for every machine with a live token and nothing minted a credential for it, so the composer left it out as a user with no password — and every enrolment since the mesh moved to this bus was refused before the mesh heard of it. The comment above the issuing code already said the account is created before the token is handed over; now it is. Recorded rather than minted, because the token's secret is the password. And 'broker accounts', which composes the same list the declaration carries and writes it to standard output. For genesis, where no declaration can reach the machine running the bus because that machine is not yet a node. It says what it composed; whoever is raising the machine places it. A control plane that wrote the file itself would have to learn where the bus keeps its configuration and how to make it reload, which is the module's knowledge. --- cmd/mesh-controller/buscertificate.go | 87 ++++++++++++++++++++++++ cmd/mesh-controller/main.go | 2 +- cmd/mesh-controller/nodes.go | 7 +- internal/broker/genesis_template_test.go | 28 +++++--- internal/inventory/busrecords_test.go | 61 +++++++++++++++++ internal/inventory/bususers.go | 44 +++++++++--- internal/inventory/nodes.go | 24 +++++++ 7 files changed, 231 insertions(+), 22 deletions(-) diff --git a/cmd/mesh-controller/buscertificate.go b/cmd/mesh-controller/buscertificate.go index 983a05e..e369ad2 100644 --- a/cmd/mesh-controller/buscertificate.go +++ b/cmd/mesh-controller/buscertificate.go @@ -1,6 +1,7 @@ package main import ( + "context" "crypto/rand" "crypto/rsa" "crypto/x509" @@ -12,7 +13,10 @@ import ( "net" "os" "path/filepath" + "strings" "time" + + "github.com/novox/mesh-controller/internal/broker" ) // The bus's own certificate, made by the mesh rather than borrowed from an image. @@ -169,3 +173,86 @@ func writeBusCertificate(crt, key string) error { busCertificateNames, template.NotAfter.Format(time.RFC3339), crt, key) return nil } + +// busAccounts writes the mesh's composed user list to a file. +// +// **For genesis, where no declaration can deliver it** (novox/hq 04-ISSUES/146). Everywhere else +// the list reaches the machine running the bus as a resource of the module that holds it — which +// requires that machine to be an enrolled node, and at genesis it is not: the first node cannot +// enrol because the account it would enrol with cannot be composed onto a bus it has no declaration +// for. The installer breaks that circle by placing the file itself, once, and the module takes the +// file over from its first push. +// +// The same composition, not a second one: this asks the store for the same records and renders them +// with the same composer the declaration uses. A genesis that hand-wrote an account would be a +// second statement of who may say what, able to disagree with the first. +// +// **It writes to standard output unless told a file**, and that is the point: the control plane +// composes and says what it composed, and whoever is raising the machine puts it where that +// machine's bus reads it. A control plane that wrote into the bus's own directory would have to +// know where that is and how to make the server re-read it — which is the module's knowledge, and +// the module is what takes this over on the first push. +// +// broker accounts > /var/lib/mesh-bus-conf/accounts.conf +func busAccounts(ctx context.Context, args []string) error { + into := "" + for i := 0; i < len(args); i++ { + switch args[i] { + case "--into": + if i+1 >= len(args) { + return errors.New("--into needs a file") + } + into = args[i+1] + i++ + default: + return fmt.Errorf("broker accounts --into : %q", args[i]) + } + } + + open, err := openStores(ctx) + if err != nil { + return err + } + defer open.Close() + + records, err := open.inventory.BusRecords(ctx) + if err != nil { + return err + } + users, err := broker.Users(records) + if err != nil { + return err + } + kept, err := open.inventory.BusUsers(ctx) + if err != nil { + return err + } + hashes := make(map[string]string, len(kept)) + for name, u := range kept { + hashes[name] = u.PasswordHash + } + filled, missing := broker.WithPasswords(users, hashes) + if len(missing) > 0 { + // To standard error, always: the composed file may be going to standard output, and a + // remark in the middle of it is a configuration the server refuses to parse. + fmt.Fprintf(os.Stderr, "leaving out %d user(s) the mesh has minted no credential for: %s\n", + len(missing), strings.Join(missing, ", ")) + } + if len(filled) == 0 { + return errors.New("not one user has a credential, so this list would refuse every " + + "connection in the mesh") + } + accounts, err := broker.ComposeAccounts(filled) + if err != nil { + return err + } + if into == "" { + fmt.Print(accounts) + return nil + } + if err := os.WriteFile(into, []byte(accounts), 0o600); err != nil { + return err + } + fmt.Printf("wrote %d user(s) to %s\n", len(filled), into) + return nil +} diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index 6b0f312..2e32e88 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -88,7 +88,7 @@ func run() error { case "identity": return identityCommand(ctx, args[1:]) case "broker": - return brokerCommand(args[1:]) + return brokerCommand(ctx, args[1:]) case "serve": return serve(ctx) case "upgrade": diff --git a/cmd/mesh-controller/nodes.go b/cmd/mesh-controller/nodes.go index 5a8dd13..3eba5c2 100644 --- a/cmd/mesh-controller/nodes.go +++ b/cmd/mesh-controller/nodes.go @@ -363,12 +363,15 @@ func identityCommand(ctx context.Context, args []string) error { return nil } -func brokerCommand(args []string) error { +func brokerCommand(ctx context.Context, args []string) error { if len(args) > 0 && args[0] == "certificate" { return busCertificate(args[1:]) } + if len(args) > 0 && args[0] == "accounts" { + return busAccounts(ctx, args[1:]) + } if len(args) == 0 || args[0] != "show" { - return errors.New("broker show | broker certificate [--check] --into ") + return errors.New("broker show | broker certificate [--check] --into | broker accounts --into ") } known, err := broker.FromEnvironment() if errors.Is(err, broker.ErrNotConfigured) { diff --git a/internal/broker/genesis_template_test.go b/internal/broker/genesis_template_test.go index 1260f85..d5c3c4e 100644 --- a/internal/broker/genesis_template_test.go +++ b/internal/broker/genesis_template_test.go @@ -62,6 +62,22 @@ func TestTheInstallersFirstUserListIsWhatTheControllerWouldCompose(t *testing.T) // theCarriedAccounts is the accounts file the installer's template writes at genesis. func theCarriedAccounts(t *testing.T) string { + t.Helper() + for _, r := range theTemplate(t) { + if r["id"] == "bus-accounts" { + content, _ := r["content"].(string) + if content == "" { + t.Fatal("the template's accounts file is empty, so the bus would refuse every connection") + } + return content + } + } + t.Fatal("the template carries no accounts file, so a mesh raised from it has a bus nobody may use") + return "" +} + +// theTemplate is the installer's bundle, as resources. +func theTemplate(t *testing.T) []map[string]any { t.Helper() path := filepath.Join("..", "..", "..", "mesh-host", "examples", "foundation-first-node-nats.lock") raw, err := os.ReadFile(path) @@ -81,17 +97,7 @@ func theCarriedAccounts(t *testing.T) string { if err := json.Unmarshal([]byte(strings.Join(lines, "\n")), &bundle); err != nil { t.Fatalf("the template is not readable: %v", err) } - for _, r := range bundle.Resources { - if r["id"] == "bus-accounts" { - content, _ := r["content"].(string) - if content == "" { - t.Fatal("the template's accounts file is empty, so the bus would refuse every connection") - } - return content - } - } - t.Fatal("the template carries no accounts file, so a mesh raised from it has a bus nobody may use") - return "" + return bundle.Resources } // subjectsIn reads one allow-list out of a composed accounts file. diff --git a/internal/inventory/busrecords_test.go b/internal/inventory/busrecords_test.go index 7f61cfe..15cf35d 100644 --- a/internal/inventory/busrecords_test.go +++ b/internal/inventory/busrecords_test.go @@ -8,6 +8,7 @@ import ( "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/catalogue" + "golang.org/x/crypto/bcrypt" ) // Reading the bus's user list out of the mesh's records, against a real store. @@ -164,3 +165,63 @@ func granted(all []string, one string) bool { } return false } + +// **A token is an account on the bus, or it is a string nothing accepts** (novox/hq 04-ISSUES/146). +// +// The composed list names an enrolment user for every machine with a live token, and nothing minted +// a credential for it — so the composer left it out as a user with no password, and every enrolment +// since the mesh moved to this bus was refused by the server before the mesh heard of it. Nothing +// caught it because nothing had enrolled since. +// +// The password cannot be minted, because it is the token's own secret: the machine will present +// exactly that string. So this checks the two halves that make the account usable — that a row +// exists under the name the composer asks for, and that the secret handed out is what that row +// accepts. +func TestIssuingATokenRecordsTheAccountItIsThePasswordOf(t *testing.T) { + inv, ctx := aMeshWith(t) + if _, err := inv.AddNode(ctx, "joiner"); err != nil { + t.Fatal(err) + } + issued, err := inv.IssueToken(ctx, "joiner", time.Hour) + if err != nil { + t.Fatal(err) + } + + name := broker.Principal{Kind: broker.KindEnrolment, Node: "joiner"}.Username() + users, err := inv.BusUsers(ctx) + if err != nil { + t.Fatal(err) + } + user, has := users[name] + if !has { + t.Fatalf("no bus account for %q; the composer would leave the enrolment out and the "+ + "machine would be refused before the mesh heard of it: %v", name, users) + } + if user.Kind != BusEnrolment || user.Node != "joiner" { + t.Errorf("the account is %+v, not this node's enrolment", user) + } + if err := bcrypt.CompareHashAndPassword([]byte(user.PasswordHash), []byte(issued.Secret)); err != nil { + t.Error("the account does not accept the secret the token carries, so presenting the " + + "token would be refused by the server") + } + + // And the composition contains it, which is the thing the server reads. + records, err := inv.BusRecords(ctx) + if err != nil { + t.Fatal(err) + } + derived, err := broker.Users(records) + if err != nil { + t.Fatal(err) + } + hashes := map[string]string{} + for n, u := range users { + hashes[n] = u.PasswordHash + } + _, missing := broker.WithPasswords(derived, hashes) + for _, m := range missing { + if m == name { + t.Fatal("the enrolment user is composed without a password, which is a user nobody can be") + } + } +} diff --git a/internal/inventory/bususers.go b/internal/inventory/bususers.go index c8c983b..eb26704 100644 --- a/internal/inventory/bususers.go +++ b/internal/inventory/bususers.go @@ -51,6 +51,40 @@ const ( // reply, into a module's sealed environment — and the mesh keeps only the hash, so a credential is // never recoverable from the store. A caller that loses it must mint again, which is a rotation and // is meant to feel like one. +// RecordBusPassword records a hash for a password the caller already holds. +// +// **For the one credential the mesh does not choose**: an enrolment token's secret is the password +// of the user that presents it (novox/hq ADR 0004, design 25 §6), so the token cannot be given a +// minted password — it already has one, and the machine will connect with exactly that string. +// Everything else goes through Mint, which chooses and returns the plaintext once. +func (i *Inventory) RecordBusPassword(ctx context.Context, u BusUser, password string) error { + if u.Username == "" || u.Kind == "" { + return errors.New("a bus user needs a username and a kind") + } + if password == "" { + return errors.New("a bus user needs a password") + } + hash, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.DefaultCost) + if err != nil { + return fmt.Errorf("cannot hash a bus password: %w", err) + } + return i.writeBusUser(ctx, u, string(hash)) +} + +// writeBusUser is the row, whoever chose the password. +func (i *Inventory) writeBusUser(ctx context.Context, u BusUser, hash string) error { + if _, err := i.store.Pool().Exec(ctx, + `insert into bus_user (username, kind, node, module, password_hash) + values ($1, $2, $3, $4, $5) + on conflict (username) do update + set kind = excluded.kind, node = excluded.node, module = excluded.module, + password_hash = excluded.password_hash, minted_at = now()`, + u.Username, u.Kind, u.Node, u.Module, hash); err != nil { + return fmt.Errorf("cannot record the bus user %s: %w", u.Username, err) + } + return nil +} + func (i *Inventory) MintBusPassword(ctx context.Context, u BusUser) (string, error) { if u.Username == "" || u.Kind == "" { return "", errors.New("a bus user needs a username and a kind") @@ -69,14 +103,8 @@ func (i *Inventory) MintBusPassword(ctx context.Context, u BusUser) (string, err return "", fmt.Errorf("cannot hash a bus password: %w", err) } - if _, err := i.store.Pool().Exec(ctx, - `insert into bus_user (username, kind, node, module, password_hash) - values ($1, $2, $3, $4, $5) - on conflict (username) do update - set kind = excluded.kind, node = excluded.node, module = excluded.module, - password_hash = excluded.password_hash, minted_at = now()`, - u.Username, u.Kind, u.Node, u.Module, string(hash)); err != nil { - return "", fmt.Errorf("cannot record the bus user %s: %w", u.Username, err) + if err := i.writeBusUser(ctx, u, string(hash)); err != nil { + return "", err } return password, nil } diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index e4dd752..971927b 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -14,6 +14,7 @@ import ( "time" "github.com/jackc/pgx/v5" + "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/store" ) @@ -263,6 +264,29 @@ func (i *Inventory) IssueToken(ctx context.Context, nodeName string, validFor ti return Issued{}, err } + // **And the account that secret is the password of** (novox/hq 04-ISSUES/146). The composed + // user list names an enrolment user for every node with a live token, and nothing minted a + // credential for it — so the composer left it out as a user with no password and every + // enrolment was refused by the server before the mesh heard of it. + // + // Recorded rather than minted: the token's secret IS the password, which is what lets a + // machine's first connection be authenticated by the thing it is enrolling with. It cannot be + // chosen here, because it has already been handed to whoever will present it. + // + // Outside the transaction on purpose. The token is what the mesh promised; a credential that + // the next composition rewrites anyway is not worth failing an issue over, and a token with no + // account is recoverable by issuing another, while an account with no token is a user nobody + // can be. + if err := i.RecordBusPassword(ctx, BusUser{ + Username: broker.Principal{Kind: broker.KindEnrolment, Node: node.Name}.Username(), + Kind: BusEnrolment, + Node: node.Name, + }, secret); err != nil { + return Issued{}, fmt.Errorf( + "the token for %s was issued and the bus account it is the password of was not "+ + "recorded, so this token cannot connect: %w", node.Name, err) + } + return Issued{Node: node, Secret: secret, Expires: expires}, nil } -- 2.54.0 From bc317456073690c1f23a72c195eb5a15ca4ee604 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 29 Sep 2026 17:45:11 +0200 Subject: [PATCH 2/4] make image reads its base from the manifest It was broken and stayed broken: the Dockerfile's fallback base is a Go older than go.mod asks for, so every hand build died at 'go mod download' with 'go.mod requires go >= 1.26.0'. The pipeline never saw it because the pipeline passes the declared base in, so the cost fell entirely on whoever built the image themselves and had to find the digest by hand (novox/hq 04-ISSUES/146). Read from module.json rather than written here as well, so the two cannot disagree, and refused outright if the manifest declares none. --- Makefile | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/Makefile b/Makefile index a41e6d7..7f50f12 100644 --- a/Makefile +++ b/Makefile @@ -27,8 +27,18 @@ build: IMAGE ?= mesh-controller:$(VERSION) DEV_TAG ?= mesh-controller:development +# The base the module declares, read from the manifest rather than written here twice. +# +# **`make image` was broken and stayed broken**, because the Dockerfile's fallback base was a Go +# older than go.mod asks for: every build died at `go mod download` with "go.mod requires go >= +# 1.26.0", and the pipeline never saw it because the pipeline passes the declared base in. Anybody +# building the image by hand hit it and had to find the digest themselves (novox/hq 04-ISSUES/146, +# what it cost). +GO_BASE ?= $(shell python3 -c "import json;print(next(o['image'] for o in json.load(open('module.json'))['build']['on'] if o['arg']=='GO_BASE'))" 2>/dev/null) + image: - docker build --build-arg VERSION=$(VERSION) -t $(IMAGE) -t $(DEV_TAG) . + @test -n "$(GO_BASE)" || { echo "module.json declares no GO_BASE; pass GO_BASE= or fix the manifest"; exit 1; } + docker build --build-arg GO_BASE=$(GO_BASE) --build-arg VERSION=$(VERSION) -t $(IMAGE) -t $(DEV_TAG) . @echo @docker image inspect $(IMAGE) --format 'built {{.RepoTags}} {{.Size}} bytes' @@ -38,7 +48,8 @@ BUILDER_IMAGE ?= mesh-builder:$(VERSION) BUILDER_DEV_TAG ?= mesh-builder:development builder-image: - docker build -f cmd/mesh-builder/Dockerfile -t $(BUILDER_IMAGE) -t $(BUILDER_DEV_TAG) . + @test -n "$(GO_BASE)" || { echo "module.json declares no GO_BASE; pass GO_BASE= or fix the manifest"; exit 1; } + docker build --build-arg GO_BASE=$(GO_BASE) -f cmd/mesh-builder/Dockerfile -t $(BUILDER_IMAGE) -t $(BUILDER_DEV_TAG) . @echo @docker image inspect $(BUILDER_IMAGE) --format 'built {{.RepoTags}} {{.Size}} bytes' -- 2.54.0 From 3fbf658c16cfe91de0774b4478dbc7bee5325a7d Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 29 Sep 2026 21:32:33 +0200 Subject: [PATCH 3/4] Two consumers may share a name; they must not share a delivery subject MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit novox/hq 04-ISSUES/146. A push consumer delivers onto an ordinary subject and everything subscribed to it gets a copy. The controller holds a consumer called 'controller' on CONTROL and another called 'controller' on EVENTS, and both were given _DELIVER.controller — so the one process, holding both subscriptions, acted on every message twice. Measured: one enrolment published, one message in the stream, one delivery, no redelivery, and the controller enrolled the machine twice — the second minting a credential that replaced the one the machine had just been handed, which is why it then reconnected for ever as a user whose password the mesh had rotated. Every report and every followed event doubled the same way, silently. The stream goes in the subject because the pair is what identifies a consumer. A subscriber's permission gains the same shape, keeping the bare name so an existing consumer keeps working until the next assertion moves it. --- internal/broker/jetstream.go | 15 +++++++++++++- internal/broker/nats.go | 10 ++++++++-- internal/broker/streams.go | 18 +++++++++++++++++ internal/broker/streams_test.go | 27 ++++++++++++++++++++++++++ internal/broker/testdata/composed.conf | 4 ++-- 5 files changed, 69 insertions(+), 5 deletions(-) diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 4df1a4a..5bbee69 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -184,7 +184,20 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { // without the other is refused by the server with a message that does not say which half is // missing. if c.Queue != "" || c.Push { - want.DeliverSubject = "_DELIVER." + c.Name + // **Per consumer, which means per stream as well as per name** (novox/hq 04-ISSUES/146). + // A push consumer delivers onto an ordinary subject, and everything subscribed to that + // subject gets a copy. The controller holds a consumer called `controller` on CONTROL and + // another called `controller` on EVENTS, and both were given `_DELIVER.controller` — so the + // one process, holding both subscriptions, acted on every message twice. It enrolled a + // joining machine twice from one request, minting a second credential that replaced the one + // the machine had just been given; the same doubling applied to every report and every + // event the controller follows. + // + // The stream is in the name because the pair is what identifies a consumer — the server + // scopes a durable's name to its stream, and this subject is the only place that scoping + // was dropped. Already within what the controller may subscribe (`_DELIVER.controller.>`), + // so no permission moves. + want.DeliverSubject = DeliverSubjectFor(c) } switch _, err := j.js.ConsumerInfo(c.Stream, c.Name); { diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 7d5a913..8e2c8dc 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -269,7 +269,13 @@ func PermissionsFor(p Principal) (Permissions, error) { "mesh.control." + p.Node + ".>", "$JS.API.CONSUMER.INFO.NODES." + p.Node, } - sub = []string{"mesh.node." + p.Node + ".declare", "_DELIVER." + p.Node} + // The deliver subject carries the stream as well as the consumer's name, so what a + // subscriber is permitted has to carry it too (novox/hq 04-ISSUES/146). The bare name + // stays: an existing consumer keeps delivering where it always did until the controller's + // next assertion moves it, and a permission that only allowed the new shape would refuse + // every node in the mesh for exactly as long as that took. + sub = []string{"mesh.node." + p.Node + ".declare", + "_DELIVER." + p.Node, "_DELIVER." + p.Node + ".>"} case KindModule: // 1. Its own namespace: it publishes its events there and serves its tools there. Nothing @@ -323,7 +329,7 @@ func PermissionsFor(p Principal) (Permissions, error) { // take work over the new bus was refused the asking (2026-09-28). worker := "SEAT_" + upperSnake(s.Name) + "_worker" stream := seatStreamName(s.Name) - sub = append(sub, "_DELIVER."+worker) + sub = append(sub, "_DELIVER."+worker, "_DELIVER."+worker+".>") pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>") for _, a := range s.Accepts { sub = append(sub, seatSubject(s, "accept", a)) diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 6408db5..f6dcd2d 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -94,6 +94,24 @@ func MeshStreams() []Stream { } } +// DeliverSubjectFor is where a push consumer's messages land. +// +// **Per consumer, which means per stream as well as per name** (novox/hq 04-ISSUES/146). A push +// consumer delivers onto an ordinary subject, and everything subscribed to that subject gets a +// copy. The controller holds a consumer called `controller` on CONTROL and another called +// `controller` on EVENTS; while both were given `_DELIVER.controller`, the one process holding +// both subscriptions acted on every message twice — a joining machine was enrolled twice from one +// request, and the second enrolment minted a credential that replaced the one the machine had just +// been handed. Every report and every followed event doubled the same way, silently: nothing is +// redelivered, no count is wrong, the work simply happens twice. +// +// The stream belongs in it because the pair is what identifies a consumer — the server scopes a +// durable's name to its stream, and this subject was the one place that scoping was dropped. It +// stays inside what a controller may already subscribe (`_DELIVER.controller.>`). +func DeliverSubjectFor(c Consumer) string { + return "_DELIVER." + c.Name + "." + c.Stream +} + // An Asserter is the part of a JetStream connection stream assertion needs. Narrow on purpose: it // keeps this testable without a server, and keeps the client library out of everything that only // wants to know what the streams are. diff --git a/internal/broker/streams_test.go b/internal/broker/streams_test.go index ca450a3..9a49289 100644 --- a/internal/broker/streams_test.go +++ b/internal/broker/streams_test.go @@ -233,3 +233,30 @@ func containsStep(steps []string, want string) bool { } return false } + +// **Two consumers may share a name, and must not share a delivery subject** (novox/hq +// 04-ISSUES/146). +// +// A push consumer delivers onto an ordinary subject and everything subscribed to it gets a copy. +// The controller holds a consumer called `controller` on CONTROL and another called `controller` on +// EVENTS; while both were given `_DELIVER.controller`, the one process holding both subscriptions +// acted on every message twice — a joining machine enrolled twice from one request, with the second +// enrolment minting a credential that replaced the one the machine had just been handed. +// +// Checked here rather than against a server because it is a property of what the mesh asks for, and +// because the failure it produces is silent: every count is right, nothing is redelivered, and the +// work simply happens twice. +func TestNoTwoConsumersDeliverOntoTheSameSubject(t *testing.T) { + seen := map[string]string{} + for _, c := range MeshConsumers() { + if !c.Push && c.Queue == "" { + continue + } + subject := DeliverSubjectFor(c) + if other, taken := seen[subject]; taken { + t.Errorf("%s on %s and %s deliver onto %s, so whoever holds both acts on every "+ + "message twice", c.Name, c.Stream, other, subject) + } + seen[subject] = c.Name + " on " + c.Stream + } +} diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 03cf35a..faebd3b 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -34,11 +34,11 @@ accounts { } } { user: "node.one", password: "$2a$11$nnnnnnnnnnnnnnnnnnnnnn", permissions: { publish: { allow: ["$JS.ACK.NODES.one.>", "$JS.API.CONSUMER.INFO.NODES.one", "mesh.control.one.>"] } - subscribe: { allow: ["_DELIVER.one", "_INBOX.node.one.>", "mesh.node.one.declare"] } + subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] } } } { user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: { publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } - subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } + subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: { -- 2.54.0 From 775df79893b4bc97a482f9abc8b7a6aa6ed35666 Mon Sep 17 00:00:00 2001 From: jochen Date: Tue, 29 Sep 2026 23:26:00 +0200 Subject: [PATCH 4/4] A machine the mesh could not read is not a machine that runs nothing MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three gatherers walk every node and pass over one whose plan will not compose, so that one broken set does not cost the rest. They read a plain error to mean that, and so read a store that was briefly unreachable as a machine running nothing. On the roster of routed names that is not a degraded answer but a false one: it states to every machine at once that another machine's names do not exist. Because the roster is part of every container's identity, a control node replaced every container it ran — its own store, the registry, the edge, mail, the bus — on a six-minute cycle for hours. The loop closed through the store this is read from: each pass restarted it, the read failed, one name left the roster, and the roster changing is every container changing. planFor now marks the two failures that really are the node's own — its set not composing, and a setting that reaches nothing — and the three gatherers pass over those and only those. Every other failure is raised, naming the machine and the read, because a mesh-wide refusal with nothing named in it is the other way to lose an evening. novox/hq 04-ISSUES/152, and 151 for why a changed roster is a changed container. --- cmd/mesh-controller/network.go | 9 +- cmd/mesh-controller/plan.go | 70 +++++++++++++-- cmd/mesh-controller/roster_failure_test.go | 99 ++++++++++++++++++++++ 3 files changed, 170 insertions(+), 8 deletions(-) create mode 100644 cmd/mesh-controller/roster_failure_test.go diff --git a/cmd/mesh-controller/network.go b/cmd/mesh-controller/network.go index fb93047..dcfccae 100644 --- a/cmd/mesh-controller/network.go +++ b/cmd/mesh-controller/network.go @@ -346,9 +346,16 @@ func whoResolves(ctx context.Context, open *stores, requirement string) ( refused := map[string]string{} for _, n := range nodes { plan, _, err := planFor(ctx, open, n.Name) - if err != nil { + switch { + case unresolvable(err): refused[n.Name] = err.Error() continue + case err != nil: + // Not a node that does not resolve — a question that went unanswered. Recording it as a + // refusal would take the machine off the private network, and the generator that reads + // this would then write a roster and a filter without it (novox/hq 04-ISSUES/152). + return nil, nil, fmt.Errorf("whether %s answers %q cannot be read: %w", + n.Name, requirement, err) } for _, m := range plan.Modules { for _, offered := range m.Offers() { diff --git a/cmd/mesh-controller/plan.go b/cmd/mesh-controller/plan.go index c6f174d..a32a4a2 100644 --- a/cmd/mesh-controller/plan.go +++ b/cmd/mesh-controller/plan.go @@ -25,7 +25,36 @@ import ( // cheapest next step. That is how novox/hq ADR 0001 records `hal/sdk` reaching 34,636: // nothing in it was wrong, and no one edit was the one that should have been a new file. +// notResolvable marks the one failure in planFor that is a statement about the node: its assigned +// modules do not compose. Every other failure means the mesh could not be *asked* — the store was +// unreachable, a key could not be read — and says nothing about the node at all. +// +// The distinction exists because three callers gather something across every machine and must carry +// on when one machine's set is broken. Each of them read a plain error as "their set does not +// resolve", and so read a store that was briefly unreachable as a machine that runs nothing. On the +// roster of routed names that is not a degraded answer but a false one: it states, to every machine +// at once, that another machine's names do not exist. A control node spent hours replacing every +// container it ran, on a six-minute cycle, because each pass restarted the store this is read from, +// the read failed, one name left the roster, and the roster is part of every container's identity +// (novox/hq 04-ISSUES/152, and 04-ISSUES/151 for why a changed roster is a changed container). +// +// So: skip a node that cannot resolve, and never a node that could not be read. +type notResolvable struct{ err error } + +func (n notResolvable) Error() string { return n.err.Error() } +func (n notResolvable) Unwrap() error { return n.err } + +// unresolvable reports whether err is a node's own set failing to compose, rather than the mesh +// being unable to answer. +func unresolvable(err error) bool { + var n notResolvable + return errors.As(err, &n) +} + // planFor works out everything a node should run, from what was assigned to it. +// +// A failure to compose the node's own modules is wrapped as notResolvable; every other failure is +// returned as it is. Callers gathering across the mesh must tell them apart — see notResolvable. func planFor(ctx context.Context, open *stores, nodeName string) (catalogue.Resolution, catalogue.SettingsBy, error) { inv := open.inventory shelf, err := inv.Catalogue(ctx) @@ -95,7 +124,9 @@ func planFor(ctx context.Context, open *stores, nodeName string) (catalogue.Reso At: onNetwork[nodeName], PublicDomain: publicDomain, Account: who.Account, AccountHome: who.AccountHome}, world) if err != nil { - return catalogue.Resolution{}, nil, err + // The node's own set does not compose. Marked, because this is the only failure here that + // a mesh-wide gatherer may pass over — see notResolvable. + return catalogue.Resolution{}, nil, notResolvable{err} } // The credential for each thing this node takes from elsewhere. Made once and kept, so the @@ -163,8 +194,13 @@ func planFor(ctx context.Context, open *stores, nodeName string) (catalogue.Reso if len(stray) > 0 { // Somebody set something that reaches no file. Said here rather than discovered by the // machine not behaving differently, which is the slowest way there is. - return catalogue.Resolution{}, nil, fmt.Errorf( - "these settings reach nothing:\n - %s", strings.Join(stray, "\n - ")) + // + // Marked like a set that will not compose, and for the same reason: it is a standing fact + // about this node's own configuration, not a question the mesh could not answer. A gatherer + // passes over it as it always did — one node's stray setting must not stop every other node + // being described (novox/hq 04-ISSUES/152). + return catalogue.Resolution{}, nil, notResolvable{fmt.Errorf( + "these settings reach nothing:\n - %s", strings.Join(stray, "\n - "))} } return resolved, settings, nil } @@ -671,11 +707,17 @@ func renderingFor(ctx context.Context, open *stores, node string, // routed name only because it carried a label the mesh composed, never because the mesh knows what // "route" means. A node that does not resolve is skipped, so one machine's broken set does not cost // the rest their names. +// +// **A node that could not be READ is a different matter and is raised.** Skipping one states, to +// every machine at once, that its names do not exist — and since the roster is part of every +// container's identity, that withdraws them and replaces every container (novox/hq 04-ISSUES/152, +// 151). So every failure here says which machine and which read, because the alternative is a +// mesh-wide refusal with nothing named in it. func routeNamesInTheMesh(ctx context.Context, open *stores) (map[string]string, error) { inv := open.inventory places, err := inv.Overlays(ctx) if err != nil { - return nil, err + return nil, fmt.Errorf("where the machines are cannot be read: %w", err) } address := map[string]string{} for _, p := range places { @@ -686,14 +728,22 @@ func routeNamesInTheMesh(ctx context.Context, open *stores) (map[string]string, nodes, err := inv.Nodes(ctx) if err != nil { - return nil, err + return nil, fmt.Errorf("which machines the mesh has cannot be read: %w", err) } out := map[string]string{} for _, n := range nodes { plan, settings, err := planFor(ctx, open, n.Name) - if err != nil { + switch { + case unresolvable(err): + // Their set does not compose, so they serve no names. Passed over, so one machine's + // broken set does not cost the rest theirs. continue + case err != nil: + // The mesh could not be asked. Returning the roster without this machine's names would + // state that they do not exist — to every machine, and indistinguishably from the + // operator having withdrawn them (novox/hq 04-ISSUES/152). + return nil, fmt.Errorf("the names %s serves cannot be read: %w", n.Name, err) } for _, m := range plan.Modules { for to := range m.Contributes { @@ -823,11 +873,17 @@ func grantsFor(ctx context.Context, open *stores, node string) ([]catalogue.Gran out := make([]catalogue.Grant, 0, len(issued)) for _, s := range issued { plan, settings, err := planFor(ctx, open, s.Consumer) - if err != nil { + switch { + case unresolvable(err): // Their set does not resolve. Skipped rather than fatal: this node is not the place // to report another machine's problem, and a grant for something that is not going to // run would have the provider create a user nothing uses. continue + case err != nil: + // The mesh could not be asked what they wanted, which is not the same as their wanting + // nothing — and withholding a grant on that reading takes a consumer's access away + // (novox/hq 04-ISSUES/152). + return nil, fmt.Errorf("what %s asked of %s cannot be read: %w", s.Consumer, s.Name, err) } values, asks, err := plan.ContributionsFrom(s.Name, s.ConsumerModule, settings) if err != nil { diff --git a/cmd/mesh-controller/roster_failure_test.go b/cmd/mesh-controller/roster_failure_test.go new file mode 100644 index 0000000..82b87ac --- /dev/null +++ b/cmd/mesh-controller/roster_failure_test.go @@ -0,0 +1,99 @@ +package main + +import ( + "context" + "strings" + "testing" +) + +// A node's own set failing to compose, and the mesh being unable to answer at all, are different +// things, and only the first may be passed over when something is gathered across every machine +// (novox/hq 04-ISSUES/152). These pin that distinction where the three gatherers rely on it. + +func TestASetThatDoesNotComposeIsMarkedAsTheNodesOwnProblem(t *testing.T) { + open := aMesh(t) + one, two := rivals() + register(t, open, one) + register(t, open, two) + for _, m := range []string{one.Module, two.Module} { + if _, err := open.inventory.Assign(t.Context(), "laptop", m); err != nil { + t.Fatal(err) + } + } + + _, _, err := planFor(t.Context(), open, "laptop") + if err == nil { + t.Fatal("two modules claiming one seat composed anyway") + } + if !unresolvable(err) { + t.Fatalf("a set that cannot compose was not marked as the node's own problem: %v", err) + } +} + +func TestAStoreThatCannotBeReadIsNotANodeThatDoesNotCompose(t *testing.T) { + open := aMesh(t) + + // Nothing is wrong with anchor. The question simply cannot be asked. + stopped, cancel := context.WithCancel(t.Context()) + cancel() + + _, _, err := planFor(stopped, open, "anchor") + if err == nil { + t.Fatal("a plan composed against a store that could not be read") + } + if unresolvable(err) { + t.Fatalf("a question the mesh could not answer was read as a node that runs nothing: %v", err) + } +} + +func TestOneIncoherentNodeDoesNotCostTheRestTheirNames(t *testing.T) { + open := aMesh(t) + one, two := rivals() + register(t, open, one) + register(t, open, two) + for _, m := range []string{one.Module, two.Module} { + if _, err := open.inventory.Assign(t.Context(), "laptop", m); err != nil { + t.Fatal(err) + } + } + + // laptop cannot compose. That is laptop's problem and nobody else's: the roster is still + // answerable, and anchor keeps whatever it serves. + if _, err := routeNamesInTheMesh(t.Context(), open); err != nil { + t.Fatalf("one node's broken set cost the whole mesh its roster: %v", err) + } +} + +func TestARosterIsNeverReturnedWithNamesItCouldNotRead(t *testing.T) { + open := aMesh(t) + + stopped, cancel := context.WithCancel(t.Context()) + cancel() + + names, err := routeNamesInTheMesh(stopped, open) + if err == nil { + t.Fatalf("a roster was composed from a store that could not be read: %v", names) + } + // The failure must be raised, not turned into an absence. A roster missing a machine's names + // is indistinguishable, on every machine that receives it, from the operator withdrawing them — + // and because the roster is part of every container's identity, it replaces all of them. + if names != nil { + t.Fatalf("a partial roster was returned beside the error: %v", names) + } +} + +// Kept so the reason survives the next person reading it: the message the gatherer raises must say +// which machine could not be read, or the operator is left with a mesh-wide failure and no name. +func TestTheRaisedFailureNamesTheMachineItCouldNotRead(t *testing.T) { + open := aMesh(t) + stopped, cancel := context.WithCancel(t.Context()) + cancel() + + _, err := routeNamesInTheMesh(stopped, open) + if err == nil { + t.Fatal("no failure was raised") + } + if !strings.Contains(err.Error(), "cannot be read") { + t.Fatalf("the failure does not say the mesh could not be read: %v", err) + } +} -- 2.54.0