diff --git a/cmd/mesh-controller/bus_waits.go b/cmd/mesh-controller/bus_waits.go new file mode 100644 index 00000000..0c02e870 --- /dev/null +++ b/cmd/mesh-controller/bus_waits.go @@ -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..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] + "…" +} diff --git a/cmd/mesh-controller/bus_waits_test.go b/cmd/mesh-controller/bus_waits_test.go new file mode 100644 index 00000000..2a462b21 --- /dev/null +++ b/cmd/mesh-controller/bus_waits_test.go @@ -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) + } +} diff --git a/cmd/mesh-controller/changeplan.go b/cmd/mesh-controller/changeplan.go index 9a866888..7a983c97 100644 --- a/cmd/mesh-controller/changeplan.go +++ b/cmd/mesh-controller/changeplan.go @@ -5,6 +5,7 @@ import ( "sort" "strings" + "github.com/novox/mesh-controller/internal/catalogue" "github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/link" ) @@ -50,7 +51,13 @@ func changePlanOf(repository, base, head string, r mergeReach, entries []invento } } 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") case waits && len(e.On) > 0: 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 } +// 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. func summaryOf(p link.ChangePlan) string { if len(p.Moved) == 0 && len(p.New) == 0 { @@ -110,7 +120,10 @@ func summaryOf(p link.ChangePlan) string { } bus := "no bus step" 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" } } diff --git a/cmd/mesh-controller/changeplan_test.go b/cmd/mesh-controller/changeplan_test.go index a2632b71..599cb20e 100644 --- a/cmd/mesh-controller/changeplan_test.go +++ b/cmd/mesh-controller/changeplan_test.go @@ -45,7 +45,7 @@ func TestAChangePlanSaysWhatEachMachineReceives(t *testing.T) { } 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))") { 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) } 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) { t.Errorf("the steps do not say %q:\n%s", want, text) } diff --git a/cmd/mesh-controller/plain_words.go b/cmd/mesh-controller/plain_words.go index faec2bbb..3d6d4ed3 100644 --- a/cmd/mesh-controller/plain_words.go +++ b/cmd/mesh-controller/plain_words.go @@ -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.", 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 { return words{Headline: "The bus upgrade failed", Needs: "decide whether to put the bus back to the version before; the details say how.", diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 8688705b..b2bd68e7 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -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 // 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" + // 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 { fmt.Printf("%s: %v\n", p.ID, err) if err := inv.SavePlan(ctx, p); err != nil { diff --git a/cmd/mesh-controller/replay336_test.go b/cmd/mesh-controller/replay336_test.go new file mode 100644 index 00000000..1b99cf71 --- /dev/null +++ b/cmd/mesh-controller/replay336_test.go @@ -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) + } +} diff --git a/cmd/mesh-controller/signals.go b/cmd/mesh-controller/signals.go index a337b0f1..1dcc3f20 100644 --- a/cmd/mesh-controller/signals.go +++ b/cmd/mesh-controller/signals.go @@ -210,6 +210,20 @@ var signalsTable = []signalRow{ newest: func(f *signalFacts) time.Time { 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 @@ -315,7 +329,9 @@ func watchReports(f *signalFacts) []conditions.Observation { func watchPlans(f *signalFacts) []conditions.Observation { var out []conditions.Observation 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 } in := f.now.Sub(p.entered) diff --git a/cmd/mesh-controller/signals_test.go b/cmd/mesh-controller/signals_test.go index 9c7f4a1d..7d34dee2 100644 --- a/cmd/mesh-controller/signals_test.go +++ b/cmd/mesh-controller/signals_test.go @@ -142,6 +142,19 @@ var suppressions = map[string]suppression{ 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. "S15": { inside: func(f *signalFacts) { diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index 88bf4434..99437153 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -97,6 +97,9 @@ type signalFacts struct { // facts is the snapshot this controller keeps for merge checks (S14). facts factsFacts + + bus busFacts + busErr error } // 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 waiting string paused bool + bus bool } // 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.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.mergesPassed, f.merges, f.mergesErr = watchedMerges.last() 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 } -// gatherPlans is every open plan, its tier's bound from what was measured, and what it waits on. -func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time) ([]planFacts, []waitFacts, error) { +// gatherBus is a new bus build waiting for its step and the walks refused for it (S17, novox/hq issue 336). +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) if err != nil { 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) 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, - waiting: planLineWith(p, now, pause, bound), paused: paused}) + waiting: planLineWith(p, now, pause, bound), paused: paused, bus: byTheBus[p.ID]}) } return out, waits, nil }