Merge pull request 'Tell the operator at once when sends wait for the bus's planned step, and say it before the merge (hq issue 336)' (#173) from fix/336-bus-step-waiting into main
This commit was merged in pull request #173.
This commit is contained in:
@@ -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] + "…"
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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.",
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user