Merge pull request 'A declaration is numbered when it is composed, and a send is recorded even by a sender being replaced (hq issue 204)' (#232) from fix/issue-204 into main

This commit was merged in pull request #232.
This commit is contained in:
2026-10-03 09:44:42 +00:00
3 changed files with 166 additions and 39 deletions
+93
View File
@@ -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)
}
}
+65 -37
View File
@@ -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 { if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil {
return err 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)) fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw))
return nil return nil
} }
@@ -348,7 +354,7 @@ func pushCommand(ctx context.Context, args []string) error {
if err != nil { if err != nil {
return err 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) plan, settings, err := planFor(held, open, node)
if err != nil { if err != nil {
return sendable{}, err return sendable{}, err
@@ -370,12 +376,8 @@ func pushCommand(ctx context.Context, args []string) error {
sentDigest := map[string]string{} sentDigest := map[string]string{}
defer release() defer release()
for _, s := range sending { for _, s := range sending {
// Numbered under the hold, one higher than the last, before the body exists — the number is // The number is inside the signed bytes, so a replayed older declaration cannot borrow a
// inside the signed bytes, so a replayed older declaration cannot borrow a newer one's // newer one's (novox/hq 04-ISSUES/107); it was taken when the composition began (issue 204).
// (novox/hq 04-ISSUES/107).
if err := number(ctx, inv, &s); err != nil {
return err
}
body, err := s.declared.Body() body, err := s.declared.Body()
if err != nil { if err != nil {
return err 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 // 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. // 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 { if err != nil {
return err return err
} }
digest := digestOf(body)
if err := inv.RecordSent(ctx, record.ID, digest); err != nil {
return err
}
sentDigest[s.node] = digest sentDigest[s.node] = digest
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) 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 { 15*time.Second); err != nil {
return err return err
} }
record, err := inv.NodeByName(ctx, s.node) if _, err := recordSent(ctx, inv, s.node, body); err != nil {
if err != nil {
return err
}
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
return err return err
} }
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) 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 // 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. // 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) { compose func(node string) (sendable, error)) ([]readyNode, []string) {
var sending []readyNode var sending []readyNode
var refusals []string var refusals []string
for _, name := range names { 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) declared, err := compose(name)
if err != nil { if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err)) refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
continue continue
} }
declared.Sequence = seq
if len(declared.Resources) == 0 { if len(declared.Resources) == 0 {
// Sent, not skipped (novox/hq issue 127). A node whose declaration composes to // 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, // 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 return nil, err
} }
defer release() 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) return compose(held, node)
}) })
for _, s := range sending { for _, s := range sending {
if err := number(ctx, open.inventory, &s); err != nil {
return refused, err
}
body, err := s.declared.Body() body, err := s.declared.Body()
if err != nil { if err != nil {
return refused, err return refused, err
@@ -670,6 +675,12 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
var sending []readyNode var sending []readyNode
var refusals []string var refusals []string
for _, name := range names { 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) plan, settings, err := planFor(ctx, open, name)
if err != nil { if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err)) 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)) refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
continue continue
} }
declared.Sequence = seq
reportLeftOut(name, declared) reportLeftOut(name, declared)
sending = append(sending, readyNode{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() defer server.Close()
for _, s := range sending { for _, s := range sending {
if err := number(ctx, inv, &s); err != nil {
return err
}
body, err := s.declared.Body() body, err := s.declared.Body()
if err != nil { if err != nil {
return err 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 { if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil {
return err return err
} }
record, err := inv.NodeByName(ctx, s.node) if _, err := recordSent(ctx, inv, s.node, body); err != nil {
if err != nil {
return err
}
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
return err return err
} }
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) 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). // 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 { // allotting is allot over one inventory, in the shape composeEach takes.
record, err := inv.NodeByName(ctx, s.node) 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 { 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 { if err != nil {
return err return "", err
} }
s.declared.Sequence = seq digest := digestOf(body)
return nil if err := inv.RecordSent(kept, record.ID, digest); err != nil {
return "", err
}
return digest, nil
} }
+8 -2
View File
@@ -17,7 +17,7 @@ import (
// the wrong machine no longer refuses the whole node), applied one level up. // the wrong machine no longer refuses the whole node), applied one level up.
func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) { func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) {
sending, refusals := composeEach( sending, refusals := composeEach(
[]string{"anchor", "home-server", "laptop"}, []string{"anchor", "home-server", "laptop"}, numbered(),
func(node string) (sendable, error) { func(node string) (sendable, error) {
if node == "anchor" { if node == "anchor" {
return sendable{}, errors.New(`nothing provides "acme-ca", wanted by route-proxy`) 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 // (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. // declaration tells it to drop what the mesh owned. It is never a refusal.
func TestAnEmptyDeclarationIsSentSoTheNodeDropsWhatItHeld(t *testing.T) { func TestAnEmptyDeclarationIsSentSoTheNodeDropsWhatItHeld(t *testing.T) {
sending, refusals := composeEach([]string{"spare"}, sending, refusals := composeEach([]string{"spare"}, numbered(),
func(string) (sendable, error) { return sendable{}, nil }) func(string) (sendable, error) { return sendable{}, nil })
if len(sending) != 1 || len(refusals) != 0 { if len(sending) != 1 || len(refusals) != 0 {
t.Errorf("an empty declaration must be sent, not skipped or refused: %v / %v", sending, refusals) 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 }
}