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) } }