diff --git a/cmd/mesh-controller/issue_204_test.go b/cmd/mesh-controller/issue_204_test.go new file mode 100644 index 0000000..9bcc5d2 --- /dev/null +++ b/cmd/mesh-controller/issue_204_test.go @@ -0,0 +1,93 @@ +package main + +import ( + "context" + "testing" + + "github.com/novox/mesh-controller/internal/inventory" +) + +// A declaration composed earlier is numbered lower than one composed later, whatever order the two +// are sent in (novox/hq issue 204). The number used to be taken at send time, after composing, so a +// declaration composed before an assignment changed and sent after a newer one carried the higher +// number — and the machine, which refuses a lower number, took the older content as the mesh's +// newest word. Taken before the composition reads anything, the order of numbers is the order of +// compositions, and the host's refusal does what it is for. +func TestADeclarationComposedEarlierIsNumberedLowerWhateverOrderItIsSent(t *testing.T) { + allot := numbered() + var composed []string + compose := func(stamp string) func(string) (sendable, error) { + return func(node string) (sendable, error) { + composed = append(composed, stamp) + return sendable{Resources: []map[string]any{{"id": node + "." + stamp}}}, nil + } + } + // Composed first — before an assignment changed — and sent last. + stale, _ := composeEach([]string{"anchor"}, allot, compose("before")) + // Composed after the change, sent first. + fresh, _ := composeEach([]string{"anchor"}, allot, compose("after")) + + if stale[0].declared.Sequence != 1 || fresh[0].declared.Sequence != 2 { + t.Fatalf("the numbers do not follow the compositions: before=%d after=%d", + stale[0].declared.Sequence, fresh[0].declared.Sequence) + } + // Sent in the other order, the numbers do not change — so the machine that has applied the + // fresh one (2) refuses the stale one (1) when it arrives late. + if !(stale[0].declared.Sequence < fresh[0].declared.Sequence) { + t.Fatal("a declaration composed earlier must carry the lower number, however late it is sent") + } + if len(composed) != 2 || composed[0] != "before" { + t.Fatalf("compositions happened in an unexpected order: %v", composed) + } +} + +// The number is taken before the first read of the composition, not after it: an allotter that +// fails leaves nothing composed for that machine, and the others are still composed. +func TestTheNumberIsTakenBeforeComposingAndItsFailureIsARefusal(t *testing.T) { + calls := 0 + allot := func(node string) (int64, error) { + if node == "anchor" { + return 0, context.DeadlineExceeded + } + return 7, nil + } + sending, refusals := composeEach([]string{"anchor", "laptop"}, allot, func(node string) (sendable, error) { + calls++ + if node == "anchor" { + t.Fatal("anchor was composed although its number could not be taken") + } + return sendable{}, nil + }) + if calls != 1 || len(sending) != 1 || sending[0].node != "laptop" || sending[0].declared.Sequence != 7 { + t.Fatalf("laptop should be composed with its number and anchor refused: %v / %v", sending, refusals) + } + if len(refusals) != 1 { + t.Fatalf("anchor's failed number should be a refusal naming it: %v", refusals) + } +} + +// What was sent is written down even when the sender's context is already cancelled (issue 204): a +// controller replaced mid-send had told the machine and never recorded it, so status read "applied, +// current" over a machine that had just been sent something else. +func TestASendIsRecordedEvenWhenTheSenderIsBeingCancelled(t *testing.T) { + inv := inventory.ForTest(t) + ctx, cancel := context.WithCancel(t.Context()) + if _, err := inv.AddNode(ctx, "anchor"); err != nil { + t.Fatal(err) + } + cancel() // the sender is going away: its context is cancelled between the send and the record + body := []byte(`{"declaration":1,"resources":[]}`) + digest, err := recordSent(ctx, inv, "anchor", body) + if err != nil { + // NodeByName on the cancelled context may itself refuse; the record must still be possible + // through the detached context, so look the node up again on a live one. + t.Fatalf("recording a send after cancellation failed: %v", err) + } + outstanding, err := inv.Outstanding(t.Context(), "anchor") + if err != nil { + t.Fatal(err) + } + if outstanding != digest || digest != digestOf(body) { + t.Fatalf("the send was not recorded: outstanding %q, sent %q", outstanding, digest) + } +} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 2a1ba11..001f235 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -205,6 +205,12 @@ func declare(ctx context.Context, args []string) error { if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil { return err } + // Written down like every other send (novox/hq issue 204): a declaration a person sent by hand + // is still what the machine was last told, and status must not read it as current for the one + // the mesh would compose. + if _, err := recordSent(ctx, inv, node, raw); err != nil { + return err + } fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw)) return nil } @@ -348,7 +354,7 @@ func pushCommand(ctx context.Context, args []string) error { if err != nil { return err } - sending, refusals := composeEach(asked, func(node string) (sendable, error) { + sending, refusals := composeEach(asked, allotting(held, inv), func(node string) (sendable, error) { plan, settings, err := planFor(held, open, node) if err != nil { return sendable{}, err @@ -370,12 +376,8 @@ func pushCommand(ctx context.Context, args []string) error { sentDigest := map[string]string{} defer release() for _, s := range sending { - // Numbered under the hold, one higher than the last, before the body exists — the number is - // inside the signed bytes, so a replayed older declaration cannot borrow a newer one's - // (novox/hq 04-ISSUES/107). - if err := number(ctx, inv, &s); err != nil { - return err - } + // 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 @@ -385,14 +387,10 @@ func pushCommand(ctx context.Context, args []string) error { } // 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. - record, err := inv.NodeByName(ctx, s.node) + digest, err := recordSent(ctx, inv, s.node, body) if err != nil { return err } - digest := digestOf(body) - if err := inv.RecordSent(ctx, record.ID, digest); err != nil { - return err - } sentDigest[s.node] = digest fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) } @@ -476,11 +474,7 @@ func pushCommand(ctx context.Context, args []string) error { 15*time.Second); err != nil { return err } - record, err := inv.NodeByName(ctx, s.node) - if err != nil { - return err - } - if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil { + 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)) @@ -569,17 +563,31 @@ type readyNode struct { // // The all-or-nothing rule is kept where it means something — sendTo, which rotates a credential // across two machines that must agree — and dropped here, where it never did. -func composeEach(names []string, +func composeEach(names []string, allot func(node string) (int64, error), compose func(node string) (sendable, error)) ([]readyNode, []string) { var sending []readyNode var refusals []string for _, name := range names { + // **Numbered before it is composed, not before it is sent** (novox/hq issue 204). The + // number says where this declaration stands against every other the mesh composed for the + // machine, and the host refuses one lower than the last it applied. Taken at send time, as + // it was, a declaration composed a minute ago — before an assignment changed — went out with + // a number higher than one composed after the change and sent before it, and the machine + // took the older content as the newer word: on 2026-10-02 a runtime assigned and applied on + // two machines was undone two seconds later by exactly that. Taken here, before the first + // read, what was composed earlier is numbered lower whatever order the sends happen in. + seq, err := allot(name) + if err != nil { + refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err)) + continue + } declared, err := compose(name) if err != nil { refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err)) continue } + declared.Sequence = seq if len(declared.Resources) == 0 { // Sent, not skipped (novox/hq issue 127). A node whose declaration composes to // nothing may have HELD something before — the broker opening a placement gave it, @@ -607,13 +615,10 @@ func sendRound(ctx context.Context, open *stores, names []string, return nil, err } defer release() - sending, refused := composeEach(names, func(node string) (sendable, error) { + sending, refused := composeEach(names, allotting(held, open.inventory), func(node string) (sendable, error) { return compose(held, node) }) for _, s := range sending { - if err := number(ctx, open.inventory, &s); err != nil { - return refused, err - } body, err := s.declared.Body() if err != nil { return refused, err @@ -670,6 +675,12 @@ func sendTo(ctx context.Context, open *stores, names []string) error { var sending []readyNode var refusals []string for _, name := range names { + // Numbered before composing, for the reason composeEach gives (novox/hq issue 204). + seq, err := allot(ctx, inv, name) + if err != nil { + refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err)) + continue + } plan, settings, err := planFor(ctx, open, name) if err != nil { refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err)) @@ -681,6 +692,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error { refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err)) continue } + declared.Sequence = seq reportLeftOut(name, declared) sending = append(sending, readyNode{name, declared}) } @@ -696,9 +708,6 @@ func sendTo(ctx context.Context, open *stores, names []string) error { defer server.Close() for _, s := range sending { - if err := number(ctx, inv, &s); err != nil { - return err - } body, err := s.declared.Body() if err != nil { return err @@ -706,11 +715,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error { if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil { return err } - record, err := inv.NodeByName(ctx, s.node) - if err != nil { - return err - } - if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil { + 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)) @@ -951,15 +956,38 @@ func seatHolders(ctx context.Context, inv *inventory.Inventory) (map[string]brok } // number gives one send the next sequence for its node (novox/hq 04-ISSUES/107). -func number(ctx context.Context, inv *inventory.Inventory, s *readyNode) error { - record, err := inv.NodeByName(ctx, s.node) +// allotting is allot over one inventory, in the shape composeEach takes. +func allotting(ctx context.Context, inv *inventory.Inventory) func(node string) (int64, error) { + return func(node string) (int64, error) { return allot(ctx, inv, node) } +} + +// allot takes the next sequence for a machine — the number its next declaration carries. +func allot(ctx context.Context, inv *inventory.Inventory, node string) (int64, error) { + record, err := inv.NodeByName(ctx, node) if err != nil { - return err + return 0, err } - seq, err := inv.NextSequence(ctx, record.ID) + return inv.NextSequence(ctx, record.ID) +} + +// recordSent writes down what a machine was just sent, and returns the digest. +// +// **On a context that outlives the caller's** (novox/hq issue 204). The record is written after the +// declaration is away, so a send that failed is never recorded as current — and a controller being +// replaced mid-send had its context cancelled between the two, so the machine was told and the mesh +// never wrote it down: status read "applied, current" over a machine that had just been sent +// something else. What was sent was sent; the record of it must not depend on the sender living +// another second. Bounded, so a store that is away does not hold a dying process open for ever. +func recordSent(ctx context.Context, inv *inventory.Inventory, node string, body []byte) (string, error) { + kept, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second) + defer cancel() + record, err := inv.NodeByName(kept, node) if err != nil { - return err + return "", err } - s.declared.Sequence = seq - return nil + digest := digestOf(body) + if err := inv.RecordSent(kept, record.ID, digest); err != nil { + return "", err + } + return digest, nil } diff --git a/cmd/mesh-controller/push_test.go b/cmd/mesh-controller/push_test.go index aeaabfa..ddaf362 100644 --- a/cmd/mesh-controller/push_test.go +++ b/cmd/mesh-controller/push_test.go @@ -17,7 +17,7 @@ import ( // the wrong machine no longer refuses the whole node), applied one level up. func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) { sending, refusals := composeEach( - []string{"anchor", "home-server", "laptop"}, + []string{"anchor", "home-server", "laptop"}, numbered(), func(node string) (sendable, error) { if node == "anchor" { return sendable{}, errors.New(`nothing provides "acme-ca", wanted by route-proxy`) @@ -43,7 +43,7 @@ func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) { // (novox/hq issue 127): it may have held something before, and only sending the empty // declaration tells it to drop what the mesh owned. It is never a refusal. func TestAnEmptyDeclarationIsSentSoTheNodeDropsWhatItHeld(t *testing.T) { - sending, refusals := composeEach([]string{"spare"}, + sending, refusals := composeEach([]string{"spare"}, numbered(), func(string) (sendable, error) { return sendable{}, nil }) if len(sending) != 1 || len(refusals) != 0 { t.Errorf("an empty declaration must be sent, not skipped or refused: %v / %v", sending, refusals) @@ -74,3 +74,9 @@ func TestASkippedMachineIsStillAnError(t *testing.T) { } } } + +// numbered is an allotter for tests: one higher per call, as the inventory's is per machine. +func numbered() func(string) (int64, error) { + var n int64 + return func(string) (int64, error) { n++; return n, nil } +}