Step 3.4. On the new bus there is no reply queue to declare and no correlation to check: each account is granted one inbox prefix and no other, so an answer cannot reach the wrong asker. That settles a cost build.go records having paid — on a shared reply exchange every asker saw every result, which is why the correlation was checked rather than assumed. And a tool nobody serves says so at once rather than after the whole wait. The difference between "that module is down" and "that tool is slow" is the first thing a person asking wants, and both tests are against a real server because both are claims about what the server does, not about this code. RequestBuild stays as it is, and is a different shape on the new bus rather than the same one: a build takes minutes, so it is work submitted to a queue with the outcome returning to a reply subject the request carries — the pattern design 25 §2 already sets for anything crossing a stream. It touches the builder too, so it goes with that conversion.
155 lines
4.8 KiB
Go
155 lines
4.8 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"os"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// **Both implementations, one fixture.** The point of the seam is not that the transport can be
|
|
// swapped — it is that the two can be held to the same envelope while both ship, so the day the
|
|
// bus moves is a configuration change rather than a discovery.
|
|
//
|
|
// Against a real server, because what the fixture pins is what reaches the wire:
|
|
//
|
|
// docker run -d --rm --name t -p 14222:4222 nats:2.10-alpine -js
|
|
// MESH_TEST_NATS=nats://127.0.0.1:14222 go test ./internal/link/ -run TestTheNatsBus
|
|
func TestTheNatsBusEmitsTheEnvelopeTheFixturePins(t *testing.T) {
|
|
url := os.Getenv("MESH_TEST_NATS")
|
|
if url == "" {
|
|
t.Skip("MESH_TEST_NATS unset")
|
|
}
|
|
f := loadFixture(t, "events/module-event.json")
|
|
|
|
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: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"},
|
|
}); err != nil && err != nats.ErrStreamNameAlreadyInUse {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
bus := OverNATS{JS: js}
|
|
body, _ := json.Marshal(f.Given.Body)
|
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
if err := EmitEvent(ctx, bus, f.Given.Key, f.Given.Module, f.Given.Node, f.Given.Body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
// Read back from the stream, not from the thing that wrote it.
|
|
raw, err := js.GetLastMsg("EVENTS", f.Wire.Subject)
|
|
if err != nil {
|
|
t.Fatalf("nothing landed on %s, which the fixture names: %v", f.Wire.Subject, err)
|
|
}
|
|
for _, h := range f.Wire.RequiredHeaders {
|
|
if raw.Header.Get(h) == "" {
|
|
t.Errorf("%s is not set on the wire, and the fixture requires it", h)
|
|
}
|
|
}
|
|
if got := raw.Header.Get("x-source"); got != f.Given.Module {
|
|
t.Errorf("x-source is %q; the subject says %q", got, f.Given.Module)
|
|
}
|
|
if string(raw.Data) != string(body) {
|
|
t.Errorf("the payload is %s, expected the body alone: %s", raw.Data, body)
|
|
}
|
|
// The envelope must not also be nested inside the payload.
|
|
var nested map[string]any
|
|
if json.Unmarshal(raw.Data, &nested) == nil {
|
|
if _, has := nested["key"]; has {
|
|
t.Error("the payload carries the envelope's own fields, which the fixture refuses")
|
|
}
|
|
}
|
|
}
|
|
|
|
// The subject a declaration lands on is one node's, and nothing else's — the state shape.
|
|
func TestADeclarationIsAddressedToOneNode(t *testing.T) {
|
|
if got := DeclareSubject("anchor"); got != "mesh.node.anchor.declare" {
|
|
t.Fatalf("a declaration would go to %q", got)
|
|
}
|
|
if DeclareSubject("anchor") == DeclareSubject("laptop") {
|
|
t.Fatal("two nodes share a declaration subject, so each would apply the other's")
|
|
}
|
|
}
|
|
|
|
// A module cannot emit under another's name: the subject is derived from the source, and the
|
|
// server's permissions make that subject the authority.
|
|
func TestAnEventsSubjectIsDerivedFromItsSource(t *testing.T) {
|
|
if got := EventSubject("shop", "order.placed"); got != "mesh.mod.shop.event.order.placed" {
|
|
t.Fatalf("an event would land on %q", got)
|
|
}
|
|
if EventSubject("shop", "x") == EventSubject("billing", "x") {
|
|
t.Fatal("two modules share an event subject, so neither owns its own name")
|
|
}
|
|
}
|
|
|
|
// A tool nobody serves says so at once. The difference between "that module is down" and "that
|
|
// tool is slow" is the first thing a person asking wants, and waiting out the timeout to say it
|
|
// is how a fast answer becomes a slow non-answer.
|
|
func TestAskingAToolNobodyServesFailsAtOnce(t *testing.T) {
|
|
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)
|
|
}
|
|
defer conn.Close()
|
|
|
|
bus := OverNATS{Conn: conn}
|
|
began := time.Now()
|
|
_, err = bus.AskTool(context.Background(), "nobody", "home", nil, 30*time.Second)
|
|
if err == nil {
|
|
t.Fatal("asking a tool nothing serves succeeded")
|
|
}
|
|
if took := time.Since(began); took > 2*time.Second {
|
|
t.Errorf("took %v to say nothing serves it; the caller waited out the timeout", took)
|
|
}
|
|
if !strings.Contains(err.Error(), "nothing serves") {
|
|
t.Errorf("the refusal does not say nobody is there: %v", err)
|
|
}
|
|
}
|
|
|
|
// And a served tool answers, with no reply queue to declare and no correlation to check.
|
|
func TestAskingAServedToolAnswers(t *testing.T) {
|
|
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)
|
|
}
|
|
defer conn.Close()
|
|
|
|
sub, err := conn.Subscribe(ToolSubject("shop", "price"), func(m *nats.Msg) {
|
|
_ = m.Respond([]byte(`{"total":12}`))
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer sub.Unsubscribe()
|
|
|
|
got, err := OverNATS{Conn: conn}.AskTool(context.Background(), "shop", "price",
|
|
[]byte(`{"qty":4}`), 5*time.Second)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if string(got) != `{"total":12}` {
|
|
t.Fatalf("the answer came back as %s", got)
|
|
}
|
|
}
|