A gate judges its own send and the build it sent, and never puts the controller back behind its store (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, 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.
This commit is contained in:
@@ -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
|
||||
|
||||
+153
-1
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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"
|
||||
|
||||
+10
@@ -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()
|
||||
);
|
||||
@@ -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)
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
+2
-1
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user