From 2421b82ad2864c794f027308cfefcea6d090573f Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 22:22:03 +0200 Subject: [PATCH] Keep a named push from sending builds a policy or a plan holds back A named push flushed every other machine whose declaration differed from what it was last sent (hq ADR 0083). Under an upgrade policy of `record`, or a plan still waiting on its first machine (ADR 0218), every machine running the module differs, so `push ` sent the held build to all of them (hq issue 259). Each send now records which build of each module it carried (node.sent_builds, migration 0061). The cascade, and the bus holder added to a named push, skip a machine any of whose modules would move to a build its policy records or an open plan has not sent it, and say which module, which build, why, and that `push ` sends it. A machine whose last send was not recorded is held until it is named. The named machine itself, a whole-mesh push and `push --behind` are unchanged. --- cmd/mesh-controller/adopting_test.go | 2 +- cmd/mesh-controller/held_back.go | 241 +++++++++++++ cmd/mesh-controller/held_back_test.go | 335 ++++++++++++++++++ cmd/mesh-controller/issue_204_test.go | 2 +- cmd/mesh-controller/plan.go | 44 ++- cmd/mesh-controller/push.go | 113 +++--- cmd/mesh-controller/sendable.go | 4 + internal/inventory/catalogue.go | 31 ++ internal/inventory/doing_test.go | 4 +- ...1-a-send-records-the-builds-it-carried.sql | 14 + internal/inventory/nodes.go | 41 ++- internal/inventory/sent_builds_test.go | 81 +++++ internal/link/heard_test.go | 2 +- 13 files changed, 831 insertions(+), 83 deletions(-) create mode 100644 cmd/mesh-controller/held_back.go create mode 100644 cmd/mesh-controller/held_back_test.go create mode 100644 internal/inventory/migrations/0061-a-send-records-the-builds-it-carried.sql create mode 100644 internal/inventory/sent_builds_test.go diff --git a/cmd/mesh-controller/adopting_test.go b/cmd/mesh-controller/adopting_test.go index 73c94aa..3589fec 100644 --- a/cmd/mesh-controller/adopting_test.go +++ b/cmd/mesh-controller/adopting_test.go @@ -85,7 +85,7 @@ func reportsReaching(t *testing.T, open *stores, reachable []link.Reach, held .. if err != nil { t.Fatal(err) } - if err := open.inventory.RecordSent(ctx, record.ID, digestOf(body)); err != nil { + if err := open.inventory.RecordSent(ctx, record.ID, digestOf(body), nil); err != nil { t.Fatal(err) } if _, err := (link.Enrolment{Inventory: open.inventory}).Heard(ctx, link.Report{ diff --git a/cmd/mesh-controller/held_back.go b/cmd/mesh-controller/held_back.go new file mode 100644 index 0000000..2d83545 --- /dev/null +++ b/cmd/mesh-controller/held_back.go @@ -0,0 +1,241 @@ +package main + +import ( + "context" + "fmt" + "io" + "sort" + "strings" + + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" +) + +// A push sends no build a policy or a plan holds back, except to the machine it names (novox/hq +// issue 259, ADR 0221). +// +// A named push ends by sending every other machine whose declaration differs from what it was last +// sent (ADR 0083), so that a grant the push's work minted reaches the provider in the same act. A +// digest cannot say why a machine differs. A module whose upgrade policy records rather than rolls +// out makes every machine running it differ from the merge on, and so did a module whose plan was +// still waiting on its first machine (ADR 0218): `push ` sent all four machines the build +// that was meant to be walked through the mesh one machine at a time, and a fault in it was met +// everywhere at once. +// +// What tells the two apart is which build of each module the machine was last sent, kept with every +// send. A machine any of whose modules would move to a build its policy or an open plan holds back +// is not sent by a push that did not name it; the push says which, and why, and how to send it. + +// heldBack is why a push that did not name a machine must not send it, empty when it may. +// +// `modules` is what the machine would be sent now; `sent` and `known` what it was last sent, as +// Inventory.SentBuilds answers. A module moves when the build it would carry is not the one the +// machine was last sent — including a module the machine was never sent at all. A move is held when +// the module's policy records rather than rolls out, or when an open plan has not yet sent this +// machine (planStillToSend). A machine whose last send was not recorded is held whole: what it carried +// is not known, so a held upgrade cannot be told from anything else. +func heldBack(node string, modules []string, sent map[string]string, known bool, + current map[string]inventory.CurrentBuild, plans []inventory.Plan) []string { + if !known { + return []string{"which builds it was last sent is not known — it was last sent before the " + + "mesh kept them, or sent a declaration by hand"} + } + var why []string + for _, m := range modules { + now := current[m] + was, carried := sent[m] + if carried && was == now.Commit { + continue + } + move := fmt.Sprintf("%s would move %sto %s", m, fromBuild(was, carried), buildName(now.Commit)) + if !now.RollOut { + why = append(why, move+", which its upgrade policy records rather than rolls out") + continue + } + if id := planStillToSend(plans, m, node); id != "" { + why = append(why, move+", which "+id+" has not sent it yet (one machine first)") + } + } + sort.Strings(why) + return why +} + +// planStillToSend is the open plan that has a module's new build still to send this machine, or empty: +// one holding the module that has neither finished sending it nor sent it here first, and has not +// failed it (novox/hq ADR 0218). The same reading rolledOutByAPlan makes for the whole module, made +// per machine. +func planStillToSend(plans []inventory.Plan, module, node string) string { + for _, p := range plans { + s, holds := p.Modules[module] + if !p.Open() || !holds { + continue + } + if s == nil { + return p.ID + } + if s.SentAt != nil || s.State == "failed" { + continue + } + first := false + for _, n := range s.First { + if n == node { + first = true + } + } + if !first { + return p.ID + } + } + return "" +} + +func fromBuild(was string, carried bool) string { + if !carried { + return "(never sent it) " + } + return "from " + buildName(was) + " " +} + +func buildName(commit string) string { + if commit == "" { + return "a build with no source" + } + return shortCommit(commit) +} + +// heldMachines reads, for each machine named, why a push that did not name it must not send it +// (heldBack), and answers only the machines held. A machine whose set cannot be worked out is left +// to the send, which says why. +func heldMachines(ctx context.Context, open *stores, names []string) (map[string][]string, error) { + out := map[string][]string{} + if len(names) == 0 { + return out, nil + } + inv := open.inventory + current, err := inv.CurrentBuilds(ctx) + if err != nil { + return nil, err + } + plans, err := inv.OpenPlans(ctx) + if err != nil { + return nil, err + } + for _, node := range names { + plan, _, err := planFor(ctx, open, node) + if err != nil { + continue + } + modules := make([]string, 0, len(plan.Modules)) + for _, m := range plan.Modules { + modules = append(modules, m.Module) + } + sent, known, err := inv.SentBuilds(ctx, node) + if err != nil { + return nil, err + } + if why := heldBack(node, modules, sent, known, current, plans); len(why) > 0 { + out[node] = why + } + } + return out, nil +} + +// sayHeld is what a push says about a machine it left behind on purpose: that it is behind, why it +// was not sent, that whatever else it is owed waits with it, and the command that sends it. +func sayHeld(w io.Writer, node string, why []string) { + fmt.Fprintf(w, "\n%s is behind and was not sent: %s. A push sends no build a policy or a plan "+ + "holds back to a machine it did not name (novox/hq ADR 0221), so anything else it is owed — a "+ + "grant from this push among it — waits with it. `push %s` sends it\n", + node, strings.Join(why, "; "), node) +} + +// flushBehind is the end of a named push: every other machine now behind is sent too, by name, over +// as many rounds as the sends take to settle (novox/hq issue 057, ADR 0083) — except a machine whose +// modules would move to a build a policy or a plan holds back, which is named and left (ADR 0221). +// +// `handled` is every machine already sent or already said; it is not considered again. Answers the +// machines that could not be composed, as refusals. +func flushBehind(ctx context.Context, open *stores, nodes []inventory.Node, handled map[string]bool, + compose func(held context.Context, node string) (sendable, error), d delivery, holder string, + w io.Writer) ([]string, error) { + inv := open.inventory + var refusals []string + // Bounded by the node count: a node is marked handled the round it is considered and is never + // considered twice, so the loop cannot run more than len(nodes) rounds. The bound is a guard + // against a logic error, not a real limit — if it were ever hit, that is a bug rather than a + // cascade legitimately still converging, so it is said rather than passed over in silence. + rounds := 0 + for { + would, err := wouldSend(ctx, open, nodes) + if err != nil { + return refusals, err + } + behind, err := inv.Waiting(ctx, would) + if err != nil { + return refusals, err + } + var also []string + for _, m := range behind { + if !handled[m.Node] { + also = append(also, m.Node) + } + } + if len(also) == 0 { + return refusals, nil + } + if rounds++; rounds > len(nodes) { + fmt.Fprintf(w, "\nstopped cascading after %d rounds with %s still behind — this "+ + "should not happen; run `push --behind` to finish\n", + rounds-1, strings.Join(also, ", ")) + return refusals, nil + } + sort.Strings(also) + held, err := heldMachines(ctx, open, also) + if err != nil { + return refusals, err + } + var sending []string + for _, name := range also { + // Every candidate this round is marked handled — the sent ones so they are not + // re-listed, the held ones because they stay held, and the refused ones so a machine + // that cannot be composed does not make the loop spin on it for ever. + handled[name] = true + if why, isHeld := held[name]; isHeld { + sayHeld(w, name, why) + continue + } + sending = append(sending, name) + } + if len(sending) == 0 { + continue + } + fmt.Fprintf(w, "\nthis push left %s behind — a provision granted from there, or a "+ + "declaration since changed; sending it too\n", strings.Join(sending, ", ")) + // Tolerantly, exactly as the named send: a machine that cannot be composed is collected as + // a refusal and reported at the end, and the others are still sent (novox/hq ADR 0066). + // Held for this round only, and after the last round's were given back, so two pushes + // cascading into each other's machines never each wait on the other. + refused, err := sendRound(ctx, open, sending, compose, d, holder) + refusals = append(refusals, refused...) + if err != nil { + return refusals, err + } + } +} + +// composeForPush is how a push composes one machine: its set resolved, what it cannot host and what +// is left out of it said, and its declaration allocated. +func composeForPush(open *stores, gens map[string]catalogue.Generator) func(held context.Context, node string) (sendable, error) { + return func(held context.Context, node string) (sendable, error) { + plan, settings, err := planFor(held, open, node) + if err != nil { + return sendable{}, err + } + reportUnhostable(node, plan) + declared, err := declarationWith(held, open, node, plan, settings, gens, Allocating) + if err == nil { + reportLeftOut(node, declared) + } + return declared, err + } +} diff --git a/cmd/mesh-controller/held_back_test.go b/cmd/mesh-controller/held_back_test.go new file mode 100644 index 0000000..5733803 --- /dev/null +++ b/cmd/mesh-controller/held_back_test.go @@ -0,0 +1,335 @@ +package main + +import ( + "bytes" + "context" + "encoding/json" + "reflect" + "slices" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/catalogue" + "github.com/novox/mesh-controller/internal/inventory" + "github.com/novox/mesh-controller/internal/overlay" +) + +// novox/hq issue 259, ADR 0221: a push that did not name a machine does not send it a build its +// upgrade policy records rather than rolls out, nor one an open plan has not sent it yet. Anything +// else that moved is still a consequence the push sends (ADR 0083). +func TestAHeldBuildHoldsAMachineANamedPushDidNotName(t *testing.T) { + current := map[string]inventory.CurrentBuild{ + "resolver": {Commit: "c2c2c2c2c2"}, + "agent": {Commit: "a2", RollOut: true}, + "network": {}, + } + modules := []string{"network", "resolver", "agent"} + sent := map[string]string{"network": "", "resolver": "c1c1c1c1c1", "agent": "a2"} + + // A module whose policy records moved: held, naming it, both builds and why. + why := heldBack("laptop", modules, sent, true, current, nil) + if len(why) != 1 || !strings.Contains(why[0], "resolver would move from c1c1c1c1 to c2c2c2c2") || + !strings.Contains(why[0], "upgrade policy records") { + t.Fatalf("a recorded upgrade did not hold the machine: %v", why) + } + + // Nothing moved — what differs is a grant, a peer, a setting: not held (issue 057). + sent["resolver"] = "c2c2c2c2c2" + if why := heldBack("laptop", modules, sent, true, current, nil); len(why) != 0 { + t.Fatalf("a machine whose builds are all current was held: %v", why) + } + + // A module whose policy rolls out moved, and no plan holds it: sent, as before. + sent["agent"] = "a1" + if why := heldBack("laptop", modules, sent, true, current, nil); len(why) != 0 { + t.Fatalf("a rolled-out upgrade no plan holds was held: %v", why) + } + + // The last send's builds are not known: held whole. + if why := heldBack("laptop", modules, nil, false, current, nil); len(why) != 1 || + !strings.Contains(why[0], "not known") { + t.Fatalf("a machine whose last send was not recorded was not held: %v", why) + } + + // A module the machine was never sent, under a recording policy: held, and said so. + delete(sent, "resolver") + sent["agent"] = "a2" + if why := heldBack("laptop", modules, sent, true, current, nil); len(why) != 1 || + !strings.Contains(why[0], "resolver would move (never sent it) to c2c2c2c2") { + t.Fatalf("a module never sent under a recording policy: %v", why) + } +} + +// ADR 0218 meets ADR 0083: a plan waiting on its first machine has not sent the rest, and a push +// naming some other machine must not send them for it. +func TestAPlanWaitingOnItsFirstMachineHoldsTheRest(t *testing.T) { + at := time.Now() + current := map[string]inventory.CurrentBuild{"agent": {Commit: "a2", RollOut: true}} + sent := map[string]string{"agent": "a1"} + waiting := []inventory.Plan{{ID: "plan-7", State: inventory.PlanRolling, Modules: map[string]*inventory.PlanModule{ + "agent": {State: "built", First: []string{"ace"}, FirstAt: &at}}}} + + why := heldBack("g14", []string{"agent"}, sent, true, current, waiting) + if len(why) != 1 || !strings.Contains(why[0], "plan-7 has not sent it yet") { + t.Fatalf("a machine the plan has not reached was not held: %v", why) + } + // The first machine itself was sent by the plan: not held by it. + if why := heldBack("ace", []string{"agent"}, sent, true, current, waiting); len(why) != 0 { + t.Fatalf("the plan's first machine was held: %v", why) + } + // Built but not yet sent anywhere, or not yet built: the plan has it still to send. + for what, s := range map[string]*inventory.PlanModule{"built, unsent": {State: "built"}, "unasked": nil} { + plans := []inventory.Plan{{ID: "plan-8", State: inventory.PlanBuilding, + Modules: map[string]*inventory.PlanModule{"agent": s}}} + if why := heldBack("ace", []string{"agent"}, sent, true, current, plans); len(why) != 1 { + t.Errorf("%s: not held: %v", what, why) + } + } + // Sent everywhere, failed, or a plan no longer open: the plan holds nothing back. + for what, plans := range map[string][]inventory.Plan{ + "sent everywhere": {{ID: "p", State: inventory.PlanRolling, Modules: map[string]*inventory.PlanModule{ + "agent": {State: "built", First: []string{"ace"}, FirstAt: &at, SentAt: &at}}}}, + "failed": {{ID: "p", State: inventory.PlanRolling, Modules: map[string]*inventory.PlanModule{ + "agent": {State: "failed"}}}}, + "closed": {{ID: "p", State: inventory.PlanDone, Modules: map[string]*inventory.PlanModule{ + "agent": {State: "built"}}}}, + } { + if why := heldBack("g14", []string{"agent"}, sent, true, current, plans); len(why) != 0 { + t.Errorf("%s: held: %v", what, why) + } + } +} + +// A send records the build of each module it carried; a module left out of it keeps the build it +// was last sent, since the machine keeps that one. +func TestASendCarriesTheCurrentBuildsAndALeftOutModuleKeepsItsOwn(t *testing.T) { + current := map[string]inventory.CurrentBuild{"a": {Commit: "a2"}, "b": {Commit: "b2"}, "c": {}} + got := carriedBuilds([]string{"a", "b", "c"}, map[string]string{"b": "a setting does not compose"}, + current, map[string]string{"a": "a1", "b": "b1"}) + if want := map[string]string{"a": "a2", "b": "b1", "c": ""}; !reflect.DeepEqual(got, want) { + t.Fatalf("carried %v, wanted %v", got, want) + } + // Not known before: the left-out module is not recorded at all, so it reads as never sent. + got = carriedBuilds([]string{"a", "b"}, map[string]string{"b": "x"}, current, nil) + if want := map[string]string{"a": "a2"}; !reflect.DeepEqual(got, want) { + t.Fatalf("carried %v, wanted %v", got, want) + } +} + +// recordedDelivery sends nothing and records each send as the mesh does, so the next comparison +// reads the machine as current — and writes down which machines it declared. +type recordedDelivery struct { + inv *inventory.Inventory + declared []string +} + +func (r *recordedDelivery) grant(context.Context, []readyNode) error { return nil } + +func (r *recordedDelivery) declare(ctx context.Context, s readyNode, body []byte) (string, error) { + r.declared = append(r.declared, s.node) + return recordSent(ctx, r.inv, s.node, body, s.declared.Builds) +} + +// aResolver is a module built from a repository, at a commit, with something on the machine that +// says which build it is. +func aResolver(t *testing.T, open *stores, commit string, asked time.Time) { + t.Helper() + m := catalogue.Manifest{Module: "resolver", Version: "1", Resources: []map[string]any{ + {"id": "zones", "type": "file", "path": "/etc/resolver/zones", "content": "built from " + commit}, + }} + if err := open.inventory.RegisterModule(t.Context(), m, inventory.Source{ + Repository: "novox/mesh-catalog", Path: "modules/resolver", BuiltFrom: commit, Asked: asked}); err != nil { + t.Fatal(err) + } +} + +// A third machine on the private network: once it is sent, every other machine's peers change with +// it, which is a consequence a push must still send — no build moved. +func aThirdMachine(t *testing.T, open *stores) { + t.Helper() + ctx := t.Context() + record, err := open.inventory.AddNode(ctx, "spare") + if err != nil { + t.Fatal(err) + } + if err := open.inventory.SetPlace(ctx, "spare", "spare.example:51820", "here", false, "10.77.0.3"); err != nil { + t.Fatal(err) + } + reported, err := json.Marshal(map[string]any{"capabilities": []map[string]any{ + {"name": "container-runtime", "present": true}, {"name": "wireguard", "present": true}, + {"name": "systemd", "present": true}}}) + if err != nil { + t.Fatal(err) + } + var profile map[string]any + if err := json.Unmarshal(reported, &profile); err != nil { + t.Fatal(err) + } + if err := open.inventory.RecordProfile(ctx, record.ID, profile); err != nil { + t.Fatal(err) + } + if err := open.inventory.RecordSealingKey(ctx, record.ID, aPublicKey(t)); err != nil { + t.Fatal(err) + } + if err := open.inventory.RecordOverlayKey(ctx, record.ID, aPublicKey(t)); err != nil { + t.Fatal(err) + } + if _, err := open.inventory.Assign(ctx, "spare", overlay.Name); err != nil { + t.Fatal(err) + } +} + +// The issue as it happened, against the real stores: a change merged with the policy `record`, a push +// naming the anchor, and the laptop — running the same module — left with what it had, by name. +func TestANamedPushLeavesAMachineAPolicyHoldsBack(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + inv := open.inventory + asked := time.Now().Add(-time.Hour) + aResolver(t, open, "c1c1c1c1c1", asked) + for _, node := range []string{"anchor", "laptop"} { + if _, err := inv.Assign(ctx, node, "resolver"); err != nil { + t.Fatal(err) + } + } + gens, err := generators(ctx, open) + if err != nil { + t.Fatal(err) + } + compose := composeForPush(open, gens) + d := &recordedDelivery{inv: inv} + if _, err := sendRound(ctx, open, []string{"anchor", "laptop"}, compose, d, ""); err != nil { + t.Fatal(err) + } + if builds, known, err := inv.SentBuilds(ctx, "laptop"); err != nil || !known || builds["resolver"] != "c1c1c1c1c1" { + t.Fatalf("the send did not record the build it carried: %v %v %v", builds, known, err) + } + digestOfLaptop := func() string { + sent, err := inv.Outstanding(ctx, "laptop") + if err != nil { + t.Fatal(err) + } + return sent + } + before := digestOfLaptop() + + // The change merges; the policy is the default, record. `push anchor` sends the anchor... + aResolver(t, open, "c2c2c2c2c2", asked.Add(time.Minute)) + d.declared = nil + if _, err := sendRound(ctx, open, []string{"anchor"}, compose, d, ""); err != nil { + t.Fatal(err) + } + // ...and its cascade leaves the laptop, saying so. + var said bytes.Buffer + d.declared = nil + refused, err := flushBehind(ctx, open, mustNodes(t, open), map[string]bool{"anchor": true}, compose, d, "", &said) + if err != nil || len(refused) != 0 { + t.Fatalf("the cascade failed: %v %v", refused, err) + } + if len(d.declared) != 0 { + t.Fatalf("the cascade sent %v a build its policy records", d.declared) + } + if digestOfLaptop() != before { + t.Fatal("the laptop's last send moved: it was sent the held build") + } + for _, want := range []string{"laptop is behind and was not sent", "resolver would move from c1c1c1c1 to c2c2c2c2", + "upgrade policy records", "`push laptop` sends it"} { + if !strings.Contains(said.String(), want) { + t.Errorf("the push did not say %q:\n%s", want, said.String()) + } + } + + // Held and owed something else at once — a peer joined: still not sent, and both said: why it + // is held, and that what else it is owed waits with it. + aThirdMachine(t, open) + d.declared = nil + if _, err := sendRound(ctx, open, []string{"spare"}, compose, d, ""); err != nil { + t.Fatal(err) + } + d.declared = nil + said.Reset() + if _, err := flushBehind(ctx, open, mustNodes(t, open), map[string]bool{"anchor": true, "spare": true}, + compose, d, "", &said); err != nil { + t.Fatal(err) + } + if len(d.declared) != 0 || digestOfLaptop() != before { + t.Fatalf("a held machine owed a consequence was sent: %v", d.declared) + } + if !strings.Contains(said.String(), "resolver would move") || !strings.Contains(said.String(), "anything else it is owed") { + t.Fatalf("the push did not say both:\n%s", said.String()) + } + + // A policy that rolls out: the laptop is a consequence like any other, and sent. + if err := inv.SetUpgradeOf(ctx, "resolver", inventory.Upgrade{RollOut: true}); err != nil { + t.Fatal(err) + } + said.Reset() + if _, err := flushBehind(ctx, open, mustNodes(t, open), map[string]bool{"anchor": true, "spare": true}, + compose, d, "", &said); err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(d.declared, []string{"laptop"}) || digestOfLaptop() == before { + t.Fatalf("a rolled-out upgrade's machine was not sent: %v\n%s", d.declared, said.String()) + } + if builds, _, _ := inv.SentBuilds(ctx, "laptop"); builds["resolver"] != "c2c2c2c2c2" { + t.Fatalf("the new send did not record the new build: %v", builds) + } +} + +// Issue 057's case is unchanged: a machine whose builds are all current and whose declaration moved +// for another reason is sent by a push that names someone else. +func TestANamedPushStillSendsAConsequenceNothingHolds(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + inv := open.inventory + aResolver(t, open, "c1c1c1c1c1", time.Now().Add(-time.Hour)) + if _, err := inv.Assign(ctx, "laptop", "resolver"); err != nil { + t.Fatal(err) + } + gens, err := generators(ctx, open) + if err != nil { + t.Fatal(err) + } + compose := composeForPush(open, gens) + d := &recordedDelivery{inv: inv} + if _, err := sendRound(ctx, open, []string{"anchor", "laptop"}, compose, d, ""); err != nil { + t.Fatal(err) + } + + // `push spare`, the machine just placed: the others' peers change with it. + aThirdMachine(t, open) + if _, err := sendRound(ctx, open, []string{"spare"}, compose, d, ""); err != nil { + t.Fatal(err) + } + d.declared = nil + var said bytes.Buffer + if _, err := flushBehind(ctx, open, mustNodes(t, open), map[string]bool{"spare": true}, compose, d, "", &said); err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(d.declared, []string{"anchor", "laptop"}) { + t.Fatalf("a consequence nothing holds was not sent: %v\n%s", d.declared, said.String()) + } + if strings.Contains(said.String(), "was not sent") { + t.Fatalf("a machine nothing holds was said to be held:\n%s", said.String()) + } + + // A machine whose last send was not recorded — a declaration sent by hand — is held until named. + record, err := inv.NodeByName(ctx, "laptop") + if err != nil { + t.Fatal(err) + } + if err := inv.RecordSent(ctx, record.ID, "sent-by-hand", nil); err != nil { + t.Fatal(err) + } + d.declared = nil + said.Reset() + if _, err := flushBehind(ctx, open, mustNodes(t, open), map[string]bool{"spare": true}, compose, d, "", &said); err != nil { + t.Fatal(err) + } + // The anchor, the hub, may still be settling from the machine placed above; the laptop is the + // question. + if slices.Contains(d.declared, "laptop") || !strings.Contains(said.String(), "laptop is behind and was not sent") { + t.Fatalf("a machine whose last send is not known was sent: %v\n%s", d.declared, said.String()) + } +} diff --git a/cmd/mesh-controller/issue_204_test.go b/cmd/mesh-controller/issue_204_test.go index 9bcc5d2..f32af93 100644 --- a/cmd/mesh-controller/issue_204_test.go +++ b/cmd/mesh-controller/issue_204_test.go @@ -77,7 +77,7 @@ func TestASendIsRecordedEvenWhenTheSenderIsBeingCancelled(t *testing.T) { } cancel() // the sender is going away: its context is cancelled between the send and the record body := []byte(`{"declaration":1,"resources":[]}`) - digest, err := recordSent(ctx, inv, "anchor", body) + digest, err := recordSent(ctx, inv, "anchor", body, nil) if err != nil { // NodeByName on the cancelled context may itself refuse; the record must still be possible // through the detached context, so look the node up again on a live one. diff --git a/cmd/mesh-controller/plan.go b/cmd/mesh-controller/plan.go index 2c66137..967c4bd 100644 --- a/cmd/mesh-controller/plan.go +++ b/cmd/mesh-controller/plan.go @@ -397,9 +397,49 @@ func declarationWith(ctx context.Context, open *stores, node string, if err != nil { return sendable{}, err } - return sendable{Resources: composed.Resources, Adoption: adoption, + out := sendable{Resources: composed.Resources, Adoption: adoption, Received: composed.Received, Mesh: with.Mesh, BusUsers: with.BusUsers, - LeftOut: sortedKeysOf(composed.LeftOut), leftOutWhy: composed.LeftOut}, nil + LeftOut: sortedKeysOf(composed.LeftOut), leftOutWhy: composed.LeftOut} + // And which build of each module it carries, for the send to record (novox/hq issue 259, ADR + // 0221). Read only on the send path: a question about what would be sent records nothing. + if choosing == Allocating { + current, err := open.inventory.CurrentBuilds(ctx) + if err != nil { + return sendable{}, err + } + before, known, err := open.inventory.SentBuilds(ctx, node) + if err != nil { + return sendable{}, err + } + if !known { + before = nil + } + names := make([]string, 0, len(plan.Modules)) + for _, m := range plan.Modules { + names = append(names, m.Module) + } + out.Builds = carriedBuilds(names, composed.LeftOut, current, before) + } + return out, nil +} + +// carriedBuilds is the build of each module a declaration carries, as a send records it (novox/hq +// issue 259): the module's current build for each module in it, and for a module left out of it +// (ADR 0163, rule 6) the build it was last sent, since the machine keeps that one — or nothing, when +// that is not known. Never nil, so a send through here always records what it knows. +func carriedBuilds(modules []string, leftOut map[string]string, current map[string]inventory.CurrentBuild, + before map[string]string) map[string]string { + out := map[string]string{} + for _, m := range modules { + if _, left := leftOut[m]; left { + if was, kept := before[m]; kept { + out[m] = was + } + continue + } + out[m] = current[m].Commit + } + return out } // sortedKeysOf is a map's keys, sorted — so what a declaration says it left out does not move diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 3e43f0a..94cddfd 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -217,8 +217,9 @@ func declare(ctx context.Context, args []string) error { } // Written down like every other send (novox/hq issue 204): a declaration a person sent by hand // is still what the machine was last told, and status must not read it as current for the one - // the mesh would compose. - if _, err := recordSent(ctx, inv, node, raw); err != nil { + // the mesh would compose. Which builds it carried is recorded as not known (novox/hq issue 259): + // the mesh did not compose it, so a push that does not name this machine treats it as held. + if _, err := recordSent(ctx, inv, node, raw, nil); err != nil { return err } fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw)) @@ -367,6 +368,23 @@ func pushCommand(ctx context.Context, args []string) error { if err != nil { return err } + // **Not when its own modules are held back** (novox/hq issue 259, ADR 0221): added rather than + // named, it is sent its whole declaration, and a build its policy records or a plan has not sent + // it yet would go with the user list. Named and left, with what that costs. + saidHeld := map[string]bool{} + if holderBehind && len(args) == 1 { + held, err := heldMachines(ctx, open, []string{holder}) + if err != nil { + return err + } + if why, isHeld := held[holder]; isHeld { + sayHeld(os.Stdout, holder, why) + fmt.Printf("%s holds the bus, and the user list it would carry has changed: until it is "+ + "sent, the bus may refuse what this push's machines were newly granted\n", holder) + holderBehind = false + saidHeld[holder] = true + } + } asked = brokerFirst(asked, holder, holderBehind) // Held from composing to sending, so a converge on one of them cannot send between the two @@ -418,76 +436,21 @@ func pushCommand(ctx context.Context, args []string) error { // // Compared against what each machine was last SENT, not against a before/after of this push: // the mint usually happened at `assign` or `module issue`, before this command ran, so the - // only durable signal is "what it should be" versus "what it last received". A machine behind - // for an unrelated reason is caught here too, which is not a cost — a named push that knew a - // machine was behind and left it so would be the very silence this removes. Bounded: a + // only durable signal is "what it should be" versus "what it last received". Bounded: a // flushed send may itself mint, so this converges over a few rounds. + // + // **Except a machine a policy or a plan holds back** (novox/hq issue 259, ADR 0221): one whose + // modules would move to a build their upgrade policy records rather than rolls out, or that an + // open plan has not sent it yet. It is named, with why, and left for a push that names it. if len(args) == 1 { - flushed := map[string]bool{args[0]: true} - // Bounded by the node count: a node is marked flushed the round it is handled and is - // never handled twice, so the loop cannot run more than len(nodes) rounds. The bound is - // a guard against a logic error, not a real limit — if it were ever hit, that is a bug - // rather than a cascade legitimately still converging, so it is said rather than passed - // over in silence, unlike the earlier fixed cap that could stop a real cascade short. - rounds := 0 - for { - would, err := wouldSend(ctx, open, nodes) - if err != nil { - return err - } - behind, err := inv.Waiting(ctx, would) - if err != nil { - return err - } - var also []string - for _, m := range behind { - if !flushed[m.Node] { - also = append(also, m.Node) - } - } - if len(also) == 0 { - break - } - if rounds++; rounds > len(nodes) { - fmt.Printf("\nstopped cascading after %d rounds with %s still behind — this "+ - "should not happen; run `push --behind` to finish\n", - rounds-1, strings.Join(also, ", ")) - break - } - sort.Strings(also) - fmt.Printf("\nthis push left %s behind — a provision granted from there, or a "+ - "declaration since changed; sending it too\n", strings.Join(also, ", ")) - // Tolerantly, exactly as the named send above: a machine that cannot be composed is - // collected as a refusal and reported at the end, and the others are still sent - // (novox/hq ADR 0066). The earlier cut routed these through sendTo, which is - // all-or-nothing — so one swept machine's compose error failed the operator's named - // push and skipped its --wait, the very intolerance the main path exists to avoid. - // Held for this round only, and after the last round's were given back, so two pushes - // cascading into each other's machines never each wait on the other. - refused, err := sendRound(ctx, open, also, - func(held context.Context, node string) (sendable, error) { - plan, settings, err := planFor(held, open, node) - if err != nil { - return sendable{}, err - } - reportUnhostable(node, plan) - declared, err := declarationWith(held, open, node, plan, settings, gens, Allocating) - if err == nil { - reportLeftOut(node, declared) - } - return declared, err - }, - bus, holder) - refusals = append(refusals, refused...) - if err != nil { - return err - } - // Every candidate this round is marked handled — the sent ones so they are not - // re-listed, and the refused ones so a machine that cannot be composed does not make - // the loop spin on it for ever. Its refusal is already in the report. - for _, name := range also { - flushed[name] = true - } + handled := map[string]bool{args[0]: true} + for n := range saidHeld { + handled[n] = true + } + refused, err := flushBehind(ctx, open, nodes, handled, composeForPush(open, gens), bus, holder, os.Stdout) + refusals = append(refusals, refused...) + if err != nil { + return err } } @@ -762,7 +725,7 @@ func (b overTheBus) declare(ctx context.Context, s readyNode, body []byte) (stri } // After it is away, not before. A digest recorded for something that failed to send would make // the machine look current for a declaration it never received. - digest, err := recordSent(ctx, b.open.inventory, s.node, body) + digest, err := recordSent(ctx, b.open.inventory, s.node, body, s.declared.Builds) if err != nil { return "", err } @@ -1244,7 +1207,11 @@ func allot(ctx context.Context, inv *inventory.Inventory, node string) (int64, e // never wrote it down: status read "applied, current" over a machine that had just been sent // something else. What was sent was sent; the record of it must not depend on the sender living // another second. Bounded, so a store that is away does not hold a dying process open for ever. -func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body []byte) (string, error) { +// +// And the build of each module it carried (novox/hq issue 259, ADR 0221), nil when that is not known: +// what tells a machine held back by a policy or a plan from one a push left behind. +func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body []byte, + builds map[string]string) (string, error) { kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second) defer cancel() record, err := inv.NodeByName(kept, node) @@ -1252,7 +1219,7 @@ func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body return "", err } digest := digestOf(body) - if err := inv.RecordSent(kept, record.ID, digest); err != nil { + if err := inv.RecordSent(kept, record.ID, digest, builds); err != nil { return "", err } return digest, nil diff --git a/cmd/mesh-controller/sendable.go b/cmd/mesh-controller/sendable.go index 4219e97..d9ff4f1 100644 --- a/cmd/mesh-controller/sendable.go +++ b/cmd/mesh-controller/sendable.go @@ -43,6 +43,10 @@ type sendable struct { LeftOut []string // leftOutWhy is why each was, for push and plan to say; never on the wire. leftOutWhy map[string]string + // Builds is the build of each module this declaration carries — module to the commit its build + // was made from — recorded with the send and never on the wire (novox/hq issue 259, ADR 0221). + // Composed only on the send path; nil records that it is not known. + Builds map[string]string } // adoptionEnvelope is what an adopted node is told about its mode. Taken is every module taken on diff --git a/internal/inventory/catalogue.go b/internal/inventory/catalogue.go index 342bb0b..c7e5be6 100644 --- a/internal/inventory/catalogue.go +++ b/internal/inventory/catalogue.go @@ -1158,6 +1158,37 @@ func (i *Inventory) SetUpgradeOf(ctx context.Context, module string, u Upgrade) return nil } +// CurrentBuild is the build a module is at and whether its policy rolls a new one out (novox/hq +// issue 259). +type CurrentBuild struct { + // Commit is the commit the module's manifest was read at — its `built_from` — and empty for a + // module held without a source. + Commit string + // RollOut is the module's upgrade policy, as Upgrade.RollOut. + RollOut bool +} + +// CurrentBuilds is every module's current build and upgrade policy, in one read: what a send records +// it carried, and what a push compares a machine's last send against. +func (i *Inventory) CurrentBuilds(ctx context.Context) (map[string]CurrentBuild, error) { + rows, err := i.store.Pool().Query(ctx, + `select name, coalesce(built_from, ''), upgrade = 'roll-out' from module`) + if err != nil { + return nil, err + } + defer rows.Close() + out := map[string]CurrentBuild{} + for rows.Next() { + var name string + var b CurrentBuild + if err := rows.Scan(&name, &b.Commit, &b.RollOut); err != nil { + return nil, err + } + out[name] = b + } + return out, rows.Err() +} + // Running is every machine assigned a module, in a stable order. // // **Assigned, not reported.** A machine that is assigned the module and has not applied it yet is diff --git a/internal/inventory/doing_test.go b/internal/inventory/doing_test.go index 8ce8f0e..76d6a20 100644 --- a/internal/inventory/doing_test.go +++ b/internal/inventory/doing_test.go @@ -189,7 +189,7 @@ func TestAMachineIsWaitingWhenWhatItWasSentIsNotWhatItShouldBe(t *testing.T) { } // Sent what it should be: not waiting. - if err := inv.RecordSent(ctx, anchor.ID, "aaa"); err != nil { + if err := inv.RecordSent(ctx, anchor.ID, "aaa", nil); err != nil { t.Fatal(err) } waiting, err = inv.Waiting(ctx, map[string]string{"anchor": "aaa", "laptop": "bbb"}) @@ -234,7 +234,7 @@ func TestAMachineWithNothingComputedForItIsNotWaiting(t *testing.T) { // It has been sent something before, which is what makes this the case the guard is for: a // machine with a digest and nothing computed for it would compare against the empty string // and look out of date, when the truth is that nobody worked out what it should be. - if err := inv.RecordSent(ctx, node.ID, "what-it-got-last-time"); err != nil { + if err := inv.RecordSent(ctx, node.ID, "what-it-got-last-time", nil); err != nil { t.Fatal(err) } waiting, err := inv.Waiting(ctx, map[string]string{}) diff --git a/internal/inventory/migrations/0061-a-send-records-the-builds-it-carried.sql b/internal/inventory/migrations/0061-a-send-records-the-builds-it-carried.sql new file mode 100644 index 0000000..0952aa1 --- /dev/null +++ b/internal/inventory/migrations/0061-a-send-records-the-builds-it-carried.sql @@ -0,0 +1,14 @@ +-- A send records the build of each module it carried (novox/hq issue 259, ADR 0221). +-- +-- A named push ends by sending every other machine whose declaration differs from what it was last +-- sent (ADR 0083). Read from the declaration's digest alone, a module whose upgrade policy records +-- rather than rolls out, or whose plan sends one machine first (ADR 0218), made every machine running +-- it differ, so `push ` sent all of them the build the policy was holding back. Which +-- build of each module a machine was last sent is what tells a held upgrade from a consequence of the +-- push, and it is not in a digest. +-- +-- Module name to the commit its build was made from — the module's `built_from` when the declaration +-- was composed, empty for a module the mesh holds without a source. NULL for a machine last sent +-- before this was kept, or sent a declaration by hand: what it carried is not known, and the push +-- treats such a machine as held until it is pushed by name. +alter table node add column sent_builds jsonb; diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index e16e323..a1da8cf 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -909,18 +909,53 @@ func (i *Inventory) LastReports(ctx context.Context) ([]Reported, error) { return out, rows.Err() } -// RecordSent keeps a digest of the declaration a machine was last sent. +// RecordSent keeps a digest of the declaration a machine was last sent, and the build of each module +// it carried. // // **A digest rather than the declaration.** The mesh can compute what a machine should be at any // moment; keeping a copy would be a second account of it, able to disagree with the first. What // cannot be recomputed is what was *actually sent*, and that is the whole difference between a // machine that is out of date and one that has never been told. -func (i *Inventory) RecordSent(ctx context.Context, node, digest string) error { +// +// **And which build of each module** (novox/hq issue 259, ADR 0221): module name to the commit its +// build was made from. A digest cannot say whether a machine differs because a module moved to a build +// its upgrade policy holds back, or because of something a push made — a grant — and only the second +// is a push's to send to a machine it did not name. Nil records that it is not known, as for a +// declaration sent by hand. +func (i *Inventory) RecordSent(ctx context.Context, node, digest string, builds map[string]string) error { + var carried *string + if builds != nil { + raw, err := json.Marshal(builds) + if err != nil { + return err + } + text := string(raw) + carried = &text + } _, err := i.store.Pool().Exec(ctx, - `update node set sent = $2, sent_at = now() where id = $1`, node, digest) + `update node set sent = $2, sent_at = now(), sent_builds = $3::jsonb where id = $1`, node, digest, carried) return err } +// SentBuilds is the build of each module a machine was last sent, by its name: module to the commit +// its build was made from (novox/hq issue 259). Known is false when that was not kept — a machine +// last sent before it was, sent a declaration by hand, or one the mesh does not know. +func (i *Inventory) SentBuilds(ctx context.Context, name string) (builds map[string]string, known bool, err error) { + var raw []byte + err = i.store.Pool().QueryRow(ctx, `select sent_builds from node where name = $1`, name).Scan(&raw) + if errors.Is(err, pgx.ErrNoRows) { + return nil, false, nil + } + if err != nil || raw == nil { + return nil, false, err + } + builds = map[string]string{} + if err := json.Unmarshal(raw, &builds); err != nil { + return nil, false, err + } + return builds, true, nil +} + // RecordSentBusUsers keeps a digest of the bus's user list a machine was just sent, by its name // (novox/hq issue 249): whether the machine holding the bus must go first is whether this differs // from the list composed now. diff --git a/internal/inventory/sent_builds_test.go b/internal/inventory/sent_builds_test.go new file mode 100644 index 0000000..9a1d150 --- /dev/null +++ b/internal/inventory/sent_builds_test.go @@ -0,0 +1,81 @@ +package inventory + +import ( + "reflect" + "testing" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// novox/hq issue 259: a send keeps which build of each module it carried. A machine never sent +// anything, or sent with nothing recorded — before this was kept, or by hand — is not known, which +// is not the same as having been sent nothing. +func TestASendKeepsTheBuildsItCarried(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + node, err := inv.AddNode(ctx, "anchor") + if err != nil { + t.Fatal(err) + } + if builds, known, err := inv.SentBuilds(ctx, "anchor"); err != nil || known || builds != nil { + t.Fatalf("a machine never sent anything has known builds %v (%v): %v", builds, known, err) + } + + carried := map[string]string{"resolver": "abc123", "network": ""} + if err := inv.RecordSent(ctx, node.ID, "d1", carried); err != nil { + t.Fatal(err) + } + builds, known, err := inv.SentBuilds(ctx, "anchor") + if err != nil || !known || !reflect.DeepEqual(builds, carried) { + t.Fatalf("the builds sent were not kept: %v %v %v", builds, known, err) + } + + // Sent with nothing carried: known, and empty. + if err := inv.RecordSent(ctx, node.ID, "d2", map[string]string{}); err != nil { + t.Fatal(err) + } + if builds, known, err := inv.SentBuilds(ctx, "anchor"); err != nil || !known || len(builds) != 0 { + t.Fatalf("an empty send: %v %v %v", builds, known, err) + } + + // Sent by hand: not known, and the digest still recorded. + if err := inv.RecordSent(ctx, node.ID, "d3", nil); err != nil { + t.Fatal(err) + } + if builds, known, err := inv.SentBuilds(ctx, "anchor"); err != nil || known || builds != nil { + t.Fatalf("a send whose builds are not known read as %v %v: %v", builds, known, err) + } + if sent, err := inv.Outstanding(ctx, "anchor"); err != nil || sent != "d3" { + t.Fatalf("the digest was not recorded with it: %q %v", sent, err) + } + + if _, known, err := inv.SentBuilds(ctx, "nobody"); err != nil || known { + t.Fatalf("a machine the mesh does not know: %v %v", known, err) + } +} + +// A module's current build is the commit its manifest was read at, with its upgrade policy. +func TestTheCurrentBuildsAreTheCatalogues(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "resolver", Version: "1"}, + Source{Repository: "novox/mesh-catalog", BuiltFrom: "c1"}); err != nil { + t.Fatal(err) + } + if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "by-hand", Version: "1"}, Source{}); err != nil { + t.Fatal(err) + } + if err := inv.SetUpgradeOf(ctx, "by-hand", Upgrade{RollOut: true}); err != nil { + t.Fatal(err) + } + current, err := inv.CurrentBuilds(ctx) + if err != nil { + t.Fatal(err) + } + if got := current["resolver"]; got != (CurrentBuild{Commit: "c1"}) { + t.Errorf("resolver is at %+v", got) + } + if got := current["by-hand"]; got != (CurrentBuild{RollOut: true}) { + t.Errorf("a module with no source is at %+v", got) + } +} diff --git a/internal/link/heard_test.go b/internal/link/heard_test.go index 8941cab..c59e9dd 100644 --- a/internal/link/heard_test.go +++ b/internal/link/heard_test.go @@ -100,7 +100,7 @@ func TestABareAliveDoesNotWipeTheDeclarationThatSaysANodeIsCurrent(t *testing.T) } // The mesh sent this node a declaration, and the node applied it and named which by digest. const digest = "d640d1b6a1b2c3d4e5f60718293a4b5c6d7e8f90a1b2c3d4e5f6071829304152" - if err := inv.RecordSent(ctx, node.ID, digest); err != nil { + if err := inv.RecordSent(ctx, node.ID, digest, nil); err != nil { t.Fatal(err) } if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{