From 541603c15c3a814dcfcd4b372fae41a9cfcb5981 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 19:26:40 +0200 Subject: [PATCH] Announce the build agent's verbs, let a waiting build hear its cancel, and retry a stopped rollout (hq ADR 0219) The holder's verbs were served and found by nothing; it now answers discovery with its machine's four, as a runtime announces a seat's verb, so the console finds /node-build-agent.kill. A cancel publishes no outcome, so a build waited for also looks at the cancelled set. A plan that stopped at its first machine is retried by sending that machine the module again, unless a newer plan holds it. --- cmd/mesh-builder/holder.go | 41 +++++++++++++ cmd/mesh-builder/holder_test.go | 19 +++++++ cmd/mesh-builder/main.go | 7 +++ cmd/mesh-controller/plan_retry.go | 95 +++++++++++++++++++++++++++++-- cmd/mesh-controller/queue_test.go | 64 +++++++++++++++++++++ internal/link/builds_nats.go | 18 +++++- internal/link/queue_test.go | 36 ++++++++++++ 7 files changed, 273 insertions(+), 7 deletions(-) diff --git a/cmd/mesh-builder/holder.go b/cmd/mesh-builder/holder.go index 95d69f4..7240ede 100644 --- a/cmd/mesh-builder/holder.go +++ b/cmd/mesh-builder/holder.go @@ -12,7 +12,10 @@ import ( "sync" "time" + "github.com/nats-io/nats.go/micro" + "github.com/novox/mesh-controller/internal/builder" + "github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/link" ) @@ -293,3 +296,41 @@ func plainRun(ctx context.Context, dir, name string, args ...string) (string, er } return string(out), nil } + +// announcement is what this holder answers discovery with (novox/hq ADR 0195, ADR 0197): this +// machine's verbs of the seat, in the shape every tool runtime announces a seat's verb — kind seat, +// the module answering, the seat, scope node, the machine, its description and argument schema — so +// the console finds `/node-build-agent.kill` by searching, as it finds any seat's verb. +// +// One service per machine, named for the seat and identified by the machine, so the answer is this +// machine's four verbs and nothing more. The verbs are the compiled seat row's: a build machine has no +// store, and what it serves is what this binary was built to serve. +func announcement(seat, module, node string) micro.Info { + s, _ := catalogue.SeatNamed(seat) + var endpoints []micro.EndpointInfo + for _, v := range s.Serves { + schema, _ := json.Marshal(v.Input) + endpoints = append(endpoints, micro.EndpointInfo{ + Name: seat + "__" + v.Name, + Subject: link.NodeSeatToolSubject(seat, v.Name, node), + Metadata: map[string]string{ + "kind": "seat", "module": module, "tool": v.Name, "seat": seat, "scope": "node", + "node": node, "interchangeable": "false", "description": v.Description, "schema": string(schema), + }, + }) + } + return micro.Info{ + ServiceIdentity: micro.ServiceIdentity{Name: seat, ID: node, Version: "0.1.0", + Metadata: map[string]string{"seat": seat, "scope": "node", "node": node, "module": module}}, + Description: "what the build running on " + node + " is, and ending, pausing and resuming it (novox/hq ADR 0219)", + Endpoints: endpoints, + } +} + +// moduleOf is the module a credential was issued for: its user is `.`. +func moduleOf(user string) string { + if _, module, ok := strings.Cut(user, "."); ok && module != "" { + return module + } + return "build-agent" +} diff --git a/cmd/mesh-builder/holder_test.go b/cmd/mesh-builder/holder_test.go index f5b5238..51a28c0 100644 --- a/cmd/mesh-builder/holder_test.go +++ b/cmd/mesh-builder/holder_test.go @@ -92,3 +92,22 @@ func TestKillEndsTheBuildRunningHereAndNoOther(t *testing.T) { t.Errorf("after the kill current says %+v", c) } } + +// A holder announces this machine's verbs of the seat as the console reads a seat's verb, and no more. +func TestAHolderAnnouncesItsMachinesVerbsForTheConsole(t *testing.T) { + info := announcement(link.TheBuildMachine, moduleOf("ace.build-agent"), "ace") + if info.Name != "node-build-agent" || info.ID != "ace" || len(info.Endpoints) != 4 { + t.Fatalf("announced %s/%s with %d endpoints", info.Name, info.ID, len(info.Endpoints)) + } + kill := info.Endpoints[1] + md := kill.Metadata + if kill.Subject != "mesh.seat.node-build-agent.tool.kill.ace" || md["kind"] != "seat" || md["seat"] != "node-build-agent" || + md["scope"] != "node" || md["node"] != "ace" || md["tool"] != "kill" || md["module"] != "build-agent" || + !strings.Contains(md["description"], "killed by hand") || !strings.Contains(md["schema"], `"id"`) { + t.Fatalf("kill is announced as %+v", kill) + } + body, err := json.Marshal(info) + if err != nil || len(body) > 8*1024 { + t.Fatalf("the answer is %d bytes (%v)", len(body), err) + } +} diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index a28b791..24fdcd8 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -141,6 +141,13 @@ func run() error { return err } defer stopServing() + // And says so, for the console to find (novox/hq ADR 0197). + stopAnnouncing, err := link.OverNATS{Conn: js.Conn()}.Announce( + announcement(seat, moduleOf(credential.User), on), log.New(os.Stderr, "", 0)) + if err != nil { + return err + } + defer stopAnnouncing() go h.stateWhilePaused(ctx, time.Hour) } diff --git a/cmd/mesh-controller/plan_retry.go b/cmd/mesh-controller/plan_retry.go index 405f465..4165250 100644 --- a/cmd/mesh-controller/plan_retry.go +++ b/cmd/mesh-controller/plan_retry.go @@ -156,17 +156,67 @@ func retryRefusal(p inventory.Plan, plans []inventory.Plan) error { 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 { + 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) + } + return nil } - 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) + stopped := stoppedRollouts(p) + if len(stopped) == 0 { + return fmt.Errorf("nothing in tier %d of %s failed to build or stopped rolling out — it stopped at: %s", + p.Tier, p.ID, p.Note) + } + // **A rollout is retried unless the module has moved on**: a newer plan holding it sends — or + // sent — a newer build, and sending this one again would put the older build back on its machines. + for _, m := range stopped { + if q, found := newerPlanFor(m, p, plans); found { + return fmt.Errorf("%s has a newer plan, %s (%s, %s at %s): sending %s's build of it again would put "+ + "the older build back", m, q.ID, q.State, q.Repository, short(q.Commit), p.ID) + } } return nil } +// stoppedRollouts is the modules of a plan's current tier whose rollout stopped at its first machine +// (issue 249, ADR 0218): built, sent to the first machine, never to the rest, and why it stopped kept. +func stoppedRollouts(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 == "built" && s.FirstAt != nil && s.SentAt == nil && + len(s.First) > 0 && s.Why != "" { + out = append(out, m) + } + } + sort.Strings(out) + return out +} + +// newerPlanFor is a plan made after this one that holds the module, superseded ones aside. +func newerPlanFor(module string, p inventory.Plan, plans []inventory.Plan) (inventory.Plan, bool) { + for _, q := range plans { + if q.ID == p.ID || q.State == inventory.PlanSuperseded || !q.Created.After(p.Created) { + continue + } + for _, tier := range q.Tiers { + for _, m := range tier { + if m == module { + return q, true + } + } + } + } + return inventory.Plan{}, false +} + +// sendRollout sends machines what the mesh would send them now, answering the ones it sent. A +// variable so a test of a retried rollout needs no machine. +var sendRollout = sendToEach + // 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 { @@ -192,9 +242,17 @@ func retryPlan(ctx context.Context, open *stores, id string) (string, error) { if err != nil { return "", err } + recent, err := inv.RecentPlans(ctx, 50) + if err != nil { + return "", err + } + plans = append(plans, recent...) if err := retryRefusal(p, plans); err != nil { return "", err } + if len(failedIn(p)) == 0 { + return retryRollouts(ctx, open, &p) + } entries, err := inv.Catalogued(ctx) if err != nil { return "", err @@ -306,3 +364,28 @@ func planHolding(module string, openPlans, recent []inventory.Plan) (inventory.P } return inventory.Plan{}, false } + +// retryRollouts sends each module whose rollout stopped to the machines it was first sent to, again, +// records that send as the first anew, and sets the plan rolling: from there it goes on as the plan +// would have — the rest sent once those report they applied it, the next tier after (ADR 0218). +func retryRollouts(ctx context.Context, open *stores, p *inventory.Plan) (string, error) { + var said []string + for _, m := range stoppedRollouts(*p) { + s := p.Modules[m] + sent, err := sendRollout(ctx, open, s.First) + if err != nil { + return "", fmt.Errorf("%s could not be sent to %s again, so %s stays failed: %w", + m, strings.Join(s.First, ", "), p.ID, err) + } + now := time.Now().UTC() + s.First, s.FirstAt, s.Why = sent, &now, "" + said = append(said, m+" to "+strings.Join(sent, ", ")) + } + p.State = inventory.PlanRolling + p.Note = fmt.Sprintf("tier %d retried by hand; sent %s first again", p.Tier, strings.Join(said, "; ")) + if err := open.inventory.SavePlan(ctx, *p); err != nil { + return "", err + } + return fmt.Sprintf("%s retried at tier %d of %d: sent %s first again; the rest follow once it reports it "+ + "applied, as the plan would have", p.ID, p.Tier, len(p.Tiers), strings.Join(said, "; ")), nil +} diff --git a/cmd/mesh-controller/queue_test.go b/cmd/mesh-controller/queue_test.go index 2401d0a..bc23f0e 100644 --- a/cmd/mesh-controller/queue_test.go +++ b/cmd/mesh-controller/queue_test.go @@ -580,3 +580,67 @@ func TestKillAndPauseAskTheHolderOnTheMachine(t *testing.T) { t.Error("a machine nothing answers on was resumed") } } + +// A plan that stopped at its first machine is retried: the module sent to that machine again, the +// send recorded as the first anew, and the plan goes on — unless a newer plan holds the module. +func TestAPlanStoppedAtItsFirstMachineIsRetried(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + asked := asksRecorded(t) + twoTiers(t, open) + var sentTo [][]string + was := sendRollout + sendRollout = func(_ context.Context, _ *stores, names []string) ([]string, error) { + sentTo = append(sentTo, names) + return names, nil + } + t.Cleanup(func() { sendRollout = was }) + + long := time.Now().UTC().Add(-2 * time.Hour) + stopped := inventory.Plan{ID: "plan-rollout", Repository: "novox/a", Branch: "main", Commit: "c0ffee", + Created: long, State: inventory.PlanFailed, Tiers: [][]string{{"a"}, {"b"}}, + Note: "a stopped at its first machine in tier 0: laptop refused what it was sent", + Modules: map[string]*inventory.PlanModule{"a": {State: "built", BuiltAt: &long, Commit: "c0ffee", + First: []string{"laptop"}, FirstAt: &long, Why: "laptop refused what it was sent"}}} + if err := open.inventory.SavePlan(ctx, stopped); err != nil { + t.Fatal(err) + } + + // A newer plan holding a refuses it: sending the older build would put it back. + newer := inventory.Plan{ID: "plan-newer", Repository: "novox/other", Commit: "d00d", Created: long.Add(time.Hour), + State: inventory.PlanDone, Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{}} + if err := open.inventory.SavePlan(ctx, newer); err != nil { + t.Fatal(err) + } + if _, err := retryPlan(ctx, open, stopped.ID); err == nil || !strings.Contains(err.Error(), "plan-newer") { + t.Fatalf("retried under a newer plan: %v", err) + } + newer.State = inventory.PlanSuperseded + if err := open.inventory.SavePlan(ctx, newer); err != nil { + t.Fatal(err) + } + + said, err := retryPlan(ctx, open, stopped.ID) + if err != nil { + t.Fatal(err) + } + if len(sentTo) != 1 || !reflect.DeepEqual(sentTo[0], []string{"laptop"}) || !strings.Contains(said, "laptop") { + t.Fatalf("sent %v; said %q", sentTo, said) + } + p, err := open.inventory.PlanByID(ctx, stopped.ID) + if err != nil { + t.Fatal(err) + } + a := p.Modules["a"] + if p.State != inventory.PlanRolling || a.FirstAt == nil || !a.FirstAt.After(long) || a.Why != "" { + t.Fatalf("after retry the plan is %s, a %+v", p.State, a) + } + // And it goes on: a records (its policy sends nothing more), so the next tier is asked. + advancePlans(ctx, open) + if p, _ = open.inventory.PlanByID(ctx, stopped.ID); p.Tier != 1 || p.Modules["b"] == nil || p.Modules["b"].State != "asked" { + t.Fatalf("the retried plan did not go on: tier %d %s %+v", p.Tier, p.State, p.Modules["b"]) + } + if len(*asked) != 1 { + t.Fatalf("asked %v", *asked) + } +} diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index 343b371..dfb7b1b 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -106,7 +106,20 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest, waiting, cancelWait := context.WithTimeout(ctx, wait) defer cancelWait() for { - msg, err := outcomes.NextMsgWithContext(waiting) + // **A wait that hears a cancel** (novox/hq ADR 0219). A person cancelling an ask records its + // failure without publishing the role's outcome — that subject is the holders', and the + // controller may not speak for them — so the waiter also looks, every while, at the cancelled + // set the cancel wrote first. A kill needs no look: the holder announces it as any outcome. + slice, endSlice := context.WithTimeout(waiting, cancelLook) + msg, err := outcomes.NextMsgWithContext(slice) + endSlice() + if errors.Is(err, context.DeadlineExceeded) && waiting.Err() == nil { + if cancelled, _ := IsCancelled(b.js.Conn(), b.role(), request.ID); cancelled { + return BuildResult{ID: request.ID, Repository: request.Repository, Path: request.Path, + Ref: request.Ref, Source: request.Source, DryRun: request.DryRun, Failed: CancelledByHand}, nil + } + continue + } switch { case errors.Is(err, context.DeadlineExceeded): return BuildResult{}, waitingFor(wait) @@ -125,6 +138,9 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest, } } +// cancelLook is how often a waiting asker looks whether its ask was cancelled. +var cancelLook = 2 * time.Second + // --- the machine's side --------------------------------------------------------------------- type natsMachine struct { diff --git a/internal/link/queue_test.go b/internal/link/queue_test.go index 6f824b3..5da7672 100644 --- a/internal/link/queue_test.go +++ b/internal/link/queue_test.go @@ -229,3 +229,39 @@ func TestNatsAHoldersPausedStateIsReadBack(t *testing.T) { t.Fatalf("read back %+v", said) } } + +// A `build` waiting on its outcome hears a cancel — which publishes no outcome — from the cancelled +// set, and a kill from the outcome the holder announces (novox/hq ADR 0219). +func TestNatsAWaitingAskerHearsItsAskCancelledOrKilled(t *testing.T) { + js := aBusWithACancelledSet(t) + was := cancelLook + cancelLook = 200 * time.Millisecond + t.Cleanup(func() { cancelLook = was }) + ask := &natsBuilds{js: js, seat: TheBuildMachine} + + go func() { + time.Sleep(500 * time.Millisecond) + _ = MarkCancelled(context.Background(), js, TheBuildMachine, "build-waited-cancelled") + }() + result, err := ask.Submit(context.Background(), BuildRequest{ID: "build-waited-cancelled", Repository: "/r"}, 10*time.Second) + if err != nil || result.Failed != CancelledByHand || result.ID != "build-waited-cancelled" { + t.Fatalf("the waiter heard %+v (%v)", result, err) + } + + // Killed: the holder announces the failure as any outcome, and the waiter has it. + ctx, stop := context.WithCancel(context.Background()) + defer stop() + machine := MachineOverNATS(js, "ace") + defer machine.Close() + go func() { + _ = machine.Take(ctx, func(ctx context.Context, work Build) { + r := work.Request() + _ = work.Announce(ctx, BuildResult{ID: r.ID, Repository: r.Repository, On: "ace", Failed: KilledByHand}) + _ = work.Done() + }) + }() + result, err = ask.Submit(context.Background(), BuildRequest{ID: "build-waited-killed", Repository: "/r"}, 10*time.Second) + if err != nil || result.Failed != KilledByHand || result.On != "ace" { + t.Fatalf("the waiter heard %+v (%v)", result, err) + } +}