From a3e8a4185ba0b28176864d4abd726bfa20e1e238 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 04:01:17 +0200 Subject: [PATCH 1/4] An assignment issues its bus credential, a push refuses one nobody issued, and what reads a secret restarts on it (hq issue 203) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `assign` recorded a module and `push` sealed a random own secret where its bus credential belongs; the process crash-looped until a person ran `module issue` and pushed again, and the only warning was one line in a list printed on every push. Now assigning a module that declares a broker secret issues the credential in the same act — kept when one exists, so re-assigning rotates nothing — and when the bus cannot be reached from here the assignment says which verb to run. A push never seals a placeholder in a credential's place: a module whose bus user is unminted is refused by name, with the verb. The control plane's own user is the installer's, seeded at genesis, which the test now says. And what reads one of a module's own secrets is restarted when it changes — composed for a container or daemon that names the secret's path in its volumes, environment or env-files, so a manifest need not say it: the build machine ran on an hour-old credential because its manifest restarted it on its environment file alone (issue 206). A scheduled or run-once process is left alone; it reads afresh. --- cmd/mesh-controller/acts.go | 38 ++++++++ cmd/mesh-controller/addresses_test.go | 6 ++ .../credential_on_assign_test.go | 90 ++++++++++++++++++ cmd/mesh-controller/plan.go | 17 ++++ internal/catalogue/declaration.go | 93 +++++++++++++++++++ internal/catalogue/secret_restart_test.go | 52 +++++++++++ 6 files changed, 296 insertions(+) create mode 100644 cmd/mesh-controller/credential_on_assign_test.go create mode 100644 internal/catalogue/secret_restart_test.go diff --git a/cmd/mesh-controller/acts.go b/cmd/mesh-controller/acts.go index 9644a76..e51e99c 100644 --- a/cmd/mesh-controller/acts.go +++ b/cmd/mesh-controller/acts.go @@ -3,6 +3,7 @@ package main import ( "context" "fmt" + "github.com/novox/mesh-controller/internal/broker" "sort" "strings" @@ -67,6 +68,13 @@ func assign(ctx context.Context, open *stores, node, module string) (string, err for _, line := range settled { said += "\n " + line } + // Its bus credential, in the same act (novox/hq issue 203): an assignment pushed before its + // credential exists delivers a process that cannot authenticate and crash-loops until somebody + // runs a second verb and a second push. Issued here when the module speaks on the bus and has + // no credential yet; kept when it has one, so re-assigning rotates nothing. + if line := issueOnAssign(ctx, open, node, module); line != "" { + said += "\n " + line + } plan, _, err := planFor(ctx, open, node) if err != nil { // Kept, and still refused. Both halves are the answer, and the rest of the mesh is still @@ -152,3 +160,33 @@ func blockedElsewhere(ctx context.Context, open *stores, except string) string { out.WriteString("\nThis may or may not be what just changed — it is what is true now.") return out.String() } + +// issueOnAssign gives a newly assigned module its bus credential, the way `module issue` does, and +// says what it did in one line. Nothing for a module that declares no broker secret; nothing for one +// whose user is already minted (a credential is rotated on purpose, never by re-assigning); and when +// the bus cannot be reached from here, the line names the verb and the push that would refuse the +// module until it is run — never a silent placeholder (novox/hq issue 203). +func issueOnAssign(ctx context.Context, open *stores, node, module string) string { + inv := open.inventory + shelf, err := inv.Catalogue(ctx) + if err != nil { + return "" + } + m, known := shelf[module] + if !known || mayIssue(m) != nil { + return "" + } + user := broker.Principal{Kind: broker.KindModule, Node: node, Module: module}.Username() + if _, minted, err := inv.BusUserHash(ctx, user); err != nil || minted { + return "" + } + busAddress, err := broker.BusAddress() + if err == nil { + err = issueOnTheNewBus(ctx, inv, m, node, busAddress) + } + if err != nil { + return fmt.Sprintf("its bus credential is not issued (%v): `module issue %s --node %s` first — "+ + "`push %s` refuses to send %s until it is", err, module, node, node, module) + } + return fmt.Sprintf("its bus credential is issued and sealed to %s, and arrives with the push", node) +} diff --git a/cmd/mesh-controller/addresses_test.go b/cmd/mesh-controller/addresses_test.go index 39ec4ec..5a23d2e 100644 --- a/cmd/mesh-controller/addresses_test.go +++ b/cmd/mesh-controller/addresses_test.go @@ -175,6 +175,12 @@ func TestTheControlPlaneIsToldWhereTheNodePutTheStoreAndTheBroker(t *testing.T) Guards: []int{15672}, Resources: []map[string]any{{"id": "server", "type": "container", "name": "mesh-broker", "ports": []any{"5671:5671", "5672:5672", "127.0.0.1:15672:15672"}, "image": "mq@" + aDigest}}}) + // The control plane's own bus user is the installer's, seeded at genesis before the controller + // runs (SeedBusUser); without it a push now refuses the credential nobody issued (issue 203). + if err := open.inventory.SeedBusUser(ctx, inventory.BusUser{Username: "anchor.mesh-controller", + Kind: inventory.BusController, Node: "anchor", Module: "mesh-controller"}, "bootstrap"); err != nil { + t.Fatal(err) + } if _, err := assign(ctx, open, "anchor", "mesh-controller"); err != nil { t.Fatal(err) } diff --git a/cmd/mesh-controller/credential_on_assign_test.go b/cmd/mesh-controller/credential_on_assign_test.go new file mode 100644 index 0000000..959ef07 --- /dev/null +++ b/cmd/mesh-controller/credential_on_assign_test.go @@ -0,0 +1,90 @@ +package main + +import ( + "strings" + "testing" + + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" +) + +// A fresh assignment is pushed before its credential exists (novox/hq issue 203): `assign` recorded +// the module, `push` sealed a random own secret where the bus credential belongs, and the process +// crash-looped until a person ran `module issue` and pushed again. Now assigning a module that speaks +// on the bus issues its credential in the same act — or, when the bus cannot be reached from here, +// says which verb to run — and a push never seals a placeholder in a credential's place. + +func aTalker() catalogue.Manifest { + return catalogue.Manifest{Module: "talker", Version: "1", + OwnSecrets: catalogue.OwnSecrets{"broker": {Path: "/var/lib/mesh/talker/broker"}}, + Resources: []map[string]any{ + {"id": "state", "type": "directory", "path": "/var/lib/mesh/talker", "mode": "0700"}, + }} +} + +func TestAssigningAModuleThatSpeaksOnTheBusNamesItsCredential(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + register(t, open, aTalker()) + + // No bus is known to this process, so the credential cannot be issued here: the assignment + // stands and says exactly what must happen before a push — never silently. + said, err := assign(ctx, open, "laptop", "talker") + if err != nil { + t.Fatal(err) + } + if !strings.Contains(said, "module issue talker --node laptop") { + t.Fatalf("an assignment whose credential could not be issued does not name the verb:\n%s", said) + } + + // And the push refuses to send it, naming the same verb, rather than sealing a placeholder. + plan, settings, err := planFor(ctx, open, "laptop") + if err != nil { + t.Fatal(err) + } + _, err = declarationFor(ctx, open, "laptop", plan, settings) + if err == nil { + t.Fatal("a push sealed a placeholder where talker's bus credential belongs") + } + if !strings.Contains(err.Error(), "module issue talker --node laptop") || !strings.Contains(err.Error(), "issue 203") { + t.Fatalf("the refusal does not say what to run: %v", err) + } + + // Once the user is minted, the push goes on to the credential the mesh sealed, and re-assigning + // does not mint again: a credential rotates on purpose, never by habit. + if _, err := open.inventory.MintBusPassword(ctx, inventory.BusUser{ + Username: "laptop.talker", Kind: inventory.BusModule, Node: "laptop", Module: "talker"}); err != nil { + t.Fatal(err) + } + hash, _, err := open.inventory.BusUserHash(ctx, "laptop.talker") + if err != nil { + t.Fatal(err) + } + said, err = assign(ctx, open, "laptop", "talker") + if err != nil { + t.Fatal(err) + } + if strings.Contains(said, "module issue") { + t.Fatalf("a module with a minted credential was told to issue one:\n%s", said) + } + again, _, err := open.inventory.BusUserHash(ctx, "laptop.talker") + if err != nil { + t.Fatal(err) + } + if again != hash { + t.Fatal("re-assigning rotated the credential") + } +} + +// A module that declares no broker secret is left alone: nothing to issue, nothing said. +func TestAssigningAModuleThatDoesNotSpeakSaysNothingOfCredentials(t *testing.T) { + open := aMesh(t) + register(t, open, helloWeb()) + said, err := assign(t.Context(), open, "laptop", "hello-web") + if err != nil { + t.Fatal(err) + } + if strings.Contains(said, "credential") { + t.Fatalf("a module without a broker secret was told about credentials:\n%s", said) + } +} diff --git a/cmd/mesh-controller/plan.go b/cmd/mesh-controller/plan.go index 4300808..ef8712c 100644 --- a/cmd/mesh-controller/plan.go +++ b/cmd/mesh-controller/plan.go @@ -533,6 +533,23 @@ func renderingFor(ctx context.Context, open *stores, node string, var sealed string var err error if choosing == Allocating { + // **The broker credential is never invented here** (novox/hq issue 203). Every other + // own secret is the mesh's to make — a password nobody else knows — but this one + // is an account on the bus, minted by `module issue` and sealed by it; a push that + // made a random one would deliver a file the process cannot read and report the + // machine applied. Refused by name, with the verb. + if name == "broker" { + user := broker.Principal{Kind: broker.KindModule, Node: node, Module: m.Module}.Username() + if _, minted, err := inv.BusUserHash(ctx, user); err != nil { + return catalogue.Rendering{}, inventory.Node{}, err + } else if !minted { + return catalogue.Rendering{}, inventory.Node{}, fmt.Errorf( + "%s on %s has no bus credential: nothing was issued for %s, and a push "+ + "would seal a placeholder its process cannot read (novox/hq issue 203). "+ + "`module issue %s --node %s`, then push again", + m.Module, node, user, m.Module, node) + } + } sealed, err = inv.SecretForModule(ctx, node, m.Module, name) } else { var held bool diff --git a/internal/catalogue/declaration.go b/internal/catalogue/declaration.go index 03846bd..bc66f65 100644 --- a/internal/catalogue/declaration.go +++ b/internal/catalogue/declaration.go @@ -872,6 +872,15 @@ func (r Resolution) compose(with Rendering, owner map[string]string, if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil { copied["reload-on"] = renamed } + // **What reads one of this module's own secrets is restarted when it changes** (novox/hq + // issue 203, issue 206). A credential is re-issued by the mesh, and a container that + // mounted the old file keeps the old one open: the build machine ran for an hour on a + // credential the mesh had replaced, because its manifest restarted it on its + // environment file and nobody had thought to name the credential too. Composed here so + // no manifest has to say it, for a container or a daemon that names the secret's path. + if reads := secretsReadBy(copied, m); len(reads) > 0 { + copied["restart-on"] = withRestartOn(copied["restart-on"], reads) + } // **A version prepares its state before it runs** (novox/hq ADR 0135). Derived from the // module's own resource rather than declared beside it: what prepares the state is the // module's own code, so what it is given has to be what that code is given — and a @@ -2079,3 +2088,87 @@ func portOfEndpoint(values map[string]any, ports map[string]int) { values["port"] = port } } + +// secretsReadBy is the file resources of this module's own secrets that a container or a daemon reads +// — named in its volumes, its environment or its env-files by the secret's placed path — as +// restart-on ids. Nothing for other shapes, and nothing for a scheduled or run-once process, which +// the host refuses a restart-on for (it runs again anyway, and reads the file afresh). +func secretsReadBy(resource map[string]any, m Manifest) []string { + kind := fmt.Sprint(resource["type"]) + if kind != "container" && kind != "process" { + return nil + } + if resource["schedule"] != nil || resource["run-once"] == true { + return nil + } + var mentioned []string + for _, key := range []string{"volumes", "env", "env-file"} { + mentioned = append(mentioned, stringsIn(resource[key])...) + } + var out []string + for _, name := range sortedKeys(m.OwnSecrets) { + path := m.OwnSecrets[name].Path + if path == "" { + continue + } + for _, s := range mentioned { + // A volume is `source:destination[:mode]`; an env value or an env-file is the path itself. + if s == path || strings.HasPrefix(s, path+":") { + out = append(out, m.Module+"."+NeedID(name)) + break + } + } + } + return out +} + +// stringsIn is every string in a list or a map's values; nothing for anything else. +func stringsIn(v any) []string { + switch x := v.(type) { + case []any: + var out []string + for _, item := range x { + if s, ok := item.(string); ok { + out = append(out, s) + } + } + return out + case []string: + return x + case map[string]any: + var out []string + for _, k := range sortedKeys(x) { + if s, ok := x[k].(string); ok { + out = append(out, s) + } + } + return out + case map[string]string: + var out []string + for _, k := range sortedKeys(x) { + out = append(out, x[k]) + } + return out + } + return nil +} + +// withRestartOn is a resource's restart-on list with these ids added once each. +func withRestartOn(have any, add []string) []any { + var out []any + seen := map[string]bool{} + for _, id := range reflectsRenamed("", have) { + s := fmt.Sprint(id) + if !seen[s] { + seen[s] = true + out = append(out, s) + } + } + for _, id := range add { + if !seen[id] { + seen[id] = true + out = append(out, id) + } + } + return out +} diff --git a/internal/catalogue/secret_restart_test.go b/internal/catalogue/secret_restart_test.go new file mode 100644 index 0000000..5460c09 --- /dev/null +++ b/internal/catalogue/secret_restart_test.go @@ -0,0 +1,52 @@ +package catalogue + +import ( + "reflect" + "strings" + "testing" +) + +// What reads one of a module's own secrets is restarted when the secret changes (novox/hq issue 203, +// issue 206): the build machine kept an hour-old credential open because its manifest restarted it +// on its environment file alone. Composed, so a manifest need not say it; a scheduled process is +// left alone, because the host refuses a restart-on for one and it reads the file afresh each run. +func TestAContainerReadingAnOwnSecretIsRestartedWhenItChanges(t *testing.T) { + m := Manifest{Module: "agent", Version: "1", + OwnSecrets: OwnSecrets{"broker": {Path: "/var/lib/mesh/agent/broker"}}, + Resources: []map[string]any{ + {"id": "mesh-state", "type": "directory", "path": "/var/lib/mesh/agent", "mode": "0700"}, + {"id": "settings", "type": "file", "path": "/var/lib/mesh/agent/agent.env", "mode": "0600", "content": "A=1\n"}, + {"id": "server", "type": "container", "name": "agent", "network": "host", + "image": "registry.example/agent@sha256:" + strings.Repeat("a", 64), + "volumes": []any{"/var/lib/mesh/agent:/run/mesh:ro", "/var/lib/mesh/agent/broker:/run/mesh/broker:ro"}, + "env-file": []any{"/var/lib/mesh/agent/agent.env"}, + "restart-on": []any{"settings"}}, + {"id": "nightly", "type": "container", "name": "agent-nightly", "schedule": "0 3 * * *", + "image": "registry.example/agent@sha256:" + strings.Repeat("a", 64), + "volumes": []any{"/var/lib/mesh/agent/broker:/run/mesh/broker:ro"}}, + {"id": "other", "type": "container", "name": "agent-other", + "image": "registry.example/agent@sha256:" + strings.Repeat("a", 64)}, + }} + got, err := Resolve(shelf(m), []string{m.Module}, + Node{Name: "anchor", At: "10.0.0.1", Capabilities: map[string]bool{"container-runtime": true}}, World{}) + if err != nil { + t.Fatal(err) + } + out, err := got.Declaration(Rendering{Needed: map[string]map[string]string{"agent": {"broker": "SEALED"}}}) + if err != nil { + t.Fatal(err) + } + by := map[string]map[string]any{} + for _, r := range out { + by[r["id"].(string)] = r + } + if want := []any{"agent.settings", "agent.needs-broker"}; !reflect.DeepEqual(by["agent.server"]["restart-on"], want) { + t.Fatalf("the server reads the credential and is not restarted on it: %v", by["agent.server"]["restart-on"]) + } + if _, has := by["agent.nightly"]["restart-on"]; has { + t.Fatalf("a scheduled container was given a restart-on, which the host refuses: %v", by["agent.nightly"]["restart-on"]) + } + if _, has := by["agent.other"]["restart-on"]; has { + t.Fatalf("a container that reads no secret was given one to restart on: %v", by["agent.other"]["restart-on"]) + } +} -- 2.54.0 From 76a8b8df9e4e9537ee33de5312a42a98927eaedd Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 04:04:41 +0200 Subject: [PATCH 2/4] The controller owns a worker's shape, type included: one of the wrong type is re-made on a work queue (hq issue 206) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A holder built for a pull worker cannot bind a push one — `cannot pull subscribe to push based consumer` — and on 2026-10-03 the build machine rolled before the controller that would have redefined its worker, restarted on that for an hour, and nothing could build the controller that would have ended it. The server cannot change a consumer's type in place, so the assertion re-makes one of the wrong type: on a work queue nothing is lost, because what was acknowledged is gone from the stream and what was not is delivered again from the start. On a stream that keeps its history it is said and left, since a re-made consumer replays what this one acknowledged (issue 156), and that is a person's call. Proven against a real bus: a push worker with one ask acknowledged and two pending is re-made as pull, a pull subscription binds, and takes exactly the two. --- internal/broker/jetstream.go | 37 ++++++ internal/broker/worker_type_change_test.go | 145 +++++++++++++++++++++ 2 files changed, 182 insertions(+) create mode 100644 internal/broker/worker_type_change_test.go diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 8b4a26e..72cf05f 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -214,6 +214,43 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { switch have, err := j.js.ConsumerInfo(c.Stream, c.Name); { case err == nil: + // **The controller owns the worker's shape, type included** (novox/hq issue 206). A holder + // built for a pull worker cannot bind a push one — `cannot pull subscribe to push based + // consumer` — and on 2026-10-03 the build machine rolled before the controller that would + // have redefined its worker, restarted on that for an hour, and nothing could build the + // controller that would have ended it. The server cannot change a consumer's type in place, + // so one of the wrong type is re-made: on a work queue nothing is lost, because what was + // acknowledged is gone from the stream and what was not is delivered again from the start. + // On any other stream a re-made consumer would replay what this one acknowledged (issue + // 156), so there it is said and left, and the person re-makes it knowing the cost. + if havePush, wantPush := have.Config.DeliverSubject != "", want.DeliverSubject != ""; havePush != wantPush { + shape := func(push bool) string { + if push { + return "push" + } + return "pull" + } + info, err := j.js.StreamInfo(c.Stream) + if err != nil { + return fmt.Errorf("asking about stream %s to re-make consumer %s: %w", c.Stream, c.Name, err) + } + if info.Config.Retention != nats.WorkQueuePolicy { + j.note("consumer %s on %s is %s and should be %s; not re-made, because %s keeps its history "+ + "and a re-made consumer replays what this one acknowledged (novox/hq issue 156). Re-make it by hand", + c.Name, c.Stream, shape(havePush), shape(wantPush), c.Stream) + return nil + } + j.note("consumer %s on %s changes from %s to %s delivery: re-made where it left off, nothing "+ + "acknowledged comes back and nothing pending is lost (novox/hq issue 206); a holder bound to "+ + "the old shape binds again", c.Name, c.Stream, shape(havePush), shape(wantPush)) + if err := j.js.DeleteConsumer(c.Stream, c.Name); err != nil { + return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err) + } + if _, err := j.js.AddConsumer(c.Stream, want); err != nil { + return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err) + } + return nil + } // Where an existing consumer starts is its history, not something an assertion may move: // the server refuses a changed deliver policy outright. Carried across, so asserting twice // is the no-op a restart depends on. diff --git a/internal/broker/worker_type_change_test.go b/internal/broker/worker_type_change_test.go new file mode 100644 index 0000000..dfe719c --- /dev/null +++ b/internal/broker/worker_type_change_test.go @@ -0,0 +1,145 @@ +package broker + +import ( + "os" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// A seat's worker that changed from push to pull delivery strands a holder built for the new shape +// (novox/hq issue 206): the server refuses a pull subscription on a push consumer, and the controller +// that would redefine it was the build that nobody could take. The controller owns the worker's +// shape, type included: on a work queue it re-makes one of the wrong type, losing nothing, and a +// pull subscription then binds and takes what was pending. +// +// docker run -d --rm --name t -p 14231:4222 nats:2.10-alpine -js +// MESH_TEST_NATS=nats://127.0.0.1:14231 go test ./internal/broker/ -run TestAWorker +func TestAWorkerOfTheWrongTypeIsRemadeOnAWorkQueueAndAPullThenBinds(t *testing.T) { + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + js, err := Dial(url) + if err != nil { + t.Fatal(err) + } + defer js.Close() + + const stream, worker, filter = "SEAT_T_SHELF", "SEAT_T_SHELF_worker", "mesh.seat.t-shelf.accept.>" + _ = js.js.DeleteStream(stream) + if _, err := js.js.AddStream(&nats.StreamConfig{ + Name: stream, Subjects: []string{filter}, Retention: nats.WorkQueuePolicy, Storage: nats.MemoryStorage, + }); err != nil { + t.Fatal(err) + } + defer func() { _ = js.js.DeleteStream(stream) }() + + // The worker as the previous controller defined it: push, in a queue group. + if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{ + Durable: worker, AckPolicy: nats.AckExplicitPolicy, AckWait: 60 * time.Second, MaxDeliver: 5, + FilterSubject: filter, DeliverSubject: "_DELIVER." + worker, DeliverGroup: "holders", + }); err != nil { + t.Fatal(err) + } + for _, body := range []string{"one", "two", "three"} { + if _, err := js.js.Publish("mesh.seat.t-shelf.accept.build", []byte(body)); err != nil { + t.Fatal(err) + } + } + // The old holder took and acknowledged the first ask, then went away. + old, err := js.js.QueueSubscribeSync(filter, "holders", nats.Bind(stream, worker)) + if err != nil { + t.Fatal(err) + } + m, err := old.NextMsg(twoSeconds) + if err != nil { + t.Fatal(err) + } + if string(m.Data) != "one" { + t.Fatalf("the first ask is %q", m.Data) + } + if err := m.AckSync(); err != nil { + t.Fatal(err) + } + if err := old.Unsubscribe(); err != nil { + t.Fatal(err) + } + + // The new controller asserts the worker as the mesh derives it now: pull. + if err := js.EnsureConsumer(Consumer{ + Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60, MaxDeliver: 5, + Why: "the test's worker", + }); err != nil { + t.Fatal(err) + } + have, err := js.js.ConsumerInfo(stream, worker) + if err != nil { + t.Fatal(err) + } + if have.Config.DeliverSubject != "" || have.Config.DeliverGroup != "" { + t.Fatalf("the worker is still push: %+v", have.Config) + } + + // A holder built for the new shape binds, and takes exactly what the old one left. + sub, err := js.js.PullSubscribe(filter, worker, nats.Bind(stream, worker), nats.ManualAck()) + if err != nil { + t.Fatalf("a pull subscription does not bind the re-made worker: %v", err) + } + got, err := sub.Fetch(3, nats.MaxWait(twoSeconds)) + if err != nil && len(got) == 0 { + t.Fatalf("nothing pending was delivered: %v", err) + } + var bodies []string + for _, g := range got { + bodies = append(bodies, string(g.Data)) + _ = g.Ack() + } + if len(bodies) != 2 || bodies[0] != "two" || bodies[1] != "three" { + t.Fatalf("the pending asks after the acknowledged one, in order: %v", bodies) + } + + // Asserted again, the pull worker is the no-op a restart depends on. + if err := js.EnsureConsumer(Consumer{ + Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60, MaxDeliver: 5, + }); err != nil { + t.Fatal(err) + } +} + +// On a stream that keeps its history, a worker of the wrong type is said and left: re-making it would +// replay what it acknowledged (novox/hq issue 156), and that is a person's call. +func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(t *testing.T) { + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + js, err := Dial(url) + if err != nil { + t.Fatal(err) + } + defer js.Close() + const stream, worker, filter = "EVENTS_T", "EVENTS_T_reader", "mesh.t.event.>" + _ = js.js.DeleteStream(stream) + if _, err := js.js.AddStream(&nats.StreamConfig{Name: stream, Subjects: []string{filter}, Storage: nats.MemoryStorage}); err != nil { + t.Fatal(err) + } + defer func() { _ = js.js.DeleteStream(stream) }() + if _, err := js.js.AddConsumer(stream, &nats.ConsumerConfig{ + Durable: worker, AckPolicy: nats.AckExplicitPolicy, AckWait: 60 * time.Second, + FilterSubject: filter, DeliverSubject: "_DELIVER." + worker, + }); err != nil { + t.Fatal(err) + } + if err := js.EnsureConsumer(Consumer{Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60}); err != nil { + t.Fatal(err) + } + have, err := js.js.ConsumerInfo(stream, worker) + if err != nil { + t.Fatal(err) + } + if have.Config.DeliverSubject == "" { + t.Fatal("a history stream's consumer was re-made, which replays what it acknowledged") + } +} -- 2.54.0 From 28853a251bdbf91cea3b0533f49e53bb0472f9d6 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 04:04:41 +0200 Subject: [PATCH 3/4] The build seat's holder follows the controller that defines its worker (hq issue 206) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A plan is ordered by artifacts and says nothing about what must be running before what (ADR 0162); on 2026-10-03 that put the build machine in tier 0 and the controller in tier 1, and the new build machine could not bind the worker the old controller had defined. One running order enters the graph, named as its own edge: a module claiming the build seat follows the control plane, and the built-by edge from the control plane to that holder yields to it — the controller is built by whichever build machine is running, as the runtime image always was. The edge orders a plan and never widens it, like built-by. --- cmd/mesh-controller/release_plan.go | 25 ++++++++++--- cmd/mesh-controller/worker_order_test.go | 45 ++++++++++++++++++++++++ internal/inventory/dependencies.go | 16 +++++++++ 3 files changed, 81 insertions(+), 5 deletions(-) create mode 100644 cmd/mesh-controller/worker_order_test.go diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 83976f5..438ea38 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -36,17 +36,29 @@ func tiersOf(set []string, edges []inventory.Edge) [][]string { for _, m := range set { deps[m] = map[string]bool{} } + // The build seat's holders follow the controller that defines their worker (EdgeWorkerOf, + // novox/hq issue 206), so the built-by edge from that controller to such a holder yields: the + // controller is built by whichever build machine is running, as the runtime image always was. + worker := map[string]map[string]bool{} + for _, e := range edges { + if e.Kind == inventory.EdgeWorkerOf && in[e.From] && in[e.To] { + if worker[e.To] == nil { + worker[e.To] = map[string]bool{} + } + worker[e.To][e.From] = true + } + } for _, e := range edges { // A code dependency — B packages A's source — rebuilds B with A, in the same tier: B's // build needs nothing of A's first. The other kinds order: stands-on and declared after // the base is built, built-by after the build machine is built and running — except for - // what the build machine itself stands on. The runtime image is built by the builder and - // the builder is built on the runtime image; the image comes first, built by the builder - // that is running, which is the only one there could be. + // what the build machine itself stands on, and for the controller whose worker the build + // machine binds. The runtime image is built by the builder and the builder is built on the + // runtime image; the image comes first, built by the builder that is running. if !in[e.From] || !in[e.To] || e.From == e.To || e.Kind == inventory.EdgePackages { continue } - if e.Kind == inventory.EdgeBuiltBy && isBaseOf(e.From, e.To, edges, in) { + if e.Kind == inventory.EdgeBuiltBy && (isBaseOf(e.From, e.To, edges, in) || worker[e.From][e.To]) { continue } deps[e.From][e.To] = true @@ -122,7 +134,10 @@ func reachableFrom(moved []string, edges []inventory.Edge) []string { for grew := true; grew; { grew = false for _, e := range edges { - if e.Kind == inventory.EdgeBuiltBy { + // Built-by and worker-of order a plan; neither widens it. A new build machine changes + // nothing it builds, and a new controller changes nothing about the holder it orders — + // what packages the controller's source is already a code edge. + if e.Kind == inventory.EdgeBuiltBy || e.Kind == inventory.EdgeWorkerOf { continue } if in[e.To] && !in[e.From] { diff --git a/cmd/mesh-controller/worker_order_test.go b/cmd/mesh-controller/worker_order_test.go new file mode 100644 index 0000000..f5b43f2 --- /dev/null +++ b/cmd/mesh-controller/worker_order_test.go @@ -0,0 +1,45 @@ +package main + +import ( + "testing" + + "github.com/novox/mesh-controller/internal/inventory" +) + +// The holder of the build seat follows the controller that defines its worker (novox/hq issue 206). +// On 2026-10-03 a plan put the build machine in tier 0 and the controller in tier 1; the new build +// machine could not bind the worker the old controller had defined, and nothing could build the +// controller that would have redefined it. The built-by edge from the controller to its build +// machine yields to that order: the controller is built by whichever build machine is running. +func TestTheBuildSeatsHolderFollowsTheControllerThatDefinesItsWorker(t *testing.T) { + edges := []inventory.Edge{ + {From: "build-agent", To: "mesh-controller", Kind: inventory.EdgePackages}, + {From: "build-agent", To: "mesh-controller", Kind: inventory.EdgeWorkerOf}, + {From: "mesh-controller", To: "build-agent", Kind: inventory.EdgeBuiltBy}, + {From: "route-proxy", To: "mesh-controller", Kind: inventory.EdgePackages}, + {From: "route-proxy", To: "build-agent", Kind: inventory.EdgeBuiltBy}, + } + set := reachableFrom([]string{"mesh-controller"}, edges) + if len(set) != 3 { + t.Fatalf("the controller, what packages it, and nothing more: %v", set) + } + tiers := tiersOf(set, edges) + pos := map[string]int{} + for i, tier := range tiers { + for _, m := range tier { + pos[m] = i + } + } + if pos["mesh-controller"] != 0 { + t.Fatalf("the controller first, built by the build machine that is running: %v", tiers) + } + if pos["build-agent"] <= pos["mesh-controller"] { + t.Fatalf("the build machine after the controller that defines its worker: %v", tiers) + } + if pos["route-proxy"] <= pos["build-agent"] { + t.Fatalf("what the build machine builds comes after it: %v", tiers) + } + if hasCycle(tiers, edges) { + t.Fatalf("no cycle here: %v", tiers) + } +} diff --git a/internal/inventory/dependencies.go b/internal/inventory/dependencies.go index 26274e7..580586c 100644 --- a/internal/inventory/dependencies.go +++ b/internal/inventory/dependencies.go @@ -18,8 +18,17 @@ const ( EdgeBuiltBy = "built-by" // EdgeDeclared: the manifest's own `build.on`. EdgeDeclared = "declared" + // EdgeWorkerOf: the module holds the build seat, whose worker the control plane defines + // (novox/hq issue 206). The one place a *running* order enters the graph: a build machine rolled + // before the controller that redefines its worker cannot bind it, and nothing can then build the + // controller that would end that — so the holder of the build seat follows the controller, and + // the controller is built by whichever build machine is running, as it always was. + EdgeWorkerOf = "worker-of" ) +// TheControlPlane is the module that defines every seat's worker on the bus. +const TheControlPlane = "mesh-controller" + // Edge is one dependency: From depends on To, in the way Kind says. type Edge struct { From string `json:"from"` @@ -106,6 +115,13 @@ func dependenciesOf(entries []Entry, against map[string][]string, read map[strin } } } + if known[TheControlPlane] { + for _, b := range builders { + if b != TheControlPlane { + add(b, TheControlPlane, EdgeWorkerOf) + } + } + } sort.Slice(out, func(a, b int) bool { if out[a].From != out[b].From { return out[a].From < out[b].From -- 2.54.0 From c294949f2ad132c2a96cc7fc3a385e9ca6b0bfd8 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 11:07:11 +0200 Subject: [PATCH 4/4] A worker of the wrong type on a history-keeping stream is re-made to deliver from now on, never from the start (hq issue 207) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Left for a hand, the hand re-made it with the server's default — everything the stream holds — and on 2026-10-03 that replayed every build ask since 1 October into the catalogue. Re-made with deliver-new instead: nothing acknowledged comes back; what was in flight is said and asked again. --- internal/broker/jetstream.go | 22 ++++++++++----- internal/broker/worker_type_change_test.go | 32 ++++++++++++++++++---- 2 files changed, 42 insertions(+), 12 deletions(-) diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 72cf05f..cf86519 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -235,14 +235,22 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { return fmt.Errorf("asking about stream %s to re-make consumer %s: %w", c.Stream, c.Name, err) } if info.Config.Retention != nats.WorkQueuePolicy { - j.note("consumer %s on %s is %s and should be %s; not re-made, because %s keeps its history "+ - "and a re-made consumer replays what this one acknowledged (novox/hq issue 156). Re-make it by hand", - c.Name, c.Stream, shape(havePush), shape(wantPush), c.Stream) - return nil + // **A stream that keeps its history is re-made from now on, never from the start.** + // Left for a hand, the hand re-makes it with the server's default — everything the + // stream holds — which on 2026-10-03 replayed every build ask since 1 October and + // re-registered nine modules from the past (novox/hq issue 207). What this consumer + // had not yet acknowledged is lost with it, and said: on a history stream that is + // the smaller cost, and the asks in flight are visible to whoever asked. + j.note("consumer %s on %s changes from %s to %s delivery on a stream that keeps its history: "+ + "re-made to deliver from now on, so nothing this one acknowledged comes back (novox/hq issue "+ + "207); %d ask(s) it had not acknowledged are not carried over and must be asked again", + c.Name, c.Stream, shape(havePush), shape(wantPush), have.NumPending+uint64(have.NumAckPending)) + want.DeliverPolicy = nats.DeliverNewPolicy + } else { + j.note("consumer %s on %s changes from %s to %s delivery: re-made where it left off, nothing "+ + "acknowledged comes back and nothing pending is lost (novox/hq issue 206); a holder bound to "+ + "the old shape binds again", c.Name, c.Stream, shape(havePush), shape(wantPush)) } - j.note("consumer %s on %s changes from %s to %s delivery: re-made where it left off, nothing "+ - "acknowledged comes back and nothing pending is lost (novox/hq issue 206); a holder bound to "+ - "the old shape binds again", c.Name, c.Stream, shape(havePush), shape(wantPush)) if err := j.js.DeleteConsumer(c.Stream, c.Name); err != nil { return fmt.Errorf("re-making consumer %s on %s as %s: %w", c.Name, c.Stream, shape(wantPush), err) } diff --git a/internal/broker/worker_type_change_test.go b/internal/broker/worker_type_change_test.go index dfe719c..641ebd5 100644 --- a/internal/broker/worker_type_change_test.go +++ b/internal/broker/worker_type_change_test.go @@ -108,9 +108,10 @@ func TestAWorkerOfTheWrongTypeIsRemadeOnAWorkQueueAndAPullThenBinds(t *testing.T } } -// On a stream that keeps its history, a worker of the wrong type is said and left: re-making it would -// replay what it acknowledged (novox/hq issue 156), and that is a person's call. -func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(t *testing.T) { +// On a stream that keeps its history, a worker of the wrong type is re-made to deliver from now on: +// re-making it from the start would replay what it acknowledged (novox/hq issue 156), and leaving it +// for a hand re-made it exactly that way on 2026-10-03 (issue 207). +func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsRemadeFromNowOn(t *testing.T) { url := os.Getenv("MESH_TEST_NATS") if url == "" { t.Skip("MESH_TEST_NATS unset") @@ -132,6 +133,12 @@ func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(t *testing.T) { }); err != nil { t.Fatal(err) } + // History the old consumer would have acknowledged long ago, and must not come back. + for i := 0; i < 3; i++ { + if _, err := js.js.Publish("mesh.t.event.old", []byte("old")); err != nil { + t.Fatal(err) + } + } if err := js.EnsureConsumer(Consumer{Name: worker, Stream: stream, Filters: []string{filter}, AckWaitSeconds: 60}); err != nil { t.Fatal(err) } @@ -139,7 +146,22 @@ func TestAWorkerOfTheWrongTypeOnAHistoryStreamIsLeftAndSaid(t *testing.T) { if err != nil { t.Fatal(err) } - if have.Config.DeliverSubject == "" { - t.Fatal("a history stream's consumer was re-made, which replays what it acknowledged") + if have.Config.DeliverSubject != "" { + t.Fatal("a history stream's consumer of the wrong type was left as it was") + } + if have.Config.DeliverPolicy != nats.DeliverNewPolicy || have.NumPending != 0 { + t.Fatalf("re-made consumer delivers %v with %d pending; it must deliver from now on with nothing of the past", have.Config.DeliverPolicy, have.NumPending) + } + // And what arrives from now on is delivered. + if _, err := js.js.Publish("mesh.t.event.new", []byte("new")); err != nil { + t.Fatal(err) + } + sub, err := js.js.PullSubscribe(filter, worker, nats.Bind(stream, worker)) + if err != nil { + t.Fatal(err) + } + got, err := sub.Fetch(1, nats.MaxWait(3*time.Second)) + if err != nil || len(got) != 1 || string(got[0].Data) != "new" { + t.Fatalf("the re-made consumer delivered %v, %v; want the one new message", got, err) } } -- 2.54.0