diff --git a/cmd/mesh-controller/delivery_order_test.go b/cmd/mesh-controller/delivery_order_test.go new file mode 100644 index 0000000..878a201 --- /dev/null +++ b/cmd/mesh-controller/delivery_order_test.go @@ -0,0 +1,145 @@ +package main + +import ( + "context" + "errors" + "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. +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 — 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")) + 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) + } + + 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 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) + } + 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 { + 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. +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..c66d769 100644 --- a/cmd/mesh-controller/hold_test.go +++ b/cmd/mesh-controller/hold_test.go @@ -18,16 +18,16 @@ 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 { - _, 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/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/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 380f627..9a2d3c5 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, holder, 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, holder) 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, holder string) ([]string, error) { held, release, err := holdNodes(ctx, open, names) if err != nil { return nil, err @@ -633,18 +616,246 @@ 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, holder, 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) +} + +// 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 +// 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. +// +// **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 + } + 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 err + } + digest, err := d.declare(ctx, s, body) + if err != nil { + 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 +} + +// 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 + } + 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 +} + +// 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 the user list it would be sent now differs +// from the one it was last sent (novox/hq issue 249). +// +// **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) + 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 + } + } + plan, _, err := planFor(ctx, open, h.Node) + if err != nil { + // It cannot be worked out: sending it would refuse the whole send, and `plan` says why. + return h.Node, false, nil + } + 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, 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. // // **After the rest have been sent, never instead of sending them.** It is still an error, because @@ -667,23 +878,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 +936,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: " "}, holder, 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,18 +979,23 @@ 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 — 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) } - // 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. - issued, failed := 0, 0 - var first error + 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 := 0 + refused := map[string]error{} for _, s := range sent { node := s.node for _, d := range records.Assigned[node] { @@ -791,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++ @@ -803,9 +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 { - 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) + 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/release_plan.go b/cmd/mesh-controller/release_plan.go index 267d2c7..8679101 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 @@ -339,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 @@ -400,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 @@ -470,6 +535,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 +553,64 @@ 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 !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 + } + } now := time.Now().UTC() - state.SentAt = &now - if len(running) == 0 { + step := nextRollout(*state, running, policy.Together, reports, now, planWaitBound) + 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 != "": + 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. + 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 +661,110 @@ 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 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{} + for _, n := range s.First { + sentFirst[n] = true + } + var rest []string + for _, n := range running { + if !sentFirst[n] { + rest = append(rest, n) + } + } + var waiting, failed []string + for _, n := range s.First { + r, said := byNode[n] + // Only a report about what it was last sent says anything about this build. + if !said || r.At == nil || !r.Current { + 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 { + 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) @@ -563,6 +788,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 := "" @@ -693,7 +920,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 @@ -702,25 +941,16 @@ 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) + 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 } - 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) @@ -795,7 +1025,13 @@ 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, 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]) + } 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..644f4be --- /dev/null +++ b/cmd/mesh-controller/rollout_first_test.go @@ -0,0 +1,152 @@ +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"} + 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 := 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, 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) + } + + state := inventory.PlanModule{First: []string{"ace"}, FirstAt: &sentAt} + 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, 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 about older": report(inventory.OutcomeApplied, false), + } { + 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 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: stop, the rest untouched. + for _, outcome := range []string{inventory.OutcomeFailed, inventory.OutcomeRefused} { + 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) + } + } + + // 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(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 := 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 := next(inventory.PlanModule{}, nil, false, nil); len(step.send) != 0 || step.first { + t.Fatalf("a module nothing runs was sent: %+v", step) + } + 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) { + 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) + } + } +} + +// 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) + } +} 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/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/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 b16542b..3fa7aab 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" @@ -55,10 +56,26 @@ 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 != "" { + // 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. // @@ -69,13 +86,34 @@ 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 && s.State != "failed")) { + 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 { @@ -323,6 +361,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], ", ")) @@ -330,6 +407,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, ", "))) @@ -423,25 +508,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 +557,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, "/") 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/0058-a-newer-plan-supersedes-an-older.sql b/internal/inventory/migrations/0058-a-newer-plan-supersedes-an-older.sql new file mode 100644 index 0000000..bf3768f --- /dev/null +++ b/internal/inventory/migrations/0058-a-newer-plan-supersedes-an-older.sql @@ -0,0 +1,22 @@ +-- 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 ''; + +-- 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/plans.go b/internal/inventory/plans.go index c68bafe..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. @@ -37,8 +40,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. @@ -47,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. @@ -63,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 } @@ -95,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 @@ -106,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) + } +} 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) + } +}