Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
905f3363c9 | ||
|
|
9f9d9b3b25 | ||
|
|
bde4b61b3b |
+16
-15
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
@@ -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))
|
||||
}
|
||||
|
||||
@@ -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.>")
|
||||
}
|
||||
|
||||
@@ -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
@@ -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: {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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):
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user