4 Commits
Author SHA1 Message Date
jochen 822e52123c Say S17 clears once the new bus is sent, as it does
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 ready: it delivers once merged
The comment and the signals table's bound said it cleared once the bus's machine runs the new build; it
clears once that machine has been sent it, when no send is refused for the bus any more.
2026-10-09 03:23:47 +02:00
jochen 9a7ae131cb TestReplay336 reads the bus call in its mesh form
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 superseded: a newer head of the same pull request
The previous commit named the upgrade as a mesh call, and the replay still looked for the command-line
words; it now asserts the seat and the upgrade, as any wording of the call says them.
2026-10-09 03:08:21 +02:00
jochen 4354d9d7c7 Name the bus upgrade as the mesh call a person makes, and say a same-source rebuild needs no step (review of #173)
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 fail: its merge-check.sh failed: --- FAIL: TestReplay336 (1.12s)
mesh/delivery superseded: a newer head of the same pull request
The condition and the plan named the verb in command-line form; the operator reaches it through the mesh
MCP server, so it is written as that call. A rebuild of the same source moves nothing (issue 280), which
the plan cannot know before the build, so its line says so.
2026-10-09 02:54:23 +02:00
jochen 1fdc68a8aa Tell the operator at once when sends wait for the bus's planned step, and say it before the merge (hq issue 336)
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 superseded: a newer head of the same pull request
For 28 minutes on 2026-10-08 every send to the control-node was refused for a new bus build that only a
person's bus upgrade moves, and no condition said so: the refusal lived only in each walk's note, and S3
would have called it lateness after half an hour, in words that named neither the bus nor the verb.

- Row S17, bus.<module>.step-waiting: raised on the first watchdog tick after a walk's send is refused for
  the bus, for the operator, naming the machines, what waits, the bus build from and to, since when and
  mesh-controller.bus upgrade. It clears once the bus's machine runs the build the mesh holds.
- S3 leaves out a walk held only by the bus's step.
- A change that builds the bus says in its delivery plan and summary (which mesh/merge-gate carries) that
  merging it needs a person's bus upgrade, and that nothing else reaches its machine until then.
- TestReplay336 fails on the commit before and passes on this one.
2026-10-09 02:33:57 +02:00
10 changed files with 533 additions and 9 deletions
+195
View File
@@ -0,0 +1,195 @@
package main
import (
"fmt"
"sort"
"strings"
"sync"
"time"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
)
// A send held for the bus's planned step (novox/hq issue 336).
//
// **A wait only a person can end is said to that person at once.** No send replaces the bus but its planned
// step (sendToEach, ADR 0236), and the step is a person's `bus upgrade`. So while a new bus build waits, every
// walk whose tier sends to the bus's machine is refused, and tries again each tick. The refusal was kept only
// as the walk's note; nothing was raised, and the operator found the hold by asking why a walk did not move.
// Row S17 raises it on the first tick after the first refusal, in the bus's scope, for the operator: what
// waits, behind which bus build, since when, and the verb. A walk held only by it is not late, so S3 leaves it
// out, as it leaves out a walk under a paused build seat. It clears once the bus's machine has been sent the
// build the mesh holds — the step was taken, and no send is refused for the bus any more — whatever the walks'
// notes still say. Whether the new bus came up healthy is the step's own condition (probe DB).
// kindBusStepWaiting is S17's kind: bus.<module>.step-waiting.
const kindBusStepWaiting = "bus-step-waiting"
// busFacts is a new bus build waiting for its planned step, and the sends held for it.
type busFacts struct {
module, to string
// from is the build each machine of the bus runs; machines those whose bus the step would replace.
from map[string]string
machines []string
waits []busWaitFacts
}
// busWaitFacts is one walk whose send was refused because it would replace the bus.
type busWaitFacts struct {
plan, repository, commit string
// modules are what its tier sends: what waits.
modules []string
// since is the first refusal: as this controller saw it, or, read back, the save that kept the refusal.
since time.Time
}
// busRefusedFirst is when this controller first saw each walk refused for the bus's step. The walk's note
// keeps the refusal across a restart; this keeps its moment more exactly than the walk's last save.
var busRefusedFirst = &firstSeen{at: map[string]time.Time{}}
type firstSeen struct {
mu sync.Mutex
at map[string]time.Time
}
// mark keeps the first moment an id was seen.
func (s *firstSeen) mark(id string, at time.Time) {
s.mu.Lock()
defer s.mu.Unlock()
if _, seen := s.at[id]; !seen {
s.at[id] = at
}
}
// of is when an id was first seen; zero when it was not.
func (s *firstSeen) of(id string) time.Time {
s.mu.Lock()
defer s.mu.Unlock()
return s.at[id]
}
// keepOnly forgets every id not given: a walk no longer held is not held since then.
func (s *firstSeen) keepOnly(ids map[string]bool) {
s.mu.Lock()
defer s.mu.Unlock()
for id := range s.at {
if !ids[id] {
delete(s.at, id)
}
}
}
// refusedForTheBus is whether an error is the refusal of a send for the bus's planned step.
func refusedForTheBus(err error) bool {
return err != nil && strings.Contains(err.Error(), errBusWaits.Error())
}
// busWaitsOf is the sends held for the bus's step, from the open walks: each walk not waiting for its
// delivery's word whose note keeps a refusal for this very bus build, while the build would still replace
// the bus on a machine. first is when this controller first saw a walk refused, zero when it did not.
func busWaitsOf(plans []inventory.Plan, b busPending, first func(string) time.Time) busFacts {
f := busFacts{module: b.module, to: b.to, from: b.from}
for _, n := range b.machines {
if b.moves(n) {
f.machines = append(f.machines, n)
}
}
if len(f.machines) == 0 {
return f
}
for _, p := range plans {
if !p.Open() || p.Waiting() || !strings.Contains(p.Note, errBusWaits.Error()) || !strings.Contains(p.Note, short(b.to)) {
continue
}
since := p.Updated
if seen := first(p.ID); !seen.IsZero() && (since.IsZero() || seen.Before(since)) {
since = seen
}
var modules []string
if p.Tier >= 0 && p.Tier < len(p.Tiers) {
modules = append(modules, p.Tiers[p.Tier]...)
} else {
modules = planModules(p)
}
sort.Strings(modules)
f.waits = append(f.waits, busWaitFacts{plan: p.ID, repository: p.Repository, commit: p.Commit, modules: modules,
since: since})
}
sort.Slice(f.waits, func(i, j int) bool { return f.waits[i].since.Before(f.waits[j].since) })
return f
}
// heldByTheBus is the walks of a bus's facts, by id: what S3 leaves out.
func (b busFacts) heldByTheBus() map[string]bool {
out := map[string]bool{}
for _, w := range b.waits {
out[w.plan] = true
}
return out
}
// busUpgradeVerb is the call that ends the wait, as the summary names it: through the mesh MCP server, with
// why, saying whether the new version can be undone ("reversible": "true") or not ("irreversible": "true").
const busUpgradeVerb = "`mesh_call mesh-controller.bus {\"upgrade\": \"true\", \"why\": \"…\", \"reversible\" or \"irreversible\": \"true\"}`"
// watchBusWaits is S17: a send held for the bus's planned step, said at once to the operator.
func watchBusWaits(f *signalFacts) []conditions.Observation {
b := f.bus
if len(b.waits) == 0 || len(b.machines) == 0 {
return nil
}
since := b.waits[0].since
var walks, what []string
for _, w := range b.waits {
walks = append(walks, fmt.Sprintf("%s (%s %s, tier: %s)", w.plan, w.repository, short(w.commit),
strings.Join(w.modules, ", ")))
what = append(what, deliveryWhat(w.modules, w.repository))
}
var from []string
for _, n := range b.machines {
from = append(from, short(orNotKnown(b.from[n])))
}
machines := namesWords(b.machines, 3)
in := f.now.Sub(since)
return []conditions.Observation{{Scope: conditions.ScopeBus, ID: b.module, Token: "step-waiting",
Kind: kindBusStepWaiting, Severity: conditions.Warning, Resolver: conditions.ResolverOperator,
Machine: b.machines[0], Also: b.machines[1:],
Summary: fmt.Sprintf("sends to %s wait for the bus's planned step: a new bus build (%s %s → %s) would replace "+
"the bus there, which only a person's %s does; waiting since %s: %s", strings.Join(b.machines, ", "),
b.module, strings.Join(sortedUnique(from), ", "), short(b.to), busUpgradeVerb,
since.UTC().Format(time.RFC3339), strings.Join(walks, "; ")),
Said: fmt.Sprintf("%d walk(s) refused since %s", len(b.waits), since.UTC().Format(time.RFC3339)),
Headline: clipWords("Sends to "+machines+" wait for a bus upgrade", conditions.HeadlineMax),
Explanation: clipWords(fmt.Sprintf("A new version of the mesh's message system is built, and only a person "+
"installs it. Until then nothing else is sent to %s: %s waits, for %s so far.", machines,
namesWords(sortedUnique(what), 3), humanDuration(in)), conditions.ExplanationMax),
Needs: "start the bus upgrade " + FromMeshMCPServer,
Resolved: "Resolved: the bus upgrade started, and sends go on"}}
}
// sortedUnique is a list sorted, each once.
func sortedUnique(xs []string) []string {
seen := map[string]bool{}
var out []string
for _, x := range xs {
if !seen[x] {
seen[x] = true
out = append(out, x)
}
}
sort.Strings(out)
return out
}
// clipWords keeps a plain sentence within a bound, at a word.
func clipWords(s string, n int) string {
if len(s) <= n {
return s
}
cut := strings.LastIndex(s[:n-1], " ")
if cut <= 0 {
cut = n - 1
}
return s[:cut] + "…"
}
+159
View File
@@ -0,0 +1,159 @@
package main
import (
"context"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
)
// novox/hq issue 336: a send held because it would replace the bus outside its planned step waits for a
// person, so a person is told at once — what waits, behind which bus build, since when, and the verb that
// ends it — and the wait is not read as a walk running late.
// aBusOnAnchor is a new bus build waiting for its step on anchor: nats 88135ad0 there, 32307bd1 held.
func aBusOnAnchor() busPending {
return busPending{module: "nats", machines: []string{"anchor"}, from: map[string]string{"anchor": "88135ad0aaaa"},
to: "32307bd1bbbb", same: map[string]bool{}}
}
// refusedNote is the note a walk keeps when its send was refused for the bus, as advanceHeld writes it.
func refusedNote(b busPending, tier int) string {
held := "sending anchor would replace the bus (nats " + short(b.from["anchor"]) + " → " + short(b.to) +
"), which is a planned step: `bus upgrade --why …` snapshots its streams first and checks them after (novox/hq ADR 0236)"
return "tier " + string(rune('0'+tier)) + ": " + errBusWaits.Error() + ": " + held + " — tried again"
}
func TestASendHeldForTheBusStepIsFoundFromTheWalksItHolds(t *testing.T) {
now := time.Date(2026, 10, 8, 18, 30, 0, 0, time.UTC)
b := aBusOnAnchor()
walk := func(id, repo, note string, updated time.Time) inventory.Plan {
return inventory.Plan{ID: id, Repository: repo, Commit: "c0ffee001122", State: inventory.PlanRolling, Tier: 0,
Tiers: [][]string{{"mesh-host"}}, Modules: map[string]*inventory.PlanModule{"mesh-host": {}}, Note: note,
Updated: updated}
}
for _, c := range []struct {
name string
bus busPending
plans []inventory.Plan
first map[string]time.Time
want []string
since time.Time
}{
{"a walk refused for the bus", b,
[]inventory.Plan{walk("plan-1", "novox/mesh-host", refusedNote(b, 0), now.Add(-27*time.Minute))}, nil,
[]string{"plan-1"}, now.Add(-27 * time.Minute)},
{"refused earlier than its last save, as this controller saw it", b,
[]inventory.Plan{walk("plan-1", "novox/mesh-host", refusedNote(b, 0), now.Add(-5*time.Minute))},
map[string]time.Time{"plan-1": now.Add(-28 * time.Minute)}, []string{"plan-1"}, now.Add(-28 * time.Minute)},
{"a walk refused for another reason", b,
[]inventory.Plan{walk("plan-1", "novox/mesh-host", "tier 0: the build seat is paused — tried again", now)}, nil,
nil, time.Time{}},
{"refused for an older bus build than the one held now", func() busPending { o := b; o.to = "99999999cccc"; return o }(),
[]inventory.Plan{walk("plan-1", "novox/mesh-host", refusedNote(b, 0), now)}, nil, nil, time.Time{}},
{"the bus already runs the build held: its step started", func() busPending {
o := aBusOnAnchor()
o.from = map[string]string{"anchor": o.to}
return o
}(), []inventory.Plan{walk("plan-1", "novox/mesh-host", refusedNote(b, 0), now)}, nil, nil, time.Time{}},
{"a walk waiting for its delivery's word asks no send", b, func() []inventory.Plan {
p := walk("plan-1", "novox/mesh-host", refusedNote(b, 0), now)
p.Delivery = &inventory.PlanDelivery{Awaits: "mesh-delivery"}
return []inventory.Plan{p}
}(), nil, nil, time.Time{}},
} {
t.Run(c.name, func(t *testing.T) {
got := busWaitsOf(c.plans, c.bus, func(id string) time.Time { return c.first[id] })
var ids []string
for _, w := range got.waits {
ids = append(ids, w.plan)
}
if strings.Join(ids, ",") != strings.Join(c.want, ",") {
t.Fatalf("held for the bus: %v, want %v", ids, c.want)
}
if len(c.want) > 0 {
if !got.waits[0].since.Equal(c.since) {
t.Errorf("waiting since %s, want %s", got.waits[0].since, c.since)
}
if strings.Join(got.machines, ",") != "anchor" || got.to != b.to || got.module != "nats" {
t.Errorf("behind %+v", got)
}
}
})
}
}
// **Said at once, not as lateness, and cleared when the step starts**: the walk held only for the bus raises
// the bus's condition on the first tick, before any tier bound, naming what waits, the bus build from and to,
// since when and `bus upgrade`, for the operator; S3 says nothing of it even past its bound; and the
// condition clears once the bus's machine runs the new build.
func TestASendHeldForTheBusStepIsSaidAtOnceAndNotAsLateness(t *testing.T) {
now := time.Date(2026, 10, 8, 18, 30, 0, 0, time.UTC)
b := aBusOnAnchor()
held := func(entered time.Duration) *signalFacts {
f := calm(now)
f.plans = []planFacts{{id: "plan-1", repository: "novox/mesh-host", commit: "c0ffee00", tier: 0, tiers: 1,
entered: now.Add(-entered), bound: 30 * time.Minute, waiting: refusedNote(b, 0), bus: true}}
f.bus = busFacts{module: "nats", to: b.to, from: b.from, machines: []string{"anchor"},
waits: []busWaitFacts{{plan: "plan-1", repository: "novox/mesh-host", commit: "c0ffee001122",
modules: []string{"mesh-host"}, since: now.Add(-time.Minute)}}}
return f
}
f := held(time.Minute)
got := watchBusWaits(f)
if len(got) != 1 {
t.Fatalf("a send held for the bus a minute ago raised %+v", got)
}
o := got[0]
if o.Key() != "bus.nats.step-waiting" || o.Kind != kindBusStepWaiting || o.Resolver != conditions.ResolverOperator {
t.Fatalf("raised %s (%s), resolver %q", o.Key(), o.Kind, o.Resolver)
}
for _, want := range []string{"anchor", "mesh-host", "88135ad0", "32307bd1", "mesh_call mesh-controller.bus",
`"upgrade": "true"`, `"reversible"`, `"irreversible"`, "2026-10-08T18:29:00Z", "plan-1"} {
if !strings.Contains(o.Summary, want) {
t.Errorf("its summary does not say %q: %s", want, o.Summary)
}
}
if !strings.Contains(o.Explanation, "mesh-host") || !strings.Contains(o.Headline, "bus upgrade") || o.Needs == "" {
t.Errorf("its words do not say what waits and what the operator does: %+v", o)
}
if why, ok := conditions.PlainWords(conditions.Words{Headline: o.Headline, Explanation: o.Explanation,
Resolved: o.Resolved, Needs: o.Needs}, "anchor"); !ok {
t.Errorf("its words are not plain: %s", why)
}
// Past S3's bound: still the bus's wait, never a stalled walk.
late := held(45 * time.Minute)
if s3 := watchPlans(late); len(s3) != 0 {
t.Fatalf("a walk held only by the bus step was said stalled: %+v", s3)
}
// A walk held for something else past its bound is still stalled.
other := held(45 * time.Minute)
other.plans[0].bus = false
if s3 := watchPlans(other); len(s3) != 1 {
t.Fatalf("a walk held for something else past its bound raised %+v", s3)
}
// Through the keeper: open at the first tick, cleared when the step started.
store := conditions.NewInMemory()
k := conditions.NewKeeper(t.Context(), conditions.Options{Store: store, History: store, Teller: &conditions.Told{},
Now: func() time.Time { return now }})
defer k.Close(context.Background())
w := &watchdogs{keeper: k, started: now.Add(-time.Hour)}
w.see(t.Context(), held(time.Minute))
open, err := k.Open(t.Context())
if err != nil || len(open) != 1 || open[0].Key != "bus.nats.step-waiting" {
t.Fatalf("after the first tick, open: %+v (%v)", open, err)
}
started := held(time.Minute)
started.bus = busFacts{module: "nats", to: b.to}
started.plans[0].bus = false
w.see(t.Context(), started)
if open, _ := k.Open(t.Context()); len(open) != 0 {
t.Fatalf("the step started, and open: %+v", open)
}
}
+15 -2
View File
@@ -5,6 +5,7 @@ import (
"sort" "sort"
"strings" "strings"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link" "github.com/novox/mesh-controller/internal/link"
) )
@@ -50,7 +51,13 @@ func changePlanOf(repository, base, head string, r mergeReach, entries []invento
} }
} }
switch { switch {
case name == "nats": case (name == "nats" || catalogue.ProvidesBus(e.Manifest)) && len(e.On) > 0:
// Said before the merge (novox/hq issue 336): a new bus build holds every send to the bus's machine
// until a person takes the step, so merging it is a promise to take it.
p.Steps = append(p.Steps, fmt.Sprintf("%s: %s is never sent by an ordinary send, and until %s runs, "+
"nothing else is sent to %s either — unless the new build turns out the same as the one running there "+
"(issue 280)", busUpgradeNeeded, name, busUpgradeVerb, strings.Join(e.On, ", ")))
case name == "nats" || catalogue.ProvidesBus(e.Manifest):
p.Steps = append(p.Steps, "a planned bus step: the bus is upgraded by `bus upgrade`, never by an ordinary send") p.Steps = append(p.Steps, "a planned bus step: the bus is upgraded by `bus upgrade`, never by an ordinary send")
case waits && len(e.On) > 0: case waits && len(e.On) > 0:
p.Steps = append(p.Steps, fmt.Sprintf("%s waits for a person: its policy records (%s)", name, p.Steps = append(p.Steps, fmt.Sprintf("%s waits for a person: its policy records (%s)", name,
@@ -81,6 +88,9 @@ func changePlanOf(repository, base, head string, r mergeReach, entries []invento
return p return p
} }
// busUpgradeNeeded is how a change plan says that merging it holds the bus's machine for a person's step.
const busUpgradeNeeded = "merging this needs a person's bus upgrade"
// summaryOf is a change plan in one line: what it builds, where it goes, and whether the bus moves. // summaryOf is a change plan in one line: what it builds, where it goes, and whether the bus moves.
func summaryOf(p link.ChangePlan) string { func summaryOf(p link.ChangePlan) string {
if len(p.Moved) == 0 && len(p.New) == 0 { if len(p.Moved) == 0 && len(p.New) == 0 {
@@ -110,7 +120,10 @@ func summaryOf(p link.ChangePlan) string {
} }
bus := "no bus step" bus := "no bus step"
for _, s := range p.Steps { for _, s := range p.Steps {
if strings.HasPrefix(s, "a planned bus step") { switch {
case strings.HasPrefix(s, busUpgradeNeeded):
bus = busUpgradeNeeded
case strings.HasPrefix(s, "a planned bus step") && bus == "no bus step":
bus = "a bus step" bus = "a bus step"
} }
} }
+6 -2
View File
@@ -45,7 +45,7 @@ func TestAChangePlanSaysWhatEachMachineReceives(t *testing.T) {
} }
p = plan("modules/nats/Dockerfile", "modules/photos/x.js") p = plan("modules/nats/Dockerfile", "modules/photos/x.js")
if !strings.Contains(p.Summary, "a bus step") || !strings.Contains(p.Summary, "2 wait(s) for a person") || if !strings.Contains(p.Summary, "needs a person's bus upgrade") || !strings.Contains(p.Summary, "2 wait(s) for a person") ||
!strings.Contains(p.Summary, "(+1 dependent(s))") { !strings.Contains(p.Summary, "(+1 dependent(s))") {
t.Errorf("the bus and a held module read %q", p.Summary) t.Errorf("the bus and a held module read %q", p.Summary)
} }
@@ -58,7 +58,11 @@ func TestAChangePlanSaysWhatEachMachineReceives(t *testing.T) {
t.Errorf("the deploy plan reads %v", got) t.Errorf("the deploy plan reads %v", got)
} }
text := strings.Join(p.Steps, "\n") text := strings.Join(p.Steps, "\n")
for _, want := range []string{"a planned bus step", "photos waits for a person", "nats provides mesh-bus"} { // Said before the merge (novox/hq issue 336): merging it holds every send to the bus's machine until a
// person runs the bus's step, and the verb that does.
for _, want := range []string{"merging this needs a person's bus upgrade", "nothing else is sent to anchor",
"mesh_call mesh-controller.bus", "the same as the one running there", "photos waits for a person",
"nats provides mesh-bus"} {
if !strings.Contains(text, want) { if !strings.Contains(text, want) {
t.Errorf("the steps do not say %q:\n%s", want, text) t.Errorf("the steps do not say %q:\n%s", want, text)
} }
+7
View File
@@ -523,6 +523,13 @@ var plainWordings = map[string]func(conditions.Observation) words{
Explanation: "The mesh's message system is being upgraded; some things pause until it is done.", Explanation: "The mesh's message system is being upgraded; some things pause until it is done.",
Resolved: "The bus upgrade is done"} Resolved: "The bus upgrade is done"}
}), }),
kindBusStepWaiting: worded(func(o conditions.Observation) words {
return words{Headline: "Sends wait for a bus upgrade",
Needs: "start the bus upgrade " + FromMeshMCPServer,
Explanation: "A new version of the mesh's message system is built, and only a person installs it. Until " +
"then nothing else is sent to the machine that runs it.",
Resolved: "Resolved: the bus upgrade started, and sends go on"}
}),
kindBusUpgradeFailed: worded(func(o conditions.Observation) words { kindBusUpgradeFailed: worded(func(o conditions.Observation) words {
return words{Headline: "The bus upgrade failed", return words{Headline: "The bus upgrade failed",
Needs: "decide whether to put the bus back to the version before; the details say how.", Needs: "decide whether to put the bus back to the version before; the details say how.",
+4
View File
@@ -533,6 +533,10 @@ func advanceHeld(ctx context.Context, open *stores) {
// the state is left as it was and the step is tried again on the next tick. Said and // the state is left as it was and the step is tried again on the next tick. Said and
// kept when it is new: the same refusal on every tick is one fact, not one per tick. // kept when it is new: the same refusal on every tick is one fact, not one per tick.
p.Note = "tier " + fmt.Sprint(p.Tier) + ": " + err.Error() + " — tried again" p.Note = "tier " + fmt.Sprint(p.Tier) + ": " + err.Error() + " — tried again"
// A refusal for the bus's planned step is a person's to end: S17 says it from its first moment.
if refusedForTheBus(err) {
busRefusedFirst.mark(p.ID, time.Now())
}
if planSnapshot(*p) != before { if planSnapshot(*p) != before {
fmt.Printf("%s: %v\n", p.ID, err) fmt.Printf("%s: %v\n", p.ID, err)
if err := inv.SavePlan(ctx, p); err != nil { if err := inv.SavePlan(ctx, p); err != nil {
+92
View File
@@ -0,0 +1,92 @@
package main
import (
"context"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
)
// TestReplay336 replays novox/hq issue 336 (2026-10-08): a catalogue merge built a new bus build for the
// control-node; from then on every send to that machine was refused for the bus's planned step, which only a
// person starts, and for 28 minutes nothing but the walks' notes said so. The outcome asserted: the first
// watchdog tick after the first refusal opens a condition for the operator that names what waits and the verb
// that ends it, and it clears once the bus's machine runs the new build. Written with only what the controller
// had before its fix — the stores, register, assign, a plan saved and advanced, the watchdogs' gathering and
// seeing — so it is laid over the commit before.
func TestReplay336(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
inv := open.inventory
bus := catalogue.Manifest{Module: "nats", Version: "2", Provides: []catalogue.Offer{{Name: "mesh-bus"}}}
if err := inv.RegisterModule(ctx, bus, inventory.Source{Repository: "novox/mesh-catalog", Seat: "git",
Path: "modules/nats", BuiltFrom: "32307bd1", Head: "32307bd1"}); err != nil {
t.Fatal(err)
}
register(t, open, catalogue.Manifest{Module: "engine", Version: "1"})
for _, m := range []string{"nats", "engine"} {
if _, err := inv.Assign(ctx, "anchor", m); err != nil {
t.Fatal(err)
}
}
// The control-node runs the bus build it was last sent; the mesh holds a newer one, which waits for its step.
if err := inv.RecordSent(ctx, nodeID(t, open, "anchor"), "d-anchor", map[string]string{"nats": "88135ad0"}); err != nil {
t.Fatal(err)
}
// A walk of the engine, built, whose tier sends to the control-node.
now := time.Now().UTC()
plan := inventory.Plan{ID: "plan-engine", Repository: "novox/mesh-host", Commit: "e1e1e1e1", Created: now,
State: inventory.PlanBuilding, Tiers: [][]string{{"engine"}},
Modules: map[string]*inventory.PlanModule{"engine": {State: "built", BuiltAt: &now, Commit: "e1e1e1e1", Build: "b"}}}
if err := inv.SavePlan(ctx, &plan); err != nil {
t.Fatal(err)
}
advancePlans(ctx, open)
p, err := inv.PlanByID(ctx, "plan-engine")
if err != nil {
t.Fatal(err)
}
if !strings.Contains(p.Note, errBusWaits.Error()) {
t.Fatalf("the walk's send to the bus's machine was not refused for the bus: %s %q", p.State, p.Note)
}
store := conditions.NewInMemory()
k := conditions.NewKeeper(ctx, conditions.Options{Store: store, History: store, Teller: &conditions.Told{}})
defer k.Close(context.Background())
w := &watchdogs{open: open, keeper: k, started: now.Add(-time.Hour)}
waiting := func() []conditions.Condition {
t.Helper()
w.see(ctx, w.gather(ctx))
all, err := k.Open(ctx)
if err != nil {
t.Fatal(err)
}
var out []conditions.Condition
for _, c := range all {
if strings.HasPrefix(c.Key, "bus.") && c.Resolver == conditions.ResolverOperator {
out = append(out, c)
}
}
return out
}
got := waiting()
if len(got) != 1 {
t.Fatalf("the first tick after the refusal told the operator nothing about the bus's step: %d condition(s)", len(got))
}
for _, want := range []string{"anchor", "engine", "mesh-controller.bus", "upgrade", "32307bd1"} {
if !strings.Contains(got[0].Summary, want) {
t.Errorf("the condition does not say %q: %s", want, got[0].Summary)
}
}
// The step taken: the control-node is sent the new bus build, and the wait is over.
if err := inv.RecordSent(ctx, nodeID(t, open, "anchor"), "d-anchor-2", map[string]string{"nats": "32307bd1"}); err != nil {
t.Fatal(err)
}
if got := waiting(); len(got) != 0 {
t.Fatalf("the bus's machine runs the new build, and the operator is still asked: %+v", got)
}
}
+17 -1
View File
@@ -210,6 +210,20 @@ var signalsTable = []signalRow{
newest: func(f *signalFacts) time.Time { newest: func(f *signalFacts) time.Time {
return newestOf(f.waits, func(w waitFacts) time.Time { return w.since }) return newestOf(f.waits, func(w waitFacts) time.Time { return w.since })
}}, }},
{Row: "S17", Signal: "a send held for the bus's planned step is told to a person", Emitter: "controller's plan",
Trigger: "each send refused because it would replace the bus outside its step (novox/hq issue 336)",
Bound: "none: raised at the first refusal, for the operator, naming what waits, the bus build from and to, " +
"since when and the mesh-controller.bus call; cleared once the bus's machine has been sent the build the mesh holds",
Kind: kindBusStepWaiting, Severity: conditions.Warning, Phase: 3,
needs: func(f *signalFacts) error {
if f.plansErr != nil {
return f.plansErr
}
return f.busErr
}, watch: watchBusWaits,
newest: func(f *signalFacts) time.Time {
return newestOf(f.bus.waits, func(w busWaitFacts) time.Time { return w.since })
}},
} }
// watchFacts is S14: the snapshot a merge check is fed is older than its bound, or none was kept since // watchFacts is S14: the snapshot a merge check is fed is older than its bound, or none was kept since
@@ -315,7 +329,9 @@ func watchReports(f *signalFacts) []conditions.Observation {
func watchPlans(f *signalFacts) []conditions.Observation { func watchPlans(f *signalFacts) []conditions.Observation {
var out []conditions.Observation var out []conditions.Observation
for _, p := range f.plans { for _, p := range f.plans {
if p.paused { // Under a paused build seat, or held only by the bus's planned step (S17, novox/hq issue 336): a wait
// for a person, said as itself, not a walk running late.
if p.paused || p.bus {
continue continue
} }
in := f.now.Sub(p.entered) in := f.now.Sub(p.entered)
+13
View File
@@ -142,6 +142,19 @@ var suppressions = map[string]suppression{
since: f.now.Add(-31 * time.Minute)}} since: f.now.Add(-31 * time.Minute)}}
}, },
}, },
// A send refused because it would replace the bus outside its planned step: said at its first refusal,
// whatever the bound (novox/hq issue 336). Inside: a new bus build waits, and no send was refused for it.
"S17": {
inside: func(f *signalFacts) {
f.bus = busFacts{module: "nats", to: "32307bd1bbbb", from: map[string]string{"anchor": "88135ad0aaaa"},
machines: []string{"anchor"}}
},
past: func(f *signalFacts) {
f.bus = busFacts{module: "nats", to: "32307bd1bbbb", from: map[string]string{"anchor": "88135ad0aaaa"},
machines: []string{"anchor"}, waits: []busWaitFacts{{plan: "plan-1", repository: "novox/app",
commit: "c0ffee001122", modules: []string{"app"}, since: f.now.Add(-time.Second)}}}
},
},
// Twice by hand within a fortnight is a healer wanted; once, or the first of two a day too old, is not. // Twice by hand within a fortnight is a healer wanted; once, or the first of two a day too old, is not.
"S15": { "S15": {
inside: func(f *signalFacts) { inside: func(f *signalFacts) {
+25 -4
View File
@@ -97,6 +97,9 @@ type signalFacts struct {
// facts is the snapshot this controller keeps for merge checks (S14). // facts is the snapshot this controller keeps for merge checks (S14).
facts factsFacts facts factsFacts
bus busFacts
busErr error
} }
// factsFacts is when the newest facts snapshot was taken, when this controller began keeping it, and // factsFacts is when the newest facts snapshot was taken, when this controller began keeping it, and
@@ -147,6 +150,7 @@ type planFacts struct {
bound time.Duration bound time.Duration
waiting string waiting string
paused bool paused bool
bus bool
} }
// waitFacts is one walk waiting for its delivery's word: since its merge opened it. // waitFacts is one walk waiting for its delivery's word: since its merge opened it.
@@ -320,7 +324,8 @@ func (w *watchdogs) gather(ctx context.Context) *signalFacts {
} }
} }
f.machines, f.machinesErr = w.gatherMachines(ctx, inv, now) f.machines, f.machinesErr = w.gatherMachines(ctx, inv, now)
f.plans, f.waits, f.plansErr = gatherPlans(ctx, inv, now) f.bus, f.busErr = gatherBus(ctx, inv)
f.plans, f.waits, f.plansErr = gatherPlans(ctx, inv, now, f.bus.heldByTheBus())
f.loop, f.loopErr = w.gatherLoop() f.loop, f.loopErr = w.gatherLoop()
f.mergesPassed, f.merges, f.mergesErr = watchedMerges.last() f.mergesPassed, f.merges, f.mergesErr = watchedMerges.last()
if f.mergesErr == nil && !f.mergesPassed.IsZero() && now.Sub(f.mergesPassed) > 3*mergeCatchUpEvery { if f.mergesErr == nil && !f.mergesPassed.IsZero() && now.Sub(f.mergesPassed) > 3*mergeCatchUpEvery {
@@ -438,8 +443,24 @@ func (w *watchdogs) gatherMachines(ctx context.Context, inv *inventory.Inventory
return out, nil return out, nil
} }
// gatherPlans is every open plan, its tier's bound from what was measured, and what it waits on. // gatherBus is a new bus build waiting for its step and the walks refused for it (S17, novox/hq issue 336).
func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time) ([]planFacts, []waitFacts, error) { func gatherBus(ctx context.Context, inv *inventory.Inventory) (busFacts, error) {
b, err := pendingBus(ctx, inv)
if err != nil {
return busFacts{}, fmt.Errorf("what a bus upgrade would do cannot be read: %w", err)
}
plans, err := inv.OpenPlans(ctx)
if err != nil {
return busFacts{}, fmt.Errorf("the open plans cannot be read: %w", err)
}
f := busWaitsOf(plans, b, busRefusedFirst.of)
busRefusedFirst.keepOnly(f.heldByTheBus())
return f, nil
}
// gatherPlans is every open plan, its tier's bound from what was measured, and what it waits on; byTheBus are
// the walks held only by the bus's planned step, which S3 leaves to S17.
func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time, byTheBus map[string]bool) ([]planFacts, []waitFacts, error) {
plans, err := inv.OpenPlans(ctx) plans, err := inv.OpenPlans(ctx)
if err != nil { if err != nil {
return nil, nil, fmt.Errorf("the open plans cannot be read: %w", err) return nil, nil, fmt.Errorf("the open plans cannot be read: %w", err)
@@ -466,7 +487,7 @@ func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time) (
bound := bounds.of(p.Repository) bound := bounds.of(p.Repository)
out = append(out, planFacts{id: p.ID, repository: p.Repository, commit: p.Commit, tier: p.Tier, out = append(out, planFacts{id: p.ID, repository: p.Repository, commit: p.Commit, tier: p.Tier,
tiers: len(p.Tiers), entered: inTierSince(p), bound: bound, tiers: len(p.Tiers), entered: inTierSince(p), bound: bound,
waiting: planLineWith(p, now, pause, bound), paused: paused}) waiting: planLineWith(p, now, pause, bound), paused: paused, bus: byTheBus[p.ID]})
} }
return out, waits, nil return out, waits, nil
} }