From e6e1e3bc895863571d45e81632d01513c4c384de Mon Sep 17 00:00:00 2001 From: jochen Date: Fri, 9 Oct 2026 17:15:40 +0200 Subject: [PATCH] A gate judges its own send and the build it sent, and never puts the controller back behind its store (hq issue 352) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit On 2026-10-09 a release's gate on the control node read the machine's report against a newer send another plan had just made there, failed three builds the machine had reported healthy, and put them back on every machine to a controller older than the store's schema; that controller then passed the newer plan's gate from its own health. - A gate keeps what its send carried (digest, sequence) and reads the report against it; a report on the last send is on it too. - A gate judges only the build the machine was last sent: another build there supersedes the judging — no verdict, nothing put back. - A controller is told its build (MESH_CONTROLLER_VERSION, ${version} in a process's env) and records how far it reads the store's schema; a put-back to a build that reaches less, or never said, is refused and the current build kept, said as urgent. - A release's open gate holds other sends of its modules there, and a plan's own first send waits on it. --- cmd/mesh-controller/acting.go | 17 +- cmd/mesh-controller/gate.go | 154 +++++++- cmd/mesh-controller/gate_test.go | 3 +- cmd/mesh-controller/issue352_test.go | 356 ++++++++++++++++++ cmd/mesh-controller/push.go | 30 ++ cmd/mesh-controller/release.go | 54 ++- cmd/mesh-controller/release_plan.go | 25 +- internal/catalogue/build.go | 31 +- internal/catalogue/version_in_a_path_test.go | 25 ++ internal/inventory/gate.go | 3 + ...records-what-its-stores-schema-reaches.sql | 10 + internal/inventory/nodes.go | 48 ++- internal/inventory/plans.go | 4 + internal/inventory/plans_test.go | 80 ++++ internal/inventory/schema.go | 47 +++ module.json | 3 +- 16 files changed, 869 insertions(+), 21 deletions(-) create mode 100644 cmd/mesh-controller/issue352_test.go create mode 100644 internal/inventory/migrations/0087-a-controller-records-what-its-stores-schema-reaches.sql create mode 100644 internal/inventory/schema.go diff --git a/cmd/mesh-controller/acting.go b/cmd/mesh-controller/acting.go index 63f7aea2..bae46dea 100644 --- a/cmd/mesh-controller/acting.go +++ b/cmd/mesh-controller/acting.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "os" + "strings" "sync" "time" @@ -179,7 +180,21 @@ func (a *actor) release() { // holderOf is this process as the lease's holder. func holderOf(instance string) lease.Holder { host, _ := os.Hostname() - return lease.Holder{Instance: instance, Host: host, Build: version} + build := runningBuild() + if build == "" { + build = version + } + return lease.Holder{Instance: instance, Host: host, Build: build} +} + +// RunningBuildVar is where the declaration tells this process which build it is (module.json, the +// controller process's env): the version its bundle is delivered as, `${version}` composed by the +// catalogue from the bundle's digest. Empty for a process placed by hand. +const RunningBuildVar = "MESH_CONTROLLER_VERSION" + +// runningBuild is the version of the build this process is, or empty when the declaration did not say. +func runningBuild() string { + return strings.TrimSpace(os.Getenv(RunningBuildVar)) } // serveUnderTheLease takes the lease for the serving controller, waiting while another holds it, and diff --git a/cmd/mesh-controller/gate.go b/cmd/mesh-controller/gate.go index d2282170..4ddb202d 100644 --- a/cmd/mesh-controller/gate.go +++ b/cmd/mesh-controller/gate.go @@ -87,6 +87,9 @@ const ( healthWaiting healthNotYet healthBroken + // healthSuperseded is a judging that cannot go on: the machine was sent another build of the module + // after the gate's send (novox/hq issue 352). No verdict on the build judged, and nothing put back. + healthSuperseded ) // served is what one machine's node tools answered the bus's discovery with. @@ -122,6 +125,39 @@ type gateFacts struct { // groupsAdded is, per module, whether the move judged puts an account in a group its previous build did // not (issue 318 review): the only move whose wait for a new login is excused. groupsAdded map[string]bool + // sent is, per machine, the declaration the gate's own send carried there (novox/hq issue 352): a + // report is held against it, never against the send made last. sentBuilds is what each machine was + // last sent of every module, and judged the commit of each module this gate judges: a machine last + // sent another build of the module is not running the build judged. + sent map[string]inventory.SentDeclaration + sentBuilds map[string]map[string]string + commits map[string]string +} + +// reportedOn says a machine's last report is on what the gate sent it (novox/hq issue 352): on that +// declaration, or one it was sent after it — or, for a gate kept before sends were kept on it, on the +// declaration last sent. On 2026-10-09 a release's gate read the control node's report against a newer +// send another plan had just made there, and failed three builds the machine had reported healthy as +// "has not reported on what it was sent". +func (f gateFacts) reportedOn(machine string, r inventory.Reported) bool { + if sent, kept := f.sent[machine]; kept { + return sent.ReportsOn(r) + } + return r.Current +} + +// supersededOn says the machine was last sent another build of the module than the one this gate judges +// (novox/hq issue 352): the judging cannot go on, whatever the machine reports. On 2026-10-09 a controller +// put back by one gate judged another gate's newer controller build passed on the same machine, reading +// the put-back build's health as the newer one's. +func (f gateFacts) supersededOn(module, machine string) (string, bool) { + judged, known := f.commits[module] + sent, has := f.sentBuilds[machine][module] + if !known || !has || judged == "" || sent == "" || sameCommit(sent, judged) { + return "", false + } + return fmt.Sprintf("%s was sent %s %s after this gate's %s: the build judged no longer runs there, and "+ + "this judging is superseded by that send's", machine, module, short(sent), short(judged)), true } // gatherGateFacts reads what a judging needs, from the store, the bus and this controller's memory. A @@ -199,9 +235,12 @@ func judgeHealth(module, component string, m catalogue.Manifest, machine string, return healthBroken, fmt.Sprintf("the witness on %s judged the %s %s and %s: %s", machine, r.Component, short(r.From), r.Outcome, r.Why) } + if why, superseded := f.supersededOn(module, machine); superseded { + return healthSuperseded, why + } r, said := f.reports[machine] switch { - case !said || r.At == nil || !r.Current: + case !said || r.At == nil || !f.reportedOn(machine, r): return healthNotYet, fmt.Sprintf("%s has not reported on what it was sent", machine) case r.Outcome == inventory.OutcomeFailed || r.Outcome == inventory.OutcomeRefused: return healthBroken, fmt.Sprintf("%s %s what it was sent", machine, r.Outcome) @@ -438,6 +477,21 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs return "", err } facts.groupsAdded = movesAddingGroups(ctx, open.inventory, g, pairs, shelf) + facts.sent = g.Sent + facts.commits, facts.sentBuilds = judgedCommits(g, pairs), map[string]map[string]string{} + // A module this gate put back at once (putBackBroken) was sent its earlier build by the gate itself: + // not another send, and not a judging superseded. + for _, m := range g.Returned { + delete(facts.commits, m) + } + for _, j := range pairs { + if _, read := facts.sentBuilds[j.node]; read { + continue + } + if builds, known, err := open.inventory.SentBuilds(ctx, j.node); err == nil && known { + facts.sentBuilds[j.node] = builds + } + } // **What is wrong with a machine itself is the machine's** (novox/hq issue 281): read once for each // machine judged, apart from what is wrong with a module there, and never pinned on the module the // gate happens to be kept on. @@ -474,6 +528,13 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs if _, seen := reading[j.module]; !seen { modules = append(modules, j.module) } + if h == healthSuperseded { + // Decided at once (novox/hq issue 352): nothing of this gate's can be judged on a machine that + // was sent another build of it, and nothing is put back — the later send is what runs there. + g.Failing, g.Last = nil, "" + decide(g, inventory.GateSuperseded, said, now) + return g.Verdict, nil + } if h == healthBroken && !slices.Contains(g.Broken, j.module) { g.Broken = append(g.Broken, j.module) if g.BrokenWhy == "" { @@ -577,6 +638,40 @@ func judgeMoves(ctx context.Context, open *stores, g *inventory.PlanGate, pairs return g.Verdict, nil } +// judgedCommits is the commit of each module a gate judges: the gate's own To for its module, and each +// carried move's. Pure. +func judgedCommits(g *inventory.PlanGate, pairs []judged) map[string]string { + out := map[string]string{} + for _, c := range g.Carried { + if c.To != "" { + out[c.Module] = c.To + } + } + if g.To != "" { + for _, j := range pairs { + if _, has := out[j.module]; !has && !slices.ContainsFunc(g.Carried, func(c inventory.CarriedMove) bool { return c.Module == j.module }) { + out[j.module] = g.To + } + } + } + return out +} + +// sentNow is what each machine was just sent, read after a send for the gate to keep (novox/hq issue +// 352): a machine whose send is not on record is left out, and its report is read as before. +func sentNow(ctx context.Context, inv *inventory.Inventory, machines []string) map[string]inventory.SentDeclaration { + out := map[string]inventory.SentDeclaration{} + for _, n := range machines { + if s, found, err := inv.SentTo(ctx, n); err == nil && found { + out[n] = s + } + } + if len(out) == 0 { + return nil + } + return out +} + // whyFor is a passing gate's why as one module's verdict says it: the send's, and that module's own wait // for a person, never another's (issue 318 review). func whyFor(g *inventory.PlanGate, module string) string { @@ -802,6 +897,16 @@ func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module str "back", module, short(state.Previous), inventory.KeptBuilds)) return } + if g.Component == lease.ComponentController { + // **The controller is never put back to a build older than the store's schema** (novox/hq issue + // 352): the build before it carries fewer migrations than the failed one applied, starts behind + // its own records, and judges the next gate with what it can read. The current build is kept and + // the condition says so; a person decides. + if why, ok := controllerSchemaAllows(ctx, inv, previous); !ok { + notBack(why) + return + } + } if err := inv.RestoreModule(ctx, previous); err != nil { notBack(err.Error()) return @@ -835,6 +940,53 @@ func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module str sayRollback(ctx, open, module, g, "") } +// controllerSchemaAllows says the store's schema lets this build of the controller be put back: the +// build recorded, when it served, a reach at or past the highest migration the store has applied. One +// that never recorded a reach is not proved safe, and is refused as such (novox/hq issue 352). Why +// says what is kept and why when it is not. +func controllerSchemaAllows(ctx context.Context, inv *inventory.Inventory, previous inventory.Build) (string, bool) { + applied, err := inv.SchemaApplied(ctx) + if err != nil { + return "what the store's schema reaches cannot be read, so whether the build before it can read it is not " + + "known; the current build is kept: " + err.Error(), false + } + build := versionOfBuild(previous) + if build == "" { + return fmt.Sprintf("the build before it (%s) names no bundle to know it by, so whether it can read the store's "+ + "schema (migration %04d) is not known; the current build is kept, and a person decides", short(previous.Commit), applied), false + } + reach, known, err := inv.SchemaReachOf(ctx, build) + if err != nil { + return "what the build before it knows of the store's schema cannot be read; the current build is kept: " + err.Error(), false + } + if !known { + return fmt.Sprintf("the build before it (%s, %s) never recorded how far it reads the store's schema — a "+ + "controller records that when it serves — so it is not proved to read migration %04d, which the store "+ + "has applied; a controller older than its store starts behind its own records and judges with what it "+ + "can read, so the current build is kept, and a person decides", short(previous.Commit), build, applied), false + } + if reach < applied { + return fmt.Sprintf("the build before it (%s, %s) reads the store's schema up to migration %04d, and the store "+ + "is at %04d: a controller older than its store starts behind its own records and judges with what it "+ + "can read, so the current build is kept, and a person decides", short(previous.Commit), build, reach, applied), false + } + return "", true +} + +// versionOfBuild is the version a build's bundle is delivered as — its archive's digest, short, as the +// catalogue names it (`${version}`) — read from the build's artifacts; empty when none is a bundle. +func versionOfBuild(b inventory.Build) string { + for _, a := range b.Made { + if a.Kind != catalogue.ArtifactBundle && a.Kind != catalogue.ArtifactArchive { + continue + } + if _, hex, found := strings.Cut(a.Reference, "sha256:"); found && len(hex) >= 12 { + return hex[:12] + } + } + return "" +} + // rollbacks is what a failed send puts back, sent together (novox/hq issue 281): a gate that judged one // send judges what it moved as one, and what it found wanting goes back in one send per machine — not // in a send for each module, which is the churn that failed the gate in the first place. diff --git a/cmd/mesh-controller/gate_test.go b/cmd/mesh-controller/gate_test.go index 9bbd46d3..8fe19c58 100644 --- a/cmd/mesh-controller/gate_test.go +++ b/cmd/mesh-controller/gate_test.go @@ -113,9 +113,10 @@ func aGateMesh(t *testing.T) *gateMesh { } for node, h := range g.health { if h == healthNotYet { + // Not reported on the send: neither the send made last, nor the gate's own (issue 352). f.rolledBack[node] = nil r := f.reports[node] - r.Current = false + r.Current, r.Declared, r.ReportedSequence = false, "", 0 f.reports[node] = r } } diff --git a/cmd/mesh-controller/issue352_test.go b/cmd/mesh-controller/issue352_test.go new file mode 100644 index 00000000..951f8514 --- /dev/null +++ b/cmd/mesh-controller/issue352_test.go @@ -0,0 +1,356 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "reflect" + "strings" + "testing" + "time" + + "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" +) + +// novox/hq issue 352: on 2026-10-09 a release's gate on the control node read the machine's report against a +// newer send another plan had just made there — not against its own send — and failed three builds the +// machine had reported healthy ("has not reported on what it was sent"), put them back on every machine, +// to a controller older than the store's schema, and that controller then judged the newer plan's +// controller passed from the put-back build's health. + +// TestReplay352 replays the walk on the backlog fixture: the release sends anchor and anchor reports; +// another send reaches anchor, unreported; the gate still passes. And a send that moves a judged module +// to another build supersedes the judging: no verdict, nothing put back. +func TestReplay352(t *testing.T) { + t.Run("a newer send to the judged machine does not unreport the gate's", testANewerSendDoesNotUnreportTheGatesOwn) + t.Run("a send that moves the module supersedes the judging", testASendThatMovesTheModuleSupersedesTheJudging) +} + +func testANewerSendDoesNotUnreportTheGatesOwn(t *testing.T) { + b := aBacklog(t) + ctx := t.Context() + inv := b.open.inventory + advancePlans(ctx, b.open) // anchor is sent, and the fixture reports it applied + // 16:31:37 — another plan sends anchor a newer declaration, which it has not reported on. + carried, _, _ := inv.SentBuilds(ctx, "anchor") + if err := inv.RecordSent(ctx, nodeID(t, b.open, "anchor"), "d-anchor-newer", carried); err != nil { + t.Fatal(err) + } + reports, _ := inv.LastReports(ctx) + for _, r := range reports { + if r.Node == "anchor" && r.Current { + t.Fatal("the fixture's newer send reads as reported") + } + } + gateEvery, gateBound = 0, 0 // past the bound at once: before the fix, "has not reported" fails it here + for i := 0; i < 4; i++ { + advancePlans(ctx, b.open) + } + p := b.release(t) + if p.State == inventory.PlanFailed || strings.Contains(p.Note, "has not reported") { + t.Fatalf("the release failed on the newer send: %s %s", p.State, p.Note) + } + if g := p.Release.Gate; g != nil && (g.Sent == nil || g.Sent["anchor"].Digest == "") { + t.Fatalf("the gate does not keep what it sent: %+v", g) + } + if v, found, err := inv.GateOf(ctx, "build-app-c2"); err != nil || !found || v.Verdict != inventory.GatePassed { + t.Fatalf("app's pass on anchor was not kept: %+v %v %v", v, found, err) + } + if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c2" { + t.Fatalf("app was put back to %s", current["app"].Commit) + } +} + +// A plan's own first send waits while a release judges the same module on that machine with another build. +func TestAPlansFirstSendWaitsForAReleaseJudgingTheModuleThere(t *testing.T) { + b := aBacklog(t) + ctx := t.Context() + advancePlans(ctx, b.open) // the release judges app c2 on anchor + _, _, err := gatedSend(ctx, b.open, "anchor", []inventory.CarriedMove{{Module: "app", Node: "anchor", From: "c2", To: "c3", Build: "build-app-c3"}}) + if !errors.Is(err, errWalkedElsewhere) || !strings.Contains(err.Error(), "release-") { + t.Fatalf("a newer build of a judged module was sent under the release's gate: %v", err) + } + if len(b.sent) != 1 { + t.Fatalf("sent %v", b.sent) + } +} + +// A merge plan's judging is superseded the same way: another send moved its module on the first machine. +func TestAPlansJudgingIsSupersededByASendThatMovesItsModule(t *testing.T) { + g := aGateMesh(t) + ctx := t.Context() + inv := g.open.inventory + advancePlans(ctx, g.open) // anchor is sent app c2 first + if err := inv.RecordSent(ctx, nodeID(t, g.open, "anchor"), "d-anchor-c3", map[string]string{"app": "c3"}); err != nil { + t.Fatal(err) + } + gateEvery = 0 + advancePlans(ctx, g.open) + p := g.plan(t) + if p.State != inventory.PlanSuperseded || !strings.Contains(p.Note, "superseded") || !strings.Contains(p.Note, "c3") { + t.Fatalf("the plan is %s: %s", p.State, p.Note) + } + // Nothing put back: the registered build stands, the build is not marked, and the plan's gate made no + // rollback (a release may walk what the other send left waiting on anchor; that is not a put-back). + if r := p.Modules["app"].Gate.Rollback; r != "" { + t.Fatalf("a superseded judging made a rollback: %q", r) + } + if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c2" { + t.Fatalf("app was put back to %s", current["app"].Commit) + } + if failed, _ := inv.GateFailed(ctx, "build-2"); failed { + t.Fatal("a superseded build was marked failed") + } +} + +func testASendThatMovesTheModuleSupersedesTheJudging(t *testing.T) { + b := aBacklog(t) + ctx := t.Context() + inv := b.open.inventory + advancePlans(ctx, b.open) + // Another send moves app on anchor to a build this gate does not judge. + carried, _, _ := inv.SentBuilds(ctx, "anchor") + carried["app"] = "c3" + if err := inv.RecordSent(ctx, nodeID(t, b.open, "anchor"), "d-anchor-c3", carried); err != nil { + t.Fatal(err) + } + gateEvery = 0 + advancePlans(ctx, b.open) + p := b.release(t) + if p.State != inventory.PlanSuperseded || !strings.Contains(p.Note, "superseded") || !strings.Contains(p.Note, "c3") { + t.Fatalf("the release is %s: %s", p.State, p.Note) + } + if !reflect.DeepEqual(b.sent, [][]string{{"anchor"}}) { + t.Fatalf("sent %v: a superseded judging puts nothing back", b.sent) + } + if _, found, _ := inv.GateOf(ctx, "build-app-c2"); found { + t.Fatal("a superseded judging kept a verdict") + } + if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c2" { + t.Fatalf("app was put back to %s", current["app"].Commit) + } +} + +// A report is on the gate's own send: the declaration itself, or one sequenced after it; a gate kept +// without its send reads the report against the send made last, as before. Pure. +func TestAReportIsHeldAgainstTheGatesOwnSend(t *testing.T) { + sent := inventory.SentDeclaration{Digest: "d-490", Sequence: 490} + for _, c := range []struct { + r inventory.Reported + want bool + }{ + {inventory.Reported{Declared: "d-490", Current: false}, true}, + {inventory.Reported{Declared: "d-491", ReportedSequence: 491, Current: true}, true}, + {inventory.Reported{Declared: "d-489", ReportedSequence: 489, Current: false}, false}, + {inventory.Reported{Declared: "other", ReportedSequence: 490}, true}, // the same sequence, said by another digest + {inventory.Reported{Declared: "d-495", ReportedSequence: 495, Current: true}, true}, // the last send: this one or a later one + {inventory.Reported{Declared: "", Current: false}, false}, + } { + if got := sent.ReportsOn(c.r); got != c.want { + t.Errorf("%+v on %+v: %v", c.r, sent, got) + } + } + byDigest := inventory.SentDeclaration{Digest: "d-1"} + if !byDigest.ReportsOn(inventory.Reported{Declared: "d-1"}) || byDigest.ReportsOn(inventory.Reported{ReportedSequence: 5}) || + !byDigest.ReportsOn(inventory.Reported{Current: true}) { + t.Fatal("a send kept without a sequence is matched by its digest and by the last send alone") + } + f := gateFacts{sent: map[string]inventory.SentDeclaration{"anchor": sent}} + if !f.reportedOn("anchor", inventory.Reported{Declared: "d-490"}) || f.reportedOn("anchor", inventory.Reported{Declared: "d-1"}) { + t.Fatal("a gate that kept its send read the report against something other than it") + } + if !f.reportedOn("laptop", inventory.Reported{Current: true}) || f.reportedOn("laptop", inventory.Reported{Current: false}) { + t.Fatal("a gate that did not keep its send does not read the report against the send made last") + } + // Through the merge plan's first-machine wait too. + at := time.Now() + state := inventory.PlanModule{First: []string{"anchor"}, FirstAt: &at, + Gate: &inventory.PlanGate{Machines: []string{"anchor"}, Sent: map[string]inventory.SentDeclaration{"anchor": sent}}} + reports := []inventory.Reported{{Node: "anchor", At: &at, Outcome: inventory.OutcomeApplied, Current: false, Declared: "d-490"}} + if step := nextRollout(state, []string{"anchor", "laptop"}, false, reports, at.Add(time.Minute), time.Hour); step.waiting != "" || step.failed != "" { + t.Fatalf("the first machine's report on the plan's own send read as none: %+v", step) + } + reports[0].Declared = "d-480" + if step := nextRollout(state, []string{"anchor", "laptop"}, false, reports, at.Add(time.Minute), time.Hour); step.waiting == "" { + t.Fatalf("a report on an older send read as the plan's: %+v", step) + } +} + +// A gate judges only the build the machine was last sent: last sent another build of the module, the +// judging is superseded, whatever the machine reports. Pure. +func TestAGateJudgesOnlyTheBuildTheMachineWasLastSent(t *testing.T) { + at := time.Now() + f := gateFacts{now: at, reports: map[string]inventory.Reported{"anchor": {Node: "anchor", Outcome: inventory.OutcomeApplied, + At: &at, Current: true}}, engines: map[string]string{}, served: map[string]served{}, rolledBack: map[string][]lease.Rollback{}, + commits: map[string]string{"mesh-controller": "e6b00e2e"}, sentBuilds: map[string]map[string]string{"anchor": {"mesh-controller": "ef26d4cb"}}} + taken := at.Add(-30 * time.Second) + f.holder = &lease.Holder{Taken: taken, Health: &lease.Health{Ready: true}} + h, why := judgeHealth("mesh-controller", lease.ComponentController, catalogue.Manifest{}, "anchor", at.Add(-time.Minute), f) + if h != healthSuperseded || !strings.Contains(why, "ef26d4cb") || !strings.Contains(why, "e6b00e2e") { + t.Fatalf("a controller build the machine no longer runs: %v %q", h, why) + } + f.sentBuilds["anchor"]["mesh-controller"] = "e6b00e2e" + if h, why := judgeHealth("mesh-controller", lease.ComponentController, catalogue.Manifest{}, "anchor", at.Add(-time.Minute), f); h == healthSuperseded { + t.Fatalf("the build sent read as another: %q", why) + } + delete(f.sentBuilds, "anchor") + if h, why := judgeHealth("mesh-controller", lease.ComponentController, catalogue.Manifest{}, "anchor", at.Add(-time.Minute), f); h == healthSuperseded { + t.Fatalf("a machine whose send is not known read as superseded: %q", why) + } + g := &inventory.PlanGate{To: "c2", Carried: []inventory.CarriedMove{{Module: "late", Node: "anchor", To: "c5"}}} + if got := judgedCommits(g, []judged{{"app", "anchor"}, {"late", "anchor"}}); got["app"] != "c2" || got["late"] != "c5" { + t.Fatalf("judged commits %v", got) + } +} + +// A move of another build of a module to a machine where a release or a plan is judging that module +// waits for that judging; the same build to that machine is already there. Pure. +func TestAWalkWaitsForAJudgingOfTheSameModuleOnThatMachine(t *testing.T) { + at := time.Now() + release := inventory.Plan{ID: "release-1", State: inventory.PlanRolling, Release: &inventory.PlanRelease{ + Gate: &inventory.PlanGate{Machines: []string{"novox"}, Carried: []inventory.CarriedMove{ + {Module: "mesh-controller", Node: "novox", From: "ef26d4cb", To: "2913c54c"}}}}} + merge := inventory.Plan{ID: "plan-1", State: inventory.PlanRolling, Modules: map[string]*inventory.PlanModule{ + "app": {First: []string{"anchor"}, FirstAt: &at, Commit: "c2", Gate: &inventory.PlanGate{Machines: []string{"anchor"}}}}} + f := moveFacts{plans: []inventory.Plan{release, merge}} + for _, c := range []struct { + module, node, to, want string + }{ + {"mesh-controller", "novox", "e6b00e2e", "release-1"}, // the day's case: a newer controller to the judged machine + {"mesh-controller", "novox", "2913c54c", ""}, // the same build: already there + {"mesh-controller", "ace", "2913c54c", "release-1"}, // another machine while the first is judged + {"mesh-host", "novox", "x", ""}, // a module the release does not carry + {"app", "anchor", "c2", ""}, + {"app", "anchor", "c3", "plan-1"}, + {"app", "laptop", "c2", "plan-1"}, + } { + if got := f.walkedBy(c.module, c.node, c.to); got != c.want { + t.Errorf("%s %s to %s: walked by %q, want %q", c.module, c.to, c.node, got, c.want) + } + } + release.Release.Gate.Verdict = inventory.GatePassed + merge.Modules["app"].Gate.Verdict = inventory.GatePassed + if f.walkedBy("mesh-controller", "novox", "e6b00e2e") != "" || f.walkedBy("app", "laptop", "c3") != "" { + t.Fatal("a passed judging still holds a move") + } + release.Release.Gate.Verdict = "" + f.plans[0].State = inventory.PlanSuperseded + if f.walkedBy("mesh-controller", "novox", "e6b00e2e") != "" { + t.Fatal("a closed release still holds a move") + } +} + +// The controller is never put back to a build that reaches less of the store's schema than the store +// has, or to one that never said what it reaches: the current build is kept, and the condition says so. +func TestTheControllerIsNotPutBackToABuildOlderThanTheStoresSchema(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + inv := open.inventory + keeper, _ := withConditionsInMemory(t) + told := &conditions.Told{} + was := doctorFrom + doctorFrom = &doctor{open: open, keeper: keeper, teller: told} + t.Cleanup(func() { doctorFrom = was }) + wasSend := sendRollout + var sent [][]string + sendRollout = func(ctx context.Context, open *stores, names []string) ([]string, error) { + sent = append(sent, names) + return names, nil + } + t.Cleanup(func() { sendRollout = wasSend }) + + build := func(id, commit, digest string, asked time.Time) inventory.Build { + manifest, _ := json.Marshal(catalogue.Manifest{Module: "mesh-controller", Version: commit}) + b := inventory.Build{ID: id, Module: "mesh-controller", Commit: commit, Repository: "novox/mesh-controller", Path: ".", + Manifest: manifest, Asked: asked, At: asked, Made: []inventory.Artifact{{Name: "controller", Kind: catalogue.ArtifactBundle, + Reference: "mesh-artifact://mesh-controller/controller/blobs/sha256:" + digest}}} + if err := inv.RecordBuild(ctx, b); err != nil { + t.Fatal(err) + } + if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "mesh-controller", Version: commit}, + inventory.Source{Repository: "novox/mesh-controller", Seat: "git", Path: ".", BuiltFrom: commit, Head: commit, Asked: asked}); err != nil { + t.Fatal(err) + } + return b + } + previous := build("build-old", "ef26d4cb", strings.Repeat("1", 64), time.Now().Add(-2*time.Hour)) + failed := build("build-new", "e6b00e2e", strings.Repeat("2", 64), time.Now().Add(-time.Minute)) + if versionOfBuild(previous) != strings.Repeat("1", 12) { + t.Fatalf("the build's version is %q", versionOfBuild(previous)) + } + applied, err := inv.SchemaApplied(ctx) + if err != nil || applied < 87 { + t.Fatalf("the store's schema reaches %d (%v)", applied, err) + } + // The build before never recorded what it reads: not proved, refused. + if why, ok := controllerSchemaAllows(ctx, inv, previous); ok || !strings.Contains(why, "never recorded") { + t.Fatalf("an unknown reach: %v %q", ok, why) + } + // It reads less than the store has: refused, naming both. + if err := inv.RecordSchemaReach(ctx, versionOfBuild(previous), applied-1); err != nil { + t.Fatal(err) + } + if why, ok := controllerSchemaAllows(ctx, inv, previous); ok || !strings.Contains(why, "is at") { + t.Fatalf("a reach behind the store: %v %q", ok, why) + } + // Through the gate: the failed build is marked, nothing is put back, nothing is sent, the condition is urgent. + at := time.Now().Add(-5 * time.Minute) + state := &inventory.PlanModule{Build: failed.ID, Commit: failed.Commit, Previous: previous.Commit, First: []string{"anchor"}, FirstAt: &at} + p := inventory.Plan{ID: "plan-352", Repository: "novox/mesh-controller", Branch: "main", Commit: failed.Commit, Created: at, + State: inventory.PlanRolling, Tiers: [][]string{{"mesh-controller"}}, Modules: map[string]*inventory.PlanModule{"mesh-controller": state}} + if err := inv.SavePlan(ctx, &p); err != nil { + t.Fatal(err) + } + gateFailed(ctx, open, &p, "mesh-controller", state, []string{"anchor"}, "not healthy within 10m0s of its apply") + if state.Gate.Rollback != inventory.NotRolledBack || !strings.Contains(p.Note, "NOT put back") || !strings.Contains(p.Note, "the current build is kept") { + t.Fatalf("rollback %q: %s", state.Gate.Rollback, p.Note) + } + if len(sent) != 0 { + t.Fatalf("sent %v: nothing is put back", sent) + } + if current, _ := inv.CurrentBuilds(ctx); current["mesh-controller"].Commit != failed.Commit { + t.Fatalf("the module was put back to %s", current["mesh-controller"].Commit) + } + if marked, _ := inv.GateFailed(ctx, failed.ID); !marked { + t.Fatal("the failed build is not marked failed at its gate") + } + open2, _ := keeper.Open(ctx) + var found bool + for _, c := range open2 { + if c.Kind == kindRollbackFailed && c.Severity == conditions.Urgent && strings.Contains(c.Summary, "current build is kept") { + found = true + } + } + if !found { + t.Fatalf("no urgent rollback-failed condition saying the current build is kept: %+v", open2) + } + // Reaching the store: allowed. + if err := inv.RecordSchemaReach(ctx, versionOfBuild(previous), applied); err != nil { + t.Fatal(err) + } + if why, ok := controllerSchemaAllows(ctx, inv, previous); !ok { + t.Fatalf("a build that reads the whole schema was refused: %q", why) + } + // A build with no bundle to know it by: refused. + if why, ok := controllerSchemaAllows(ctx, inv, inventory.Build{Commit: "x"}); ok || !strings.Contains(why, "names no bundle") { + t.Fatalf("a build without a bundle: %v %q", ok, why) + } +} + +// The lease's holder names the build the declaration told it it is, and the version stamp only without one. +func TestTheHolderNamesTheBuildTheDeclarationToldIt(t *testing.T) { + t.Setenv(RunningBuildVar, " ad62528c47c7 ") + if h := holderOf("x"); h.Build != "ad62528c47c7" { + t.Fatalf("the holder's build is %q", h.Build) + } + t.Setenv(RunningBuildVar, "") + if h := holderOf("x"); h.Build != version { + t.Fatalf("without a declared version the holder's build is %q", h.Build) + } + if reach, err := schemaReach(); err != nil || reach < 87 { + t.Fatalf("this build's reach is %d (%v)", reach, err) + } +} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 230f4611..5b0e2605 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -76,6 +76,21 @@ func serve(ctx context.Context) (err error) { } defer open.Close() inv := open.inventory + // **What this build reads of the store's schema, on record** (novox/hq issue 352): the highest migration + // it carries, by its version, so a gate that would put this build back later knows it reads the store + // as it is then. A build that does not know its version records nothing, and is never put back. + if build := runningBuild(); build != "" { + if reach, err := schemaReach(); err != nil { + fmt.Printf("what this build reads of the store's schema is not recorded: %v\n", err) + } else if err := inv.RecordSchemaReach(ctx, build, reach); err != nil { + fmt.Printf("what this build (%s) reads of the store's schema is not recorded: %v\n", build, err) + } else { + fmt.Printf("this build (%s) reads the store's schema up to migration %04d; recorded\n", build, reach) + } + } else { + fmt.Printf("this process was not told which build it is (%s), so what it reads of the store's schema is not "+ + "recorded, and a gate will never put it back\n", RunningBuildVar) + } ident, err := openIdentity(ctx) if err != nil { @@ -1598,3 +1613,18 @@ func reportUnheldPushed(w io.Writer, named bool, asked []string, unheld map[stri fmt.Fprintf(w, "%s: %d unmet seat dependenc(ies) — see `status`\n", node, len(lines)) } } + +// schemaReach is the highest migration this build carries for the inventory's store (novox/hq issue 352). +func schemaReach() (int, error) { + migrations, err := inventory.Migrations() + if err != nil { + return 0, err + } + reach := 0 + for _, m := range migrations { + if m.Number > reach { + reach = m.Number + } + } + return reach, nil +} diff --git a/cmd/mesh-controller/release.go b/cmd/mesh-controller/release.go index 421a01f8..56b702c8 100644 --- a/cmd/mesh-controller/release.go +++ b/cmd/mesh-controller/release.go @@ -180,12 +180,35 @@ func (f moveFacts) moves(node string, modules []string, sent map[string]string, return out } -// walkedBy is the open plan that has started walking a module's build — sent it to a first machine, -// not yet passed — other than to this machine; empty when none does. -func (f moveFacts) walkedBy(module, node string) string { +// walkedBy is the open plan that has started walking a module's build — sent it to a first machine, not +// yet passed — and whose walk this move would cross; empty when none does. A move of the same build to a +// machine that plan already sent it is not a crossing: the build is there. A move of **another** build of +// the module to that machine is (novox/hq issue 352): on 2026-10-09 a merge's plan sent the control node a +// newer controller while a release's gate was judging the controller there, the release's gate read the +// machine's report against the newer send, failed three builds and put them back under the new plan's +// feet. A release keeps no module records: its open gate's carried moves are its walk. +func (f moveFacts) walkedBy(module, node, to string) string { for _, p := range f.plans { + if !p.Open() { + continue + } + if p.Release != nil { + g := p.Release.Gate + if g == nil || g.Verdict != "" { + continue + } + for _, c := range g.Carried { + if c.Module == module && !(slices.Contains(g.Machines, node) && sameCommit(c.To, to)) { + return p.ID + } + } + continue + } s, holds := p.Modules[module] - if !p.Open() || !holds || s == nil || s.FirstAt == nil || s.SentAt != nil || slices.Contains(s.First, node) { + if !holds || s == nil || s.FirstAt == nil || s.SentAt != nil { + continue + } + if slices.Contains(s.First, node) && sameCommit(s.Commit, to) { continue } if s.Gate != nil && s.Gate.Verdict == inventory.GatePassed { @@ -277,11 +300,19 @@ func gatedSend(ctx context.Context, open *stores, node string, owns []inventory. if own(mv.Module) { continue } - if id := f.walkedBy(mv.Module, node); id != "" { + if id := f.walkedBy(mv.Module, node, mv.To); id != "" { return nil, nil, fmt.Errorf("%w: %s's build %s waits on %s, which %s is walking", errWalkedElsewhere, mv.Module, short(mv.To), node, id) } } + // **And the plan's own modules wait too** (novox/hq issue 352): a walk of the same module by another + // plan, or a release, on this machine is not crossed with a newer build; this send waits for its gate. + for _, o := range owns { + if id := f.walkedBy(o.Module, node, o.To); id != "" { + return nil, nil, fmt.Errorf("%w: %s's build %s waits on %s, which %s is walking", errWalkedElsewhere, + o.Module, short(o.To), node, id) + } + } for _, o := range owns { i := slices.IndexFunc(moves, func(mv inventory.CarriedMove) bool { return mv.Module == o.Module }) switch { @@ -526,7 +557,7 @@ func waitingMoves(ctx context.Context, open *stores, all bool) (map[string][]inv return nil, err } for _, mv := range moves { - if f.walkedBy(mv.Module, n.Name) == "" { + if f.walkedBy(mv.Module, n.Name, mv.To) == "" { out[n.Name] = append(out[n.Name], mv) } } @@ -658,7 +689,7 @@ func advanceRelease(ctx context.Context, open *stores, p *inventory.Plan) (bool, r.Next++ continue } - r.Gate = &inventory.PlanGate{Machines: sent, Since: &now, Carried: moves} + r.Gate = &inventory.PlanGate{Machines: sent, Since: &now, Carried: moves, Sent: sentNow(ctx, open.inventory, sent)} p.Note = fmt.Sprintf("sent %s %d build(s) that waited for a gate; judging them there", node, len(moves)) if said := recreationsSaid(moves); said != "" { p.Note += "; " + said @@ -693,6 +724,15 @@ func advanceRelease(ctx context.Context, open *stores, p *inventory.Plan) (bool, r.Next++ r.Gate = nil return true, nil + case inventory.GateSuperseded: + // Another send moved a judged module on the judged machine (novox/hq issue 352): no verdict on what + // was carried, nothing put back, and this release ends; the builds still waiting are released again + // by the next pass, judged afresh. + p.State = inventory.PlanSuperseded + p.Note = fmt.Sprintf("superseded on %s: %s — nothing judged, nothing put back; what still waits is released again", + strings.Join(g.Machines, ", "), g.Why) + fmt.Printf("%s: %s\n", p.ID, p.Note) + return true, nil } p.Note = "" batched, back := batchingRollbacks(ctx) diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index c6b5dfb1..31c90340 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -793,6 +793,15 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, case inventory.GateFailed: failFirstSend(ctx, open, p, m, state, state.Gate.Machines, state.Gate.Why, step.rest) return true, nil + case inventory.GateSuperseded: + // Another send moved this module on its first machine (novox/hq issue 352): the build is not + // judged, not marked, not put back; the plan ends here, said, and a newer plan carries on. + state.Why = "superseded: " + state.Gate.Why + p.State = inventory.PlanSuperseded + p.Note = fmt.Sprintf("%s's judging on %s was superseded: %s", m, strings.Join(state.Gate.Machines, ", "), + state.Gate.Why) + fmt.Printf("%s: %s\n", p.ID, p.Note) + return true, nil case inventory.GatePassed: if !state.Gate.Kept { gatePassed(ctx, open, p, m, state) @@ -964,6 +973,9 @@ func firstSend(ctx context.Context, open *stores, p *inventory.Plan, node string } now := time.Now().UTC() lead := modules[0] + // What each machine was just sent, kept on the gate (novox/hq issue 352): its report is held against + // this send, whatever it is sent after. + sentWhat := sentNow(ctx, inv, sent) for _, m := range modules { s := p.Modules[m] // What the first machine ran before: what a failed gate puts back (ADR 0236). @@ -972,7 +984,7 @@ func firstSend(ctx context.Context, open *stores, p *inventory.Plan, node string } s.First, s.FirstAt = sent, &now s.Gate = &inventory.PlanGate{Component: coreComponent(m), Machines: firstRunning(sent, runningOf[m]), - From: s.Previous, To: s.Commit, Since: &now} + From: s.Previous, To: s.Commit, Since: &now, Sent: sentWhat} s.GatedBy = "" if m == lead { s.Gate.Carried = carried @@ -1131,8 +1143,15 @@ func nextRollout(s inventory.PlanModule, running []string, together bool, report var waiting, failed []string for _, n := range s.First { r, said := byNode[n] - // Only a report about what it was last sent says anything about this build. - if !said || r.At == nil || !r.Current { + // Only a report about what this plan sent it — or what it was sent after that — says anything about + // this build (novox/hq issue 352); a plan from before sends were kept on the gate reads the report + // against the send made last, as before. + reported := r.Current + if s.Gate != nil && s.Gate.Sent != nil { + sent, kept := s.Gate.Sent[n] + reported = kept && sent.ReportsOn(r) + } + if !said || r.At == nil || !reported { waiting = append(waiting, n) continue } diff --git a/internal/catalogue/build.go b/internal/catalogue/build.go index 388b769c..23067e2f 100644 --- a/internal/catalogue/build.go +++ b/internal/catalogue/build.go @@ -173,11 +173,7 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) { // had — so re-composing a declaration moves nothing, where a commit would move the // path of an identical binary and recreate everything that reads it. for key, value := range filled { - text, isText := value.(string) - if !isText || !strings.Contains(text, versionRef) { - continue - } - filled[key] = strings.ReplaceAll(text, versionRef, versionOf(artifact.Digest)) + filled[key] = withVersion(value, versionOf(artifact.Digest)) } default: return Manifest{}, fmt.Errorf("%s: %q is a %q, and an artifact is %q, %q, %q or %q", @@ -189,6 +185,31 @@ func (m Manifest) Resolve(built []Built) (Manifest, error) { return out, nil } +// withVersion is a resource's value with `${version}` filled: in a string, and in each string of a map — +// a process's env, where the controller is told which build it is (novox/hq issue 352). Anything else is +// left as it is. +func withVersion(value any, version string) any { + switch v := value.(type) { + case string: + if strings.Contains(v, versionRef) { + return strings.ReplaceAll(v, versionRef, version) + } + case map[string]any: + out := make(map[string]any, len(v)) + for k, x := range v { + out[k] = withVersion(x, version) + } + return out + case map[string]string: + out := make(map[string]string, len(v)) + for k, x := range v { + out[k] = strings.ReplaceAll(x, versionRef, version) + } + return out + } + return value +} + // checkBuild is the manifest's own account of what it builds. func (b *Build) problems(module string) []string { if b == nil { diff --git a/internal/catalogue/version_in_a_path_test.go b/internal/catalogue/version_in_a_path_test.go index 7f536754..2c2e4d44 100644 --- a/internal/catalogue/version_in_a_path_test.go +++ b/internal/catalogue/version_in_a_path_test.go @@ -95,3 +95,28 @@ func TestAResourceWithoutAVersionReferenceIsUntouched(t *testing.T) { t.Fatalf("a path naming no version became %q", path) } } + +// A process's env can name the build's own version too (novox/hq issue 352): the controller is told which +// build it is, and records what that build reads of the store's schema under it. +func TestAProcessEnvCanNameTheBuildsOwnVersion(t *testing.T) { + m := Manifest{ + Module: "mesh-controller", + Build: &Build{Artifacts: []Artifact{{Name: "controller", Kind: ArtifactBundle, Language: "go", System: "arch"}}}, + Resources: []map[string]any{{ + "id": "controller", "type": "process", "artifact": "controller", "run": []any{"./mesh-controller", "serve"}, + "env": map[string]any{"MESH_CONTROLLER_VERSION": "${version}", "OTHER": "kept"}, + }}, + } + got, err := m.Resolve([]Built{{Name: "controller", Kind: ArtifactBundle, + Reference: "artifact-store://mesh-controller/controller", Digest: aDigest}}) + if err != nil { + t.Fatal(err) + } + env, _ := got.Resources[0]["env"].(map[string]any) + if env["MESH_CONTROLLER_VERSION"] != "ad62528c47c7" || env["OTHER"] != "kept" { + t.Fatalf("the env resolved to %v", env) + } + if run, _ := got.Resources[0]["run"].([]any); len(run) != 2 { + t.Fatalf("the run was changed: %v", run) + } +} diff --git a/internal/inventory/gate.go b/internal/inventory/gate.go index 7baa82e7..b5b76768 100644 --- a/internal/inventory/gate.go +++ b/internal/inventory/gate.go @@ -23,6 +23,9 @@ import ( const ( GatePassed = "passed" GateFailed = "failed" + // GateSuperseded is a judging ended by a later send to the judged machine that moved the module to + // another build (novox/hq issue 352): no verdict on the build, and nothing put back. + GateSuperseded = "superseded" RollingBack = "rolling-back" RolledBack = "rolled-back" diff --git a/internal/inventory/migrations/0087-a-controller-records-what-its-stores-schema-reaches.sql b/internal/inventory/migrations/0087-a-controller-records-what-its-stores-schema-reaches.sql new file mode 100644 index 00000000..4f86d122 --- /dev/null +++ b/internal/inventory/migrations/0087-a-controller-records-what-its-stores-schema-reaches.sql @@ -0,0 +1,10 @@ +-- A controller records, when it serves, how far the store's schema reaches in the build it is (novox/hq +-- issue 352): the highest migration it carries, by its build's version. A gate that fails the controller's +-- build puts the build before it back, and on 2026-10-09 that build was older than the migrations the +-- failed one had applied: it started, said it was behind its own row, and judged the next gate half-blind. +-- A put-back now reads this and keeps the current build when the one before it reaches less than the store. +create table controller_schema ( + build text primary key, + reach integer not null, + recorded timestamptz not null default now() +); diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index 03f970b6..4c88e0e3 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -1054,13 +1054,57 @@ type Reported struct { // acted on the current words, not merely spoken after they were written. False also covers // a machine that has not said which, which is every host from before reports carried it. Current bool + // Declared is the digest of the declaration the last report was about, and ReportedSequence that + // declaration's sequence as the report claimed it (zero from an engine that claims none): what a + // gate holds against the send it made, rather than against the send made last (novox/hq issue 352). + Declared string + ReportedSequence int64 +} + +// SentDeclaration is what one send carried to a machine, as a gate keeps it: the declaration's digest and +// its sequence (novox/hq issue 352). A report about this declaration, or about one sequenced after it, is a +// report on what the gate sent — whatever the machine was sent since. +type SentDeclaration struct { + Digest string `json:"digest"` + Sequence int64 `json:"sequence,omitempty"` +} + +// ReportsOn says a report is about this send: the declaration itself; one the same machine was sequenced +// after it; or the declaration the machine was sent last (Current), which is this send or a later one — +// sends to a machine are made one after another. A send kept without a sequence is matched by its digest +// and by the last send alone. +func (s SentDeclaration) ReportsOn(r Reported) bool { + if r.Current || (s.Digest != "" && r.Declared == s.Digest) { + return true + } + return s.Sequence > 0 && r.ReportedSequence >= s.Sequence +} + +// SentTo is the declaration a machine was last sent, by name: its digest and sequence, and false when it +// was never sent one. +func (i *Inventory) SentTo(ctx context.Context, name string) (SentDeclaration, bool, error) { + var digest *string + var seq *int64 + err := i.store.Pool().QueryRow(ctx, `select sent, sequence from node where name = $1`, name).Scan(&digest, &seq) + if errors.Is(err, pgx.ErrNoRows) { + return SentDeclaration{}, false, fmt.Errorf("%w: %s", ErrNoSuchNode, name) + } + if err != nil || digest == nil || *digest == "" { + return SentDeclaration{}, false, err + } + s := SentDeclaration{Digest: *digest} + if seq != nil { + s.Sequence = *seq + } + return s, true, nil } // LastReports is every machine's last report beside when it was last sent a declaration. func (i *Inventory) LastReports(ctx context.Context) ([]Reported, error) { rows, err := i.store.Pool().Query(ctx, `select n.name, coalesce(r.outcome, ''), r.at, n.sent_at, - r.declared is not null and r.declared <> '' and r.declared = n.sent + r.declared is not null and r.declared <> '' and r.declared = n.sent, + coalesce(r.declared, ''), coalesce(r.reported_sequence, 0) from node n left join node_report r on r.node = n.id order by n.name`) if err != nil { @@ -1070,7 +1114,7 @@ func (i *Inventory) LastReports(ctx context.Context) ([]Reported, error) { var out []Reported for rows.Next() { var r Reported - if err := rows.Scan(&r.Node, &r.Outcome, &r.At, &r.Sent, &r.Current); err != nil { + if err := rows.Scan(&r.Node, &r.Outcome, &r.At, &r.Sent, &r.Current, &r.Declared, &r.ReportedSequence); err != nil { return nil, err } out = append(out, r) diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index bf1d1532..59f832b8 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -126,6 +126,10 @@ type PlanGate struct { To string `json:"to,omitempty"` // Since is when the judging began: the first machine reported the new build applied. Since *time.Time `json:"since,omitempty"` + // Sent is, per machine, the declaration the gate's send carried there (novox/hq issue 352): what a + // machine's report is held against. Absent on a gate kept before it was, which reads the report + // against the send made last, as before. + Sent map[string]SentDeclaration `json:"sent,omitempty"` // Passes counts the consecutive judgings that found it healthy, LastPass the newest; a judging that // does not resets them. Passes int `json:"passes,omitempty"` diff --git a/internal/inventory/plans_test.go b/internal/inventory/plans_test.go index b92c5195..ea84f117 100644 --- a/internal/inventory/plans_test.go +++ b/internal/inventory/plans_test.go @@ -115,3 +115,83 @@ func TestTheNewestMergeOfABranchIsTheOneMergedLast(t *testing.T) { t.Fatalf("one merge time, two plans: %s", p.ID) } } + +// novox/hq issue 352: what a machine was last sent is read back by name with its sequence, a report keeps +// the declaration it was about and that declaration's sequence, and a gate's sends are kept with the plan. +func TestASendAndAReportAreKnownByTheirDeclaration(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + record, err := inv.AddNode(ctx, "anchor") + if err != nil { + t.Fatal(err) + } + if _, found, err := inv.SentTo(ctx, "anchor"); err != nil || found { + t.Fatalf("a machine never sent anything: %v %v", found, err) + } + if _, _, err := inv.SentTo(ctx, "nobody"); err == nil { + t.Fatal("a machine that does not exist was answered") + } + seq, err := inv.NextSequence(ctx, record.ID) + if err != nil { + t.Fatal(err) + } + if err := inv.RecordSent(ctx, record.ID, "d-1", map[string]string{"app": "c1"}); err != nil { + t.Fatal(err) + } + sent, found, err := inv.SentTo(ctx, "anchor") + if err != nil || !found || sent.Digest != "d-1" || sent.Sequence != seq { + t.Fatalf("sent %+v %v %v", sent, found, err) + } + if _, err := inv.RecordOrderedDoing(ctx, record.ID, Doing{Node: "anchor", Outcome: OutcomeApplied, Declared: "d-1", Applied: 1}, + ReportOrder{Sequence: seq, ReportSequence: 1}, func(ReportOrder) bool { return false }); err != nil { + t.Fatal(err) + } + reports, err := inv.LastReports(ctx) + if err != nil || len(reports) != 1 || reports[0].Declared != "d-1" || reports[0].ReportedSequence != seq || !reports[0].Current { + t.Fatalf("reports %+v %v", reports, err) + } + // Sent again, unreported: the report is no longer on the last send, and is still on the first. + if err := inv.RecordSent(ctx, record.ID, "d-2", map[string]string{"app": "c1"}); err != nil { + t.Fatal(err) + } + reports, _ = inv.LastReports(ctx) + if reports[0].Current || !sent.ReportsOn(reports[0]) { + t.Fatalf("after a newer send: %+v", reports[0]) + } + at := time.Now().UTC() + p := Plan{ID: "plan-352", Repository: "novox/x", Commit: "c", Created: at, State: PlanRolling, Tiers: [][]string{{"app"}}, + Modules: map[string]*PlanModule{"app": {Gate: &PlanGate{Machines: []string{"anchor"}, Since: &at, + Sent: map[string]SentDeclaration{"anchor": sent}}}}} + if err := inv.SavePlan(ctx, &p); err != nil { + t.Fatal(err) + } + kept, err := inv.PlanByID(ctx, "plan-352") + if err != nil || kept.Modules["app"].Gate.Sent["anchor"] != sent { + t.Fatalf("the gate's send was not kept with the plan: %+v %v", kept.Modules["app"].Gate, err) + } +} + +// A controller build's reach of the store's schema is kept by its version, and the store's own is read. +func TestASchemaReachIsKeptByBuild(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + applied, err := inv.SchemaApplied(ctx) + if err != nil || applied < 87 { + t.Fatalf("applied %d %v", applied, err) + } + if _, known, err := inv.SchemaReachOf(ctx, "ad62528c47c7"); err != nil || known { + t.Fatalf("an unrecorded build: %v %v", known, err) + } + if err := inv.RecordSchemaReach(ctx, "", 87); err == nil { + t.Fatal("a reach without a build was recorded") + } + if err := inv.RecordSchemaReach(ctx, "ad62528c47c7", 86); err != nil { + t.Fatal(err) + } + if err := inv.RecordSchemaReach(ctx, "ad62528c47c7", 87); err != nil { + t.Fatal(err) + } + if reach, known, err := inv.SchemaReachOf(ctx, "ad62528c47c7"); err != nil || !known || reach != 87 { + t.Fatalf("reach %d %v %v", reach, known, err) + } +} diff --git a/internal/inventory/schema.go b/internal/inventory/schema.go new file mode 100644 index 00000000..fa971828 --- /dev/null +++ b/internal/inventory/schema.go @@ -0,0 +1,47 @@ +package inventory + +import ( + "context" + "errors" + + "github.com/jackc/pgx/v5" +) + +// What a controller build knows of the store's schema (novox/hq issue 352): the highest migration it +// carries, recorded by its version when it serves, and the highest migration the store has applied. A +// put-back of the controller to a build that reaches less than the store is refused (gate.go), because +// such a controller starts behind its own records and judges with what it can read. + +// RecordSchemaReach keeps the highest migration the build serving now carries. +func (i *Inventory) RecordSchemaReach(ctx context.Context, build string, reach int) error { + if build == "" { + return errors.New("a schema reach is recorded by a build's version, and this controller has none") + } + _, err := i.store.Pool().Exec(ctx, + `insert into controller_schema (build, reach) values ($1, $2) + on conflict (build) do update set reach = excluded.reach, recorded = now()`, build, reach) + return err +} + +// SchemaReachOf is the highest migration a build carries, as it recorded when it served; false for a +// build that never did. +func (i *Inventory) SchemaReachOf(ctx context.Context, build string) (int, bool, error) { + var reach int + err := i.store.Pool().QueryRow(ctx, `select reach from controller_schema where build = $1`, build).Scan(&reach) + if errors.Is(err, pgx.ErrNoRows) { + return 0, false, nil + } + return reach, err == nil, err +} + +// SchemaApplied is the highest migration the store has applied. +func (i *Inventory) SchemaApplied(ctx context.Context) (int, error) { + var n *int + if err := i.store.Pool().QueryRow(ctx, `select max(number) from migration`).Scan(&n); err != nil { + return 0, err + } + if n == nil { + return 0, nil + } + return *n, nil +} diff --git a/module.json b/module.json index 3d49f65a..e0af4cbb 100644 --- a/module.json +++ b/module.json @@ -118,7 +118,8 @@ "MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}", "MESH_BROKER_MANAGEMENT_PORT": "${seat:mesh-broker:15672}", "MESH_BROKER_ADDRESS_PORT": "${seat:mesh-broker:5671}", - "MESH_BUS_NATS_FILE": "${dir:mesh-state}/bus" + "MESH_BUS_NATS_FILE": "${dir:mesh-state}/bus", + "MESH_CONTROLLER_VERSION": "${version}" }, "replaces": [ "server"