Compare commits

...
3 Commits
Author SHA1 Message Date
jochen 905f3363c9 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.
2026-10-02 22:35:39 +02:00
jochen 9f9d9b3b25 A seat's holders pull one ask at a time from one shared worker (hq ADR 0190, issue 186)
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.
2026-10-02 22:34:44 +02:00
jochen bde4b61b3b The build role is the node-scoped seat node-build-agent, and its work is shared by every holder (hq ADR 0190)
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.
2026-10-02 22:31:55 +02:00
14 changed files with 220 additions and 64 deletions
+16 -15
View File
@@ -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
}
+25 -12
View File
@@ -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)
}
}
+14 -8
View File
@@ -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
@@ -368,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))
}
+21
View File
@@ -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.>")
}
+1 -1
View File
@@ -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"),
+4 -4
View File
@@ -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: {
@@ -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: {
+10 -1
View File
@@ -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.<build id>` 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
+4 -3
View File
@@ -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())
}
}
+1 -1
View File
@@ -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)
}
}
+1 -1
View File
@@ -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{
+11 -2
View File
@@ -19,8 +19,17 @@ 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.
//
// **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
// the seat, so both sides name the role and neither names the other.
+1 -1
View File
@@ -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.
+28 -12
View File
@@ -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
+83 -3
View File
@@ -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)
}
@@ -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):
}
}