Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
905f3363c9 | ||
|
|
9f9d9b3b25 | ||
|
|
bde4b61b3b |
+16
-15
@@ -144,12 +144,18 @@ func ConsumerFor(p Principal) (Consumer, bool) {
|
|||||||
}, true
|
}, 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
|
// **One worker for every holder, and each holder pulls one ask when it is idle** (novox/hq ADR
|
||||||
// be the telegram sender — and the queue group is *delivery*. Tie delivery to the seat and the
|
// 0190). The seat is *authority* — who may be the telegram sender — and the worker is *delivery*,
|
||||||
// day somebody allows two holders for throughput, every message is processed twice with nothing
|
// kept separate so that relaxing one changes nothing about the other: a node-scoped seat has a
|
||||||
// reporting it. Kept separate, relaxing one changes nothing about the other.
|
// 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) {
|
func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool) {
|
||||||
if len(seat.Accepts) == 0 {
|
if len(seat.Accepts) == 0 {
|
||||||
return Consumer{}, false
|
return Consumer{}, false
|
||||||
@@ -158,18 +164,13 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool)
|
|||||||
Name: "SEAT_" + upperSnake(seat.Name) + "_worker",
|
Name: "SEAT_" + upperSnake(seat.Name) + "_worker",
|
||||||
Stream: seatStreamName(seat.Name),
|
Stream: seatStreamName(seat.Name),
|
||||||
Filters: []string{"mesh.seat." + seat.Name + ".accept.>"},
|
Filters: []string{"mesh.seat." + seat.Name + ".accept.>"},
|
||||||
Queue: "holders",
|
|
||||||
AckWaitSeconds: 60,
|
AckWaitSeconds: 60,
|
||||||
MaxDeliver: 5,
|
MaxDeliver: 5,
|
||||||
// **One in flight.** A holder works one ask at a time, so the server hands it one at a
|
// As many in flight as there are holders working, which pulling bounds by itself: a holder
|
||||||
// time: with the default of many, every ask behind the one being worked was delivered,
|
// fetches one and fetches again only after it acknowledged. The server's default stands.
|
||||||
// left unacknowledged for the length of the work, redelivered after the ack wait, and
|
Why: fmt.Sprintf("%s on %s holds %s; every holder pulls one ask at a time from this worker "+
|
||||||
// after the fifth time dropped — on 2026-10-01 twenty-six of forty-three builds asked in
|
"and acknowledges after the work is done, so a crash mid-work redelivers rather than "+
|
||||||
// two minutes were never built, and the queue read as empty (novox/hq issue 186).
|
"loses and an idle holder is the one that takes the next ask", module, node, seat.Name),
|
||||||
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),
|
|
||||||
}, true
|
}, 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
|
// The seat is authority and the worker is delivery (novox/hq ADR 0190): one worker per seat, shared
|
||||||
// allows two holders, every message is processed twice with nothing reporting it.
|
// by every holder and pulled from, so a second holder takes the next ask rather than a copy of the
|
||||||
func TestAHoldersWorkerUsesAQueueGroupAnyway(t *testing.T) {
|
// 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())
|
c, ok := HolderConsumerFor("one", "telegram", telegramSeat())
|
||||||
if !ok {
|
if !ok {
|
||||||
t.Fatal("the holder of a seat with inbound work got no worker")
|
t.Fatal("the holder of a seat with inbound work got no worker")
|
||||||
}
|
}
|
||||||
if c.Queue == "" {
|
two, _ := HolderConsumerFor("two", "telegram", telegramSeat())
|
||||||
t.Fatal("the worker is not in a queue group, so a second holder would double-process")
|
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" {
|
if c.Stream != "SEAT_TELEGRAM_SENDER" {
|
||||||
t.Fatalf("the worker reads %q, not the seat's own stream", c.Stream)
|
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])
|
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):
|
// Every holder of a seat shares one worker and pulls from it (novox/hq ADR 0190): no queue group
|
||||||
// asks queued behind the one being worked wait in the stream rather than being delivered,
|
// and no delivery subject, because a push consumer hands the next ask to whichever subscriber the
|
||||||
// left to expire and dropped after the fifth redelivery.
|
// server picks, busy or not; and no cap of one in flight, because pulling bounds the asks in flight
|
||||||
func TestAHoldersWorkerTakesOneAskAtATime(t *testing.T) {
|
// by the holders that are free — which is what ended the race of issue 186, where asks delivered
|
||||||
c, found := HolderConsumerFor("anchor", "builder", DeclaredSeat{Name: "mesh-build-machine", Accepts: []string{"build"}})
|
// 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 {
|
if !found {
|
||||||
t.Fatal("a seat that accepts work has no worker")
|
t.Fatal("a seat that accepts work has no worker")
|
||||||
}
|
}
|
||||||
if c.MaxAckPending != 1 {
|
if c.Queue != "" || c.Push {
|
||||||
t.Fatalf("the worker may have %d asks in flight; one, so a queue is a queue", c.MaxAckPending)
|
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
|
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
|
// 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.
|
// 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
|
// 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.
|
// 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
|
// 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
|
// 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.
|
// 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.>")
|
pub = append(pub, "mesh.seat."+seat+".accept.>")
|
||||||
}
|
}
|
||||||
// **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own
|
// **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.
|
// 3. Seats it holds: full participation.
|
||||||
for _, s := range p.Holds {
|
for _, s := range p.Holds {
|
||||||
// Taking work from the role's queue: the worker consumer it binds (asked about,
|
// Taking work from the role's queue: the worker consumer every holder shares (asked
|
||||||
// delivered on, acknowledged), each on the seat's own stream. The first machine to
|
// about, pulled from, acknowledged), on the seat's own stream (novox/hq ADR 0190). A
|
||||||
// take work over the new bus was refused the asking (2026-09-28).
|
// 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"
|
worker := "SEAT_" + upperSnake(s.Name) + "_worker"
|
||||||
stream := seatStreamName(s.Name)
|
stream := seatStreamName(s.Name)
|
||||||
sub = append(sub, "_DELIVER."+worker, "_DELIVER."+worker+".>")
|
pub = append(pub,
|
||||||
pub = append(pub, "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>")
|
"$JS.API.CONSUMER.INFO."+stream+"."+worker,
|
||||||
|
"$JS.API.CONSUMER.MSG.NEXT."+stream+"."+worker,
|
||||||
|
"$JS.ACK."+stream+"."+worker+".>")
|
||||||
for _, a := range s.Accepts {
|
for _, a := range s.Accepts {
|
||||||
sub = append(sub, seatSubject(s, "accept", a))
|
sub = append(sub, seatSubject(s, "accept", a))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -441,3 +441,24 @@ func contains(list []string, want string) bool {
|
|||||||
}
|
}
|
||||||
return false
|
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
|
// 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
|
// message on the control branch. Same three audiences, one publish: whoever asked, this, and the
|
||||||
// catalogue.
|
// 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
|
// 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.
|
// without anybody telling it (novox/hq 04-ISSUES/131). Appended, because the index is a name.
|
||||||
moduleEventSubject("gitea", "pull.merged"),
|
moduleEventSubject("gitea", "pull.merged"),
|
||||||
|
|||||||
+4
-4
@@ -24,8 +24,8 @@ accounts {
|
|||||||
jetstream: enabled
|
jetstream: enabled
|
||||||
users = [
|
users = [
|
||||||
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
|
{ 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"] }
|
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-build-machine.event.built", "mesh.seat.mesh-controller.tool.>"] }
|
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" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
|
{ 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"] }
|
subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] }
|
||||||
} }
|
} }
|
||||||
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
|
{ 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"] }
|
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: ["_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"] }
|
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" }
|
allow_responses: { max: 1, ttl: "1m" }
|
||||||
} }
|
} }
|
||||||
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
|
{ 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,
|
// 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
|
// `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.
|
// 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,
|
{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"},
|
{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
|
// 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
|
// 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
|
delivered[s.Delivers] = s.Name
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// Sixteen since node-service-manager (novox/hq ADR 0177).
|
// Seventeen since node-build-agent (novox/hq ADR 0190) — sixteen once the retired
|
||||||
if len(Seats()) != 16 {
|
// mesh-build-machine row goes, when no registered manifest claims it any more.
|
||||||
t.Errorf("the mesh defines %d seats rather than 16; the set is closed, so a change here is "+
|
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())
|
"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 != "" {
|
if r := repositoryKey(e.Source.Repository); r != "" {
|
||||||
byRepository[r] = append(byRepository[r], name)
|
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)
|
builders = append(builders, name)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,7 +12,7 @@ func TestDependenciesAreOneRelationWithTheirKinds(t *testing.T) {
|
|||||||
return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}}
|
return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}}
|
||||||
}
|
}
|
||||||
builder := entry("builder", "http://forge/novox/mesh-catalog.git")
|
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 := entry("shop-plugin", "http://forge/novox/mesh-catalog.git")
|
||||||
plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}}
|
plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}}
|
||||||
entries := []Entry{
|
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
|
// 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.
|
// exactly one answer. Too long for request/reply, too particular to be an event.
|
||||||
|
|
||||||
// TheBuildMachine is the role a build is submitted to.
|
// TheBuildMachine is the role a build is submitted to: node-scoped, held on every machine that
|
||||||
const TheBuildMachine = "mesh-build-machine"
|
// 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
|
// 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.
|
// 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" {
|
if KeyRoleBuilt != "built" {
|
||||||
t.Fatalf("the role's event is %q, and a holder emits its verbs bare", KeyRoleBuilt)
|
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)
|
t.Fatalf("the role is %q", TheBuildMachine)
|
||||||
}
|
}
|
||||||
// The two must differ, or one publish would serve both and this doubling would be pointless.
|
// 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
|
// **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
|
// (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.
|
// 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 {
|
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"}})
|
broker.DeclaredSeat{Name: m.seat, Accepts: []string{"build"}})
|
||||||
if !found {
|
if !found {
|
||||||
return fmt.Errorf("%s accepts no work, so there is nothing for this machine to take", m.seat)
|
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
|
// **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 —
|
// 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
|
// "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
|
// on `…accept.>` is rejected even though it is narrower. Learned twice now, on two different
|
||||||
// consumers, which is why it is written down here.
|
// consumers, which is why it is written down here.
|
||||||
filter := worker.Filters[0]
|
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())
|
nats.Bind(worker.Stream, worker.Name), nats.ManualAck())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf(
|
return fmt.Errorf(
|
||||||
@@ -159,13 +161,27 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
|
|||||||
m.sub = sub
|
m.sub = sub
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
if ctx.Err() != nil {
|
||||||
case <-ctx.Done():
|
|
||||||
return 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
|
var request BuildRequest
|
||||||
if err := json.Unmarshal(msg.Data, &request); err != nil {
|
if err := json.Unmarshal(msg.Data, &request); err != nil {
|
||||||
// Unreadable: terminated rather than retried, because the next attempt reads the same
|
// 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)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
clean := func() {
|
clean := func() {
|
||||||
_ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE")
|
_ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT")
|
||||||
for _, s := range broker.MeshStreams() {
|
for _, s := range broker.MeshStreams() {
|
||||||
_ = js.Context().PurgeStream(s.Name)
|
_ = 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.
|
// 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)
|
deadline := time.Now().Add(5 * time.Second)
|
||||||
for time.Now().Before(deadline) {
|
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 {
|
if err == nil && info.State.Msgs == 0 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -177,7 +177,7 @@ func TestNatsABuildWaitsForAMachineRatherThanFailing(t *testing.T) {
|
|||||||
if _, err := js.Context().Publish(BuildWork(), body); err != nil {
|
if _, err := js.Context().Publish(BuildWork(), body); err != nil {
|
||||||
t.Fatal(err)
|
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 {
|
if err != nil || info.State.Msgs != 1 {
|
||||||
t.Fatalf("the work did not queue: %+v %v", info, err)
|
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) }
|
func quietLog() *log.Logger { return log.New(io.Discard, "", 0) }
|
||||||
|
|
||||||
var _ = quietLog
|
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