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 }