Files
mesh-controller/cmd/mesh-controller/tier_send_test.go
T
jochen a44dc01c65
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer delivery to the same trunk took over its walk
Excuse no wait for a move from a build not known, and count a module put back only once there is one (hq issue 318 review)
2026-10-08 15:49:07 +02:00

438 lines
18 KiB
Go

package main
import (
"context"
"encoding/json"
"fmt"
"reflect"
"strings"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/lease"
)
// A plan's tier is sent to each machine once, and judged there by one gate (novox/hq issue 281).
//
// On 2026-10-06 a merge rebuilt about seventy modules in one tier, and the plan sent the first machine
// one module at a time — one whole declaration per module, seconds apart. The node-engine set aside one
// declaration for the next, the node tools never settled, the mesh said the machine's core was behind,
// and the gate failed the module it happened to be judging for that machine-level condition; it then
// could not put the module back, having read as "the build it ran before" the very build an earlier send
// of the same tier had already carried there.
// tierMesh is a mesh with every module of one tier running on anchor and laptop at build c1, a newer
// build c2 of each registered, a plan whose one tier built them all, and every send recorded and
// answered: the machine applies what it is sent and reports it. A judging reads the store's reports, the
// node tools as answering everywhere, and the conditions the test opens.
type tierMesh struct {
open *stores
modules []string
sent [][]string
conds []conditions.Condition
keeper *conditions.Keeper
}
func aTierMesh(t *testing.T, modules ...string) *tierMesh {
t.Helper()
open := aMesh(t)
ctx := t.Context()
inv := open.inventory
tm := &tierMesh{open: open, modules: modules}
tm.keeper, _ = withConditionsInMemory(t)
was := doctorFrom
doctorFrom = &doctor{open: open, keeper: tm.keeper, teller: &conditions.Told{}}
t.Cleanup(func() { doctorFrom = was })
register := func(m, commit string, asked time.Time) {
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: m, Version: commit},
inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/" + m, BuiltFrom: commit,
Head: commit, Asked: asked}); err != nil {
t.Fatal(err)
}
}
first := map[string]string{}
for _, m := range modules {
for i, commit := range []string{"c1", "c2"} {
at := time.Now().Add(-2 * time.Hour)
if commit == "c2" {
at = time.Now().Add(-time.Minute)
}
b := inventory.Build{ID: fmt.Sprintf("build-%s-%d", m, i+1), Module: m, Commit: commit,
Repository: "novox/mesh-catalog", Path: "modules/" + m, Asked: at, At: at}
b.Manifest, _ = json.Marshal(catalogue.Manifest{Module: m, Version: commit})
if err := inv.RecordBuild(ctx, b); err != nil {
t.Fatal(err)
}
}
register(m, "c1", time.Now().Add(-2*time.Hour))
first[m] = "c1"
for _, n := range []string{"anchor", "laptop"} {
if _, err := inv.Assign(ctx, n, m); err != nil {
t.Fatal(err)
}
}
}
for _, n := range []string{"anchor", "laptop"} {
if err := inv.RecordSent(ctx, nodeID(t, open, n), "d-"+n+"-c1", first); err != nil {
t.Fatal(err)
}
}
for _, m := range modules {
register(m, "c2", time.Now().Add(-time.Minute))
}
n := 0
wasSend := sendRollout
sendRollout = func(ctx context.Context, open *stores, names []string) ([]string, error) {
tm.sent = append(tm.sent, append([]string(nil), names...))
current, err := open.inventory.CurrentBuilds(ctx)
if err != nil {
return nil, err
}
builds := map[string]string{}
for _, m := range tm.modules {
builds[m] = current[m].Commit
}
for _, node := range names {
n++
digest := fmt.Sprintf("d-%s-%d", node, n)
if err := open.inventory.RecordSent(ctx, nodeID(t, open, node), digest, builds); err != nil {
return nil, err
}
if _, err := open.inventory.RecordDoing(ctx, nodeID(t, open, node), inventory.Doing{Node: node,
Outcome: inventory.OutcomeApplied, Declared: digest, Applied: 1, At: time.Now()}); err != nil {
return nil, err
}
}
return names, nil
}
t.Cleanup(func() { sendRollout = wasSend })
wasGather := gatherGateFacts
gatherGateFacts = func(ctx context.Context, open *stores, component string) (gateFacts, error) {
f := gateFacts{now: time.Now(), reports: map[string]inventory.Reported{}, engines: map[string]string{},
rolledBack: map[string][]lease.Rollback{}, served: map[string]served{}, judged: true,
open: append([]conditions.Condition(nil), tm.conds...)}
reports, err := open.inventory.LastReports(ctx)
if err != nil {
return f, err
}
for _, r := range reports {
f.reports[r.Node] = r
f.served[r.Node] = served{runtime: true, tools: map[string]bool{}}
}
return f, nil
}
t.Cleanup(func() { gatherGateFacts = wasGather })
wasSettle, wasEvery, wasBound := gateSettle, gateEvery, gateBound
// One judging per advance until a test says otherwise: a pass waits gateEvery for the next.
gateSettle, gateEvery = 0, time.Hour
t.Cleanup(func() { gateSettle, gateEvery, gateBound = wasSettle, wasEvery, wasBound })
built := time.Now().UTC()
states := map[string]*inventory.PlanModule{}
for _, m := range modules {
states[m] = &inventory.PlanModule{State: "built", BuiltAt: &built, Commit: "c2", Build: "build-" + m + "-2"}
}
plan := inventory.Plan{ID: "plan-tier", Repository: "novox/mesh-catalog", Branch: "main", Commit: "c2",
Created: built, State: inventory.PlanBuilding, Tiers: [][]string{modules}, Modules: states}
if err := inv.SavePlan(ctx, &plan); err != nil {
t.Fatal(err)
}
return tm
}
func (tm *tierMesh) plan(t *testing.T) inventory.Plan {
t.Helper()
p, err := tm.open.inventory.PlanByID(t.Context(), "plan-tier")
if err != nil {
t.Fatal(err)
}
return p
}
// raise opens a condition about a machine itself, raised now: after the send.
func (tm *tierMesh) raise(machine, kind, summary string) {
tm.conds = append(tm.conds, conditions.Condition{Key: "machine." + machine + "." + kind, Kind: kind,
Subject: conditions.Subject{Scope: conditions.ScopeMachine, ID: machine, Machine: machine},
Raised: time.Now().Add(time.Second), Summary: summary, Source: "D10"})
}
// The incident's tier, small: a dozen modules whose first machine is the same are sent there in one
// send, judged by one gate, and the rest of their machines are sent once — two sends, not two dozen.
func TestATierOfManyModulesIsSentToEachMachineOnce(t *testing.T) {
var modules []string
for i := 1; i <= 12; i++ {
modules = append(modules, fmt.Sprintf("m%02d", i))
}
tm := aTierMesh(t, modules...)
ctx := t.Context()
advancePlans(ctx, tm.open)
if !reflect.DeepEqual(tm.sent, [][]string{{"anchor"}}) {
t.Fatalf("sent %v: the first machine once, for the whole tier", tm.sent)
}
p := tm.plan(t)
lead := p.Modules["m01"]
if lead.Gate == nil || len(lead.Gate.Carried) != len(modules) || lead.GatedBy != "" {
t.Fatalf("the gate on the first module does not judge the whole send: %+v", lead.Gate)
}
for _, m := range modules {
s := p.Modules[m]
if s.FirstAt == nil || !reflect.DeepEqual(s.First, []string{"anchor"}) || s.Previous != "c1" {
t.Fatalf("%s's first send is not recorded, or what anchor ran before is not c1: %+v", m, s)
}
if m != "m01" && s.GatedBy != "m01" {
t.Fatalf("%s is not gated by the send's gate: %+v", m, s)
}
}
gateEvery = 0
for i := 0; i < 6; i++ {
advancePlans(ctx, tm.open)
}
p = tm.plan(t)
if !reflect.DeepEqual(tm.sent, [][]string{{"anchor"}, {"laptop"}}) {
t.Fatalf("sent %v: one send per machine for the tier", tm.sent)
}
if p.State != inventory.PlanDone {
t.Fatalf("the plan is %s: %s", p.State, p.Note)
}
for _, m := range modules {
if g := p.Modules[m].Gate; g == nil || g.Verdict != inventory.GatePassed {
t.Fatalf("%s's gate: %+v", m, g)
}
if v, found, err := tm.open.inventory.GateOf(ctx, "build-"+m+"-2"); err != nil || !found || v.Verdict != inventory.GatePassed {
t.Fatalf("%s's pass was not kept: %+v %v %v", m, v, found, err)
}
}
}
// The core behind about the node tools, while the node tools are themselves among what was sent and
// settling, is the node tools' — inside their settle window, the gate's bound — and no failure of the
// modules sent with them: at the bound only the node tools are put back.
func TestACoreBehindWhileTheNodeToolsSettleFailsNoOtherModule(t *testing.T) {
tm := aTierMesh(t, "app1", "app2", broker.RuntimeModule)
ctx := t.Context()
inv := tm.open.inventory
advancePlans(ctx, tm.open)
if len(tm.sent) != 1 {
t.Fatalf("sent %v", tm.sent)
}
tm.raise("anchor", kindCoreBehind, "anchor runs the node tools it was last sent, not yet reported applied, "+
"and no plan is rolling them out")
gateEvery = 0
advancePlans(ctx, tm.open)
g := tm.plan(t).Modules["app1"].Gate
if g.Verdict != "" || !reflect.DeepEqual(g.Failing, []string{broker.RuntimeModule}) {
t.Fatalf("inside the settle window the gate is %+v: only the node tools wait on it", g)
}
gateBound = -time.Second
advancePlans(ctx, tm.open)
p := tm.plan(t)
g = p.Modules["app1"].Gate
if p.State != inventory.PlanFailed || g.Verdict != inventory.GateFailed ||
!reflect.DeepEqual(g.Failing, []string{broker.RuntimeModule}) {
t.Fatalf("the plan is %s (%s), the gate %+v", p.State, p.Note, g)
}
current, _ := inv.CurrentBuilds(ctx)
if current[broker.RuntimeModule].Commit != "c1" {
t.Fatalf("the node tools were not put back: %s", current[broker.RuntimeModule].Commit)
}
for _, m := range []string{"app1", "app2"} {
if current[m].Commit != "c2" {
t.Fatalf("%s was put back for the node tools' condition: %s", m, current[m].Commit)
}
if failed, _ := inv.GateFailed(ctx, "build-"+m+"-2"); failed {
t.Fatalf("%s's build was marked failed for the node tools' condition", m)
}
}
// Healthy for the passes asked, the module the gate is kept on keeps its own pass (novox/hq issue 318
// review): it is not blamed, and not left without a verdict.
if !strings.Contains(p.Note, "app1 passed on its own and is kept") {
t.Fatalf("the plan blames the module its gate was kept on, or leaves it unjudged: %s", p.Note)
}
if v, found, _ := inv.GateOf(ctx, "build-app1-2"); !found || v.Verdict != inventory.GatePassed {
t.Fatalf("app1's own pass was not kept: %+v", v)
}
}
// A machine-level condition about anything else holds the machine back as a whole: everything the send
// moved there fails together, once, with that reason — not one module for it.
func TestAMachineUnhealthyAsAWholeFailsTheWholeSendTogether(t *testing.T) {
tm := aTierMesh(t, "app1", "app2", "app3")
ctx := t.Context()
inv := tm.open.inventory
advancePlans(ctx, tm.open)
tm.raise("anchor", "disk-full", "anchor's root file system is full")
gateEvery, gateBound = 0, -time.Second
advancePlans(ctx, tm.open)
p := tm.plan(t)
g := p.Modules["app1"].Gate
if p.State != inventory.PlanFailed || !strings.Contains(g.Why, "anchor as a whole") ||
!reflect.DeepEqual(g.Failing, []string{"app1", "app2", "app3"}) {
t.Fatalf("the plan is %s (%s), the gate %+v", p.State, p.Note, g)
}
current, _ := inv.CurrentBuilds(ctx)
for _, m := range []string{"app1", "app2", "app3"} {
if current[m].Commit != "c1" {
t.Fatalf("%s was not put back with the send: %s", m, current[m].Commit)
}
}
if !reflect.DeepEqual(tm.sent, [][]string{{"anchor"}, {"anchor"}}) {
t.Fatalf("sent %v: the first send and one send putting all of it back, and nothing to laptop", tm.sent)
}
if !strings.Contains(p.Note, "put back app1 to c1, app2 to c1, app3 to c1 on anchor, in one send") {
t.Fatalf("the plan does not say what was put back: %s", p.Note)
}
}
// The machine-level conditions, sorted: one naming a module is that module's; the core behind is the
// core's when the send moved it and nobody's when it did not; anything else is the machine's as a whole.
func TestTheConditionsAboutAMachineItself(t *testing.T) {
since := time.Now().Add(-time.Minute)
about := func(kind, id string) conditions.Condition {
return conditions.Condition{Key: "machine." + id + "." + kind, Kind: kind, Raised: time.Now(),
Subject: conditions.Subject{Scope: conditions.ScopeMachine, ID: id, Machine: "anchor"}}
}
f := gateFacts{judged: true, open: []conditions.Condition{about(kindCoreBehind, "anchor")}}
if w := aboutTheMachine("anchor", []string{"app", "blueman"}, since, f); w.whole != "" || len(w.on) != 0 || len(w.facts.open) != 0 {
t.Errorf("the core behind, with no core component sent, counted against what was sent: %+v", w)
}
if w := aboutTheMachine("anchor", []string{"app", broker.RuntimeModule}, since, f); w.whole != "" ||
w.on[broker.RuntimeModule] == "" || w.on["app"] != "" {
t.Errorf("the core behind, with the node tools sent, is not theirs alone: %+v", w)
}
f.open = []conditions.Condition{about("x", "anchor.app")}
if w := aboutTheMachine("anchor", []string{"app", "other"}, since, f); w.on["app"] == "" || w.on["other"] != "" || w.whole != "" {
t.Errorf("a condition naming a module is not that module's: %+v", w)
}
f.open = []conditions.Condition{about("unreachable", "anchor")}
if w := aboutTheMachine("anchor", []string{"app"}, since, f); !strings.Contains(w.whole, "anchor as a whole") {
t.Errorf("a condition about the machine is not the machine's as a whole: %+v", w)
}
f.open[0].Raised = since.Add(-time.Hour)
if w := aboutTheMachine("anchor", []string{"app"}, since, f); w.whole != "" || len(w.facts.open) != 1 {
t.Errorf("a condition older than the send counted: %+v", w)
}
}
// A module its first send changed nothing of — an earlier send already carried this very build there —
// fails a gate as no verdict on its build: left as it was, never marked failed, nothing said as a
// rollback that failed; and `plans retry` asks it again rather than refusing a build that was never
// judged (the incident's "NOT put back: no build kept").
func TestAModuleItsSendChangedNothingOfIsLeftAsItWasAndRetried(t *testing.T) {
g := aGateMesh(t)
ctx := t.Context()
inv := g.open.inventory
// An earlier send carried c2 to anchor before the plan's own.
if err := inv.RecordSent(ctx, nodeID(t, g.open, "anchor"), "d-anchor-early", map[string]string{"app": "c2"}); err != nil {
t.Fatal(err)
}
advancePlans(ctx, g.open)
if p := g.plan(t); p.Modules["app"].Previous != "c2" {
t.Fatalf("what anchor ran before the send: %+v", p.Modules["app"])
}
g.health["anchor"] = healthNotYet
gateEvery, gateBound = 0, -time.Second
advancePlans(ctx, g.open)
p := g.plan(t)
gate := p.Modules["app"].Gate
if p.State != inventory.PlanFailed || gate.Verdict != inventory.GateFailed || gate.Rollback != gateUnchanged ||
!strings.Contains(p.Note, "app was left as it was") {
t.Fatalf("the plan is %s (%s), the gate %+v", p.State, p.Note, gate)
}
if failed, _ := inv.GateFailed(ctx, "build-2"); failed {
t.Fatal("a build its send changed nothing of was marked failed")
}
if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c2" {
t.Fatalf("the module was moved: %s", current["app"].Commit)
}
open, _ := g.keeper.Open(ctx)
for _, c := range open {
if c.Kind == kindRollbackFailed || c.Kind == kindRolledBack {
t.Fatalf("said as a rollback: %+v", c)
}
}
was := askABuild
askABuild = func(context.Context, buildSource, string, string) (string, error) { return "build-3", nil }
t.Cleanup(func() { askABuild = was })
g.health["anchor"] = healthGood
said, err := retryPlan(ctx, g.open, "plan-gate")
if err != nil {
t.Fatalf("a plan stopped at a gate that judged no build was not retried: %v", err)
}
p = g.plan(t)
s := p.Modules["app"]
if p.State != inventory.PlanBuilding || s.State != "asked" || s.Build != "build-3" || s.FirstAt != nil || s.Gate != nil {
t.Fatalf("retried as %q: the plan is %s, app %+v", said, p.State, s)
}
}
// **In a plan's tier path, a broken module is put back at once too** (novox/hq issue 318 review): the node
// tools' witness reverts them on the first machine. Whether the gate is kept on them or on the module beside
// them, they are registered back and marked in the judging that finds them broken, the plan goes on judging
// the other module to its own pass, and only then fails.
func TestATierPutsABrokenModuleBackAtOnce(t *testing.T) {
for _, order := range [][]string{{broker.RuntimeModule, "app1"}, {"app1", broker.RuntimeModule}} {
t.Run("gate kept on "+order[0], func(t *testing.T) {
tm := aTierMesh(t, order...)
ctx := t.Context()
inv := tm.open.inventory
tierFacts := gatherGateFacts
var seen []string
gatherGateFacts = func(ctx context.Context, open *stores, component string) (gateFacts, error) {
f, err := tierFacts(ctx, open, component)
current, _ := open.inventory.CurrentBuilds(ctx)
seen = append(seen, current[broker.RuntimeModule].Commit)
if current[broker.RuntimeModule].Commit == "c2" {
f.rolledBack["anchor"] = []lease.Rollback{{Component: lease.ComponentNodeTools,
Outcome: lease.OutcomeRolledBack, From: "c2", To: "c1", At: time.Now(), Why: "the node tools did not answer"}}
}
return f, err
}
gateEvery, gateSettle, gateBound = 0, 100*time.Millisecond, 10*time.Second
for i := 0; i < 20 && len(seen) < 2; i++ {
advancePlans(ctx, tm.open)
}
if len(seen) < 2 || seen[0] != "c2" || seen[1] != "c1" {
t.Fatalf("the node tools' registered build at each judging: %v; want c2, then c1 at once", seen)
}
if failed, _ := inv.GateFailed(ctx, "build-"+broker.RuntimeModule+"-2"); !failed {
t.Fatal("the node tools' build is not marked failed at once")
}
if p := tm.plan(t); p.State == inventory.PlanFailed {
t.Fatalf("the plan failed at once: %s; want it judging app1 first", p.Note)
}
for deadline := time.Now().Add(5 * time.Second); time.Now().Before(deadline); {
advancePlans(ctx, tm.open)
if tm.plan(t).State == inventory.PlanFailed {
break
}
time.Sleep(20 * time.Millisecond)
}
p := tm.plan(t)
if p.State != inventory.PlanFailed {
t.Fatalf("the plan is %s: %s; want failed, for the node tools", p.State, p.Note)
}
if v, found, _ := inv.GateOf(ctx, "build-app1-2"); !found || v.Verdict != inventory.GatePassed {
t.Fatalf("app1's verdict: %+v; want its own pass (plan: %s)", v, p.Note)
}
if current, _ := inv.CurrentBuilds(ctx); current["app1"].Commit != "c2" || current[broker.RuntimeModule].Commit != "c1" {
t.Fatalf("registered app1 %s, node tools %s; want c2 and c1", current["app1"].Commit,
current[broker.RuntimeModule].Commit)
}
})
}
}