Show and change the build queue through the controller, and have plans follow it (hq ADR 0219)

Nothing showed what waited for a build machine, and an ask could not be
dropped without leaving the plan that made it waiting for ever. New verbs:
queue, cancel, clear, rebuild, replay, kill, pause, resume, and plans retry.
Every ask a person drops is recorded failed through the same take-in as a
failed build; a plan keeps the id it asked each module under and matches
its outcome by it. replay is a dry run unless registered, and registering
an older commit than one registered since needs --older (hq issue 207).
A plan waiting on a seat paused on every holder says so and is not late;
a failed plan can be retried, and a rebuild joins the plan holding the
module instead of running beside it.
This commit is contained in:
jochen
2026-10-05 19:17:56 +02:00
parent e610f2d92c
commit 106507b1d3
13 changed files with 1772 additions and 58 deletions
+583
View File
@@ -0,0 +1,583 @@
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"sort"
"strings"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// The build queue, controlled by hand (novox/hq ADR 0219).
//
// Seen whole — what waits, what runs where and for how long, what was handed out as often as it
// may be and never settled — and changed through the controller: an ask cancelled, the queue
// cleared, a module rebuilt, a build replayed. What one machine is running is that machine's
// holder's to end or pause, and the controller asks it (`kill`, `pause`, `resume`).
//
// **Everything that drops an ask leaves a failed outcome for it**, taken in exactly as a build that
// failed is (takeIn, then the plan): a plan waiting on an ask a person removed fails, saying so,
// rather than waiting for ever on an answer nobody will give.
// holderAsks is how long a holder's verb is waited for.
const holderAsks = 15 * time.Second
// dialTheBus opens the controller's own connection, for a command that reads or changes the queue.
func dialTheBus() (*broker.JetStream, error) {
address, err := broker.BusAddress()
if err != nil {
return nil, err
}
js, err := broker.Dial(address)
if err != nil {
return nil, fmt.Errorf("cannot reach the bus: %w", err)
}
return js, nil
}
// queueCommand prints every ask in the build seat's work queue.
func queueCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("queue", flag.ContinueOnError)
asJSON := set.Bool("json", false, "the queue as JSON")
if _, err := parseAround(set, args); err != nil {
return err
}
js, err := dialTheBus()
if err != nil {
return err
}
defer js.Close()
seat := buildSeatHeld(ctx)
q, err := link.ReadQueue(ctx, js, seat)
if err != nil {
return err
}
if *asJSON {
body, err := json.MarshalIndent(q, "", " ")
if err != nil {
return err
}
fmt.Println(string(body))
return nil
}
fmt.Print(queueText(q, time.Now()))
return nil
}
// queueText is the queue as a person reads it: a summary line, then each kind in queue order.
// Never what an ask carries as `held` — that is every artifact the mesh has built.
func queueText(q link.Queue, now time.Time) string {
var b strings.Builder
waiting, running, dead := q.Of(link.AskWaiting), q.Of(link.AskInFlight), q.Of(link.AskDead)
fmt.Fprintf(&b, "%s: %d waiting, %d in flight, %d dead\n", q.Seat, len(waiting), len(running), len(dead))
ago := func(t time.Time) string {
if t.IsZero() {
return "asked at an unknown time"
}
return "asked " + now.Sub(t).Round(time.Second).String() + " ago"
}
what := func(a link.QueuedAsk) string {
s := a.Repository
if a.Path != "" {
s += " at " + a.Path
}
if a.Ref != "" {
s += " on " + a.Ref
}
return s
}
if len(running) > 0 {
fmt.Fprintln(&b, "\nin flight:")
for _, a := range running {
where := "taken, not yet said where"
if a.On != "" {
where = fmt.Sprintf("on %s for %s", a.On, now.Sub(a.Started).Round(time.Second))
}
fmt.Fprintf(&b, " %-26s %s — %s, %s (seq %d)\n", a.ID, what(a), where, ago(a.AskedAt), a.Seq)
}
}
if len(waiting) > 0 {
fmt.Fprintln(&b, "\nwaiting:")
for _, a := range waiting {
fmt.Fprintf(&b, " %-26s %s — %s (seq %d)\n", a.ID, what(a), ago(a.AskedAt), a.Seq)
}
}
if len(dead) > 0 {
fmt.Fprintf(&b, "\ndead — handed out as often as the worker allows (%d) and never settled, still held:\n", q.MaxDeliver)
for _, a := range dead {
fmt.Fprintf(&b, " %-26s %s — %s (seq %d)\n", a.ID, what(a), ago(a.AskedAt), a.Seq)
}
}
switch {
case len(q.Asks) == 0:
fmt.Fprintln(&b, "nothing is asked of it")
default:
fmt.Fprintln(&b, "\n`cancel <id>` drops a waiting or dead ask, `clear` every waiting one, `kill <id>` ends one in flight")
}
return b.String()
}
// cancelCommand drops one waiting or dead ask.
func cancelCommand(ctx context.Context, args []string) error {
if len(args) != 1 {
return errors.New("cancel <build id>")
}
js, err := dialTheBus()
if err != nil {
return err
}
defer js.Close()
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
seat := buildSeatHeld(ctx)
q, err := link.ReadQueue(ctx, js, seat)
if err != nil {
return err
}
ask, found := q.Find(args[0])
if !found {
return fmt.Errorf("%s is not in the %s queue: it was built, cancelled, or never asked — `builds` says "+
"what came of it", args[0], seat)
}
said, err := cancelAsk(ctx, js, open, seat, ask)
if err != nil {
return err
}
fmt.Println(said)
return nil
}
// cancelAsk drops one ask from the queue and records it failed, cancelled by hand.
//
// **In this order, and each step for a reason** (novox/hq ADR 0219):
// 1. an ask in flight is refused: deleting its message ends nothing a machine is running — that is
// `kill`, on the machine running it;
// 2. its id goes into the seat's cancelled set, so a holder that fetches it from now on ends it;
// 3. for a waiting ask, the worker is read again: one handed out since the queue was read was
// taken in the moment of cancelling — the holder that took it either saw the set and ends it
// as cancelled, or did not and is building it, and which is for `queue` to say; nothing is
// deleted or recorded here, so a build that is running is not recorded as cancelled;
// 4. the message is deleted by its sequence;
// 5. the failed outcome is taken in as any failed build's is, and the plan that asked fails.
func cancelAsk(ctx context.Context, js *broker.JetStream, open *stores, seat string, ask link.QueuedAsk) (string, error) {
// **Refused only where a machine said it is building it.** An ask the worker counts pending that no
// machine said it started is either in the moment between a holder's fetch and its start — and
// that holder checks the cancelled set first — or one past its deliveries the server has not yet
// stopped counting, which happens only when some holder pulls again: with every holder paused it
// would stay, cancellable by nothing and killable by nobody.
if ask.State == link.AskInFlight && ask.On != "" {
return "", fmt.Errorf("%s is in flight on %s: cancelling drops an ask nobody is building. `kill %s` "+
"ends the build where it runs", ask.ID, ask.On, ask.ID)
}
if err := link.MarkCancelled(ctx, js, seat, ask.ID); err != nil {
return "", err
}
worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: seat, Accepts: []string{"build"}})
if ask.State == link.AskWaiting {
if taken, err := takenSince(js, worker, ask.Seq); err != nil {
return "", err
} else if taken {
return "", fmt.Errorf("%s was taken by a holder as it was cancelled. If the holder saw the cancel "+
"it ends it as %s and its outcome says so; otherwise it is building — `queue` says which, "+
"and `kill %s` ends it", ask.ID, link.CancelledByHand, ask.ID)
}
}
if err := js.Context().DeleteMsg(worker.Stream, ask.Seq); err != nil && !errors.Is(err, nats.ErrMsgNotFound) &&
!errors.Is(err, jetstream.ErrMsgNotFound) && !strings.Contains(err.Error(), "no message found") {
return "", fmt.Errorf("%s is marked cancelled and could not be deleted from the queue (seq %d): %w — a "+
"holder taking it ends it as cancelled", ask.ID, ask.Seq, err)
}
recordCancelled(ctx, open, ask)
return fmt.Sprintf("cancelled %s (%s, %s): deleted from the %s queue and recorded failed, %s — a plan "+
"that asked for it fails with that", ask.ID, ask.Repository, ask.State, seat, link.CancelledByHand), nil
}
// takenSince is whether the worker has handed out the ask at this sequence.
func takenSince(js *broker.JetStream, worker broker.Consumer, seq uint64) (bool, error) {
info, err := js.Context().ConsumerInfo(worker.Stream, worker.Name)
if errors.Is(err, nats.ErrConsumerNotFound) {
return false, nil
}
if err != nil {
return false, fmt.Errorf("cannot read the worker of the queue again: %w", err)
}
return info.Delivered.Stream >= seq, nil
}
// recordCancelled takes the failed outcome in the way the daemon takes in any build's: recorded,
// then the plan that asked for it — by the id it asked with, or by repository and path.
func recordCancelled(ctx context.Context, open *stores, ask link.QueuedAsk) {
r := ask.Request
result := link.BuildResult{ID: ask.ID, Repository: ask.Repository, Path: ask.Path, Ref: ask.Ref,
Source: r.Source, DryRun: r.DryRun, Failed: link.CancelledByHand}
_ = builds{open.inventory, open}.Built(ctx, result)
}
// clearCommand cancels every waiting ask, and with --dead every dead one too.
func clearCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("clear", flag.ContinueOnError)
dead := set.Bool("dead", false, "the dead asks too")
if _, err := parseAround(set, args); err != nil {
return err
}
js, err := dialTheBus()
if err != nil {
return err
}
defer js.Close()
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
seat := buildSeatHeld(ctx)
said, err := clearQueue(ctx, js, open, seat, *dead)
fmt.Print(said)
return err
}
// clearQueue cancels what clear names, each as `cancel` would, and never anything in flight.
func clearQueue(ctx context.Context, js *broker.JetStream, open *stores, seat string, dead bool) (string, error) {
q, err := link.ReadQueue(ctx, js, seat)
if err != nil {
return "", err
}
var b strings.Builder
cancelled, refused := 0, 0
for _, a := range q.Asks {
if a.State == link.AskInFlight || (a.State == link.AskDead && !dead) {
continue
}
said, err := cancelAsk(ctx, js, open, seat, a)
if err != nil {
refused++
fmt.Fprintf(&b, " %s: %v\n", a.ID, err)
continue
}
cancelled++
fmt.Fprintf(&b, " %s\n", said)
}
what := "waiting"
if dead {
what = "waiting and dead"
}
fmt.Fprintf(&b, "%d %s ask(s) cancelled", cancelled, what)
if refused > 0 {
fmt.Fprintf(&b, ", %d could not be", refused)
}
fmt.Fprintln(&b)
if n := len(q.Of(link.AskInFlight)); n > 0 {
fmt.Fprintf(&b, "%d in flight, left running: `kill <id>` ends one where it runs\n", n)
}
if n := len(q.Of(link.AskDead)); n > 0 && !dead {
fmt.Fprintf(&b, "%d dead, left: `clear --dead` cancels them too\n", n)
}
if refused > 0 {
return b.String(), fmt.Errorf("%d ask(s) could not be cancelled", refused)
}
return b.String(), nil
}
// rebuildCommand asks the module's current source again — or, given a build's id, that build's
// repository, path and ref — under a new id.
func rebuildCommand(ctx context.Context, args []string) error {
if len(args) != 1 {
return errors.New("rebuild <module | build id>")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
var source buildSource
var path, ref, module string
if b, found, err := inv.BuildByID(ctx, args[0]); err != nil {
return err
} else if found {
source, module = sourceOfBuild(b, entries)
path, ref = b.Path, b.Ref
} else if e, known := entryNamed(entries, args[0]); known {
// As a plan asks it: the branch it follows, never a commit a build once named (issue 215).
source = buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
path, ref, module = e.Source.Path, followedBranch(e.Source.Ref), e.Manifest.Module
} else {
return fmt.Errorf("%s is neither a module the catalogue holds nor a build the mesh recorded", args[0])
}
// **A module a plan holds unbuilt, or failed, joins that plan** rather than running beside it: the
// plan would ask it again, or stay failed on an outcome a rebuild has replaced.
if module != "" {
joined, said, err := joinAPlan(ctx, open, module)
if err != nil {
return err
}
if joined {
fmt.Println(said)
return nil
}
}
id, err := askABuild(ctx, source, path, ref)
if err != nil {
return err
}
fmt.Printf("rebuild asked as %s\n", id)
return nil
}
// entryNamed is the catalogue's entry for a module.
func entryNamed(entries []inventory.Entry, name string) (inventory.Entry, bool) {
for _, e := range entries {
if e.Manifest.Module == name {
return e, true
}
}
return inventory.Entry{}, false
}
// sourceOfBuild is where a recorded build's repository is asked from: the catalogued module's own
// source when the build is of one — on its seat, so the outcome registers it as the mesh records it
// (ADR 0111) — else the repository as it was cloned.
func sourceOfBuild(b inventory.Build, entries []inventory.Entry) (buildSource, string) {
for _, e := range entries {
if (b.Module != "" && e.Manifest.Module == b.Module) ||
(b.Module == "" && repositoryMatches(e.Source.Repository, b.Repository) && e.Source.Path == b.Path) {
return buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}, e.Manifest.Module
}
}
return buildSource{Repository: b.Repository}, b.Module
}
// replayCommand asks a recorded build's repository and path again at the commit it built.
func replayCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("replay", flag.ContinueOnError)
register := set.Bool("register", false, "register what it builds, as any build is")
older := set.Bool("older", false, "with --register: even though a newer build of the module is registered")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) != 1 {
return errors.New("replay <build id> [--register [--older]]")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
b, found, err := inv.BuildByID(ctx, positionals[0])
if err != nil {
return err
}
if !found {
return fmt.Errorf("no build %s is recorded", positionals[0])
}
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
source, module := sourceOfBuild(b, entries)
var history []inventory.Build
if module != "" {
if history, err = inv.Builds(ctx, module, 100); err != nil {
return err
}
}
if err := replayRefusal(b, module, history, *register, *older); err != nil {
return err
}
id, err := buildOneAsked(ctx, source, b.Path, b.Commit, 0, !*register)
if err != nil {
return err
}
fmt.Println(replaySaid(b, module, id, *register))
return nil
}
// replayRefusal is why a replay is not asked, or nothing (novox/hq ADR 0219, issue 207).
//
// A replay re-asks a recorded build at the commit it built, under a new id — and an id is when it
// was asked, which is what orders builds of one module (issue 219). So a registered replay of an
// older commit is newer than everything since, and would roll that older commit out as the module's
// current version: what issue 207 recorded happening by accident. It is therefore a dry run unless
// --register says otherwise, and --register is refused when a newer build of a different commit is
// registered, unless --older says that is the point.
func replayRefusal(b inventory.Build, module string, history []inventory.Build, register, older bool) error {
if b.Commit == "" {
return fmt.Errorf("%s recorded no commit — it failed before it knew what it was building, so there is "+
"nothing to replay; `rebuild %s` asks its repository, path and ref again", b.ID, b.ID)
}
if older && !register {
return errors.New("--older only says what --register may do; a dry run registers nothing")
}
if !register || older {
return nil
}
for _, h := range history {
if h.ID == b.ID || !h.Worked() || h.Commit == b.Commit {
continue
}
if h.AskedOrAt().After(b.AskedOrAt()) {
return fmt.Errorf("a newer build of %s is registered — %s, from %s — and registering a replay of %s "+
"would roll that older commit out as %s's current version (novox/hq issue 207). `replay %s` "+
"without --register looks at it; --register --older registers it anyway",
module, h.ID, short(h.Commit), short(b.Commit), module, b.ID)
}
}
return nil
}
// replaySaid is what a replay prints once asked: what it builds, and what becomes of the outcome.
func replaySaid(b inventory.Build, module, id string, register bool) string {
what := b.Repository
if module != "" {
what = module
}
said := fmt.Sprintf("replaying %s (%s) at %s as %s", b.ID, what, short(b.Commit), id)
if !register {
return said + "\n a dry run: the outcome is looked at and not taken in — nothing is recorded or " +
"registered, and nothing is sent. `builds --log " + id + "` follows it; `--register` registers it"
}
return said + "\n registered when it is built, as the module's current version — newer than every build " +
"asked before now — and rolled out as its policy says. `builds --log " + id + "` follows it"
}
// killCommand finds the machine running a build and asks its holder to end it.
func killCommand(ctx context.Context, args []string) error {
if len(args) != 1 {
return errors.New("kill <build id>")
}
id := args[0]
js, err := dialTheBus()
if err != nil {
return err
}
defer js.Close()
seat := buildSeatHeld(ctx)
since := time.Now().Add(-7 * 24 * time.Hour)
if at, ok := link.BuildAskedAt(id); ok {
since = at.Add(-time.Minute)
}
heard, err := link.ReadBuildEvents(ctx, js, seat, since)
if err != nil {
return err
}
started, ok := heard.Started[id]
if !ok {
return fmt.Errorf("no machine has said it started %s: if it waits in the queue, `cancel %s` drops it", id, id)
}
if heard.Outcomes[id] {
return fmt.Errorf("%s has already ended on %s — `builds` says how", id, started.On)
}
answer, err := link.AskSeatTool(ctx, js.Conn(), seat, "kill", started.On, map[string]string{"id": id}, 45*time.Second)
if err != nil {
return err
}
return printHolderAnswer(started.On, answer)
}
// pauseCommand asks one machine's holder, or every holder's, to pause or resume.
func pauseCommand(ctx context.Context, verb string, args []string) error {
if len(args) > 1 {
return fmt.Errorf("%s [node]", verb)
}
js, err := dialTheBus()
if err != nil {
return err
}
defer js.Close()
seat := buildSeatHeld(ctx)
nodes := args
if len(nodes) == 0 {
if nodes, err = buildSeatHolders(ctx, seat); err != nil {
return err
}
if len(nodes) == 0 {
return fmt.Errorf("nothing holds %s, so there is nothing to %s", seat, verb)
}
}
failed := 0
for _, node := range nodes {
answer, err := link.AskSeatTool(ctx, js.Conn(), seat, verb, node, map[string]string{}, holderAsks)
if err != nil {
fmt.Printf("%s: %v\n", node, err)
failed++
continue
}
if err := printHolderAnswer(node, answer); err != nil {
failed++
}
}
if failed > 0 {
return fmt.Errorf("%d of %d machine(s) did not %s", failed, len(nodes), verb)
}
return nil
}
// printHolderAnswer prints what a holder said, or its refusal as the command's failure.
func printHolderAnswer(node string, answer link.Answer) error {
if answer.Error != "" {
fmt.Printf("%s: %s\n", node, answer.Error)
return errors.New(answer.Error)
}
var said struct {
Said string `json:"said"`
}
if json.Unmarshal(answer.Result, &said) == nil && said.Said != "" {
fmt.Printf("%s: %s\n", node, said.Said)
return nil
}
fmt.Printf("%s: %s\n", node, string(answer.Result))
return nil
}
// buildSeatHolders is every machine an assigned module holding the build seat runs on.
func buildSeatHolders(ctx context.Context, seat string) ([]string, error) {
open, err := openStores(ctx)
if err != nil {
return nil, err
}
defer open.Close()
entries, err := open.inventory.Catalogued(ctx)
if err != nil {
return nil, err
}
return holdersAmong(entries, seat), nil
}
// holdersAmong is the machines of every catalogued module claiming the seat, sorted, once each.
func holdersAmong(entries []inventory.Entry, seat string) []string {
seen := map[string]bool{}
var out []string
for _, e := range entries {
if !e.Manifest.ClaimsSeat(seat) {
continue
}
for _, n := range e.On {
if !seen[n] {
seen[n] = true
out = append(out, n)
}
}
}
sort.Strings(out)
return out
}