package announce import ( "encoding/json" "strings" "testing" "time" "github.com/nats-io/nats.go/micro" "github.com/novox/mesh-tools/node-tools/internal/bus" mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" ) func connect(t *testing.T, module, node string) *bus.Conn { t.Helper() c, err := bus.Connect(bus.Credential{URL: mt.URL(t), Module: module, Node: node}) if err != nil { t.Fatal(err) } c.Logf = func(string, ...any) {} t.Cleanup(c.Close) return c } // lagging answers PING at once and INFO after `lag` — or never, with a negative lag: a runtime whose // one large answer is slow to cross the bus, as the laptop's was (2026-10-05). func lagging(t *testing.T, conn *bus.Conn, s Service, lag time.Duration, endpoints []Endpoint) { t.Helper() ping, _ := json.Marshal(micro.Ping{ServiceIdentity: s.identity(), Type: micro.PingResponseType}) stop, err := conn.Raw("$SRV.PING", func(string, []byte) []byte { return ping }) if err != nil { t.Fatal(err) } t.Cleanup(stop) stop, err = conn.Raw("$SRV.INFO", func(string, []byte) []byte { if lag < 0 { return nil } time.Sleep(lag) body, _ := json.Marshal(Info(s, endpoints)) return body }) if err != nil { t.Fatal(err) } t.Cleanup(stop) conn.Flush() } // A runtime whose answer arrives after the window is waited for, because it said it was there; one // that never says what it serves is named, not dropped. With the fixed 750 ms window the laptop's // modules were said to run nowhere whenever its answer was late (2026-10-05). func TestDiscoveryWaitsForEveryRuntimeThatSaidItIsThere(t *testing.T) { was := Patience Patience = 2500 * time.Millisecond t.Cleanup(func() { Patience = was }) fast := connect(t, "node-tools", "desk") stopFast, err := Serve(fast, Service{Name: "node-tools", ID: "desk", Metadata: map[string]string{"node": "desk"}}, func() []Endpoint { return []Endpoint{{Kind: KindTool, Module: "alpha", Tool: "one", Node: "desk", Subject: "mesh.mod.alpha.tool.one.desk"}} }) if err != nil { t.Fatal(err) } t.Cleanup(stopFast) fast.Flush() slow := connect(t, "node-tools", "laptop") lagging(t, slow, Service{Name: "node-tools", ID: "laptop", Metadata: map[string]string{"node": "laptop"}}, Window+400*time.Millisecond, []Endpoint{{Kind: KindTool, Module: "slack", Tool: "slack_check", Node: "laptop", Subject: "mesh.mod.slack.tool.slack_check.laptop"}}) mute := connect(t, "node-tools", "sleeper") lagging(t, mute, Service{Name: "node-tools", ID: "sleeper", Metadata: map[string]string{"node": "sleeper"}}, -1, nil) asker := connect(t, "console", "desk") start := time.Now() d, err := Gather(asker) if err != nil { t.Fatal(err) } took := time.Since(start) heard := map[string]Heard{} for _, h := range d.Heard { heard[h.Info.ID] = h } if h, ok := heard["laptop"]; !ok { t.Fatalf("the slow runtime was not heard: %+v", d.Heard) } else if h.Took <= Window || h.Machine() != "laptop" || h.Bytes == 0 { t.Errorf("the slow runtime's answer: took %s (window %s), machine %q, %d bytes", h.Took, Window, h.Machine(), h.Bytes) } if _, ok := heard["desk"]; !ok { t.Errorf("the fast runtime was not heard: %+v", d.Heard) } if len(d.Silent) != 1 || d.Silent[0].ID != "sleeper" || d.Silent[0].Machine != "sleeper" { t.Errorf("the runtime that never said what it serves is not named: %+v", d.Silent) } if took < Patience || took > Patience+time.Second { t.Errorf("waited %s with a silent runtime; patience is %s", took, Patience) } } // With nobody slow, discovery takes the window and no longer. func TestDiscoveryDoesNotWaitWhenEveryoneAnswered(t *testing.T) { fast := connect(t, "node-tools", "desk") stop, err := Serve(fast, Service{Name: "node-tools", ID: "desk"}, func() []Endpoint { return nil }) if err != nil { t.Fatal(err) } t.Cleanup(stop) fast.Flush() start := time.Now() d, err := Gather(connect(t, "console", "desk")) if err != nil { t.Fatal(err) } if took := time.Since(start); took > Window+500*time.Millisecond || len(d.Silent) != 0 || len(d.Heard) == 0 { t.Errorf("took %s, heard %d, silent %+v", took, len(d.Heard), d.Silent) } } // Too large even with first lines, descriptions are left out — and the answer says which, without // touching the service's own metadata. func TestAnAnswerTooLargeIsShortenedUntilItFits(t *testing.T) { long := strings.Repeat("A sentence that is long. ", 20) var endpoints []Endpoint for i := 0; i < 50; i++ { endpoints = append(endpoints, Endpoint{Kind: KindTool, Module: "m", Tool: "t" + strings.Repeat("x", i), Node: "laptop", Description: long, Subject: "mesh.mod.m.tool.t"}) } own := map[string]string{"node": "laptop"} info := Info(Service{Name: "node-tools", ID: "laptop", Metadata: own}, endpoints) full, _ := json.Marshal(info) firstLines := 0 for range endpoints { firstLines += len(firstLine(long)) } limit := int64(len(full) - len(long)*len(endpoints) + firstLines/2) // first lines alone do not fit var said []string got := shorten(func(f string, a ...any) { said = append(said, f) }, info, full, limit) if int64(len(got)) > limit { t.Fatalf("still %d bytes, beyond %d", len(got), limit) } var back micro.Info if err := json.Unmarshal(got, &back); err != nil { t.Fatal(err) } if back.Metadata[Shortened] != "descriptions left out" || back.Metadata["node"] != "laptop" { t.Errorf("metadata: %v", back.Metadata) } if _, has := own[Shortened]; has { t.Error("the service's own metadata was changed") } if len(said) != 2 { t.Errorf("each step is said: %v", said) } }