diff --git a/README.md b/README.md index 63c63cf..b6fc533 100644 --- a/README.md +++ b/README.md @@ -61,10 +61,22 @@ The same package carries the client (novox/hq design 25 §7, design 34): `mesh t refuses to bind anything but loopback. With `--console `, `tools` and `call` go through a console already on the machine and need no credential. -Discovery asks the modules: every runtime answers a `tools` verb for each module it serves, with names, -descriptions and schemas from the code that answers them, and the console asks the catalogue which -modules the mesh holds and each module what it serves. A module that does not answer is named, never -dropped. A module may not name a tool of its own `tools`; the runtime refuses it at load. +The console announces six tools and reaches everything else by address (novox/hq ADR 0195): +`mesh_overview` (the mesh's seats and machines), `mesh_machine` (one machine's seats and modules), +`mesh_search` (a tool by words), `mesh_describe` (one tool's arguments), `mesh_call` (call one by +address) and `mesh_runtimes` (which runtimes answered discovery: per runtime its machine, how long its +answer took, its size in bytes, how many modules and tools it announced, whether it was shortened to +fit the bus, when it was last heard — and who was expected and not heard). + +Discovery asks the bus (ADR 0197): every runtime answers the NATS services protocol's `$SRV.PING` and +`$SRV.INFO` with what it serves at that moment, and the console reads where the controller's records +place each module. The console waits at least 750 ms, and up to 5 s for every runtime that answered +PING to send what it serves — a large runtime's answer crossing the bus to a broker on another machine +can arrive well after the rest. A runtime that said it is there and did not say what it serves, or +that answered earlier and not now, is named: the console never calls a module missing, or on another +machine, while a runtime that might serve it was not heard. A runtime whose answer outgrows the bus +announces first-line descriptions, then none, and says so. A module may not name a tool of its own +`tools`; the runtime refuses it at load. ## Verified diff --git a/node-tools/internal/announce/announce.go b/node-tools/internal/announce/announce.go index c0613b3..f013cad 100644 --- a/node-tools/internal/announce/announce.go +++ b/node-tools/internal/announce/announce.go @@ -22,10 +22,22 @@ import ( // Version is the announced service version (semver, as the protocol requires). const Version = "0.1.0" -// Window is how long the console gathers discovery answers: every instance answers one request, and -// how many will is what is being found out. +// Window is how long the console gathers discovery answers at the least: every instance answers one +// request, and how many will is what is being found out. var Window = 750 * time.Millisecond +// Patience is how long the console waits, at the most, for the full answer of a service that has +// said it is there. **A fixed window lost the largest runtime** (2026-10-05): the laptop's answer — +// 48 modules, 341 endpoints, 164 kB — went up to the broker on another machine and back, arrived +// last every time (a median of 365 ms, one in 25 after 813 ms on a quiet link, later still while the +// runtime re-served), and when it missed the 750 ms window the console said its modules ran nowhere. +// A PING answer is a hundred bytes and arrives at once, so who is there is known early; what each +// serves is waited for until it arrives or this passes, and a service that never sends it is named. +var Patience = 5 * time.Second + +// Shortened is the service metadata key a runtime sets when its announcement was cut to fit the bus. +const Shortened = "shortened" + // Kinds of endpoint. const ( KindTool = "tool" // a module's own tool @@ -163,17 +175,10 @@ func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error) return nil } // Too large for the bus is no answer at all, and was silent: shorten every description to its - // first line and say so, rather than vanish. + // first line and say so, rather than vanish; still too large, leave descriptions out. if limit := conn.MaxPayload(); limit > 0 && int64(len(body)) > limit { if info, ok := v.(micro.Info); ok { - for i := range info.Endpoints { - info.Endpoints[i].Metadata["description"] = firstLine(info.Endpoints[i].Metadata["description"]) - } - if shorter, err := json.Marshal(info); err == nil { - conn.Logf("[mesh-tools] what this runtime serves is %d bytes, beyond the bus's %d; announced with each description cut to its first line (%d bytes)", - len(body), limit, len(shorter)) - body = shorter - } + body = shorten(conn.Logf, info, body, limit) } } return body @@ -200,22 +205,149 @@ func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error) }, nil } -// Gather asks every service on the bus what it serves and answers what came back within the window. -// An answer that is not an info_response is skipped. -func Gather(conn *bus.Conn) ([]micro.Info, error) { - raw, err := conn.Gather("$SRV.INFO", nil, Window) - if err != nil { - return nil, err +// shorten is an info_response cut to fit the bus: descriptions to their first line, then none, the +// service's metadata saying which (so discovery can tell a person what it was given). +func shorten(logf func(string, ...any), info micro.Info, body []byte, limit int64) []byte { + md := map[string]string{} + for k, v := range info.Metadata { + md[k] = v } - var out []micro.Info - for _, b := range raw { - var i micro.Info - if json.Unmarshal(b, &i) != nil || i.Type != micro.InfoResponseType { - continue + info.Metadata = md + for _, step := range []string{"descriptions cut to their first line", "descriptions left out"} { + for i := range info.Endpoints { + e := info.Endpoints[i].Metadata + if step == "descriptions left out" { + delete(e, "description") + } else { + e["description"] = firstLine(e["description"]) + } + } + md[Shortened] = step + shorter, err := json.Marshal(info) + if err != nil { + return body + } + logf("[mesh-tools] what this runtime serves is %d bytes, beyond the bus's %d; announced with %s (%d bytes)", + len(body), limit, step, len(shorter)) + if int64(len(shorter)) <= limit { + return shorter + } + body = shorter + } + return body +} + +// Heard is one service instance's answer to discovery: what it serves, how long it took to arrive, +// and how large it was. +type Heard struct { + Info micro.Info + Took time.Duration + Bytes int + Shortened string // why its descriptions are not whole, or "" +} + +// Machine is the machine an instance runs on: its metadata's, or the one its endpoints name. +func (h Heard) Machine() string { + if n := h.Info.Metadata["node"]; n != "" { + return n + } + for _, e := range h.Info.Endpoints { + if n := e.Metadata["node"]; n != "" { + return n + } + } + return "" +} + +// Instance names one service instance. +type Instance struct { + Name string + ID string + Machine string +} + +// Discovery is what one round of discovery heard, and who said it was there but did not say what it +// serves in time. +type Discovery struct { + At time.Time + Waited time.Duration + Heard []Heard + Silent []Instance +} + +// Gather asks every service on the bus who it is and what it serves, and waits at least the window +// and at most the patience: until every instance that answered PING has answered INFO. An answer +// that is not the protocol's is skipped. +func Gather(conn *bus.Conn) (Discovery, error) { + start := time.Now() + d := Discovery{At: start.UTC()} + pings, stopPings, err := conn.Collect("$SRV.PING", nil) + if err != nil { + return d, err + } + defer stopPings() + infos, stopInfos, err := conn.Collect("$SRV.INFO", nil) + if err != nil { + return d, err + } + defer stopInfos() + + type key struct{ name, id string } + pinged := map[key]int{} + answered := map[key]int{} + where := map[key]string{} + var order []key + complete := func() bool { + for k, n := range pinged { + if answered[k] < n { + return false + } + } + return true + } + least := time.NewTimer(Window) + defer least.Stop() + most := time.NewTimer(Patience) + defer most.Stop() + windowOver := false + for { + select { + case a := <-pings: + var p micro.Ping + if json.Unmarshal(a.Data, &p) != nil || p.Type != micro.PingResponseType { + continue + } + k := key{p.Name, p.ID} + if _, seen := pinged[k]; !seen { + order = append(order, k) + } + pinged[k]++ + if where[k] == "" { + where[k] = p.Metadata["node"] + } + case a := <-infos: + var i micro.Info + if json.Unmarshal(a.Data, &i) != nil || i.Type != micro.InfoResponseType { + continue + } + answered[key{i.Name, i.ID}]++ + d.Heard = append(d.Heard, Heard{Info: i, Took: a.At.Sub(start), Bytes: len(a.Data), Shortened: i.Metadata[Shortened]}) + case <-least.C: + windowOver = true + case <-most.C: + d.Waited = time.Since(start) + for _, k := range order { + if answered[k] < pinged[k] { + d.Silent = append(d.Silent, Instance{Name: k.name, ID: k.id, Machine: where[k]}) + } + } + return d, nil + } + if windowOver && complete() { + d.Waited = time.Since(start) + return d, nil } - out = append(out, i) } - return out, nil } // Endpoints reads an info_response's endpoints back into what they announce. diff --git a/node-tools/internal/announce/gather_test.go b/node-tools/internal/announce/gather_test.go new file mode 100644 index 0000000..27d836f --- /dev/null +++ b/node-tools/internal/announce/gather_test.go @@ -0,0 +1,156 @@ +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) + } +} diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index dd3e650..425d96d 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -649,37 +649,51 @@ func (c *Conn) Raw(subject string, answer func(subject string, data []byte) []by return func() { _ = sub.Unsubscribe() }, nil } -// Gather publishes one request and collects every answer that arrives within the window: a -// discovery request every service instance answers (ADR 0197). It never stops early — how many will -// answer is what it is finding out. -func (c *Conn) Gather(subject string, body []byte, window time.Duration) ([][]byte, error) { +// Arrival is one answer to a request, and when it arrived. +type Arrival struct { + Data []byte + At time.Time +} + +// Collect publishes one request and hands every answer to the channel as it arrives, until stopped: +// for a caller that decides for itself when it has heard enough — discovery waits for every service +// that said it is there, not for a fixed moment (2026-10-05). Answers beyond what the caller has +// taken are kept up to the channel's room; an empty answer is no answer. +func (c *Conn) Collect(subject string, body []byte) (<-chan Arrival, func(), error) { inbox := c.nc.NewRespInbox() - sub, err := c.nc.SubscribeSync(inbox) + out := make(chan Arrival, 1024) + var mu sync.Mutex + stopped := false + sub, err := c.nc.Subscribe(inbox, func(msg *nats.Msg) { + if len(msg.Data) == 0 { + return + } + a := Arrival{Data: msg.Data, At: time.Now()} + mu.Lock() + defer mu.Unlock() + if stopped { + return + } + select { + case out <- a: + default: + c.Logf("[mesh-tools] more answers to %s than are kept; one (%d bytes) dropped", subject, len(msg.Data)) + } + }) if err != nil { - return nil, err + return nil, func() {}, err + } + stop := func() { + _ = sub.Unsubscribe() + mu.Lock() + stopped = true + mu.Unlock() } - defer func() { _ = sub.Unsubscribe() }() if err := c.nc.PublishRequest(subject, inbox, body); err != nil { - return nil, err - } - var out [][]byte - deadline := time.Now().Add(window) - for { - left := time.Until(deadline) - if left <= 0 { - return out, nil - } - msg, err := sub.NextMsg(left) - if err != nil { - if errors.Is(err, nats.ErrTimeout) { - return out, nil - } - return out, err - } - if len(msg.Data) > 0 { - out = append(out, msg.Data) - } + stop() + return nil, func() {}, err } + return out, stop, nil } // MaxPayload is the largest message the bus carries, as the server told this connection. diff --git a/node-tools/internal/console/address.go b/node-tools/internal/console/address.go index 2409e6a..2d53464 100644 --- a/node-tools/internal/console/address.go +++ b/node-tools/internal/console/address.go @@ -2,8 +2,8 @@ package console // The mesh's tools found by address, not announced whole (novox/hq ADR 0195, to-be 34 §3a). // -// The console announces five tools. Everything the mesh answers is reached through them by one -// address per layer: +// The console announces six tools — five to find and call, and mesh_runtimes to say which runtimes +// discovery heard. Everything the mesh answers is reached through them by one address per layer: // // . a seat held once for the mesh — its holder answers // /. a seat held once per machine — that machine's holder answers @@ -32,13 +32,14 @@ import ( // one agent turn's search, describe and call ask once, short enough that nothing goes stale. var IndexKept = 5 * time.Second -// The five tools the console announces. Names of the API's kind — letters, digits, `_`, `-`. +// The six tools the console announces. Names of the API's kind — letters, digits, `_`, `-`. const ( verbOverview = "mesh_overview" verbMachine = "mesh_machine" verbSearch = "mesh_search" verbDescribe = "mesh_describe" verbCall = "mesh_call" + verbRuntimes = "mesh_runtimes" ) // searchCap is how many matches a search answers before it says how many more there were. @@ -49,7 +50,7 @@ const grammar = "Addresses: `.` for a seat held once for the mesh (e "`/.` for a module on one machine (e.g. `novox/postgres.postgres_list_databases`); " + "and `.` also for a module whose instances are interchangeable." -// discovery is the five tools as tools/list announces them. +// discovery is the six tools as tools/list announces them. func discovery() []map[string]any { str := func(desc string) map[string]any { return map[string]any{"type": "string", "description": desc} } obj := func(props map[string]any, required ...string) map[string]any { @@ -77,12 +78,17 @@ func discovery() []map[string]any { "arguments": map[string]any{"type": "object", "description": "the tool's arguments, as mesh_describe gives its schema"}, }, "address"), "description": "Call one tool by its address with its arguments; the answer says which machine gave it. " + grammar}, + {"name": verbRuntimes, "inputSchema": obj(map[string]any{}), + "description": "Which runtimes answered discovery, asked now: for each, its machine, how long its answer took to " + + "arrive, its size in bytes, how many modules and tools it announced, whether it was shortened to fit the " + + "bus, and when it was last heard — and every runtime or machine expected that was not heard. Read-only; " + + "for when an address is said to be missing or on another machine."}, } } func isDiscovery(name string) bool { switch name { - case verbOverview, verbMachine, verbSearch, verbDescribe, verbCall: + case verbOverview, verbMachine, verbSearch, verbDescribe, verbCall, verbRuntimes: return true } return false @@ -117,6 +123,44 @@ type index struct { Modules map[string]*moduleInfo NotAnswering []string Listing *Listing // the flat catalogue the call path resolves subjects with + Discovery announce.Discovery + // Recorded is where the controller's records place each module that declares tools. + Recorded map[string][]string + // Unheard is every runtime whose answer this index lacks: it said it was there and did not say + // what it serves in time, or it answered before and not now. An address it might answer is never + // called missing while it is here (2026-10-05). + Unheard []unheard +} + +// unheard is one runtime discovery did not hear in full, and why. +type unheard struct { + Runtime string `json:"runtime"` + Machine string `json:"machine,omitempty"` + Why string `json:"why"` +} + +// unheardOn is what is unheard on one machine; with no machine, everything unheard. +func (x *index) unheardOn(node string) []unheard { + var out []unheard + for _, u := range x.Unheard { + if node == "" || u.Machine == node || u.Machine == "" { + out = append(out, u) + } + } + return out +} + +// sayUnheard is the unheard, as a sentence's end. +func sayUnheard(us []unheard) string { + var parts []string + for _, u := range us { + who := u.Runtime + if u.Machine != "" { + who += " on " + u.Machine + } + parts = append(parts, who+" ("+u.Why+")") + } + return strings.Join(parts, "; ") } func (x *index) seat(name string) *seatInfo { @@ -193,20 +237,24 @@ func indexOn(conn *bus.Conn) (*index, error) { wg.Add(2) go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }() go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }() - infos, err := announce.Gather(conn) + d, err := announce.Gather(conn) wg.Wait() if err != nil { return nil, err } l := &Listing{Tools: []Tool{}, NotAnswering: []string{}} - x := &index{Modules: map[string]*moduleInfo{}, Listing: l} + x := &index{Modules: map[string]*moduleInfo{}, Listing: l, Discovery: d, Recorded: map[string][]string{}} + for _, s := range d.Silent { + x.Unheard = append(x.Unheard, unheard{Runtime: s.Name, Machine: s.Machine, + Why: fmt.Sprintf("it answered PING; what it serves did not arrive within %s", announce.Patience)}) + } announced := map[string]map[string]bool{} // module → node → announced something seats := map[string]*seatInfo{} toolAt := map[string]int{} // . or . → index in l.Tools machines := map[string]bool{} - for _, info := range infos { - for _, e := range announce.Endpoints(info) { + for _, h := range d.Heard { + for _, e := range announce.Endpoints(h.Info) { if e.Node != "" { machines[e.Node] = true } @@ -298,6 +346,7 @@ func indexOn(conn *bus.Conn) (*index, error) { if !m.Tools { continue } + x.Recorded[m.Module] = append([]string{}, m.On...) for _, n := range m.On { // A holder of a seat held once for the mesh announces no machine — the seat is the mesh's, // not a machine's — so what it announced without one answers for wherever it is assigned. @@ -309,6 +358,9 @@ func indexOn(conn *bus.Conn) (*index, error) { } else { l.NotAnswering = append(l.NotAnswering, "mesh-controller (its records of the modules did not answer, so what is missing cannot be said)") } + for _, u := range x.Unheard { + l.NotAnswering = append(l.NotAnswering, "the runtime "+sayUnheard([]unheard{u})) + } sort.Strings(l.NotAnswering) x.NotAnswering = l.NotAnswering @@ -336,9 +388,12 @@ func containsHolder(hs []holder, h holder) bool { return false } -func (s *Surface) index() (*index, error) { +func (s *Surface) index() (*index, error) { return s.indexAsked(false) } + +// indexAsked is what the mesh answered, asked again when `fresh` or when the kept answer is old. +func (s *Surface) indexAsked(fresh bool) (*index, error) { s.mu.Lock() - if s.idx != nil && time.Since(s.idxAt) <= IndexKept { + if !fresh && s.idx != nil && time.Since(s.idxAt) <= IndexKept { x := s.idx s.mu.Unlock() return x, nil @@ -349,6 +404,7 @@ func (s *Surface) index() (*index, error) { return nil, err } s.mu.Lock() + s.remember(x) s.idx, s.idxAt = x, time.Now() s.mu.Unlock() return x, nil @@ -397,6 +453,18 @@ func resolve(x *index, address string) (target, error) { m := x.Modules[prefix] if m == nil { + if on := x.Recorded[prefix]; len(on) > 0 { + why := "it did not announce itself on the bus" + if us := x.unheardOnAny(on); len(us) > 0 { + why = "discovery did not hear in full from " + sayUnheard(us) + } + return target{}, fmt.Errorf("%s is assigned to %s in the mesh's records, but nothing that answered discovery serves it: %s. "+ + "Ask again in a moment; mesh_runtimes shows which runtimes answered", prefix, strings.Join(on, ", "), why) + } + if len(x.Unheard) > 0 { + return target{}, fmt.Errorf("nothing that answered discovery is called %s, but not every runtime answered: %s — "+ + "%s may be theirs. Ask again in a moment; mesh_runtimes shows which runtimes answered", prefix, sayUnheard(x.Unheard), prefix) + } return target{}, fmt.Errorf("nothing in the mesh is called %s: no seat, and no module that answers tools. "+ "mesh_search finds a tool by words", prefix) } @@ -412,11 +480,30 @@ func resolve(x *index, address string) (target, error) { return target{Address: rest, Key: rest, Name: rest, Tool: *tool}, nil } if len(m.On) > 0 && !contains(m.On, node) { + if us := x.unheardOn(node); len(us) > 0 || contains(x.Recorded[prefix], node) { + why := "it did not announce itself there" + if len(us) > 0 { + why = "discovery did not hear in full from " + sayUnheard(us) + } + return target{}, fmt.Errorf("%s on %s did not answer discovery: %s. It was heard on %s. "+ + "Ask again in a moment; mesh_runtimes shows which runtimes answered", prefix, node, why, orNobody(m.On)) + } return target{}, fmt.Errorf("%s does not run on %s; it runs on %s", prefix, node, orNobody(m.On)) } return target{Address: node + "/" + rest, Key: rest + "@" + node, Name: rest, Node: node, Tool: *tool}, nil } +// unheardOnAny is what is unheard on any of these machines. +func (x *index) unheardOnAny(nodes []string) []unheard { + var out []unheard + for _, u := range x.Unheard { + if u.Machine == "" || contains(nodes, u.Machine) { + out = append(out, u) + } + } + return out +} + func contains(xs []string, s string) bool { for _, x := range xs { if x == s { @@ -511,8 +598,11 @@ func failure(text string) map[string]any { return map[string]any{"content": []map[string]any{{"type": "text", "text": text}}, "isError": true} } -// discover answers one of the five. +// discover answers one of the discovery verbs. func (s *Surface) discover(name string, args map[string]any) map[string]any { + if name == verbRuntimes { + return s.runtimesAnswer() + } x, err := s.index() if err != nil { return failure("the mesh's discovery failed: " + err.Error()) @@ -606,7 +696,18 @@ func (s *Surface) discover(name string, args map[string]any) map[string]any { } modules = append(modules, entry) } - return answerText(map[string]any{"machine": node, "seats": seats, "modules": modules}) + out := map[string]any{"machine": node, "seats": seats, "modules": modules} + if seats == nil { + out["seats"] = []map[string]any{} + } + if modules == nil { + out["modules"] = []map[string]any{} + } + if us := x.unheardOn(node); len(us) > 0 { + out["incomplete"] = "discovery did not hear in full from " + sayUnheard(us) + + "; what is listed may be missing what it serves. Ask again in a moment; mesh_runtimes shows which runtimes answered" + } + return answerText(out) case verbSearch: query := strings.ToLower(str("query")) @@ -684,6 +785,10 @@ func (s *Surface) discover(name string, args map[string]any) map[string]any { out["matches"] = []hit{} out["hint"] = "nothing matched every word; try fewer words, or mesh_overview and mesh_machine to browse" } + if len(x.Unheard) > 0 { + out["incomplete"] = "discovery did not hear in full from " + sayUnheard(x.Unheard) + + "; what they serve is not searched. mesh_runtimes shows which runtimes answered" + } return answerText(out) case verbDescribe: diff --git a/node-tools/internal/console/address_test.go b/node-tools/internal/console/address_test.go index 4a0cfb3..60b6903 100644 --- a/node-tools/internal/console/address_test.go +++ b/node-tools/internal/console/address_test.go @@ -36,7 +36,7 @@ func call(t *testing.T, endpoint, tool string, args map[string]any) (string, boo return text(t, post(t, endpoint, string(body))) } -// novox/hq ADR 0195: the console announces five tools, and everything the mesh answers is reached +// novox/hq ADR 0195: the console announces six tools, and everything the mesh answers is reached // through them by one address per layer. func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { was := IndexKept @@ -113,7 +113,7 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { defer up.Close() endpoint := "http://" + up.Address + "/mcp" - // Five tools, nothing else. + // Six tools, nothing else. listed := post(t, endpoint, `{"jsonrpc":"2.0","id":2,"method":"tools/list"}`)["result"].(map[string]any) var names []string for _, x := range listed["tools"].([]any) { @@ -125,7 +125,7 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { } } } - if got := strings.Join(names, ","); got != "mesh_overview,mesh_machine,mesh_search,mesh_describe,mesh_call" { + if got := strings.Join(names, ","); got != "mesh_overview,mesh_machine,mesh_search,mesh_describe,mesh_call,mesh_runtimes" { t.Errorf("announced %s", got) } diff --git a/node-tools/internal/console/mcp.go b/node-tools/internal/console/mcp.go index aa4f801..0183e83 100644 --- a/node-tools/internal/console/mcp.go +++ b/node-tools/internal/console/mcp.go @@ -50,6 +50,8 @@ type Surface struct { at time.Time idx *index idxAt time.Time + // runtimes is every runtime discovery has heard, as it last heard it (mesh_runtimes). + runtimes map[string]*runtimeSeen // Flat announces the whole catalogue, as the console did before ADR 0195: for a person reading // it or a client that wants it. Off by default; MESH_CONSOLE_FLAT=1 turns it on. Flat bool @@ -300,7 +302,8 @@ func (s *Surface) instructions() string { return "The tools of a Novox mesh, reached as " + s.who + ", found by address rather than listed " + "whole (novox/hq ADR 0195). mesh_overview shows the mesh's seats and machines; mesh_machine one " + "machine's seats and modules; mesh_search finds a tool by words; mesh_describe gives one tool's " + - "arguments; mesh_call calls it. " + grammar + " What may be called was fixed when this account " + + "arguments; mesh_call calls it; mesh_runtimes says which machines' runtimes answered discovery, how fast " + + "and how much — ask it when an address is said to be missing or on another machine. " + grammar + " What may be called was fixed when this account " + "was issued, so a refusal means the account, not the tool." } diff --git a/node-tools/internal/console/runtimes.go b/node-tools/internal/console/runtimes.go new file mode 100644 index 0000000..855b28f --- /dev/null +++ b/node-tools/internal/console/runtimes.go @@ -0,0 +1,200 @@ +package console + +// What discovery heard from each runtime (2026-10-05): the console said a module ran nowhere when the +// largest runtime's answer arrived after the window closed, and nothing showed that it had. Every +// round of discovery is remembered per runtime, so a runtime heard before and not now is named, and +// `mesh_runtimes` says who answered, how fast and how much. + +import ( + "fmt" + "sort" + "time" + + "github.com/novox/mesh-tools/node-tools/internal/announce" +) + +// RuntimeForgotten is how long a runtime heard once is expected again: within it, its silence is said +// wherever an address it might answer is asked for; after it, it is listed only by mesh_runtimes. +var RuntimeForgotten = 15 * time.Minute + +// runtimeSeen is one runtime instance as discovery last heard it. +type runtimeSeen struct { + Name string + ID string + Machine string + LastHeard time.Time + Took time.Duration + Bytes int + Modules int + Tools int + SeatVerbs int + Shortened string + Answered bool // in the latest round + Why string // why not, when it did not +} + +func runtimeKey(name, id string) string { return name + "/" + id } + +// seenOf is what one answer announced, counted. +func seenOf(h announce.Heard, at time.Time) *runtimeSeen { + modules := map[string]bool{} + tools := map[string]bool{} + verbs := map[string]bool{} + for _, e := range announce.Endpoints(h.Info) { + switch e.Kind { + case announce.KindTool: + modules[e.Module] = true + tools[e.Module+"."+e.Tool] = true + case announce.KindSeat: + verbs[e.Seat+"."+e.Tool] = true + } + } + return &runtimeSeen{Name: h.Info.Name, ID: h.Info.ID, Machine: h.Machine(), LastHeard: at, Took: h.Took, + Bytes: h.Bytes, Modules: len(modules), Tools: len(tools), SeatVerbs: len(verbs), Shortened: h.Shortened, Answered: true} +} + +// remember records a round of discovery and adds to the index every runtime heard lately that did not +// answer this time. Called with s.mu held. +func (s *Surface) remember(x *index) { + if s.runtimes == nil { + s.runtimes = map[string]*runtimeSeen{} + } + at := x.Discovery.At + if at.IsZero() { + at = time.Now().UTC() + } + heard := map[string]bool{} + for _, h := range x.Discovery.Heard { + k := runtimeKey(h.Info.Name, h.Info.ID) + seen := seenOf(h, at) + if heard[k] { + // Two instances under one name and id: both counted, the slower answer's time kept. + prev := s.runtimes[k] + seen.Modules += prev.Modules + seen.Tools += prev.Tools + seen.SeatVerbs += prev.SeatVerbs + seen.Bytes += prev.Bytes + if prev.Took > seen.Took { + seen.Took = prev.Took + } + } + heard[k] = true + s.runtimes[k] = seen + } + // A runtime restarted answers under a new instance id: the old one is gone, not silent, and saying + // it was missed would hide nothing but cry wolf after every restart. + now := map[string]bool{} + for k := range heard { + now[s.runtimes[k].Name+"@"+s.runtimes[k].Machine] = true + } + for k, r := range s.runtimes { + if !heard[k] && r.Machine != "" && now[r.Name+"@"+r.Machine] { + delete(s.runtimes, k) + } + } + silent := map[string]string{} + for _, u := range x.Discovery.Silent { + silent[runtimeKey(u.Name, u.ID)] = fmt.Sprintf("it answered PING; what it serves did not arrive within %s", announce.Patience) + } + for k, r := range s.runtimes { + if heard[k] { + continue + } + r.Answered = false + if why, ok := silent[k]; ok { + r.Why = why + continue // already in the index's unheard + } + r.Why = "it did not answer this round of discovery" + if time.Since(r.LastHeard) <= RuntimeForgotten { + x.Unheard = append(x.Unheard, unheard{Runtime: r.Name, Machine: r.Machine, + Why: fmt.Sprintf("it answered discovery %s ago, not now", time.Since(r.LastHeard).Round(time.Second))}) + } + } + for _, u := range x.Discovery.Silent { + k := runtimeKey(u.Name, u.ID) + if _, known := s.runtimes[k]; !known { + s.runtimes[k] = &runtimeSeen{Name: u.Name, ID: u.ID, Machine: u.Machine, Why: silent[k]} + } + } +} + +// runtimesAnswer is mesh_runtimes: discovery asked now, and every runtime as it was heard. +func (s *Surface) runtimesAnswer() map[string]any { + x, err := s.indexAsked(true) + if err != nil { + return failure("the mesh's discovery failed: " + err.Error()) + } + s.mu.Lock() + var all []runtimeSeen + for _, r := range s.runtimes { + all = append(all, *r) + } + s.mu.Unlock() + sort.Slice(all, func(i, j int) bool { + if all[i].Machine != all[j].Machine { + return all[i].Machine < all[j].Machine + } + return runtimeKey(all[i].Name, all[i].ID) < runtimeKey(all[j].Name, all[j].ID) + }) + now := time.Now() + var answered, notHeard []map[string]any + onMachine := map[string]bool{} + for _, r := range all { + entry := map[string]any{"runtime": r.Name, "instance": r.ID} + if r.Machine != "" { + entry["machine"] = r.Machine + } + if !r.LastHeard.IsZero() { + entry["last heard"] = r.LastHeard.Format(time.RFC3339) + entry["last heard, ago"] = now.Sub(r.LastHeard).Round(time.Second).String() + entry["took"] = r.Took.Round(time.Millisecond).String() + entry["bytes"] = r.Bytes + entry["modules"] = r.Modules + entry["tools"] = r.Tools + entry["seat verbs"] = r.SeatVerbs + entry["shortened"] = r.Shortened != "" + if r.Shortened != "" { + entry["shortened, how"] = r.Shortened + } + } + if r.Answered { + onMachine[r.Machine] = true + answered = append(answered, entry) + } else { + entry["why"] = r.Why + if r.LastHeard.IsZero() { + entry["last heard"] = "never in full, since this console started" + } + notHeard = append(notHeard, entry) + } + } + for _, m := range x.Machines { + if !onMachine[m] && !machineListed(notHeard, m) { + notHeard = append(notHeard, map[string]any{"machine": m, + "why": "the mesh knows this machine and no runtime on it has answered discovery since this console started"}) + } + } + if answered == nil { + answered = []map[string]any{} + } + if notHeard == nil { + notHeard = []map[string]any{} + } + return answerText(map[string]any{ + "asked": x.Discovery.At.Format(time.RFC3339), + "waited": x.Discovery.Waited.Round(time.Millisecond).String(), + "window": fmt.Sprintf("at least %s; up to %s for a runtime that answered PING", announce.Window, announce.Patience), + "answered": answered, + "not heard": notHeard, + }) +} + +func machineListed(entries []map[string]any, machine string) bool { + for _, e := range entries { + if e["machine"] == machine { + return true + } + } + return false +} diff --git a/node-tools/internal/console/runtimes_test.go b/node-tools/internal/console/runtimes_test.go new file mode 100644 index 0000000..c82d48a --- /dev/null +++ b/node-tools/internal/console/runtimes_test.go @@ -0,0 +1,219 @@ +package console + +import ( + "encoding/json" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go/micro" + + "github.com/novox/mesh-tools/node-tools/internal/announce" + "github.com/novox/mesh-tools/node-tools/internal/bus" + mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" +) + +// slowRuntime answers PING at once and INFO after `lag`, or never with a negative lag. +func slowRuntime(t *testing.T, conn *bus.Conn, node string, lag time.Duration, endpoints []announce.Endpoint) { + t.Helper() + id := micro.ServiceIdentity{Name: "node-tools", ID: node, Version: announce.Version, Metadata: map[string]string{"node": node}} + ping, _ := json.Marshal(micro.Ping{ServiceIdentity: id, 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(announce.Info(announce.Service{Name: id.Name, ID: id.ID, Metadata: id.Metadata}, endpoints)) + return body + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stop) + conn.Flush() +} + +// The laptop's symptom (2026-10-05), reproduced: its runtime's one large answer reached the console +// after the 750 ms window, and the console said its modules ran nowhere — "nothing in the mesh is +// called slack" — or ran only on another machine. Now a runtime that said it is there is waited for, +// one that never says what it serves is named wherever an address it might answer is asked for, and +// mesh_runtimes shows who answered, how fast and how much. +func TestARuntimeThatAnswersLateIsFoundAndOneThatNeverDoesIsNamed(t *testing.T) { + keptWas, patienceWas := IndexKept, announce.Patience + // Kept for the test: one round of discovery answers every question but mesh_runtimes, which asks anew. + IndexKept, announce.Patience = time.Minute, 2500*time.Millisecond + t.Cleanup(func() { IndexKept, announce.Patience = keptWas, patienceWas }) + mt.New(t) + + laptop := connect(t, "node-tools", "laptop") + stopTool, err := laptop.HandleSubject("mesh.mod.slack.tool.slack_check.laptop", func(json.RawMessage) (any, error) { + return map[string]any{"ok": true}, nil + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stopTool) + slowRuntime(t, laptop, "laptop", announce.Window+400*time.Millisecond, []announce.Endpoint{ + {Kind: announce.KindTool, Module: "slack", Tool: "slack_check", Node: "laptop", Description: "Is Slack started once?", + Subject: "mesh.mod.slack.tool.slack_check.laptop"}, + {Kind: announce.KindTool, Module: "nextcloud-client", Tool: "nextcloud_check", Node: "laptop", + Subject: "mesh.mod.nextcloud-client.tool.nextcloud_check.laptop"}, + }) + sleeper := connect(t, "node-tools", "sleeper") + slowRuntime(t, sleeper, "sleeper", -1, nil) + desktop := connect(t, "node-tools", "desktop") + slowRuntime(t, desktop, "desktop", 0, []announce.Endpoint{ + {Kind: announce.KindTool, Module: "nextcloud-client", Tool: "nextcloud_check", Node: "desktop", + Subject: "mesh.mod.nextcloud-client.tool.nextcloud_check.desktop"}, + }) + + controller := connect(t, "mesh-controller", "bench") + for verb, output := range map[string]string{ + "nodes": `[{"name":"laptop"},{"name":"sleeper"},{"name":"desktop"},{"name":"bench"}]`, + "modules": `[{"module":"slack","on":["laptop"],"tools":true},{"module":"nap","on":["sleeper"],"tools":true},` + + `{"module":"nextcloud-client","on":["laptop","desktop","sleeper"],"tools":true}]`, + } { + output := output + stop, err := controller.HandleSubject("mesh.seat.mesh-controller.tool."+verb, func(json.RawMessage) (any, error) { + return map[string]any{"output": output, "ok": true}, nil + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stop) + } + controller.Flush() + + console := connect(t, "node-tools", "laptop-console") + up, err := Serve(NewSurface(console, "laptop.node-tools"), "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer up.Close() + endpoint := "http://" + up.Address + "/mcp" + + // The late answer is waited for: the laptop's module is found and called. + if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "laptop/slack.slack_check"}); isErr || !strings.Contains(got, `"ok": true`) { + t.Errorf("a late runtime's module: %s", got) + } + if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "laptop/nextcloud-client.nextcloud_check"}); isErr && strings.Contains(got, "does not run on laptop") { + t.Errorf("a late runtime's module was said to run elsewhere: %s", got) + } + machine, _ := call(t, endpoint, "mesh_machine", map[string]any{"node": "laptop"}) + if !strings.Contains(machine, "laptop/slack.slack_check") { + t.Errorf("the late runtime's machine lists nothing: %s", machine) + } + + // The silent one is named, never called missing. + for address, want := range map[string]string{ + "sleeper/nap.doze": "nap is assigned to sleeper in the mesh's records, but nothing that answered discovery serves it: discovery did not hear in full from node-tools on sleeper (it answered PING", + "sleeper/nextcloud-client.nextcloud_check": "nextcloud-client on sleeper did not answer discovery: discovery did not hear in full from node-tools on sleeper", + } { + got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": address}) + if !isErr || !strings.Contains(got, want) || strings.Contains(got, "nothing in the mesh is called") || strings.Contains(got, "does not run on") { + t.Errorf("%s: %s\nwant %q", address, got, want) + } + } + machine, _ = call(t, endpoint, "mesh_machine", map[string]any{"node": "sleeper"}) + if !strings.Contains(machine, `"incomplete"`) || !strings.Contains(machine, `"modules": []`) { + t.Errorf("a machine whose runtime did not answer: %s", machine) + } + overview, _ := call(t, endpoint, "mesh_overview", nil) + if !strings.Contains(overview, "the runtime node-tools on sleeper") { + t.Errorf("overview does not name the silent runtime: %s", overview) + } + + // mesh_runtimes: who answered, how fast, how much; who did not. + got, isErr := call(t, endpoint, "mesh_runtimes", nil) + if isErr { + t.Fatal(got) + } + var report struct { + Answered []map[string]any `json:"answered"` + NotHeard []map[string]any `json:"not heard"` + } + if err := json.Unmarshal([]byte(got), &report); err != nil { + t.Fatalf("%v: %s", err, got) + } + var lap map[string]any + for _, r := range report.Answered { + if r["machine"] == "laptop" { + lap = r + } + } + if lap == nil { + t.Fatalf("the laptop's runtime is not listed as answering: %s", got) + } + took, _ := time.ParseDuration(lap["took"].(string)) + if took < announce.Window || lap["bytes"].(float64) <= 0 || lap["modules"].(float64) != 2 || lap["tools"].(float64) != 2 || + lap["shortened"] != false || lap["last heard"] == nil { + t.Errorf("the laptop's runtime: %v", lap) + } + if !strings.Contains(got, `"machine": "sleeper"`) || !strings.Contains(got, "answered PING") { + t.Errorf("the silent runtime is not listed: %s", got) + } + silent := false + for _, r := range report.NotHeard { + silent = silent || r["machine"] == "sleeper" + } + if !silent { + t.Errorf("sleeper is not among the not heard: %s", got) + } +} + +// A runtime heard once and silent now is named in the answers about its machine. +func TestARuntimeHeardBeforeAndNotNowIsNamed(t *testing.T) { + s := &Surface{} + before := &index{Modules: map[string]*moduleInfo{}, Recorded: map[string][]string{}, Discovery: announce.Discovery{At: time.Now().UTC(), + Heard: []announce.Heard{{Info: micro.Info{ServiceIdentity: micro.ServiceIdentity{Name: "node-tools", ID: "laptop", + Metadata: map[string]string{"node": "laptop"}}}, Took: 300 * time.Millisecond, Bytes: 1000}}}} + s.remember(before) + now := &index{Modules: map[string]*moduleInfo{"nextcloud-client": {Module: "nextcloud-client", On: []string{"desktop"}, + Tools: []Tool{{Module: "nextcloud-client", Name: "nextcloud_check"}}}}, Recorded: map[string][]string{}, + Discovery: announce.Discovery{At: time.Now().UTC()}} + s.remember(now) + if len(now.Unheard) != 1 || now.Unheard[0].Machine != "laptop" { + t.Fatalf("unheard: %+v", now.Unheard) + } + _, err := resolve(now, "laptop/nextcloud-client.nextcloud_check") + if err == nil || !strings.Contains(err.Error(), "nextcloud-client on laptop did not answer discovery") || + !strings.Contains(err.Error(), "answered discovery") { + t.Errorf("%v", err) + } + _, err = resolve(now, "laptop/slack.slack_check") + if err == nil || strings.Contains(err.Error(), "nothing in the mesh is called") || !strings.Contains(err.Error(), "slack may be theirs") { + t.Errorf("%v", err) + } + // Nobody unheard: the plain answers stand. + quiet := &index{Modules: now.Modules, Recorded: map[string][]string{}} + if _, err := resolve(quiet, "laptop/nextcloud-client.nextcloud_check"); err == nil || !strings.Contains(err.Error(), "does not run on laptop; it runs on desktop") { + t.Errorf("%v", err) + } + if _, err := resolve(quiet, "laptop/slack.slack_check"); err == nil || !strings.Contains(err.Error(), "nothing in the mesh is called slack") { + t.Errorf("%v", err) + } +} + +// A runtime restarted answers under a new instance id; its old instance is gone, not unheard. +func TestARestartedRuntimeIsNotNamedUnheard(t *testing.T) { + heard := func(id string) *index { + return &index{Modules: map[string]*moduleInfo{}, Recorded: map[string][]string{}, Discovery: announce.Discovery{At: time.Now().UTC(), + Heard: []announce.Heard{{Info: micro.Info{ServiceIdentity: micro.ServiceIdentity{Name: "node-tools", ID: id, + Metadata: map[string]string{"node": "laptop"}}}}}}} + } + s := &Surface{} + s.remember(heard("before-restart")) + after := heard("after-restart") + s.remember(after) + if len(after.Unheard) != 0 { + t.Errorf("a restarted runtime was named unheard: %+v", after.Unheard) + } + if len(s.runtimes) != 1 { + t.Errorf("the old instance is still kept: %d runtimes", len(s.runtimes)) + } +}