Give each pending assignment its own condition, and keep a raised row until it clears
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery-group group fix/assign-says-why-a-module-is-not-there delivering: 1 of 2 delivered
mesh/delivery superseded: a newer delivery to the same trunk took over its walk

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.
This commit is contained in:
jochen
2026-10-08 16:45:43 +02:00
parent 249d97d1c8
commit 8adb7f1a05
7 changed files with 367 additions and 35 deletions
+14 -8
View File
@@ -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
}
+62 -8
View File
@@ -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
}
// `<scope>.<node>.<module>.<row>.<kind>`: 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
}
+225 -4
View File
@@ -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)
}
}