The live tests reached one shared bus and assert, read and remove the mesh's own objects by their fixed names, so packages run in parallel deleted what each other read and the suite passed only one package at a time; a red suite read as noise. internal/testbus starts a server per test, linked in at the nats-server release go.mod pins, and a test holds that pin to the catalogue's bus image and to the facts snapshot's bus when there is one, so the tests never run a bus the mesh does not. The waiter test read a timing (the most connections held at one look) and now reads the state it means (the fewest held across the wait). make check runs the packages in parallel under the race detector, with a timeout.
79 lines
2.8 KiB
Go
79 lines
2.8 KiB
Go
package link
|
|
|
|
import (
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/testbus"
|
|
)
|
|
|
|
// **Verified 2026-09-27**: the caller asked for a reply to `_INBOX.LCr3M83q…` and the consumer
|
|
// saw `$JS.ACK.PROBE.probe_consumer.1.1.1…`. The design was right, and enrolment's payload-borne
|
|
// reply subject is necessary rather than defensive.
|
|
//
|
|
// Design 25 §2 asserts that a reply address is **eaten** when a message crosses a stream: core
|
|
// request/reply puts the requester's inbox in the message's Reply field, but a JetStream consumer
|
|
// has already claimed that field for its own ack address by the time a handler sees it. The whole
|
|
// enrolment design rests on it — the reply subject travels in the payload instead — so it is
|
|
// checked rather than believed.
|
|
//
|
|
// go test ./internal/link/ -run TestAReplyAddress (each test on a bus of its own: internal/testbus)
|
|
func TestAReplyAddressDoesNotSurviveAStream(t *testing.T) {
|
|
url := testbus.URL(t)
|
|
conn, err := nats.Connect(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer conn.Close()
|
|
js, err := conn.JetStream()
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := js.AddStream(&nats.StreamConfig{
|
|
Name: "PROBE", Subjects: []string{"probe.>"}, Retention: nats.WorkQueuePolicy,
|
|
}); err != nil && err != nats.ErrStreamNameAlreadyInUse {
|
|
t.Fatal(err)
|
|
}
|
|
defer js.DeleteStream("PROBE")
|
|
|
|
seen := make(chan *nats.Msg, 1)
|
|
sub, err := js.Subscribe("probe.enrol", func(m *nats.Msg) { seen <- m },
|
|
nats.Durable("probe_consumer"), nats.ManualAck())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer sub.Unsubscribe()
|
|
|
|
// A caller doing what core request/reply does: publish with its own inbox as the reply.
|
|
inbox := nats.NewInbox()
|
|
if err := conn.PublishMsg(&nats.Msg{Subject: "probe.enrol", Reply: inbox, Data: []byte("{}")}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
select {
|
|
case m := <-seen:
|
|
t.Logf("the caller asked for a reply to %s", inbox)
|
|
t.Logf("the consumer sees a Reply field of %s", m.Reply)
|
|
if m.Reply == inbox {
|
|
t.Fatalf("the reply address SURVIVED the stream. Design 25 §2 says it does not, and " +
|
|
"builds enrolment around carrying the reply subject in the payload to work " +
|
|
"around it. If this holds generally, that work is unnecessary and the design " +
|
|
"should say so.")
|
|
}
|
|
if m.Reply == "" {
|
|
t.Fatal("the Reply field is empty rather than claimed; the design says it carries " +
|
|
"the consumer's ack address, which is a different fact")
|
|
}
|
|
// Answering it would send the enrolling node's credentials to an ack subject.
|
|
if len(m.Reply) < 7 || m.Reply[:7] != "$JS.ACK" {
|
|
t.Errorf("the Reply field is %q, which is neither the caller's inbox nor an ack "+
|
|
"address — the design's reasoning assumes one of the two", m.Reply)
|
|
}
|
|
m.Ack()
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("nothing was delivered")
|
|
}
|
|
}
|