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 <one>` 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 <node>` 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.
This commit is contained in:
jochen
2026-10-05 22:22:03 +02:00
parent 8a400d165e
commit 2421b82ad2
13 changed files with 831 additions and 83 deletions
+1 -1
View File
@@ -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{
+241
View File
@@ -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 <anchor>` 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
}
}
+335
View File
@@ -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())
}
}
+1 -1
View File
@@ -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.
+42 -2
View File
@@ -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
+40 -73
View File
@@ -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
+4
View File
@@ -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
+31
View File
@@ -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
+2 -2
View File
@@ -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{})
@@ -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 <one machine>` 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;
+38 -3
View File
@@ -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.
+81
View File
@@ -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)
}
}
+1 -1
View File
@@ -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{