From e610f2d92c2129e1573ac13b325c0bac0f0466e6 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 19:17:42 +0200 Subject: [PATCH] Let a build agent be paused, have a build killed, and end an ask cancelled as it took it (hq ADR 0219) A queued ask could only be waited out and a running build only ended by stopping the machine, which redelivered it elsewhere. The holder now serves current, kill, pause and resume on its own machine's subjects; a kill ends the build's process group and labelled containers and settles the ask as failed, killed by hand; pause is kept in the workspace across a restart and said on the bus. The controller writes cancelled ids to a cancelled set the holder reads on taking an ask, closing the race a delete alone leaves. --- cmd/mesh-builder/holder.go | 295 +++++++++++++++++ cmd/mesh-builder/holder_test.go | 94 ++++++ cmd/mesh-builder/main.go | 87 ++++- cmd/mesh-controller/push.go | 5 + internal/broker/cancelled.go | 92 ++++++ internal/broker/cancelled_test.go | 71 ++++ internal/broker/nats.go | 13 + internal/broker/state.go | 3 +- internal/broker/testdata/composed.conf | 2 +- internal/builder/builder.go | 3 + internal/builder/kill.go | 78 +++++ internal/builder/kill_test.go | 103 ++++++ internal/catalogue/seats.go | 30 +- internal/link/builds_nats.go | 57 +++- internal/link/queue.go | 439 +++++++++++++++++++++++++ internal/link/queue_test.go | 231 +++++++++++++ internal/link/seattools.go | 14 +- 17 files changed, 1601 insertions(+), 16 deletions(-) create mode 100644 cmd/mesh-builder/holder.go create mode 100644 cmd/mesh-builder/holder_test.go create mode 100644 internal/broker/cancelled.go create mode 100644 internal/broker/cancelled_test.go create mode 100644 internal/builder/kill.go create mode 100644 internal/builder/kill_test.go create mode 100644 internal/link/queue.go create mode 100644 internal/link/queue_test.go diff --git a/cmd/mesh-builder/holder.go b/cmd/mesh-builder/holder.go new file mode 100644 index 0000000..95d69f4 --- /dev/null +++ b/cmd/mesh-builder/holder.go @@ -0,0 +1,295 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "strings" + "sync" + "time" + + "github.com/novox/mesh-controller/internal/builder" + "github.com/novox/mesh-controller/internal/link" +) + +// What a holder of the build seat answers for, on its own machine (novox/hq ADR 0219). +// +// **The queue is the controller's; the build running here is this machine's.** The controller can +// see and change what waits in the seat's queue, but an ask a holder already took is a process tree +// on this machine and containers in this machine's runtime, and only this machine can end them. So +// the holder serves four verbs on the seat's subjects for this machine: what it is building, kill +// it, pause, resume. +// +// **One build at a time** (ADR 0190), which is also what builder.Said assumes — a package global +// set per build — so "the build running here" is one or none, and kill names it by id so a call +// that arrives as one build ends and the next begins cannot end the wrong one. + +// holder is this machine's state as a holder of the build seat. +type holder struct { + on, seat string + // workspace is where the paused flag is kept, so a holder restarted while paused stays paused + // rather than silently taking work again. + workspace string + // say publishes this machine's state: whether it takes work (link.HolderState). + say func(link.HolderState) error + // remove runs what a kill needs outside the build: the containers left behind. + remove builder.Runner + + mu sync.Mutex + paused bool + running *running +} + +// running is the build this machine is doing. +type running struct { + request link.BuildRequest + step string + started time.Time + cancel context.CancelFunc + killed bool + done chan struct{} +} + +// pausedFile is where the flag lives in the workspace. +func pausedFile(workspace string) string { return filepath.Join(workspace, ".mesh-builder-paused") } + +// newHolder reads the paused flag the workspace keeps. +func newHolder(on, seat, workspace string, say func(link.HolderState) error) *holder { + h := &holder{on: on, seat: seat, workspace: workspace, say: say, remove: plainRun} + if _, err := os.Stat(pausedFile(workspace)); err == nil { + h.paused = true + } + return h +} + +// Paused is asked by the taking loop before every fetch. +func (h *holder) Paused() bool { + h.mu.Lock() + defer h.mu.Unlock() + return h.paused +} + +// setPaused records the flag in the workspace first and then in memory, so what this holder says +// it is and what it would be after a restart never differ. +func (h *holder) setPaused(paused bool) error { + path := pausedFile(h.workspace) + if paused { + if err := os.MkdirAll(h.workspace, 0o755); err != nil { + return err + } + if err := os.WriteFile(path, []byte(time.Now().UTC().Format(time.RFC3339)+"\n"), 0o644); err != nil { + return fmt.Errorf("cannot keep the paused flag in %s: %w", path, err) + } + } else if err := os.Remove(path); err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("cannot remove the paused flag %s: %w", path, err) + } + h.mu.Lock() + h.paused = paused + h.mu.Unlock() + h.announce() + return nil +} + +// announce says whether this machine takes work. Never fatal: the flag is kept either way, and the +// controller reading an older state is a plan read as late rather than a build lost. +func (h *holder) announce() { + if h.say == nil { + return + } + if err := h.say(link.HolderState{On: h.on, Paused: h.Paused(), At: time.Now().UTC().Format(time.RFC3339Nano)}); err != nil { + fmt.Fprintf(os.Stderr, "cannot say whether this machine takes builds: %v\n", err) + } +} + +// stateWhilePaused says this machine's state at start, and again every while it stays paused, so +// a pause outlives the events stream's retention. +func (h *holder) stateWhilePaused(ctx context.Context, every time.Duration) { + h.announce() + tick := time.NewTicker(every) + defer tick.Stop() + for { + select { + case <-ctx.Done(): + return + case <-tick.C: + if h.Paused() { + h.announce() + } + } + } +} + +// begin records the build this machine took, with the cancel that ends it. +func (h *holder) begin(request link.BuildRequest, cancel context.CancelFunc) *running { + r := &running{request: request, started: time.Now(), cancel: cancel, done: make(chan struct{})} + h.mu.Lock() + h.running = r + h.mu.Unlock() + return r +} + +// end clears it, and lets a kill waiting on it know it is over. +func (h *holder) end(r *running) { + h.mu.Lock() + if h.running == r { + h.running = nil + } + h.mu.Unlock() + close(r.done) +} + +// stepped records the step a build is at, for `current`. +func (h *holder) stepped(r *running, step string) { + h.mu.Lock() + r.step = step + h.mu.Unlock() +} + +// wasKilled is whether a person killed this build. +func (h *holder) wasKilled(r *running) bool { + h.mu.Lock() + defer h.mu.Unlock() + return r.killed +} + +// currentBuild is what `current` answers. +type currentBuild struct { + On string `json:"on"` + Paused bool `json:"paused"` + Running *struct { + ID string `json:"id"` + Repository string `json:"repository"` + Path string `json:"path,omitempty"` + Ref string `json:"ref,omitempty"` + Step string `json:"step,omitempty"` + Started string `json:"started"` + Elapsed string `json:"elapsed"` + } `json:"running,omitempty"` + Said string `json:"said"` +} + +func (h *holder) current() currentBuild { + h.mu.Lock() + defer h.mu.Unlock() + out := currentBuild{On: h.on, Paused: h.paused} + taking := "taking builds" + if h.paused { + taking = "paused, taking no new build" + } + if h.running == nil { + out.Said = fmt.Sprintf("%s is building nothing; %s", h.on, taking) + return out + } + r := h.running + elapsed := time.Since(r.started).Round(time.Second) + out.Running = &struct { + ID string `json:"id"` + Repository string `json:"repository"` + Path string `json:"path,omitempty"` + Ref string `json:"ref,omitempty"` + Step string `json:"step,omitempty"` + Started string `json:"started"` + Elapsed string `json:"elapsed"` + }{r.request.ID, r.request.Repository, r.request.Path, r.request.Ref, r.step, + r.started.UTC().Format(time.RFC3339), elapsed.String()} + out.Said = fmt.Sprintf("%s is building %s (%s) at %s for %s; %s", + h.on, r.request.Repository, r.request.ID, orNothing(r.step), elapsed, taking) + return out +} + +// kill ends the build with this id, if it is the one running here: its context cancelled — which +// kills each command's process group — and the containers it started removed by their label. The +// build's own goroutine then announces it failed, killed by hand, and settles the ask, so it is not +// redelivered; this waits a while for that, to say it happened. +func (h *holder) kill(id string) (string, error) { + h.mu.Lock() + r := h.running + if r == nil || r.request.ID != id { + doing := "nothing" + if r != nil { + doing = r.request.ID + } + h.mu.Unlock() + return "", fmt.Errorf("%s is not building %s; it is building %s. `queue` says where an ask is", h.on, id, doing) + } + r.killed = true + cancel, done := r.cancel, r.done + h.mu.Unlock() + + cancel() + cleanup, stop := context.WithTimeout(context.Background(), 30*time.Second) + defer stop() + removed, err := builder.RemoveContainersOf(cleanup, h.remove, id) + containers := fmt.Sprintf("%d container(s) it started removed", removed) + if err != nil { + containers = "its containers could not be listed or removed: " + err.Error() + } + select { + case <-done: + return fmt.Sprintf("killed %s (%s) on %s: its commands ended, %s, and its outcome announced as "+ + "failed, %s — settled, so it is not handed to another machine", id, r.request.Repository, h.on, + containers, link.KilledByHand), nil + case <-time.After(30 * time.Second): + return fmt.Sprintf("killed %s on %s: %s; the build has not finished ending yet — `builds` says when "+ + "its outcome is in", id, h.on, containers), nil + } +} + +// handlers are the seat's verbs, as this machine answers them. +func (h *holder) handlers() map[string]link.ToolHandler { + return map[string]link.ToolHandler{ + "current": func(context.Context, json.RawMessage) (any, error) { return h.current(), nil }, + "kill": func(_ context.Context, raw json.RawMessage) (any, error) { + var args struct { + ID string `json:"id"` + } + if err := json.Unmarshal(raw, &args); err != nil || strings.TrimSpace(args.ID) == "" { + return nil, errors.New("kill needs the build's id") + } + said, err := h.kill(strings.TrimSpace(args.ID)) + if err != nil { + return nil, err + } + return map[string]any{"said": said}, nil + }, + "pause": func(context.Context, json.RawMessage) (any, error) { + if err := h.setPaused(true); err != nil { + return nil, err + } + said := h.on + " is paused: it takes no new build until resumed" + if c := h.current(); c.Running != nil { + said += "; " + c.Running.ID + " runs on and finishes" + } + return map[string]any{"said": said, "paused": true}, nil + }, + "resume": func(context.Context, json.RawMessage) (any, error) { + if err := h.setPaused(false); err != nil { + return nil, err + } + return map[string]any{"said": h.on + " takes builds again", "paused": false}, nil + }, + } +} + +func orNothing(s string) string { + if s == "" { + return "its start" + } + return s +} + +// plainRun runs a command outside any build: no line reaches a build's log, because the build whose +// log it would be is the one being ended. +func plainRun(ctx context.Context, dir, name string, args ...string) (string, error) { + cmd := exec.CommandContext(ctx, name, args...) + cmd.Dir = dir + out, err := cmd.CombinedOutput() + if err != nil { + return string(out), fmt.Errorf("%s %s: %w: %s", name, strings.Join(args, " "), err, strings.TrimSpace(string(out))) + } + return string(out), nil +} diff --git a/cmd/mesh-builder/holder_test.go b/cmd/mesh-builder/holder_test.go new file mode 100644 index 0000000..f5b5238 --- /dev/null +++ b/cmd/mesh-builder/holder_test.go @@ -0,0 +1,94 @@ +package main + +import ( + "context" + "encoding/json" + "strings" + "testing" + + "github.com/novox/mesh-controller/internal/link" +) + +// A holder paused stays paused across its restart: the flag is kept in its workspace (novox/hq ADR +// 0219), and what it says about itself follows. +func TestAPausedHolderStaysPausedAcrossARestart(t *testing.T) { + workspace := t.TempDir() + var said []link.HolderState + h := newHolder("ace", link.TheBuildMachine, workspace, func(s link.HolderState) error { + said = append(said, s) + return nil + }) + if h.Paused() { + t.Fatal("a new holder starts paused") + } + if _, err := h.handlers()["pause"](context.Background(), json.RawMessage(`{}`)); err != nil { + t.Fatal(err) + } + if !h.Paused() || len(said) != 1 || !said[0].Paused || said[0].On != "ace" { + t.Fatalf("paused: %v, said %+v", h.Paused(), said) + } + again := newHolder("ace", link.TheBuildMachine, workspace, nil) + if !again.Paused() { + t.Fatal("restarted, the holder forgot it was paused") + } + if _, err := again.handlers()["resume"](context.Background(), json.RawMessage(`{}`)); err != nil { + t.Fatal(err) + } + if newHolder("ace", link.TheBuildMachine, workspace, nil).Paused() { + t.Fatal("resumed, the holder came back paused") + } +} + +// kill ends the build with that id — its context, which ends its commands — removes what it left by +// label, and refuses an id it is not building. +func TestKillEndsTheBuildRunningHereAndNoOther(t *testing.T) { + h := newHolder("ace", link.TheBuildMachine, t.TempDir(), nil) + var removed []string + h.remove = func(_ context.Context, _ string, name string, args ...string) (string, error) { + removed = append(removed, name+" "+strings.Join(args, " ")) + if args[0] == "ps" { + return "c1\n", nil + } + return "", nil + } + kill := h.handlers()["kill"] + if _, err := kill(context.Background(), json.RawMessage(`{"id":"build-1"}`)); err == nil { + t.Fatal("killed a build while none ran") + } + + building, cancel := context.WithCancel(context.Background()) + r := h.begin(link.BuildRequest{ID: "build-1", Repository: "novox/a"}, cancel) + h.stepped(r, "image") + if c := h.current(); c.Running == nil || c.Running.ID != "build-1" || c.Running.Step != "image" { + t.Fatalf("current says %+v", c) + } + if _, err := kill(context.Background(), json.RawMessage(`{"id":"build-2"}`)); err == nil || + !strings.Contains(err.Error(), "build-1") { + t.Fatalf("killed another id: %v", err) + } + // The build's own goroutine: it ends when its context does, as a build's commands do. + go func() { + <-building.Done() + h.end(r) + }() + answer, err := kill(context.Background(), json.RawMessage(`{"id":"build-1"}`)) + if err != nil { + t.Fatal(err) + } + if building.Err() == nil { + t.Fatal("the build's context was not cancelled") + } + if !h.wasKilled(r) { + t.Error("the build is not marked killed, so it would be announced as an ordinary failure") + } + said := answer.(map[string]any)["said"].(string) + if !strings.Contains(said, "1 container(s)") || !strings.Contains(said, link.KilledByHand) { + t.Errorf("kill said %q", said) + } + if len(removed) != 2 || !strings.Contains(removed[0], "label=mesh.build=build-1") { + t.Errorf("removed %v", removed) + } + if c := h.current(); c.Running != nil { + t.Errorf("after the kill current says %+v", c) + } +} diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index 5928a12..a28b791 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -19,11 +19,13 @@ import ( "context" "encoding/json" "fmt" + "log" "net/url" "os" "os/signal" "strings" "syscall" + "time" "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/builder" @@ -107,48 +109,89 @@ func run() error { ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer stop() - machine, err := takeWorkFrom(credential, on) + js, seat, err := dialFor(credential) if err != nil { return err } + defer js.Close() + + // **This machine's holder: paused or not, and the build it is running** (novox/hq ADR 0219). The + // paused flag is read from the workspace before anything is taken, so a holder restarted while + // paused takes nothing. + h := newHolder(on, seat, workspace, func(state link.HolderState) error { + body, err := json.Marshal(state) + if err != nil { + return err + } + _, err = js.Context().Publish(link.BuildPausedOf(seat, on), body) + return err + }) + if h.Paused() { + fmt.Fprintf(os.Stderr, "paused (kept in %s): taking no build until resumed\n", pausedFile(workspace)) + } + machine := link.MachineOverNATSWith(js, on, seat, link.MachineOptions{Paused: h.Paused}) defer machine.Close() + // The seat's verbs, on this machine's subjects. Only the seat that declares them: the retired + // one serves none, and a subscription its holder has no grant for would be refused for ever. + if seat == link.TheBuildMachine { + stopServing, err := link.OverNATS{Conn: js.Conn()}.ServeNodeSeatTools(seat, on, h.handlers(), + log.New(os.Stderr, "", 0)) + if err != nil { + return err + } + defer stopServing() + go h.stateWhilePaused(ctx, time.Hour) + } + fmt.Fprintf(os.Stderr, "building for the mesh, publishing to %s\n", registry) publisher := builder.Registry{Address: registry, Run: builder.Command} return machine.Take(ctx, func(ctx context.Context, work link.Build) { - answer(ctx, publisher, on, workspace, work) + answer(ctx, publisher, on, workspace, work, h) }) } -// takeWorkFrom opens this machine's link to whichever bus the mesh is on. +// dialFor opens this machine's link to whichever bus the mesh is on, and says which build seat it +// holds. // // **One place chooses**, as everywhere else the bus change went (novox/hq ADR 0116 step 5): a build // machine told about both would take work from one and answer on the other, and every log line would // say it was fine. -func takeWorkFrom(credential Credential, on string) (link.BuildMachine, error) { +func dialFor(credential Credential) (*broker.JetStream, string, error) { // **The credential names the bus, and there is one** (novox/hq ADR 0131, design 28 task 5.5). // A credential for the mesh's bus carries user, password and fingerprint beside the address, // and that is enough to dial it, pinned. if !credential.onTheNewBus() { - return nil, fmt.Errorf("the credential at hand names %q, which is not the mesh's bus", credential.URL) + return nil, "", fmt.Errorf("the credential at hand names %q, which is not the mesh's bus", credential.URL) } js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint) if err != nil { - return nil, err + return nil, "", err } // **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 + return js, seat, nil } // answer does one build and says what happened, whichever way it went. -func answer(ctx context.Context, publisher builder.Publisher, on, workspace string, work link.Build) { +func answer(ctx context.Context, publisher builder.Publisher, on, workspace string, work link.Build, h *holder) { request := work.Request() + // **Its own context, so it can be killed alone** (novox/hq ADR 0219): cancelled by `kill`, it ends + // this build's commands and nothing else; the machine's own context ending — a SIGTERM — still + // reaches it through the parent, and that keeps today's meaning below. + building, cancel := context.WithCancel(ctx) + defer cancel() + var mine *running + if h != nil { + mine = h.begin(request, cancel) + defer h.end(mine) + } + // **First thing, and to stdout.** A build request that arrives and produces no visible line until // it either finishes or fails is indistinguishable from one that never arrived — which cost a long // diagnosis against a running mesh, chasing "the handler never fired" when the truth was only that @@ -162,6 +205,9 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri say := func(step, message string) { fmt.Fprintf(os.Stderr, " [%s] %s\n", step, message) work.Say(step, message) + if mine != nil && step != "run" && step != "output" { + h.stepped(mine, step) + } } builder.Said = say defer func() { builder.Said = nil }() @@ -188,11 +234,20 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri // The package-registry credential is a build input, so it is resolved before the clone: a // build that could not have resolved its dependencies is refused in front of the reason, not // after a clone that then fails at npm ci. - built, err = builder.Build(ctx, builder.Command, publisher, + // Every container it starts is labelled with its id, so a kill finds what outlived the + // docker client (ADR 0219). + built, err = builder.Build(building, builder.Labelled(builder.Command, request.ID), publisher, request.Repository, request.Path, request.Ref, workspace, request.Held, npmrc, forgeFrom(), say, request.Seats) } - if err != nil { + // Only a build the kill ended: one that finished in the moment the kill arrived built, and says so. + killed := err != nil && mine != nil && h.wasKilled(mine) + if killed { + // **Killed by hand is the outcome, whatever the build was doing** (novox/hq ADR 0219): the + // error it ended with is the kill's consequence, not a fault of the source. + result.Failed = link.KilledByHand + say("failed", link.KilledByHand) + } else if err != nil { // A failure is a result. A build that fails and says nothing is indistinguishable from a // builder that is not running, and those want completely different responses. result.Failed = err.Error() @@ -217,7 +272,17 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri } } - if err := work.Announce(ctx, result); err != nil { + // A killed build is announced on a context of its own: the build's was the one cancelled, and the + // outcome must go out and the ask be settled — acknowledged, never redelivered to another machine + // to be built again. A machine being stopped is the other case and keeps its meaning: the + // parent's context is gone, nothing is announced or settled, and the ask is redelivered. + announcing := ctx + if killed { + fresh, stop := context.WithTimeout(context.Background(), 30*time.Second) + defer stop() + announcing = fresh + } + if err := work.Announce(announcing, result); err != nil { // Said, not fatal: the build happened. A build reported as failed because announcing it // failed is a lie about work that was done — and the request stays unsettled below only if // nothing was said at all, so another machine can try. diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 9a2d3c5..3e43f0a 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -1142,6 +1142,11 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string) if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil { return err } + // And each work queue's cancelled set (novox/hq ADR 0219), so a holder taking an ask can ask + // whether it was cancelled the moment it took it. + if err := broker.RaiseCancelledSets(js, inventory.MeshSeats()); err != nil { + return err + } // Every module's state (novox/hq ADR 0201), from the catalogue: a bucket exists from // registration, so a module reading one may watch it before its owner runs anywhere. One that // nothing declares any more is said and kept — what it holds is data. diff --git a/internal/broker/cancelled.go b/internal/broker/cancelled.go new file mode 100644 index 0000000..afd0008 --- /dev/null +++ b/internal/broker/cancelled.go @@ -0,0 +1,92 @@ +package broker + +import ( + "context" + "fmt" + "strings" + "time" + + "github.com/nats-io/nats.go/jetstream" +) + +// A seat's cancelled set (novox/hq ADR 0219). +// +// **Cancelling an ask is deleting its message from the seat's work queue** — and that alone has a +// race no ordering on one side closes: a holder may fetch the ask in the moment between the +// controller reading the queue and deleting the message, and build what a person cancelled. So the +// controller first writes the ask's id into the seat's cancelled set, and a holder looks its ask up +// there after taking it and before building: listed, it terminates the ask and announces it failed +// rather than starting it. +// +// **A key-value bucket, because that is the shape the bus already has for "a small set the +// controller writes and many read cheaply"** (ADR 0201's state buckets): one direct read by key per +// ask taken — no consumer, no subscription a holder must keep, nothing that grows a holder's grants +// beyond one read subject. Kept beside the seat's own stream and worker and named like them, so the +// three objects of one work queue read as one family; asserted by the controller with the queue, +// and aged out with it — an id cancelled a week ago names an ask the queue no longer holds. + +// CancelledSetName is the bucket holding a seat's cancelled asks. +func CancelledSetName(seat string) string { return "SEAT_" + upperSnake(seat) + "_cancelled" } + +// hasCancelledSet is whether a seat's queue can be cancelled from: the work queues the controller +// asks, whose asks it alone shows and changes. A module's own seat's queue is that module's affair. +func hasCancelledSet(seat string) bool { + for _, s := range seatsTheControllerAsks { + if s == seat { + return true + } + } + return false +} + +// IsCancelledSet says a bucket is a seat's cancelled set rather than a module's state — the mesh's +// own, and never one to report as state nothing declares. +func IsCancelledSet(bucket string) bool { + return strings.HasPrefix(bucket, "SEAT_") && strings.HasSuffix(bucket, "_cancelled") +} + +// cancelledSetAge is how long a cancelled id is kept: as long as the seat's queue keeps an ask. +const cancelledSetAge = 7 * 24 * time.Hour + +// A CancelledSetAsserter is what raising the cancelled sets needs of a connection. +type CancelledSetAsserter interface { + EnsureCancelledSet(seat string) error +} + +// RaiseCancelledSets asserts the cancelled set of every work queue the controller asks. +func RaiseCancelledSets(a CancelledSetAsserter, seats []DeclaredSeat) error { + for _, s := range seats { + if len(s.Accepts) == 0 || !hasCancelledSet(s.Name) { + continue + } + if err := a.EnsureCancelledSet(s.Name); err != nil { + return fmt.Errorf("asserting the cancelled set of %s: %w", s.Name, err) + } + } + return nil +} + +// EnsureCancelledSet creates a seat's cancelled set if absent and brings its options to match. +// Direct reads on, which is how a holder looks an id up with one request. +func (j *JetStream) EnsureCancelledSet(seat string) error { + js, err := jetstream.New(j.conn) + if err != nil { + return err + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ + Bucket: CancelledSetName(seat), + Description: fmt.Sprintf("the asks of the %s seat cancelled by hand (novox/hq ADR 0219): written "+ + "by the controller before it deletes an ask from the queue, read by a holder on taking one, "+ + "so an ask fetched in that moment is ended rather than built", seat), + History: 1, + TTL: cancelledSetAge, + MaxValueSize: 4 * 1024, + MaxBytes: 4 * 1024 * 1024, + Storage: jetstream.FileStorage, + }); err != nil { + return fmt.Errorf("asserting bucket %s: %w", CancelledSetName(seat), err) + } + return nil +} diff --git a/internal/broker/cancelled_test.go b/internal/broker/cancelled_test.go new file mode 100644 index 0000000..f9e63f7 --- /dev/null +++ b/internal/broker/cancelled_test.go @@ -0,0 +1,71 @@ +package broker + +import "testing" + +// A build agent reads its seat's cancelled set by key and nothing more, answers its verbs on its own +// machine's subjects, and says whether it is paused under its own machine's name; the controller may +// ask any machine's holder its verbs (novox/hq ADR 0219). +func TestTheBuildQueueIsControlledWithTheGrantsItNeedsAndNoMore(t *testing.T) { + seat := Seat{Name: "node-build-agent", Scope: "node", Accepts: []string{"build"}, + Emits: []string{"started", "built", "log.*", "paused.*"}, + Serves: []string{"current", "kill", "pause", "resume"}} + holder, err := PermissionsFor(Principal{Kind: KindModule, Node: "ace", Module: "build-agent", + Holds: []Seat{seat}, PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + has(t, holder.Publish, "$JS.API.DIRECT.GET.KV_SEAT_NODE_BUILD_AGENT_cancelled.$KV.SEAT_NODE_BUILD_AGENT_cancelled.>") + hasNot(t, holder.Publish, "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>") + has(t, holder.Publish, "mesh.seat.node-build-agent.event.paused.*") + has(t, holder.Subscribe, "mesh.seat.node-build-agent.tool.kill.ace") + hasNot(t, holder.Subscribe, "mesh.seat.node-build-agent.tool.kill.g14") + + // A module's own seat's queue is its own affair: no cancelled set, no grant for one. + other, _ := PermissionsFor(Principal{Kind: KindModule, Node: "ace", Module: "telegram", + Holds: []Seat{{Name: "telegram-sender", Accepts: []string{"send"}}}, PasswordHash: "x"}) + hasNot(t, other.Publish, "$JS.API.DIRECT.GET.KV_SEAT_TELEGRAM_SENDER_cancelled.$KV.SEAT_TELEGRAM_SENDER_cancelled.>") + + controller, _ := PermissionsFor(Principal{Kind: KindController, PasswordHash: "x"}) + has(t, controller.Publish, "mesh.seat.node-build-agent.tool.>") + + if CancelledSetName("node-build-agent") != "SEAT_NODE_BUILD_AGENT_cancelled" || + !IsCancelledSet("SEAT_NODE_BUILD_AGENT_cancelled") || IsCancelledSet("build-agent_cancelled") { + t.Error("the cancelled set is not named as its seat's family") + } +} + +// Only the work queues the controller asks get a cancelled set. +func TestOnlyTheControllersQueuesHaveACancelledSet(t *testing.T) { + var asserted []string + a := asserterFunc(func(seat string) error { asserted = append(asserted, seat); return nil }) + if err := RaiseCancelledSets(a, []DeclaredSeat{ + {Name: "node-build-agent", Accepts: []string{"build"}}, + {Name: "telegram-sender", Accepts: []string{"send"}}, + {Name: "mesh-controller"}, + }); err != nil { + t.Fatal(err) + } + if len(asserted) != 1 || asserted[0] != "node-build-agent" { + t.Fatalf("asserted %v", asserted) + } +} + +type asserterFunc func(string) error + +func (f asserterFunc) EnsureCancelledSet(seat string) error { return f(seat) } + +// A seat's cancelled set is never reported as state nothing declares. +func TestACancelledSetIsNotUndeclaredState(t *testing.T) { + undeclared, err := RaiseBuckets(fakeBuckets{names: []string{"SEAT_NODE_BUILD_AGENT_cancelled", "gone_state"}}, nil) + if err != nil { + t.Fatal(err) + } + if len(undeclared) != 1 || undeclared[0] != "gone_state" { + t.Fatalf("undeclared %v", undeclared) + } +} + +type fakeBuckets struct{ names []string } + +func (f fakeBuckets) EnsureBucket(Bucket) error { return nil } +func (f fakeBuckets) BucketNames() ([]string, error) { return f.names, nil } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 556f8e3..9264429 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -221,6 +221,10 @@ func PermissionsFor(p Principal) (Permissions, error) { // the role, and whichever machine holding it is idle takes it. for _, seat := range seatsTheControllerAsks { pub = append(pub, "mesh.seat."+seat+".accept.>") + // **And its holders' verbs, on every machine** (novox/hq ADR 0219): what a holder is + // building, kill it, pause it, resume it. The queue is the controller's to show and to + // change, and what one machine is doing with an ask it took only that machine can say. + pub = append(pub, "mesh.seat."+seat+".tool.>") } // **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own // facts under the seat it holds, because a role's events belong to the role and keep their @@ -404,6 +408,15 @@ func PermissionsFor(p Principal) (Permissions, error) { "$JS.API.CONSUMER.INFO."+stream+"."+worker, "$JS.API.CONSUMER.MSG.NEXT."+stream+"."+worker, "$JS.ACK."+stream+"."+worker+".>") + // **And whether an ask it took was cancelled** (novox/hq ADR 0219): one key of the + // seat's cancelled set, read directly by its id, so a holder that fetched an ask in the + // moment the controller cancelled it ends it instead of building it. Read, never written: + // the set is the controller's, and only the work queues the controller asks — and so may + // cancel from — have one. + if hasCancelledSet(s.Name) { + set := CancelledSetName(s.Name) + pub = append(pub, "$JS.API.DIRECT.GET.KV_"+set+".$KV."+set+".>") + } for _, a := range s.Accepts { sub = append(sub, seatSubject(s, "accept", a)) } diff --git a/internal/broker/state.go b/internal/broker/state.go index 0324543..9c1e594 100644 --- a/internal/broker/state.go +++ b/internal/broker/state.go @@ -160,7 +160,8 @@ func RaiseBuckets(a BucketAsserter, buckets []Bucket) (undeclared []string, err return nil, fmt.Errorf("listing the bus's state: %w", err) } for _, n := range names { - if !declared[n] { + // A seat's cancelled set is the mesh's own (novox/hq ADR 0219), not a module's state. + if !declared[n] && !IsCancelledSet(n) { undeclared = append(undeclared, n) } } diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index 4c7b32e..af658e6 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ 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", "mesh.seat.node-build-agent.accept.>"] } + 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-build-machine.tool.>", "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.>", "mesh.seat.node-build-agent.tool.>"] } subscribe: { allow: ["$JS.API.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_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" } } } diff --git a/internal/builder/builder.go b/internal/builder/builder.go index 5f26afa..1441265 100644 --- a/internal/builder/builder.go +++ b/internal/builder/builder.go @@ -761,6 +761,9 @@ func Command(ctx context.Context, dir, name string, args ...string) (string, err tell("run", "$ (%s) %s %s", short(filepath.Base(dir)), name, strings.Join(args, " ")) cmd := exec.CommandContext(ctx, name, args...) cmd.Dir = dir + // **Ended whole when the build is** (novox/hq ADR 0219): its own process group, killed as one + // when the context ends, so a build killed by hand leaves no `git` or `docker` child running. + inItsOwnGroup(cmd) out, err := cmd.CombinedOutput() if err != nil { tell("run", "! %s %s failed after %s", name, args[0], since(started)) diff --git a/internal/builder/kill.go b/internal/builder/kill.go new file mode 100644 index 0000000..a46dec2 --- /dev/null +++ b/internal/builder/kill.go @@ -0,0 +1,78 @@ +package builder + +import ( + "context" + "os/exec" + "strings" + "syscall" + "time" +) + +// A build killed by hand (novox/hq ADR 0219). +// +// A build is a tree of commands — git, then docker, and docker's own children — and a container the +// docker client started keeps running when the client dies. So ending one by hand is three things: +// the build's context is cancelled, which kills each command's whole process group rather than the +// one process the context knows; every container the build started carries the build's id as a +// label, so what outlived its client is found and removed by that label; and the holder announces +// the outcome itself, since nothing else will. + +// BuildLabel is the label every container a build starts carries, valued with the build's id. +const BuildLabel = "mesh.build" + +// KillWait is how long a command killed with its build may take to let go of its output before it +// is abandoned: a grandchild holding the pipe open must not hold the build open with it. +const KillWait = 10 * time.Second + +// inItsOwnGroup makes a command the leader of its own process group, killed as a group when its +// context ends. +func inItsOwnGroup(cmd *exec.Cmd) { + cmd.SysProcAttr = &syscall.SysProcAttr{Setpgid: true} + cmd.Cancel = func() error { + if cmd.Process == nil { + return nil + } + // The negative pid is the group: the command and everything it started. + if err := syscall.Kill(-cmd.Process.Pid, syscall.SIGKILL); err != nil { + return cmd.Process.Kill() + } + return nil + } + cmd.WaitDelay = KillWait +} + +// Labelled is a Runner that marks every container a build starts with the build's id, so a kill +// finds what outlived the docker client (novox/hq ADR 0219). Applied at the one place every command +// passes rather than at each `docker run` the builder composes: a run added later is labelled too. +func Labelled(run Runner, id string) Runner { + return func(ctx context.Context, dir string, name string, args ...string) (string, error) { + return run(ctx, dir, name, LabelledArgs(name, args, id)...) + } +} + +// LabelledArgs is a command's arguments with the build's label added, when it is a `docker run`. +func LabelledArgs(name string, args []string, id string) []string { + if name != "docker" || len(args) == 0 || args[0] != "run" || id == "" { + return args + } + out := make([]string, 0, len(args)+2) + out = append(out, "run", "--label", BuildLabel+"="+id) + return append(out, args[1:]...) +} + +// RemoveContainersOf removes every container labelled with the build's id, running or not, and +// says how many. Run with a context of its own: the build's is the one that was just cancelled. +func RemoveContainersOf(ctx context.Context, run Runner, id string) (int, error) { + out, err := run(ctx, "", "docker", "ps", "-aq", "--filter", "label="+BuildLabel+"="+id) + if err != nil { + return 0, err + } + ids := strings.Fields(out) + if len(ids) == 0 { + return 0, nil + } + if _, err := run(ctx, "", "docker", append([]string{"rm", "-f"}, ids...)...); err != nil { + return 0, err + } + return len(ids), nil +} diff --git a/internal/builder/kill_test.go b/internal/builder/kill_test.go new file mode 100644 index 0000000..03a7c90 --- /dev/null +++ b/internal/builder/kill_test.go @@ -0,0 +1,103 @@ +package builder + +import ( + "context" + "errors" + "os" + "path/filepath" + "reflect" + "strconv" + "strings" + "syscall" + "testing" + "time" +) + +// A build killed by hand (novox/hq ADR 0219) ends every command it started: the context's end kills +// the command's whole process group, not the one process the context knows about. +func TestAKilledBuildEndsItsWholeProcessGroup(t *testing.T) { + dir := t.TempDir() + child := filepath.Join(dir, "child") + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + began := time.Now() + go func() { + // A shell that starts a grandchild and waits on it: killing the shell alone would leave the + // grandchild running, holding the output open. + _, err := Command(ctx, dir, "sh", "-c", `sleep 60 & echo $! > child; wait`) + done <- err + }() + var pid int + for deadline := time.Now().Add(5 * time.Second); time.Now().Before(deadline); time.Sleep(20 * time.Millisecond) { + if raw, err := os.ReadFile(child); err == nil && strings.TrimSpace(string(raw)) != "" { + pid, _ = strconv.Atoi(strings.TrimSpace(string(raw))) + break + } + } + if pid == 0 { + t.Fatal("the command never started its child") + } + cancel() + select { + case err := <-done: + if err == nil { + t.Fatal("a killed command reported success") + } + case <-time.After(KillWait + 5*time.Second): + t.Fatal("the killed command never returned") + } + if took := time.Since(began); took > KillWait { + t.Errorf("ending the command took %s: the group was not killed, the wait ran out", took) + } + for deadline := time.Now().Add(2 * time.Second); time.Now().Before(deadline); time.Sleep(20 * time.Millisecond) { + if err := syscall.Kill(pid, 0); errors.Is(err, syscall.ESRCH) { + return + } + } + _ = syscall.Kill(pid, syscall.SIGKILL) + t.Fatalf("the grandchild %d outlived the kill", pid) +} + +// Every `docker run` a build starts carries its id as a label; nothing else is touched. +func TestEveryContainerABuildStartsCarriesItsID(t *testing.T) { + got := LabelledArgs("docker", []string{"run", "--rm", "img", "sh"}, "build-1") + if want := []string{"run", "--label", "mesh.build=build-1", "--rm", "img", "sh"}; !reflect.DeepEqual(got, want) { + t.Errorf("docker run became %v", got) + } + for _, c := range []struct { + name string + args []string + }{{"docker", []string{"build", "."}}, {"git", []string{"run"}}, {"docker", nil}} { + if got := LabelledArgs(c.name, c.args, "build-1"); !reflect.DeepEqual(got, c.args) { + t.Errorf("%s %v became %v", c.name, c.args, got) + } + } + var ran [][]string + run := Labelled(func(_ context.Context, _ string, name string, args ...string) (string, error) { + ran = append(ran, append([]string{name}, args...)) + return "", nil + }, "build-2") + _, _ = run(context.Background(), "", "docker", "run", "img") + if want := [][]string{{"docker", "run", "--label", "mesh.build=build-2", "img"}}; !reflect.DeepEqual(ran, want) { + t.Errorf("ran %v", ran) + } +} + +// What a kill removes is found by the label, and only what it finds. +func TestAKillRemovesTheContainersLabelledWithTheBuild(t *testing.T) { + var ran []string + run := func(_ context.Context, _ string, name string, args ...string) (string, error) { + ran = append(ran, name+" "+strings.Join(args, " ")) + if args[0] == "ps" { + return "c1\nc2\n", nil + } + return "", nil + } + n, err := RemoveContainersOf(context.Background(), run, "build-3") + if err != nil || n != 2 { + t.Fatalf("%d %v", n, err) + } + if want := []string{"docker ps -aq --filter label=mesh.build=build-3", "docker rm -f c1 c2"}; !reflect.DeepEqual(ran, want) { + t.Errorf("ran %v", ran) + } +} diff --git a/internal/catalogue/seats.go b/internal/catalogue/seats.go index 39a35ef..7764096 100644 --- a/internal/catalogue/seats.go +++ b/internal/catalogue/seats.go @@ -117,8 +117,16 @@ var defaultSeats = append([]Seat{ // **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. + // + // **And its holder answers for the build it is running** (novox/hq ADR 0219): what it is + // building, kill it, take nothing new, take again. The queue as a whole is the controller's to + // show and change (`queue`, `cancel`, `clear`); what one machine does with an ask it already + // took only that machine can do. `paused.` is each holder saying whether it takes work, so + // the controller can tell a plan waiting on a paused seat from one that is late. {Name: "node-build-agent", Scope: ScopeNode, - Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0190"}, + Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*", "paused.*"}, + Serves: buildAgentVerbs(), + Decision: "novox/hq ADR 0190, ADR 0219"}, // **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 @@ -581,3 +589,23 @@ func loginShellVerbs() []Verb { }, []string{"command"})}, } } + +// buildAgentVerbs are what a holder of node-build-agent answers on its own machine (novox/hq ADR +// 0219). One build at a time per holder (ADR 0190), so "the build running here" is one or none. +func buildAgentVerbs() []Verb { + return []Verb{ + {Name: "current", Description: "The build this machine is running — its id, repository, path, " + + "the step it is at, when it started and for how long — or none; and whether this machine is " + + "paused, taking no new build.", + Input: schema(map[string]string{}, nil)}, + {Name: "kill", Description: "End the build with this id, running here: its commands and the " + + "containers it started are stopped, and its outcome is announced as failed, killed by hand — " + + "settled, so it is not handed to another machine.", + Input: schema(map[string]string{"id": "the build's id, as `queue` or `current` says it"}, []string{"id"})}, + {Name: "pause", Description: "Take no new build on this machine until resumed; a build running " + + "here finishes. Kept across a restart of the holder.", + Input: schema(map[string]string{}, nil)}, + {Name: "resume", Description: "Take builds again on this machine.", + Input: schema(map[string]string{}, nil)}, + } +} diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index 21bdca4..343b371 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -5,6 +5,7 @@ import ( "encoding/json" "errors" "fmt" + "os" "time" "github.com/nats-io/nats.go" @@ -131,6 +132,14 @@ type natsMachine struct { on string seat string sub *nats.Subscription + opts MachineOptions +} + +// MachineOptions is what a holder tells the taking loop about itself (novox/hq ADR 0219). +type MachineOptions struct { + // Paused is asked before every fetch: while it says so, nothing new is taken, and a build + // already running finishes. Nil is never paused. + Paused func() bool } // MachineOverNATS takes build work from the current build role. @@ -145,6 +154,14 @@ func MachineOverNATSOn(js *broker.JetStream, on, seat string) BuildMachine { return &natsMachine{js: js, on: on, seat: seat} } +// MachineOverNATSWith is MachineOverNATSOn for a holder that can be paused (novox/hq ADR 0219). +func MachineOverNATSWith(js *broker.JetStream, on, seat string, opts MachineOptions) BuildMachine { + return &natsMachine{js: js, on: on, seat: seat, opts: opts} +} + +// pausedPoll is how often a paused holder looks again whether it was resumed. +const pausedPoll = 2 * time.Second + func (m *natsMachine) Close() { if m.sub != nil { _ = m.sub.Unsubscribe() @@ -189,6 +206,17 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) if ctx.Err() != nil { return nil } + // **Paused takes nothing new** (novox/hq ADR 0219). Asked before the fetch, never during a + // build: what this machine already took it finishes, and what it has not taken stays in the + // queue for another holder — or for this one, resumed. + if m.opts.Paused != nil && m.opts.Paused() { + select { + case <-ctx.Done(): + return nil + case <-time.After(pausedPoll): + } + continue + } // 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)) @@ -220,12 +248,25 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) _ = msg.Term() continue } + // **Cancelled in the moment it was fetched** (novox/hq ADR 0219): the controller wrote the + // id into the seat's cancelled set before deleting the ask, so one taken in between is + // ended here, terminated rather than redelivered, and its outcome said as failed — never + // built. A set that cannot be read is said and the ask built: a cancel is a person's + // exception, and a holder that refused every build while its grant was missing would + // stop the mesh building for it. + build := &natsBuild{request: request, msg: msg, on: m.on, js: m.js, seat: m.seat} + if cancelled, err := IsCancelled(m.js.Conn(), m.seat, request.ID); err != nil { + fmt.Fprintf(os.Stderr, "could not read whether %s was cancelled, so it is built: %v\n", request.ID, err) + } else if cancelled { + endCancelled(ctx, build) + continue + } // A build outlives the acknowledgement window many times over; said while it runs, // as the controller says it for its own long handlers, so the server neither hands // 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, seat: m.seat}) + do(ctx, build) close(working) } } @@ -307,3 +348,17 @@ func (b *natsBuild) Say(step, message string) { } func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) } + +// endCancelled settles an ask that was cancelled as it was taken: its outcome said as failed, as the +// controller already recorded it, so anybody still waiting on it hears the answer — then +// terminated, so the queue neither keeps nor redelivers it. +func endCancelled(ctx context.Context, b *natsBuild) { + r := b.request + fmt.Fprintf(os.Stderr, "%s was cancelled by hand as this machine took it; not built\n", r.ID) + result := BuildResult{ID: r.ID, Repository: r.Repository, Path: r.Path, Ref: r.Ref, On: b.on, + Source: r.Source, DryRun: r.DryRun, Failed: CancelledByHand} + if err := b.Announce(ctx, result); err != nil { + fmt.Fprintf(os.Stderr, "cannot say %s was cancelled: %v\n", r.ID, err) + } + _ = b.msg.Term() +} diff --git a/internal/link/queue.go b/internal/link/queue.go new file mode 100644 index 0000000..2855129 --- /dev/null +++ b/internal/link/queue.go @@ -0,0 +1,439 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "regexp" + "sort" + "time" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" + + "github.com/novox/mesh-controller/internal/broker" +) + +// The build queue, read and changed by hand (novox/hq ADR 0219). +// +// **A queue nobody can see is a queue nobody can trust.** Until this, the build seat's work queue +// was visible only as its consequences: a plan LATE with nothing to say why, a build that never +// started because five asks ahead of it were waiting on a laptop that was shut. ADR 0219 makes the +// queue the controller's to show and change, and the build running on one machine that machine's +// holder's to end. +// +// Three kinds of ask are in the seat's stream at any moment, told apart by the worker consumer +// every holder pulls from: +// +// - **waiting** — after the last one the worker handed out (stream seq > delivered.stream_seq); +// - **in flight** — handed out and not yet settled, which the holder that took it said with its +// `started` event, and which the worker counts among its acks pending; +// - **dead** — handed out as many times as the worker allows and never settled. A work queue +// drops what it settles, so one still in the stream behind the worker is either in flight or +// dead; the started events and the worker's pending count say which. +// +// Everything that drops an ask leaves a failed outcome for it, recorded the way a failed build's is +// (`cancelled by hand`, `killed by hand`), so a plan waiting on it fails visibly rather than waiting +// for ever on an ask nobody will answer. + +// The words an outcome says when a person ended the ask. +const ( + CancelledByHand = "cancelled by hand" + KilledByHand = "killed by hand" +) + +// Ask states in a queue. +const ( + AskWaiting = "waiting" + AskInFlight = "in flight" + AskDead = "dead" +) + +// QueuedAsk is one ask in the seat's work queue. +type QueuedAsk struct { + Seq uint64 `json:"seq"` + Request BuildRequest `json:"-"` + ID string `json:"id"` + // Repository, Path and Ref are the ask's, as the request carried them; Held never travels out of + // here — it is every artifact the mesh has built, and a queue listing is not where that belongs. + Repository string `json:"repository"` + Path string `json:"path,omitempty"` + Ref string `json:"ref,omitempty"` + AskedAt time.Time `json:"asked-at,omitempty"` + State string `json:"state"` + // On and Started are the holder that took an in-flight ask and when, from its started event. + On string `json:"on,omitempty"` + Started time.Time `json:"started,omitempty"` +} + +// Queue is the seat's work queue as it stands. +type Queue struct { + Seat string `json:"seat"` + // Delivered is the last stream sequence the worker handed out; AckPending how many it handed out + // and has not had settled; MaxDeliver how many times it hands one ask out. + Delivered uint64 `json:"delivered"` + AckPending int `json:"ack-pending"` + MaxDeliver int `json:"max-deliver"` + Asks []QueuedAsk `json:"asks"` +} + +// Of is the asks in one state, in queue order. +func (q Queue) Of(state string) []QueuedAsk { + var out []QueuedAsk + for _, a := range q.Asks { + if a.State == state { + out = append(out, a) + } + } + return out +} + +// Find is the ask with this id. +func (q Queue) Find(id string) (QueuedAsk, bool) { + for _, a := range q.Asks { + if a.ID == id { + return a, true + } + } + return QueuedAsk{}, false +} + +// BuildHeard is what the seat's events say about builds: the latest start of each id, and every id +// whose outcome was heard. +type BuildHeard struct { + Started map[string]BuildStart + Outcomes map[string]bool +} + +// ClassifyQueue says which state each ask is in. Pure: the stream's asks in sequence order, what the +// worker says, and what the seat's events said. +// +// **In flight is what a machine is building now.** A holder builds one ask at a time (ADR 0190), so +// what a machine is building is its latest start, if no outcome has been heard for it. An ask +// behind the worker that is some machine's current build is in flight; at most AckPending of them +// are, newest start first — a machine that died mid-build has a latest start forever, and the +// worker's own count is what bounds that. An ask behind the worker that no machine has said it +// started is in flight only while the worker counts more pending than starts explain (taken, and +// not yet said); otherwise it is dead. +// +// **The worker's count is the server's, and it lags in one case**: an ask handed out its last time +// and handed back is still counted pending until some holder pulls again, when the server finds it +// past its deliveries and lets it go. Holders pull whenever they are idle, so this is a moment — +// except while every holder is paused, when such an ask reads as in flight with no machine named. +// Cancel takes an ask in flight that no machine said it started, for exactly that reason. +func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard BuildHeard) []QueuedAsk { + // Each machine's latest start. + latest := map[string]BuildStart{} + for _, s := range heard.Started { + at := startedAt(s) + if was, ok := latest[s.On]; !ok || at.After(startedAt(was)) { + latest[s.On] = s + } + } + current := map[string]BuildStart{} + for _, s := range latest { + if !heard.Outcomes[s.ID] { + current[s.ID] = s + } + } + + out := make([]QueuedAsk, len(asks)) + copy(out, asks) + var running, unsaid []int + for i := range out { + a := &out[i] + if a.Seq > delivered { + a.State = AskWaiting + continue + } + a.State = AskDead + if s, ok := current[a.ID]; ok { + a.On, a.Started = s.On, startedAt(s) + running = append(running, i) + continue + } + if _, ever := heard.Started[a.ID]; !ever { + unsaid = append(unsaid, i) + } + } + sort.SliceStable(running, func(x, y int) bool { return out[running[x]].Started.After(out[running[y]].Started) }) + slots := ackPending + for _, i := range running { + if slots == 0 { + // A machine's latest start the worker no longer counts: it died with the ask, and the + // ask was handed out as often as it may be. + out[i].On, out[i].Started = "", time.Time{} + continue + } + out[i].State = AskInFlight + slots-- + } + for _, i := range unsaid { + if slots == 0 { + break + } + out[i].State = AskInFlight + slots-- + } + return out +} + +func startedAt(s BuildStart) time.Time { + t, _ := time.Parse(time.RFC3339Nano, s.At) + return t +} + +// ReadQueue reads a build seat's work queue: the stream's asks, the worker's position, and what the +// seat's events said about the builds behind it. Through the JetStream API alone (STREAM.INFO, +// CONSUMER.INFO, STREAM.MSG.GET), which only the controller's account reaches. +func ReadQueue(ctx context.Context, js *broker.JetStream, seat string) (Queue, error) { + worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: seat, Accepts: []string{"build"}}) + q := Queue{Seat: seat} + api, err := jetstream.New(js.Conn()) + if err != nil { + return q, err + } + stream, err := api.Stream(ctx, worker.Stream) + if err != nil { + return q, fmt.Errorf("the %s work queue (%s) cannot be read: %w", seat, worker.Stream, err) + } + consumer, err := stream.Consumer(ctx, worker.Name) + switch { + case errors.Is(err, jetstream.ErrConsumerNotFound): + // No holder has ever been given a worker: everything in the stream is waiting. + case err != nil: + return q, fmt.Errorf("the worker of %s cannot be read: %w", seat, err) + default: + info, err := consumer.Info(ctx) + if err != nil { + return q, fmt.Errorf("the worker of %s cannot be read: %w", seat, err) + } + q.Delivered = info.Delivered.Stream + q.AckPending = info.NumAckPending + q.MaxDeliver = info.Config.MaxDeliver + } + + // Every ask, by next-by-subject from the first: the stream holds only what is unsettled, so this + // is a handful of reads, never the history. + var asks []QueuedAsk + var seq uint64 = 1 + for n := 0; n < 10000; n++ { + msg, err := stream.GetMsg(ctx, seq, jetstream.WithGetMsgSubject(BuildWorkOf(seat))) + if errors.Is(err, jetstream.ErrMsgNotFound) { + break + } + if err != nil { + return q, fmt.Errorf("reading the %s work queue: %w", seat, err) + } + ask := QueuedAsk{Seq: msg.Sequence} + if json.Unmarshal(msg.Data, &ask.Request) == nil { + r := ask.Request + ask.ID, ask.Repository, ask.Path, ask.Ref = r.ID, r.Repository, r.Path, r.Ref + if at, ok := BuildAskedAt(r.ID); ok { + ask.AskedAt = at + } + } + asks = append(asks, ask) + seq = msg.Sequence + 1 + } + + // What the events say, from shortly before the oldest ask behind the worker: its start, if any, + // came after it was asked. + var since time.Time + for _, a := range asks { + if a.Seq <= q.Delivered && (since.IsZero() || (!a.AskedAt.IsZero() && a.AskedAt.Before(since))) { + since = a.AskedAt + } + } + heard := BuildHeard{Started: map[string]BuildStart{}, Outcomes: map[string]bool{}} + if len(asks) > 0 && q.Delivered > 0 { + if since.IsZero() { + since = time.Now().Add(-7 * 24 * time.Hour) + } + if heard, err = ReadBuildEvents(ctx, js, seat, since.Add(-time.Minute)); err != nil { + return q, err + } + } + q.Asks = ClassifyQueue(asks, q.Delivered, q.AckPending, heard) + return q, nil +} + +// ReadBuildEvents is every start and outcome the seat announced since a moment, from the events +// stream, with a consumer of its own that is gone when this returns. +func ReadBuildEvents(ctx context.Context, js *broker.JetStream, seat string, since time.Time) (BuildHeard, error) { + heard := BuildHeard{Started: map[string]BuildStart{}, Outcomes: map[string]bool{}} + // One token after `event`: started and built, never a build's log lines, which are two. + subject := "mesh.seat." + seat + ".event.*" + sub, err := js.Context().PullSubscribe(subject, "", + nats.BindStream(broker.EventsStream), nats.StartTime(since), nats.AckNone()) + if err != nil { + return heard, fmt.Errorf("cannot read %s from the bus: %w", subject, err) + } + defer func() { _ = sub.Unsubscribe() }() + for { + batch, err := sub.Fetch(500, nats.MaxWait(2*time.Second)) + if err != nil && !errors.Is(err, nats.ErrTimeout) && !errors.Is(err, context.DeadlineExceeded) { + return heard, fmt.Errorf("reading %s: %w", subject, err) + } + for _, msg := range batch { + switch msg.Subject { + case BuildStartedOf(seat): + var s BuildStart + if json.Unmarshal(msg.Data, &s) == nil && s.ID != "" { + if was, ok := heard.Started[s.ID]; !ok || startedAt(s).After(startedAt(was)) { + heard.Started[s.ID] = s + } + } + case BuildOutcomeOf(seat): + var r BuildResult + if json.Unmarshal(msg.Data, &r) == nil && r.ID != "" { + heard.Outcomes[r.ID] = true + } + } + } + if len(batch) < 500 { + break + } + if ctx.Err() != nil { + return heard, ctx.Err() + } + } + return heard, nil +} + +// --- the cancelled set ---------------------------------------------------------------------- + +// safeKey is an id that can be a key, and a subject token, as it is: a build id, never a pattern. +var safeKey = regexp.MustCompile(`^[A-Za-z0-9_-]+$`) + +// Cancelled is what the controller writes for an ask it cancelled. +type Cancelled struct { + At string `json:"at"` + Why string `json:"why"` +} + +// MarkCancelled puts an ask's id into the seat's cancelled set (novox/hq ADR 0219). The controller's +// act, before it deletes the ask from the queue. +func MarkCancelled(ctx context.Context, js *broker.JetStream, seat, id string) error { + if !safeKey.MatchString(id) { + return fmt.Errorf("%q is not a build id", id) + } + api, err := jetstream.New(js.Conn()) + if err != nil { + return err + } + kv, err := api.KeyValue(ctx, broker.CancelledSetName(seat)) + if err != nil { + return fmt.Errorf("the cancelled set of %s is not on the bus — the controller asserts it at its "+ + "start, so one older than this has not: %w", seat, err) + } + body, err := json.Marshal(Cancelled{At: time.Now().UTC().Format(time.RFC3339Nano), Why: CancelledByHand}) + if err != nil { + return err + } + _, err = kv.Put(ctx, id, body) + return err +} + +// IsCancelled asks the seat's cancelled set whether an ask was cancelled: one direct read of one key, +// which is all a holder is granted of it. +func IsCancelled(conn *nats.Conn, seat, id string) (bool, error) { + if !safeKey.MatchString(id) { + // Not a key the controller could have written. + return false, nil + } + set := broker.CancelledSetName(seat) + reply, err := conn.Request("$JS.API.DIRECT.GET.KV_"+set+".$KV."+set+"."+id, nil, 3*time.Second) + if err != nil { + return false, err + } + return cancelledFrom(reply) +} + +// cancelledFrom reads a direct get's answer: a value is a cancel; not found, or a deletion marker, +// is not. +func cancelledFrom(reply *nats.Msg) (bool, error) { + if reply.Header != nil { + switch reply.Header.Get("Status") { + case "": + case "404": + return false, nil + default: + return false, fmt.Errorf("the cancelled set answered %s %s", + reply.Header.Get("Status"), reply.Header.Get("Description")) + } + if op := reply.Header.Get("KV-Operation"); op == "DEL" || op == "PURGE" { + return false, nil + } + } + return len(reply.Data) > 0, nil +} + +// --- a holder's verbs ----------------------------------------------------------------------- + +// NodeSeatToolSubject is where a node-scoped seat's verb is asked of one machine's holder (design 33 +// §4): the seat's verb subject with the machine as its last token. +func NodeSeatToolSubject(seat, verb, node string) string { + return SeatToolSubject(seat, verb) + "." + node +} + +// AskSeatTool asks one machine's holder of a node seat one of its verbs and reads its answer. +func AskSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string, args any, + timeout time.Duration) (Answer, error) { + body, err := json.Marshal(args) + if err != nil { + return Answer{}, err + } + asking, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + reply, err := conn.RequestWithContext(asking, NodeSeatToolSubject(seat, verb, node), body) + switch { + case errors.Is(err, nats.ErrNoResponders): + return Answer{}, fmt.Errorf("nothing on %s answers %s.%s: its holder is not running, or is "+ + "older than the verb", node, seat, verb) + case errors.Is(err, context.DeadlineExceeded), errors.Is(err, nats.ErrTimeout): + return Answer{}, fmt.Errorf("%s did not answer %s.%s within %s", node, seat, verb, timeout) + case err != nil: + return Answer{}, err + } + var answer Answer + if err := json.Unmarshal(reply.Data, &answer); err != nil { + return Answer{}, fmt.Errorf("%s answered %s.%s with something unreadable: %w", node, seat, verb, err) + } + return answer, nil +} + +// --- a holder saying whether it takes work ---------------------------------------------------- + +// HolderState is what a holder says about itself whenever it is paused or resumed, at its start, +// and while it stays paused: whether it takes work (novox/hq ADR 0219). Retained on the events +// stream under the machine's own subject, so the last one is the machine's state. +type HolderState struct { + On string `json:"on"` + Paused bool `json:"paused"` + At string `json:"at"` +} + +// BuildPausedOf is where one machine's holder of a build seat says whether it takes work. +func BuildPausedOf(seat, node string) string { return "mesh.seat." + seat + ".event.paused." + node } + +// PausedSaid reads what each machine last said about taking work. A machine that never said is not +// in the answer — a holder older than ADR 0219 takes work and never pauses. +func PausedSaid(js *broker.JetStream, seat string, nodes []string) (map[string]HolderState, error) { + out := map[string]HolderState{} + for _, node := range nodes { + msg, err := js.Context().GetLastMsg(broker.EventsStream, BuildPausedOf(seat, node)) + if errors.Is(err, nats.ErrMsgNotFound) { + continue + } + if err != nil { + return out, err + } + var s HolderState + if json.Unmarshal(msg.Data, &s) == nil { + out[node] = s + } + } + return out, nil +} diff --git a/internal/link/queue_test.go b/internal/link/queue_test.go new file mode 100644 index 0000000..6f824b3 --- /dev/null +++ b/internal/link/queue_test.go @@ -0,0 +1,231 @@ +package link + +import ( + "context" + "encoding/json" + "testing" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" +) + +// The build queue read and changed by hand (novox/hq ADR 0219). + +// Which ask is waiting, in flight and dead, from the stream, the worker and the seat's events. +func TestTheQueueTellsWaitingInFlightAndDeadApart(t *testing.T) { + at := func(m int) string { + return time.Date(2026, 10, 5, 12, m, 0, 0, time.UTC).Format(time.RFC3339Nano) + } + asks := []QueuedAsk{ + {Seq: 1, ID: "dead-1"}, // handed out five times, its machine moved on: dead + {Seq: 2, ID: "running-g"}, // g14's latest start, no outcome: in flight + {Seq: 3, ID: "running-a"}, // ace's latest start, no outcome: in flight + {Seq: 4, ID: "unsaid"}, // taken, not yet said: in flight while the worker counts it + {Seq: 5, ID: "waiting-1"}, // after the worker's last delivery + {Seq: 6, ID: "waiting-2"}, + } + heard := BuildHeard{ + Started: map[string]BuildStart{ + "dead-1": {ID: "dead-1", On: "g14", At: at(1)}, + "running-g": {ID: "running-g", On: "g14", At: at(5)}, + "running-a": {ID: "running-a", On: "ace", At: at(6)}, + "done": {ID: "done", On: "novox", At: at(7)}, + }, + Outcomes: map[string]bool{"done": true}, + } + got := ClassifyQueue(asks, 4, 3, heard) + want := map[string]string{"dead-1": AskDead, "running-g": AskInFlight, "running-a": AskInFlight, + "unsaid": AskInFlight, "waiting-1": AskWaiting, "waiting-2": AskWaiting} + for _, a := range got { + if a.State != want[a.ID] { + t.Errorf("%s is %s, want %s", a.ID, a.State, want[a.ID]) + } + } + q := Queue{Asks: got} + if a, _ := q.Find("running-a"); a.On != "ace" || a.Started.IsZero() { + t.Errorf("an ask in flight does not say where: %+v", a) + } + + // **The worker's count bounds it**: with two pending, the machine whose start is oldest died with + // its ask, and the ask nobody said is not in flight either. + got = ClassifyQueue(asks, 4, 2, heard) + q = Queue{Asks: got} + if a, _ := q.Find("running-g"); a.State != AskInFlight { + t.Errorf("running-g is %s", a.State) + } + if a, _ := q.Find("unsaid"); a.State != AskDead { + t.Errorf("an ask no start explains, past the worker's count, is %s", a.State) + } + // None pending: everything behind the worker is dead, and names no machine. + for _, a := range ClassifyQueue(asks, 4, 0, heard) { + if a.Seq <= 4 && (a.State != AskDead || a.On != "") { + t.Errorf("%s with nothing pending is %s on %q", a.ID, a.State, a.On) + } + } +} + +// What a direct read of the cancelled set answers. +func TestACancelledSetReadSaysWhatWasCancelled(t *testing.T) { + found := &nats.Msg{Header: nats.Header{"Nats-Subject": []string{"$KV.x.build-1"}}, Data: []byte(`{"why":"cancelled by hand"}`)} + if c, err := cancelledFrom(found); err != nil || !c { + t.Errorf("a value read as %v %v", c, err) + } + missing := &nats.Msg{Header: nats.Header{"Status": []string{"404"}}} + if c, err := cancelledFrom(missing); err != nil || c { + t.Errorf("not found read as %v %v", c, err) + } + deleted := &nats.Msg{Header: nats.Header{"KV-Operation": []string{"DEL"}}} + if c, err := cancelledFrom(deleted); err != nil || c { + t.Errorf("a deletion marker read as %v %v", c, err) + } + broken := &nats.Msg{Header: nats.Header{"Status": []string{"408"}, "Description": []string{"Request Timeout"}}} + if _, err := cancelledFrom(broken); err == nil { + t.Error("an error answer read as an answer") + } + if c, err := IsCancelled(nil, TheBuildMachine, "a.b.>"); err != nil || c { + t.Error("a pattern was looked up as a key") + } +} + +// --- against a real server ---------------------------------------------------------------------- + +func aBusWithACancelledSet(t *testing.T) *broker.JetStream { + t.Helper() + js := aBusWithTheBuildRole(t) + seat := broker.DeclaredSeat{Name: TheBuildMachine, Accepts: []string{"build"}} + if err := broker.RaiseCancelledSets(js, []broker.DeclaredSeat{seat}); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = js.Context().DeleteKeyValue(broker.CancelledSetName(TheBuildMachine)) }) + return js +} + +// **The race a delete alone leaves open**: an ask fetched in the moment it was cancelled is ended by +// the holder that took it — terminated, never built, its outcome said as failed. +func TestNatsAnAskCancelledAsItWasTakenIsNotBuilt(t *testing.T) { + js := aBusWithACancelledSet(t) + ctx, stop := context.WithCancel(context.Background()) + defer stop() + + if err := MarkCancelled(ctx, js, TheBuildMachine, "build-cancelled"); err != nil { + t.Fatal(err) + } + outcomes, err := js.Conn().SubscribeSync(BuildOutcome()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = outcomes.Unsubscribe() }() + _ = js.Conn().Flush() + for _, id := range []string{"build-cancelled", "build-kept"} { + body, _ := json.Marshal(BuildRequest{ID: id, Repository: "/r"}) + if _, err := js.Context().Publish(BuildWork(), body); err != nil { + t.Fatal(err) + } + } + + built := make(chan string, 2) + machine := MachineOverNATS(js, "anchor") + defer machine.Close() + go func() { + _ = machine.Take(ctx, func(ctx context.Context, work Build) { + built <- work.Request().ID + _ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: "anchor", Failed: "no"}) + _ = work.Done() + }) + }() + select { + case id := <-built: + if id != "build-kept" { + t.Fatalf("%s was built", id) + } + case <-time.After(10 * time.Second): + t.Fatal("the ask that was not cancelled was never built") + } + msg, err := outcomes.NextMsg(5 * time.Second) + if err != nil { + t.Fatal(err) + } + var said BuildResult + if err := json.Unmarshal(msg.Data, &said); err != nil || said.ID != "build-cancelled" || + said.Failed != CancelledByHand || said.On != "anchor" { + t.Fatalf("the cancelled ask's outcome is %+v (%v)", said, err) + } + deadline := time.Now().Add(5 * time.Second) + for time.Now().Before(deadline) { + if info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT"); err == nil && info.State.Msgs == 0 { + return + } + time.Sleep(20 * time.Millisecond) + } + t.Fatal("the cancelled ask is still queued: terminated, it would not be") +} + +// A paused holder takes nothing; resumed, it takes what waited. +func TestNatsAPausedHolderTakesNothingNew(t *testing.T) { + js := aBusWithACancelledSet(t) + ctx, stop := context.WithCancel(context.Background()) + defer stop() + paused := make(chan bool, 1) + paused <- true + isPaused := func() bool { + p := <-paused + paused <- p + return p + } + took := make(chan string, 1) + machine := MachineOverNATSWith(js, "anchor", TheBuildMachine, MachineOptions{Paused: isPaused}) + defer machine.Close() + go func() { + _ = machine.Take(ctx, func(ctx context.Context, work Build) { + took <- work.Request().ID + _ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: "anchor", Failed: "no"}) + _ = work.Done() + }) + }() + body, _ := json.Marshal(BuildRequest{ID: "build-waits", Repository: "/r"}) + if _, err := js.Context().Publish(BuildWork(), body); err != nil { + t.Fatal(err) + } + select { + case id := <-took: + t.Fatalf("a paused holder took %s", id) + case <-time.After(3 * time.Second): + } + q, err := ReadQueue(ctx, js, TheBuildMachine) + if err != nil { + t.Fatal(err) + } + if len(q.Asks) != 1 || q.Asks[0].State != AskWaiting { + t.Fatalf("while paused the queue reads %+v", q.Asks) + } + <-paused + paused <- false + select { + case id := <-took: + if id != "build-waits" { + t.Fatalf("took %s", id) + } + case <-time.After(10 * time.Second): + t.Fatal("resumed, the holder never took what waited") + } +} + +// What a holder last said about taking work is read back per machine. +func TestNatsAHoldersPausedStateIsReadBack(t *testing.T) { + js := aBusWithTheBuildRole(t) + for _, s := range []HolderState{{On: "ace", Paused: true}, {On: "g14", Paused: true}, {On: "g14", Paused: false}} { + body, _ := json.Marshal(s) + if _, err := js.Context().Publish(BuildPausedOf(TheBuildMachine, s.On), body); err != nil { + t.Fatal(err) + } + } + said, err := PausedSaid(js, TheBuildMachine, []string{"ace", "g14", "novox"}) + if err != nil { + t.Fatal(err) + } + if !said["ace"].Paused || said["g14"].Paused || len(said) != 2 { + t.Fatalf("read back %+v", said) + } +} diff --git a/internal/link/seattools.go b/internal/link/seattools.go index f579c53..89b24a5 100644 --- a/internal/link/seattools.go +++ b/internal/link/seattools.go @@ -42,6 +42,18 @@ const RebindAfter = 30 * time.Second // tried again (2026-09-30). A refused subscription is therefore retried until it holds: the server // says so asynchronously and invalidates the subscription, which is what is checked. func (b OverNATS) ServeSeatTools(seat string, handlers map[string]ToolHandler, logger *log.Logger) (func(), error) { + return b.serveTools(seat, func(verb string) string { return SeatToolSubject(seat, verb) }, handlers, logger) +} + +// ServeNodeSeatTools is ServeSeatTools for one machine's holder of a node-scoped seat (novox/hq ADR +// 0159, ADR 0219): each verb on the seat's subject for this machine and no other, so a call names +// the machine it is for and only that machine's holder answers it. +func (b OverNATS) ServeNodeSeatTools(seat, node string, handlers map[string]ToolHandler, logger *log.Logger) (func(), error) { + return b.serveTools(seat, func(verb string) string { return NodeSeatToolSubject(seat, verb, node) }, handlers, logger) +} + +func (b OverNATS) serveTools(seat string, subjectOf func(string) string, handlers map[string]ToolHandler, + logger *log.Logger) (func(), error) { var subs []*nats.Subscription done := make(chan struct{}) stop := func() { @@ -52,7 +64,7 @@ func (b OverNATS) ServeSeatTools(seat string, handlers map[string]ToolHandler, l } for verb, handle := range handlers { verb, handle := verb, handle - subject := SeatToolSubject(seat, verb) + subject := subjectOf(verb) bind := func() (*nats.Subscription, error) { return b.Conn.QueueSubscribe(subject, "seat."+seat, func(msg *nats.Msg) { // Its own goroutine per call: a slow `push` must not hold up a `status` asked beside it,