Issue grants before code, and the bus's machine first (hq issue 249)

A module's new state reached every machine before the permissions to use
it, which came only with a later push. Memberships and buckets now go
before declarations, the machine holding mesh-broker goes first when its
user list must change, and a grant that fails sends nothing and is an
error so the rollout is retried.
This commit is contained in:
jochen
2026-10-05 18:17:52 +02:00
parent d7f359c498
commit 22660dc274
3 changed files with 317 additions and 78 deletions
@@ -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)
}
}
}
+2 -2
View File
@@ -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 {
+217 -76
View File
@@ -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
}