From 106507b1d3e2c1f6cc16319df93bccb97a1734a3 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 19:17:56 +0200 Subject: [PATCH] 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. --- cmd/mesh-controller/build.go | 40 +- cmd/mesh-controller/main.go | 31 +- cmd/mesh-controller/plan_retry.go | 308 ++++++++++++ cmd/mesh-controller/queue.go | 583 +++++++++++++++++++++++ cmd/mesh-controller/queue_test.go | 582 ++++++++++++++++++++++ cmd/mesh-controller/readable.go | 2 +- cmd/mesh-controller/release_plan.go | 166 +++++-- cmd/mesh-controller/release_plan_test.go | 8 +- cmd/mesh-controller/seatverbs.go | 38 ++ cmd/mesh-controller/status.go | 6 +- internal/catalogue/verbs.go | 31 ++ internal/inventory/builds.go | 30 ++ internal/inventory/plans.go | 5 + 13 files changed, 1772 insertions(+), 58 deletions(-) create mode 100644 cmd/mesh-controller/plan_retry.go create mode 100644 cmd/mesh-controller/queue.go create mode 100644 cmd/mesh-controller/queue_test.go diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 6172e39..915fcfc 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -392,22 +392,34 @@ func buildBehind(ctx context.Context, wait time.Duration) error { // // Separated from the command so `--behind` can walk a list without a second path to the same act. func buildOne(ctx context.Context, source buildSource, path, ref string, wait time.Duration) error { + _, err := buildOneAsked(ctx, source, path, ref, wait, false) + return err +} + +// buildOneAsked is buildOne answering the id it asked with — what a plan keeps to match the outcome +// by (novox/hq ADR 0219) — and, for an ask not waited for, optionally a dry run: built and looked +// at, never taken in (issue 240), which is what `replay` asks unless told to register. +func buildOneAsked(ctx context.Context, source buildSource, path, ref string, wait time.Duration, + dryRun bool) (string, error) { + if dryRun && wait != 0 { + return "", errors.New("a dry run waited for is `build --dry-run`") + } // Before anything is asked of a builder: a source on a seat nobody holds is refused here, with // the reason, rather than sent to a machine to fail at `git clone`. repository, err := cloneFrom(ctx, source) if err != nil { - return err + return "", err } ident, err := openIdentity(ctx) if err != nil { - return err + return "", err } defer ident.Close() server, err := connectLink(ctx, nil, nil, nil) if err != nil { - return err + return "", err } defer server.Close() @@ -420,6 +432,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti Ref: ref, Held: heldBy(ctx), Seats: seatBases(ctx), + DryRun: dryRun, } fmt.Printf("asked for %s", source) if source.Seat != "" { @@ -438,7 +451,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti seat := buildSeatHeld(ctx) ask, err := askOverOn(seat) if err != nil { - return err + return "", err } defer ask.Close() fmt.Printf(" of %s\n", seat) @@ -449,26 +462,31 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti // is still here. A tool call cannot hold a connection for the minutes a build takes; it // follows the build by its id instead. if err := ask.Ask(ctx, request); err != nil { - return err + return "", err + } + if dryRun { + fmt.Printf("asked as a dry run, not waited for: `builds --log %s` follows it as it runs; "+ + "its outcome is not taken in\n", request.ID) + return request.ID, nil } fmt.Printf("asked, not waited for: `builds --log %s` follows it as it runs, and `builds` "+ "shows what came of it; the module is registered when the outcome comes\n", request.ID) - return nil + return request.ID, nil } result, err := ask.Submit(ctx, request, wait) if err != nil { - return err + return request.ID, err } open, err := openStores(ctx) if err != nil { - return err + return request.ID, err } defer open.Close() manifest, kept, err := takeIn(ctx, open.inventory, result) if err != nil { - return err + return request.ID, err } // Said as recorded: what each artifact is, not where this builder happened to push it. for _, made := range kept.Made { @@ -478,7 +496,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti manifest.Module, manifest.Version, result.On, short(result.Commit)) saysWhenThePolicyActs(ctx, open.inventory, manifest.Module) fmt.Printf(" run `assign %s` to put it somewhere\n", manifest.Module) - return nil + return request.ID, nil } // saysWhenThePolicyActs tells whoever built a module that its upgrade policy will send the @@ -628,6 +646,8 @@ type answers struct { reported []inventory.Reported // plans is what the last merges produced and where each stands (novox/hq ADR 0162). plans []inventory.Plan + // paused is whether the build seat takes work, which a plan waiting on it says (ADR 0219). + paused pauseView // refused is why a machine cannot be worked out at all, by name. A different thing from every // other answer here: those are about a machine that was told something, and this is about one // that cannot be told anything — it never reaches waiting, because nothing was computed for it diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index b15f40d..38fe897 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -73,6 +73,21 @@ func run() error { return askCommand(ctx, args[1:]) case "builds": return buildsCommand(ctx, args[1:]) + // The build queue, controlled by hand (novox/hq ADR 0219). + case "queue": + return queueCommand(ctx, args[1:]) + case "cancel": + return cancelCommand(ctx, args[1:]) + case "clear": + return clearCommand(ctx, args[1:]) + case "rebuild": + return rebuildCommand(ctx, args[1:]) + case "replay": + return replayCommand(ctx, args[1:]) + case "kill": + return killCommand(ctx, args[1:]) + case "pause", "resume": + return pauseCommand(ctx, args[0], args[1:]) case "collection": return collectionCommand(ctx, args[1:]) case "plans": @@ -202,6 +217,14 @@ func usage() { build --behind build every module the mesh holds older than its source build --on rebuild every module that stands on this module's artifacts, bases first builds [] what has been built lately, and what came of it + queue [--json] every ask in the build queue: waiting, in flight (where, how long), dead + cancel drop a waiting or dead ask; recorded failed, cancelled by hand + clear [--dead] cancel every waiting ask (and the dead ones); never one in flight + rebuild ask the module's source again, or that build's, under a new id + replay [--register [--older]] that build's commit again; a dry run unless --register + kill end a build where it runs; recorded failed, killed by hand + pause [] / resume [] the build seat's holder there, or every holder, takes nothing new / again + plans retry ask a failed plan's failed builds again, and carry the plan on collection [--json] kept archives held/unheld by a manifest, and what the sweep may let go builder issue a broker account for a build machine, scoped to build work, delivered as the builder module's broker secret (module add it first) @@ -271,7 +294,7 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error { case err != nil && result.Failed != "": fmt.Printf("%s: %v\n", result.ID, err) if result.Module != "" { - planBuilt(ctx, b.open, result.Module, result.Commit, result.Failed, asked) + planBuilt(ctx, b.open, result.Module, result.Commit, result.Failed, asked, result.ID) } else { planFailedBuild(ctx, b.open, result) } @@ -280,18 +303,18 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error { // Not a failure: the module is already at what a later request built. A plan that asked // before that later request is answered by it; one that asked after it ignores this. fmt.Printf("%s: %v\n", result.ID, err) - planBuilt(ctx, b.open, manifest.Module, result.Commit, "", asked) + planBuilt(ctx, b.open, manifest.Module, result.Commit, "", asked, result.ID) return nil case err != nil: fmt.Printf("%s: heard and recorded, and not registered: %v\n", result.ID, err) if manifest.Module != "" { - planBuilt(ctx, b.open, manifest.Module, result.Commit, err.Error(), asked) + planBuilt(ctx, b.open, manifest.Module, result.Commit, err.Error(), asked, result.ID) } return nil } fmt.Printf("%s: %s %s registered, built on %s from %s\n", result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit)) saysWhenThePolicyActs(ctx, b.inv, manifest.Module) - planBuilt(ctx, b.open, manifest.Module, result.Commit, "", asked) + planBuilt(ctx, b.open, manifest.Module, result.Commit, "", asked, result.ID) return nil } diff --git a/cmd/mesh-controller/plan_retry.go b/cmd/mesh-controller/plan_retry.go new file mode 100644 index 0000000..405f465 --- /dev/null +++ b/cmd/mesh-controller/plan_retry.go @@ -0,0 +1,308 @@ +package main + +import ( + "context" + "fmt" + "os" + "sort" + "strings" + "time" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// A plan follows what is done to the build queue (novox/hq ADR 0219). +// +// Two things a person does to builds by hand would otherwise leave a plan saying something untrue: +// +// - **Pausing the build seat.** A plan whose builds wait in the queue of a seat whose every holder +// is paused is not late — nothing will take its asks until somebody resumes — and saying LATE +// sends a reader looking for a fault that is a decision. It says what it waits on instead. +// - **A build that failed, been cancelled or killed, and is asked again.** `plans retry` re-asks a +// failed plan's failed modules and the plan goes on from that tier as if they had built the +// first time; `rebuild` of a module a plan holds unbuilt joins that plan rather than running +// beside it — beside it, the plan would either ask it again or stay failed on an outcome the +// rebuild has already replaced. + +// pauseView is whether the build seat takes work: the holders that said they are paused, and +// whether that is every holder. +type pauseView struct { + Nodes []string + All bool +} + +// pausedWaiting is what a plan waits on when the build seat is paused under it, or false: a plan +// building, whose current tier has an ask outstanding, while every holder of the seat is paused. +func pausedWaiting(p inventory.Plan, pause pauseView, now time.Time) (string, bool) { + if p.State != inventory.PlanBuilding || !pause.All || len(pause.Nodes) == 0 || p.Tier >= len(p.Tiers) { + return "", false + } + var asked *time.Time + for _, m := range p.Tiers[p.Tier] { + if s := p.Modules[m]; s != nil && s.State == "asked" && s.AskedAt != nil { + if asked == nil || s.AskedAt.Before(*asked) { + asked = s.AskedAt + } + } + } + if asked == nil { + return "", false + } + return fmt.Sprintf("waiting: the build seat is paused on %s (asked %s ago)", + strings.Join(pause.Nodes, ", "), now.Sub(*asked).Round(time.Second)), true +} + +// awaitsABuild is whether any open plan has an ask outstanding in its current tier — the only case +// where the seat being paused changes what a plan says. +func awaitsABuild(plans []inventory.Plan) bool { + for _, p := range plans { + if p.State != inventory.PlanBuilding || p.Tier >= len(p.Tiers) { + continue + } + for _, m := range p.Tiers[p.Tier] { + if s := p.Modules[m]; s != nil && s.State == "asked" { + return true + } + } + } + return false +} + +// buildSeatPause reads whether the build seat's holders take work, from what each last said on the +// bus (link.HolderState) — read only when a plan waits on a build, so a mesh with nothing building +// does not dial the bus to say so. Anything unreadable is said and read as not paused: a plan then +// reads as late, which is what it said before this existed. +func buildSeatPause(ctx context.Context, inv *inventory.Inventory, plans []inventory.Plan) pauseView { + if !awaitsABuild(plans) { + return pauseView{} + } + entries, err := inv.Catalogued(ctx) + if err != nil { + return pauseView{} + } + seat := buildSeatAmong(entries) + holders := holdersAmong(entries, seat) + if len(holders) == 0 { + return pauseView{} + } + address, err := broker.BusAddress() + if err != nil { + return pauseView{} + } + js, err := broker.Dial(address) + if err != nil { + fmt.Fprintf(os.Stderr, "could not reach the bus to read whether the build seat is paused: %v\n", err) + return pauseView{} + } + defer js.Close() + said, err := link.PausedSaid(js, seat, holders) + if err != nil { + fmt.Fprintf(os.Stderr, "could not read whether the build seat is paused: %v\n", err) + return pauseView{} + } + return pauseOf(holders, said) +} + +// pauseOf is the view from what each holder said. +func pauseOf(holders []string, said map[string]link.HolderState) pauseView { + var v pauseView + for _, n := range holders { + if said[n].Paused { + v.Nodes = append(v.Nodes, n) + } + } + sort.Strings(v.Nodes) + v.All = len(holders) > 0 && len(v.Nodes) == len(holders) + return v +} + +// failedIn is the modules of a plan's current tier that failed to build, sorted. +func failedIn(p inventory.Plan) []string { + if p.Tier >= len(p.Tiers) { + return nil + } + var out []string + for _, m := range p.Tiers[p.Tier] { + if s := p.Modules[m]; s != nil && s.State == "failed" { + out = append(out, m) + } + } + sort.Strings(out) + return out +} + +// newerOpenPlan is an open plan of the same repository and branch made after this one: the plan +// that holds what this one held now (issue 254, ADR 0218). +func newerOpenPlan(p inventory.Plan, plans []inventory.Plan) (inventory.Plan, bool) { + for _, q := range plans { + if q.ID == p.ID || !q.Open() || !strings.EqualFold(q.Repository, p.Repository) || + (p.Branch != "" && q.Branch != "" && p.Branch != q.Branch) || !q.Created.After(p.Created) { + continue + } + return q, true + } + return inventory.Plan{}, false +} + +// retryRefusal is why a plan cannot be retried, or nothing. +func retryRefusal(p inventory.Plan, plans []inventory.Plan) error { + switch { + case p.State == inventory.PlanDone: + return fmt.Errorf("%s is done; there is nothing to retry", p.ID) + case p.State == inventory.PlanSuperseded: + return fmt.Errorf("%s was %s — what it had not built is in that plan", p.ID, p.Note) + case p.Open(): + return fmt.Errorf("%s is still %s; nothing in it failed to retry — `rebuild ` asks one module again", p.ID, p.State) + } + if q, found := newerOpenPlan(p, plans); found { + return fmt.Errorf("%s supersedes it: a newer merge of %s (%s at %s) is open, and retrying %s would build "+ + "what that one replaced", q.ID, q.Repository, q.ID, short(q.Commit), p.ID) + } + if len(failedIn(p)) == 0 { + return fmt.Errorf("nothing in tier %d of %s failed to build — it stopped at: %s. Retry asks failed builds "+ + "again", p.Tier, p.ID, p.Note) + } + return nil +} + +// resumed sets a failed plan building again once nothing in its tier is failed. +func resumed(p *inventory.Plan, why string) { + if p.State == inventory.PlanFailed && len(failedIn(*p)) == 0 { + p.State = inventory.PlanBuilding + p.Note = why + } +} + +// retryPlan asks a failed plan's failed modules again, under new ids, and sets it building again at +// that tier: what follows is the plan going on as if they had built the first time. +func retryPlan(ctx context.Context, open *stores, id string) (string, error) { + inv := open.inventory + release, err := inv.HoldPlans(ctx, true) + if err != nil { + return "", err + } + defer release() + p, err := inv.PlanByID(ctx, id) + if err != nil { + return "", err + } + plans, err := inv.OpenPlans(ctx) + if err != nil { + return "", err + } + if err := retryRefusal(p, plans); err != nil { + return "", err + } + entries, err := inv.Catalogued(ctx) + if err != nil { + return "", err + } + byName := map[string]inventory.Entry{} + for _, e := range entries { + byName[e.Manifest.Module] = e + } + failed := failedIn(p) + var asked []string + for _, m := range failed { + askModule(ctx, &p, m, byName) + if s := p.Modules[m]; s.State == "asked" { + asked = append(asked, m+" as "+s.Build) + } + } + resumed(&p, fmt.Sprintf("tier %d retried by hand: %s asked again", p.Tier, strings.Join(failed, ", "))) + if err := inv.SavePlan(ctx, p); err != nil { + return "", err + } + if p.State != inventory.PlanBuilding { + return "", fmt.Errorf("%s could not be resumed: %s", p.ID, p.Note) + } + return fmt.Sprintf("%s retried at tier %d of %d: asked %s; the plan goes on from there as any plan does", + p.ID, p.Tier, len(p.Tiers), strings.Join(asked, ", ")), nil +} + +// joinAPlan asks a module again for the plan that holds it unbuilt or failed, if one does: open plans +// first, then the most recent failed one that nothing newer supersedes. Says whether it joined one. +func joinAPlan(ctx context.Context, open *stores, module string) (bool, string, error) { + inv := open.inventory + release, err := inv.HoldPlans(ctx, true) + if err != nil { + return false, "", err + } + defer release() + openPlans, err := inv.OpenPlans(ctx) + if err != nil { + return false, "", err + } + recent, err := inv.RecentPlans(ctx, 20) + if err != nil { + return false, "", err + } + p, found := planHolding(module, openPlans, recent) + if !found { + return false, "", nil + } + entries, err := inv.Catalogued(ctx) + if err != nil { + return false, "", err + } + byName := map[string]inventory.Entry{} + for _, e := range entries { + byName[e.Manifest.Module] = e + } + was := p.State + askModule(ctx, &p, module, byName) + s := p.Modules[module] + resumed(&p, fmt.Sprintf("tier %d: %s rebuilt by hand", p.Tier, module)) + if err := inv.SavePlan(ctx, p); err != nil { + return false, "", err + } + if s.State != "asked" { + return true, "", fmt.Errorf("%s could not be asked for %s: %s", module, p.ID, s.Why) + } + said := fmt.Sprintf("rebuild asked as %s, joining %s at tier %d (%s at %s): the plan takes this build as %s's outcome", + s.Build, p.ID, p.Tier, p.Repository, short(p.Commit), module) + switch { + case was == inventory.PlanFailed && p.State == inventory.PlanBuilding: + said += "; the plan had failed and builds again from this tier" + case was == inventory.PlanFailed: + said += "; the plan stays failed while " + strings.Join(failedIn(p), ", ") + " failed too — `plans retry " + p.ID + "` asks them" + } + return true, said, nil +} + +// planHolding is the plan a rebuild of a module joins: an open plan whose current tier holds it not +// yet built, else the newest failed plan whose current tier does, and that no open plan of its +// repository supersedes. +func planHolding(module string, openPlans, recent []inventory.Plan) (inventory.Plan, bool) { + holds := func(p inventory.Plan) bool { + if p.Tier >= len(p.Tiers) { + return false + } + for _, m := range p.Tiers[p.Tier] { + if m == module { + s := p.Modules[m] + return s == nil || s.State != "built" + } + } + return false + } + for _, p := range openPlans { + if holds(p) { + return p, true + } + } + sorted := append([]inventory.Plan(nil), recent...) + sort.SliceStable(sorted, func(i, j int) bool { return sorted[i].Created.After(sorted[j].Created) }) + for _, p := range sorted { + if p.State != inventory.PlanFailed || !holds(p) { + continue + } + if _, superseded := newerOpenPlan(p, openPlans); superseded { + continue + } + return p, true + } + return inventory.Plan{}, false +} diff --git a/cmd/mesh-controller/queue.go b/cmd/mesh-controller/queue.go new file mode 100644 index 0000000..018ff2e --- /dev/null +++ b/cmd/mesh-controller/queue.go @@ -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 ` drops a waiting or dead ask, `clear` every waiting one, `kill ` 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 ") + } + 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 ` 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 ") + } + 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 [--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 ") + } + 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 +} diff --git a/cmd/mesh-controller/queue_test.go b/cmd/mesh-controller/queue_test.go new file mode 100644 index 0000000..2401d0a --- /dev/null +++ b/cmd/mesh-controller/queue_test.go @@ -0,0 +1,582 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "os" + "reflect" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/link" +) + +// The build queue, controlled by hand (novox/hq ADR 0219), and the plans that follow it. +// +// docker run -d --rm --name bq-nats -p 14294:4222 nats:2.10-alpine -js +// make postgres PG_PORT=55566 PG_CONTAINER=bq-pg +// MESH_TEST_NATS=nats://127.0.0.1:14294 \ +// MESH_TEST_POSTGRES='postgres://postgres:check@127.0.0.1:55566/postgres?sslmode=disable' \ +// go test ./cmd/mesh-controller/ -run 'Queue|Cancel|Clear|Retry|Rebuild|Replay|Paused' + +// asksRecorded makes every ask a plan makes return the next id in a row, and says which were asked. +func asksRecorded(t *testing.T) *[]string { + t.Helper() + var asked []string + was := askABuild + askABuild = func(_ context.Context, source buildSource, path, ref string) (string, error) { + id := link.NewBuildID(time.Now().Add(time.Duration(len(asked)) * time.Millisecond)) + asked = append(asked, source.Repository+"#"+id) + return id, nil + } + t.Cleanup(func() { askABuild = was }) + return &asked +} + +// twoTiers registers two modules, b standing on a, each with a source a plan asks. +func twoTiers(t *testing.T, open *stores) { + t.Helper() + for _, name := range []string{"a", "b"} { + if err := open.inventory.RegisterModule(t.Context(), catalogue.Manifest{Module: name, Version: "1"}, + inventory.Source{Repository: "novox/" + name, Seat: "git", Ref: "main", BuiltFrom: "c0ffee", Head: "c0ffee"}); err != nil { + t.Fatal(err) + } + } +} + +// A failed plan is retried: its failed module asked again under a new id, the plan building at that +// tier, and when that build comes in the plan goes on and asks its next tier. +func TestAFailedPlanIsRetriedAndGoesOnThroughItsLaterTiers(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + asked := asksRecorded(t) + twoTiers(t, open) + before := time.Now().UTC().Add(-time.Hour) + failed := inventory.Plan{ID: "plan-retry", Repository: "novox/a", Branch: "main", Commit: "c0ffee", + Created: before, State: inventory.PlanFailed, Tier: 0, Tiers: [][]string{{"a"}, {"b"}}, + Note: "a failed to build in tier 0", + Modules: map[string]*inventory.PlanModule{"a": {State: "failed", AskedAt: &before, Build: "build-1", + Why: link.KilledByHand}}} + if err := open.inventory.SavePlan(ctx, failed); err != nil { + t.Fatal(err) + } + + said, err := retryPlan(ctx, open, failed.ID) + if err != nil { + t.Fatal(err) + } + p, err := open.inventory.PlanByID(ctx, failed.ID) + if err != nil { + t.Fatal(err) + } + a := p.Modules["a"] + if p.State != inventory.PlanBuilding || a.State != "asked" || a.Build == "" || a.Build == "build-1" || a.Why != "" { + t.Fatalf("after retry the plan is %s and a is %+v", p.State, a) + } + if !strings.Contains(said, a.Build) { + t.Errorf("retry does not say the new id: %q", said) + } + if len(*asked) != 1 { + t.Fatalf("asked %v", *asked) + } + + // The killed build's own late outcome is not this ask's; the new one's is, and tier 1 follows. + planBuilt(ctx, open, "a", "c0ffee", link.KilledByHand, before, "build-1") + if p, _ = open.inventory.PlanByID(ctx, failed.ID); p.State != inventory.PlanBuilding { + t.Fatalf("the old ask's outcome failed the retried plan: %s %q", p.State, p.Note) + } + asking, _ := link.BuildAskedAt(a.Build) + planBuilt(ctx, open, "a", "c0ffee", "", asking, a.Build) + p, err = open.inventory.PlanByID(ctx, failed.ID) + if err != nil { + t.Fatal(err) + } + if p.Tier != 1 || p.Modules["b"] == nil || p.Modules["b"].State != "asked" || p.Modules["b"].Build == "" { + t.Fatalf("the retried plan did not go on to tier 1: tier %d, %s, b %+v", p.Tier, p.State, p.Modules["b"]) + } + if len(*asked) != 2 { + t.Fatalf("asked %v", *asked) + } +} + +// What retry refuses, and says why. +func TestRetryRefusesWhatItCannotResume(t *testing.T) { + at := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + plan := func(id, state string, created time.Time) inventory.Plan { + return inventory.Plan{ID: id, Repository: "novox/mesh-catalog", Branch: "main", Commit: id + "c0ffee", + Created: created, State: state, Tiers: [][]string{{"a"}}, + Modules: map[string]*inventory.PlanModule{"a": {State: "failed"}}} + } + failed := plan("plan-1", inventory.PlanFailed, at) + for _, c := range []struct { + p inventory.Plan + others []inventory.Plan + says string + }{ + {plan("plan-d", inventory.PlanDone, at), nil, "is done"}, + {plan("plan-s", inventory.PlanSuperseded, at), nil, "was"}, + {plan("plan-o", inventory.PlanBuilding, at), nil, "still building"}, + {failed, []inventory.Plan{plan("plan-2", inventory.PlanRolling, at.Add(time.Hour))}, "plan-2 supersedes it"}, + } { + err := retryRefusal(c.p, c.others) + if err == nil || !strings.Contains(err.Error(), c.says) { + t.Errorf("%s: %v, wanted it to say %q", c.p.ID, err, c.says) + } + } + stopped := failed + stopped.Modules = map[string]*inventory.PlanModule{"a": {State: "built"}} + stopped.Note = "a stopped at its first machine" + if err := retryRefusal(stopped, nil); err == nil || !strings.Contains(err.Error(), "nothing in tier 0") { + t.Errorf("a plan with nothing failed to build was retried: %v", err) + } + // Another branch's newer plan, an older one, and a failed one do not supersede it. + other := plan("plan-3", inventory.PlanBuilding, at.Add(time.Hour)) + other.Branch = "release" + if err := retryRefusal(failed, []inventory.Plan{other, plan("plan-0", inventory.PlanBuilding, at.Add(-time.Hour)), + plan("plan-4", inventory.PlanFailed, at.Add(time.Hour))}); err != nil { + t.Errorf("refused for a plan that does not supersede it: %v", err) + } +} + +// `rebuild` of a module a failed plan holds joins that plan; one held by nothing runs alone. +func TestARebuildJoinsThePlanHoldingTheModule(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + asked := asksRecorded(t) + twoTiers(t, open) + before := time.Now().UTC().Add(-time.Hour) + failed := inventory.Plan{ID: "plan-join", Repository: "novox/a", Branch: "main", Commit: "c0ffee", + Created: before, State: inventory.PlanFailed, Tiers: [][]string{{"a"}, {"b"}}, + Note: "a failed to build in tier 0", + Modules: map[string]*inventory.PlanModule{"a": {State: "failed", AskedAt: &before, Build: "build-1", Why: link.CancelledByHand}}} + if err := open.inventory.SavePlan(ctx, failed); err != nil { + t.Fatal(err) + } + if err := rebuildCommand(ctx, []string{"a"}); err != nil { + t.Fatal(err) + } + p, err := open.inventory.PlanByID(ctx, failed.ID) + if err != nil { + t.Fatal(err) + } + if p.State != inventory.PlanBuilding || p.Modules["a"].State != "asked" || p.Modules["a"].Build == "build-1" { + t.Fatalf("the rebuild did not join the failed plan: %s %+v", p.State, p.Modules["a"]) + } + if len(*asked) != 1 { + t.Fatalf("asked %v", *asked) + } + // b is in a tier not yet reached: nothing holds it in its current tier, so it is asked alone. + if err := rebuildCommand(ctx, []string{"b"}); err != nil { + t.Fatal(err) + } + if p, _ = open.inventory.PlanByID(ctx, failed.ID); p.Modules["b"] != nil { + t.Fatalf("a module the plan has not reached joined it: %+v", p.Modules["b"]) + } + if len(*asked) != 2 { + t.Fatalf("asked %v", *asked) + } +} + +// A failed plan an open plan of its repository supersedes is not joined: the open one holds the module. +func TestARebuildJoinsTheOpenPlanBeforeASupersededFailedOne(t *testing.T) { + at := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + failed := inventory.Plan{ID: "plan-old", Repository: "novox/a", Branch: "main", Created: at, State: inventory.PlanFailed, + Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{"a": {State: "failed"}}} + newer := inventory.Plan{ID: "plan-new", Repository: "novox/a", Branch: "main", Created: at.Add(time.Hour), + State: inventory.PlanBuilding, Tiers: [][]string{{"x"}, {"a"}}, Modules: map[string]*inventory.PlanModule{}} + if _, found := planHolding("a", []inventory.Plan{newer}, []inventory.Plan{newer, failed}); found { + t.Fatal("joined a failed plan a newer open one supersedes") + } + if p, found := planHolding("a", nil, []inventory.Plan{failed}); !found || p.ID != "plan-old" { + t.Fatalf("did not join the failed plan holding it: %v %s", found, p.ID) + } + built := failed + built.Modules = map[string]*inventory.PlanModule{"a": {State: "built"}} + if _, found := planHolding("a", nil, []inventory.Plan{built}); found { + t.Fatal("joined a plan that has the module built") + } +} + +// replay is a dry run unless registered, and registering an older commit than one registered since +// is refused unless --older says it is meant (novox/hq issue 207). +func TestReplayRefusesToRollAnOlderCommitOutUnlessToldTo(t *testing.T) { + asked := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + old := inventory.Build{ID: "build-old", Module: "a", Commit: "0ldc0mm1t", Asked: asked} + newer := inventory.Build{ID: "build-new", Module: "a", Commit: "n3wc0mm1t", Asked: asked.Add(time.Hour)} + history := []inventory.Build{newer, old} + + if err := replayRefusal(old, "a", history, false, false); err != nil { + t.Errorf("a dry run was refused: %v", err) + } + err := replayRefusal(old, "a", history, true, false) + if err == nil || !strings.Contains(err.Error(), "build-new") || !strings.Contains(err.Error(), "--older") { + t.Errorf("registering an older commit than the one registered was not refused: %v", err) + } + if err := replayRefusal(old, "a", history, true, true); err != nil { + t.Errorf("--register --older was refused: %v", err) + } + if err := replayRefusal(newer, "a", history, true, false); err != nil { + t.Errorf("registering the newest build again was refused: %v", err) + } + // A newer build of the same commit, or one that failed, is not a newer version to roll back from. + same := inventory.Build{ID: "build-same", Module: "a", Commit: old.Commit, Asked: asked.Add(2 * time.Hour)} + broken := inventory.Build{ID: "build-broken", Module: "a", Failed: "no", Asked: asked.Add(3 * time.Hour)} + if err := replayRefusal(old, "a", []inventory.Build{broken, same, old}, true, false); err != nil { + t.Errorf("refused for a newer build of the same commit or a failed one: %v", err) + } + if err := replayRefusal(inventory.Build{ID: "build-x", Failed: "clone"}, "", nil, false, false); err == nil { + t.Error("a build that recorded no commit was replayed") + } + if err := replayRefusal(old, "a", history, false, true); err == nil { + t.Error("--older without --register was taken") + } + if said := replaySaid(old, "a", "build-1", false); !strings.Contains(said, "dry run") || !strings.Contains(said, "--register") { + t.Errorf("a dry replay does not say what it is: %q", said) + } + if said := replaySaid(old, "a", "build-1", true); !strings.Contains(said, "registered") || !strings.Contains(said, "rolled out") { + t.Errorf("a registered replay does not say what it does: %q", said) + } +} + +// A plan waiting on builds of a seat whose every holder is paused says so and is not late; a seat +// paused on some holders only is not a reason the plan is waiting. +func TestAPlanWaitingOnAPausedSeatSaysSoAndIsNotLate(t *testing.T) { + now := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + asked := now.Add(-12 * time.Minute) + p := inventory.Plan{ID: "plan-p", Repository: "novox/a", Commit: "c0ffee", State: inventory.PlanBuilding, + Updated: now.Add(-2 * time.Hour), Tiers: [][]string{{"a"}}, + Modules: map[string]*inventory.PlanModule{"a": {State: "asked", AskedAt: &asked}}} + + all := pauseOf([]string{"g14", "ace"}, map[string]link.HolderState{"ace": {Paused: true}, "g14": {Paused: true}}) + line := planLineWith(p, now, all) + if !strings.Contains(line, "waiting: the build seat is paused on ace, g14 (asked 12m0s ago)") || strings.Contains(line, "LATE") { + t.Errorf("a plan on a paused seat reads %q", line) + } + if st := planStatuses([]inventory.Plan{p}, now, all)[0]; st.Late || !strings.Contains(st.Waiting, "paused on ace, g14") { + t.Errorf("status --json says %+v", st) + } + if _, late := openPlans([]inventory.Plan{p}, all); late != 0 { + t.Errorf("a plan waiting on a paused seat counted late") + } + + some := pauseOf([]string{"g14", "ace"}, map[string]link.HolderState{"ace": {Paused: true}}) + if some.All || !reflect.DeepEqual(some.Nodes, []string{"ace"}) { + t.Fatalf("%+v", some) + } + if line := planLineWith(p, now, some); !strings.Contains(line, "LATE") { + t.Errorf("a seat paused on one holder of two made the plan not late: %q", line) + } + if st := planStatuses([]inventory.Plan{p}, now, some)[0]; !st.Late { + t.Errorf("status --json: %+v", st) + } + // A plan rolling out, or with nothing asked, is not waiting on the seat. + rolling := p + rolling.State = inventory.PlanRolling + if _, paused := pausedWaiting(rolling, all, now); paused { + t.Error("a rolling plan reads as waiting on the build seat") + } +} + +// The queue's verbs are the controller seat's, each to the command it names. +func TestTheQueueVerbsRunTheirCommands(t *testing.T) { + for _, c := range []struct { + verb string + args map[string]any + want []string + }{ + {"queue", nil, []string{"queue"}}, + {"cancel", map[string]any{"id": "build-1"}, []string{"cancel", "build-1"}}, + {"clear", nil, []string{"clear"}}, + {"clear", map[string]any{"dead": "true"}, []string{"clear", "--dead"}}, + {"rebuild", map[string]any{"what": "gitea"}, []string{"rebuild", "gitea"}}, + {"replay", map[string]any{"id": "build-1"}, []string{"replay", "build-1"}}, + {"replay", map[string]any{"id": "build-1", "register": "true", "older": "true"}, []string{"replay", "build-1", "--register", "--older"}}, + {"kill", map[string]any{"id": "build-1"}, []string{"kill", "build-1"}}, + {"pause", nil, []string{"pause"}}, + {"resume", map[string]any{"node": "ace"}, []string{"resume", "ace"}}, + {"plans", map[string]any{"retry": "plan-1"}, []string{"plans", "retry", "plan-1"}}, + } { + got, err := argvFor(c.verb, c.args) + if err != nil || !reflect.DeepEqual(got, c.want) { + t.Errorf("%s %v: %v %v, want %v", c.verb, c.args, got, err, c.want) + } + } + if _, err := argvFor("cancel", nil); err == nil { + t.Error("cancel without an id was taken") + } + declared := map[string]bool{} + for _, v := range catalogue.ControllerVerbs { + declared[v.Name] = true + } + for _, v := range []string{"queue", "cancel", "clear", "rebuild", "replay", "kill", "pause", "resume"} { + if !declared[v] { + t.Errorf("%s is not a verb of the controller seat", v) + } + } +} + +// --- against a real bus ----------------------------------------------------------------------- + +// aBuildQueue is the build seat's queue, worker and cancelled set on a real server, the controller +// pointed at it, and a function that asks the seat one build. +func aBuildQueue(t *testing.T) (*broker.JetStream, func(id, repository string) link.BuildRequest) { + t.Helper() + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + t.Setenv(broker.NATSVar, url) + js, err := broker.Dial(url) + if err != nil { + t.Fatal(err) + } + t.Cleanup(js.Close) + seat := broker.DeclaredSeat{Name: link.TheBuildMachine, Accepts: []string{"build"}, + Emits: []string{"started", "built", "log.*", "paused.*"}} + if err := broker.AssertMeshStreams(js); err != nil { + t.Fatal(err) + } + _ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT") + if err := broker.RaiseSeats(js, []broker.DeclaredSeat{seat}, map[string]broker.Holder{ + link.TheBuildMachine: {Node: "anchor", Module: "build-agent"}}); err != nil { + t.Fatal(err) + } + if err := broker.RaiseCancelledSets(js, []broker.DeclaredSeat{seat}); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + _ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT") + _ = js.Context().DeleteKeyValue(broker.CancelledSetName(link.TheBuildMachine)) + _ = js.Context().PurgeStream(broker.EventsStream) + }) + ask := func(id, repository string) link.BuildRequest { + r := link.BuildRequest{ID: id, Repository: repository, Held: map[string]string{"x/y": "secret-ish"}} + body, _ := json.Marshal(r) + if _, err := js.Context().Publish(link.BuildWork(), body); err != nil { + t.Fatal(err) + } + return r + } + return js, ask +} + +// deadOne takes the oldest ask from the worker and hands it back as often as the worker allows. +func deadOne(t *testing.T, js *broker.JetStream) { + t.Helper() + worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: link.TheBuildMachine, Accepts: []string{"build"}}) + sub, err := js.Context().PullSubscribe(worker.Filters[0], worker.Name, nats.Bind(worker.Stream, worker.Name), nats.ManualAck()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = sub.Unsubscribe() }() + for i := 0; i < worker.MaxDeliver; i++ { + msgs, err := sub.Fetch(1, nats.MaxWait(3*time.Second)) + if err != nil { + t.Fatalf("delivery %d: %v", i+1, err) + } + _ = msgs[0].Nak() + } + // One more pull, as a holder always has one waiting: the server finds the ask past its deliveries + // then, and stops counting it pending. + if msgs, _ := sub.Fetch(1, nats.MaxWait(time.Second)); len(msgs) > 0 { + t.Fatalf("an ask past its deliveries was delivered again") + } +} + +// cancel drops a waiting ask from the bus and records it failed, cancelled by hand — and the plan +// that asked for it, matched by the id, fails with it; an ask the plan's records name no module for +// is still found. +func TestCancelDeletesTheAskAndFailsThePlanThatAskedIt(t *testing.T) { + js, ask := aBuildQueue(t) + open := aMesh(t) + ctx := t.Context() + twoTiers(t, open) + id := link.NewBuildID(time.Now()) + ask(id, "https://forge.example/novox/a.git") + asked, _ := link.BuildAskedAt(id) + plan := inventory.Plan{ID: "plan-cancel", Repository: "novox/a", Commit: "c0ffee", Created: asked, + State: inventory.PlanBuilding, Tiers: [][]string{{"a"}, {"b"}}, + Modules: map[string]*inventory.PlanModule{"a": {State: "asked", AskedAt: &asked, Build: id}}} + if err := open.inventory.SavePlan(ctx, plan); err != nil { + t.Fatal(err) + } + + q, err := link.ReadQueue(ctx, js, link.TheBuildMachine) + if err != nil { + t.Fatal(err) + } + if len(q.Asks) != 1 || q.Asks[0].State != link.AskWaiting || q.Asks[0].ID != id { + t.Fatalf("the queue reads %+v", q) + } + if text := queueText(q, time.Now()); !strings.Contains(text, "1 waiting, 0 in flight, 0 dead") || + strings.Contains(text, "secret-ish") { + t.Errorf("the queue says:\n%s", text) + } + + if err := cancelCommand(ctx, []string{id}); err != nil { + t.Fatal(err) + } + if q, _ = link.ReadQueue(ctx, js, link.TheBuildMachine); len(q.Asks) != 0 { + t.Fatalf("the ask is still queued: %+v", q.Asks) + } + if cancelled, err := link.IsCancelled(js.Conn(), link.TheBuildMachine, id); err != nil || !cancelled { + t.Errorf("the cancelled set does not hold it: %v %v", cancelled, err) + } + b, found, err := open.inventory.BuildByID(ctx, id) + if err != nil || !found || b.Failed != link.CancelledByHand { + t.Fatalf("the cancel is recorded as %+v (%v %v)", b, found, err) + } + p, err := open.inventory.PlanByID(ctx, plan.ID) + if err != nil { + t.Fatal(err) + } + if p.State != inventory.PlanFailed || p.Modules["a"].State != "failed" || p.Modules["a"].Why != link.CancelledByHand { + t.Fatalf("the plan that asked is %s, a %+v", p.State, p.Modules["a"]) + } + if err := cancelCommand(ctx, []string{id}); err == nil { + t.Error("an ask cancelled already was cancelled again") + } +} + +// clear cancels every waiting ask and leaves the dead ones unless told, and never one in flight. +func TestClearCancelsTheWaitingAndTheDeadOnlyWhenTold(t *testing.T) { + js, ask := aBuildQueue(t) + open := aMesh(t) + ctx := t.Context() + start := time.Now() + dead := link.NewBuildID(start) + ask(dead, "https://forge.example/novox/dead.git") + deadOne(t, js) + var waiting []string + for i := 1; i <= 2; i++ { + id := link.NewBuildID(start.Add(time.Duration(i) * time.Millisecond)) + waiting = append(waiting, id) + ask(id, fmt.Sprintf("https://forge.example/novox/w%d.git", i)) + } + + q, err := link.ReadQueue(ctx, js, link.TheBuildMachine) + if err != nil { + t.Fatal(err) + } + if len(q.Of(link.AskDead)) != 1 || len(q.Of(link.AskWaiting)) != 2 || q.MaxDeliver != 5 { + t.Fatalf("the queue reads %+v", q) + } + said, err := clearQueue(ctx, js, open, link.TheBuildMachine, false) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(said, "2 waiting ask(s) cancelled") || !strings.Contains(said, "1 dead, left") { + t.Errorf("clear said:\n%s", said) + } + q, _ = link.ReadQueue(ctx, js, link.TheBuildMachine) + if len(q.Asks) != 1 || q.Asks[0].ID != dead || q.Asks[0].State != link.AskDead { + t.Fatalf("after clear the queue is %+v", q.Asks) + } + if _, err := clearQueue(ctx, js, open, link.TheBuildMachine, true); err != nil { + t.Fatal(err) + } + if q, _ = link.ReadQueue(ctx, js, link.TheBuildMachine); len(q.Asks) != 0 { + t.Fatalf("clear --dead left %+v", q.Asks) + } + for _, id := range append(waiting, dead) { + if b, found, _ := open.inventory.BuildByID(ctx, id); !found || b.Failed != link.CancelledByHand { + t.Errorf("%s is recorded as %+v", id, b) + } + } +} + +// An ask in flight is not cancelled: kill ends it where it runs. +func TestCancelRefusesAnAskInFlight(t *testing.T) { + js, ask := aBuildQueue(t) + open := aMesh(t) + ctx := t.Context() + id := link.NewBuildID(time.Now()) + ask(id, "https://forge.example/novox/a.git") + worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: link.TheBuildMachine, Accepts: []string{"build"}}) + sub, err := js.Context().PullSubscribe(worker.Filters[0], worker.Name, nats.Bind(worker.Stream, worker.Name), nats.ManualAck()) + if err != nil { + t.Fatal(err) + } + defer func() { _ = sub.Unsubscribe() }() + if _, err := sub.Fetch(1, nats.MaxWait(3*time.Second)); err != nil { + t.Fatal(err) + } + started, _ := json.Marshal(link.BuildStart{ID: id, On: "ace", At: time.Now().UTC().Format(time.RFC3339Nano)}) + if _, err := js.Context().Publish(link.BuildStarted(), started); err != nil { + t.Fatal(err) + } + q, err := link.ReadQueue(ctx, js, link.TheBuildMachine) + if err != nil { + t.Fatal(err) + } + a, _ := q.Find(id) + if a.State != link.AskInFlight || a.On != "ace" { + t.Fatalf("the taken ask reads %+v", a) + } + _, err = cancelAsk(ctx, js, open, link.TheBuildMachine, a) + if err == nil || !strings.Contains(err.Error(), "kill "+id) { + t.Fatalf("an ask in flight was cancelled: %v", err) + } + if _, found, _ := open.inventory.BuildByID(ctx, id); found { + t.Error("a refused cancel recorded an outcome") + } +} + +// kill finds the machine from the build's start and asks that machine's holder; pause asks the +// machine named. Each prints what the holder answered, and a refusal is the command's failure. +func TestKillAndPauseAskTheHolderOnTheMachine(t *testing.T) { + js, _ := aBuildQueue(t) + ctx := t.Context() + asked := map[string]string{} + handlers := map[string]link.ToolHandler{ + "kill": func(_ context.Context, raw json.RawMessage) (any, error) { + var args struct{ ID string } + _ = json.Unmarshal(raw, &args) + asked["kill"] = args.ID + if args.ID != "build-running" { + return nil, fmt.Errorf("ace is not building %s", args.ID) + } + return map[string]any{"said": "killed " + args.ID}, nil + }, + "pause": func(context.Context, json.RawMessage) (any, error) { + asked["pause"] = "ace" + return map[string]any{"said": "ace is paused"}, nil + }, + } + stop, err := link.OverNATS{Conn: js.Conn()}.ServeNodeSeatTools(link.TheBuildMachine, "ace", handlers, nil) + if err != nil { + t.Fatal(err) + } + defer stop() + for _, id := range []string{"build-running", "build-other"} { + started, _ := json.Marshal(link.BuildStart{ID: id, On: "ace", At: time.Now().UTC().Format(time.RFC3339Nano)}) + if _, err := js.Context().Publish(link.BuildStarted(), started); err != nil { + t.Fatal(err) + } + } + if err := killCommand(ctx, []string{"build-running"}); err != nil { + t.Fatal(err) + } + if asked["kill"] != "build-running" { + t.Fatalf("the holder was asked %v", asked) + } + if err := killCommand(ctx, []string{"build-other"}); err == nil { + t.Error("the holder's refusal was not the command's") + } + if err := killCommand(ctx, []string{"build-never"}); err == nil || !strings.Contains(err.Error(), "cancel build-never") { + t.Errorf("a build nobody started: %v", err) + } + if err := pauseCommand(ctx, "pause", []string{"ace"}); err != nil || asked["pause"] != "ace" { + t.Fatalf("pause: %v %v", err, asked) + } + if err := pauseCommand(ctx, "resume", []string{"g14"}); err == nil { + t.Error("a machine nothing answers on was resumed") + } +} diff --git a/cmd/mesh-controller/readable.go b/cmd/mesh-controller/readable.go index 6049fd0..119f9f7 100644 --- a/cmd/mesh-controller/readable.go +++ b/cmd/mesh-controller/readable.go @@ -176,7 +176,7 @@ func statusAsJSON(asked answers) ([]byte, error) { out := meshStatus{Machines: len(nodes), Wrong: []machineDoing{}, Quiet: []machineQuiet{}, Behind: []moduleBehind{}, Waiting: []machineWaiting{}, Reported: []machineReported{}, Unresolved: []machineUnresolved{}, - Network: asked.network, Adopted: adoptedNodes(nodes), Plans: planStatuses(asked.plans, time.Now())} + Network: asked.network, Adopted: adoptedNodes(nodes), Plans: planStatuses(asked.plans, time.Now(), asked.paused)} // In a stated order, so two readings of an unchanged mesh are the same document. untakenNodes := make([]string, 0, len(asked.untaken)) for name := range asked.untaken { diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 499760a..9e9d79b 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -324,37 +324,52 @@ func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) e for _, e := range entries { byName[e.Manifest.Module] = e } - now := time.Now().UTC() for _, name := range p.Tiers[p.Tier] { - state := p.Modules[name] - if state == nil { - state = &inventory.PlanModule{} - p.Modules[name] = state - } - e, known := byName[name] - if !known { - state.State = "failed" - state.Why = "no longer in the catalogue" - p.State = inventory.PlanFailed - p.Note = name + " is no longer in the catalogue" - continue - } - source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat} - fmt.Printf(" tier %d: ", p.Tier) - // The branch it follows, never a commit a build once named (novox/hq 04-ISSUES/215). - if err := buildOne(ctx, source, e.Source.Path, followedBranch(e.Source.Ref), 0); err != nil { - state.State = "failed" - state.Why = err.Error() - p.State = inventory.PlanFailed - p.Note = fmt.Sprintf("%s could not be asked for: %v", name, err) - continue - } - state.State = "asked" - state.AskedAt = &now + askModule(ctx, p, name, byName) } return nil } +// askABuild is how a plan asks for one build, not waited for, and learns the id it asked under. A +// variable so a test of what a plan does around an ask needs no build machine. +var askABuild = func(ctx context.Context, source buildSource, path, ref string) (string, error) { + return buildOneAsked(ctx, source, path, ref, 0, false) +} + +// askModule asks the build machine for one module of a plan and marks it asked, with the id it was +// asked under (novox/hq ADR 0219) — or failed, with the plan, when it could not be asked. +func askModule(ctx context.Context, p *inventory.Plan, name string, byName map[string]inventory.Entry) { + now := time.Now().UTC() + state := p.Modules[name] + if state == nil { + state = &inventory.PlanModule{} + p.Modules[name] = state + } + e, known := byName[name] + if !known { + state.State = "failed" + state.Why = "no longer in the catalogue" + p.State = inventory.PlanFailed + p.Note = name + " is no longer in the catalogue" + return + } + source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat} + fmt.Printf(" tier %d: ", p.Tier) + // The branch it follows, never a commit a build once named (novox/hq 04-ISSUES/215). + id, err := askABuild(ctx, source, e.Source.Path, followedBranch(e.Source.Ref)) + if err != nil { + state.State = "failed" + state.Why = err.Error() + p.State = inventory.PlanFailed + p.Note = fmt.Sprintf("%s could not be asked for: %v", name, err) + return + } + state.State = "asked" + state.AskedAt = &now + state.Build = id + state.Why, state.Commit, state.BuiltAt = "", "", nil +} + // planBuilt marks a module built (or failed) in every open plan whose current tier holds it, and // advances what that completes. Called from the daemon's take-in of every outcome. // @@ -363,7 +378,7 @@ func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) e // the later plan's answer — it stood on the bases from before the later plan's merge, and taking it // would send machines, and the next tier, what the later merge replaced. asked is zero when the // build's request time is not known, and such an outcome is taken as before. -func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time) { +func planBuilt(ctx context.Context, open *stores, module, commit, failed string, asked time.Time, id string) { inv := open.inventory // One controller works the plans at a time (novox/hq issue 213); an outcome waits its turn rather // than write over what the holder is about to save. Not taken, it is still in the build records, @@ -399,7 +414,9 @@ func planBuilt(ctx context.Context, open *stores, module, commit, failed string, state = &inventory.PlanModule{} p.Modules[module] = state } - if askedBefore(asked, state.AskedAt) { + // **The plan's own ask is its outcome, by id** (novox/hq ADR 0219); another build of the module + // is, as before, when it was asked at or after the plan's ask (issue 219). + if !(id != "" && state.Build == id) && askedBefore(asked, state.AskedAt) { continue } if failed != "" { @@ -523,6 +540,7 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, // plan never hears it. The record is the fact; a build recorded after the ask is that tier's // outcome, whoever was listening. recorded := map[string][]inventory.Build{} + byID := map[string]inventory.Build{} for _, m := range tier { if s := p.Modules[m]; s != nil && s.State == "asked" { builds, err := inv.Builds(ctx, m, 5) @@ -530,9 +548,19 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, return false, err } recorded[m] = builds + // Its own ask's record, by id — found even when the outcome named no module (ADR 0219). + if s.Build != "" { + b, found, err := inv.BuildByID(ctx, s.Build) + if err != nil { + return false, err + } + if found { + byID[s.Build] = b + } + } } } - if settleFromRecords(p, tier, recorded) { + if settleFromRecords(p, tier, recorded, byID) { return true, nil } // Asked: wait for every build. @@ -811,7 +839,11 @@ func planTicker(ctx context.Context, open *stores) { } // planLine is one plan as `status` says it. -func planLine(p inventory.Plan, now time.Time) string { +func planLine(p inventory.Plan, now time.Time) string { return planLineWith(p, now, pauseView{}) } + +// planLineWith is planLine knowing whether the build seat is paused (novox/hq ADR 0219): a plan +// waiting on builds nobody will take until a person resumes the seat says so, and is not late. +func planLineWith(p inventory.Plan, now time.Time, pause pauseView) string { where := fmt.Sprintf("tier %d of %d", min(p.Tier+1, len(p.Tiers)), len(p.Tiers)) switch p.State { case inventory.PlanDone: @@ -822,6 +854,9 @@ func planLine(p inventory.Plan, now time.Time) string { return fmt.Sprintf("%s %s %s", p.Repository, short(p.Commit), p.Note) } since := now.Sub(p.Updated).Round(time.Second) + if waiting, paused := pausedWaiting(p, pause, now); paused { + return fmt.Sprintf("%s %s %s, %s", p.Repository, short(p.Commit), where, waiting) + } late := "" if since > planWaitBound { late = " — LATE" @@ -835,20 +870,49 @@ func planLine(p inventory.Plan, now time.Time) string { // planFailedBuild marks the module a failed build was for when the result names no module: by the // repository and path the plan's modules were asked at. +// +// **By the id first** (novox/hq ADR 0219): a plan keeps the id it asked each module under, so an +// outcome that never learnt its module's name — cancelled, killed, failed at the clone — is matched +// to the module it was asked for exactly. Repository and path remain for a plan from before ids +// were kept. func planFailedBuild(ctx context.Context, open *stores, result link.BuildResult) { + asked, _ := link.BuildAskedAt(result.ID) + if plans, err := open.inventory.OpenPlans(ctx); err == nil { + if module := moduleAskedAs(plans, result.ID); module != "" { + planBuilt(ctx, open, module, result.Commit, result.Failed, asked, result.ID) + return + } + } entries, err := open.inventory.Catalogued(ctx) if err != nil { return } for _, e := range entries { if repositoryMatches(e.Source.Repository, result.Repository) && e.Source.Path == result.Path { - asked, _ := link.BuildAskedAt(result.ID) - planBuilt(ctx, open, e.Manifest.Module, result.Commit, result.Failed, asked) + planBuilt(ctx, open, e.Manifest.Module, result.Commit, result.Failed, asked, result.ID) return } } } +// moduleAskedAs is the module an open plan's current tier asked for under this id, or nothing. +func moduleAskedAs(plans []inventory.Plan, id string) string { + if id == "" { + return "" + } + for _, p := range plans { + if p.Tier >= len(p.Tiers) { + continue + } + for _, m := range p.Tiers[p.Tier] { + if s := p.Modules[m]; s != nil && s.Build == id { + return m + } + } + } + return "" +} + func repositoryMatches(a, b string) bool { trim := func(s string) string { return strings.ToLower(strings.TrimSuffix(s, ".git")) } return trim(a) == trim(b) || strings.HasSuffix(trim(a), "/"+trim(b)) || strings.HasSuffix(trim(b), "/"+trim(a)) @@ -867,7 +931,7 @@ type planStatus struct { Late bool `json:"late"` } -func planStatuses(plans []inventory.Plan, now time.Time) []planStatus { +func planStatuses(plans []inventory.Plan, now time.Time, pause pauseView) []planStatus { out := make([]planStatus, 0, len(plans)) for _, p := range plans { ps := planStatus{ID: p.ID, Repository: p.Repository, Commit: p.Commit, State: p.State, @@ -878,6 +942,10 @@ func planStatuses(plans []inventory.Plan, now time.Time) []planStatus { ps.Waiting = "builds of tier " + fmt.Sprint(p.Tier) } ps.Late = now.Sub(p.Updated) > planWaitBound + // Paused is a person's decision, not lateness (novox/hq ADR 0219). + if waiting, paused := pausedWaiting(p, pause, now); paused { + ps.Waiting, ps.Late = waiting, false + } } out = append(out, ps) } @@ -885,12 +953,15 @@ func planStatuses(plans []inventory.Plan, now time.Time) []planStatus { } // openPlans is the open plans among the recent ones, and how many have waited past the bound. -func openPlans(plans []inventory.Plan) ([]inventory.Plan, int) { +func openPlans(plans []inventory.Plan, pause pauseView) ([]inventory.Plan, int) { var open []inventory.Plan late := 0 for _, p := range plans { if p.Open() { open = append(open, p) + if _, paused := pausedWaiting(p, pause, time.Now()); paused { + continue + } if time.Since(p.Updated) > planWaitBound { late++ } @@ -923,7 +994,7 @@ func plansCommand(ctx context.Context, args []string) error { if err != nil { return err } - fmt.Printf("%s — %s\n", p.ID, planLine(p, now)) + fmt.Printf("%s — %s\n", p.ID, planLineWith(p, now, buildSeatPause(ctx, inv, []inventory.Plan{p}))) for i, tier := range p.Tiers { marker := " " if i == p.Tier && p.Open() { @@ -938,6 +1009,9 @@ func plansCommand(ctx context.Context, args []string) error { if s.Commit != "" { state += " from " + short(s.Commit) } + if s.Build != "" && s.State != "built" { + state += " (" + s.Build + ")" + } if s.Why != "" { state += ": " + s.Why } @@ -950,6 +1024,15 @@ func plansCommand(ctx context.Context, args []string) error { if *whatIf != "" { return planWhatIf(ctx, inv, *whatIf, splitList(*paths), splitList(*modules)) } + // `retry` (novox/hq ADR 0219): a failed plan's failed builds asked again, and the plan goes on. + if len(positionals) == 2 && positionals[0] == "retry" { + said, err := retryPlan(ctx, open, positionals[1]) + if err != nil { + return err + } + fmt.Println(said) + return nil + } // `stop`, or `close` (novox/hq issue 254): a person ending a plan that will not move again — one // waiting on a report that cannot come — so it stops reading as work in progress. Marked failed // with who ended it; what it asked still builds and registers. @@ -991,8 +1074,9 @@ func plansCommand(ctx context.Context, args []string) error { fmt.Println("no merge has produced a plan yet") return nil } + pause := buildSeatPause(ctx, inv, plans) for _, p := range plans { - fmt.Printf("%-28s %s\n", p.ID, planLine(p, now)) + fmt.Printf("%-28s %s\n", p.ID, planLineWith(p, now, pause)) } return nil } @@ -1094,7 +1178,12 @@ func splitList(s string) []string { // settleFromRecords marks every module of the tier still `asked` built — or failed — from a build // recorded after it was asked, and says whether it changed anything (novox/hq 04-ISSUES/214). // Newest first, as Builds answers: the first record after the ask is the outcome of that ask. -func settleFromRecords(p *inventory.Plan, tier []string, recorded map[string][]inventory.Build) bool { +// +// The record of the plan's own ask, by its id, is that outcome before anything else (novox/hq ADR +// 0219): the plan asked under it, and a failure recorded without a module — cancelled, killed — is +// found by nothing else. +func settleFromRecords(p *inventory.Plan, tier []string, recorded map[string][]inventory.Build, + byID map[string]inventory.Build) bool { changed := false for _, m := range tier { s := p.Modules[m] @@ -1114,6 +1203,9 @@ func settleFromRecords(p *inventory.Plan, tier []string, recorded map[string][]i } outcome = &b } + if own, found := byID[s.Build]; s.Build != "" && found { + outcome = &own + } if outcome == nil { continue } diff --git a/cmd/mesh-controller/release_plan_test.go b/cmd/mesh-controller/release_plan_test.go index 365cdfc..e42c214 100644 --- a/cmd/mesh-controller/release_plan_test.go +++ b/cmd/mesh-controller/release_plan_test.go @@ -145,7 +145,7 @@ func TestAPlanSettlesAnAskedBuildFromTheRecords(t *testing.T) { // Only a build from before the ask: not this ask's outcome. "builder": {{ID: "build-0", Commit: "06ea2168", At: asked.Add(-time.Hour)}}, } - if !settleFromRecords(&p, p.Tiers[0], records) { + if !settleFromRecords(&p, p.Tiers[0], records, nil) { t.Fatal("nothing settled, though the controller's build is recorded after the ask") } if s := p.Modules["mesh-controller"]; s.State != "built" || s.Commit != "2ebbb799" || s.BuiltAt == nil { @@ -162,13 +162,13 @@ func TestAPlanSettlesAnAskedBuildFromTheRecords(t *testing.T) { late := map[string][]inventory.Build{"postgres": { {ID: "build-old", Commit: "efff5415", Asked: asked.Add(-18 * time.Minute), At: asked.Add(12 * time.Minute)}, }} - if settleFromRecords(&r, r.Tiers[0], late) || r.Modules["postgres"].State != "asked" { + if settleFromRecords(&r, r.Tiers[0], late, nil) || r.Modules["postgres"].State != "asked" { t.Errorf("an earlier ask's late outcome settled this ask: %+v", r.Modules["postgres"]) } // Newest heard first: the earlier ask's late outcome, then this ask's own, heard before it. late["postgres"] = append(late["postgres"], inventory.Build{ID: "build-mine", Commit: "4bcd5f73", Asked: asked.Add(time.Second), At: asked.Add(5 * time.Minute)}) - if !settleFromRecords(&r, r.Tiers[0], late) || r.Modules["postgres"].State != "built" || + if !settleFromRecords(&r, r.Tiers[0], late, nil) || r.Modules["postgres"].State != "built" || r.Modules["postgres"].Commit != "4bcd5f73" { t.Errorf("this ask's own outcome, heard before the earlier ask's, did not settle it: %+v", r.Modules["postgres"]) } @@ -176,7 +176,7 @@ func TestAPlanSettlesAnAskedBuildFromTheRecords(t *testing.T) { // A failure recorded after the ask fails the plan, as hearing it would have. q := inventory.Plan{ID: "plan-2", Tiers: [][]string{{"x"}}, Modules: map[string]*inventory.PlanModule{"x": {State: "asked", AskedAt: &asked}}} - settleFromRecords(&q, q.Tiers[0], map[string][]inventory.Build{"x": {{ID: "b", Failed: "no", At: asked.Add(time.Minute)}}}) + settleFromRecords(&q, q.Tiers[0], map[string][]inventory.Build{"x": {{ID: "b", Failed: "no", At: asked.Add(time.Minute)}}}, nil) if q.State != inventory.PlanFailed || q.Modules["x"].State != "failed" { t.Errorf("a recorded failure did not fail the plan: %+v %+v", q, q.Modules["x"]) } diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index d7dec42..a3fdfc0 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -103,10 +103,48 @@ func argvFor(verb string, args map[string]any) ([]string, error) { if id := str("close"); id != "" { return []string{"plans", "close", id}, nil } + if id := str("retry"); id != "" { + return []string{"plans", "retry", id}, nil + } if id := str("id"); id != "" { return []string{"plans", id}, nil } return []string{"plans"}, nil + // The build queue (novox/hq ADR 0219). + case "queue": + return []string{"queue"}, nil + case "cancel", "kill": + if err := need("id"); err != nil { + return nil, err + } + return []string{verb, str("id")}, nil + case "clear": + if str("dead") == "true" { + return []string{"clear", "--dead"}, nil + } + return []string{"clear"}, nil + case "rebuild": + if err := need("what"); err != nil { + return nil, err + } + return []string{"rebuild", str("what")}, nil + case "replay": + if err := need("id"); err != nil { + return nil, err + } + argv := []string{"replay", str("id")} + if str("register") == "true" { + argv = append(argv, "--register") + } + if str("older") == "true" { + argv = append(argv, "--older") + } + return argv, nil + case "pause", "resume": + if n := str("node"); n != "" { + return []string{verb, n}, nil + } + return []string{verb}, nil case "plan": if err := need("node"); err != nil { return nil, err diff --git a/cmd/mesh-controller/status.go b/cmd/mesh-controller/status.go index 9756347..7cf2141 100644 --- a/cmd/mesh-controller/status.go +++ b/cmd/mesh-controller/status.go @@ -138,14 +138,14 @@ func printStatus(asked answers) error { len(quiet), strings.Join(said, "\n ")) } - if open, late := openPlans(asked.plans); len(open) > 0 { + if open, late := openPlans(asked.plans, asked.paused); len(open) > 0 { fmt.Printf("%d plan(s) open", len(open)) if late > 0 { fmt.Printf(", %d waiting past %s", late, planWaitBound) } fmt.Println(":") for _, p := range open { - fmt.Printf(" %s\n", planLine(p, time.Now())) + fmt.Printf(" %s\n", planLineWith(p, time.Now(), asked.paused)) } fmt.Println() } @@ -420,6 +420,8 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) { if err != nil { return answers{}, err } + // Whether the build seat takes work, for a plan waiting on it (novox/hq ADR 0219). + out.paused = buildSeatPause(ctx, inv, out.plans) // And which machines are not running what the mesh would send them. The same question as a // module being behind its source, one level down: that one says the catalogue is out of date, diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index 8898585..decb038 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -97,6 +97,7 @@ var ControllerVerbs = []Verb{ "id": "a plan's id (as `plans` lists them): that plan, tier by tier", "stop": "a plan's id: stop it — what was asked still builds, nothing further is asked", "close": "a plan's id: close a plan that will not move again, as failed by hand (novox/hq issue 254)", + "retry": "a failed plan's id: ask its failed builds again under new ids, and carry the plan on from that tier (novox/hq ADR 0219)", "repository": "owner/repository: the plan a merge there would produce, saving nothing (what-if); with paths or modules", "paths": "with repository: the files the merge would change, comma-separated, from the repository's root", "modules": "with repository: or the modules it would change, comma-separated", @@ -156,6 +157,36 @@ var ControllerVerbs = []Verb{ Input: schema(map[string]string{ "command": "the command line, as the controller's binary takes it; quotes group a word with spaces", }, []string{"command"})}, + // The build queue, controlled by hand (novox/hq ADR 0219). Every verb that drops an ask leaves a + // failed outcome for it, so a plan waiting on it fails visibly instead of hanging. + {Name: "queue", Description: "Every ask in the build queue: waiting, in flight (on which machine, for how long), " + + "and dead (handed out as often as allowed and never settled) — each with its id, repository, path, ref and when it was asked.", + Input: schema(nil, nil)}, + {Name: "cancel", Description: "Drop one waiting or dead ask from the build queue; its outcome is recorded failed, " + + "cancelled by hand, and a plan that asked for it fails. One in flight is refused: `kill` ends it where it runs.", + Input: schema(map[string]string{"id": "the ask's build id, as `queue` lists it"}, []string{"id"})}, + {Name: "clear", Description: "Cancel every waiting ask in the build queue — and with dead, every dead one too — each " + + "recorded failed, cancelled by hand. Never touches one in flight.", + Input: schema(map[string]string{"dead": "\"true\" to cancel the dead asks as well"}, nil)}, + {Name: "rebuild", Description: "Ask a module's current source again under a new id — the branch it follows — or, given a " + + "build's id, that build's repository, path and ref. A module a plan holds unbuilt or failed joins that plan. Answers the new id.", + Input: schema(map[string]string{"what": "a module's name, or a build's id"}, []string{"what"})}, + {Name: "replay", Description: "Ask a recorded build's repository and path again at the commit it built, under a new id. " + + "A dry run unless register: nothing recorded or registered. Registering is refused when a newer build of the module " + + "is registered — it would roll the older commit out (novox/hq issue 207) — unless older says so.", + Input: schema(map[string]string{ + "id": "the build's id", + "register": "\"true\" to register what it builds", + "older": "\"true\", with register: even though a newer build of the module is registered", + }, []string{"id"})}, + {Name: "kill", Description: "End a build where it runs: the machine that took it stops its commands and containers and " + + "announces it failed, killed by hand — settled, never handed to another machine.", + Input: schema(map[string]string{"id": "the build's id"}, []string{"id"})}, + {Name: "pause", Description: "The build seat's holder on one machine — or every holder — takes no new build until resumed; " + + "a build running finishes. Kept across a restart of the holder. A plan waiting on a paused seat says so and is not late.", + Input: schema(map[string]string{"node": "one machine; every holder when absent"}, nil)}, + {Name: "resume", Description: "The build seat's holder on one machine — or every holder — takes builds again.", + Input: schema(map[string]string{"node": "one machine; every holder when absent"}, nil)}, {Name: "build", Description: "Have the build machine build a repository. Answers at once with the build's id: " + "`builds` with that id follows it line by line, and the module is registered when the outcome comes.", Input: schema(map[string]string{ diff --git a/internal/inventory/builds.go b/internal/inventory/builds.go index 5400acf..082fce2 100644 --- a/internal/inventory/builds.go +++ b/internal/inventory/builds.go @@ -3,8 +3,11 @@ package inventory import ( "context" "encoding/json" + "errors" "sort" "time" + + "github.com/jackc/pgx/v5" ) // What has been built. @@ -161,6 +164,33 @@ func (i *Inventory) Builds(ctx context.Context, module string, limit int) ([]Bui return out, rows.Err() } +// BuildByID is one build's record, by the id its request carried — what a plan matches its outcome +// by when the outcome names no module (novox/hq ADR 0219), and what `rebuild` and `replay` read the +// repository, path, ref and commit from. False when no outcome with that id was recorded. +func (i *Inventory) BuildByID(ctx context.Context, id string) (Build, bool, error) { + var b Build + var made []byte + var asked *time.Time + err := i.store.Pool().QueryRow(ctx, + `select id, repository, ref, coalesce(module,''), commit_hash, built_on, failed, made, + coalesce(source_path,''), asked, at + from build where id = $1`, id).Scan(&b.ID, &b.Repository, &b.Ref, &b.Module, &b.Commit, + &b.On, &b.Failed, &made, &b.Path, &asked, &b.At) + if errors.Is(err, pgx.ErrNoRows) { + return Build{}, false, nil + } + if err != nil { + return Build{}, false, err + } + if asked != nil { + b.Asked = *asked + } + if err := json.Unmarshal(made, &b.Made); err != nil { + return Build{}, false, err + } + return b, true, nil +} + // Held is every artifact this mesh has built, keyed "/". // // **The successful build of each module asked last wins**, which is the same rule the rest of the diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index cf90380..50e8b89 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -49,6 +49,11 @@ type PlanModule struct { FirstAt *time.Time `json:"first_at,omitempty"` Commit string `json:"commit,omitempty"` Why string `json:"why,omitempty"` + // Build is the id of the build the plan asked for this module (novox/hq ADR 0219), so the plan + // matches its outcome by id — the one thing every outcome echoes, a failed one that never learnt + // its module's name included. Empty in a plan from before it was kept, which is matched by + // module, or by repository and path, as before. + Build string `json:"build,omitempty"` } // The states a plan passes through.