From bde4b61b3bc90dfcca1ddaed86da04c18fd546da Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 2 Oct 2026 22:31:55 +0200 Subject: [PATCH 1/5] The build role is the node-scoped seat node-build-agent, and its work is shared by every holder (hq ADR 0190) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit One build machine built everything, in a queue of one, because the seat was mesh-scoped and a mesh seat has one holder. ADR 0190 makes building a node role: node-build-agent, held on every machine that builds, with the work asked of the role and taken by whichever holder is idle. The work subject of a node-scoped seat carries no node — that token is for a seat's tools, asked of one machine (design 33 §4) — so holders on several machines read one queue; a test now says so. The retired mesh-build-machine row stays while the builder module's registered manifest claims it: a claim to a seat the mesh no longer defines is refused, and the machine holding it would be unresolvable until build-agent replaces it. Removed once no manifest claims it. The installer's genesis template (in the host's repository) still grants the controller the old seat's subjects; its test here says so until that template names node-build-agent. --- internal/broker/nats.go | 8 +++++--- internal/broker/nats_test.go | 21 +++++++++++++++++++++ internal/broker/streams.go | 2 +- internal/broker/testdata/composed.conf | 4 ++-- internal/catalogue/seats.go | 11 ++++++++++- internal/catalogue/seats_test.go | 7 ++++--- internal/inventory/dependencies.go | 2 +- internal/inventory/dependencies_test.go | 2 +- internal/link/builds.go | 6 ++++-- internal/link/builds_current_test.go | 2 +- internal/link/builds_nats_test.go | 6 +++--- 11 files changed, 53 insertions(+), 18 deletions(-) diff --git a/internal/broker/nats.go b/internal/broker/nats.go index e81a32d..b0272bf 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -110,10 +110,10 @@ type Principal struct { PasswordHash string } -// meshSeatsTheControllerUses are the roles the mesh's own flows submit work to. Named rather than +// seatsTheControllerAsks are the roles the mesh's own flows submit work to. Named rather than // derived from the seat set: the controller is not a module and declares no `uses`, so its side of a // seat has to be stated, and a list is what makes "which roles does the mesh itself talk to" answerable. -var meshSeatsTheControllerUses = []string{"mesh-build-machine"} +var seatsTheControllerAsks = []string{"node-build-agent"} // enrolmentPrefix is the space every enrolling node's user and inbox live under, so the one place the // controller may answer an enrolment is derived from the same constant the user is named from. @@ -208,7 +208,9 @@ func PermissionsFor(p Principal) (Permissions, error) { // Work the mesh's own flows submit to a role, and the outcomes they wait on (ADR 0121). A // build is the one today: the controller asks, and reads the answer from the seat's event // like the catalogue does — which is why no holder needs to publish into anybody's inbox. - for _, seat := range meshSeatsTheControllerUses { + // A node-scoped seat's work subject carries no node (novox/hq ADR 0190): the ask goes to + // the role, and whichever machine holding it is idle takes it. + for _, seat := range seatsTheControllerAsks { pub = append(pub, "mesh.seat."+seat+".accept.>") } // **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own diff --git a/internal/broker/nats_test.go b/internal/broker/nats_test.go index 4479a21..1a60c2e 100644 --- a/internal/broker/nats_test.go +++ b/internal/broker/nats_test.go @@ -441,3 +441,24 @@ func contains(list []string, want string) bool { } return false } + +// A node-scoped seat's work is shared (novox/hq ADR 0190): its holder on any machine subscribes the +// seat's one work subject, with no node in it, so holders on several machines read one queue. The +// node token belongs to a seat's tools, which are asked of one machine (design 33 §4), not to its work. +func TestANodeSeatsWorkSubjectCarriesNoNode(t *testing.T) { + seat := Seat{Name: "node-build-agent", Scope: "node", Accepts: []string{"build"}, Serves: []string{"status"}} + perms, err := PermissionsFor(Principal{Kind: KindModule, Node: "anchor", Module: "build-agent", Holds: []Seat{seat}}) + if err != nil { + t.Fatal(err) + } + has(t, perms.Subscribe, "mesh.seat.node-build-agent.accept.build") + hasNot(t, perms.Subscribe, "mesh.seat.node-build-agent.accept.build.anchor") + // And its tools still carry the machine. + has(t, perms.Subscribe, "mesh.seat.node-build-agent.tool.status.anchor") + // The controller asks the role, not a machine. + controller, err := PermissionsFor(Principal{Kind: KindController}) + if err != nil { + t.Fatal(err) + } + has(t, controller.Publish, "mesh.seat.node-build-agent.accept.>") +} diff --git a/internal/broker/streams.go b/internal/broker/streams.go index fd9f7f9..210359b 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -217,7 +217,7 @@ var ControllerFollows = []string{ // A build's outcome, which is the build-machine role's own event now (ADR 0121) rather than a // message on the control branch. Same three audiences, one publish: whoever asked, this, and the // catalogue. - seatEventSubject("mesh-build-machine", "built"), + seatEventSubject("node-build-agent", "built"), // The forge's merges: what moved a source, so the mesh builds what that source produces // without anybody telling it (novox/hq 04-ISSUES/131). Appended, because the index is a name. moduleEventSubject("gitea", "pull.merged"), diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 623f8c9..79d73f3 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,8 +24,8 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused"] } - subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>"] } + subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/catalogue/seats.go b/internal/catalogue/seats.go index 78248e0..536f9ec 100644 --- a/internal/catalogue/seats.go +++ b/internal/catalogue/seats.go @@ -96,8 +96,17 @@ var defaultSeats = []Seat{ // A build says what it does as it does it (novox/hq ADR 0157): `started` when work is taken, // `log.` for every line, `built` for the outcome. The log's tail token is the build's // id, so a reader follows one build by subject alone. + // **Node-scoped, and every holder takes from one queue** (novox/hq ADR 0190): a build is asked of + // the role, and whichever machine holding the seat is idle pulls it. One holder per machine is + // what the scope says; sharing the work is what a seat's queue has always done. + {Name: "node-build-agent", Scope: ScopeNode, + Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0190"}, + // **Retired by ADR 0190, kept while a manifest still claims it.** The one build machine's seat. + // A claim to a seat the mesh no longer defines is refused, and the module holding this one is + // assigned on a live machine until build-agent replaces it — removing the row first would make + // that machine unresolvable in the meantime. Deleted once no registered manifest claims it. {Name: "mesh-build-machine", Scope: ScopeMesh, - Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0121"}, + Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0190"}, {Name: "node-dns-resolver", Scope: ScopeNode, Decision: "novox/hq ADR 0121"}, // The intrusion prevention's verbs (novox/hq ADR 0179): what a person asks a machine's ban list // whatever keeps it — who is banned and why, ban one address, let one go. Every holder serves all diff --git a/internal/catalogue/seats_test.go b/internal/catalogue/seats_test.go index a3e7047..22f3db7 100644 --- a/internal/catalogue/seats_test.go +++ b/internal/catalogue/seats_test.go @@ -44,9 +44,10 @@ func TestTheSeatsAreAClosedSetAndEachNamesItsDecision(t *testing.T) { delivered[s.Delivers] = s.Name } } - // Sixteen since node-service-manager (novox/hq ADR 0177). - if len(Seats()) != 16 { - t.Errorf("the mesh defines %d seats rather than 16; the set is closed, so a change here is "+ + // Seventeen since node-build-agent (novox/hq ADR 0190) — sixteen once the retired + // mesh-build-machine row goes, when no registered manifest claims it any more. + if len(Seats()) != 17 { + t.Errorf("the mesh defines %d seats rather than 17; the set is closed, so a change here is "+ "a decision (novox/hq ADR 0110): %s", len(Seats()), seatNames()) } } diff --git a/internal/inventory/dependencies.go b/internal/inventory/dependencies.go index 0633526..26274e7 100644 --- a/internal/inventory/dependencies.go +++ b/internal/inventory/dependencies.go @@ -63,7 +63,7 @@ func dependenciesOf(entries []Entry, against map[string][]string, read map[strin if r := repositoryKey(e.Source.Repository); r != "" { byRepository[r] = append(byRepository[r], name) } - if e.Manifest.ClaimsSeat("mesh-build-machine") { + if e.Manifest.ClaimsSeat("node-build-agent") || e.Manifest.ClaimsSeat("mesh-build-machine") { builders = append(builders, name) } } diff --git a/internal/inventory/dependencies_test.go b/internal/inventory/dependencies_test.go index 198e193..1109d34 100644 --- a/internal/inventory/dependencies_test.go +++ b/internal/inventory/dependencies_test.go @@ -12,7 +12,7 @@ func TestDependenciesAreOneRelationWithTheirKinds(t *testing.T) { return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}} } builder := entry("builder", "http://forge/novox/mesh-catalog.git") - builder.Manifest.Claims = []catalogue.Claim{{Name: "mesh-build-machine", Scope: catalogue.ScopeMesh}} + builder.Manifest.Claims = []catalogue.Claim{{Name: "node-build-agent", Scope: catalogue.ScopeNode}} plugin := entry("shop-plugin", "http://forge/novox/mesh-catalog.git") plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}} entries := []Entry{ diff --git a/internal/link/builds.go b/internal/link/builds.go index 28465fe..d71c3be 100644 --- a/internal/link/builds.go +++ b/internal/link/builds.go @@ -19,8 +19,10 @@ import ( // act on or a declaration a node reconciles toward; a build is a request that takes minutes and has // exactly one answer. Too long for request/reply, too particular to be an event. -// TheBuildMachine is the role a build is submitted to. -const TheBuildMachine = "mesh-build-machine" +// TheBuildMachine is the role a build is submitted to: node-scoped, held on every machine that +// builds, and the work shared among them (novox/hq ADR 0190). The name stays for every caller; what +// it names moved from the mesh's one build machine to whichever build agent is idle. +const TheBuildMachine = "node-build-agent" // BuildWork is where a build request lands, and BuildOutcome is where its result does. Derived from // the seat, so both sides name the role and neither names the other. diff --git a/internal/link/builds_current_test.go b/internal/link/builds_current_test.go index 7f43f2c..12406fe 100644 --- a/internal/link/builds_current_test.go +++ b/internal/link/builds_current_test.go @@ -26,7 +26,7 @@ func TestTheOldBusAnnouncesABuildUnderBothNames(t *testing.T) { if KeyRoleBuilt != "built" { t.Fatalf("the role's event is %q, and a holder emits its verbs bare", KeyRoleBuilt) } - if TheBuildMachine != "mesh-build-machine" { + if TheBuildMachine != "node-build-agent" { t.Fatalf("the role is %q", TheBuildMachine) } // The two must differ, or one publish would serve both and this doubling would be pointless. diff --git a/internal/link/builds_nats_test.go b/internal/link/builds_nats_test.go index a87d6a7..38971a1 100644 --- a/internal/link/builds_nats_test.go +++ b/internal/link/builds_nats_test.go @@ -48,7 +48,7 @@ func aBusWithTheBuildRole(t *testing.T) *broker.JetStream { t.Fatal(err) } clean := func() { - _ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") + _ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT") for _, s := range broker.MeshStreams() { _ = js.Context().PurgeStream(s.Name) } @@ -158,7 +158,7 @@ func TestNatsABuildIsTakenAndItsOutcomeReachesEverybody(t *testing.T) { // And the work left the queue: a request a machine took and settled must not be given to another. deadline := time.Now().Add(5 * time.Second) for time.Now().Before(deadline) { - info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE") + info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT") if err == nil && info.State.Msgs == 0 { return } @@ -177,7 +177,7 @@ func TestNatsABuildWaitsForAMachineRatherThanFailing(t *testing.T) { if _, err := js.Context().Publish(BuildWork(), body); err != nil { t.Fatal(err) } - info, err := js.Context().StreamInfo("SEAT_MESH_BUILD_MACHINE") + info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT") if err != nil || info.State.Msgs != 1 { t.Fatalf("the work did not queue: %+v %v", info, err) } -- 2.54.0 From 9f9d9b3b25934f62fd267577f690ff649c0feb8c Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 2 Oct 2026 22:34:44 +0200 Subject: [PATCH 2/5] A seat's holders pull one ask at a time from one shared worker (hq ADR 0190, issue 186) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The worker a holder bound was a push consumer in a queue group with one ask in flight: right for one holder, and with two it would still be a queue of one — the server hands a pushed ask to whichever subscriber it picks, busy or not, and the in-flight cap is per consumer, not per holder. Now the worker is pulled: every machine holding the seat binds the same durable and fetches one ask when it has finished the last, so an idle machine is the one that takes the next, the asks in flight are bounded by the holders working, and nothing is delivered that nobody asked for — which is also what ended the race issue 186 describes. A holder's grants trade the delivery subject for MSG.NEXT on the worker; the ack grant and the heartbeat that keeps a long build alive stay. Proven against a real bus: the build round trip, a backlog taken by a machine that arrives later, and work handed back by one machine coming round again. --- internal/broker/derived.go | 31 ++++++++++---------- internal/broker/derived_test.go | 37 ++++++++++++++++-------- internal/broker/nats.go | 14 +++++---- internal/broker/testdata/composed.conf | 4 +-- internal/link/builds_nats.go | 40 ++++++++++++++++++-------- 5 files changed, 80 insertions(+), 46 deletions(-) diff --git a/internal/broker/derived.go b/internal/broker/derived.go index 2471c3e..466b88a 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -144,12 +144,18 @@ func ConsumerFor(p Principal) (Consumer, bool) { }, true } -// HolderConsumerFor is the worker a seat's holder gets on that seat's work queue. +// HolderConsumerFor is the worker a seat's holders share on that seat's work queue. // -// **A queue group even though the seat guarantees one holder.** The seat is *authority* — who may -// be the telegram sender — and the queue group is *delivery*. Tie delivery to the seat and the -// day somebody allows two holders for throughput, every message is processed twice with nothing -// reporting it. Kept separate, relaxing one changes nothing about the other. +// **One worker for every holder, and each holder pulls one ask when it is idle** (novox/hq ADR +// 0190). The seat is *authority* — who may be the telegram sender — and the worker is *delivery*, +// kept separate so that relaxing one changes nothing about the other: a node-scoped seat has a +// holder per machine, and all of them take from this one consumer, so the work is shared without +// any holder knowing about the others. Pulled rather than pushed because a push consumer hands the +// next ask to whichever subscriber the server picks, busy or not, and a pulled one is asked for by +// a holder that has just become free. Which is also what ends the race issue 186 describes — asks +// delivered behind the one being worked, expiring unacknowledged and dropped after the fifth +// redelivery: nothing is delivered that nobody asked for. A long build keeps its own ask alive +// (stillWorking); the ack wait is for a holder that died. func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) { if len(seat.Accepts) == 0 { return Consumer{}, false @@ -158,18 +164,13 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) Name: "SEAT_" + upperSnake(seat.Name) + "_worker", Stream: seatStreamName(seat.Name), Filters: []string{"mesh.seat." + seat.Name + ".accept.>"}, - Queue: "holders", AckWaitSeconds: 60, MaxDeliver: 5, - // **One in flight.** A holder works one ask at a time, so the server hands it one at a - // time: with the default of many, every ask behind the one being worked was delivered, - // left unacknowledged for the length of the work, redelivered after the ack wait, and - // after the fifth time dropped — on 2026-10-01 twenty-six of forty-three builds asked in - // two minutes were never built, and the queue read as empty (novox/hq issue 186). - MaxAckPending: 1, - Why: fmt.Sprintf("%s on %s holds %s; it acknowledges after the work is done, so a "+ - "crash mid-work redelivers rather than loses; one in flight, so a queue of asks is a "+ - "queue and not a race against the ack wait", module, node, seat.Name), + // As many in flight as there are holders working, which pulling bounds by itself: a holder + // fetches one and fetches again only after it acknowledged. The server's default stands. + Why: fmt.Sprintf("%s on %s holds %s; every holder pulls one ask at a time from this worker "+ + "and acknowledges after the work is done, so a crash mid-work redelivers rather than "+ + "loses and an idle holder is the one that takes the next ask", module, node, seat.Name), }, true } diff --git a/internal/broker/derived_test.go b/internal/broker/derived_test.go index 939d078..c7257d5 100644 --- a/internal/broker/derived_test.go +++ b/internal/broker/derived_test.go @@ -88,15 +88,20 @@ func TestAModuleThatConsumesNothingGetsNoConsumer(t *testing.T) { } } -// The seat is authority and the queue group is delivery. Tie them together and the day somebody -// allows two holders, every message is processed twice with nothing reporting it. -func TestAHoldersWorkerUsesAQueueGroupAnyway(t *testing.T) { +// The seat is authority and the worker is delivery (novox/hq ADR 0190): one worker per seat, shared +// by every holder and pulled from, so a second holder takes the next ask rather than a copy of the +// same one — which is what a queue group used to guard, and what pulling one durable gives outright. +func TestAHoldersWorkerIsOneSharedByItsHolders(t *testing.T) { c, ok := HolderConsumerFor("one", "telegram", telegramSeat()) if !ok { t.Fatal("the holder of a seat with inbound work got no worker") } - if c.Queue == "" { - t.Fatal("the worker is not in a queue group, so a second holder would double-process") + two, _ := HolderConsumerFor("two", "telegram", telegramSeat()) + if c.Name != two.Name || c.Stream != two.Stream { + t.Fatal("two holders got two workers, so each would process every ask") + } + if c.Push || c.Queue != "" { + t.Fatal("the worker is pushed, so the server would hand an ask to a busy holder") } if c.Stream != "SEAT_TELEGRAM_SENDER" { t.Fatalf("the worker reads %q, not the seat's own stream", c.Stream) @@ -154,15 +159,23 @@ func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) { has(t, perms.Subscribe, c.Filters[0]) } -// A holder works one ask at a time, so the server hands it one at a time (novox/hq issue 186): -// asks queued behind the one being worked wait in the stream rather than being delivered, -// left to expire and dropped after the fifth redelivery. -func TestAHoldersWorkerTakesOneAskAtATime(t *testing.T) { - c, found := HolderConsumerFor("anchor", "builder", DeclaredSeat{Name: "mesh-build-machine", Accepts: []string{"build"}}) +// Every holder of a seat shares one worker and pulls from it (novox/hq ADR 0190): no queue group +// and no delivery subject, because a push consumer hands the next ask to whichever subscriber the +// server picks, busy or not; and no cap of one in flight, because pulling bounds the asks in flight +// by the holders that are free — which is what ended the race of issue 186, where asks delivered +// behind the one being worked expired and were dropped. +func TestAHoldersWorkerIsPulledByEveryHolder(t *testing.T) { + c, found := HolderConsumerFor("anchor", "build-agent", DeclaredSeat{Name: "node-build-agent", Accepts: []string{"build"}}) if !found { t.Fatal("a seat that accepts work has no worker") } - if c.MaxAckPending != 1 { - t.Fatalf("the worker may have %d asks in flight; one, so a queue is a queue", c.MaxAckPending) + if c.Queue != "" || c.Push { + t.Fatalf("the worker is pushed (queue %q, push %v); a holder pulls when it is free", c.Queue, c.Push) + } + if c.MaxAckPending != 0 { + t.Fatalf("the worker caps asks in flight at %d; pulling bounds them by the holders working", c.MaxAckPending) + } + if c.Name != "SEAT_NODE_BUILD_AGENT_worker" || c.Stream != "SEAT_NODE_BUILD_AGENT" { + t.Fatalf("the worker is %s on %s; one per seat, shared by its holders", c.Name, c.Stream) } } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index b0272bf..b52cb34 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -370,13 +370,17 @@ func PermissionsFor(p Principal) (Permissions, error) { // 3. Seats it holds: full participation. for _, s := range p.Holds { - // Taking work from the role's queue: the worker consumer it binds (asked about, - // delivered on, acknowledged), each on the seat's own stream. The first machine to - // take work over the new bus was refused the asking (2026-09-28). + // Taking work from the role's queue: the worker consumer every holder shares (asked + // about, pulled from, acknowledged), on the seat's own stream (novox/hq ADR 0190). A + // holder pulls — asks the consumer for its next message, answered on its own inbox — + // so what it needs is MSG.NEXT on that worker and nothing delivered to it. The first + // machine to 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, "_DELIVER."+worker+".>") - pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>") + pub = append(pub, + "$JS.API.CONSUMER.INFO."+stream+"."+worker, + "$JS.API.CONSUMER.MSG.NEXT."+stream+"."+worker, + "$JS.ACK."+stream+"."+worker+".>") for _, a := range s.Accepts { sub = append(sub, seatSubject(s, "accept", a)) } diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 79d73f3..3994266 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -37,8 +37,8 @@ accounts { 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", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } - subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] } + 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", "$JS.API.CONSUMER.MSG.NEXT.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] } + subscribe: { allow: ["_INBOX.one.telegram.>", "mesh.assignment.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: { diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index 85e32c3..0355857 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -126,29 +126,31 @@ func (m *natsMachine) Close() { } } -// Take binds to the role's worker and hands each request over, one at a time. +// Take binds to the role's worker and pulls one request at a time, handing each over. // // **Bound, never created.** The work queue and the worker on it are the controller's to define // (design 25 §3), and a build machine reaches no part of the JetStream API — so a missing one is said // as the mesh's to answer rather than quietly created with whatever this client defaults to. +// +// **Pulled, one at a time, by whichever holder is free** (novox/hq ADR 0190). Every machine holding +// the role binds this same worker; a machine asks for the next request only when it has finished +// the last, so a slow machine never holds an ask an idle one could take, and a machine that took +// five at once would run five container builds against one runtime and finish all of them slower +// than the first. func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) error { - worker, found := broker.HolderConsumerFor(m.on, "builder", + worker, found := broker.HolderConsumerFor(m.on, "build-agent", broker.DeclaredSeat{Name: m.seat, Accepts: []string{"build"}}) if !found { return fmt.Errorf("%s accepts no work, so there is nothing for this machine to take", m.seat) } - // One at a time, which the consumer's own ack-pending limit enforces rather than a prefetch - // setting: a machine that took five requests at once would run five container builds against one - // runtime and finish all of them slower than the first. - work := make(chan *nats.Msg, 1) // **The consumer's own filter, not the one subject this machine cares about.** The client checks // what is asked for against the consumer's filter and refuses anything that is not the same — // "subject does not match consumer" — so subscribing `…accept.build` against a consumer filtered // on `…accept.>` is rejected even though it is narrower. Learned twice now, on two different // consumers, which is why it is written down here. filter := worker.Filters[0] - sub, err := m.js.Context().ChanQueueSubscribe(filter, worker.Queue, work, + sub, err := m.js.Context().PullSubscribe(filter, worker.Name, nats.Bind(worker.Stream, worker.Name), nats.ManualAck()) if err != nil { return fmt.Errorf( @@ -159,13 +161,27 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) m.sub = sub for { - select { - case <-ctx.Done(): + if ctx.Err() != nil { return nil - case msg, ok := <-work: - if !ok { - return errors.New("the bus stopped delivering build work") + } + // One, and wait a while for it; an empty queue is a timeout, which is the normal state of a + // machine with nothing to build, and is asked again. + fetched, err := sub.Fetch(1, nats.Context(ctx)) + switch { + case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded): + return nil + case errors.Is(err, nats.ErrTimeout): + continue + case err != nil: + if sub.IsValid() { + // A transient fault in asking — a reconnect, a slow server — is asked past rather + // than ending the machine; one that outlasts the ack wait redelivers nothing lost. + time.Sleep(time.Second) + continue } + return fmt.Errorf("the bus stopped delivering build work: %w", err) + } + for _, msg := range fetched { var request BuildRequest if err := json.Unmarshal(msg.Data, &request); err != nil { // Unreadable: terminated rather than retried, because the next attempt reads the same -- 2.54.0 From 905f3363c9eddd9bd6b76df2ede268f9c677c375 Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 2 Oct 2026 22:35:39 +0200 Subject: [PATCH 3/5] Two machines holding the build role share one queue, and neither is handed an ask while busy (hq ADR 0190) Against a real bus: three asks, two machines; each takes one, the third waits until one is free and then goes to that one; a machine that stops leaves nothing taken twice. The redelivery of an ask a dead machine held is the ack wait's, proven by the hand-back test beside this one. And the order a live mesh switches over in, written where the role is named: queued builds first, then this controller, then build-agent assigned where machines build, then the builder and the old seat's stream forgotten. --- internal/link/builds.go | 7 +++ internal/link/builds_nats_test.go | 80 +++++++++++++++++++++++++++++++ 2 files changed, 87 insertions(+) diff --git a/internal/link/builds.go b/internal/link/builds.go index d71c3be..2f0b599 100644 --- a/internal/link/builds.go +++ b/internal/link/builds.go @@ -22,6 +22,13 @@ import ( // TheBuildMachine is the role a build is submitted to: node-scoped, held on every machine that // builds, and the work shared among them (novox/hq ADR 0190). The name stays for every caller; what // it names moved from the mesh's one build machine to whichever build agent is idle. +// +// **Switching a live mesh over, in order.** The old seat's stream and worker +// (SEAT_MESH_BUILD_MACHINE, SEAT_MESH_BUILD_MACHINE_worker) stay on the bus until removed by hand, +// and the builder module keeps draining them while it is assigned. From the moment a controller +// with this name runs, new asks go to node-build-agent and wait in its stream until some machine +// holds the seat. So: let the queued builds finish; roll this controller; register and assign +// build-agent to the machines that build; unassign builder and forget it and its seat's stream. const TheBuildMachine = "node-build-agent" // BuildWork is where a build request lands, and BuildOutcome is where its result does. Derived from diff --git a/internal/link/builds_nats_test.go b/internal/link/builds_nats_test.go index 38971a1..9bac7bd 100644 --- a/internal/link/builds_nats_test.go +++ b/internal/link/builds_nats_test.go @@ -245,3 +245,83 @@ func TestNatsWorkAMachineDidNotAnswerGoesBackToTheQueue(t *testing.T) { func quietLog() *log.Logger { return log.New(io.Discard, "", 0) } var _ = quietLog + +// Two machines holding the role share one queue (novox/hq ADR 0190): three asks, each machine takes +// one and the third waits until one of them is done; an ask is never handed to a machine that is +// busy; and a machine that stops mid-ask leaves its ask to the other. +func TestNatsTwoMachinesShareTheWorkAndNeitherIsHandedMoreThanItCanTake(t *testing.T) { + js := aBusWithTheBuildRole(t) + ctx, stop := context.WithCancel(context.Background()) + defer stop() + + for _, id := range []string{"w-1", "w-2", "w-3"} { + body, _ := json.Marshal(BuildRequest{ID: id, Repository: "/r"}) + if _, err := js.Context().Publish(BuildWork(), body); err != nil { + t.Fatal(err) + } + } + + type taken struct{ machine, id string } + took := make(chan taken, 8) + release := map[string]chan struct{}{"anchor": make(chan struct{}), "laptop": make(chan struct{})} + machines := map[string]BuildMachine{} + for _, name := range []string{"anchor", "laptop"} { + name := name + m := MachineOverNATS(js, name) + machines[name] = m + defer m.Close() + go func() { + _ = m.Take(ctx, func(ctx context.Context, work Build) { + took <- taken{name, work.Request().ID} + <-release[name] + _ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: name}) + _ = work.Done() + }) + }() + } + + // Each machine took exactly one, and they are different asks. + first := map[string]string{} + for i := 0; i < 2; i++ { + select { + case got := <-took: + if _, twice := first[got.machine]; twice { + t.Fatalf("%s was handed a second ask while busy with its first", got.machine) + } + first[got.machine] = got.id + case <-time.After(10 * time.Second): + t.Fatalf("only %d machine(s) took work; two idle holders should both have", len(first)) + } + } + if first["anchor"] == first["laptop"] { + t.Fatalf("both machines took %q: the queue is not shared, it is copied", first["anchor"]) + } + // The third waits: nobody is free. + select { + case got := <-took: + t.Fatalf("%s was handed %s while both machines were busy", got.machine, got.id) + case <-time.After(2 * time.Second): + } + // One finishes, and only then is the third taken — by that machine, the one that is free. + close(release["anchor"]) + release["anchor"] = make(chan struct{}) + select { + case got := <-took: + if got.machine != "anchor" { + t.Fatalf("the third ask went to %s, which is still busy", got.machine) + } + case <-time.After(10 * time.Second): + t.Fatal("the third ask was never taken after a machine became free") + } + // A machine that stops mid-ask leaves its ask unacknowledged, and the ack wait brings it round + // to whoever is left — the path TestNatsWorkAMachineDidNotAnswerGoesBackToTheQueue proves with + // an explicit hand-back, because the real wait is a minute. Here: the laptop goes, anchor + // finishes, and with nothing queued nothing more is taken by the machine that is left. + machines["laptop"].Close() + close(release["anchor"]) + select { + case got := <-took: + t.Fatalf("%s took %s; the queue should be empty", got.machine, got.id) + case <-time.After(2 * time.Second): + } +} -- 2.54.0 From a5d6a1187c7110fb738c005e9973bca795c482b6 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 02:47:49 +0200 Subject: [PATCH 4/5] A build machine serves the seat its credential claims (hq ADR 0190, the handover) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit After the build role moved to node-build-agent, nothing would hold it until build-agent is registered — and registering build-agent needs a build outcome that only the running builder could produce, bound as it was to the old seat by name. One binary, two roles: the seat a machine serves is the first its credential claims, as the mesh writes the claims beside the credential it issues (ADR 0159); the old builder keeps draining mesh-build-machine, a build-agent takes node-build-agent, and what each says about a build goes out as that seat's events, so an outcome is heard where the asker of that seat listens. A credential naming no claim serves the current role. --- cmd/mesh-builder/main.go | 22 +++++++++++++++++- cmd/mesh-builder/seat_test.go | 27 +++++++++++++++++++++++ internal/link/build_seat_test.go | 32 +++++++++++++++++++++++++++ internal/link/builds.go | 38 +++++++++++++++++++++++++++----- internal/link/builds_nats.go | 23 +++++++++++++------ 5 files changed, 129 insertions(+), 13 deletions(-) create mode 100644 cmd/mesh-builder/seat_test.go create mode 100644 internal/link/build_seat_test.go diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index f7d6946..0eb8215 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -137,7 +137,12 @@ func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) { if err != nil { return nil, err } - return link.MachineOverNATS(js, on), nil + // **The seat this machine serves is the one its credential claims** (novox/hq ADR 0190, the + // handover): the mesh issues a build machine's credential naming the seat its module claims, + // and one binary serves the old role as `builder` and the new as `build-agent` from that alone. + seat := link.BuildSeatClaimed(credential.seatsClaimed()) + fmt.Fprintf(os.Stderr, "taking build work as a holder of %s\n", seat) + return link.MachineOverNATSOn(js, on, seat), nil } // answer does one build and says what happened, whichever way it went. @@ -453,6 +458,21 @@ type Credential struct { // as two fields and this machine joins them once, here, to dial. User string `json:"user,omitempty"` Password string `json:"password,omitempty"` + // Claims are the seats the module this credential was issued for claims, as the mesh writes + // them beside the credential (novox/hq ADR 0159). The first is the build role this machine + // serves; a credential naming none is from before claims travelled in it. + Claims []struct { + Seat string `json:"seat"` + } `json:"claims,omitempty"` +} + +// seatsClaimed is the seats the credential names, in order. +func (c Credential) seatsClaimed() []string { + out := make([]string, 0, len(c.Claims)) + for _, claim := range c.Claims { + out = append(out, claim.Seat) + } + return out } // onTheNewBus is whether a credential is for the bus being built: its address says so, and the diff --git a/cmd/mesh-builder/seat_test.go b/cmd/mesh-builder/seat_test.go new file mode 100644 index 0000000..8518646 --- /dev/null +++ b/cmd/mesh-builder/seat_test.go @@ -0,0 +1,27 @@ +package main + +import ( + "encoding/json" + "testing" + + "github.com/novox/mesh-controller/internal/link" +) + +// The seat a build machine serves comes from its credential (novox/hq ADR 0190 handover). +func TestTheCredentialSaysWhichBuildRoleThisMachineServes(t *testing.T) { + var held Credential + if err := json.Unmarshal([]byte(`{"url":"nats://bus:4222","user":"anchor.builder","password":"x", + "claims":[{"seat":"mesh-build-machine","scope":"mesh","serves":[]}]}`), &held); err != nil { + t.Fatal(err) + } + if got := link.BuildSeatClaimed(held.seatsClaimed()); got != "mesh-build-machine" { + t.Errorf("the old builder's credential serves %q", got) + } + var bare Credential + if err := json.Unmarshal([]byte(`{"url":"nats://bus:4222","user":"anchor.build-agent","password":"x"}`), &bare); err != nil { + t.Fatal(err) + } + if got := link.BuildSeatClaimed(bare.seatsClaimed()); got != link.TheBuildMachine { + t.Errorf("a credential without claims serves %q, want %s", got, link.TheBuildMachine) + } +} diff --git a/internal/link/build_seat_test.go b/internal/link/build_seat_test.go new file mode 100644 index 0000000..4969798 --- /dev/null +++ b/internal/link/build_seat_test.go @@ -0,0 +1,32 @@ +package link + +import "testing" + +// A build machine serves the seat its credential claims (novox/hq ADR 0190 handover): the old +// `builder` keeps the old role, a `build-agent` takes the new, from one binary and no flag. +func TestABuildMachineServesTheSeatItsCredentialClaims(t *testing.T) { + if got := BuildSeatClaimed([]string{"mesh-build-machine"}); got != "mesh-build-machine" { + t.Errorf("a credential claiming the old role serves %q", got) + } + if got := BuildSeatClaimed([]string{"node-build-agent"}); got != TheBuildMachine { + t.Errorf("a credential claiming the new role serves %q", got) + } + if got := BuildSeatClaimed(nil); got != TheBuildMachine { + t.Errorf("a credential claiming nothing serves %q, want the current role", got) + } + if got := BuildSeatClaimed([]string{"", "node-build-agent"}); got != TheBuildMachine { + t.Errorf("an empty claim is skipped; got %q", got) + } +} + +// What a machine says about a build is the event of the seat it took the build from, so an outcome +// is heard where the asker of that seat listens. +func TestABuildsEventsAreItsSeats(t *testing.T) { + if BuildOutcomeOf(TheBuildMachineBefore) != "mesh.seat.mesh-build-machine.event.built" { + t.Error(BuildOutcomeOf(TheBuildMachineBefore)) + } + if BuildWorkOf(TheBuildMachine) != BuildWork() || BuildOutcomeOf(TheBuildMachine) != BuildOutcome() || + BuildStartedOf(TheBuildMachine) != BuildStarted() || BuildLogOf(TheBuildMachine, "b1") != BuildLog("b1") { + t.Error("the no-argument forms must name the current role") + } +} diff --git a/internal/link/builds.go b/internal/link/builds.go index 2f0b599..f43cf5b 100644 --- a/internal/link/builds.go +++ b/internal/link/builds.go @@ -31,10 +31,36 @@ import ( // build-agent to the machines that build; unassign builder and forget it and its seat's stream. const TheBuildMachine = "node-build-agent" +// TheBuildMachineBefore is the role a build was submitted to until ADR 0190: the mesh's one build +// machine, mesh-scoped. Kept named while the handover runs — a machine whose credential claims it +// still serves it, and the controller still hears its outcomes — and dropped with the retired seat +// row once nothing claims it. +const TheBuildMachineBefore = "mesh-build-machine" + +// BuildSeatClaimed is the build role a machine serves: the first seat its credential claims, or the +// current role when the credential names none (a credential from before claims travelled in it, or +// one written by hand). **The credential decides, not the binary** (ADR 0190 handover): one build +// machine binary runs as the old `builder` on the old seat and as a `build-agent` on the new one, +// each taking the work the mesh issued it a credential for, so neither drains the other's queue +// and the switch needs no flag day. +func BuildSeatClaimed(claimed []string) string { + for _, seat := range claimed { + if seat != "" { + return seat + } + } + return TheBuildMachine +} + // BuildWork is where a build request lands, and BuildOutcome is where its result does. Derived from -// the seat, so both sides name the role and neither names the other. -func BuildWork() string { return "mesh.seat." + TheBuildMachine + ".accept.build" } -func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.built" } +// the seat, so both sides name the role and neither names the other. The no-argument forms name the +// current role; the `Of` forms take the seat, for the handover during which two roles exist. +func BuildWork() string { return BuildWorkOf(TheBuildMachine) } +func BuildOutcome() string { return BuildOutcomeOf(TheBuildMachine) } +func BuildWorkOf(seat string) string { return "mesh.seat." + seat + ".accept.build" } +func BuildOutcomeOf(seat string) string { + return "mesh.seat." + seat + ".event.built" +} // BuildStarted is where a build machine says it has taken a build, and BuildLog is where it says // what it is doing, one line per message, under the build's own id (novox/hq ADR 0157). @@ -44,8 +70,10 @@ func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.bui // lived in one container's stderr on one machine. Every line is now an event of the role, retained // with the rest of the mesh's events, so a reader follows a build live by subscribing its subject, // or reads it back afterwards from the stream, and a viewer is a subscriber and nothing more. -func BuildStarted() string { return "mesh.seat." + TheBuildMachine + ".event.started" } -func BuildLog(id string) string { return "mesh.seat." + TheBuildMachine + ".event.log." + id } +func BuildStarted() string { return BuildStartedOf(TheBuildMachine) } +func BuildLog(id string) string { return BuildLogOf(TheBuildMachine, id) } +func BuildStartedOf(seat string) string { return "mesh.seat." + seat + ".event.started" } +func BuildLogOf(seat, id string) string { return "mesh.seat." + seat + ".event.log." + id } // BuildStart is what a build machine says the moment it takes a build. type BuildStart struct { diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index 0355857..9e158b6 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -115,9 +115,16 @@ type natsMachine struct { sub *nats.Subscription } -// MachineOverNATS takes build work from the role this machine holds. +// MachineOverNATS takes build work from the current build role. func MachineOverNATS(js *broker.JetStream, on string) BuildMachine { - return &natsMachine{js: js, on: on, seat: TheBuildMachine} + return MachineOverNATSOn(js, on, TheBuildMachine) +} + +// MachineOverNATSOn takes build work from the role named — the one this machine's credential claims +// (ADR 0190 handover): its asks come from that seat's worker, and what it says about a build goes +// out as that seat's events, so an outcome is heard where the asker listens. +func MachineOverNATSOn(js *broker.JetStream, on, seat string) BuildMachine { + return &natsMachine{js: js, on: on, seat: seat} } func (m *natsMachine) Close() { @@ -194,7 +201,7 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) // the ask to a second machine nor counts the wait against its deliveries. working := make(chan struct{}) go stillWorking(msg, working) - do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js}) + do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js, seat: m.seat}) close(working) } } @@ -205,7 +212,9 @@ type natsBuild struct { msg *nats.Msg on string js *broker.JetStream - seq int + // seat is the role this build was taken from; what the machine says about it is that role's. + seat string + seq int } func (b *natsBuild) Request() BuildRequest { return b.request } @@ -232,7 +241,7 @@ func (b *natsBuild) Announce(ctx context.Context, result BuildResult) error { } publish, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() - if _, err := b.js.Context().Publish(BuildOutcome(), body, nats.Context(publish)); err != nil { + if _, err := b.js.Context().Publish(BuildOutcomeOf(b.seat), body, nats.Context(publish)); err != nil { return fmt.Errorf("cannot announce a build's outcome: %w", err) } return nil @@ -252,7 +261,7 @@ func (b *natsBuild) Began(ctx context.Context) error { } publish, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() - if _, err := b.js.Context().Publish(BuildStarted(), body, nats.Context(publish)); err != nil { + if _, err := b.js.Context().Publish(BuildStartedOf(b.seat), body, nats.Context(publish)); err != nil { return fmt.Errorf("cannot say a build started: %w", err) } return nil @@ -270,7 +279,7 @@ func (b *natsBuild) Say(step, message string) { if err != nil { return } - _ = b.js.Conn().Publish(BuildLog(b.request.ID), body) + _ = b.js.Conn().Publish(BuildLogOf(b.seat, b.request.ID), body) } func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) } -- 2.54.0 From ff5ef0ab6005a618d183bba383d02cf5d95f5d9f Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 02:50:27 +0200 Subject: [PATCH 5/5] The controller asks the build role that has a holder, and hears both roles' outcomes (hq ADR 0190, the handover) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A controller that asked node-build-agent from its first run would queue every build where nothing pulls, and the build that registers build-agent — the first holder — would be among them. So the role is chosen at ask time from the catalogue: the current role when any assigned module claims it, the retired one while only the builder does, the current one when neither. Outcomes are followed on both seats, the controller may publish to both, and a build's log is read under whichever role did it; a machine on the retired role is proven on the bus to take that role's asks. The switch order is written where the role is named, and the retired half is marked for removal with the seat row. --- cmd/mesh-controller/build.go | 61 +++++++++++++++++++++++--- cmd/mesh-controller/build_seat_test.go | 43 ++++++++++++++++++ internal/broker/nats.go | 5 ++- internal/broker/streams.go | 5 +++ internal/broker/testdata/composed.conf | 4 +- internal/link/builds.go | 21 ++++++--- internal/link/builds_nats.go | 30 ++++++++++--- internal/link/builds_nats_test.go | 36 +++++++++++++++ internal/link/receive_nats.go | 5 ++- 9 files changed, 187 insertions(+), 23 deletions(-) create mode 100644 cmd/mesh-controller/build_seat_test.go diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 60695c3..56f7afb 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -430,11 +430,13 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti } fmt.Println() - ask, err := askOver(server) + seat := buildSeatHeld(ctx) + ask, err := askOverOn(seat) if err != nil { return err } defer ask.Close() + fmt.Printf(" of %s\n", seat) if wait == 0 { // Asked and not waited for (novox/hq issue 176): the outcome is the role's event, and the @@ -555,7 +557,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai } defer server.Close() - ask, err := askOver(server) + ask, err := askOverOn(buildSeatHeld(ctx)) if err != nil { return err } @@ -668,12 +670,56 @@ 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(_ *link.Server) (link.Builders, error) { +func askOverOn(seat string) (link.Builders, error) { address, err := broker.BusAddress() if err != nil { return nil, err } - return link.BuildsOverNATS(address) + return link.BuildsOverNATSOn(address, seat) +} + +// buildSeatHeld is the build role to ask: the one some assigned module claims (novox/hq ADR 0190, +// the handover). Read from the catalogue at ask time, because the answer changes exactly once, the +// moment the first build-agent is assigned — and a controller that asked the new role before then +// would queue work nothing takes, while the outcome that registers build-agent itself has to come +// from the old builder. When the catalogue cannot be read the current role is asked, said aloud. +func buildSeatHeld(ctx context.Context) string { + open, err := openStores(ctx) + if err != nil { + fmt.Fprintf(os.Stderr, "could not read what is assigned, so the build is asked of %s: %v\n", + link.TheBuildMachine, err) + return link.TheBuildMachine + } + defer open.Close() + entries, err := open.inventory.Catalogued(ctx) + if err != nil { + fmt.Fprintf(os.Stderr, "could not read the catalogue, so the build is asked of %s: %v\n", + link.TheBuildMachine, err) + return link.TheBuildMachine + } + return buildSeatAmong(entries) +} + +// buildSeatAmong is the rule, over what the catalogue holds: the current build role when any +// assigned module claims it; else the retired role while an assigned module still claims that; else +// the current role, which is where every ask goes once the handover is done. +func buildSeatAmong(entries []inventory.Entry) string { + heldBefore := false + for _, e := range entries { + if len(e.On) == 0 { + continue + } + if e.Manifest.ClaimsSeat(link.TheBuildMachine) { + return link.TheBuildMachine + } + if e.Manifest.ClaimsSeat(link.TheBuildMachineBefore) { + heldBefore = true + } + } + if heldBefore { + return link.TheBuildMachineBefore + } + return link.TheBuildMachine } // buildLog prints everything a build machine said about one build, read back from the bus. @@ -693,10 +739,13 @@ func buildLog(ctx context.Context, id string) error { } defer js.Close() - sub, err := js.Context().PullSubscribe(link.BuildLog(id), "", + // Under whichever build role did it: a build asked of the retired role during the handover + // (ADR 0190) said its lines as that role's events, and a reader should not have to know which. + lines := link.BuildLogOf("*", id) + sub, err := js.Context().PullSubscribe(lines, "", nats.BindStream(broker.EventsStream), nats.DeliverAll(), nats.AckNone()) if err != nil { - return fmt.Errorf("cannot read %s from the bus: %w", link.BuildLog(id), err) + return fmt.Errorf("cannot read %s from the bus: %w", lines, err) } defer func() { _ = sub.Unsubscribe() }() diff --git a/cmd/mesh-controller/build_seat_test.go b/cmd/mesh-controller/build_seat_test.go new file mode 100644 index 0000000..e40ca0a --- /dev/null +++ b/cmd/mesh-controller/build_seat_test.go @@ -0,0 +1,43 @@ +package main + +import ( + "testing" + + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +func claiming(module, seat string, on ...string) inventory.Entry { + return inventory.Entry{ + Manifest: catalogue.Manifest{Module: module, Claims: []catalogue.Claim{{Name: seat}}}, + On: on, + } +} + +// The controller asks the build role that has a holder (novox/hq ADR 0190 handover): the retired +// one while only the builder is assigned, the current one from the first build-agent on, and the +// current one when nothing holds either — where every ask goes once the handover is done. +func TestTheControllerAsksTheBuildRoleThatHasAHolder(t *testing.T) { + onlyTheBuilder := []inventory.Entry{ + claiming("builder", link.TheBuildMachineBefore, "anchor"), + claiming("build-agent", link.TheBuildMachine), // registered, assigned nowhere yet + } + if got := buildSeatAmong(onlyTheBuilder); got != link.TheBuildMachineBefore { + t.Errorf("with only the builder assigned, asked %q", got) + } + bothHeld := []inventory.Entry{ + claiming("builder", link.TheBuildMachineBefore, "anchor"), + claiming("build-agent", link.TheBuildMachine, "home-server"), + } + if got := buildSeatAmong(bothHeld); got != link.TheBuildMachine { + t.Errorf("with a build-agent assigned anywhere, asked %q", got) + } + neither := []inventory.Entry{claiming("builder", link.TheBuildMachineBefore)} + if got := buildSeatAmong(neither); got != link.TheBuildMachine { + t.Errorf("with no holder of either, asked %q, want the current role", got) + } + if got := buildSeatAmong(nil); got != link.TheBuildMachine { + t.Errorf("an empty catalogue asks %q", got) + } +} diff --git a/internal/broker/nats.go b/internal/broker/nats.go index b52cb34..025bcd5 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -113,7 +113,10 @@ type Principal struct { // seatsTheControllerAsks are the roles the mesh's own flows submit work to. Named rather than // derived from the seat set: the controller is not a module and declares no `uses`, so its side of a // seat has to be stated, and a list is what makes "which roles does the mesh itself talk to" answerable. -var seatsTheControllerAsks = []string{"node-build-agent"} +// Both build roles while the handover runs (novox/hq ADR 0190): the controller asks whichever has a +// holder, and the retired one has one until build-agent replaces the builder. The second entry +// goes with the retired seat row. +var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"} // enrolmentPrefix is the space every enrolling node's user and inbox live under, so the one place the // controller may answer an enrolment is derived from the same constant the user is named from. diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 210359b..526cd65 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -221,6 +221,11 @@ var ControllerFollows = []string{ // The forge's merges: what moved a source, so the mesh builds what that source produces // without anybody telling it (novox/hq 04-ISSUES/131). Appended, because the index is a name. moduleEventSubject("gitea", "pull.merged"), + // The retired build role's outcome too, while the handover runs (novox/hq ADR 0190): the one + // build machine keeps answering on its seat until build-agent replaces it, and the outcome that + // registers build-agent itself comes from there. Appended, for the same reason as above; goes + // with the retired seat row. + seatEventSubject("mesh-build-machine", "built"), } // moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 3994266..ba18138 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,8 +24,8 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>"] } - subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>"] } + subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/link/builds.go b/internal/link/builds.go index f43cf5b..5757671 100644 --- a/internal/link/builds.go +++ b/internal/link/builds.go @@ -23,12 +23,21 @@ import ( // builds, and the work shared among them (novox/hq ADR 0190). The name stays for every caller; what // it names moved from the mesh's one build machine to whichever build agent is idle. // -// **Switching a live mesh over, in order.** The old seat's stream and worker -// (SEAT_MESH_BUILD_MACHINE, SEAT_MESH_BUILD_MACHINE_worker) stay on the bus until removed by hand, -// and the builder module keeps draining them while it is assigned. From the moment a controller -// with this name runs, new asks go to node-build-agent and wait in its stream until some machine -// holds the seat. So: let the queued builds finish; roll this controller; register and assign -// build-agent to the machines that build; unassign builder and forget it and its seat's stream. +// **Switching a live mesh over, in order** — and why no step strands a build. The old seat's +// stream and worker (SEAT_MESH_BUILD_MACHINE, SEAT_MESH_BUILD_MACHINE_worker) stay on the bus until +// removed by hand, and the builder keeps draining them while it is assigned, because a machine +// serves the seat its credential claims (BuildSeatClaimed) and the controller asks the seat that +// has a holder (buildSeatAmong in the command) and hears both seats' outcomes: +// +// 1. Merge the controller and the host's first user list together; the new controller rolls and, +// seeing only the builder assigned, still asks mesh-build-machine — which the builder holds. +// 2. Merge the catalogue's build-agent; the builder builds it and the controller registers it. +// 3. On each machine that builds: `module issue build-agent --node `, then `assign`, then +// `push`. The first holder appears, and from then on asks go to node-build-agent. +// 4. Unassign builder everywhere and `module forget` it. +// 5. By hand: delete SEAT_MESH_BUILD_MACHINE and its worker, drop the retired seat row and +// TheBuildMachineBefore with it, and the second entries in seatsTheControllerAsks and +// ControllerFollows. const TheBuildMachine = "node-build-agent" // TheBuildMachineBefore is the role a build was submitted to until ADR 0190: the mesh's one build diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index 9e158b6..c3b1011 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -25,16 +25,34 @@ import ( type natsBuilds struct { js *broker.JetStream owned bool + // seat is the build role asked: the one that has a holder (ADR 0190 handover), chosen by the + // controller from what is assigned, so an ask lands where a machine is pulling. + seat string } -// BuildsOverNATS is the asking side on the bus being built. It dials, because the command that asks -// for a build is a one-shot and holds nothing else. +// BuildsOverNATS is the asking side on the bus being built, asking the current build role. It dials, +// because the command that asks for a build is a one-shot and holds nothing else. func BuildsOverNATS(address string) (Builders, error) { + return BuildsOverNATSOn(address, TheBuildMachine) +} + +// BuildsOverNATSOn is the asking side for one named build role — during the handover from the one +// build machine to build agents, the role that has a holder (ADR 0190). +func BuildsOverNATSOn(address, seat string) (Builders, error) { js, err := broker.Dial(address) if err != nil { return nil, fmt.Errorf("cannot reach the bus at %s to ask for a build: %w", address, err) } - return &natsBuilds{js: js, owned: true}, nil + return &natsBuilds{js: js, owned: true, seat: seat}, nil +} + +// role is the seat asked: what the asker was made for, or the current build role for one made +// without saying (a test building the struct by hand). +func (b *natsBuilds) role() string { + if b.seat == "" { + return TheBuildMachine + } + return b.seat } func (b *natsBuilds) Close() { @@ -51,7 +69,7 @@ func (b *natsBuilds) Ask(ctx context.Context, request BuildRequest) error { } publish, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() - if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil { + if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil { return fmt.Errorf("cannot submit a build: %w", err) } return nil @@ -63,7 +81,7 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest, // Subscribed before the ask, so an outcome cannot arrive before there is anywhere for it to // land. Core, not the stream: the asker is waiting now, and the durable copy of this outcome is // the same event on EVENTS, which the controller records. - outcomes, err := b.js.Conn().SubscribeSync(BuildOutcome()) + outcomes, err := b.js.Conn().SubscribeSync(BuildOutcomeOf(b.role())) if err != nil { return BuildResult{}, fmt.Errorf("cannot listen for a build's outcome: %w", err) } @@ -80,7 +98,7 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest, // be assumed, because nothing else will ever say so. publish, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() - if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil { + if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil { return BuildResult{}, fmt.Errorf("cannot submit a build: %w", err) } diff --git a/internal/link/builds_nats_test.go b/internal/link/builds_nats_test.go index 9bac7bd..b0feaa2 100644 --- a/internal/link/builds_nats_test.go +++ b/internal/link/builds_nats_test.go @@ -325,3 +325,39 @@ func TestNatsTwoMachinesShareTheWorkAndNeitherIsHandedMoreThanItCanTake(t *testi case <-time.After(2 * time.Second): } } + +// During the handover (ADR 0190) two build roles exist. A machine whose credential claims the retired +// one takes an ask published to that seat and answers as that seat; the asker of that seat hears it. +func TestNatsAMachineOnTheRetiredBuildRoleTakesThatRolesAsks(t *testing.T) { + js := aBusWithTheBuildRole(t) + seats := []broker.DeclaredSeat{{Name: TheBuildMachineBefore, Accepts: []string{"build"}, Emits: []string{"built"}}} + if err := broker.RaiseSeats(js, seats, map[string]broker.Holder{TheBuildMachineBefore: {Node: "anchor", Module: "builder"}}); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") }) + ctx, stop := context.WithCancel(context.Background()) + defer stop() + + machine := MachineOverNATSOn(js, "anchor", TheBuildMachineBefore) + defer machine.Close() + go func() { + _ = machine.Take(ctx, func(ctx context.Context, work Build) { + _ = work.Began(ctx) + _ = work.Announce(ctx, BuildResult{ID: work.Request().ID, Repository: work.Request().Repository, On: "anchor", Commit: "abc"}) + _ = work.Done() + }) + }() + + asker, err := BuildsOverNATSOn(os.Getenv("MESH_TEST_NATS"), TheBuildMachineBefore) + if err != nil { + t.Fatal(err) + } + defer asker.Close() + result, err := asker.Submit(ctx, BuildRequest{ID: "build-old-seat", Repository: "r"}, 20*time.Second) + if err != nil { + t.Fatal(err) + } + if result.On != "anchor" || result.ID != "build-old-seat" { + t.Errorf("the retired role's holder did not answer: %+v", result) + } +} diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go index 0d812b8..a1889ee 100644 --- a/internal/link/receive_nats.go +++ b/internal/link/receive_nats.go @@ -223,10 +223,11 @@ func kindOfSubject(subject string) (string, bool) { return KindCatchUp, true case broker.ControllerFollows[3]: return KindSourceMoved, true - case BuildOutcome(): + case BuildOutcome(), BuildOutcomeOf(TheBuildMachineBefore): // A build's outcome is the role's event now, so it arrives on the events stream rather than // the control branch — and is acted on by the same handler, because what the controller does - // with it did not change (novox/hq ADR 0121). + // with it did not change (novox/hq ADR 0121). From either build role while the handover + // runs (ADR 0190): the old builder still answers on the retired seat until it is unassigned. return KindBuilt, true } return "", false -- 2.54.0