From 9873b3bf137a68a96e84a5c3b62171a1718c63df Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 19:45:45 +0200 Subject: [PATCH] Keep a replay from moving the mesh backwards, and tighten the queue's edges (review of hq ADR 0219) A registered replay asked now outranked newer asks of its module, and a plan took any later outcome as its answer. A replay is now refused while the module is asked anywhere, a plan module asked under an id is answered by that id alone, and a rebuild of a commit asks what the module follows. An ask handed back after a restart no longer reads as dead; a cancel that meets a start is withdrawn; pause holds for an ask fetched as it lands; kill removes containers before and after the build ends and says whether its outcome went out; a holder may say only its own machine is paused. --- cmd/mesh-builder/holder.go | 96 ++++++++++++---- cmd/mesh-builder/holder_test.go | 45 +++++++- cmd/mesh-builder/main.go | 14 ++- cmd/mesh-controller/plan_retry.go | 14 ++- cmd/mesh-controller/queue.go | 170 ++++++++++++++++++++++++---- cmd/mesh-controller/queue_test.go | 140 +++++++++++++++++++++-- cmd/mesh-controller/release_plan.go | 14 ++- internal/broker/cancelled_test.go | 4 +- internal/broker/nats.go | 11 ++ internal/link/builds_nats.go | 24 +++- internal/link/builds_nats_test.go | 7 +- internal/link/queue.go | 82 +++++++++----- internal/link/queue_test.go | 85 ++++++++++++-- 13 files changed, 600 insertions(+), 106 deletions(-) diff --git a/cmd/mesh-builder/holder.go b/cmd/mesh-builder/holder.go index 7240ede..1c3edf3 100644 --- a/cmd/mesh-builder/holder.go +++ b/cmd/mesh-builder/holder.go @@ -54,7 +54,12 @@ type running struct { started time.Time cancel context.CancelFunc killed bool - done chan struct{} + // returned is the build's own work having ended, before a kill or not; outcome is what was then + // announced, and announced whether it went out. + returned bool + outcome string + announced bool + done chan struct{} } // pausedFile is where the flag lives in the workspace. @@ -152,13 +157,6 @@ func (h *holder) stepped(r *running, step string) { h.mu.Unlock() } -// wasKilled is whether a person killed this build. -func (h *holder) wasKilled(r *running) bool { - h.mu.Lock() - defer h.mu.Unlock() - return r.killed -} - // currentBuild is what `current` answers. type currentBuild struct { On string `json:"on"` @@ -204,10 +202,19 @@ func (h *holder) current() currentBuild { return out } -// kill ends the build with this id, if it is the one running here: its context cancelled — which -// kills each command's process group — and the containers it started removed by their label. The -// build's own goroutine then announces it failed, killed by hand, and settles the ask, so it is not -// redelivered; this waits a while for that, to say it happened. +// The bounds a kill keeps: each pass removing containers, and the wait for the build to end between +// them. The worst case — 15s, 20s, 15s — is inside what the controller waits for the answer +// (killAnswer in the controller's queue.go, 75s). +var ( + killRemoves = 15 * time.Second + killWaits = 20 * time.Second +) + +// kill ends the build with this id, if it is the one running here and still working: its context +// cancelled — which kills each command's process group — and the containers it started removed by +// their label, once at once and again after the build has ended, so one created while it was being +// killed is not left. The build's own goroutine announces it failed, killed by hand, and settles the +// ask so it is not redelivered; the answer says whether that happened. func (h *holder) kill(id string) (string, error) { h.mu.Lock() r := h.running @@ -219,27 +226,68 @@ func (h *holder) kill(id string) (string, error) { h.mu.Unlock() return "", fmt.Errorf("%s is not building %s; it is building %s. `queue` says where an ask is", h.on, id, doing) } + if r.returned { + h.mu.Unlock() + return "", fmt.Errorf("%s on %s has already ended on its own and is saying how; `builds` shows it", id, h.on) + } r.killed = true cancel, done := r.cancel, r.done h.mu.Unlock() cancel() - cleanup, stop := context.WithTimeout(context.Background(), 30*time.Second) - defer stop() - removed, err := builder.RemoveContainersOf(cleanup, h.remove, id) - containers := fmt.Sprintf("%d container(s) it started removed", removed) - if err != nil { - containers = "its containers could not be listed or removed: " + err.Error() - } + removed, removeErr := h.removeContainers(id) + ended := false select { case <-done: - return fmt.Sprintf("killed %s (%s) on %s: its commands ended, %s, and its outcome announced as "+ - "failed, %s — settled, so it is not handed to another machine", id, r.request.Repository, h.on, - containers, link.KilledByHand), nil - case <-time.After(30 * time.Second): + ended = true + case <-time.After(killWaits): + } + again, againErr := h.removeContainers(id) + removed += again + if removeErr == nil { + removeErr = againErr + } + containers := fmt.Sprintf("%d container(s) it started removed", removed) + if removeErr != nil { + containers = "its containers could not all be listed or removed: " + removeErr.Error() + } + if !ended { return fmt.Sprintf("killed %s on %s: %s; the build has not finished ending yet — `builds` says when "+ "its outcome is in", id, h.on, containers), nil } + h.mu.Lock() + outcome, announced := r.outcome, r.announced + h.mu.Unlock() + if !announced { + return fmt.Sprintf("killed %s (%s) on %s: its commands ended and %s, and its outcome could not be "+ + "announced — the ask is not settled and will be handed out again", id, r.request.Repository, h.on, + containers), nil + } + return fmt.Sprintf("killed %s (%s) on %s: its commands ended, %s, and its outcome announced as failed, %s — "+ + "settled, so it is not handed to another machine", id, r.request.Repository, h.on, containers, outcome), nil +} + +// removeContainers is one pass of removing what the build left, bounded. +func (h *holder) removeContainers(id string) (int, error) { + cleanup, stop := context.WithTimeout(context.Background(), killRemoves) + defer stop() + return builder.RemoveContainersOf(cleanup, h.remove, id) +} + +// returned records that the build's work ended, and says whether a kill came first — only then is +// the build killed; an error it ended with on its own is its own outcome. +func (h *holder) returned(r *running) bool { + h.mu.Lock() + defer h.mu.Unlock() + r.returned = true + return r.killed +} + +// said records the outcome announced, and whether it went out, for a kill to answer with. +func (h *holder) said(r *running, outcome string, announced bool) { + h.mu.Lock() + r.outcome, r.announced = outcome, announced + h.mu.Unlock() } // handlers are the seat's verbs, as this machine answers them. @@ -313,6 +361,8 @@ func announcement(seat, module, node string) micro.Info { endpoints = append(endpoints, micro.EndpointInfo{ Name: seat + "__" + v.Name, Subject: link.NodeSeatToolSubject(seat, v.Name, node), + // The queue group the verbs are served in, as every runtime announces its own. + QueueGroup: "seat." + seat, 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), diff --git a/cmd/mesh-builder/holder_test.go b/cmd/mesh-builder/holder_test.go index 51a28c0..c27e1d7 100644 --- a/cmd/mesh-builder/holder_test.go +++ b/cmd/mesh-builder/holder_test.go @@ -66,9 +66,13 @@ func TestKillEndsTheBuildRunningHereAndNoOther(t *testing.T) { !strings.Contains(err.Error(), "build-1") { t.Fatalf("killed another id: %v", err) } - // The build's own goroutine: it ends when its context does, as a build's commands do. + // The build's own goroutine: it ends when its context does, as a build's commands do, and + // announces what came of it. go func() { <-building.Done() + if h.returned(r) { + h.said(r, link.KilledByHand, true) + } h.end(r) }() answer, err := kill(context.Background(), json.RawMessage(`{"id":"build-1"}`)) @@ -78,14 +82,15 @@ func TestKillEndsTheBuildRunningHereAndNoOther(t *testing.T) { if building.Err() == nil { t.Fatal("the build's context was not cancelled") } - if !h.wasKilled(r) { + if !r.killed { t.Error("the build is not marked killed, so it would be announced as an ordinary failure") } said := answer.(map[string]any)["said"].(string) - if !strings.Contains(said, "1 container(s)") || !strings.Contains(said, link.KilledByHand) { + if !strings.Contains(said, "2 container(s)") || !strings.Contains(said, "announced as failed") || !strings.Contains(said, link.KilledByHand) { t.Errorf("kill said %q", said) } - if len(removed) != 2 || !strings.Contains(removed[0], "label=mesh.build=build-1") { + // Removed at the kill and again once the build had ended: a container made in between is caught. + if len(removed) != 4 || !strings.Contains(removed[0], "label=mesh.build=build-1") || !strings.Contains(removed[2], "label=mesh.build=build-1") { t.Errorf("removed %v", removed) } if c := h.current(); c.Running != nil { @@ -101,7 +106,7 @@ func TestAHolderAnnouncesItsMachinesVerbsForTheConsole(t *testing.T) { } 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" || + if kill.Subject != "mesh.seat.node-build-agent.tool.kill.ace" || kill.QueueGroup != "seat.node-build-agent" || 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) @@ -111,3 +116,33 @@ func TestAHolderAnnouncesItsMachinesVerbsForTheConsole(t *testing.T) { t.Fatalf("the answer is %d bytes (%v)", len(body), err) } } + +// A build that ended on its own as the kill arrived says what it did: the kill is refused, and its +// own error is its outcome. One whose outcome could not be announced is not said to be settled. +func TestAKillArrivingAfterTheBuildEndedIsRefusedAndAnUnannouncedKillSaysSo(t *testing.T) { + h := newHolder("ace", link.TheBuildMachine, t.TempDir(), nil) + h.remove = func(context.Context, string, string, ...string) (string, error) { return "", nil } + _, cancel := context.WithCancel(context.Background()) + r := h.begin(link.BuildRequest{ID: "build-1"}, cancel) + if killed := h.returned(r); killed { + t.Fatal("a build nobody killed reads as killed") + } + if _, err := h.kill("build-1"); err == nil || !strings.Contains(err.Error(), "ended on its own") { + t.Fatalf("killed a build that had ended: %v", err) + } + h.end(r) + + building, cancel := context.WithCancel(context.Background()) + r = h.begin(link.BuildRequest{ID: "build-2"}, cancel) + go func() { + <-building.Done() + if h.returned(r) { + h.said(r, link.KilledByHand, false) + } + h.end(r) + }() + said, err := h.kill("build-2") + if err != nil || strings.Contains(said, "announced as failed") || !strings.Contains(said, "could not be announced") { + t.Fatalf("kill said %q (%v)", said, err) + } +} diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index 24fdcd8..0d6bdf1 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -247,8 +247,12 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri request.Repository, request.Path, request.Ref, workspace, request.Held, npmrc, forgeFrom(), say, request.Seats) } - // Only a build the kill ended: one that finished in the moment the kill arrived built, and says so. - killed := err != nil && mine != nil && h.wasKilled(mine) + // Only a build the kill ended: the kill came before its work did. One that finished — built, or + // failed on its own — in the moment the kill arrived says what it did, and the kill is refused. + killed := false + if mine != nil { + killed = h.returned(mine) && err != nil + } if killed { // **Killed by hand is the outcome, whatever the build was doing** (novox/hq ADR 0219): the // error it ended with is the kill's consequence, not a fault of the source. @@ -289,7 +293,11 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri defer stop() announcing = fresh } - if err := work.Announce(announcing, result); err != nil { + announceErr := work.Announce(announcing, result) + if mine != nil { + h.said(mine, result.Failed, announceErr == nil) + } + if err := announceErr; err != nil { // Said, not fatal: the build happened. A build reported as failed because announcing it // failed is a lie about work that was done — and the request stays unsettled below only if // nothing was said at all, so another machine can try. diff --git a/cmd/mesh-controller/plan_retry.go b/cmd/mesh-controller/plan_retry.go index 4165250..4bbf851 100644 --- a/cmd/mesh-controller/plan_retry.go +++ b/cmd/mesh-controller/plan_retry.go @@ -50,10 +50,20 @@ func pausedWaiting(p inventory.Plan, pause pauseView, now time.Time) (string, bo 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 + waited := now.Sub(*asked) + said := fmt.Sprintf("waiting: the build seat is paused on %s (asked %s ago)", + strings.Join(pause.Nodes, ", "), waited.Round(time.Second)) + // Not late — a pause is a decision — but a pause forgotten is a plan that never moves, so one held + // longer than a day is named. + if waited > pausedTooLong { + said += " — PAUSED OVER A DAY: `resume` takes builds again" + } + return said, true } +// pausedTooLong is how long a plan may wait on a paused build seat before it says the pause is long. +const pausedTooLong = 24 * time.Hour + // 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 { diff --git a/cmd/mesh-controller/queue.go b/cmd/mesh-controller/queue.go index 018ff2e..949ec9f 100644 --- a/cmd/mesh-controller/queue.go +++ b/cmd/mesh-controller/queue.go @@ -6,6 +6,7 @@ import ( "errors" "flag" "fmt" + "os" "sort" "strings" "time" @@ -29,8 +30,12 @@ import ( // 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 +// holderAsks is how long a holder's verb is waited for; killAnswer how long its kill is — longer +// than the holder's worst case (its two passes removing containers and its wait between, 50s). +const ( + holderAsks = 15 * time.Second + killAnswer = 75 * time.Second +) // dialTheBus opens the controller's own connection, for a command that reads or changes the queue. func dialTheBus() (*broker.JetStream, error) { @@ -100,8 +105,11 @@ func queueText(q link.Queue, now time.Time) string { fmt.Fprintln(&b, "\nin flight:") for _, a := range running { where := "taken, not yet said where" - if a.On != "" { + switch { + case a.On != "": where = fmt.Sprintf("on %s for %s", a.On, now.Sub(a.Started).Round(time.Second)) + case a.Was != "": + where = fmt.Sprintf("handed back by %s, to be handed out again", a.Was) } fmt.Fprintf(&b, " %-26s %s — %s, %s (seq %d)\n", a.ID, what(a), where, ago(a.AskedAt), a.Seq) } @@ -163,21 +171,19 @@ func cancelCommand(ctx context.Context, args []string) error { // 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; +// 1. an ask a machine says it is building 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. +// taken in the moment of cancelling, and the cancel is withdrawn and refused — the holder either +// read the mark first and ends it as cancelled, or is building it; +// 4. for an ask the worker counts handed out that no machine is building — taken and not yet said, +// handed back after a restart, or past its deliveries while nobody pulls — the holder's own look +// at the set is waited out, and a start heard since withdraws and refuses the cancel: a holder +// that took it before the mark is building it, and only `kill` ends that; +// 5. the message is deleted by its sequence; +// 6. 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) @@ -185,14 +191,31 @@ func cancelAsk(ctx context.Context, js *broker.JetStream, open *stores, seat str if err := link.MarkCancelled(ctx, js, seat, ask.ID); err != nil { return "", err } + withdraw := func() { + if err := link.UnmarkCancelled(ctx, js, seat, ask.ID); err != nil { + fmt.Fprintf(os.Stderr, "could not withdraw the cancel of %s: %v — a holder taking it ends it as cancelled\n", ask.ID, err) + } + } worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: seat, Accepts: []string{"build"}}) - if ask.State == link.AskWaiting { + switch ask.State { + case link.AskWaiting: if taken, err := takenSince(js, worker, ask.Seq); err != nil { + withdraw() 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) + withdraw() + return "", fmt.Errorf("%s was taken by a holder as it was cancelled, so the cancel is withdrawn. If the "+ + "holder read it first it ends the ask 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) + } + case link.AskInFlight: + if started, err := startedSince(ctx, js, seat, ask); err != nil { + withdraw() + return "", err + } else if started != "" { + withdraw() + return "", fmt.Errorf("%s has started on %s since it was read, so the cancel is withdrawn: `kill %s` "+ + "ends it there", ask.ID, started, ask.ID) } } if err := js.Context().DeleteMsg(worker.Stream, ask.Seq); err != nil && !errors.Is(err, nats.ErrMsgNotFound) && @@ -205,6 +228,38 @@ func cancelAsk(ctx context.Context, js *broker.JetStream, open *stores, seat str "that asked for it fails with that", ask.ID, ask.Repository, ask.State, seat, link.CancelledByHand), nil } +// holderLooks is how long a cancel waits for a holder that took the ask before the mark to say it +// started: longer than a holder's look at the cancelled set (3s) and its start that follows. +var holderLooks = 5 * time.Second + +// startedSince waits out a holder's look at the cancelled set and says the machine that started the +// ask since it was read, or nothing. A start it ended as cancelled comes with its outcome and is not +// a start of a build. +func startedSince(ctx context.Context, js *broker.JetStream, seat string, ask link.QueuedAsk) (string, error) { + select { + case <-ctx.Done(): + return "", ctx.Err() + case <-time.After(holderLooks): + } + since := time.Now().Add(-7 * 24 * time.Hour) + if !ask.AskedAt.IsZero() { + since = ask.AskedAt.Add(-time.Minute) + } + heard, err := link.ReadBuildEvents(ctx, js, seat, since) + if err != nil { + return "", err + } + s, ok := heard.Started[ask.ID] + if !ok || heard.Outcomes[ask.ID] { + return "", nil + } + at, _ := time.Parse(time.RFC3339Nano, s.At) + if ask.Started.IsZero() || at.After(ask.Started) { + return s.On, nil + } + return "", 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) @@ -314,6 +369,19 @@ func rebuildCommand(ctx context.Context, args []string) error { } else if found { source, module = sourceOfBuild(b, entries) path, ref = b.Path, b.Ref + // **A build at a commit is rebuilt at what its module follows now** (novox/hq ADR 0219, issue + // 219): asked now, a rebuild of an old commit would be the newest ask of the module and roll + // that commit out over everything since. Building that commit again is `replay`, which says + // what it would do and refuses to register it while anything newer is asked or registered. + if e, known := entryNamed(entries, module); known && ref != "" && followedBranch(ref) == "" { + ref = followedBranch(e.Source.Ref) + follows := ref + if follows == "" { + follows = "the repository's default branch" + } + fmt.Printf("%s was built at commit %s; a rebuild asks what %s follows now, %s — `replay %s` "+ + "builds that commit\n", b.ID, short(b.Ref), module, follows, b.ID) + } } 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} @@ -401,7 +469,25 @@ func replayCommand(ctx context.Context, args []string) error { return err } } - if err := replayRefusal(b, module, history, *register, *older); err != nil { + // What is asked of the module and not yet answered, when the replay would be registered. + var outstanding []string + if *register { + js, err := dialTheBus() + if err != nil { + return err + } + q, err := link.ReadQueue(ctx, js, buildSeatHeld(ctx)) + js.Close() + if err != nil { + return err + } + plans, err := inv.OpenPlans(ctx) + if err != nil { + return err + } + outstanding = outstandingFor(module, b.Repository, b.Path, q, plans) + } + if err := replayRefusal(b, module, history, *register, *older, outstanding); err != nil { return err } id, err := buildOneAsked(ctx, source, b.Path, b.Commit, 0, !*register) @@ -420,7 +506,13 @@ func replayCommand(ctx context.Context, args []string) error { // 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 { +// +// **And refused while anything newer is outstanding**, --older or not: an ask of the module in the +// queue, or an open plan holding it unbuilt. Asked now, the replay would be newer than those asks, +// and their builds — made from newer source — would be recorded and never registered (issue 219 +// orders by ask), the plan waiting on them sending the older commit instead. +func replayRefusal(b inventory.Build, module string, history []inventory.Build, register, older bool, + outstanding []string) 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) @@ -428,6 +520,11 @@ func replayRefusal(b inventory.Build, module string, history []inventory.Build, if older && !register { return errors.New("--older only says what --register may do; a dry run registers nothing") } + if register && len(outstanding) > 0 { + return fmt.Errorf("%s is asked and not yet answered — %s — and a registered replay asked now would "+ + "replace what those build (novox/hq issue 219). Wait for them, or `replay %s` without --register to look", + orNone(module), strings.Join(outstanding, "; "), b.ID) + } if !register || older { return nil } @@ -487,7 +584,7 @@ func killCommand(ctx context.Context, args []string) error { 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) + answer, err := link.AskSeatTool(ctx, js.Conn(), seat, "kill", started.On, map[string]string{"id": id}, killAnswer) if err != nil { return err } @@ -581,3 +678,32 @@ func holdersAmong(entries []inventory.Entry, seat string) []string { sort.Strings(out) return out } + +// outstandingFor is everything asked of a module and not yet answered: its asks in the build queue, +// by repository and path, and every open plan holding it not yet built. +func outstandingFor(module, repository, path string, q link.Queue, plans []inventory.Plan) []string { + var out []string + for _, a := range q.Asks { + if repositoryMatches(a.Repository, repository) && a.Path == path { + out = append(out, fmt.Sprintf("%s %s in the queue", a.ID, a.State)) + } + } + if module == "" { + return out + } + for _, p := range plans { + if !p.Open() { + continue + } + for _, tier := range p.Tiers { + for _, m := range tier { + if m == module { + if st := p.Modules[m]; st == nil || st.State != "built" { + out = append(out, fmt.Sprintf("%s holds it not yet built", p.ID)) + } + } + } + } + } + return out +} diff --git a/cmd/mesh-controller/queue_test.go b/cmd/mesh-controller/queue_test.go index bc23f0e..d78bb67 100644 --- a/cmd/mesh-controller/queue_test.go +++ b/cmd/mesh-controller/queue_test.go @@ -212,31 +212,53 @@ func TestReplayRefusesToRollAnOlderCommitOutUnlessToldTo(t *testing.T) { 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 { + if err := replayRefusal(old, "a", history, false, false, nil); err != nil { t.Errorf("a dry run was refused: %v", err) } - err := replayRefusal(old, "a", history, true, false) + err := replayRefusal(old, "a", history, true, false, nil) 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 { + if err := replayRefusal(old, "a", history, true, true, nil); err != nil { t.Errorf("--register --older was refused: %v", err) } - if err := replayRefusal(newer, "a", history, true, false); err != nil { + if err := replayRefusal(newer, "a", history, true, false, nil); 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 { + if err := replayRefusal(old, "a", []inventory.Build{broken, same, old}, true, false, nil); 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 { + if err := replayRefusal(inventory.Build{ID: "build-x", Failed: "clone"}, "", nil, false, false, nil); err == nil { t.Error("a build that recorded no commit was replayed") } - if err := replayRefusal(old, "a", history, false, true); err == nil { + if err := replayRefusal(old, "a", history, false, true, nil); err == nil { t.Error("--older without --register was taken") } + // **Nothing newer outstanding**, --older or not: an ask of the module in the queue, or an open plan + // holding it unbuilt, would be replaced by a replay asked now (issue 219). + q := link.Queue{Asks: []link.QueuedAsk{ + {ID: "build-queued", Repository: "https://forge.example/novox/a.git", State: link.AskWaiting}, + {ID: "build-elsewhere", Repository: "https://forge.example/novox/b.git", State: link.AskWaiting}, + {ID: "build-other-path", Repository: "https://forge.example/novox/a.git", Path: "sub", State: link.AskWaiting}, + }} + plans := []inventory.Plan{ + {ID: "plan-holds", State: inventory.PlanBuilding, Tiers: [][]string{{"x"}, {"a"}}, Modules: map[string]*inventory.PlanModule{}}, + {ID: "plan-built", State: inventory.PlanRolling, Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{"a": {State: "built"}}}, + {ID: "plan-failed", State: inventory.PlanFailed, Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{}}, + } + outstanding := outstandingFor("a", "novox/a", "", q, plans) + if len(outstanding) != 2 || !strings.Contains(outstanding[0], "build-queued") || !strings.Contains(outstanding[1], "plan-holds") { + t.Fatalf("outstanding: %v", outstanding) + } + if err := replayRefusal(old, "a", history, true, true, outstanding); err == nil || !strings.Contains(err.Error(), "build-queued") { + t.Errorf("registered over an outstanding ask: %v", err) + } + if err := replayRefusal(old, "a", history, false, false, outstanding); err != nil { + t.Errorf("a dry run was refused for what is outstanding: %v", err) + } 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) } @@ -266,6 +288,13 @@ func TestAPlanWaitingOnAPausedSeatSaysSoAndIsNotLate(t *testing.T) { t.Errorf("a plan waiting on a paused seat counted late") } + longAgo := now.Add(-25 * time.Hour) + forgotten := p + forgotten.Modules = map[string]*inventory.PlanModule{"a": {State: "asked", AskedAt: &longAgo}} + if line := planLineWith(forgotten, now, all); !strings.Contains(line, "PAUSED OVER A DAY") || strings.Contains(line, "LATE") { + t.Errorf("a plan paused over a day reads %q", line) + } + 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) @@ -644,3 +673,100 @@ func TestAPlanStoppedAtItsFirstMachineIsRetried(t *testing.T) { t.Fatalf("asked %v", *asked) } } + +// A plan module asked under an id is settled by that id's outcome alone: a replay or a rebuild beside +// the plan, asked later, never answers it (novox/hq ADR 0219). +func TestAPlanIsAnsweredOnlyByTheBuildItAskedFor(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + asked := time.Now().UTC().Add(-time.Minute) + plan := inventory.Plan{ID: "plan-own", Repository: "novox/a", Commit: "c0ffee", Created: asked, + State: inventory.PlanBuilding, Tiers: [][]string{{"a"}}, + Modules: map[string]*inventory.PlanModule{"a": {State: "asked", AskedAt: &asked, Build: "build-own"}}} + if err := open.inventory.SavePlan(ctx, plan); err != nil { + t.Fatal(err) + } + planBuilt(ctx, open, "a", "0ldc0mm1t", "", time.Now().UTC(), "build-replay") + p, _ := open.inventory.PlanByID(ctx, plan.ID) + if p.Modules["a"].State != "asked" { + t.Fatalf("a replay asked after the plan settled it: %+v", p.Modules["a"]) + } + // From the records too: a later build of the module recorded is not the plan's. + recorded := map[string][]inventory.Build{"a": {{ID: "build-replay", Commit: "0ldc0mm1t", Asked: time.Now(), At: time.Now()}}} + if settleFromRecords(&p, p.Tiers[0], recorded, nil) { + t.Fatalf("the records settled it with another build: %+v", p.Modules["a"]) + } + planBuilt(ctx, open, "a", "c0ffee", "", asked, "build-own") + if p, _ = open.inventory.PlanByID(ctx, plan.ID); p.Modules["a"].State != "built" || p.Modules["a"].Commit != "c0ffee" { + t.Fatalf("its own build did not settle it: %+v", p.Modules["a"]) + } +} + +// rebuild of a build made at a commit asks what the module follows now, never the commit. +func TestARebuildOfACommitAsksWhatTheModuleFollows(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + asked := asksRecorded(t) + twoTiers(t, open) + var refs []string + was := askABuild + askABuild = func(c context.Context, source buildSource, path, ref string) (string, error) { + refs = append(refs, ref) + return was(c, source, path, ref) + } + if err := open.inventory.RecordBuild(ctx, inventory.Build{ID: "build-at-commit", Repository: "novox/a", + Ref: "0123456789abcdef0123456789abcdef01234567", Module: "a", Commit: "0123456789abcdef0123456789abcdef01234567"}); err != nil { + t.Fatal(err) + } + if err := rebuildCommand(ctx, []string{"build-at-commit"}); err != nil { + t.Fatal(err) + } + if len(refs) != 1 || refs[0] != "main" || len(*asked) != 1 { + t.Fatalf("asked %v at %v", *asked, refs) + } +} + +// An ask the worker counts out that no machine said it started is cancelled only if no start comes +// while a holder looks: one that does withdraws the cancel, and kill is what ends it. +func TestCancelOfAnAskNobodySaidIsWithdrawnWhenItStarts(t *testing.T) { + js, ask := aBuildQueue(t) + open := aMesh(t) + ctx := t.Context() + was := holderLooks + holderLooks = 300 * time.Millisecond + t.Cleanup(func() { holderLooks = was }) + 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) + } + 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 != "" { + t.Fatalf("the taken ask reads %+v", a) + } + // The holder that took it says it started, while the cancel waits. + go func() { + time.Sleep(100 * time.Millisecond) + started, _ := json.Marshal(link.BuildStart{ID: id, On: "ace", At: time.Now().UTC().Format(time.RFC3339Nano)}) + _, _ = js.Context().Publish(link.BuildStarted(), started) + }() + if _, err := cancelAsk(ctx, js, open, link.TheBuildMachine, a); err == nil || !strings.Contains(err.Error(), "started on ace") { + t.Fatalf("a build that started was cancelled: %v", err) + } + if cancelled, _ := link.IsCancelled(js.Conn(), link.TheBuildMachine, id); cancelled { + t.Error("the withdrawn cancel is still marked") + } + if _, found, _ := open.inventory.BuildByID(ctx, id); found { + t.Error("a withdrawn cancel recorded an outcome") + } +} diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 9e9d79b..0f527fa 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -416,7 +416,15 @@ func planBuilt(ctx context.Context, open *stores, module, commit, failed string, } // **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) { + // **A module asked under an id is answered by that id's outcome and no other** (novox/hq ADR + // 0219): a replay, a rebuild beside the plan, or an older ask finishing late is somebody else's + // build, made from other source, and settling the plan with it would send that. A plan from + // before ids were kept is matched as it was: by when the build was asked (issue 219). + if state.Build != "" { + if state.Build != id { + continue + } + } else if askedBefore(asked, state.AskedAt) { continue } if failed != "" { @@ -1192,6 +1200,10 @@ func settleFromRecords(p *inventory.Plan, tier []string, recorded map[string][]i } var outcome *inventory.Build for i := range recorded[m] { + // Asked under an id, only that id's record answers (ADR 0219), and it is looked up below. + if s.Build != "" { + break + } b := recorded[m][i] if b.At.Before(*s.AskedAt) { break diff --git a/internal/broker/cancelled_test.go b/internal/broker/cancelled_test.go index f9e63f7..5bf4a0a 100644 --- a/internal/broker/cancelled_test.go +++ b/internal/broker/cancelled_test.go @@ -16,7 +16,9 @@ func TestTheBuildQueueIsControlledWithTheGrantsItNeedsAndNoMore(t *testing.T) { } has(t, holder.Publish, "$JS.API.DIRECT.GET.KV_SEAT_NODE_BUILD_AGENT_cancelled.$KV.SEAT_NODE_BUILD_AGENT_cancelled.>") hasNot(t, holder.Publish, "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>") - has(t, holder.Publish, "mesh.seat.node-build-agent.event.paused.*") + has(t, holder.Publish, "mesh.seat.node-build-agent.event.paused.ace") + hasNot(t, holder.Publish, "mesh.seat.node-build-agent.event.paused.*") + has(t, holder.Publish, "mesh.seat.node-build-agent.event.log.*") has(t, holder.Subscribe, "mesh.seat.node-build-agent.tool.kill.ace") hasNot(t, holder.Subscribe, "mesh.seat.node-build-agent.tool.kill.g14") diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 9264429..2cca06d 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -124,6 +124,10 @@ type Principal struct { // goes with the retired seat row. var seatsTheControllerAsks = []string{"node-build-agent", "mesh-build-machine"} +// perMachineEvents are a node-scoped seat's events about the holder itself, whose last token is the +// holder's machine (novox/hq ADR 0219): `paused.`, the build agent saying whether it takes work. +var perMachineEvents = map[string]bool{"paused.*": true} + // enrolmentPrefix is the space every enrolling node's user and inbox live under, so the one place the // controller may answer an enrolment is derived from the same constant the user is named from. const enrolmentPrefix = "enrol" @@ -421,6 +425,13 @@ func PermissionsFor(p Principal) (Permissions, error) { sub = append(sub, seatSubject(s, "accept", a)) } for _, e := range s.Emits { + // **A machine says its own state and no other's** (novox/hq ADR 0219): on a node-scoped + // seat, an event about the holder itself carries the machine as its last token, and + // each holder is granted its own machine's alone. + if s.Scope == "node" && p.Node != "" && perMachineEvents[e] { + pub = append(pub, seatSubject(s, "event", strings.TrimSuffix(e, "*")+p.Node)) + continue + } pub = append(pub, seatSubject(s, "event", e)) } for _, t := range s.Serves { diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index dfb7b1b..f5d1c4c 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -175,8 +175,12 @@ func MachineOverNATSWith(js *broker.JetStream, on, seat string, opts MachineOpti return &natsMachine{js: js, on: on, seat: seat, opts: opts} } -// pausedPoll is how often a paused holder looks again whether it was resumed. -const pausedPoll = 2 * time.Second +// pausedPoll is how often a paused holder looks again whether it was resumed, and pausedFetch how +// long one pull of a holder that can be paused waits. +const ( + pausedPoll = 2 * time.Second + pausedFetch = 5 * time.Second +) func (m *natsMachine) Close() { if m.sub != nil { @@ -235,7 +239,14 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) } // One, and wait a while for it; an empty queue is a timeout, which is the normal state of a // machine with nothing to build, and is asked again. - fetched, err := sub.Fetch(1, nats.Context(ctx)) + // Asked for a few seconds at a time when the holder can be paused, so a pause reaches a pull + // already waiting within that, rather than when the client's own wait runs out. + asking, endAsking := ctx, func() {} + if m.opts.Paused != nil { + asking, endAsking = context.WithTimeout(ctx, pausedFetch) + } + fetched, err := sub.Fetch(1, nats.Context(asking)) + endAsking() switch { case ctx.Err() != nil: // Ours ended: the machine is being stopped. @@ -257,6 +268,13 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build)) return fmt.Errorf("the bus stopped delivering build work: %w", err) } for _, msg := range fetched { + // **Paused while the pull was answered** (ADR 0219): "takes no new build" holds even for + // the one that arrived in that moment. Handed back at once for another holder — counted as + // a delivery, which a pause landing exactly then costs and nothing else does. + if m.opts.Paused != nil && m.opts.Paused() { + _ = msg.Nak() + continue + } var request BuildRequest if err := json.Unmarshal(msg.Data, &request); err != nil { // Unreadable: terminated rather than retried, because the next attempt reads the same diff --git a/internal/link/builds_nats_test.go b/internal/link/builds_nats_test.go index b0feaa2..fcbe211 100644 --- a/internal/link/builds_nats_test.go +++ b/internal/link/builds_nats_test.go @@ -303,8 +303,9 @@ func TestNatsTwoMachinesShareTheWorkAndNeitherIsHandedMoreThanItCanTake(t *testi case <-time.After(2 * time.Second): } // One finishes, and only then is the third taken — by that machine, the one that is free. - close(release["anchor"]) - release["anchor"] = make(chan struct{}) + // Released by a send, never by replacing the channel: the map is read by both machines' builds + // while this runs, and writing it raced them. + release["anchor"] <- struct{}{} select { case got := <-took: if got.machine != "anchor" { @@ -318,7 +319,7 @@ func TestNatsTwoMachinesShareTheWorkAndNeitherIsHandedMoreThanItCanTake(t *testi // an explicit hand-back, because the real wait is a minute. Here: the laptop goes, anchor // finishes, and with nothing queued nothing more is taken by the machine that is left. machines["laptop"].Close() - close(release["anchor"]) + release["anchor"] <- struct{}{} select { case got := <-took: t.Fatalf("%s took %s; the queue should be empty", got.machine, got.id) diff --git a/internal/link/queue.go b/internal/link/queue.go index 2855129..254dc52 100644 --- a/internal/link/queue.go +++ b/internal/link/queue.go @@ -62,8 +62,11 @@ type QueuedAsk struct { Ref string `json:"ref,omitempty"` AskedAt time.Time `json:"asked-at,omitempty"` State string `json:"state"` - // On and Started are the holder that took an in-flight ask and when, from its started event. + // On and Started are the holder building an in-flight ask and when it started, from its started + // event. Was is the holder that started an in-flight ask and is no longer building it — handed + // back, to be handed out again — with Started its start there. On string `json:"on,omitempty"` + Was string `json:"was-on,omitempty"` Started time.Time `json:"started,omitempty"` } @@ -109,19 +112,24 @@ type BuildHeard struct { // ClassifyQueue says which state each ask is in. Pure: the stream's asks in sequence order, what the // worker says, and what the seat's events said. // -// **In flight is what a machine is building now.** A holder builds one ask at a time (ADR 0190), so -// what a machine is building is its latest start, if no outcome has been heard for it. An ask -// behind the worker that is some machine's current build is in flight; at most AckPending of them -// are, newest start first — a machine that died mid-build has a latest start forever, and the -// worker's own count is what bounds that. An ask behind the worker that no machine has said it -// started is in flight only while the worker counts more pending than starts explain (taken, and -// not yet said); otherwise it is dead. +// **In flight is what the worker counts handed out and unsettled**, at most AckPending of the asks +// behind it, given out in this order: +// +// 1. what a machine is building now — its latest start, no outcome heard (a holder builds one ask +// at a time, ADR 0190), newest start first; +// 2. what a machine started and is no longer building — a holder restarted mid-build, the ask +// handed back and waiting to be handed out again while it has deliveries left — newest start +// first, naming the machine it was last on; +// 3. what was taken and not yet said started. +// +// Only what is left after that is dead: so nothing behind the worker is dead while there are no more +// of them than the worker counts pending. A machine that died mid-build keeps a latest start for +// ever; the worker's count is what says the ask was handed out as often as it may be. // // **The worker's count is the server's, and it lags in one case**: an ask handed out its last time // and handed back is still counted pending until some holder pulls again, when the server finds it // past its deliveries and lets it go. Holders pull whenever they are idle, so this is a moment — // except while every holder is paused, when such an ask reads as in flight with no machine named. -// Cancel takes an ask in flight that no machine said it started, for exactly that reason. func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard BuildHeard) []QueuedAsk { // Each machine's latest start. latest := map[string]BuildStart{} @@ -140,7 +148,7 @@ func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard Bui out := make([]QueuedAsk, len(asks)) copy(out, asks) - var running, unsaid []int + var running, handedBack, unsaid []int for i := range out { a := &out[i] if a.Seq > delivered { @@ -153,28 +161,29 @@ func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard Bui running = append(running, i) continue } - if _, ever := heard.Started[a.ID]; !ever { - unsaid = append(unsaid, i) - } - } - sort.SliceStable(running, func(x, y int) bool { return out[running[x]].Started.After(out[running[y]].Started) }) - slots := ackPending - for _, i := range running { - if slots == 0 { - // A machine's latest start the worker no longer counts: it died with the ask, and the - // ask was handed out as often as it may be. - out[i].On, out[i].Started = "", time.Time{} + if s, ever := heard.Started[a.ID]; ever { + a.Was, a.Started = s.On, startedAt(s) + handedBack = append(handedBack, i) continue } - out[i].State = AskInFlight - slots-- + unsaid = append(unsaid, i) } - for _, i := range unsaid { - if slots == 0 { - break + newestFirst := func(ix []int) { + sort.SliceStable(ix, func(x, y int) bool { return out[ix[x]].Started.After(out[ix[y]].Started) }) + } + newestFirst(running) + newestFirst(handedBack) + slots := ackPending + for _, group := range [][]int{running, handedBack, unsaid} { + for _, i := range group { + if slots == 0 { + // Past what the worker counts: handed out as often as it may be. + out[i].On, out[i].Started = "", time.Time{} + continue + } + out[i].State = AskInFlight + slots-- } - out[i].State = AskInFlight - slots-- } return out } @@ -336,6 +345,23 @@ func MarkCancelled(ctx context.Context, js *broker.JetStream, seat, id string) e return err } +// UnmarkCancelled takes an ask's id out of the cancelled set: a cancel withdrawn, because the ask +// was taken as it was being cancelled. +func UnmarkCancelled(ctx context.Context, js *broker.JetStream, seat, id string) error { + if !safeKey.MatchString(id) { + return nil + } + api, err := jetstream.New(js.Conn()) + if err != nil { + return err + } + kv, err := api.KeyValue(ctx, broker.CancelledSetName(seat)) + if err != nil { + return err + } + return kv.Delete(ctx, id) +} + // IsCancelled asks the seat's cancelled set whether an ask was cancelled: one direct read of one key, // which is all a holder is granted of it. func IsCancelled(conn *nats.Conn, seat, id string) (bool, error) { diff --git a/internal/link/queue_test.go b/internal/link/queue_test.go index 5da7672..4ff3a74 100644 --- a/internal/link/queue_test.go +++ b/internal/link/queue_test.go @@ -3,6 +3,7 @@ package link import ( "context" "encoding/json" + "sync" "testing" "time" @@ -35,9 +36,9 @@ func TestTheQueueTellsWaitingInFlightAndDeadApart(t *testing.T) { }, Outcomes: map[string]bool{"done": true}, } - got := ClassifyQueue(asks, 4, 3, heard) + got := ClassifyQueue(asks, 4, 2, heard) want := map[string]string{"dead-1": AskDead, "running-g": AskInFlight, "running-a": AskInFlight, - "unsaid": AskInFlight, "waiting-1": AskWaiting, "waiting-2": AskWaiting} + "unsaid": AskDead, "waiting-1": AskWaiting, "waiting-2": AskWaiting} for _, a := range got { if a.State != want[a.ID] { t.Errorf("%s is %s, want %s", a.ID, a.State, want[a.ID]) @@ -48,16 +49,45 @@ func TestTheQueueTellsWaitingInFlightAndDeadApart(t *testing.T) { t.Errorf("an ask in flight does not say where: %+v", a) } - // **The worker's count bounds it**: with two pending, the machine whose start is oldest died with - // its ask, and the ask nobody said is not in flight either. - got = ClassifyQueue(asks, 4, 2, heard) + // **The worker's count bounds it**: with one pending, the machine whose start is oldest died + // with its ask. + got = ClassifyQueue(asks, 4, 1, heard) q = Queue{Asks: got} - if a, _ := q.Find("running-g"); a.State != AskInFlight { - t.Errorf("running-g is %s", a.State) + if a, _ := q.Find("running-a"); a.State != AskInFlight { + t.Errorf("running-a is %s", a.State) + } + if a, _ := q.Find("running-g"); a.State != AskDead || a.On != "" { + t.Errorf("past the worker's count running-g is %s on %q", a.State, a.On) + } + // **Handed back after a holder restarted mid-build**: started on g14, g14 then started something + // else, and the worker still counts it — it waits to be handed out again, and is not dead. With + // no more asks behind the worker than it counts pending, none is dead. + restarted := []QueuedAsk{{Seq: 1, ID: "dead-1"}, {Seq: 2, ID: "running-g"}} + for _, a := range ClassifyQueue(restarted, 2, 2, heard) { + if a.State != AskInFlight { + t.Errorf("%s is %s with two behind the worker and two pending", a.ID, a.State) + } + if a.ID == "dead-1" && (a.On != "" || a.Was != "g14") { + t.Errorf("an ask handed back reads on %q, was on %q", a.On, a.Was) + } + } + // The handed-back one takes a leftover slot before an ask nobody said started. + got = ClassifyQueue(asks, 4, 4, heard) + q = Queue{Asks: got} + for _, id := range []string{"dead-1", "running-g", "running-a", "unsaid"} { + if a, _ := q.Find(id); a.State != AskInFlight { + t.Errorf("with four pending %s is %s", id, a.State) + } + } + got = ClassifyQueue(asks, 4, 3, heard) + q = Queue{Asks: got} + if a, _ := q.Find("dead-1"); a.State != AskInFlight { + t.Errorf("the handed-back ask lost its slot to one nobody said: %s", a.State) } if a, _ := q.Find("unsaid"); a.State != AskDead { - t.Errorf("an ask no start explains, past the worker's count, is %s", a.State) + t.Errorf("unsaid is %s", a.State) } + // None pending: everything behind the worker is dead, and names no machine. for _, a := range ClassifyQueue(asks, 4, 0, heard) { if a.Seq <= 4 && (a.State != AskDead || a.On != "") { @@ -265,3 +295,42 @@ func TestNatsAWaitingAskerHearsItsAskCancelledOrKilled(t *testing.T) { t.Fatalf("the waiter heard %+v (%v)", result, err) } } + +// Paused in the moment a pull was answered: the ask is handed back, not built (ADR 0219). +func TestNatsAHolderPausedAsItsPullWasAnsweredHandsTheAskBack(t *testing.T) { + js := aBusWithACancelledSet(t) + ctx, stop := context.WithCancel(context.Background()) + defer stop() + // Not paused when it asks; paused by the time the answer is read. + var looks int + var mu sync.Mutex + paused := func() bool { + mu.Lock() + defer mu.Unlock() + looks++ + return looks > 1 + } + took := make(chan string, 1) + machine := MachineOverNATSWith(js, "anchor", TheBuildMachine, MachineOptions{Paused: paused}) + defer machine.Close() + go func() { + _ = machine.Take(ctx, func(ctx context.Context, work Build) { took <- work.Request().ID }) + }() + time.Sleep(200 * time.Millisecond) + body, _ := json.Marshal(BuildRequest{ID: "build-handed-back", Repository: "/r"}) + if _, err := js.Context().Publish(BuildWork(), body); err != nil { + t.Fatal(err) + } + select { + case id := <-took: + t.Fatalf("a holder paused as its pull was answered built %s", id) + case <-time.After(2 * time.Second): + } + q, err := ReadQueue(ctx, js, TheBuildMachine) + if err != nil { + t.Fatal(err) + } + if len(q.Asks) != 1 || q.Asks[0].ID != "build-handed-back" { + t.Fatalf("the ask is not back in the queue: %+v", q.Asks) + } +}