Files
mesh-host/internal/link/hearing_nats_test.go
T
jochen 31804bd8c2 Apply through one queue, and order what is applied and reported (hq to-be 45 Phase 2)
A delivery and the five-minute reconcile were two paths that applied, ordered only by a lock, and
each order it allowed was met live (issues 257, 261, 267). Now both only enqueue: one worker takes
the newest declaration held when it starts, applies it once and makes one report, and reports leave
in the order they are made.

A declaration may carry the controller's lease epoch beside its sequence; one older than what this
node applied is refused before anything is touched, counted, logged and reported. A report carries
the declaration's epoch and sequence and the host's own report sequence, kept on disk so it goes on
increasing across restarts and self-updates. Without an epoch, today's behaviour stands.
2026-10-06 11:54:10 +02:00

343 lines
11 KiB
Go

package link
import (
"context"
"crypto/ed25519"
"encoding/json"
"fmt"
"os"
"reflect"
"testing"
"time"
"github.com/nats-io/nats.go"
)
// The host's link against a real server, because every claim here is about one.
//
// Whether binding to a consumer the host did not create works, whether a declaration on the node's
// own subject arrives, whether acknowledging it removes it from the consumer's pending — none of
// that can be reasoned out, and the first two are the ones that would leave a node silently hearing
// nothing:
//
// docker run -d --rm --name t -p 14223:4222 nats:2.10-alpine -js
// MESH_TEST_NATS=nats://127.0.0.1:14223 go test ./internal/link/ -run TestNats
func aBus(t *testing.T) (*nats.Conn, nats.JetStreamContext) {
t.Helper()
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
conn, err := nats.Connect(url)
if err != nil {
t.Fatal(err)
}
t.Cleanup(conn.Close)
js, err := conn.JetStream()
if err != nil {
t.Fatal(err)
}
// **Ensured and purged, not deleted and recreated.** Delete-then-add looked like a reset and is
// not one: a test that did that inherited the previous test's messages, and the symptom was a
// declaration counted as delivered twice — which reads as a redelivery bug in the code under
// test rather than as a dirty stream. Purge is defined to empty a stream; recreating one is a
// race with the server's own teardown.
for _, want := range []*nats.StreamConfig{
{Name: "NODES", Subjects: []string{"mesh.node.*.declare"}, MaxMsgsPerSubject: 1},
{Name: "CONTROL", Subjects: []string{"mesh.control.*.report", "mesh.control.enrol"},
Retention: nats.WorkQueuePolicy},
} {
if _, err := js.StreamInfo(want.Name); err != nil {
if _, err := js.AddStream(want); err != nil {
t.Fatal(err)
}
}
if err := js.PurgeStream(want.Name); err != nil {
t.Fatal(err)
}
}
return conn, js
}
// theMeshMakes is the consumer the controller creates when a node enrols. Made here by the test
// because the host may not: its account reaches no part of the JetStream API, which is the whole
// reason this binds rather than subscribes.
//
// Removed afterwards, and each test names its own node: two tests sharing a consumer name share its
// delivery count and its pending list, and the first thing that goes wrong reads as a fault in the
// host rather than in the test beside it.
func theMeshMakes(t *testing.T, js nats.JetStreamContext, node string) {
t.Helper()
t.Cleanup(func() { _ = js.DeleteConsumer("NODES", node) })
if _, err := js.AddConsumer("NODES", &nats.ConsumerConfig{
Durable: node,
FilterSubject: DeclareSubject(node),
AckPolicy: nats.AckExplicitPolicy,
AckWait: 300 * time.Second,
DeliverSubject: "_DELIVER." + node,
}); err != nil {
t.Fatal(err)
}
}
func signedBy(t *testing.T, key ed25519.PrivateKey, declaration []byte) []byte {
t.Helper()
body, err := json.Marshal(Signed{
Declaration: declaration, Signature: ed25519.Sign(key, declaration),
})
if err != nil {
t.Fatal(err)
}
return body
}
// A declaration on this node's own subject reaches the host, is applied, and acknowledging it
// empties the consumer — which is what tells the mesh the node has it.
func TestNatsADeclarationReachesTheHostAndIsSettled(t *testing.T) {
conn, js := aBus(t)
const node = "settling"
theMeshMakes(t, js, node)
public, private, _ := ed25519.GenerateKey(nil)
m := Membership{Node: node, Signer: public}
// Dialled directly rather than through Open: the test server has no TLS, and what is being
// checked is the subscription and the settling, not the pin — which PinnedConfig owns and its
// own tests cover.
l := &natsLink{conn: conn, js: js, node: node,
arrived: make(chan Declaration, drainDepth), lost: make(chan error, 1)}
feed := make(chan *nats.Msg, drainDepth)
sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node))
if err != nil {
t.Fatalf("the host could not bind to the consumer the mesh made for it: %v", err)
}
defer func() { _ = sub.Unsubscribe() }()
go func() {
for msg := range feed {
l.arrived <- natsDeclaration{msg}
}
}()
if _, err := js.Publish(DeclareSubject(node),
signedBy(t, private, []byte(`{"declared":"d1"}`))); err != nil {
t.Fatal(err)
}
select {
case d := <-l.Declarations():
report := handleBody(context.Background(), m, d.Body(),
func(context.Context, []byte, []byte) Report {
return Report{Applied: []string{"store"}}
})
if report.Refused != "" {
t.Fatalf("a declaration the mesh signed was refused: %s", report.Refused)
}
if err := d.Handled(); err != nil {
t.Fatalf("the node could not acknowledge its own declaration: %v", err)
}
case <-time.After(8 * time.Second):
t.Fatal("no declaration reached the host")
}
// **Nothing pending is the property**; a delivery count is not. Delivery is at-least-once by
// design, so pinning "delivered exactly once" would be asserting something the mesh does not
// rely on. What matters is that the acknowledgement landed, so the mesh can tell the node has
// it — and that no redelivery was needed to get there, which is what would say the node was
// too slow to answer for its own ack wait.
deadline := time.Now().Add(5 * time.Second)
var last string
for time.Now().Before(deadline) {
info, err := js.ConsumerInfo("NODES", node)
switch {
case err != nil:
last = err.Error()
case info.NumAckPending == 0 && info.NumRedelivered == 0:
return
default:
last = fmt.Sprintf("pending %d, redelivered %d", info.NumAckPending, info.NumRedelivered)
}
time.Sleep(20 * time.Millisecond)
}
t.Fatalf("the declaration was not settled, so the mesh cannot tell the node has it: %s", last)
}
// **A node that was away gets exactly the current declaration and nothing older.** Three pushed
// while nothing is listening leave one on the stream, and it is the newest — the wire-level answer
// to novox/hq issue 107, and the half of the drain that stops being the host's problem.
func TestNatsANodeThatWasAwayGetsOnlyTheNewest(t *testing.T) {
_, js := aBus(t)
const node = "returning"
_, private, _ := ed25519.GenerateKey(nil)
for _, id := range []string{"d1", "d2", "d3"} {
if _, err := js.Publish(DeclareSubject(node),
signedBy(t, private, []byte(`{"declared":"`+id+`"}`))); err != nil {
t.Fatal(err)
}
}
info, err := js.StreamInfo("NODES")
if err != nil {
t.Fatal(err)
}
if info.State.Msgs != 1 {
t.Fatalf("%d declarations survived for one node; a node that was away would apply a backlog "+
"of things nobody wants any more", info.State.Msgs)
}
theMeshMakes(t, js, node)
feed := make(chan *nats.Msg, drainDepth)
sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node))
if err != nil {
t.Fatal(err)
}
defer func() { _ = sub.Unsubscribe() }()
select {
case msg := <-feed:
if want := declaredIn(signedBy(t, private, []byte(`{"declared":"d3"}`))); declaredIn(msg.Data) != want {
t.Fatalf("the node was given %q rather than the newest", msg.Data)
}
case <-time.After(8 * time.Second):
t.Fatal("the node that was away was given nothing")
}
}
// A report goes through the stream and a heartbeat does not: the one that must survive the
// controller's store restarting is kept, and the one that must not is not.
func TestNatsAReportIsKeptAndAHeartbeatIsNot(t *testing.T) {
conn, js := aBus(t)
bus := OverNATS{Conn: conn, JS: js}
ctx := context.Background()
body, _ := json.Marshal(Report{Node: "anchor", Declared: "d1"})
if err := bus.Report(ctx, "anchor", body); err != nil {
t.Fatal(err)
}
beat, _ := json.Marshal(Alive{Node: "anchor"})
if err := bus.Alive(ctx, "anchor", beat); err != nil {
t.Fatal(err)
}
info, err := js.StreamInfo("CONTROL")
if err != nil {
t.Fatal(err)
}
if info.State.Msgs != 1 {
t.Fatalf("%d messages were kept; a report must be and a heartbeat must not", info.State.Msgs)
}
}
// **One apply queue, against a real bus** (novox/hq to-be 45 §6, R2). A declaration is being applied
// when two more are pushed and the five-minute reconcile comes due: the worker applies the one in
// hand, then only the newest — the reconcile is that apply — and the mesh's stream holds one account
// per apply, each naming the order it applied, numbered in the order they were made. Every
// declaration delivered is settled.
func TestNatsTheQueueAppliesTheNewestOnceAndReportsInOrder(t *testing.T) {
conn, js := aBus(t)
const node = "queueing"
theMeshMakes(t, js, node)
public, private, _ := ed25519.GenerateKey(nil)
m := Membership{Node: node, Signer: public}
ctx, stop := context.WithCancel(context.Background())
defer stop()
l := &natsLink{conn: conn, js: js, node: node,
arrived: make(chan Declaration, drainDepth), lost: make(chan error, 1)}
feed := make(chan *nats.Msg, drainDepth)
sub, err := js.ChanSubscribe(DeclareSubject(node), feed, nats.Bind("NODES", node))
if err != nil {
t.Fatal(err)
}
defer func() { _ = sub.Unsubscribe() }()
go func() {
for msg := range feed {
l.arrived <- natsDeclaration{msg}
}
}()
reports, err := js.SubscribeSync(ReportSubject(node), nats.BindStream("CONTROL"))
if err != nil {
t.Fatal(err)
}
defer func() { _ = reports.Unsubscribe() }()
a := &applies{}
inFirst, release := make(chan struct{}, 1), make(chan struct{})
reconciled := 0
q := &Queue{Membership: m, Timeout: 5 * time.Second, Unsaid: &keptInMemory{},
Apply: func(ctx context.Context, raw, sig []byte) Report {
r := a.apply(ctx, raw, sig)
if r.Sequence == 1 {
inFirst <- struct{}{}
<-release
}
return r
},
Reconcile: func(context.Context) (Report, bool) { reconciled++; return Report{}, false },
}
worker := runWorker(ctx, q)
go func() { _ = serve(ctx, l, m, q, nil, 5*time.Second) }()
push := func(id string, sequence int64) {
inner, _ := json.Marshal(map[string]any{"declaration": 1, "epoch": 7, "sequence": sequence,
"resources": []any{map[string]any{"id": id}}})
if _, err := js.Publish(DeclareSubject(node), signedBy(t, private, inner)); err != nil {
t.Fatal(err)
}
}
push("one", 1)
select {
case <-inFirst:
case <-time.After(8 * time.Second):
t.Fatal("the first declaration never reached the worker")
}
push("two", 2)
push("three", 3)
q.ReconcileDue()
time.Sleep(2 * drainWindow) // the link gathers both and hands them to the queue
close(release)
var heard []Report
for len(accounts(heard)) < 2 {
msg, err := reports.NextMsg(8 * time.Second)
if err != nil {
t.Fatalf("the mesh heard %+v and then nothing: %v", heard, err)
}
var r Report
if err := json.Unmarshal(msg.Data, &r); err != nil {
t.Fatal(err)
}
heard = append(heard, r)
_ = msg.Ack()
}
stop()
<-worker
if got := a.all(); !reflect.DeepEqual(got, []string{"one", "three"}) {
t.Fatalf("applied %v; the one in hand and then only the newest were wanted", got)
}
if reconciled != 0 {
t.Fatal("the reconcile applied beside the delivery it was due with")
}
accounted := accounts(heard)
if accounted[0].Sequence != 1 || accounted[1].Sequence != 3 || accounted[1].Epoch != 7 {
t.Fatalf("the accounts do not name what was applied: %+v", accounted)
}
for i := 1; i < len(heard); i++ {
if heard[i].ReportSequence <= heard[i-1].ReportSequence {
t.Fatalf("the reports left out of the order they were made: %+v", heard)
}
}
deadline := time.Now().Add(5 * time.Second)
for {
info, err := js.ConsumerInfo("NODES", node)
if err == nil && info.NumAckPending == 0 {
break
}
if time.Now().After(deadline) {
t.Fatalf("a delivered declaration was left unsettled: %+v %v", info, err)
}
time.Sleep(20 * time.Millisecond)
}
}