From 8adb7f1a051941e943cc962dd268ea7585709cb7 Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 16:45:43 +0200 Subject: [PATCH] Give each pending assignment its own condition, and keep a raised row until it clears Second review of #150: one key per machine and module let a newer failure be cleared in the tick that raised it, pruning could orphan an open condition, and a build no longer waited for read as never asked although it may still run. --- cmd/mesh-controller/build.go | 22 +- cmd/mesh-controller/pending.go | 70 +++++- cmd/mesh-controller/pending_test.go | 229 +++++++++++++++++- internal/inventory/fortest.go | 10 + ...0082-an-assignment-waits-for-its-build.sql | 9 +- internal/inventory/pending.go | 49 +++- internal/link/builds_nats.go | 13 +- 7 files changed, 367 insertions(+), 35 deletions(-) diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 28f5e7ec..ec5d2b2f 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -415,16 +415,22 @@ func recordAsked(ctx context.Context, r inventory.BuildRequest) { recordBuildRequest(ctx, open.inventory, r) } -// markNotAsked says a kept build request was not handed over, or not waited for. -func markNotAsked(ctx context.Context, id string, why error) { +// markWaitFailed says what became of a kept build request whose waited ask failed: never handed over, so +// nothing runs; or handed over and no longer waited for, so asked with its outcome unknown — still read as in +// flight until its outcome or its bound (novox/hq issue 325). +func markWaitFailed(ctx context.Context, id string, why error) { open, err := openStores(ctx) if err != nil { - fmt.Fprintf(os.Stderr, "build %s could not be marked as not asked: %v\n", id, err) + fmt.Fprintf(os.Stderr, "build %s could not be marked: %v\n", id, err) return } defer open.Close() - if err := open.inventory.MarkNotAsked(ctx, id, why.Error()); err != nil { - fmt.Fprintf(os.Stderr, "build %s could not be marked as not asked: %v\n", id, err) + mark := open.inventory.MarkOutcomeUnknown + if errors.Is(why, link.ErrNotHandedOver) { + mark = open.inventory.MarkNotAsked + } + if err := mark(ctx, id, why.Error()); err != nil { + fmt.Fprintf(os.Stderr, "build %s could not be marked: %v\n", id, err) } } @@ -518,15 +524,15 @@ func buildOneAsked(ctx context.Context, source buildSource, path, ref string, wa return request.ID, nil } - // Waited for: kept while it is waited for, since that can take minutes, and marked when the hand-over or - // the wait failed — an outcome heard later is still the last word. + // Waited for: kept while it is waited for, since that can take minutes, and marked when the hand-over + // failed (not asked) or the wait did (asked, outcome unknown) — an outcome heard later is the last word. if keep { recordAsked(ctx, asked) } result, err := ask.Submit(ctx, request, wait) if err != nil { if keep { - markNotAsked(ctx, request.ID, err) + markWaitFailed(ctx, request.ID, err) } return request.ID, err } diff --git a/cmd/mesh-controller/pending.go b/cmd/mesh-controller/pending.go index af8fc44c..3627eb68 100644 --- a/cmd/mesh-controller/pending.go +++ b/cmd/mesh-controller/pending.go @@ -281,8 +281,11 @@ func whyNotBuilt(r inventory.RequestOutcome, now time.Time) string { return fmt.Sprintf("%s; the merge that added it (%s) could not ask for its build: %s", where, short(r.Commit), firstLine(r.NotAsked)) case r.NotAsked != "": - return fmt.Sprintf("%s; build %s, asked by %s, was not handed over or not waited for: %s", where, r.ID, + return fmt.Sprintf("%s; build %s, asked by %s, was not handed over: %s", where, r.ID, askerOr(r.For), firstLine(r.NotAsked)) + case r.OutcomeUnknown != "": + return fmt.Sprintf("%s; build %s was asked at %s, its asker stopped waiting (%s), and no outcome has "+ + "been heard in %s", where, r.ID, clock(r.At), firstLine(r.OutcomeUnknown), now.Sub(r.At).Round(time.Minute)) default: return fmt.Sprintf("%s; build %s was asked at %s, and no outcome has been heard in %s", where, r.ID, clock(r.At), now.Sub(r.At).Round(time.Minute)) @@ -545,7 +548,7 @@ func settlePending(ctx context.Context, open *stores, now time.Time) []string { switch { case r == nil && now.Sub(p.Since) > buildRequestBound: why = fmt.Sprintf("build %s is no longer on record", p.Build) - case r != nil && !stillComing(*r, now) && (r.Heard || r.NotAsked != ""): + case r != nil && !stillComing(*r, now) && (r.Heard || r.NotAsked != "" || r.OutcomeUnknown != ""): why = whyNotBuilt(*r, now) case r != nil && !stillComing(*r, now): why = fmt.Sprintf("no outcome of build %s was heard within %s of asking", p.Build, buildRequestBound) @@ -576,9 +579,10 @@ func containsString(xs []string, x string) bool { return false } -// pendingObservation is the condition for a pending assignment that was not made. +// pendingObservation is the condition for a pending assignment that was not made: one per pending +// assignment, its row's id the last part of its id, so one row's condition is never another's. func pendingObservation(p inventory.PendingAssignment) conditions.Observation { - return conditions.Observation{Scope: conditions.ScopeMachine, ID: p.Node + "." + p.Module, + return conditions.Observation{Scope: conditions.ScopeMachine, ID: fmt.Sprintf("%s.%s.%d", p.Node, p.Module, p.ID), Token: kindPendingEnded, Kind: kindPendingEnded, Machine: p.Node, Severity: conditions.Warning, Resolver: conditions.ResolverOperator, Source: "pending assignments", Summary: p.Note} @@ -628,7 +632,7 @@ func raisePendingEnded(ctx context.Context, inv *inventory.Inventory) []string { said = append(said, fmt.Sprintf("the condition for %s on %s was cleared and not recorded: %v", p.Module, p.Node, err)) } } - return said + return append(said, clearOrphaned(ctx, keeper, inv)...) } // pendingAnswered says whether a pending assignment that was not made has been answered since, and how. @@ -647,14 +651,60 @@ func pendingAnswered(ctx context.Context, inv *inventory.Inventory, p inventory. if err != nil { return "", false, err } + // A newer pending assignment answers it only while it is open or once it was made: one that itself + // ended unmade is its own condition, and answers nothing. for _, r := range rows { - if r.ID != p.ID && r.Since.After(p.Since) { + if r.ID != p.ID && r.Since.After(p.Since) && (r.Open() || r.State == inventory.PendingApplied) { return "assigned again, pending on another build", true, nil } } return "", false, nil } +// clearOrphaned clears every open condition of an assignment not made whose pending assignment is no longer +// on record: its machine was removed, which takes its pending assignments with it. +func clearOrphaned(ctx context.Context, keeper *conditions.Keeper, inv *inventory.Inventory) []string { + open, err := keeper.Open(ctx) + if err != nil { + return []string{fmt.Sprintf("the open conditions could not be read to clear orphaned ones: %v", err)} + } + byID := map[int64]string{} + var ids []int64 + for _, c := range open { + if c.Kind != kindPendingEnded { + continue + } + // `....`: the row is the part before the kind. + parts := strings.Split(c.Key, ".") + if len(parts) < 2 { + continue + } + var id int64 + if _, err := fmt.Sscan(parts[len(parts)-2], &id); err != nil { + continue + } + byID[id] = c.Key + ids = append(ids, id) + } + if len(ids) == 0 { + return nil + } + known, err := inv.PendingKnown(ctx, ids) + if err != nil { + return []string{fmt.Sprintf("the pending assignments could not be read to clear orphaned conditions: %v", err)} + } + var said []string + for id, key := range byID { + if known[id] { + continue + } + if _, err := keeper.Clear(ctx, key, "its pending assignment is no longer on record: its machine was removed"); err != nil { + said = append(said, fmt.Sprintf("the condition %s could not be cleared: %v", key, err)) + } + } + return said +} + // withdrawPending is `unassign` of a module only pending on a machine, under the machine's hold: a waiting // pending assignment is withdrawn; one being made now is refused, never overridden; one that expired or was // refused is taken back, which answers its condition. False when there is none. @@ -663,6 +713,7 @@ func withdrawPending(ctx context.Context, inv *inventory.Inventory, node, module if err != nil { return "", false, err } + var takenBack []string for _, p := range rows { switch p.State { case inventory.PendingApplying: @@ -687,9 +738,12 @@ func withdrawPending(ctx context.Context, inv *inventory.Inventory, node, module if err := inv.MarkPending(ctx, p.ID, "acknowledged"); err != nil { return "", false, err } - return fmt.Sprintf("the pending assignment of %s to %s, which %s, is taken back; its condition clears", - module, node, p.State), true, nil + takenBack = append(takenBack, fmt.Sprintf("build %s (%s)", p.Build, p.State)) } } + if len(takenBack) > 0 { + return fmt.Sprintf("the pending assignment(s) of %s to %s that were not made are taken back: %s; their "+ + "conditions clear", module, node, strings.Join(takenBack, ", ")), true, nil + } return "", false, nil } diff --git a/cmd/mesh-controller/pending_test.go b/cmd/mesh-controller/pending_test.go index acaf93ae..3efc8053 100644 --- a/cmd/mesh-controller/pending_test.go +++ b/cmd/mesh-controller/pending_test.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "fmt" "slices" "strings" "testing" @@ -15,6 +16,39 @@ import ( "github.com/novox/mesh-controller/internal/link" ) +// failedBuild settles the merge's build of modules/sensors as failed. +func failedBuild(t *testing.T, open *stores, id string) { + t.Helper() + if err := (builds{open.inventory, open}).Built(t.Context(), outcome(id, "modules/sensors", "3da80a4b00aa", + "", "boom")); err != nil { + t.Fatal(err) + } +} + +func conditionOpen(t *testing.T, p inventory.PendingAssignment) bool { + t.Helper() + _, found, err := conditionsFrom.Get(t.Context(), pendingObservation(p).Key()) + if err != nil { + t.Fatal(err) + } + return found +} + +func rowOf(t *testing.T, open *stores, id int64) inventory.PendingAssignment { + t.Helper() + rows, err := open.inventory.Pending(t.Context(), time.Unix(0, 0)) + if err != nil { + t.Fatal(err) + } + for _, r := range rows { + if r.ID == id { + return r + } + } + t.Fatalf("pending assignment %d is not on record", id) + return inventory.PendingAssignment{} +} + // novox/hq issue 325: an assignment of a module the catalogue does not hold says which case it is in — a // build in flight (kept pending), known and not built (said, with build), or unknown (refused, with the // closest names) — and a pending assignment is made when its build registers the module, or ends with why. @@ -370,7 +404,7 @@ func TestAFailedAskIsNotABuildInFlight(t *testing.T) { t.Fatal(err) } _, err = assign(ctx, open, "laptop", "gauges") - if err == nil || !strings.Contains(err.Error(), "build build-waited, asked by build, was not handed over or not waited for") { + if err == nil || !strings.Contains(err.Error(), "build build-waited, asked by build, was not handed over") { t.Fatalf("a build not handed over: %v", err) } // Heard first, the outcome stands: marking it afterwards changes nothing. @@ -602,7 +636,8 @@ func TestAPendingAssignmentNeverWaitsSilently(t *testing.T) { t.Fatalf("a module registered meanwhile: %+v", p) } - // Ended rows are deleted after KeptFor (review point 7); open ones never. + // Ended rows are deleted after KeptFor (review point 7); open ones never, nor one whose condition is + // still open (second review, bug B). if _, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: "anchor", Module: "lost", Build: "build-lost"}); err != nil { t.Fatal(err) } @@ -612,10 +647,13 @@ func TestAPendingAssignmentNeverWaitsSilently(t *testing.T) { t.Fatal(err) } for _, r := range rows { - if !r.Open() { - t.Fatalf("an ended pending assignment outlived %s: %+v (%v)", inventory.KeptFor, r, said) + if !r.Open() && (r.Raised == nil || r.Cleared != nil) { + t.Fatalf("an ended pending assignment with no open condition outlived %s: %+v (%v)", inventory.KeptFor, r, said) } } + if !strings.Contains(strings.Join(said, "\n"), "deleted 1 pending assignment(s)") { + t.Fatalf("the applied row was not deleted: %v", said) + } } // The seat's assign passes build on; unassign takes none. @@ -641,3 +679,186 @@ func TestAnAssignmentNotMadeIsSaidPlainly(t *testing.T) { t.Fatalf("headline %q", w.Headline) } } + +// **One row's condition is never another's** (second review, bug A). The sequence: the first pending +// assignment's build fails and the tick raises it; the person assigns again with build "true", a second +// pending assignment; its build fails too before the next tick. That tick raises the second, and must not +// clear it by judging the first answered by the second: each row has its own condition, and a newer row that +// itself ended unmade answers nothing. +func TestASecondAssignmentNotMadeKeepsItsCondition(t *testing.T) { + open := aCatalogueMesh(t) + ctx := t.Context() + asked := 0 + was := askABuild + askABuild = func(context.Context, buildSource, string, string) (string, error) { + asked++ + return fmt.Sprintf("b-%d", asked), nil + } + t.Cleanup(func() { askABuild = was }) + mergeAdding(t, open, "modules/sensors", "3da80a4b00aa") + if _, err := assign(ctx, open, "laptop", "sensors"); err != nil { + t.Fatal(err) + } + failedBuild(t, open, "b-1") + settlePending(ctx, open, time.Now()) + first := pendingOf(t, open)[0] + if !conditionOpen(t, first) { + t.Fatal("the first assignment not made raised nothing") + } + + if _, err := assignWith(ctx, open, "laptop", assignOptions{Build: true}, "sensors"); err != nil { + t.Fatal(err) + } + failedBuild(t, open, "b-2") + settlePending(ctx, open, time.Now()) + + var second inventory.PendingAssignment + for _, p := range pendingOf(t, open) { + if p.Build == "b-2" { + second = p + } + } + if second.State != inventory.PendingExpired { + t.Fatalf("the second: %+v", second) + } + if !conditionOpen(t, second) { + t.Fatal("the second assignment's condition was cleared in the tick that raised it") + } + if !conditionOpen(t, rowOf(t, open, first.ID)) { + t.Fatal("the first was judged answered by a second that itself was not made") + } + // unassign takes both back, and the next tick clears both. + if said, err := unassign(ctx, open, "laptop", "sensors"); err != nil || !strings.Contains(said, "b-1") || + !strings.Contains(said, "b-2") { + t.Fatalf("unassign: %q %v", said, err) + } + settlePending(ctx, open, time.Now()) + if conditionOpen(t, first) || conditionOpen(t, second) { + t.Fatal("a condition stayed open after both were taken back") + } +} + +// **A raised row outlives the pruning until its condition clears** (second review, bug B); and a machine's +// removal, which takes its rows with it, clears their conditions at the next tick. +func TestARaisedRowIsKeptUntilItsConditionClears(t *testing.T) { + open := aCatalogueMesh(t) + ctx := t.Context() + inv := open.inventory + asksWithPaths(t) + mergeAdding(t, open, "modules/sensors", "3da80a4b00aa") + if _, err := assign(ctx, open, "laptop", "sensors"); err != nil { + t.Fatal(err) + } + failedBuild(t, open, "b-modules/sensors") + settlePending(ctx, open, time.Now()) + row := pendingOf(t, open)[0] + + settlePending(ctx, open, time.Now().Add(inventory.KeptFor+time.Hour)) + if rowOf(t, open, row.ID).Raised == nil || !conditionOpen(t, row) { + t.Fatal("a raised row was pruned, or its condition closed, while the condition was open") + } + if _, err := unassign(ctx, open, "laptop", "sensors"); err != nil { + t.Fatal(err) + } + settlePending(ctx, open, time.Now()) + settlePending(ctx, open, time.Now().Add(inventory.KeptFor+time.Hour)) + if known, err := inv.PendingKnown(ctx, []int64{row.ID}); err != nil || known[row.ID] { + t.Fatalf("a cleared row outlived %s: %v", inventory.KeptFor, err) + } + + // On anchor, then anchor removed: the row goes, and its condition with it at the next tick. + if _, err := assign(ctx, open, "anchor", "sensors"); err == nil { + t.Fatal("known and not built was not refused") + } + q, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: "anchor", Module: "sensors", Build: "b-x"}) + if err != nil { + t.Fatal(err) + } + if _, err := inv.SettlePending(ctx, q.ID, inventory.PendingWaiting, inventory.PendingExpired, "boom"); err != nil { + t.Fatal(err) + } + settlePending(ctx, open, time.Now()) + q.Node = "anchor" + if !conditionOpen(t, q) { + t.Fatal("not raised") + } + if _, err := inv.RemoveNodeForTest(ctx, "anchor"); err != nil { + t.Fatal(err) + } + settlePending(ctx, open, time.Now()) + if conditionOpen(t, q) { + t.Fatal("a removed machine's condition stayed open") + } +} + +// **A claim left by a controller that stopped** (second review): the tick settles it from what the mesh +// holds — back to waiting when the module is not assigned there, applied when it is — and leaves a fresh +// claim alone. +func TestAStaleClaimIsSettledFromWhatTheMeshHolds(t *testing.T) { + open := aCatalogueMesh(t) + ctx := t.Context() + inv := open.inventory + claimed := func(node string) inventory.PendingAssignment { + p, err := inv.RecordPending(ctx, inventory.PendingAssignment{Node: node, Module: "sensors", Build: "b-x"}) + if err != nil { + t.Fatal(err) + } + if ok, err := inv.ClaimPending(ctx, p.ID); err != nil || !ok { + t.Fatalf("claim: %v %v", ok, err) + } + return p + } + notAssigned, assigned := claimed("laptop"), claimed("anchor") + register(t, open, catalogue.Manifest{Module: "unrelated", Version: "1"}) + if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "sensors", Version: "1"}, inventory.Source{}); err != nil { + t.Fatal(err) + } + if _, err := inv.Assign(ctx, "anchor", "sensors"); err != nil { + t.Fatal(err) + } + + settlePending(ctx, open, time.Now()) + if rowOf(t, open, notAssigned.ID).State != inventory.PendingApplying || rowOf(t, open, assigned.ID).State != inventory.PendingApplying { + t.Fatal("a fresh claim was settled") + } + settlePending(ctx, open, time.Now().Add(claimStaleAfter+time.Minute)) + if got := rowOf(t, open, assigned.ID); got.State != inventory.PendingApplied { + t.Fatalf("a stale claim of an assigned module: %+v", got) + } + // Released to waiting, and in the same tick made, since sensors is registered now. + if got := rowOf(t, open, notAssigned.ID); got.State != inventory.PendingApplied || + !slices.Contains(assignedTo(t, open, "laptop"), "sensors") { + t.Fatalf("a stale claim of a module not assigned: %+v", got) + } +} + +// **A waited build whose asker stopped waiting is asked, outcome unknown** (second review): it may still run, +// so it reads as in flight until its outcome or its bound, never as not asked; a build never handed over +// is not asked. +func TestABuildNoLongerWaitedForIsStillInFlight(t *testing.T) { + open := aCatalogueMesh(t) + ctx := t.Context() + inv := open.inventory + for _, id := range []string{"build-timeout", "build-nothandedover"} { + dir := "modules/" + strings.TrimPrefix(id, "build-") + if err := inv.RecordBuildRequest(ctx, inventory.BuildRequest{ID: id, Repository: "novox/mesh-catalog", + Seat: "git", Path: dir, Ref: "main", For: "build"}); err != nil { + t.Fatal(err) + } + } + markWaitFailed(ctx, "build-timeout", errors.New("no build machine answered within 10m0s")) + markWaitFailed(ctx, "build-nothandedover", fmt.Errorf("%w: cannot submit a build: nats: timeout", link.ErrNotHandedOver)) + + said, err := assign(ctx, open, "laptop", "timeout") + if err != nil || !strings.Contains(said, "being built") || !strings.Contains(said, "kept as pending") { + t.Fatalf("a build no longer waited for: %q %v", said, err) + } + if _, err := assign(ctx, open, "laptop", "nothandedover"); err == nil || !strings.Contains(err.Error(), "was not handed over") { + t.Fatalf("a build never handed over: %v", err) + } + settlePending(ctx, open, time.Now().Add(buildRequestBound+time.Minute)) + rows, _ := inv.PendingFor(ctx, "laptop", "timeout") + if len(rows) != 1 || rows[0].State != inventory.PendingExpired || !strings.Contains(rows[0].Note, "its asker stopped waiting") { + t.Fatalf("past the bound: %+v", rows) + } +} diff --git a/internal/inventory/fortest.go b/internal/inventory/fortest.go index dddca8f6..74d58805 100644 --- a/internal/inventory/fortest.go +++ b/internal/inventory/fortest.go @@ -68,3 +68,13 @@ func ForTest(t *testing.T) *Inventory { } return inv } + +// RemoveNodeForTest removes a machine's record as a person removing it from the store would: no verb of the +// mesh removes one yet, and what hangs off it must go with it. +func (i *Inventory) RemoveNodeForTest(ctx context.Context, name string) (bool, error) { + tag, err := i.store.Pool().Exec(ctx, `delete from node where name = $1`, name) + if err != nil { + return false, err + } + return tag.RowsAffected() == 1, nil +} diff --git a/internal/inventory/migrations/0082-an-assignment-waits-for-its-build.sql b/internal/inventory/migrations/0082-an-assignment-waits-for-its-build.sql index 09676837..c7cdb947 100644 --- a/internal/inventory/migrations/0082-an-assignment-waits-for-its-build.sql +++ b/internal/inventory/migrations/0082-an-assignment-waits-for-its-build.sql @@ -8,8 +8,10 @@ -- Every build the controller asked for: a merge's new module, a plan's tier, a person's `build`, an `assign` -- with build. Kept once the ask is made, never before. Its outcome is the `build` row of the same id; --- an ask without one is still running, or was lost. not_asked is why the ask could not be made or waited --- for; for a merge that could not ask, the id is the controller's and no build's. +-- an ask without one is still running, or was lost. not_asked is why the ask could not be handed over; for a +-- merge that could not ask, the id is the controller's and no build's. outcome_unknown is why an asker that +-- handed it over stopped waiting: the build may still run, and is read as in flight until its outcome or +-- its bound. create table build_request ( id text primary key, repository text not null, @@ -19,6 +21,7 @@ create table build_request ( commit_hash text not null default '', asked_for text not null default '', not_asked text, + outcome_unknown text, asked_at timestamptz not null default now() ); create index build_request_asked_at on build_request (asked_at); @@ -28,7 +31,7 @@ create index build_request_asked_at on build_request (asked_at); -- assignment was refused then), expired (the build failed, was not registered, or said nothing within its -- bound) or withdrawn (`unassign`). raised_at is when an expiry or refusal was raised as a condition, -- acknowledged_at when a person took it back with `unassign`, cleared_at when its condition was cleared. --- Gone with its machine. Ended rows are deleted after 30 days. +-- Gone with its machine. Ended rows are deleted after 30 days, once any condition raised for them is cleared. create table pending_assignment ( id bigserial primary key, node uuid not null references node (id) on delete cascade, diff --git a/internal/inventory/pending.go b/internal/inventory/pending.go index da22226c..e9d82f57 100644 --- a/internal/inventory/pending.go +++ b/internal/inventory/pending.go @@ -32,9 +32,12 @@ type BuildRequest struct { Commit string // For says who asked: "merge", "plan", "build" (a person), "assign". For string - // NotAsked is why the ask could not be made, or not waited for; empty when it was. + // NotAsked is why the ask could not be handed over; empty when it was. NotAsked string - At time.Time + // OutcomeUnknown is why an asker that handed the build over stopped waiting for it: asked, outcome + // unknown. Read as in flight until its outcome or its bound. + OutcomeUnknown string + At time.Time } // Name is the module this request is expected to register, read from its directory: the last element of @@ -80,8 +83,8 @@ func (i *Inventory) RecordBuildRequest(ctx context.Context, a BuildRequest) erro return err } -// MarkNotAsked says a build request was asked and its asker could not hand it over, or stopped waiting for -// it: the words are kept, unless its outcome was heard first. +// MarkNotAsked says a kept build request was never handed over: the words are kept, unless its outcome was +// heard first. func (i *Inventory) MarkNotAsked(ctx context.Context, id, why string) error { _, err := i.store.Pool().Exec(ctx, `update build_request set not_asked = $2 @@ -89,6 +92,15 @@ func (i *Inventory) MarkNotAsked(ctx context.Context, id, why string) error { return err } +// MarkOutcomeUnknown says the asker of a build it handed over stopped waiting for it: asked, outcome unknown. +// The build may still run; its outcome, when heard, is the last word. +func (i *Inventory) MarkOutcomeUnknown(ctx context.Context, id, why string) error { + _, err := i.store.Pool().Exec(ctx, + `update build_request set outcome_unknown = $2 + where id = $1 and not exists (select 1 from build where build.id = $1)`, id, why) + return err +} + // RequestsNamed is every build request kept whose directory names the module, or whose build said it built // it, newest first, each with its outcome. func (i *Inventory) RequestsNamed(ctx context.Context, name string) ([]RequestOutcome, error) { @@ -126,7 +138,7 @@ func (i *Inventory) RequestedNames(ctx context.Context) ([]string, error) { func (i *Inventory) requests(ctx context.Context) ([]RequestOutcome, error) { rows, err := i.store.Pool().Query(ctx, `select a.id, a.repository, a.seat, a.source_path, a.ref, a.commit_hash, a.asked_for, - coalesce(a.not_asked, ''), a.asked_at, + coalesce(a.not_asked, ''), coalesce(a.outcome_unknown, ''), a.asked_at, b.id is not null, coalesce(b.failed, ''), coalesce(b.module, ''), coalesce(b.at, a.asked_at) from build_request a left join build b on b.id = a.id order by a.asked_at desc, a.id desc`) @@ -138,7 +150,7 @@ func (i *Inventory) requests(ctx context.Context) ([]RequestOutcome, error) { for rows.Next() { var a RequestOutcome if err := rows.Scan(&a.ID, &a.Repository, &a.Seat, &a.Path, &a.Ref, &a.Commit, &a.For, &a.NotAsked, - &a.At, &a.Heard, &a.Failed, &a.Module, &a.HeardAt); err != nil { + &a.OutcomeUnknown, &a.At, &a.Heard, &a.Failed, &a.Module, &a.HeardAt); err != nil { return nil, err } out = append(out, a) @@ -342,13 +354,34 @@ func (i *Inventory) SettlePending(ctx context.Context, id int64, from, state, no return tag.RowsAffected() == 1, nil } -// ForgetEndedPending deletes ended pending assignments settled more than KeptFor ago, and says how many. +// ForgetEndedPending deletes ended pending assignments settled more than KeptFor ago, and says how many. One +// raised as a condition is kept until that condition is cleared: deleting it would leave the condition +// with nothing to answer it. func (i *Inventory) ForgetEndedPending(ctx context.Context, now time.Time) (int64, error) { tag, err := i.store.Pool().Exec(ctx, - `delete from pending_assignment where state not in ('waiting', 'applying') and settled_at < $1`, + `delete from pending_assignment where state not in ('waiting', 'applying') and settled_at < $1 + and (raised_at is null or cleared_at is not null)`, now.Add(-KeptFor)) if err != nil { return 0, err } return tag.RowsAffected(), nil } + +// PendingKnown says which of these pending assignments are still on record. +func (i *Inventory) PendingKnown(ctx context.Context, ids []int64) (map[int64]bool, error) { + rows, err := i.store.Pool().Query(ctx, `select id from pending_assignment where id = any($1)`, ids) + if err != nil { + return nil, err + } + defer rows.Close() + out := map[int64]bool{} + for rows.Next() { + var id int64 + if err := rows.Scan(&id); err != nil { + return nil, err + } + out[id] = true + } + return out, rows.Err() +} diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index f5d1c4ca..a2cbd1f6 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -76,6 +76,11 @@ func (b *natsBuilds) Ask(ctx context.Context, request BuildRequest) error { return nil } +// ErrNotHandedOver is a Submit that failed before the build was handed to the build seat: nothing runs. Any +// other error from Submit came after: the build may still run, and its outcome may still come (novox/hq +// issue 325). +var ErrNotHandedOver = errors.New("the build was not handed over") + func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest, wait time.Duration) (BuildResult, error) { @@ -84,23 +89,23 @@ func (b *natsBuilds) Submit(ctx context.Context, request BuildRequest, // the same event on EVENTS, which the controller records. outcomes, err := b.js.Conn().SubscribeSync(BuildOutcomeOf(b.role())) if err != nil { - return BuildResult{}, fmt.Errorf("cannot listen for a build's outcome: %w", err) + return BuildResult{}, fmt.Errorf("%w: cannot listen for a build's outcome: %w", ErrNotHandedOver, err) } defer func() { _ = outcomes.Unsubscribe() }() if err := b.js.Conn().Flush(); err != nil { - return BuildResult{}, err + return BuildResult{}, fmt.Errorf("%w: %w", ErrNotHandedOver, err) } body, err := json.Marshal(request) if err != nil { - return BuildResult{}, err + return BuildResult{}, fmt.Errorf("%w: %w", ErrNotHandedOver, err) } // Into the role's work queue and awaited: work the bus never accepted must fail here rather than // be assumed, because nothing else will ever say so. publish, cancel := context.WithTimeout(ctx, 30*time.Second) defer cancel() if _, err := b.js.Context().Publish(BuildWorkOf(b.role()), body, nats.Context(publish)); err != nil { - return BuildResult{}, fmt.Errorf("cannot submit a build: %w", err) + return BuildResult{}, fmt.Errorf("%w: cannot submit a build: %w", ErrNotHandedOver, err) } waiting, cancelWait := context.WithTimeout(ctx, wait)