Merge pull request 'A seat's work is shared by its holders: node-build-agent, pulled one ask at a time (hq ADR 0190)' (#228) from feat/a-seats-work-is-shared-by-its-holders into main
This commit was merged in pull request #228.
This commit is contained in:
@@ -137,7 +137,12 @@ func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return link.MachineOverNATS(js, on), nil
|
||||
// **The seat this machine serves is the one its credential claims** (novox/hq ADR 0190, the
|
||||
// handover): the mesh issues a build machine's credential naming the seat its module claims,
|
||||
// and one binary serves the old role as `builder` and the new as `build-agent` from that alone.
|
||||
seat := link.BuildSeatClaimed(credential.seatsClaimed())
|
||||
fmt.Fprintf(os.Stderr, "taking build work as a holder of %s\n", seat)
|
||||
return link.MachineOverNATSOn(js, on, seat), nil
|
||||
}
|
||||
|
||||
// answer does one build and says what happened, whichever way it went.
|
||||
@@ -453,6 +458,21 @@ type Credential struct {
|
||||
// as two fields and this machine joins them once, here, to dial.
|
||||
User string `json:"user,omitempty"`
|
||||
Password string `json:"password,omitempty"`
|
||||
// Claims are the seats the module this credential was issued for claims, as the mesh writes
|
||||
// them beside the credential (novox/hq ADR 0159). The first is the build role this machine
|
||||
// serves; a credential naming none is from before claims travelled in it.
|
||||
Claims []struct {
|
||||
Seat string `json:"seat"`
|
||||
} `json:"claims,omitempty"`
|
||||
}
|
||||
|
||||
// seatsClaimed is the seats the credential names, in order.
|
||||
func (c Credential) seatsClaimed() []string {
|
||||
out := make([]string, 0, len(c.Claims))
|
||||
for _, claim := range c.Claims {
|
||||
out = append(out, claim.Seat)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// onTheNewBus is whether a credential is for the bus being built: its address says so, and the
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// The seat a build machine serves comes from its credential (novox/hq ADR 0190 handover).
|
||||
func TestTheCredentialSaysWhichBuildRoleThisMachineServes(t *testing.T) {
|
||||
var held Credential
|
||||
if err := json.Unmarshal([]byte(`{"url":"nats://bus:4222","user":"anchor.builder","password":"x",
|
||||
"claims":[{"seat":"mesh-build-machine","scope":"mesh","serves":[]}]}`), &held); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := link.BuildSeatClaimed(held.seatsClaimed()); got != "mesh-build-machine" {
|
||||
t.Errorf("the old builder's credential serves %q", got)
|
||||
}
|
||||
var bare Credential
|
||||
if err := json.Unmarshal([]byte(`{"url":"nats://bus:4222","user":"anchor.build-agent","password":"x"}`), &bare); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := link.BuildSeatClaimed(bare.seatsClaimed()); got != link.TheBuildMachine {
|
||||
t.Errorf("a credential without claims serves %q, want %s", got, link.TheBuildMachine)
|
||||
}
|
||||
}
|
||||
@@ -430,11 +430,13 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
|
||||
}
|
||||
fmt.Println()
|
||||
|
||||
ask, err := askOver(server)
|
||||
seat := buildSeatHeld(ctx)
|
||||
ask, err := askOverOn(seat)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer ask.Close()
|
||||
fmt.Printf(" of %s\n", seat)
|
||||
|
||||
if wait == 0 {
|
||||
// Asked and not waited for (novox/hq issue 176): the outcome is the role's event, and the
|
||||
@@ -555,7 +557,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai
|
||||
}
|
||||
defer server.Close()
|
||||
|
||||
ask, err := askOver(server)
|
||||
ask, err := askOverOn(buildSeatHeld(ctx))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -668,12 +670,56 @@ func heldBy(ctx context.Context) map[string]string {
|
||||
// **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5). On the bus
|
||||
// the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus
|
||||
// being built it dials, because a build request is a one-shot and holds nothing else.
|
||||
func askOver(_ *link.Server) (link.Builders, error) {
|
||||
func askOverOn(seat string) (link.Builders, error) {
|
||||
address, err := broker.BusAddress()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return link.BuildsOverNATS(address)
|
||||
return link.BuildsOverNATSOn(address, seat)
|
||||
}
|
||||
|
||||
// buildSeatHeld is the build role to ask: the one some assigned module claims (novox/hq ADR 0190,
|
||||
// the handover). Read from the catalogue at ask time, because the answer changes exactly once, the
|
||||
// moment the first build-agent is assigned — and a controller that asked the new role before then
|
||||
// would queue work nothing takes, while the outcome that registers build-agent itself has to come
|
||||
// from the old builder. When the catalogue cannot be read the current role is asked, said aloud.
|
||||
func buildSeatHeld(ctx context.Context) string {
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "could not read what is assigned, so the build is asked of %s: %v\n",
|
||||
link.TheBuildMachine, err)
|
||||
return link.TheBuildMachine
|
||||
}
|
||||
defer open.Close()
|
||||
entries, err := open.inventory.Catalogued(ctx)
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "could not read the catalogue, so the build is asked of %s: %v\n",
|
||||
link.TheBuildMachine, err)
|
||||
return link.TheBuildMachine
|
||||
}
|
||||
return buildSeatAmong(entries)
|
||||
}
|
||||
|
||||
// buildSeatAmong is the rule, over what the catalogue holds: the current build role when any
|
||||
// assigned module claims it; else the retired role while an assigned module still claims that; else
|
||||
// the current role, which is where every ask goes once the handover is done.
|
||||
func buildSeatAmong(entries []inventory.Entry) string {
|
||||
heldBefore := false
|
||||
for _, e := range entries {
|
||||
if len(e.On) == 0 {
|
||||
continue
|
||||
}
|
||||
if e.Manifest.ClaimsSeat(link.TheBuildMachine) {
|
||||
return link.TheBuildMachine
|
||||
}
|
||||
if e.Manifest.ClaimsSeat(link.TheBuildMachineBefore) {
|
||||
heldBefore = true
|
||||
}
|
||||
}
|
||||
if heldBefore {
|
||||
return link.TheBuildMachineBefore
|
||||
}
|
||||
return link.TheBuildMachine
|
||||
}
|
||||
|
||||
// buildLog prints everything a build machine said about one build, read back from the bus.
|
||||
@@ -693,10 +739,13 @@ func buildLog(ctx context.Context, id string) error {
|
||||
}
|
||||
defer js.Close()
|
||||
|
||||
sub, err := js.Context().PullSubscribe(link.BuildLog(id), "",
|
||||
// Under whichever build role did it: a build asked of the retired role during the handover
|
||||
// (ADR 0190) said its lines as that role's events, and a reader should not have to know which.
|
||||
lines := link.BuildLogOf("*", id)
|
||||
sub, err := js.Context().PullSubscribe(lines, "",
|
||||
nats.BindStream(broker.EventsStream), nats.DeliverAll(), nats.AckNone())
|
||||
if err != nil {
|
||||
return fmt.Errorf("cannot read %s from the bus: %w", link.BuildLog(id), err)
|
||||
return fmt.Errorf("cannot read %s from the bus: %w", lines, err)
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
func claiming(module, seat string, on ...string) inventory.Entry {
|
||||
return inventory.Entry{
|
||||
Manifest: catalogue.Manifest{Module: module, Claims: []catalogue.Claim{{Name: seat}}},
|
||||
On: on,
|
||||
}
|
||||
}
|
||||
|
||||
// The controller asks the build role that has a holder (novox/hq ADR 0190 handover): the retired
|
||||
// one while only the builder is assigned, the current one from the first build-agent on, and the
|
||||
// current one when nothing holds either — where every ask goes once the handover is done.
|
||||
func TestTheControllerAsksTheBuildRoleThatHasAHolder(t *testing.T) {
|
||||
onlyTheBuilder := []inventory.Entry{
|
||||
claiming("builder", link.TheBuildMachineBefore, "anchor"),
|
||||
claiming("build-agent", link.TheBuildMachine), // registered, assigned nowhere yet
|
||||
}
|
||||
if got := buildSeatAmong(onlyTheBuilder); got != link.TheBuildMachineBefore {
|
||||
t.Errorf("with only the builder assigned, asked %q", got)
|
||||
}
|
||||
bothHeld := []inventory.Entry{
|
||||
claiming("builder", link.TheBuildMachineBefore, "anchor"),
|
||||
claiming("build-agent", link.TheBuildMachine, "home-server"),
|
||||
}
|
||||
if got := buildSeatAmong(bothHeld); got != link.TheBuildMachine {
|
||||
t.Errorf("with a build-agent assigned anywhere, asked %q", got)
|
||||
}
|
||||
neither := []inventory.Entry{claiming("builder", link.TheBuildMachineBefore)}
|
||||
if got := buildSeatAmong(neither); got != link.TheBuildMachine {
|
||||
t.Errorf("with no holder of either, asked %q, want the current role", got)
|
||||
}
|
||||
if got := buildSeatAmong(nil); got != link.TheBuildMachine {
|
||||
t.Errorf("an empty catalogue asks %q", got)
|
||||
}
|
||||
}
|
||||
+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)
|
||||
}
|
||||
}
|
||||
|
||||
+17
-8
@@ -110,10 +110,13 @@ 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"}
|
||||
// Both build roles while the handover runs (novox/hq ADR 0190): the controller asks whichever has a
|
||||
// holder, and the retired one has one until build-agent replaces the builder. The second entry
|
||||
// goes with the retired seat row.
|
||||
var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"}
|
||||
|
||||
// enrolmentPrefix is the space every enrolling node's user and inbox live under, so the one place the
|
||||
// controller may answer an enrolment is derived from the same constant the user is named from.
|
||||
@@ -208,7 +211,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 +373,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,10 +217,15 @@ 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"),
|
||||
// The retired build role's outcome too, while the handover runs (novox/hq ADR 0190): the one
|
||||
// build machine keeps answering on its seat until build-agent replaces it, and the outcome that
|
||||
// registers build-agent itself comes from there. Appended, for the same reason as above; goes
|
||||
// with the retired seat row.
|
||||
seatEventSubject("mesh-build-machine", "built"),
|
||||
}
|
||||
|
||||
// moduleEventSubject is where one module's event lands. The same derivation PermissionsFor uses, so
|
||||
|
||||
+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-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused", "mesh.seat.node-build-agent.accept.>"] }
|
||||
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
{ user: "enrol.one", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
|
||||
@@ -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{
|
||||
|
||||
@@ -0,0 +1,32 @@
|
||||
package link
|
||||
|
||||
import "testing"
|
||||
|
||||
// A build machine serves the seat its credential claims (novox/hq ADR 0190 handover): the old
|
||||
// `builder` keeps the old role, a `build-agent` takes the new, from one binary and no flag.
|
||||
func TestABuildMachineServesTheSeatItsCredentialClaims(t *testing.T) {
|
||||
if got := BuildSeatClaimed([]string{"mesh-build-machine"}); got != "mesh-build-machine" {
|
||||
t.Errorf("a credential claiming the old role serves %q", got)
|
||||
}
|
||||
if got := BuildSeatClaimed([]string{"node-build-agent"}); got != TheBuildMachine {
|
||||
t.Errorf("a credential claiming the new role serves %q", got)
|
||||
}
|
||||
if got := BuildSeatClaimed(nil); got != TheBuildMachine {
|
||||
t.Errorf("a credential claiming nothing serves %q, want the current role", got)
|
||||
}
|
||||
if got := BuildSeatClaimed([]string{"", "node-build-agent"}); got != TheBuildMachine {
|
||||
t.Errorf("an empty claim is skipped; got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
// What a machine says about a build is the event of the seat it took the build from, so an outcome
|
||||
// is heard where the asker of that seat listens.
|
||||
func TestABuildsEventsAreItsSeats(t *testing.T) {
|
||||
if BuildOutcomeOf(TheBuildMachineBefore) != "mesh.seat.mesh-build-machine.event.built" {
|
||||
t.Error(BuildOutcomeOf(TheBuildMachineBefore))
|
||||
}
|
||||
if BuildWorkOf(TheBuildMachine) != BuildWork() || BuildOutcomeOf(TheBuildMachine) != BuildOutcome() ||
|
||||
BuildStartedOf(TheBuildMachine) != BuildStarted() || BuildLogOf(TheBuildMachine, "b1") != BuildLog("b1") {
|
||||
t.Error("the no-argument forms must name the current role")
|
||||
}
|
||||
}
|
||||
+53
-7
@@ -19,13 +19,57 @@ 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** — and why no step strands a build. The old seat's
|
||||
// stream and worker (SEAT_MESH_BUILD_MACHINE, SEAT_MESH_BUILD_MACHINE_worker) stay on the bus until
|
||||
// removed by hand, and the builder keeps draining them while it is assigned, because a machine
|
||||
// serves the seat its credential claims (BuildSeatClaimed) and the controller asks the seat that
|
||||
// has a holder (buildSeatAmong in the command) and hears both seats' outcomes:
|
||||
//
|
||||
// 1. Merge the controller and the host's first user list together; the new controller rolls and,
|
||||
// seeing only the builder assigned, still asks mesh-build-machine — which the builder holds.
|
||||
// 2. Merge the catalogue's build-agent; the builder builds it and the controller registers it.
|
||||
// 3. On each machine that builds: `module issue build-agent --node <n>`, then `assign`, then
|
||||
// `push`. The first holder appears, and from then on asks go to node-build-agent.
|
||||
// 4. Unassign builder everywhere and `module forget` it.
|
||||
// 5. By hand: delete SEAT_MESH_BUILD_MACHINE and its worker, drop the retired seat row and
|
||||
// TheBuildMachineBefore with it, and the second entries in seatsTheControllerAsks and
|
||||
// ControllerFollows.
|
||||
const TheBuildMachine = "node-build-agent"
|
||||
|
||||
// TheBuildMachineBefore is the role a build was submitted to until ADR 0190: the mesh's one build
|
||||
// machine, mesh-scoped. Kept named while the handover runs — a machine whose credential claims it
|
||||
// still serves it, and the controller still hears its outcomes — and dropped with the retired seat
|
||||
// row once nothing claims it.
|
||||
const TheBuildMachineBefore = "mesh-build-machine"
|
||||
|
||||
// BuildSeatClaimed is the build role a machine serves: the first seat its credential claims, or the
|
||||
// current role when the credential names none (a credential from before claims travelled in it, or
|
||||
// one written by hand). **The credential decides, not the binary** (ADR 0190 handover): one build
|
||||
// machine binary runs as the old `builder` on the old seat and as a `build-agent` on the new one,
|
||||
// each taking the work the mesh issued it a credential for, so neither drains the other's queue
|
||||
// and the switch needs no flag day.
|
||||
func BuildSeatClaimed(claimed []string) string {
|
||||
for _, seat := range claimed {
|
||||
if seat != "" {
|
||||
return seat
|
||||
}
|
||||
}
|
||||
return TheBuildMachine
|
||||
}
|
||||
|
||||
// BuildWork is where a build request lands, and BuildOutcome is where its result does. Derived from
|
||||
// the seat, so both sides name the role and neither names the other.
|
||||
func BuildWork() string { return "mesh.seat." + TheBuildMachine + ".accept.build" }
|
||||
func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.built" }
|
||||
// the seat, so both sides name the role and neither names the other. The no-argument forms name the
|
||||
// current role; the `Of` forms take the seat, for the handover during which two roles exist.
|
||||
func BuildWork() string { return BuildWorkOf(TheBuildMachine) }
|
||||
func BuildOutcome() string { return BuildOutcomeOf(TheBuildMachine) }
|
||||
func BuildWorkOf(seat string) string { return "mesh.seat." + seat + ".accept.build" }
|
||||
func BuildOutcomeOf(seat string) string {
|
||||
return "mesh.seat." + seat + ".event.built"
|
||||
}
|
||||
|
||||
// BuildStarted is where a build machine says it has taken a build, and BuildLog is where it says
|
||||
// what it is doing, one line per message, under the build's own id (novox/hq ADR 0157).
|
||||
@@ -35,8 +79,10 @@ func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.bui
|
||||
// lived in one container's stderr on one machine. Every line is now an event of the role, retained
|
||||
// with the rest of the mesh's events, so a reader follows a build live by subscribing its subject,
|
||||
// or reads it back afterwards from the stream, and a viewer is a subscriber and nothing more.
|
||||
func BuildStarted() string { return "mesh.seat." + TheBuildMachine + ".event.started" }
|
||||
func BuildLog(id string) string { return "mesh.seat." + TheBuildMachine + ".event.log." + id }
|
||||
func BuildStarted() string { return BuildStartedOf(TheBuildMachine) }
|
||||
func BuildLog(id string) string { return BuildLogOf(TheBuildMachine, id) }
|
||||
func BuildStartedOf(seat string) string { return "mesh.seat." + seat + ".event.started" }
|
||||
func BuildLogOf(seat, id string) string { return "mesh.seat." + seat + ".event.log." + id }
|
||||
|
||||
// BuildStart is what a build machine says the moment it takes a build.
|
||||
type BuildStart struct {
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -25,16 +25,34 @@ import (
|
||||
type natsBuilds struct {
|
||||
js *broker.JetStream
|
||||
owned bool
|
||||
// seat is the build role asked: the one that has a holder (ADR 0190 handover), chosen by the
|
||||
// controller from what is assigned, so an ask lands where a machine is pulling.
|
||||
seat string
|
||||
}
|
||||
|
||||
// BuildsOverNATS is the asking side on the bus being built. It dials, because the command that asks
|
||||
// for a build is a one-shot and holds nothing else.
|
||||
// BuildsOverNATS is the asking side on the bus being built, asking the current build role. It dials,
|
||||
// because the command that asks for a build is a one-shot and holds nothing else.
|
||||
func BuildsOverNATS(address string) (Builders, error) {
|
||||
return BuildsOverNATSOn(address, TheBuildMachine)
|
||||
}
|
||||
|
||||
// BuildsOverNATSOn is the asking side for one named build role — during the handover from the one
|
||||
// build machine to build agents, the role that has a holder (ADR 0190).
|
||||
func BuildsOverNATSOn(address, seat string) (Builders, error) {
|
||||
js, err := broker.Dial(address)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot reach the bus at %s to ask for a build: %w", address, err)
|
||||
}
|
||||
return &natsBuilds{js: js, owned: true}, nil
|
||||
return &natsBuilds{js: js, owned: true, seat: seat}, nil
|
||||
}
|
||||
|
||||
// role is the seat asked: what the asker was made for, or the current build role for one made
|
||||
// without saying (a test building the struct by hand).
|
||||
func (b *natsBuilds) role() string {
|
||||
if b.seat == "" {
|
||||
return TheBuildMachine
|
||||
}
|
||||
return b.seat
|
||||
}
|
||||
|
||||
func (b *natsBuilds) Close() {
|
||||
@@ -51,7 +69,7 @@ func (b *natsBuilds) Ask(ctx context.Context, request BuildRequest) error {
|
||||
}
|
||||
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||
defer cancel()
|
||||
if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil {
|
||||
if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil {
|
||||
return fmt.Errorf("cannot submit a build: %w", err)
|
||||
}
|
||||
return nil
|
||||
@@ -63,7 +81,7 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest,
|
||||
// Subscribed before the ask, so an outcome cannot arrive before there is anywhere for it to
|
||||
// land. Core, not the stream: the asker is waiting now, and the durable copy of this outcome is
|
||||
// the same event on EVENTS, which the controller records.
|
||||
outcomes, err := b.js.Conn().SubscribeSync(BuildOutcome())
|
||||
outcomes, err := b.js.Conn().SubscribeSync(BuildOutcomeOf(b.role()))
|
||||
if err != nil {
|
||||
return BuildResult{}, fmt.Errorf("cannot listen for a build's outcome: %w", err)
|
||||
}
|
||||
@@ -80,7 +98,7 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest,
|
||||
// be assumed, because nothing else will ever say so.
|
||||
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||
defer cancel()
|
||||
if _, err := b.js.Context().Publish(BuildWork(), body, nats.Context(publish)); err != nil {
|
||||
if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil {
|
||||
return BuildResult{}, fmt.Errorf("cannot submit a build: %w", err)
|
||||
}
|
||||
|
||||
@@ -115,9 +133,16 @@ type natsMachine struct {
|
||||
sub *nats.Subscription
|
||||
}
|
||||
|
||||
// MachineOverNATS takes build work from the role this machine holds.
|
||||
// MachineOverNATS takes build work from the current build role.
|
||||
func MachineOverNATS(js *broker.JetStream, on string) BuildMachine {
|
||||
return &natsMachine{js: js, on: on, seat: TheBuildMachine}
|
||||
return MachineOverNATSOn(js, on, TheBuildMachine)
|
||||
}
|
||||
|
||||
// MachineOverNATSOn takes build work from the role named — the one this machine's credential claims
|
||||
// (ADR 0190 handover): its asks come from that seat's worker, and what it says about a build goes
|
||||
// out as that seat's events, so an outcome is heard where the asker listens.
|
||||
func MachineOverNATSOn(js *broker.JetStream, on, seat string) BuildMachine {
|
||||
return &natsMachine{js: js, on: on, seat: seat}
|
||||
}
|
||||
|
||||
func (m *natsMachine) Close() {
|
||||
@@ -126,29 +151,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 +186,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
|
||||
@@ -178,7 +219,7 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
|
||||
// the ask to a second machine nor counts the wait against its deliveries.
|
||||
working := make(chan struct{})
|
||||
go stillWorking(msg, working)
|
||||
do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js})
|
||||
do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js, seat: m.seat})
|
||||
close(working)
|
||||
}
|
||||
}
|
||||
@@ -189,7 +230,9 @@ type natsBuild struct {
|
||||
msg *nats.Msg
|
||||
on string
|
||||
js *broker.JetStream
|
||||
seq int
|
||||
// seat is the role this build was taken from; what the machine says about it is that role's.
|
||||
seat string
|
||||
seq int
|
||||
}
|
||||
|
||||
func (b *natsBuild) Request() BuildRequest { return b.request }
|
||||
@@ -216,7 +259,7 @@ func (b *natsBuild) Announce(ctx context.Context, result BuildResult) error {
|
||||
}
|
||||
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||
defer cancel()
|
||||
if _, err := b.js.Context().Publish(BuildOutcome(), body, nats.Context(publish)); err != nil {
|
||||
if _, err := b.js.Context().Publish(BuildOutcomeOf(b.seat), body, nats.Context(publish)); err != nil {
|
||||
return fmt.Errorf("cannot announce a build's outcome: %w", err)
|
||||
}
|
||||
return nil
|
||||
@@ -236,7 +279,7 @@ func (b *natsBuild) Began(ctx context.Context) error {
|
||||
}
|
||||
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||
defer cancel()
|
||||
if _, err := b.js.Context().Publish(BuildStarted(), body, nats.Context(publish)); err != nil {
|
||||
if _, err := b.js.Context().Publish(BuildStartedOf(b.seat), body, nats.Context(publish)); err != nil {
|
||||
return fmt.Errorf("cannot say a build started: %w", err)
|
||||
}
|
||||
return nil
|
||||
@@ -254,7 +297,7 @@ func (b *natsBuild) Say(step, message string) {
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_ = b.js.Conn().Publish(BuildLog(b.request.ID), body)
|
||||
_ = b.js.Conn().Publish(BuildLogOf(b.seat, b.request.ID), body)
|
||||
}
|
||||
|
||||
func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) }
|
||||
|
||||
@@ -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,119 @@ 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):
|
||||
}
|
||||
}
|
||||
|
||||
// During the handover (ADR 0190) two build roles exist. A machine whose credential claims the retired
|
||||
// one takes an ask published to that seat and answers as that seat; the asker of that seat hears it.
|
||||
func TestNatsAMachineOnTheRetiredBuildRoleTakesThatRolesAsks(t *testing.T) {
|
||||
js := aBusWithTheBuildRole(t)
|
||||
seats := []broker.DeclaredSeat{{Name: TheBuildMachineBefore, Accepts: []string{"build"}, Emits: []string{"built"}}}
|
||||
if err := broker.RaiseSeats(js, seats, map[string]broker.Holder{TheBuildMachineBefore: {Node: "anchor", Module: "builder"}}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") })
|
||||
ctx, stop := context.WithCancel(context.Background())
|
||||
defer stop()
|
||||
|
||||
machine := MachineOverNATSOn(js, "anchor", TheBuildMachineBefore)
|
||||
defer machine.Close()
|
||||
go func() {
|
||||
_ = machine.Take(ctx, func(ctx context.Context, work Build) {
|
||||
_ = work.Began(ctx)
|
||||
_ = work.Announce(ctx, BuildResult{ID: work.Request().ID, Repository: work.Request().Repository, On: "anchor", Commit: "abc"})
|
||||
_ = work.Done()
|
||||
})
|
||||
}()
|
||||
|
||||
asker, err := BuildsOverNATSOn(os.Getenv("MESH_TEST_NATS"), TheBuildMachineBefore)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer asker.Close()
|
||||
result, err := asker.Submit(ctx, BuildRequest{ID: "build-old-seat", Repository: "r"}, 20*time.Second)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if result.On != "anchor" || result.ID != "build-old-seat" {
|
||||
t.Errorf("the retired role's holder did not answer: %+v", result)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -223,10 +223,11 @@ func kindOfSubject(subject string) (string, bool) {
|
||||
return KindCatchUp, true
|
||||
case broker.ControllerFollows[3]:
|
||||
return KindSourceMoved, true
|
||||
case BuildOutcome():
|
||||
case BuildOutcome(), BuildOutcomeOf(TheBuildMachineBefore):
|
||||
// A build's outcome is the role's event now, so it arrives on the events stream rather than
|
||||
// the control branch — and is acted on by the same handler, because what the controller does
|
||||
// with it did not change (novox/hq ADR 0121).
|
||||
// with it did not change (novox/hq ADR 0121). From either build role while the handover
|
||||
// runs (ADR 0190): the old builder still answers on the retired seat until it is unassigned.
|
||||
return KindBuilt, true
|
||||
}
|
||||
return "", false
|
||||
|
||||
Reference in New Issue
Block a user