From d7f359c4980a01e9fbcd655381d995b23e058d27 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 17:53:03 +0200 Subject: [PATCH 1/7] Read a new module's directory as its own, not as shared code (hq issue 252) A merge adding a module the mesh has not registered rebuilt every module built from the repository. A path under a directory known to hold modules belongs to that module when its module.json is among the changed files. --- cmd/mesh-controller/order_test.go | 8 +++++ cmd/mesh-controller/upgrades.go | 57 +++++++++++++++++++++++++++---- 2 files changed, 58 insertions(+), 7 deletions(-) diff --git a/cmd/mesh-controller/order_test.go b/cmd/mesh-controller/order_test.go index 2acac30..f4a2c62 100644 --- a/cmd/mesh-controller/order_test.go +++ b/cmd/mesh-controller/order_test.go @@ -147,6 +147,14 @@ func TestAMergeRebuildsTheModulesItChanged(t *testing.T) { {"a module the mesh does not hold", merge([]string{"modules/plex/index.ts"}, false), ""}, {"nothing said about the files", merge(nil, false), "gitea,keycloak"}, {"more files than were listed", merge([]string{"modules/gitea/index.ts"}, true), "gitea,keycloak"}, + // novox/hq issue 252: a module the mesh has never registered is still a module, when the merge + // shows it is one — and a directory that may be shared code is still shared. + {"a new module beside a held one", merge([]string{"modules/gitea/x", "modules/newmod/module.json"}, false), "gitea"}, + {"a new module's other files", merge([]string{"modules/newmod/index.ts", "modules/newmod/module.json"}, false), ""}, + {"a module removed", merge([]string{"modules/gone/module.json"}, false), ""}, + {"a directory with no manifest", merge([]string{"modules/lib/x.go"}, false), "gitea,keycloak"}, + {"a file directly among the modules", merge([]string{"modules/README.md"}, false), "gitea,keycloak"}, + {"the root's files still", merge([]string{"tsconfig.json"}, false), "gitea,keycloak"}, } { if got := named(whatTheMergeTouched(candidates, known, c.m)); got != c.want { t.Errorf("%s: rebuilt %q, wanted %q", c.what, got, c.want) diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index b16542b..f82ccb2 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -5,6 +5,7 @@ import ( "errors" "flag" "fmt" + "path" "regexp" "strings" "time" @@ -423,25 +424,45 @@ func lastLookAt(entries []inventory.Entry, m link.SourceMoved) time.Time { // // A change inside *another* module's directory is that module's business and not this one's, even // when the mesh does not hold that module: `known` is every module this repository is known to hold, -// whatever branch it was registered from. That is also the limit of this — a repository whose shared -// code sits inside a directory the mesh has never seen a module in reads as shared, and everything -// is rebuilt. Rebuilding too much is the safe direction: the fault this whole path exists for is a -// mesh that believes it is current and is not (novox/hq 04-ISSUES/131). +// whatever branch it was registered from. +// +// **And a module the mesh has never seen is still a module** (novox/hq issue 252). A merge adding a +// new module to the catalogue repository — `modules/newmod/module.json` and its files — read as a +// change to shared code, because `modules/newmod` was nobody's known directory, and every module +// built from the repository was rebuilt and rolled out for a module none of them is. So the +// directories that hold modules are known too: the parents of the known modules' directories +// (`modules`, never the root). A changed path `//…` belongs to the module at +// `/` — held or not — and rebuilds nothing else, **provided it is shown to be a +// module**: its `module.json` is among the changed files (added, changed, or removed with it). A +// directory under the same parent whose manifest the merge did not touch may as well be a shared +// library (`modules/lib`), and that is still read as shared. Rebuilding too much remains the safe +// direction: the fault this whole path exists for is a mesh that believes it is current and is not +// (novox/hq 04-ISSUES/131). A file at the root, or directly in a parent, is shared as it always was. func whatTheMergeTouched(candidates, known []inventory.Entry, m link.SourceMoved) []inventory.Entry { // Nothing said about the files, or not all of them said: everything built from it is affected. if len(m.Paths) == 0 || m.PathsTruncated { return candidates } var dirs []string + parents := map[string]bool{} for _, e := range known { if e.Source.Path != "" && sameRepository(e.Source.Repository, m) { - dirs = append(dirs, e.Source.Path) + dir := strings.Trim(e.Source.Path, "/") + dirs = append(dirs, dir) + if parent := path.Dir(dir); parent != "." && parent != "/" { + parents[parent] = true + } } } + changed := map[string]bool{} for _, p := range m.Paths { - if !insideAny(p, dirs) { - return candidates + changed[strings.TrimPrefix(p, "/")] = true + } + for _, p := range m.Paths { + if insideAny(p, dirs) || inAModuleOfItsOwn(p, parents, changed) { + continue } + return candidates } var out []inventory.Entry for _, e := range candidates { @@ -452,6 +473,28 @@ func whatTheMergeTouched(candidates, known []inventory.Entry, m link.SourceMoved return out } +// inAModuleOfItsOwn is whether a changed file is inside a module directory the mesh does not know — +// `//…` under a directory known to hold modules, whose `module.json` the same merge +// changed (novox/hq issue 252). Such a file is that module's business and nobody else's. +func inAModuleOfItsOwn(p string, parents, changed map[string]bool) bool { + p = strings.TrimPrefix(p, "/") + for parent := range parents { + rest, under := strings.CutPrefix(p, parent+"/") + if !under { + continue + } + name, _, inADirectory := strings.Cut(rest, "/") + if !inADirectory || name == "" { + // A file directly in the parent — `modules/README.md` — is about all of them. + continue + } + if changed[parent+"/"+name+"/module.json"] { + return true + } + } + return false +} + // inside is whether a changed file is in a directory: that directory itself, or under it. func inside(path, dir string) bool { dir = strings.Trim(dir, "/") From 22660dc274e7ab74bb84407e1057806a0943582e Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 17:55:01 +0200 Subject: [PATCH 2/7] Issue grants before code, and the bus's machine first (hq issue 249) A module's new state reached every machine before the permissions to use it, which came only with a later push. Memberships and buckets now go before declarations, the machine holding mesh-broker goes first when its user list must change, and a grant that fails sends nothing and is an error so the rollout is retried. --- cmd/mesh-controller/delivery_order_test.go | 98 +++++++ cmd/mesh-controller/hold_test.go | 4 +- cmd/mesh-controller/push.go | 293 +++++++++++++++------ 3 files changed, 317 insertions(+), 78 deletions(-) create mode 100644 cmd/mesh-controller/delivery_order_test.go diff --git a/cmd/mesh-controller/delivery_order_test.go b/cmd/mesh-controller/delivery_order_test.go new file mode 100644 index 0000000..7c29c00 --- /dev/null +++ b/cmd/mesh-controller/delivery_order_test.go @@ -0,0 +1,98 @@ +package main + +import ( + "context" + "errors" + "reflect" + "strings" + "testing" +) + +// recordingDelivery is a delivery that writes down what was done, in order, and fails where told. +type recordingDelivery struct { + did []string + grantErr error + declareErr error +} + +func (r *recordingDelivery) grant(_ context.Context, sending []readyNode) error { + for _, s := range sending { + r.did = append(r.did, "grant "+s.node) + } + return r.grantErr +} + +func (r *recordingDelivery) declare(_ context.Context, s readyNode, _ []byte) (string, error) { + if r.declareErr != nil { + return "", r.declareErr + } + r.did = append(r.did, "declare "+s.node) + return "digest-" + s.node, nil +} + +func ready(names ...string) []readyNode { + var out []readyNode + for _, n := range names { + out = append(out, readyNode{node: n, declared: sendable{Resources: []map[string]any{{"id": "x"}}}}) + } + return out +} + +// novox/hq issue 249: the grants that come with a module's new declarations are issued before any +// machine is sent the code that uses them. +func TestGrantsAreIssuedBeforeTheDeclarations(t *testing.T) { + d := &recordingDelivery{} + digests, err := deliver(t.Context(), d, ready("anchor", "laptop")) + if err != nil { + t.Fatal(err) + } + want := []string{"grant anchor", "grant laptop", "declare anchor", "declare laptop"} + if !reflect.DeepEqual(d.did, want) { + t.Fatalf("delivered in the order %v, wanted %v", d.did, want) + } + if digests["anchor"] != "digest-anchor" || digests["laptop"] != "digest-laptop" { + t.Fatalf("the digests sent were not answered: %v", digests) + } +} + +// A grant that cannot be issued sends nothing and is an error — never "until the next push". +func TestAGrantThatFailsSendsNothingAndSaysSo(t *testing.T) { + d := &recordingDelivery{grantErr: errors.New("the bus refused the membership")} + _, err := deliver(t.Context(), d, ready("anchor", "laptop")) + if err == nil || !strings.Contains(err.Error(), "the bus refused the membership") { + t.Fatalf("a failed grant was not said as the send's failure: %v", err) + } + for _, did := range d.did { + if strings.HasPrefix(did, "declare") { + t.Fatalf("code was sent after its grant failed: %v", d.did) + } + } + // Nothing to send is nothing granted either. + none := &recordingDelivery{grantErr: errors.New("never asked")} + if _, err := deliver(t.Context(), none, nil); err != nil || len(none.did) != 0 { + t.Fatalf("an empty send granted or failed: %v %v", none.did, err) + } +} + +// The machine holding the bus goes first: its declaration carries the user list the new grants are +// checked against. Among the machines it is moved to the front; not among them it is added only +// when it is behind. +func TestTheMachineHoldingTheBusIsSentFirst(t *testing.T) { + for _, c := range []struct { + what string + names []string + holder string + behind bool + want []string + }{ + {"among them", []string{"ace", "g14", "novox"}, "novox", false, []string{"novox", "ace", "g14"}}, + {"not among them, behind", []string{"ace", "g14"}, "novox", true, []string{"novox", "ace", "g14"}}, + {"not among them, current", []string{"ace", "g14"}, "novox", false, []string{"ace", "g14"}}, + {"nothing holds the bus", []string{"ace", "g14"}, "", true, []string{"ace", "g14"}}, + {"only it", []string{"novox"}, "novox", false, []string{"novox"}}, + } { + if got := brokerFirst(c.names, c.holder, c.behind); !reflect.DeepEqual(got, c.want) { + t.Errorf("%s: sent in the order %v, wanted %v", c.what, got, c.want) + } + } +} diff --git a/cmd/mesh-controller/hold_test.go b/cmd/mesh-controller/hold_test.go index 9c2b480..62660c7 100644 --- a/cmd/mesh-controller/hold_test.go +++ b/cmd/mesh-controller/hold_test.go @@ -18,8 +18,8 @@ func TestASendRoundGivesItsHoldBackOnEveryWayOut(t *testing.T) { plain := func(context.Context, string) (sendable, error) { return sendable{Resources: []map[string]any{{"id": "x"}}}, nil } - failing := func(readyNode, []byte) error { return errors.New("the broker went away") } - fine := func(readyNode, []byte) error { return nil } + failing := &recordingDelivery{declareErr: errors.New("the broker went away")} + fine := &recordingDelivery{} for name, round := range map[string]func() error{ "a body that cannot be marshalled": func() error { diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 380f627..65ec086 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -358,6 +358,17 @@ func pushCommand(ctx context.Context, args []string) error { asked = append(asked, n.Name) } + // **The machine holding the bus first** (novox/hq issue 249): its declaration carries the bus's + // user list, and a module's new grants are refused by the bus until that list says them. Among + // the machines asked it goes first; not among them and behind, it is added — a named push whose + // module gained a state would otherwise send the code and leave the right to use it for the + // cascade below, after. + holder, holderBehind, err := brokerBehind(ctx, open, asked) + if err != nil { + return err + } + asked = brokerFirst(asked, holder, holderBehind) + // Held from composing to sending, so a converge on one of them cannot send between the two // and be overtaken by what was composed before it (novox/hq ADR 0100). held, release, err := holdNodes(ctx, open, asked) @@ -387,35 +398,17 @@ func pushCommand(ctx context.Context, args []string) error { return declared, err }) - sentDigest := map[string]string{} defer release() - for _, s := range sending { - // The number is inside the signed bytes, so a replayed older declaration cannot borrow a - // newer one's (novox/hq 04-ISSUES/107); it was taken when the composition began (issue 204). - body, err := s.declared.Body() - if err != nil { - return err - } - if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil { - return err - } - // After it is away, not before. A digest recorded for something that failed to send would - // make the machine look current for a declaration it never received. - digest, err := recordSent(ctx, inv, s.node, body) - if err != nil { - return err - } - sentDigest[s.node] = digest - fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) + // Each machine's memberships first, then the declarations (novox/hq issue 249, ADR 0160): a push + // is the one most operators run, and on 2026-10-01 it was the one path that issued none. + bus := overTheBus{open: open, server: server, signer: ident} + sentDigest, err := deliver(ctx, bus, sending) + if err != nil { + return err } release() fmt.Printf("\n%d node(s) told\n", len(sending)) reportUnheldPushed(os.Stdout, len(args) == 1, asked, unheld) - // And each machine's memberships, as every other send does (ADR 0160): a push is the one most - // operators run, and on 2026-10-01 it was the one path that issued none. - if err := issueMemberships(ctx, open, server, sending); err != nil { - return err - } // **A named push leaves the mesh consistent, not just the machine it named** (novox/hq // issue 057, ADR 0083). Assigning a cross-node consumer mints a provision, and the PROVIDER's @@ -484,17 +477,7 @@ func pushCommand(ctx context.Context, args []string) error { } return declared, err }, - func(s readyNode, body []byte) error { - if err := link.Declare(ctx, server.Bus(), ident, s.node, body, - 15*time.Second); err != nil { - return err - } - if _, err := recordSent(ctx, inv, s.node, body); err != nil { - return err - } - fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) - return nil - }) + bus) refusals = append(refusals, refused...) if err != nil { return err @@ -621,10 +604,10 @@ func composeEach(names []string, allot func(node string) (int64, error), // sendRound holds the named nodes, composes each and sends each that composed, and gives the hold // back on every way out — a body that cannot be marshalled and a send that fails included // (novox/hq ADR 0100). A node that cannot be composed is a refusal, not an error: the others are -// still sent. +// still sent. Their memberships go before their declarations, as every send's do (issue 249). func sendRound(ctx context.Context, open *stores, names []string, compose func(held context.Context, node string) (sendable, error), - send func(s readyNode, body []byte) error) ([]string, error) { + d delivery) ([]string, error) { held, release, err := holdNodes(ctx, open, names) if err != nil { return nil, err @@ -633,18 +616,162 @@ func sendRound(ctx context.Context, open *stores, names []string, sending, refused := composeEach(names, allotting(held, open.inventory), func(node string) (sendable, error) { return compose(held, node) }) - for _, s := range sending { - body, err := s.declared.Body() - if err != nil { - return refused, err - } - if err := send(s, body); err != nil { - return refused, err - } + if _, err := deliver(held, d, sending); err != nil { + return refused, err } return refused, nil } +// delivery is the two acts of sending machines what they should be, apart, so the order between +// them is one function's and can be read and tested there (novox/hq issue 249). +type delivery interface { + // grant issues what the machines' modules may do — each module's state raised and its + // membership issued — for every machine about to be sent. + grant(ctx context.Context, sending []readyNode) error + // declare sends one machine its declaration and records it sent, answering the digest. + declare(ctx context.Context, s readyNode, body []byte) (string, error) +} + +// deliver sends the machines their declarations, **the grants first** (novox/hq issue 249). +// +// A merge gave a module a new state; its bundle reached every machine within a minute, and the +// machines' permissions on the bus did not include the state until somebody pushed by hand: the code +// arrived before the right to use it. A module that read its new state on start failed its start; the +// one that was there retried for two minutes. The memberships were issued after the declarations — +// "because the runtime it is for arrives with it" — and a membership is retained last-per-subject on +// the bus (internal/link/bus.go), so issued first it waits for the runtime that arrives after it. A +// runtime still on the old code merely holds a grant it does not use yet. +// +// **A grant that cannot be issued sends nothing**, and says so as an error. It used to be said and +// passed over — "the machines keep what they derive until the next push" — which reported a rollout +// as done that had delivered code its machines could not run, and left the remedy to a push nobody +// knew to make. Returned, the caller's rollout is not marked sent and is tried again. +func deliver(ctx context.Context, d delivery, sending []readyNode) (map[string]string, error) { + digests := map[string]string{} + if len(sending) == 0 { + return digests, nil + } + if err := d.grant(ctx, sending); err != nil { + return digests, fmt.Errorf("nothing was sent: what the machines' modules may do on the bus "+ + "could not be issued, and code sent before its grants is refused there (novox/hq issue 249): %w", err) + } + for _, s := range sending { + // The number is inside the signed bytes, so a replayed older declaration cannot borrow a + // newer one's (novox/hq 04-ISSUES/107); it was taken when the composition began (issue 204). + body, err := s.declared.Body() + if err != nil { + return digests, err + } + digest, err := d.declare(ctx, s, body) + if err != nil { + return digests, err + } + digests[s.node] = digest + } + return digests, nil +} + +// overTheBus is delivery as the mesh does it: memberships on the bus, declarations signed. +type overTheBus struct { + open *stores + server *link.Server + signer link.Signer + // indent is put before each "sent" line, for the callers whose output is nested. + indent string +} + +func (b overTheBus) grant(ctx context.Context, sending []readyNode) error { + return issueMemberships(ctx, b.open, b.server, sending) +} + +func (b overTheBus) declare(ctx context.Context, s readyNode, body []byte) (string, error) { + if err := link.Declare(ctx, b.server.Bus(), b.signer, s.node, body, 15*time.Second); err != nil { + return "", err + } + // After it is away, not before. A digest recorded for something that failed to send would make + // the machine look current for a declaration it never received. + digest, err := recordSent(ctx, b.open.inventory, s.node, body) + if err != nil { + return "", err + } + fmt.Printf("%ssent %s %d resource(s)\n", b.indent, s.node, len(s.declared.Resources)) + return digest, nil +} + +// brokerFirst is the machines to send in the order a grant needs (novox/hq issue 249): the machine +// holding the bus first — its declaration carries the bus's user list (composeBusUsers), and a +// module's new permissions are refused by the bus until that list says them. Among the machines it +// is moved to the front; not among them, it is added only when it is behind. +func brokerFirst(names []string, holder string, behind bool) []string { + if holder == "" { + return names + } + present := false + rest := make([]string, 0, len(names)) + for _, n := range names { + if n == holder { + present = true + continue + } + rest = append(rest, n) + } + if !present && !behind { + return names + } + return append([]string{holder}, rest...) +} + +// brokerBehind is the machine holding the bus — the one whose declaration carries the user list — +// and, when it is not among the machines named, whether what it should be differs from what it was +// last sent. +// +// **Its whole declaration, because that is what is recorded** (novox/hq issue 249). The mesh keeps +// a digest of what each machine was last sent and not of the user list inside it (ADR 0043: the +// list is composed on each push, never kept), so "the user list changed" is read as "the bus's +// machine is behind". That is the safe direction: a machine sent what it should be is never wrong, +// and whatever else it was behind on is what any push would have sent it. +func brokerBehind(ctx context.Context, open *stores, names []string) (string, bool, error) { + inv := open.inventory + holders, err := seatHolders(ctx, inv) + if err != nil { + return "", false, err + } + h, held := holders[theBrokerSeat] + if !held || h.Node == "" { + return "", false, nil + } + shelf, err := inv.Catalogue(ctx) + if err != nil { + return "", false, err + } + if m, known := shelf[h.Module]; !known || m.BusUsers == "" { + // A holder that is sent no user list carries no grant: nothing to send first. + return "", false, nil + } + for _, n := range names { + if n == h.Node { + return h.Node, false, nil + } + } + node, err := inv.NodeByName(ctx, h.Node) + if err != nil { + return "", false, err + } + would, err := wouldSend(ctx, open, []inventory.Node{node}) + if err != nil { + return "", false, err + } + if would[h.Node] == "" { + // It cannot be worked out: sending it would refuse the whole send, and `plan` says why. + return h.Node, false, nil + } + sent, err := inv.Outstanding(ctx, h.Node) + if err != nil { + return "", false, err + } + return h.Node, would[h.Node] != sent, nil +} + // couldNotBeResolved is what a push ends with when some machines could not be worked out. // // **After the rest have been sent, never instead of sending them.** It is still an error, because @@ -667,23 +794,36 @@ func couldNotBeResolved(refusals []string, sent int) error { // consumer and refused on the provider would leave one end holding a credential the other has // never heard of — which is the state this whole mechanism exists to make impossible. func sendTo(ctx context.Context, open *stores, names []string) error { + _, err := sendToEach(ctx, open, names) + return err +} + +// sendToEach is sendTo, answering the machines it sent: those named, and before them the machine +// holding the bus when its user list must go first (novox/hq issue 249) — so a caller that waits for +// the machines it sent waits for that one too. +func sendToEach(ctx context.Context, open *stores, names []string) ([]string, error) { inv := open.inventory ident, err := openIdentity(ctx) if err != nil { - return err + return nil, err } defer ident.Close() gens, err := generators(ctx, open) if err != nil { - return err + return nil, err } + holder, behind, err := brokerBehind(ctx, open, names) + if err != nil { + return nil, err + } + names = brokerFirst(names, holder, behind) // Held from composing to sending (novox/hq ADR 0100); a caller that holds them already — // converge, which flips the node and then sends it — is not made to wait on itself. ctx, release, err := holdNodes(ctx, open, names) if err != nil { - return err + return nil, err } defer release() @@ -712,33 +852,29 @@ func sendTo(ctx context.Context, open *stores, names []string) error { sending = append(sending, readyNode{name, declared}) } if len(refusals) > 0 { - return fmt.Errorf("nothing was sent. %d machine(s) could not be resolved:\n\n%s", + return nil, fmt.Errorf("nothing was sent. %d machine(s) could not be resolved:\n\n%s", len(refusals), strings.Join(refusals, "\n\n")) } server, err := connectLink(ctx, nil, nil, nil) if err != nil { - return err + return nil, err } defer server.Close() - for _, s := range sending { - body, err := s.declared.Body() - if err != nil { - return err - } - if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil { - return err - } - if _, err := recordSent(ctx, inv, s.node, body); err != nil { - return err - } - fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) - } // And every assignment on those machines its membership (novox/hq ADR 0160): composed from the // same records the bus's accounts are, so what a runtime serves and what its account may are one - // composition. Issued after the declaration, because the runtime it is for arrives with it. - return issueMemberships(ctx, open, server, sending) + // composition. **Issued before the declarations** (novox/hq issue 249): the runtime the + // membership is for arrives with the declaration, and a membership waits for it on the bus; the + // code arriving first was refused its own state until somebody pushed. + if _, err := deliver(ctx, overTheBus{open: open, server: server, signer: ident, indent: " "}, sending); err != nil { + return nil, err + } + sent := make([]string, 0, len(sending)) + for _, s := range sending { + sent = append(sent, s.node) + } + return sent, nil } // issueMemberships publishes the membership of every module on the machines just sent. @@ -759,16 +895,20 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se // **Every declared state's bucket, before the memberships that name it** (novox/hq ADR 0201). The // raise at start asserts them too, but a module registered and assigned since would otherwise have // its bucket only after the control plane next restarts — found the first time a module declared - // state: its bundle asked for a bucket that did not exist. Idempotent and cheap; a failure is said - // and the push stands, as a membership's is. - if buckets, err := open.inventory.DeclaredBuckets(ctx); err != nil { - fmt.Printf(" the modules' state could not be read, so no bucket was asserted: %v\n", err) - } else if _, err := broker.RaiseBuckets(broker.OnConn(bus.Conn), buckets); err != nil { - fmt.Printf(" the modules' state could not be asserted on the bus: %v — the next push tries again\n", err) + // state: its bundle asked for a bucket that did not exist. Idempotent and cheap. + // + // **A failure here is the send's failure** (novox/hq issue 249). It was said and the push stood, + // because the declarations were already away; they are sent after this now, and a module whose + // state does not exist is a module that fails its start, so nothing is sent and the caller tries + // again rather than reporting the rollout done. + buckets, err := open.inventory.DeclaredBuckets(ctx) + if err != nil { + return fmt.Errorf("the modules' state could not be read, so no bucket was asserted: %w", err) } - // The declarations are sent and recorded by now; a membership that cannot be issued is said - // and does not unsay them. Every runtime without one serves the shape it derives (ADR 0160), so - // the push stands, the first failure is named once, and the next push tries again. + if _, err := broker.RaiseBuckets(broker.OnConn(bus.Conn), buckets); err != nil { + return fmt.Errorf("the modules' state could not be asserted on the bus: %w", err) + } + // Every membership is tried, and the first failure named once. issued, failed := 0, 0 var first error for _, s := range sent { @@ -804,8 +944,9 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se fmt.Printf(" issued %d membership(s)\n", issued) } if failed > 0 { - fmt.Printf(" %d membership(s) could not be issued; the first: %v — the machines keep what "+ - "they derive until the next push\n", failed, first) + // Returned, never passed over (novox/hq issue 249): the declarations that need these are not + // sent, and the rollout that asked is tried again rather than waiting for a push. + return fmt.Errorf("%d membership(s) could not be issued; the first: %w", failed, first) } return nil } From 4ac5cfe3a7942ecb9c3cb01c2b55330be38bdea6 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 17:56:25 +0200 Subject: [PATCH 3/7] Roll a plan's module out to one machine first (hq issue 249, ADR 0218) A plan sent every machine running a module at once, ignoring the module's upgrade policy. Unless the policy says together, the first machine by name is sent, recorded in the plan, and the rest follow only once its report after the send says it applied; a failed first machine stops the plan. --- cmd/mesh-controller/release_plan.go | 145 +++++++++++++++++++++- cmd/mesh-controller/rollout_first_test.go | 84 +++++++++++++ internal/inventory/plans.go | 11 +- 3 files changed, 233 insertions(+), 7 deletions(-) create mode 100644 cmd/mesh-controller/rollout_first_test.go diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 267d2c7..3efb023 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -470,6 +470,15 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, // a module that packages another repository's source, keeps its commit; the catalogue announces // no move for it and its machines would keep the old image until somebody pushed (novox/hq // issue 189). A module whose policy records is built and left, as its policy says. + // + // **One machine first, unless the module's policy says together** (novox/hq issue 249, ADR + // 0218). The plan sent every machine running the module at once, and the operator's policy — + // one at a time, stopping at the first that fails, which an announced upgrade honours — was not + // read here at all: a module whose new declarations broke it broke everywhere in the same + // minute. Now the first machine is sent, the plan records it and waits for that machine's report + // after the send to say it applied what it was sent; only then are the rest sent. A first machine + // that fails or refuses stops the module's rollout and the plan with it, the rest untouched. + var pending []string for _, m := range tier { state := p.Modules[m] if state == nil || state.SentAt != nil || !rollsOut(m) { @@ -479,17 +488,67 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, if err != nil { return false, err } + policy, err := inv.UpgradeOf(ctx, m) + if err != nil { + return false, err + } + var reports []inventory.Reported + if state.FirstAt != nil && !policy.Together { + if reports, err = inv.LastReports(ctx); err != nil { + return false, err + } + } + step := nextRollout(*state, running, policy.Together, reports) now := time.Now().UTC() - state.SentAt = &now - if len(running) == 0 { + switch { + case step.failed != "": + state.Why = step.failed + p.State = inventory.PlanFailed + p.Note = fmt.Sprintf("%s stopped at its first machine in tier %d: %s; %s left as it was", + m, p.Tier, step.failed, orNone(strings.Join(step.rest, ", "))) + fmt.Printf("%s: %s\n", p.ID, p.Note) + return true, nil + case step.waiting != "": + wait := fmt.Sprintf("%s on %s, sent first at %s", m, step.waiting, state.FirstAt.Format("15:04")) + if now.Sub(*state.FirstAt) > planWaitBound { + wait += " — LATE" + } + pending = append(pending, wait) + continue + case len(step.send) == 0: + // No machine runs it: nothing to send, and nothing to wait for. + state.SentAt = &now continue } - if err := sendTo(ctx, open, running); err != nil { - return false, fmt.Errorf("sending %s to %s after tier %d: %w", m, strings.Join(running, ", "), p.Tier, err) + // What sendToEach answers, not what was asked: the machine holding the bus is sent before + // the first when its user list must change (issue 249), and the plan waits for it too. + sent, err := sendToEach(ctx, open, step.send) + if err != nil { + // Not marked sent, so the next step tries again (issue 249): a grant that could not be + // issued is a send that did not happen. + return false, fmt.Errorf("sending %s to %s after tier %d: %w", m, strings.Join(step.send, ", "), p.Tier, err) } - fmt.Printf("%s: tier %d built; sent %s to %s\n", p.ID, p.Tier, m, strings.Join(running, ", ")) + if step.first { + state.First = sent + state.FirstAt = &now + p.State = inventory.PlanRolling + p.Note = fmt.Sprintf("tier %d built; sent %s to %s first", p.Tier, m, strings.Join(sent, ", ")) + fmt.Printf("%s: tier %d built; sent %s to %s first, the rest once it reports it applied\n", + p.ID, p.Tier, m, strings.Join(sent, ", ")) + return true, nil + } + state.SentAt = &now + fmt.Printf("%s: tier %d built; sent %s to %s\n", p.ID, p.Tier, m, strings.Join(sent, ", ")) return true, nil } + if len(pending) > 0 { + note := "tier " + fmt.Sprint(p.Tier) + " built; waiting for " + strings.Join(pending, "; ") + + " to report it applied before the rest are sent" + changed := p.State != inventory.PlanRolling || p.Note != note + p.State = inventory.PlanRolling + p.Note = note + return changed, nil + } // And wait for what the next tier needs running. needed := gates(*p, edges, rollsOut) if len(needed) > 0 { @@ -540,6 +599,77 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, return true, nil } +// rolloutStep is what a plan does next with one built module's machines (novox/hq issue 249). +type rolloutStep struct { + // send is the machines to send now; first, whether they are the first machine's send. + send []string + first bool + // waiting names the first machines whose report the rest wait for. + waiting string + // failed says how a first machine did not take it; rest is what is then left alone. + failed string + rest []string +} + +// nextRollout is the next step of one module's rollout in a plan (novox/hq issue 249, ADR 0218). +// +// Together, every machine running it at once, as the policy says. Otherwise one machine first — +// the first by name, so the choice is the same on every controller and every resume — and the rest +// once each machine the first send reached has reported, after that send, that it applied the +// declaration it was last sent. A report after the send that failed or refused it is the rollout's +// end: the rest are not sent. A machine that has not reported since is waited for; how long it has +// been is the plan's to say. +func nextRollout(s inventory.PlanModule, running []string, together bool, reports []inventory.Reported) rolloutStep { + if len(running) == 0 { + return rolloutStep{} + } + if together { + return rolloutStep{send: running} + } + if s.FirstAt == nil { + sorted := append([]string{}, running...) + sort.Strings(sorted) + return rolloutStep{send: sorted[:1], first: true} + } + sentFirst := map[string]bool{} + for _, n := range s.First { + sentFirst[n] = true + } + var rest []string + for _, n := range running { + if !sentFirst[n] { + rest = append(rest, n) + } + } + byNode := map[string]inventory.Reported{} + for _, r := range reports { + byNode[r.Node] = r + } + var waiting, failed []string + for _, n := range s.First { + r, said := byNode[n] + // Only a report after the send, about what it was last sent, says anything about this build. + if !said || r.At == nil || r.At.Before(*s.FirstAt) || !r.Current { + waiting = append(waiting, n) + continue + } + switch r.Outcome { + case inventory.OutcomeApplied: + case inventory.OutcomeFailed, inventory.OutcomeRefused: + failed = append(failed, n+" "+r.Outcome+" what it was sent") + default: + waiting = append(waiting, n) + } + } + if len(failed) > 0 { + return rolloutStep{failed: strings.Join(failed, "; "), rest: rest} + } + if len(waiting) > 0 { + return rolloutStep{waiting: strings.Join(waiting, ", ")} + } + return rolloutStep{send: rest} +} + // planTicker advances open plans on a timer, for the steps outcomes alone cannot take. func planTicker(ctx context.Context, open *stores) { advancePlans(ctx, open) @@ -796,6 +926,11 @@ func planWhatIf(ctx context.Context, inv *inventory.Inventory, repository string if u, err := inv.UpgradeOf(ctx, name); err == nil && u.RollOut { running, _ := inv.Running(ctx, name) how = "built, then sent to " + orNone(strings.Join(running, ", ")) + // One machine first unless the policy says together (novox/hq issue 249). + if first := nextRollout(inventory.PlanModule{}, running, u.Together, nil); first.first && len(running) > 1 { + how = fmt.Sprintf("built, then sent to %s first and to the rest once it has applied it", + first.send[0]) + } rolls[name] = how } fmt.Printf(" %-22s %s\n", name, how) diff --git a/cmd/mesh-controller/rollout_first_test.go b/cmd/mesh-controller/rollout_first_test.go new file mode 100644 index 0000000..15587bb --- /dev/null +++ b/cmd/mesh-controller/rollout_first_test.go @@ -0,0 +1,84 @@ +package main + +import ( + "reflect" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/inventory" +) + +// novox/hq issue 249, ADR 0218: a plan rolls a module out to one machine first and the rest only +// once that machine has reported it applied; a module whose policy says together goes everywhere at +// once, as before. +func TestAPlanSendsOneMachineFirstAndTheRestAfterItsReport(t *testing.T) { + running := []string{"novox", "ace", "g14"} + + // Together: every machine at once. + if step := nextRollout(inventory.PlanModule{}, running, true, nil); !reflect.DeepEqual(step.send, running) || step.first { + t.Fatalf("a together policy did not send every machine at once: %+v", step) + } + + // Otherwise the first by name, alone. + step := nextRollout(inventory.PlanModule{}, running, false, nil) + if !step.first || !reflect.DeepEqual(step.send, []string{"ace"}) { + t.Fatalf("the first send was %+v, wanted ace alone", step) + } + + sentAt := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + before, after := sentAt.Add(-time.Minute), sentAt.Add(time.Minute) + state := inventory.PlanModule{First: []string{"ace"}, FirstAt: &sentAt} + report := func(at time.Time, outcome string, current bool) []inventory.Reported { + return []inventory.Reported{{Node: "ace", At: &at, Outcome: outcome, Current: current}, + {Node: "g14", At: &after, Outcome: inventory.OutcomeApplied, Current: true}} + } + + // No report yet, a report from before the send, or one about an older declaration: wait. + for what, reports := range map[string][]inventory.Reported{ + "no report": nil, + "a report before the send": report(before, inventory.OutcomeApplied, true), + "a report about older": report(after, inventory.OutcomeApplied, false), + } { + step := nextRollout(state, running, false, reports) + if len(step.send) != 0 || step.waiting != "ace" || step.failed != "" { + t.Errorf("%s: %+v, wanted to wait for ace", what, step) + } + } + + // Applied after the send: the rest, and only the rest. + step = nextRollout(state, running, false, report(after, inventory.OutcomeApplied, true)) + if step.first || !reflect.DeepEqual(step.send, []string{"novox", "g14"}) { + t.Fatalf("after ace applied it the plan sent %+v, wanted novox and g14", step) + } + + // Failed or refused after the send: stop, the rest untouched. + for _, outcome := range []string{inventory.OutcomeFailed, inventory.OutcomeRefused} { + step := nextRollout(state, running, false, report(after, outcome, true)) + if len(step.send) != 0 || !strings.Contains(step.failed, "ace "+outcome) || + !reflect.DeepEqual(step.rest, []string{"novox", "g14"}) { + t.Errorf("a first machine that %s it: %+v", outcome, step) + } + } + + // The machine holding the bus went with the first send: the rest wait for it too, and it is + // not sent again. + both := inventory.PlanModule{First: []string{"novox", "ace"}, FirstAt: &sentAt} + half := report(after, inventory.OutcomeApplied, true) + if step := nextRollout(both, running, false, half); step.waiting != "novox" { + t.Fatalf("the plan did not wait for the bus's machine sent first: %+v", step) + } + all := append(half, inventory.Reported{Node: "novox", At: &after, Outcome: inventory.OutcomeApplied, Current: true}) + if step := nextRollout(both, running, false, all); !reflect.DeepEqual(step.send, []string{"g14"}) { + t.Fatalf("after both applied it the plan sent %+v, wanted g14 alone", step) + } + + // One machine, or none: nothing is waited for that cannot come. + if step := nextRollout(inventory.PlanModule{}, nil, false, nil); len(step.send) != 0 || step.first { + t.Fatalf("a module nothing runs was sent: %+v", step) + } + only := inventory.PlanModule{First: []string{"ace"}, FirstAt: &sentAt} + if step := nextRollout(only, []string{"ace"}, false, report(after, inventory.OutcomeApplied, true)); len(step.send) != 0 || step.waiting != "" { + t.Fatalf("a module on one machine waited for more: %+v", step) + } +} diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index c68bafe..170ef31 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -37,8 +37,15 @@ type PlanModule struct { // later tier is built by it (ADR 0163's gate): the reports that open the gate are the ones // after this. SentAt *time.Time `json:"sent_at,omitempty"` - Commit string `json:"commit,omitempty"` - Why string `json:"why,omitempty"` + // First is the machines the plan sent the new build to first, and FirstAt when (novox/hq issue + // 249, ADR 0218): unless the module's policy rolls it out together, one machine takes it before + // the rest, and the rest are sent once that one reports it applied. Kept so a controller + // replaced while the plan waits on that report resumes the wait rather than sending again. The + // machine holding the bus is among them when its user list had to go first. + First []string `json:"first,omitempty"` + FirstAt *time.Time `json:"first_at,omitempty"` + Commit string `json:"commit,omitempty"` + Why string `json:"why,omitempty"` } // The states a plan passes through. From 208901c6cc0a335fca4e019fdfed853f50774b16 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 18:00:27 +0200 Subject: [PATCH 4/7] Let a newer plan supersede the older open plans of its repository (hq issue 254, ADR 0218) A merge planned without looking at open plans, so two plans worked the same modules and a stuck plan stayed open for ever. The newer plan folds in what older plans of the same repository and branch had not built or sent, and closes them as superseded. A person can close a stuck plan by id with `plans close `. --- cmd/mesh-controller/release_plan.go | 87 +++++++++++++---- cmd/mesh-controller/seatverbs.go | 3 + cmd/mesh-controller/supersede_test.go | 95 +++++++++++++++++++ cmd/mesh-controller/upgrades.go | 47 +++++++++ internal/catalogue/verbs.go | 1 + .../0057-a-newer-plan-supersedes-an-older.sql | 14 +++ internal/inventory/plans.go | 38 ++++---- internal/inventory/plans_test.go | 27 ++++++ 8 files changed, 279 insertions(+), 33 deletions(-) create mode 100644 cmd/mesh-controller/supersede_test.go create mode 100644 internal/inventory/migrations/0057-a-newer-plan-supersedes-an-older.sql diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 3efb023..74455b1 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -184,6 +184,7 @@ func planOfMerge(m link.SourceMoved, moved []string, edges []inventory.Edge) inv return inventory.Plan{ ID: fmt.Sprintf("plan-%d", time.Now().UnixNano()), Repository: m.Owner + "/" + m.Repo, + Branch: m.Base, Commit: m.Commit, Created: time.Now().UTC(), State: inventory.PlanBuilding, @@ -192,6 +193,57 @@ func planOfMerge(m link.SourceMoved, moved []string, edges []inventory.Edge) inv } } +// supersededBy is what a newer plan takes over from the open plans it supersedes (novox/hq issue +// 254, ADR 0218): the modules they had not finished, and those plans closed as superseded. +// +// **A merge looked at no plan but its own.** Two merges of one repository a few minutes apart were +// two open plans asking for the same modules, each sending machines what it built; and a plan that +// would never move again — waiting on a report that could not come, at 97b1b2b — stayed open for +// ever beside the newer ones, read as work in progress by everyone who looked. The newer merge is the +// newer intent for that repository and branch, so its plan takes over: every open plan of the same +// repository and branch **created before it** — by the time the plans were made, never by comparing +// commits, which have no order of their own — gives up the modules it had not built, and those are +// planned again in the newer plan beside what the newer merge moved. +// +// "Not built" is a module not yet asked, or asked and not answered; **and a module built and not +// yet sent to its machines**, where its policy rolls it out: closed, the older plan would never send +// it, and the catalogue announces no move for a rebuild (issue 189), so the newer plan builds and +// sends it. A build the older plan asked still finishes and registers as any build does — ordered by +// when it was asked (issue 219), so the newer plan's ask, made later, is the one that stands. +// +// A plan with no branch recorded is from before branches were kept, and is superseded by the next +// plan of its repository: what it had not built is folded in, so nothing is lost by it. +func supersededBy(newer inventory.Plan, open []inventory.Plan, rollsOut func(string) bool) ([]string, []inventory.Plan) { + folded := map[string]bool{} + var closed []inventory.Plan + for _, old := range open { + if old.ID == newer.ID || !old.Open() || !strings.EqualFold(old.Repository, newer.Repository) || + (old.Branch != "" && old.Branch != newer.Branch) || !old.Created.Before(newer.Created) { + continue + } + var took []string + for name, s := range old.Modules { + if s == nil || s.State != "built" || (s.SentAt == nil && rollsOut(name)) { + folded[name] = true + took = append(took, name) + } + } + sort.Strings(took) + old.State = inventory.PlanSuperseded + old.Note = fmt.Sprintf("superseded at tier %d by %s (%s at %s)", old.Tier, newer.ID, newer.Repository, short(newer.Commit)) + if len(took) > 0 { + old.Note += "; " + strings.Join(took, ", ") + " planned there again" + } + closed = append(closed, old) + } + out := make([]string, 0, len(folded)) + for name := range folded { + out = append(out, name) + } + sort.Strings(out) + return out, closed +} + // gates is what the next tier needs running from this one: a module of the tier that a later // tier is built by — the runtime dependency — and whose policy rolls it out, must be applied by // the machines running it before the next tier is asked. A base an image stands on need only be @@ -693,6 +745,8 @@ func planLine(p inventory.Plan, now time.Time) string { return fmt.Sprintf("%s %s done, %d tier(s)", p.Repository, short(p.Commit), len(p.Tiers)) case inventory.PlanFailed: return fmt.Sprintf("%s %s FAILED at %s: %s", p.Repository, short(p.Commit), where, p.Note) + case inventory.PlanSuperseded: + return fmt.Sprintf("%s %s %s", p.Repository, short(p.Commit), p.Note) } since := now.Sub(p.Updated).Round(time.Minute) late := "" @@ -823,7 +877,19 @@ func plansCommand(ctx context.Context, args []string) error { if *whatIf != "" { return planWhatIf(ctx, inv, *whatIf, splitList(*paths), splitList(*modules)) } - if len(positionals) == 2 && positionals[0] == "stop" { + // `stop`, or `close` (novox/hq issue 254): a person ending a plan that will not move again — one + // waiting on a report that cannot come — so it stops reading as work in progress. Marked failed + // with who ended it; what it asked still builds and registers. + if len(positionals) == 2 && (positionals[0] == "stop" || positionals[0] == "close") { + how := "stopped" + if positionals[0] == "close" { + how = "closed" + } + release, err := inv.HoldPlans(ctx, true) + if err != nil { + return err + } + defer release() p, err := inv.PlanByID(ctx, positionals[1]) if err != nil { return err @@ -832,25 +898,12 @@ func plansCommand(ctx context.Context, args []string) error { return fmt.Errorf("%s is already %s", p.ID, p.State) } p.State = inventory.PlanFailed - p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier) - release, err := inv.HoldPlans(ctx, true) - if err != nil { - return err - } - defer release() - if p, err = inv.PlanByID(ctx, positionals[1]); err != nil { - return err - } - if !p.Open() { - return fmt.Errorf("%s is already %s", p.ID, p.State) - } - p.State = inventory.PlanFailed - p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier) + p.Note = how + " by hand at tier " + fmt.Sprint(p.Tier) if err := inv.SavePlan(ctx, p); err != nil { return err } - fmt.Printf("%s stopped at tier %d of %d; what was asked still builds and registers, nothing further is asked\n", - p.ID, p.Tier, len(p.Tiers)) + fmt.Printf("%s %s at tier %d of %d; what was asked still builds and registers, nothing further is asked\n", + p.ID, how, p.Tier, len(p.Tiers)) return nil } plans, err := inv.RecentPlans(ctx, *limit) diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 3591da5..d7dec42 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -100,6 +100,9 @@ func argvFor(verb string, args map[string]any) ([]string, error) { if id := str("stop"); id != "" { return []string{"plans", "stop", id}, nil } + if id := str("close"); id != "" { + return []string{"plans", "close", id}, nil + } if id := str("id"); id != "" { return []string{"plans", id}, nil } diff --git a/cmd/mesh-controller/supersede_test.go b/cmd/mesh-controller/supersede_test.go new file mode 100644 index 0000000..2c269fb --- /dev/null +++ b/cmd/mesh-controller/supersede_test.go @@ -0,0 +1,95 @@ +package main + +import ( + "reflect" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/inventory" +) + +// novox/hq issue 254, ADR 0218: a newer plan takes over what the older open plans of its repository +// and branch had not built, and closes them as superseded; another repository's plan, another +// branch's, and a plan made after it are left alone. +func TestANewerPlanSupersedesTheOlderOpenPlansOfItsRepository(t *testing.T) { + at := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + sent := at.Add(time.Minute) + plan := func(id, repository, branch string, created time.Time, modules map[string]*inventory.PlanModule) inventory.Plan { + return inventory.Plan{ID: id, Repository: repository, Branch: branch, Commit: id + "-commit", + Created: created, State: inventory.PlanRolling, Modules: modules} + } + older := plan("plan-1", "novox/mesh-catalog", "main", at, map[string]*inventory.PlanModule{ + "gitea": {State: "built", SentAt: &sent}, // done with: stays done + "keycloak": {State: "asked"}, // asked, not answered: folded + "plex": {}, // not yet asked: folded + "agent": {State: "built"}, // built, rolls out, not sent: folded + "notes": {State: "built"}, // built, records: nothing to send + }) + stuck := plan("plan-0", "Novox/Mesh-Catalog", "", at.Add(-time.Hour), map[string]*inventory.PlanModule{ + "runtime": {State: "asked"}, + }) + other := plan("plan-2", "novox/mesh-controller", "main", at, map[string]*inventory.PlanModule{"mesh-controller": {}}) + release := plan("plan-3", "novox/mesh-catalog", "release", at, map[string]*inventory.PlanModule{"lemurs": {}}) + later := plan("plan-5", "novox/mesh-catalog", "main", at.Add(2*time.Hour), map[string]*inventory.PlanModule{"later": {}}) + done := plan("plan-6", "novox/mesh-catalog", "main", at, map[string]*inventory.PlanModule{"finished": {}}) + done.State = inventory.PlanDone + + newer := plan("plan-4", "novox/mesh-catalog", "main", at.Add(time.Hour), nil) + newer.Commit = "97b1b2b0c0ffee" + rollsOut := func(m string) bool { return m != "notes" } + folded, closed := supersededBy(newer, []inventory.Plan{stuck, older, other, release, later, done, newer}, rollsOut) + + if want := []string{"agent", "keycloak", "plex", "runtime"}; !reflect.DeepEqual(folded, want) { + t.Fatalf("folded %v, wanted %v", folded, want) + } + var ids []string + for _, p := range closed { + ids = append(ids, p.ID) + if p.State != inventory.PlanSuperseded || p.Open() { + t.Errorf("%s was left %s", p.ID, p.State) + } + if !strings.Contains(p.Note, "plan-4") || !strings.Contains(p.Note, "97b1b2b0") { + t.Errorf("%s does not name the plan that superseded it: %q", p.ID, p.Note) + } + } + if want := []string{"plan-0", "plan-1"}; !reflect.DeepEqual(ids, want) { + t.Fatalf("superseded %v, wanted %v — another repository, another branch, a later plan and a "+ + "finished one are left alone", ids, want) + } + if other.State != inventory.PlanRolling { + t.Fatal("the plan handed in was changed in place") + } + if line := planLine(closed[1], time.Now()); !strings.Contains(line, "superseded") { + t.Fatalf("a superseded plan reads %q", line) + } +} + +// novox/hq issue 254: a person closes a plan that will not move again, by its id. +func TestAPersonClosesAStuckPlan(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + stuck := inventory.Plan{ID: "plan-97b1b2b", Repository: "novox/mesh-catalog", Commit: "97b1b2b", + Created: time.Now().UTC(), State: inventory.PlanRolling, Tier: 1, Tiers: [][]string{{"a"}, {"b"}}, + Modules: map[string]*inventory.PlanModule{"a": {State: "built"}, "b": {}}} + if err := open.inventory.SavePlan(ctx, stuck); err != nil { + t.Fatal(err) + } + if err := plansCommand(ctx, []string{"close", stuck.ID}); err != nil { + t.Fatal(err) + } + closed, err := open.inventory.PlanByID(ctx, stuck.ID) + if err != nil { + t.Fatal(err) + } + if closed.State != inventory.PlanFailed || !strings.Contains(closed.Note, "closed by hand") { + t.Fatalf("the plan was left %s: %q", closed.State, closed.Note) + } + if err := plansCommand(ctx, []string{"close", stuck.ID}); err == nil { + t.Fatal("a plan already closed was closed again") + } + if argv, err := argvFor("plans", map[string]any{"close": stuck.ID}); err != nil || + !reflect.DeepEqual(argv, []string{"plans", "close", stuck.ID}) { + t.Fatalf("the seat's verb does not close a plan: %v %v", argv, err) + } +} diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index f82ccb2..50389b6 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -324,6 +324,45 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error { } defer release() plan := planOfMerge(m, movedNames, edges) + // **A newer plan supersedes the older open plans of this repository and branch** (novox/hq issue + // 254, ADR 0218): what they had not built is planned here again, and they are closed, so one plan + // works a repository's modules at a time and a stuck one ends at the next merge. + working, err := inv.OpenPlans(ctx) + if err != nil { + return notNow(err) + } + rollsOut := func(module string) bool { + u, err := inv.UpgradeOf(ctx, module) + return err == nil && u.RollOut + } + folded, superseded := supersededBy(plan, working, rollsOut) + if len(folded) > 0 { + held := map[string]bool{} + for _, e := range entries { + held[e.Manifest.Module] = true + } + names := map[string]bool{} + for _, name := range movedNames { + names[name] = true + } + var also []string + for _, name := range folded { + // One the catalogue no longer holds would fail the newer plan's ask; it is not this + // merge's to build. + if held[name] && !names[name] { + names[name] = true + movedNames = append(movedNames, name) + also = append(also, name) + } + } + if len(also) > 0 { + again := planOfMerge(m, movedNames, edges) + again.ID, again.Created = plan.ID, plan.Created + plan = again + fmt.Printf(" %s, left unbuilt by an older plan of %s, are planned here again\n", + strings.Join(also, ", "), plan.Repository) + } + } if hasCycle(plan.Tiers, edges) { fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n", strings.Join(plan.Tiers[len(plan.Tiers)-1], ", ")) @@ -331,6 +370,14 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error { if err := inv.SavePlan(ctx, plan); err != nil { return notNow(err) } + // Closed after the newer plan is kept, never before: a controller replaced between the two leaves + // both open, which the next merge settles, rather than neither. + for _, old := range superseded { + if err := inv.SavePlan(ctx, old); err != nil { + return notNow(err) + } + fmt.Printf(" %s (%s at %s) is %s\n", old.ID, old.Repository, short(old.Commit), old.Note) + } var tiers []string for i, t := range plan.Tiers { tiers = append(tiers, fmt.Sprintf("%d: %s", i, strings.Join(t, ", "))) diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index f4a1352..8898585 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -96,6 +96,7 @@ var ControllerVerbs = []Verb{ Input: schema(map[string]string{ "id": "a plan's id (as `plans` lists them): that plan, tier by tier", "stop": "a plan's id: stop it — what was asked still builds, nothing further is asked", + "close": "a plan's id: close a plan that will not move again, as failed by hand (novox/hq issue 254)", "repository": "owner/repository: the plan a merge there would produce, saving nothing (what-if); with paths or modules", "paths": "with repository: the files the merge would change, comma-separated, from the repository's root", "modules": "with repository: or the modules it would change, comma-separated", diff --git a/internal/inventory/migrations/0057-a-newer-plan-supersedes-an-older.sql b/internal/inventory/migrations/0057-a-newer-plan-supersedes-an-older.sql new file mode 100644 index 0000000..8a83eb6 --- /dev/null +++ b/internal/inventory/migrations/0057-a-newer-plan-supersedes-an-older.sql @@ -0,0 +1,14 @@ +-- A newer plan supersedes the older open plans of the same repository and branch (novox/hq issue 254, +-- ADR 0218). +-- +-- A merge produced a plan without looking at the plans still open, so two merges a few minutes apart +-- were two plans working the same modules, and a plan stuck waiting on something that would never +-- come stayed open for ever beside the newer ones. The newer plan now takes over what the older had +-- not yet built and the older is closed as `superseded` — a state of its own, so `plans` can say +-- which plan replaced it rather than reading as a failure. +-- +-- `branch` is the branch the merge went into, so only a plan of the same branch is superseded. Empty +-- for every plan from before this was kept: which branch it answered is not known, and such a plan +-- is superseded by the next plan of its repository, whichever branch — nothing is lost by it, since +-- what it had not built is folded into the plan that supersedes it. +alter table release_plan add column branch text not null default ''; diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index 170ef31..cf90380 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -15,16 +15,19 @@ import ( // the store so a controller replaced mid-plan resumes it, and so `status` can say what a merge // still waits for. type Plan struct { - ID string `json:"id"` - Repository string `json:"repository"` - Commit string `json:"commit"` - Created time.Time `json:"created"` - Updated time.Time `json:"updated"` - State string `json:"state"` - Tier int `json:"tier"` - Tiers [][]string `json:"tiers"` - Modules map[string]*PlanModule `json:"modules"` - Note string `json:"note,omitempty"` + ID string `json:"id"` + Repository string `json:"repository"` + // Branch is the branch the merge went into (novox/hq issue 254): a newer plan supersedes the open + // ones of the same repository and branch. Empty for a plan from before it was kept. + Branch string `json:"branch,omitempty"` + Commit string `json:"commit"` + Created time.Time `json:"created"` + Updated time.Time `json:"updated"` + State string `json:"state"` + Tier int `json:"tier"` + Tiers [][]string `json:"tiers"` + Modules map[string]*PlanModule `json:"modules"` + Note string `json:"note,omitempty"` } // PlanModule is one module's state within a plan. @@ -54,6 +57,9 @@ const ( PlanRolling = "rolling" PlanDone = "done" PlanFailed = "failed" + // PlanSuperseded is a plan a newer merge of the same repository and branch took over (novox/hq + // issue 254, ADR 0218): what it had not built is in the newer plan, and its note names it. + PlanSuperseded = "superseded" ) // Open says whether the plan is still being worked. @@ -70,11 +76,11 @@ func (i *Inventory) SavePlan(ctx context.Context, p Plan) error { return err } _, err = i.store.Pool().Exec(ctx, - `insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note) - values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9) + `insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch) + values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9, $10) on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier, - tiers = excluded.tiers, modules = excluded.modules, note = excluded.note`, - p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note) + tiers = excluded.tiers, modules = excluded.modules, note = excluded.note, branch = excluded.branch`, + p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note, p.Branch) return err } @@ -102,7 +108,7 @@ func (i *Inventory) PlanByID(ctx context.Context, id string) (Plan, error) { func (i *Inventory) plans(ctx context.Context, tail string) ([]Plan, error) { rows, err := i.store.Pool().Query(ctx, - `select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note + `select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note, branch from release_plan `+tail) if err != nil { return nil, err @@ -113,7 +119,7 @@ func (i *Inventory) plans(ctx context.Context, tail string) ([]Plan, error) { var p Plan var tiers, modules []byte if err := rows.Scan(&p.ID, &p.Repository, &p.Commit, &p.Created, &p.Updated, &p.State, - &p.Tier, &tiers, &modules, &p.Note); err != nil { + &p.Tier, &tiers, &modules, &p.Note, &p.Branch); err != nil { return nil, err } if err := json.Unmarshal(tiers, &p.Tiers); err != nil { diff --git a/internal/inventory/plans_test.go b/internal/inventory/plans_test.go index a9153a6..e4804a7 100644 --- a/internal/inventory/plans_test.go +++ b/internal/inventory/plans_test.go @@ -46,3 +46,30 @@ func TestAPlanIsKeptAdvancedAndResumedFromTheStore(t *testing.T) { t.Fatalf("a done plan is still among the recent ones: %+v", recent) } } + +// novox/hq issue 254: a plan keeps the branch its merge went into, and a superseded plan is not open. +func TestASupersededPlanIsNotOpen(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + p := Plan{ID: "plan-1", Repository: "novox/mesh-catalog", Branch: "main", Commit: "abc", + Created: time.Now().UTC(), State: PlanBuilding, Tiers: [][]string{{"gitea"}}, + Modules: map[string]*PlanModule{"gitea": {}}} + if err := inv.SavePlan(ctx, p); err != nil { + t.Fatal(err) + } + kept, err := inv.PlanByID(ctx, "plan-1") + if err != nil || kept.Branch != "main" { + t.Fatalf("the branch was not kept: %v %+v", err, kept) + } + kept.State = PlanSuperseded + kept.Note = "superseded at tier 0 by plan-2" + if err := inv.SavePlan(ctx, kept); err != nil { + t.Fatal(err) + } + if open, err := inv.OpenPlans(ctx); err != nil || len(open) != 0 { + t.Fatalf("a superseded plan is still open: %v %+v", err, open) + } + if recent, _ := inv.RecentPlans(ctx, 5); len(recent) != 1 || recent[0].State != PlanSuperseded { + t.Fatalf("a superseded plan is not among the recent ones as superseded: %+v", recent) + } +} From 672d1f4ca104b1185ff0babc999b4e1535b67014 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 18:01:11 +0200 Subject: [PATCH 5/7] Leave an announced move to the plan rolling the module out (hq issue 249) The catalogue's upgrade announcement sent every machine one after another without waiting for any to apply, beside the plan that now sends one machine first. A module an open plan has not finished sending is the plan's to roll out. --- cmd/mesh-controller/rollout_first_test.go | 19 +++++++++++++++++++ cmd/mesh-controller/upgrades.go | 23 +++++++++++++++++++++++ 2 files changed, 42 insertions(+) diff --git a/cmd/mesh-controller/rollout_first_test.go b/cmd/mesh-controller/rollout_first_test.go index 15587bb..cb8ce79 100644 --- a/cmd/mesh-controller/rollout_first_test.go +++ b/cmd/mesh-controller/rollout_first_test.go @@ -82,3 +82,22 @@ func TestAPlanSendsOneMachineFirstAndTheRestAfterItsReport(t *testing.T) { t.Fatalf("a module on one machine waited for more: %+v", step) } } + +// An announced move of a module an open plan is still rolling out is left to the plan: sending it +// here as well put the bundle on every machine at once (novox/hq issue 249). +func TestAnAnnouncedMoveIsLeftToThePlanRollingItOut(t *testing.T) { + sent := time.Now() + plans := []inventory.Plan{ + {ID: "plan-done", State: inventory.PlanDone, Modules: map[string]*inventory.PlanModule{"agent": {}}}, + {ID: "plan-1", State: inventory.PlanRolling, Modules: map[string]*inventory.PlanModule{ + "agent": {State: "built", First: []string{"ace"}, FirstAt: &sent}, "gitea": {State: "built", SentAt: &sent}}}, + } + if got := rolledOutByAPlan(plans, "agent"); got != "plan-1" { + t.Fatalf("a module the plan is rolling out was not left to it: %q", got) + } + for _, m := range []string{"gitea", "keycloak"} { + if got := rolledOutByAPlan(plans, m); got != "" { + t.Errorf("%s, which no plan will send, was left to %s", m, got) + } + } +} diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index 50389b6..3c71189 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -56,6 +56,18 @@ func (f following) Upgraded(ctx context.Context, u link.Upgraded) error { return nil } + // **A plan that holds the module rolls it out, and this does not** (novox/hq issue 249, ADR + // 0218). A merge's plan builds the module and sends it one machine first, the rest once that one + // has applied it; this announcement arrives as the build registers, and sending here too — one + // machine after another without waiting for any to apply — put the new bundle on every machine in + // the same minute, whatever the plan was waiting for. A move no plan answers (a build asked by + // hand) is still this handler's. + if plans, err := inv.OpenPlans(ctx); err != nil { + return notNow(err) + } else if id := rolledOutByAPlan(plans, u.Module); id != "" { + fmt.Printf("%s moved to %s; %s rolls it out to %s\n", u.Module, shortCommit(u.Commit), id, readableList(on)) + return nil + } if decision.Together { fmt.Printf("%s moved to %s; sending %s together\n", u.Module, shortCommit(u.Commit), readableList(on)) @@ -77,6 +89,17 @@ func (f following) Upgraded(ctx context.Context, u link.Upgraded) error { return nil } +// rolledOutByAPlan is the open plan that will send a module's machines its new build — one holding +// the module that has not finished sending it — or empty when none will (novox/hq issue 249). +func rolledOutByAPlan(plans []inventory.Plan, module string) string { + for _, p := range plans { + if s, holds := p.Modules[module]; p.Open() && holds && (s == nil || s.SentAt == nil) { + return p.ID + } + } + return "" +} + // readableList names machines the way a sentence does, because this is read by a person deciding // whether an upgrade went where they expected. func readableList(names []string) string { From ed907713828d7beb5b5fed3e2ca01211394861d4 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 18:17:52 +0200 Subject: [PATCH 6/7] Send the bus's machine before the grants, and only when its user list moved (hq issue 249) Grants first could hold back the very declaration that lets the controller issue them. The holder now goes first, then buckets and memberships, then the rest; a membership that fails holds back only its own machine, and an announcement whose send stopped at its grants is asked again. Whether the holder must go first is read from a digest of the user list it was last sent, not its whole declaration. Migration renumbered to 0058. --- cmd/mesh-controller/delivery_order_test.go | 71 +++++-- cmd/mesh-controller/hold_test.go | 4 +- cmd/mesh-controller/plan.go | 34 ++-- cmd/mesh-controller/push.go | 178 +++++++++++++----- cmd/mesh-controller/sendable.go | 4 + cmd/mesh-controller/upgrades.go | 24 ++- ...0058-a-newer-plan-supersedes-an-older.sql} | 8 + internal/inventory/nodes.go | 20 ++ internal/inventory/sent_bus_users_test.go | 25 +++ 9 files changed, 290 insertions(+), 78 deletions(-) rename internal/inventory/migrations/{0057-a-newer-plan-supersedes-an-older.sql => 0058-a-newer-plan-supersedes-an-older.sql} (60%) create mode 100644 internal/inventory/sent_bus_users_test.go diff --git a/cmd/mesh-controller/delivery_order_test.go b/cmd/mesh-controller/delivery_order_test.go index 7c29c00..878a201 100644 --- a/cmd/mesh-controller/delivery_order_test.go +++ b/cmd/mesh-controller/delivery_order_test.go @@ -6,6 +6,8 @@ import ( "reflect" "strings" "testing" + + "github.com/novox/mesh-controller/internal/link" ) // recordingDelivery is a delivery that writes down what was done, in order, and fails where told. @@ -39,10 +41,11 @@ func ready(names ...string) []readyNode { } // novox/hq issue 249: the grants that come with a module's new declarations are issued before any -// machine is sent the code that uses them. +// machine is sent the code that uses them — except the machine holding the bus, whose declaration +// carries the controller's own right to issue them, and goes first. func TestGrantsAreIssuedBeforeTheDeclarations(t *testing.T) { d := &recordingDelivery{} - digests, err := deliver(t.Context(), d, ready("anchor", "laptop")) + digests, err := deliver(t.Context(), d, "", ready("anchor", "laptop")) if err != nil { t.Fatal(err) } @@ -53,27 +56,71 @@ func TestGrantsAreIssuedBeforeTheDeclarations(t *testing.T) { if digests["anchor"] != "digest-anchor" || digests["laptop"] != "digest-laptop" { t.Fatalf("the digests sent were not answered: %v", digests) } + + held := &recordingDelivery{} + if _, err := deliver(t.Context(), held, "broker", ready("broker", "anchor")); err != nil { + t.Fatal(err) + } + want = []string{"declare broker", "grant broker", "grant anchor", "declare anchor"} + if !reflect.DeepEqual(held.did, want) { + t.Fatalf("with the bus's machine in the send: %v, wanted %v", held.did, want) + } } -// A grant that cannot be issued sends nothing and is an error — never "until the next push". -func TestAGrantThatFailsSendsNothingAndSaysSo(t *testing.T) { - d := &recordingDelivery{grantErr: errors.New("the bus refused the membership")} - _, err := deliver(t.Context(), d, ready("anchor", "laptop")) - if err == nil || !strings.Contains(err.Error(), "the bus refused the membership") { +// A grant that cannot be issued holds back the machines it concerns and is an error the caller +// retries on — never "until the next push" — and the bus's own machine is sent regardless, so the +// grant that would let the controller issue memberships is never held behind them. +func TestAGrantThatFailsHoldsBackWhatItConcerns(t *testing.T) { + // A failure naming no machine (the buckets): everything but the bus's machine. + d := &recordingDelivery{grantErr: errors.New("the bus refused the bucket")} + _, err := deliver(t.Context(), d, "broker", ready("broker", "anchor", "laptop")) + if err == nil || !errors.Is(err, errGrants) || !strings.Contains(err.Error(), "the bus refused the bucket") || + !strings.Contains(err.Error(), "anchor, laptop not sent") { t.Fatalf("a failed grant was not said as the send's failure: %v", err) } - for _, did := range d.did { - if strings.HasPrefix(did, "declare") { - t.Fatalf("code was sent after its grant failed: %v", d.did) - } + if want := []string{"declare broker", "grant broker", "grant anchor", "grant laptop"}; !reflect.DeepEqual(d.did, want) { + t.Fatalf("delivered %v, wanted the bus's machine alone", d.did) } + + // A membership that failed for one machine: that machine alone. + one := &recordingDelivery{grantErr: &grantsRefused{nodes: map[string]error{"laptop": errors.New("no")}}} + _, err = deliver(t.Context(), one, "", ready("anchor", "laptop")) + if !errors.Is(err, errGrants) || !strings.Contains(err.Error(), "laptop not sent") { + t.Fatalf("one machine's refused membership was not said: %v", err) + } + if want := []string{"grant anchor", "grant laptop", "declare anchor"}; !reflect.DeepEqual(one.did, want) { + t.Fatalf("delivered %v, wanted anchor sent and laptop held back", one.did) + } + + // The announced upgrade that hit it is asked again. + if !errors.Is(askAgainOnGrants(err), link.ErrTryAgain) { + t.Fatal("an announcement whose send stopped at its grants is not asked again") + } + if other := errors.New("laptop could not be resolved"); errors.Is(askAgainOnGrants(other), link.ErrTryAgain) { + t.Fatal("any failure is asked again, not only a grant's") + } + // Nothing to send is nothing granted either. none := &recordingDelivery{grantErr: errors.New("never asked")} - if _, err := deliver(t.Context(), none, nil); err != nil || len(none.did) != 0 { + if _, err := deliver(t.Context(), none, "", nil); err != nil || len(none.did) != 0 { t.Fatalf("an empty send granted or failed: %v %v", none.did, err) } } +// Whether the bus's machine goes first is read from the user list alone, by its digest. +func TestTheBusMachineIsBehindByItsUserListAlone(t *testing.T) { + list := "users: [a, b]" + if userListBehind(list, digestOf([]byte(list))) { + t.Fatal("the list it was sent reads as behind") + } + if !userListBehind(list, digestOf([]byte("users: [a]"))) || !userListBehind(list, "") { + t.Fatal("a changed or never-sent list reads as current") + } + if userListBehind("", "") { + t.Fatal("a machine sent no list reads as behind") + } +} + // The machine holding the bus goes first: its declaration carries the user list the new grants are // checked against. Among the machines it is moved to the front; not among them it is added only // when it is behind. diff --git a/cmd/mesh-controller/hold_test.go b/cmd/mesh-controller/hold_test.go index 62660c7..c66d769 100644 --- a/cmd/mesh-controller/hold_test.go +++ b/cmd/mesh-controller/hold_test.go @@ -23,11 +23,11 @@ func TestASendRoundGivesItsHoldBackOnEveryWayOut(t *testing.T) { for name, round := range map[string]func() error{ "a body that cannot be marshalled": func() error { - _, err := sendRound(ctx, open, []string{"anchor"}, unmarshallable, fine) + _, err := sendRound(ctx, open, []string{"anchor"}, unmarshallable, fine, "") return err }, "a send that fails": func() error { - _, err := sendRound(ctx, open, []string{"anchor"}, plain, failing) + _, err := sendRound(ctx, open, []string{"anchor"}, plain, failing, "") return err }, } { diff --git a/cmd/mesh-controller/plan.go b/cmd/mesh-controller/plan.go index a5b773d..2c66137 100644 --- a/cmd/mesh-controller/plan.go +++ b/cmd/mesh-controller/plan.go @@ -398,7 +398,7 @@ func declarationWith(ctx context.Context, open *stores, node string, return sendable{}, err } return sendable{Resources: composed.Resources, Adoption: adoption, - Received: composed.Received, Mesh: with.Mesh, + Received: composed.Received, Mesh: with.Mesh, BusUsers: with.BusUsers, LeftOut: sortedKeysOf(composed.LeftOut), leftOutWhy: composed.LeftOut}, nil } @@ -1341,6 +1341,20 @@ func composeBusUsers(ctx context.Context, inv *inventory.Inventory, // // Asked of what this push resolves to rather than of the seat's holder mesh-wide: the file is a // resource of that module, so the question is whether it is here. + list, missing, err := busUserList(ctx, inv, onThisNode) + if len(missing) > 0 { + fmt.Printf("the bus's user list leaves out %d user(s) the mesh has minted no credential "+ + "for: %s. Each is a user that cannot connect until one is issued\n", + len(missing), strings.Join(missing, ", ")) + } + return list, err +} + +// busUserList is composeBusUsers without saying anything: the list, and the users left out of it +// for want of a credential. Asked on every send to decide whether the machine holding the bus must +// go first (novox/hq issue 249), where saying the same missing users each time would bury them. +func busUserList(ctx context.Context, inv *inventory.Inventory, + onThisNode []catalogue.Manifest) (string, []string, error) { holdsTheBus := false for _, m := range onThisNode { if m.BusUsers != "" && m.ClaimsSeat("mesh-broker") { @@ -1348,37 +1362,33 @@ func composeBusUsers(ctx context.Context, inv *inventory.Inventory, } } if !holdsTheBus { - return "", nil + return "", nil, nil } records, err := inv.BusRecords(ctx) if err != nil { - return "", err + return "", nil, err } users, err := broker.Users(records) if err != nil { - return "", err + return "", nil, err } kept, err := inv.BusUsers(ctx) if err != nil { - return "", err + return "", nil, err } hashes := make(map[string]string, len(kept)) for name, u := range kept { hashes[name] = u.PasswordHash } filled, missing := broker.WithPasswords(users, hashes) - if len(missing) > 0 { - fmt.Printf("the bus's user list leaves out %d user(s) the mesh has minted no credential "+ - "for: %s. Each is a user that cannot connect until one is issued\n", - len(missing), strings.Join(missing, ", ")) - } if len(filled) == 0 { - return "", fmt.Errorf( + return "", missing, fmt.Errorf( "this machine runs the bus and not one user has a credential, so the composed list " + "would refuse every connection in the mesh") } - return broker.ComposeAccounts(filled) + list, err := broker.ComposeAccounts(filled) + return list, missing, err } // providerModuleOf is which module answers a need on the providing node: the one in this node's diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 65ec086..9a2d3c5 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -402,7 +402,7 @@ func pushCommand(ctx context.Context, args []string) error { // Each machine's memberships first, then the declarations (novox/hq issue 249, ADR 0160): a push // is the one most operators run, and on 2026-10-01 it was the one path that issued none. bus := overTheBus{open: open, server: server, signer: ident} - sentDigest, err := deliver(ctx, bus, sending) + sentDigest, err := deliver(ctx, bus, holder, sending) if err != nil { return err } @@ -477,7 +477,7 @@ func pushCommand(ctx context.Context, args []string) error { } return declared, err }, - bus) + bus, holder) refusals = append(refusals, refused...) if err != nil { return err @@ -607,7 +607,7 @@ func composeEach(names []string, allot func(node string) (int64, error), // still sent. Their memberships go before their declarations, as every send's do (issue 249). func sendRound(ctx context.Context, open *stores, names []string, compose func(held context.Context, node string) (sendable, error), - d delivery) ([]string, error) { + d delivery, holder string) ([]string, error) { held, release, err := holdNodes(ctx, open, names) if err != nil { return nil, err @@ -616,7 +616,7 @@ func sendRound(ctx context.Context, open *stores, names []string, sending, refused := composeEach(names, allotting(held, open.inventory), func(node string) (sendable, error) { return compose(held, node) }) - if _, err := deliver(held, d, sending); err != nil { + if _, err := deliver(held, d, holder, sending); err != nil { return refused, err } return refused, nil @@ -632,7 +632,28 @@ type delivery interface { declare(ctx context.Context, s readyNode, body []byte) (string, error) } -// deliver sends the machines their declarations, **the grants first** (novox/hq issue 249). +// errGrants marks a send that stopped because what the machines' modules may do could not be issued +// (novox/hq issue 249). Nothing about the machines is wrong; asked again, it is likely to work, so an +// announcement that hits it is held and asked again. +var errGrants = errors.New("what the machines' modules may do on the bus could not be issued, and code " + + "sent before its grants is refused there") + +// grantsRefused is a grant that failed for some machines and not others: their memberships could not +// be issued, by machine, and only those machines are held back. +type grantsRefused struct{ nodes map[string]error } + +func (g *grantsRefused) Error() string { + names := make([]string, 0, len(g.nodes)) + for n := range g.nodes { + names = append(names, n) + } + sort.Strings(names) + return fmt.Sprintf("the memberships of %s could not be issued; the first: %v", + strings.Join(names, ", "), g.nodes[names[0]]) +} + +// deliver sends the machines their declarations: **the machine holding the bus, then the grants, +// then the rest** (novox/hq issue 249). // // A merge gave a module a new state; its bundle reached every machine within a minute, and the // machines' permissions on the bus did not include the state until somebody pushed by hand: the code @@ -642,31 +663,82 @@ type delivery interface { // the bus (internal/link/bus.go), so issued first it waits for the runtime that arrives after it. A // runtime still on the old code merely holds a grant it does not use yet. // -// **A grant that cannot be issued sends nothing**, and says so as an error. It used to be said and -// passed over — "the machines keep what they derive until the next push" — which reported a rollout -// as done that had delivered code its machines could not run, and left the remedy to a push nobody -// knew to make. Returned, the caller's rollout is not marked sent and is tried again. -func deliver(ctx context.Context, d delivery, sending []readyNode) (map[string]string, error) { +// **The holder's declaration before the grants, though.** The bus's user list travels in it, and the +// controller's own right to publish memberships and raise buckets is in that list (the precedent of +// issue 183): grants first, and a grant the controller is not yet allowed to make would hold the very +// declaration that allows it — a lock only a hand on the broker could open. A runtime already running +// on that machine follows a membership issued after its declaration, as it always has. +// +// **A grant that cannot be issued holds back what it concerns, and says so as an error.** It used to +// be said and passed over — "the machines keep what they derive until the next push" — which reported +// a rollout done that had delivered code its machines could not run. A membership that failed holds +// back its own machine; a failure that names no machine (the buckets) holds back every machine but the +// holder, already sent. The error carries errGrants, so the caller's rollout is not marked sent and is +// tried again. +// +// The grants are issued, not waited on: a membership is a retained message the runtime reads when it +// comes, and the bus answers its publication; nothing here waits for a runtime to have read one. +func deliver(ctx context.Context, d delivery, holder string, sending []readyNode) (map[string]string, error) { digests := map[string]string{} if len(sending) == 0 { return digests, nil } - if err := d.grant(ctx, sending); err != nil { - return digests, fmt.Errorf("nothing was sent: what the machines' modules may do on the bus "+ - "could not be issued, and code sent before its grants is refused there (novox/hq issue 249): %w", err) - } - for _, s := range sending { + send := func(s readyNode) error { // The number is inside the signed bytes, so a replayed older declaration cannot borrow a // newer one's (novox/hq 04-ISSUES/107); it was taken when the composition began (issue 204). body, err := s.declared.Body() if err != nil { - return digests, err + return err } digest, err := d.declare(ctx, s, body) if err != nil { - return digests, err + return err } digests[s.node] = digest + return nil + } + var rest []readyNode + for _, s := range sending { + if holder != "" && s.node == holder { + if err := send(s); err != nil { + return digests, err + } + continue + } + rest = append(rest, s) + } + held := map[string]error{} + if err := d.grant(ctx, sending); err != nil { + var some *grantsRefused + if !errors.As(err, &some) { + var names []string + for _, s := range rest { + names = append(names, s.node) + } + if len(names) == 0 { + return digests, fmt.Errorf("%w: %w", errGrants, err) + } + return digests, fmt.Errorf("%w; %s not sent: %w", errGrants, strings.Join(names, ", "), err) + } + held = some.nodes + } + var notSent []string + for _, s := range rest { + if _, refused := held[s.node]; refused { + notSent = append(notSent, s.node) + continue + } + if err := send(s); err != nil { + return digests, err + } + } + if len(notSent) > 0 { + return digests, fmt.Errorf("%w; %s not sent: %w", errGrants, strings.Join(notSent, ", "), + &grantsRefused{nodes: held}) + } + if len(held) > 0 { + // Only the holder's own memberships failed, and it was sent before them. + return digests, fmt.Errorf("%w: %w", errGrants, &grantsRefused{nodes: held}) } return digests, nil } @@ -694,6 +766,16 @@ func (b overTheBus) declare(ctx context.Context, s readyNode, body []byte) (stri if err != nil { return "", err } + if s.declared.BusUsers != "" { + // And the user list it carried, so the next send reads whether it must go first from the + // list alone (novox/hq issue 249). On the same outliving context as the send's record. + kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second) + err := b.open.inventory.RecordSentBusUsers(kept, s.node, digestOf([]byte(s.declared.BusUsers))) + cancel() + if err != nil { + return "", err + } + } fmt.Printf("%ssent %s %d resource(s)\n", b.indent, s.node, len(s.declared.Resources)) return digest, nil } @@ -722,14 +804,13 @@ func brokerFirst(names []string, holder string, behind bool) []string { } // brokerBehind is the machine holding the bus — the one whose declaration carries the user list — -// and, when it is not among the machines named, whether what it should be differs from what it was -// last sent. +// and, when it is not among the machines named, whether the user list it would be sent now differs +// from the one it was last sent (novox/hq issue 249). // -// **Its whole declaration, because that is what is recorded** (novox/hq issue 249). The mesh keeps -// a digest of what each machine was last sent and not of the user list inside it (ADR 0043: the -// list is composed on each push, never kept), so "the user list changed" is read as "the bus's -// machine is behind". That is the safe direction: a machine sent what it should be is never wrong, -// and whatever else it was behind on is what any push would have sent it. +// **The user list alone, not the whole declaration.** Read from the whole declaration, any change +// pending on that machine — an upgrade its policy records rather than rolls out — went with every +// send anywhere, and a module running there always put it in its first wave. A digest of the list +// last sent is kept for this (ADR 0043: the list is composed on each push, never kept itself). func brokerBehind(ctx context.Context, open *stores, names []string) (string, bool, error) { inv := open.inventory holders, err := seatHolders(ctx, inv) @@ -753,23 +834,26 @@ func brokerBehind(ctx context.Context, open *stores, names []string) (string, bo return h.Node, false, nil } } - node, err := inv.NodeByName(ctx, h.Node) + plan, _, err := planFor(ctx, open, h.Node) if err != nil { - return "", false, err - } - would, err := wouldSend(ctx, open, []inventory.Node{node}) - if err != nil { - return "", false, err - } - if would[h.Node] == "" { // It cannot be worked out: sending it would refuse the whole send, and `plan` says why. return h.Node, false, nil } - sent, err := inv.Outstanding(ctx, h.Node) + list, _, err := busUserList(ctx, inv, plan.Modules) + if err != nil { + return h.Node, false, nil + } + sent, err := inv.SentBusUsers(ctx, h.Node) if err != nil { return "", false, err } - return h.Node, would[h.Node] != sent, nil + return h.Node, userListBehind(list, sent), nil +} + +// userListBehind is whether the user list composed now is not the one last sent, by its digest. An +// empty list composed is never behind: there is nothing for it to carry. +func userListBehind(now, sentDigest string) bool { + return now != "" && digestOf([]byte(now)) != sentDigest } // couldNotBeResolved is what a push ends with when some machines could not be worked out. @@ -867,7 +951,7 @@ func sendToEach(ctx context.Context, open *stores, names []string) ([]string, er // composition. **Issued before the declarations** (novox/hq issue 249): the runtime the // membership is for arrives with the declaration, and a membership waits for it on the bus; the // code arriving first was refused its own state until somebody pushed. - if _, err := deliver(ctx, overTheBus{open: open, server: server, signer: ident, indent: " "}, sending); err != nil { + if _, err := deliver(ctx, overTheBus{open: open, server: server, signer: ident, indent: " "}, holder, sending); err != nil { return nil, err } sent := make([]string, 0, len(sending)) @@ -898,9 +982,10 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se // state: its bundle asked for a bucket that did not exist. Idempotent and cheap. // // **A failure here is the send's failure** (novox/hq issue 249). It was said and the push stood, - // because the declarations were already away; they are sent after this now, and a module whose - // state does not exist is a module that fails its start, so nothing is sent and the caller tries - // again rather than reporting the rollout done. + // because the declarations were already away; they are sent after this now — all but the bus's + // own machine, sent before it (deliver) — and a module whose state does not exist is a module + // that fails its start, so they are not sent and the caller tries again rather than reporting the + // rollout done. buckets, err := open.inventory.DeclaredBuckets(ctx) if err != nil { return fmt.Errorf("the modules' state could not be read, so no bucket was asserted: %w", err) @@ -909,8 +994,8 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se return fmt.Errorf("the modules' state could not be asserted on the bus: %w", err) } // Every membership is tried, and the first failure named once. - issued, failed := 0, 0 - var first error + issued := 0 + refused := map[string]error{} for _, s := range sent { node := s.node for _, d := range records.Assigned[node] { @@ -931,10 +1016,9 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se return err } if err := bus.PublishMembership(ctx, node, d.Module, body); err != nil { - if first == nil { - first = err + if refused[node] == nil { + refused[node] = fmt.Errorf("%s: %w", d.Module, err) } - failed++ continue } issued++ @@ -943,10 +1027,10 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se if issued > 0 { fmt.Printf(" issued %d membership(s)\n", issued) } - if failed > 0 { - // Returned, never passed over (novox/hq issue 249): the declarations that need these are not - // sent, and the rollout that asked is tried again rather than waiting for a push. - return fmt.Errorf("%d membership(s) could not be issued; the first: %w", failed, first) + if len(refused) > 0 { + // Returned, never passed over (novox/hq issue 249): the declarations of the machines they are + // for are not sent, and the rollout that asked is tried again rather than waiting for a push. + return &grantsRefused{nodes: refused} } return nil } diff --git a/cmd/mesh-controller/sendable.go b/cmd/mesh-controller/sendable.go index 6722732..4219e97 100644 --- a/cmd/mesh-controller/sendable.go +++ b/cmd/mesh-controller/sendable.go @@ -31,6 +31,10 @@ type sendable struct { // the same composition as its received files, and every machine's private-network address. Received map[string]map[string][]catalogue.Contribution Mesh []string + // BusUsers is the bus's user list this declaration carries, empty for every machine but the one + // holding the bus; not sent apart from the file it is in. Its digest is recorded once sent, so + // whether that machine must go first is read from the list alone (novox/hq issue 249). + BusUsers string // LeftOut is every module of the machine's set left out of this declaration because a stored // setting cannot compose with its definition (novox/hq ADR 0163, rule 6), sorted. The host // keeps that module's held things and touches none of its containers; a machine is told diff --git a/cmd/mesh-controller/upgrades.go b/cmd/mesh-controller/upgrades.go index 3c71189..3fa7aab 100644 --- a/cmd/mesh-controller/upgrades.go +++ b/cmd/mesh-controller/upgrades.go @@ -65,13 +65,17 @@ func (f following) Upgraded(ctx context.Context, u link.Upgraded) error { if plans, err := inv.OpenPlans(ctx); err != nil { return notNow(err) } else if id := rolledOutByAPlan(plans, u.Module); id != "" { - fmt.Printf("%s moved to %s; %s rolls it out to %s\n", u.Module, shortCommit(u.Commit), id, readableList(on)) + // Said with its remedy: a plan that ends without sending it — failed, or closed by hand — leaves + // these machines behind, which `status` lists and `push --behind` sends (novox/hq issue 249). + fmt.Printf("%s moved to %s; %s rolls it out to %s — if that plan ends without sending it, "+ + "`status` lists them as behind and `push --behind` sends it\n", + u.Module, shortCommit(u.Commit), id, readableList(on)) return nil } if decision.Together { fmt.Printf("%s moved to %s; sending %s together\n", u.Module, shortCommit(u.Commit), readableList(on)) - return sendTo(ctx, f.open, on) + return askAgainOnGrants(sendTo(ctx, f.open, on)) } // One at a time, and stopping at the first that fails. // @@ -82,18 +86,28 @@ func (f following) Upgraded(ctx context.Context, u link.Upgraded) error { u.Module, shortCommit(u.Commit), readableList(on)) for _, node := range on { if err := sendTo(ctx, f.open, []string{node}); err != nil { - return fmt.Errorf("%s did not take %s, so the machines after it were left alone: %w", - node, u.Module, err) + return askAgainOnGrants(fmt.Errorf("%s did not take %s, so the machines after it were left alone: %w", + node, u.Module, err)) } } return nil } +// askAgainOnGrants marks a send that stopped at its grants as one to ask again (novox/hq issue 249): +// an announcement handled by a send whose memberships could not be issued is held and redelivered, +// rather than taken as handled with the machines left on the old version. +func askAgainOnGrants(err error) error { + if err != nil && errors.Is(err, errGrants) && !errors.Is(err, link.ErrTryAgain) { + return fmt.Errorf("%w: %w", link.ErrTryAgain, err) + } + return err +} + // rolledOutByAPlan is the open plan that will send a module's machines its new build — one holding // the module that has not finished sending it — or empty when none will (novox/hq issue 249). func rolledOutByAPlan(plans []inventory.Plan, module string) string { for _, p := range plans { - if s, holds := p.Modules[module]; p.Open() && holds && (s == nil || s.SentAt == nil) { + if s, holds := p.Modules[module]; p.Open() && holds && (s == nil || (s.SentAt == nil && s.State != "failed")) { return p.ID } } diff --git a/internal/inventory/migrations/0057-a-newer-plan-supersedes-an-older.sql b/internal/inventory/migrations/0058-a-newer-plan-supersedes-an-older.sql similarity index 60% rename from internal/inventory/migrations/0057-a-newer-plan-supersedes-an-older.sql rename to internal/inventory/migrations/0058-a-newer-plan-supersedes-an-older.sql index 8a83eb6..bf3768f 100644 --- a/internal/inventory/migrations/0057-a-newer-plan-supersedes-an-older.sql +++ b/internal/inventory/migrations/0058-a-newer-plan-supersedes-an-older.sql @@ -12,3 +12,11 @@ -- is superseded by the next plan of its repository, whichever branch — nothing is lost by it, since -- what it had not built is folded into the plan that supersedes it. alter table release_plan add column branch text not null default ''; + +-- And the bus's user list the machine holding the bus was last sent, as a digest (novox/hq issue +-- 249). A module's new grants are refused by the bus until its user list says them, so that machine +-- is sent first whenever the list it would be sent differs from the one it was. Read from its whole +-- declaration, every pending change on it — a recorded upgrade the operator chose not to roll out — +-- went with every send anywhere. A digest and never the list (ADR 0043: the list is composed on each +-- push, never kept). Empty for a machine never sent one, which reads as behind once. +alter table node add column sent_bus_users text not null default ''; diff --git a/internal/inventory/nodes.go b/internal/inventory/nodes.go index fef70aa..e16e323 100644 --- a/internal/inventory/nodes.go +++ b/internal/inventory/nodes.go @@ -921,6 +921,26 @@ func (i *Inventory) RecordSent(ctx context.Context, node, digest string) error { return err } +// RecordSentBusUsers keeps a digest of the bus's user list a machine was just sent, by its name +// (novox/hq issue 249): whether the machine holding the bus must go first is whether this differs +// from the list composed now. +func (i *Inventory) RecordSentBusUsers(ctx context.Context, name, digest string) error { + _, err := i.store.Pool().Exec(ctx, + `update node set sent_bus_users = $2 where name = $1`, name, digest) + return err +} + +// SentBusUsers is the digest of the bus's user list a machine was last sent, empty for none. +func (i *Inventory) SentBusUsers(ctx context.Context, name string) (string, error) { + var sent string + err := i.store.Pool().QueryRow(ctx, + `select sent_bus_users from node where name = $1`, name).Scan(&sent) + if errors.Is(err, pgx.ErrNoRows) { + return "", nil + } + return sent, err +} + // Outstanding is the digest of the declaration a machine was last sent, by its name, and empty // for one that has never been sent anything. // diff --git a/internal/inventory/sent_bus_users_test.go b/internal/inventory/sent_bus_users_test.go new file mode 100644 index 0000000..5717571 --- /dev/null +++ b/internal/inventory/sent_bus_users_test.go @@ -0,0 +1,25 @@ +package inventory + +import "testing" + +// novox/hq issue 249: the digest of the user list a machine was last sent is kept, by its name, and +// is empty for a machine never sent one. +func TestTheUserListAMachineWasSentIsKept(t *testing.T) { + inv := ForTest(t) + ctx := t.Context() + if _, err := inv.AddNode(ctx, "anchor"); err != nil { + t.Fatal(err) + } + if sent, err := inv.SentBusUsers(ctx, "anchor"); err != nil || sent != "" { + t.Fatalf("a machine never sent a list has %q: %v", sent, err) + } + if err := inv.RecordSentBusUsers(ctx, "anchor", "abc"); err != nil { + t.Fatal(err) + } + if sent, err := inv.SentBusUsers(ctx, "anchor"); err != nil || sent != "abc" { + t.Fatalf("the list sent was not kept: %q %v", sent, err) + } + if sent, err := inv.SentBusUsers(ctx, "nobody"); err != nil || sent != "" { + t.Fatalf("a machine the mesh does not know: %q %v", sent, err) + } +} From 16c5e78fa88c9e21842ef4780ca812e91a661f66 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 5 Oct 2026 18:17:52 +0200 Subject: [PATCH 7/7] Stop a rollout whose first machine is silent, and choose one that is heard from (hq issue 249, ADR 0218) A first machine that does not report within the bound now stops the module's rollout, naming it. The first machine is the first by name heard from lately; reports are judged by what the store says was last sent. A plan that ends says what it built and never sent, and a failed send's error is kept in the plan's note. --- cmd/mesh-controller/release_plan.go | 90 +++++++++++++++++----- cmd/mesh-controller/rollout_first_test.go | 93 +++++++++++++++++------ 2 files changed, 140 insertions(+), 43 deletions(-) diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 74455b1..8679101 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -391,6 +391,10 @@ func planBuilt(ctx context.Context, open *stores, module, commit, failed string, state.Why = failed p.State = inventory.PlanFailed p.Note = fmt.Sprintf("%s failed to build in tier %d", module, p.Tier) + sayUnsent(p, func(m string) bool { + u, err := inv.UpgradeOf(ctx, m) + return err == nil && u.RollOut + }) } else { state.State = "built" state.BuiltAt = &now @@ -452,8 +456,17 @@ func advanceHeld(ctx context.Context, open *stores) { moved, err := advanceOnce(ctx, open, p, edges, rollsOut) if err != nil { fmt.Printf("%s: %v\n", p.ID, err) + // Kept in the plan, so `plans` says why it has not moved rather than the log alone; + // the state is left as it was and the step is tried again on the next tick. + p.Note = "tier " + fmt.Sprint(p.Tier) + ": " + err.Error() + " — tried again" + if err := inv.SavePlan(ctx, *p); err != nil { + fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err) + } break } + if p.State == inventory.PlanFailed { + sayUnsent(p, rollsOut) + } if err := inv.SavePlan(ctx, *p); err != nil { fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err) break @@ -545,13 +558,14 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, return false, err } var reports []inventory.Reported - if state.FirstAt != nil && !policy.Together { + if !policy.Together { + // Read for the choice of the first machine as well as for its report. if reports, err = inv.LastReports(ctx); err != nil { return false, err } } - step := nextRollout(*state, running, policy.Together, reports) now := time.Now().UTC() + step := nextRollout(*state, running, policy.Together, reports, now, planWaitBound) switch { case step.failed != "": state.Why = step.failed @@ -561,11 +575,7 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, fmt.Printf("%s: %s\n", p.ID, p.Note) return true, nil case step.waiting != "": - wait := fmt.Sprintf("%s on %s, sent first at %s", m, step.waiting, state.FirstAt.Format("15:04")) - if now.Sub(*state.FirstAt) > planWaitBound { - wait += " — LATE" - } - pending = append(pending, wait) + pending = append(pending, fmt.Sprintf("%s on %s, sent first at %s", m, step.waiting, state.FirstAt.Format("15:04"))) continue case len(step.send) == 0: // No machine runs it: nothing to send, and nothing to wait for. @@ -665,22 +675,36 @@ type rolloutStep struct { // nextRollout is the next step of one module's rollout in a plan (novox/hq issue 249, ADR 0218). // -// Together, every machine running it at once, as the policy says. Otherwise one machine first — -// the first by name, so the choice is the same on every controller and every resume — and the rest -// once each machine the first send reached has reported, after that send, that it applied the -// declaration it was last sent. A report after the send that failed or refused it is the rollout's -// end: the rest are not sent. A machine that has not reported since is waited for; how long it has -// been is the plan's to say. -func nextRollout(s inventory.PlanModule, running []string, together bool, reports []inventory.Reported) rolloutStep { +// Together, every machine running it at once, as the policy says. Otherwise one machine first — the +// first by name among those that have reported within the bound, so a laptop that is away is not +// the one the rest wait on; the first by name when none has; the same choice on every controller and +// every resume — and the rest once each machine the first send reached reports, about the +// declaration it was last sent, that it applied it. Compared by the store's own record of what was +// sent (Reported.Current), never by this controller's clock against the machine's. +// +// **A first machine that fails, refuses, or does not report within the bound stops the rollout +// there** (ADR 0218 §2), naming the machine; the rest are not sent. A wait with no end is not a +// rollout: it held the plan open for ever, read as work in progress (issue 254). +func nextRollout(s inventory.PlanModule, running []string, together bool, reports []inventory.Reported, + now time.Time, bound time.Duration) rolloutStep { if len(running) == 0 { return rolloutStep{} } if together { return rolloutStep{send: running} } + byNode := map[string]inventory.Reported{} + for _, r := range reports { + byNode[r.Node] = r + } if s.FirstAt == nil { sorted := append([]string{}, running...) sort.Strings(sorted) + for _, n := range sorted { + if r, said := byNode[n]; said && r.At != nil && now.Sub(*r.At) <= bound { + return rolloutStep{send: []string{n}, first: true} + } + } return rolloutStep{send: sorted[:1], first: true} } sentFirst := map[string]bool{} @@ -693,15 +717,11 @@ func nextRollout(s inventory.PlanModule, running []string, together bool, report rest = append(rest, n) } } - byNode := map[string]inventory.Reported{} - for _, r := range reports { - byNode[r.Node] = r - } var waiting, failed []string for _, n := range s.First { r, said := byNode[n] - // Only a report after the send, about what it was last sent, says anything about this build. - if !said || r.At == nil || r.At.Before(*s.FirstAt) || !r.Current { + // Only a report about what it was last sent says anything about this build. + if !said || r.At == nil || !r.Current { waiting = append(waiting, n) continue } @@ -717,11 +737,34 @@ func nextRollout(s inventory.PlanModule, running []string, together bool, report return rolloutStep{failed: strings.Join(failed, "; "), rest: rest} } if len(waiting) > 0 { + if now.Sub(*s.FirstAt) > bound { + return rolloutStep{failed: fmt.Sprintf("%s did not report it applied within %s", + strings.Join(waiting, ", "), bound), rest: rest} + } return rolloutStep{waiting: strings.Join(waiting, ", ")} } return rolloutStep{send: rest} } +// sayUnsent adds to an ended plan's note the modules it built and never sent (novox/hq issue 249). +// An announced move of a module a plan held was left to that plan; a plan that ends without +// sending it — failed elsewhere, or closed by hand — would leave its machines behind with nothing +// saying so. A module whose rollout stopped at its first machine is not among them: that stop was +// the point. Said once. +func sayUnsent(p *inventory.Plan, rollsOut func(string) bool) { + var unsent []string + for name, s := range p.Modules { + if s != nil && s.State == "built" && s.SentAt == nil && s.FirstAt == nil && rollsOut(name) { + unsent = append(unsent, name) + } + } + if len(unsent) == 0 || strings.Contains(p.Note, "built and never sent") { + return + } + sort.Strings(unsent) + p.Note += "; built and never sent: " + strings.Join(unsent, ", ") + " — `push --behind` sends them" +} + // planTicker advances open plans on a timer, for the steps outcomes alone cannot take. func planTicker(ctx context.Context, open *stores) { advancePlans(ctx, open) @@ -899,6 +942,10 @@ func plansCommand(ctx context.Context, args []string) error { } p.State = inventory.PlanFailed p.Note = how + " by hand at tier " + fmt.Sprint(p.Tier) + sayUnsent(&p, func(m string) bool { + u, err := inv.UpgradeOf(ctx, m) + return err == nil && u.RollOut + }) if err := inv.SavePlan(ctx, p); err != nil { return err } @@ -978,9 +1025,10 @@ func planWhatIf(ctx context.Context, inv *inventory.Inventory, repository string how := "built; its policy records, so nothing is sent" if u, err := inv.UpgradeOf(ctx, name); err == nil && u.RollOut { running, _ := inv.Running(ctx, name) + reports, _ := inv.LastReports(ctx) how = "built, then sent to " + orNone(strings.Join(running, ", ")) // One machine first unless the policy says together (novox/hq issue 249). - if first := nextRollout(inventory.PlanModule{}, running, u.Together, nil); first.first && len(running) > 1 { + if first := nextRollout(inventory.PlanModule{}, running, u.Together, reports, time.Now(), planWaitBound); first.first && len(running) > 1 { how = fmt.Sprintf("built, then sent to %s first and to the rest once it has applied it", first.send[0]) } diff --git a/cmd/mesh-controller/rollout_first_test.go b/cmd/mesh-controller/rollout_first_test.go index cb8ce79..644f4be 100644 --- a/cmd/mesh-controller/rollout_first_test.go +++ b/cmd/mesh-controller/rollout_first_test.go @@ -14,47 +14,51 @@ import ( // once, as before. func TestAPlanSendsOneMachineFirstAndTheRestAfterItsReport(t *testing.T) { running := []string{"novox", "ace", "g14"} + sentAt := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + now := sentAt.Add(5 * time.Minute) + bound := 30 * time.Minute + after := sentAt.Add(time.Minute) + next := func(s inventory.PlanModule, running []string, together bool, reports []inventory.Reported) rolloutStep { + return nextRollout(s, running, together, reports, now, bound) + } // Together: every machine at once. - if step := nextRollout(inventory.PlanModule{}, running, true, nil); !reflect.DeepEqual(step.send, running) || step.first { + if step := next(inventory.PlanModule{}, running, true, nil); !reflect.DeepEqual(step.send, running) || step.first { t.Fatalf("a together policy did not send every machine at once: %+v", step) } - // Otherwise the first by name, alone. - step := nextRollout(inventory.PlanModule{}, running, false, nil) + // Otherwise the first by name, alone, when none has reported lately. + step := next(inventory.PlanModule{}, running, false, nil) if !step.first || !reflect.DeepEqual(step.send, []string{"ace"}) { t.Fatalf("the first send was %+v, wanted ace alone", step) } - sentAt := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) - before, after := sentAt.Add(-time.Minute), sentAt.Add(time.Minute) state := inventory.PlanModule{First: []string{"ace"}, FirstAt: &sentAt} - report := func(at time.Time, outcome string, current bool) []inventory.Reported { - return []inventory.Reported{{Node: "ace", At: &at, Outcome: outcome, Current: current}, + report := func(outcome string, current bool) []inventory.Reported { + return []inventory.Reported{{Node: "ace", At: &after, Outcome: outcome, Current: current}, {Node: "g14", At: &after, Outcome: inventory.OutcomeApplied, Current: true}} } - // No report yet, a report from before the send, or one about an older declaration: wait. + // No report yet, or one about an older declaration than it was last sent: wait. for what, reports := range map[string][]inventory.Reported{ - "no report": nil, - "a report before the send": report(before, inventory.OutcomeApplied, true), - "a report about older": report(after, inventory.OutcomeApplied, false), + "no report": nil, + "a report about older": report(inventory.OutcomeApplied, false), } { - step := nextRollout(state, running, false, reports) + step := next(state, running, false, reports) if len(step.send) != 0 || step.waiting != "ace" || step.failed != "" { t.Errorf("%s: %+v, wanted to wait for ace", what, step) } } - // Applied after the send: the rest, and only the rest. - step = nextRollout(state, running, false, report(after, inventory.OutcomeApplied, true)) + // Applied what it was last sent: the rest, and only the rest. + step = next(state, running, false, report(inventory.OutcomeApplied, true)) if step.first || !reflect.DeepEqual(step.send, []string{"novox", "g14"}) { t.Fatalf("after ace applied it the plan sent %+v, wanted novox and g14", step) } - // Failed or refused after the send: stop, the rest untouched. + // Failed or refused: stop, the rest untouched. for _, outcome := range []string{inventory.OutcomeFailed, inventory.OutcomeRefused} { - step := nextRollout(state, running, false, report(after, outcome, true)) + step := next(state, running, false, report(outcome, true)) if len(step.send) != 0 || !strings.Contains(step.failed, "ace "+outcome) || !reflect.DeepEqual(step.rest, []string{"novox", "g14"}) { t.Errorf("a first machine that %s it: %+v", outcome, step) @@ -64,25 +68,50 @@ func TestAPlanSendsOneMachineFirstAndTheRestAfterItsReport(t *testing.T) { // The machine holding the bus went with the first send: the rest wait for it too, and it is // not sent again. both := inventory.PlanModule{First: []string{"novox", "ace"}, FirstAt: &sentAt} - half := report(after, inventory.OutcomeApplied, true) - if step := nextRollout(both, running, false, half); step.waiting != "novox" { + half := report(inventory.OutcomeApplied, true) + if step := next(both, running, false, half); step.waiting != "novox" { t.Fatalf("the plan did not wait for the bus's machine sent first: %+v", step) } all := append(half, inventory.Reported{Node: "novox", At: &after, Outcome: inventory.OutcomeApplied, Current: true}) - if step := nextRollout(both, running, false, all); !reflect.DeepEqual(step.send, []string{"g14"}) { + if step := next(both, running, false, all); !reflect.DeepEqual(step.send, []string{"g14"}) { t.Fatalf("after both applied it the plan sent %+v, wanted g14 alone", step) } // One machine, or none: nothing is waited for that cannot come. - if step := nextRollout(inventory.PlanModule{}, nil, false, nil); len(step.send) != 0 || step.first { + if step := next(inventory.PlanModule{}, nil, false, nil); len(step.send) != 0 || step.first { t.Fatalf("a module nothing runs was sent: %+v", step) } - only := inventory.PlanModule{First: []string{"ace"}, FirstAt: &sentAt} - if step := nextRollout(only, []string{"ace"}, false, report(after, inventory.OutcomeApplied, true)); len(step.send) != 0 || step.waiting != "" { + if step := next(state, []string{"ace"}, false, report(inventory.OutcomeApplied, true)); len(step.send) != 0 || step.waiting != "" { t.Fatalf("a module on one machine waited for more: %+v", step) } } +// ADR 0218 §2: a first machine that does not report within the bound stops the rollout there, +// naming the machine and the bound; the rest are left alone. +func TestAFirstMachineThatDoesNotReportStopsTheRollout(t *testing.T) { + sentAt := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + state := inventory.PlanModule{First: []string{"ace"}, FirstAt: &sentAt} + step := nextRollout(state, []string{"ace", "g14"}, false, nil, sentAt.Add(31*time.Minute), 30*time.Minute) + if !strings.Contains(step.failed, "ace did not report it applied within 30m") || + !reflect.DeepEqual(step.rest, []string{"g14"}) || len(step.send) != 0 { + t.Fatalf("a silent first machine: %+v", step) + } +} + +// The first machine is the first by name among those heard from lately: a laptop that is away is +// not the one the rest wait on. When none has been heard from, the first by name. +func TestTheFirstMachineIsOneThatHasReportedLately(t *testing.T) { + now := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + lately, long := now.Add(-time.Minute), now.Add(-3*time.Hour) + reports := []inventory.Reported{ + {Node: "ace", At: &long}, {Node: "g14", At: &lately}, {Node: "novox", At: &lately}, + } + step := nextRollout(inventory.PlanModule{}, []string{"novox", "ace", "g14"}, false, reports, now, 30*time.Minute) + if !step.first || !reflect.DeepEqual(step.send, []string{"g14"}) { + t.Fatalf("the first send was %+v, wanted g14, the first heard from lately", step) + } +} + // An announced move of a module an open plan is still rolling out is left to the plan: sending it // here as well put the bundle on every machine at once (novox/hq issue 249). func TestAnAnnouncedMoveIsLeftToThePlanRollingItOut(t *testing.T) { @@ -101,3 +130,23 @@ func TestAnAnnouncedMoveIsLeftToThePlanRollingItOut(t *testing.T) { } } } + +// A plan that ends without sending what it built says so, with the remedy; a module whose rollout +// stopped at its first machine is not among them. +func TestAnEndedPlanSaysWhatItBuiltAndNeverSent(t *testing.T) { + at := time.Now() + p := inventory.Plan{State: inventory.PlanFailed, Note: "closed by hand at tier 1", + Modules: map[string]*inventory.PlanModule{ + "agent": {State: "built"}, + "stopped": {State: "built", First: []string{"ace"}, FirstAt: &at}, + "sent": {State: "built", SentAt: &at}, + "notes": {State: "built"}, + "later": {}, + }} + rollsOut := func(m string) bool { return m != "notes" } + sayUnsent(&p, rollsOut) + sayUnsent(&p, rollsOut) + if p.Note != "closed by hand at tier 1; built and never sent: agent — `push --behind` sends them" { + t.Fatalf("the note reads %q", p.Note) + } +}