Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
822e52123c | ||
|
|
9a7ae131cb | ||
|
|
4354d9d7c7 | ||
|
|
1fdc68a8aa |
@@ -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"
|
"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"
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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.",
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
@@ -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 {
|
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)
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1759,7 +1759,6 @@ func ParseManifest(raw []byte) (Manifest, error) {
|
|||||||
}
|
}
|
||||||
problems = append(problems, invokeProblems(m)...)
|
problems = append(problems, invokeProblems(m)...)
|
||||||
problems = append(problems, replacesProblems(m)...)
|
problems = append(problems, replacesProblems(m)...)
|
||||||
problems = append(problems, UnitSettingProblems(m)...)
|
|
||||||
problems = append(problems, endpointNameProblems(m)...)
|
problems = append(problems, endpointNameProblems(m)...)
|
||||||
problems = append(problems, RouteProblems(m)...)
|
problems = append(problems, RouteProblems(m)...)
|
||||||
for _, port := range m.Guards {
|
for _, port := range m.Guards {
|
||||||
|
|||||||
@@ -1,159 +0,0 @@
|
|||||||
package catalogue
|
|
||||||
|
|
||||||
import (
|
|
||||||
"strings"
|
|
||||||
"testing"
|
|
||||||
)
|
|
||||||
|
|
||||||
// A setting may name the unit a service resource holds: a module that stops the distribution's own timer
|
|
||||||
// for the pool the operator names cannot write the pool into its definition (novox/hq ADR 0112), and a
|
|
||||||
// unit name with ${setting:…} left in it is a unit no machine has, so the apply fails far from its cause.
|
|
||||||
|
|
||||||
func scrubber() Manifest {
|
|
||||||
return Manifest{Module: "zfs",
|
|
||||||
Settings: map[string]SettingDeclaration{
|
|
||||||
"scrub-cadence": {Kind: KindPreference, Default: "weekly", Why: "what the distribution's timer did"},
|
|
||||||
},
|
|
||||||
Resources: []map[string]any{
|
|
||||||
{"id": "distribution-weekly", "type": "service", "unit": "zfs-scrub-weekly@${setting:scrub-pool}.timer",
|
|
||||||
"state": "stopped", "boot": "disabled"},
|
|
||||||
{"id": "distribution-monthly", "type": "service", "unit": "zfs-scrub-monthly@${setting:scrub-pool}.timer",
|
|
||||||
"state": "stopped", "boot": "disabled"},
|
|
||||||
{"id": "scrub-config", "type": "file", "path": "/etc/zfs-tools/scrub.conf",
|
|
||||||
"content": "pool=${setting:scrub-pool}\ncadence=${setting:scrub-cadence}\n"},
|
|
||||||
},
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestASettingNamesAServicesUnit(t *testing.T) {
|
|
||||||
m := scrubber()
|
|
||||||
layers := WithDefaults(m, []Layer{{From: "ace", Values: map[string]any{"scrub-pool": "storage", "scrub-cadence": "monthly"}}})
|
|
||||||
for i, want := range []string{"zfs-scrub-weekly@storage.timer", "zfs-scrub-monthly@storage.timer"} {
|
|
||||||
r := map[string]any{}
|
|
||||||
for k, v := range m.Resources[i] {
|
|
||||||
r[k] = v
|
|
||||||
}
|
|
||||||
if err := settingInto(r, layers, m.Module); err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if r["unit"] != want {
|
|
||||||
t.Errorf("composed as %q, not %q", r["unit"], want)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Composed for a machine, the way a push writes it: the operator's pool and cadence reach the units' names
|
|
||||||
// and the file.
|
|
||||||
func TestAComposedMachineGetsTheUnitTheSettingNames(t *testing.T) {
|
|
||||||
r := anAdoptedAnchor()
|
|
||||||
r.Modules = append(r.Modules, scrubber())
|
|
||||||
with := anchorRendering(false)
|
|
||||||
with.Settings["zfs"] = []Layer{{From: "anchor", Values: map[string]any{"scrub-pool": "storage", "scrub-cadence": "monthly"}}}
|
|
||||||
composed, err := r.Compose(with)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatal(err)
|
|
||||||
}
|
|
||||||
if why, left := composed.LeftOut["zfs"]; left {
|
|
||||||
t.Fatalf("left out: %s", why)
|
|
||||||
}
|
|
||||||
got := byID(composed.Resources)
|
|
||||||
for id, want := range map[string]string{"zfs.distribution-weekly": "zfs-scrub-weekly@storage.timer",
|
|
||||||
"zfs.distribution-monthly": "zfs-scrub-monthly@storage.timer"} {
|
|
||||||
if got[id]["unit"] != want {
|
|
||||||
t.Errorf("%s composed as %v, not %s", id, got[id]["unit"], want)
|
|
||||||
}
|
|
||||||
t.Logf("%s: %v", id, got[id]["unit"])
|
|
||||||
}
|
|
||||||
if c := got["zfs.scrub-config"]["content"]; c != "pool=storage\ncadence=monthly\n" {
|
|
||||||
t.Errorf("the file: %q", c)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// A value that would not make a unit's name is refused by name, never written: a space, a slash, a
|
|
||||||
// newline that would begin a directive, or a second suffix.
|
|
||||||
func TestASettingThatMakesNoUnitNameIsRefused(t *testing.T) {
|
|
||||||
m := scrubber()
|
|
||||||
for _, bad := range []string{"", "stor age", "a/b", "storage\nExecStart=/bin/sh", "a@b", "$(x)", "*", `a\\x2d`} {
|
|
||||||
r := map[string]any{}
|
|
||||||
for k, v := range m.Resources[0] {
|
|
||||||
r[k] = v
|
|
||||||
}
|
|
||||||
err := settingInto(r, []Layer{{From: "ace", Values: map[string]any{"scrub-pool": bad}}}, m.Module)
|
|
||||||
if err == nil || !strings.Contains(err.Error(), "scrub-pool") || !strings.Contains(err.Error(), "unit") {
|
|
||||||
t.Errorf("%q: %v", bad, err)
|
|
||||||
}
|
|
||||||
if r["unit"] != m.Resources[0]["unit"] {
|
|
||||||
t.Errorf("%q: the unit was changed on refusal: %v", bad, r["unit"])
|
|
||||||
}
|
|
||||||
}
|
|
||||||
r := map[string]any{}
|
|
||||||
for k, v := range m.Resources[0] {
|
|
||||||
r[k] = v
|
|
||||||
}
|
|
||||||
if err := settingInto(r, nil, m.Module); err == nil || !strings.Contains(err.Error(), "${setting:scrub-pool}") {
|
|
||||||
t.Errorf("nothing set: %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// A key a unit's name asks for is a destination, so setting it is not called stray, and judging the
|
|
||||||
// settings sees it.
|
|
||||||
func TestASettingAUnitAsksForIsNotStrayAndIsJudged(t *testing.T) {
|
|
||||||
m := scrubber()
|
|
||||||
stray := strings.Join(UnusedSettings(m, []Layer{{From: "ace", Values: map[string]any{"scrub-pool": "storage"}}}), "; ")
|
|
||||||
if strings.Contains(stray, "scrub-pool") {
|
|
||||||
t.Errorf("called stray: %s", stray)
|
|
||||||
}
|
|
||||||
if err := JudgeSettings(m, nil, false); err == nil || !strings.Contains(err.Error(), "scrub-pool") {
|
|
||||||
t.Errorf("a unit's setting nothing sets is not refused when judged: %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Third review: a setting may name only a template's instance — right after the `@`, before the final
|
|
||||||
// suffix, and the whole instance — so a value can never make the unit another unit, nor another kind.
|
|
||||||
// Refused where the manifest is read, near its author.
|
|
||||||
func TestASettingInAUnitNameIsOnlyATemplatesInstance(t *testing.T) {
|
|
||||||
ok := []string{"zfs-scrub-weekly@${setting:scrub-pool}.timer", "getty@${setting:tty}.service"}
|
|
||||||
bad := []string{
|
|
||||||
"${setting:unit}",
|
|
||||||
"${setting:name}.timer",
|
|
||||||
"zfs-scrub-${setting:cadence}@storage.timer",
|
|
||||||
"zfs-scrub@${setting:pool}-x.timer",
|
|
||||||
"zfs-scrub@x-${setting:pool}.timer",
|
|
||||||
"zfs-scrub@${setting:pool}.${setting:kind}",
|
|
||||||
"zfs-scrub@${setting:pool}",
|
|
||||||
"zfs-scrub@${setting:a}${setting:b}.timer",
|
|
||||||
}
|
|
||||||
for _, unit := range append(ok, bad...) {
|
|
||||||
m := Manifest{Module: "zfs", Resources: []map[string]any{{"id": "t", "type": "service", "unit": unit, "state": "stopped"}}}
|
|
||||||
err := parsed(t, m)
|
|
||||||
refused := err != nil && strings.Contains(err.Error(), "instance")
|
|
||||||
want := false
|
|
||||||
for _, b := range bad {
|
|
||||||
want = want || b == unit
|
|
||||||
}
|
|
||||||
if refused != want {
|
|
||||||
t.Errorf("%s: refused %v, want %v (%v)", unit, refused, want, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestAnInstanceValueMayNotLeadWithADashNorBeLong(t *testing.T) {
|
|
||||||
m := scrubber()
|
|
||||||
for _, bad := range []string{"-storage", "--help", strings.Repeat("a", 65)} {
|
|
||||||
r := map[string]any{}
|
|
||||||
for k, v := range m.Resources[0] {
|
|
||||||
r[k] = v
|
|
||||||
}
|
|
||||||
err := settingInto(r, []Layer{{From: "ace", Values: map[string]any{"scrub-pool": bad}}}, m.Module)
|
|
||||||
if err == nil || !strings.Contains(err.Error(), "scrub-pool") {
|
|
||||||
t.Errorf("%q: %v", bad, err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
r := map[string]any{}
|
|
||||||
for k, v := range m.Resources[0] {
|
|
||||||
r[k] = v
|
|
||||||
}
|
|
||||||
if err := settingInto(r, []Layer{{From: "ace", Values: map[string]any{"scrub-pool": strings.Repeat("a", 64)}}}, m.Module); err != nil {
|
|
||||||
t.Errorf("64 characters: %v", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -15,7 +15,7 @@ import (
|
|||||||
// (novox/hq issues 122, 134). ADR 0112 names the operator as one of the four providers; this is the
|
// (novox/hq issues 122, 134). ADR 0112 names the operator as one of the four providers; this is the
|
||||||
// operator answering.
|
// operator answering.
|
||||||
//
|
//
|
||||||
// `${setting:<key>}` in a file's content, and in a service's unit name, is filled from the module's settings layers — the mesh's,
|
// `${setting:<key>}` in a file's content is filled from the module's settings layers — the mesh's,
|
||||||
// then this node's — the same layers a mergeable JSON file and a contribution already take, so
|
// then this node's — the same layers a mergeable JSON file and a contribution already take, so
|
||||||
// `settings set <module>` is the one place a person's values go. **Refused when no layer sets it**,
|
// `settings set <module>` is the one place a person's values go. **Refused when no layer sets it**,
|
||||||
// naming the key and the remedy: a definition that carried a default for a mail domain would be
|
// naming the key and the remedy: a definition that carried a default for a mail domain would be
|
||||||
@@ -45,9 +45,6 @@ func settingsUsed(content string) []string {
|
|||||||
// the defaults under the layers with WithDefaults. A value that is not a string is written the way a program would read
|
// the defaults under the layers with WithDefaults. A value that is not a string is written the way a program would read
|
||||||
// it (a number without a trailing .000000, a boolean as true/false).
|
// it (a number without a trailing .000000, a boolean as true/false).
|
||||||
func settingInto(resource map[string]any, layers []Layer, module string) error {
|
func settingInto(resource map[string]any, layers []Layer, module string) error {
|
||||||
if fmt.Sprint(resource["type"]) == "service" {
|
|
||||||
return settingIntoUnit(resource, layers, module)
|
|
||||||
}
|
|
||||||
if fmt.Sprint(resource["type"]) != "file" {
|
if fmt.Sprint(resource["type"]) != "file" {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -71,67 +68,6 @@ func settingInto(resource map[string]any, layers []Layer, module string) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// unitPart is what a setting may put into a unit's name: the characters systemd allows in a unit name,
|
|
||||||
// less the instance's `@` and the escape's `\`, and at least one of them. Anything else — a space, a slash,
|
|
||||||
// a newline that would begin a directive in the unit file the name ends up in — is refused, never written.
|
|
||||||
var unitPart = regexp.MustCompile(`^[A-Za-z0-9:_.][A-Za-z0-9:_.-]{0,63}$`)
|
|
||||||
|
|
||||||
// unitInstance is the one place a setting may stand in a unit's name: the whole instance of a template,
|
|
||||||
// after its `@` and before its suffix — `zfs-scrub-weekly@${setting:scrub-pool}.timer` — so a value can
|
|
||||||
// make the unit another instance of the same template and nothing else (third review of mesh-catalog #147).
|
|
||||||
var unitInstance = regexp.MustCompile(`^[A-Za-z0-9:_.-]+@\$\{setting:[a-z0-9][a-z0-9_.-]*\}\.[a-z]+$`)
|
|
||||||
|
|
||||||
// UnitSettingProblems refuses, where a manifest is read, a service whose unit name carries a setting
|
|
||||||
// anywhere but as a template's whole instance.
|
|
||||||
func UnitSettingProblems(m Manifest) []string {
|
|
||||||
var problems []string
|
|
||||||
for _, r := range m.Resources {
|
|
||||||
if fmt.Sprint(r["type"]) != "service" {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
unit, _ := r["unit"].(string)
|
|
||||||
if len(settingsUsed(unit)) == 0 && !strings.Contains(unit, "${setting:") {
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
if !unitInstance.MatchString(unit) {
|
|
||||||
problems = append(problems, fmt.Sprintf("%s: the service %v's unit %q carries a setting outside a template's "+
|
|
||||||
"instance; a setting may stand only as the whole instance, after the @ and before the suffix "+
|
|
||||||
"(name@${setting:key}.timer)", m.Module, r["id"], unit))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return problems
|
|
||||||
}
|
|
||||||
|
|
||||||
// settingIntoUnit fills ${setting:…} in a service resource's unit name: a module that holds the
|
|
||||||
// distribution's timer for the pool the operator names cannot write the pool into its definition
|
|
||||||
// (novox/hq ADR 0112). Refused, with the key and why, when nothing sets it or the value would not make a
|
|
||||||
// unit's name; the resource is left as it was.
|
|
||||||
func settingIntoUnit(resource map[string]any, layers []Layer, module string) error {
|
|
||||||
unit, ok := resource["unit"].(string)
|
|
||||||
if !ok {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
for _, key := range settingsUsed(unit) {
|
|
||||||
value, set := settingValue(layers, key)
|
|
||||||
if !set {
|
|
||||||
return fmt.Errorf(
|
|
||||||
"%s has a service whose unit says ${setting:%s}, and nothing sets %q for it — an operator's "+
|
|
||||||
"value is the assignment's, never the definition's (novox/hq ADR 0112): "+
|
|
||||||
"`settings set %s <file>` with {%q: …}%s",
|
|
||||||
module, key, key, module, key, orNoSettings(layers))
|
|
||||||
}
|
|
||||||
v := plainly(value)
|
|
||||||
if !unitPart.MatchString(v) {
|
|
||||||
return fmt.Errorf("%s: the setting %q is %q, which cannot be the instance of the unit %s: an instance "+
|
|
||||||
"takes letters, digits, ':', '_', '.' and '-', does not begin with '-' (a systemctl option), and is "+
|
|
||||||
"at most 64 characters", module, key, v, unit)
|
|
||||||
}
|
|
||||||
unit = strings.ReplaceAll(unit, "${setting:"+key+"}", v)
|
|
||||||
}
|
|
||||||
resource["unit"] = unit
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func settingValue(layers []Layer, key string) (any, bool) {
|
func settingValue(layers []Layer, key string) (any, bool) {
|
||||||
var value any
|
var value any
|
||||||
set := false
|
set := false
|
||||||
@@ -170,15 +106,11 @@ func settingKeysUsedBy(m Manifest) map[string]bool {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
for _, r := range m.Resources {
|
for _, r := range m.Resources {
|
||||||
switch fmt.Sprint(r["type"]) {
|
if fmt.Sprint(r["type"]) != "file" {
|
||||||
case "file":
|
continue
|
||||||
if content, ok := r["content"].(string); ok {
|
}
|
||||||
note(content)
|
if content, ok := r["content"].(string); ok {
|
||||||
}
|
note(content)
|
||||||
case "service":
|
|
||||||
if unit, ok := r["unit"].(string); ok {
|
|
||||||
note(unit)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
inValues := func(values map[string]any) {
|
inValues := func(values map[string]any) {
|
||||||
|
|||||||
Reference in New Issue
Block a user