From a269a76c87797eebff70a36686020fbd0fc7470b Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 21:54:19 +0200 Subject: [PATCH 1/3] The console announces five tools and reaches everything by address (hq ADR 0195) mesh_overview, mesh_machine, mesh_search, mesh_describe and mesh_call walk the mesh's structure; every tool has one address per layer: ., /., /., and . for a module the mesh issued a plain subject. A stateful module called without its machine, a node seat without one, a mesh seat with one, or a module on the wrong machine is refused naming what would work. Answers come from the mesh when asked, kept five seconds, so a tool that arrives mid-session is found. The flat catalogue stays behind MESH_CONSOLE_FLAT=1 and old . names still answer. --- node-tools/internal/console/address.go | 735 ++++++++++++++++++++ node-tools/internal/console/address_test.go | 222 ++++++ node-tools/internal/console/console_test.go | 4 +- node-tools/internal/console/mcp.go | 43 +- 4 files changed, 996 insertions(+), 8 deletions(-) create mode 100644 node-tools/internal/console/address.go create mode 100644 node-tools/internal/console/address_test.go diff --git a/node-tools/internal/console/address.go b/node-tools/internal/console/address.go new file mode 100644 index 0000000..3909b59 --- /dev/null +++ b/node-tools/internal/console/address.go @@ -0,0 +1,735 @@ +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: +// +// . a seat held once for the mesh — its holder answers +// /. a seat held once per machine — that machine's holder answers +// /. a module assigned to a machine — that assignment answers +// . also, for a module whose instances are interchangeable (ADR 0160) +// +// A module that is not interchangeable is called with its machine or refused, naming the machines +// it runs on: "whichever answers" is no answer for state a machine holds. +// +// Every discovery verb asks the mesh when it is called — kept a few seconds at most, never for a +// session — so a tool that arrived a minute ago is found without the client reconnecting. + +import ( + "encoding/json" + "fmt" + "regexp" + "sort" + "strings" + "sync" + "time" + + "github.com/novox/mesh-tools/node-tools/internal/bus" +) + +// IndexKept is how long what the mesh answered is kept before it is asked again: long enough that +// 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, `_`, `-`. +const ( + verbOverview = "mesh_overview" + verbMachine = "mesh_machine" + verbSearch = "mesh_search" + verbDescribe = "mesh_describe" + verbCall = "mesh_call" +) + +// searchCap is how many matches a search answers before it says how many more there were. +const searchCap = 25 + +const grammar = "Addresses: `.` for a seat held once for the mesh (e.g. `mesh-controller.nodes`); " + + "`/.` for a seat every machine holds (e.g. `ace/node-packet-filter.rules`); " + + "`/.` 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. +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 { + s := map[string]any{"type": "object", "properties": props} + if len(required) > 0 { + s["required"] = required + } + return s + } + return []map[string]any{ + {"name": verbOverview, "inputSchema": obj(map[string]any{}), + "description": "The mesh at a glance: the seats it holds once for the whole mesh with their verbs, the seats " + + "every machine holds, and its machines. Start here, then `mesh_machine` for one machine. " + grammar}, + {"name": verbMachine, "inputSchema": obj(map[string]any{"node": str("the machine, as mesh_overview names it")}, "node"), + "description": "One machine: the seats it holds with their verbs, and the modules assigned to it with their tools — " + + "each with the address to describe or call it by. " + grammar}, + {"name": verbSearch, "inputSchema": obj(map[string]any{"query": str("words to find in tool names and descriptions, e.g. `postgres databases`")}, "query"), + "description": "Find tools anywhere in the mesh by words: every match's address and a line of what it does, across the " + + "mesh's seats, the machines' seats and every module on every machine. " + grammar}, + {"name": verbDescribe, "inputSchema": obj(map[string]any{"address": str("the tool's address")}, "address"), + "description": "What one tool does and the arguments it takes, as a JSON schema. The machine is in the address, " + + "never an argument. " + grammar}, + {"name": verbCall, "inputSchema": obj(map[string]any{ + "address": str("the tool's address"), + "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}, + } +} + +func isDiscovery(name string) bool { + switch name { + case verbOverview, verbMachine, verbSearch, verbDescribe, verbCall: + return true + } + return false +} + +// seatInfo is a seat as the mesh's records define it, and who holds it where. +type seatInfo struct { + Seat string + Scope string // "mesh" or "node" + Verbs []Tool + Holders []holder +} + +type holder struct { + Module string `json:"module"` + Node string `json:"node"` +} + +// moduleInfo is a module that answers tools: where it runs, whether any instance will do, and its +// tools as one of its instances described them. +type moduleInfo struct { + Module string + On []string + Interchangeable bool + Tools []Tool +} + +// index is what the mesh answered about itself, at one moment. +type index struct { + Seats []seatInfo + Machines []string + Modules map[string]*moduleInfo + NotAnswering []string + Listing *Listing // the flat catalogue the call path resolves subjects with +} + +func (x *index) seat(name string) *seatInfo { + for i := range x.Seats { + if x.Seats[i].Seat == name { + return &x.Seats[i] + } + } + return nil +} + +func findTool(tools []Tool, name string) *Tool { + for i := range tools { + if tools[i].Name == name { + return &tools[i] + } + } + return nil +} + +// controllerOutput is a controller seat verb's answer: the command's printed output. +func controllerOutput(conn *bus.Conn, verb string) (string, error) { + got, err := conn.Ask("seat:mesh-controller."+verb, map[string]any{}, "") + if err != nil { + return "", err + } + var r struct { + Output string `json:"output"` + OK *bool `json:"ok"` + } + if json.Unmarshal(got.Result, &r) != nil { + return "", fmt.Errorf("mesh-controller.%s answered something that is not its output", verb) + } + if r.OK != nil && !*r.OK { + return "", fmt.Errorf("mesh-controller.%s: %s", verb, strings.TrimSpace(r.Output)) + } + return r.Output, nil +} + +// jsonIn is the JSON document a command printed, after any lines it said first: a seat verb runs +// the controller's command, and a command may warn before it answers. +func jsonIn(output string) string { + if strings.HasPrefix(strings.TrimSpace(output), "{") { + return strings.TrimSpace(output) + } + if i := strings.Index(output, "\n{"); i >= 0 { + return strings.TrimSpace(output[i+1:]) + } + return "" +} + +// A machine as `node list` prints it: its name, when it was last heard from, its mode — converged +// or adopted, the only two the command prints — and its id. Lines the command says around them +// (a warning about the bus's users, "no node records yet") are not machines and are skipped. +var nodeLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+.*\s(converged|adopted)\s+\S+\s*$`) + +// machinesIn reads the machines from `node list`'s output. +func machinesIn(output string) []string { + var out []string + for _, line := range strings.Split(output, "\n") { + if m := nodeLine.FindStringSubmatch(line); m != nil { + out = append(out, m[1]) + } + } + sort.Strings(out) + return out +} + +// A module as `module list` prints it: name, version, how it was built, and where it runs. +var moduleLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+\S+\s+.*?\s+on (.+)$`) + +// assignmentsIn reads, from `module list`'s output, the machines each module runs on. +func assignmentsIn(output string) map[string][]string { + out := map[string][]string{} + for _, line := range strings.Split(output, "\n") { + if line == "" || line[0] == ' ' || line[0] == '\t' { + continue + } + m := moduleLine.FindStringSubmatch(strings.TrimRight(line, " ")) + if m == nil { + continue + } + on := strings.TrimSpace(m[2]) + if on == "nothing" { + out[m[1]] = []string{} + continue + } + var nodes []string + for _, n := range strings.Split(on, ",") { + if n = strings.TrimSpace(n); n != "" { + nodes = append(nodes, n) + } + } + sort.Strings(nodes) + out[m[1]] = nodes + } + return out +} + +// interchangeable is whether the mesh issued the module a plain subject for this tool — one any of +// its instances answers (ADR 0160): the module's own subject with no machine after it. +func interchangeable(module string, t Tool) bool { + plain := "mesh.mod." + module + ".tool." + t.Name + for _, s := range t.Subjects { + if s == plain { + return true + } + } + return false +} + +// indexOn asks the mesh what it holds: the flat listing (catalogue, modules, the seats' tools), +// and from the controller the seats' holders, the machines and the assignments — at once. +func indexOn(conn *bus.Conn) (*index, error) { + var wg sync.WaitGroup + var seatsOut, nodesOut, modulesOut string + wg.Add(3) + go func() { defer wg.Done(); seatsOut, _ = controllerOutput(conn, "seats") }() + go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }() + go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }() + l, err := toolsOn(conn) + wg.Wait() + if err != nil { + return nil, err + } + + x := &index{Modules: map[string]*moduleInfo{}, NotAnswering: l.NotAnswering, Listing: l} + + // The seats: their verbs from the mesh's records, their holders from the controller. + var held struct { + Seats []struct { + Seat string `json:"seat"` + Scope string `json:"scope"` + Holders []holder `json:"holders"` + } `json:"seats"` + } + _ = json.Unmarshal([]byte(jsonIn(seatsOut)), &held) + holders := map[string][]holder{} + scopes := map[string]string{} + for _, s := range held.Seats { + holders[s.Seat] = s.Holders + scopes[s.Seat] = s.Scope + } + bySeat := map[string]*seatInfo{} + for _, t := range l.Tools { + if !t.Seat { + continue + } + s := bySeat[t.Module] + if s == nil { + scope := t.Scope + if scope == "" { + scope = scopes[t.Module] + } + if scope != "node" { + scope = "mesh" + } + s = &seatInfo{Seat: t.Module, Scope: scope, Holders: holders[t.Module]} + bySeat[t.Module] = s + } + s.Verbs = append(s.Verbs, t) + } + for _, s := range bySeat { + x.Seats = append(x.Seats, *s) + } + sort.Slice(x.Seats, func(i, j int) bool { return x.Seats[i].Seat < x.Seats[j].Seat }) + + // The machines: what the controller knows, and any that hold a seat. + seen := map[string]bool{} + for _, n := range machinesIn(nodesOut) { + seen[n] = true + } + for _, s := range x.Seats { + for _, h := range s.Holders { + if h.Node != "" { + seen[h.Node] = true + } + } + } + + // The modules that answer tools, where they run, and whether any instance will do. + on := assignmentsIn(modulesOut) + for _, t := range l.Tools { + if t.Seat { + continue + } + m := x.Modules[t.Module] + if m == nil { + m = &moduleInfo{Module: t.Module, On: on[t.Module]} + x.Modules[t.Module] = m + } + m.Tools = append(m.Tools, t) + if interchangeable(t.Module, t) { + m.Interchangeable = true + } + } + for _, m := range x.Modules { + // A module the controller could not place is placed where its own answer says it runs. + if len(m.On) == 0 { + nodes := map[string]bool{} + for _, t := range m.Tools { + for _, s := range t.Subjects { + base := "mesh.mod." + m.Module + ".tool." + t.Name + "." + if strings.HasPrefix(s, base) { + nodes[strings.TrimPrefix(s, base)] = true + } + } + } + for n := range nodes { + m.On = append(m.On, n) + } + sort.Strings(m.On) + } + for _, n := range m.On { + seen[n] = true + } + } + for n := range seen { + x.Machines = append(x.Machines, n) + } + sort.Strings(x.Machines) + return x, nil +} + +func (s *Surface) index() (*index, error) { + s.mu.Lock() + if s.idx != nil && time.Since(s.idxAt) <= IndexKept { + x := s.idx + s.mu.Unlock() + return x, nil + } + s.mu.Unlock() + x, err := indexOn(s.conn) + if err != nil { + return nil, err + } + s.mu.Lock() + s.idx, s.idxAt = x, time.Now() + s.mu.Unlock() + return x, nil +} + +// target is what an address resolves to. +type target struct { + Address string + Key string // the call key the existing path takes: seat:.[@node] or .[@node] + Name string // . or ., for the listing's subject lookup + Node string + Tool Tool + Seat bool +} + +// resolve turns an address into exactly one target, or says why it cannot. +func resolve(x *index, address string) (target, error) { + address = strings.TrimSpace(address) + node, rest, hasNode := strings.Cut(address, "/") + if !hasNode { + rest, node = address, "" + } + dot := strings.Index(rest, ".") + if dot <= 0 || dot == len(rest)-1 || strings.Contains(rest, "/") { + return target{}, fmt.Errorf("%q is not an address. %s", address, grammar) + } + prefix, name := rest[:dot], rest[dot+1:] + + if s := x.seat(prefix); s != nil { + verb := findTool(s.Verbs, name) + if verb == nil { + return target{}, fmt.Errorf("the seat %s has no verb %s; it has %s", prefix, name, toolNames(s.Verbs)) + } + if s.Scope == "node" { + if node == "" { + return target{}, fmt.Errorf("%s is held once per machine: write /%s — it is held on %s", + prefix, rest, orNobody(nodesOf(s.Holders))) + } + return target{Address: node + "/" + rest, Key: "seat:" + rest + "@" + node, Name: rest, Node: node, Tool: *verb, Seat: true}, nil + } + if node != "" { + return target{}, fmt.Errorf("%s is held once for the whole mesh: write %s, without a machine", prefix, rest) + } + return target{Address: rest, Key: "seat:" + rest, Name: rest, Tool: *verb, Seat: true}, nil + } + + m := x.Modules[prefix] + if m == nil { + 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) + } + tool := findTool(m.Tools, name) + if tool == nil { + return target{}, fmt.Errorf("%s has no tool %s; it has %s", prefix, name, toolNames(m.Tools)) + } + if node == "" { + if !m.Interchangeable { + return target{}, fmt.Errorf("%s keeps state on each machine it runs on, so a call names the machine: "+ + "write /%s — it runs on %s", prefix, rest, orNobody(m.On)) + } + return target{Address: rest, Key: rest, Name: rest, Tool: *tool}, nil + } + if len(m.On) > 0 && !contains(m.On, node) { + 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 +} + +func contains(xs []string, s string) bool { + for _, x := range xs { + if x == s { + return true + } + } + return false +} + +func toolNames(ts []Tool) string { + names := make([]string, 0, len(ts)) + for _, t := range ts { + names = append(names, t.Name) + } + sort.Strings(names) + return strings.Join(names, ", ") +} + +func nodesOf(hs []holder) []string { + var out []string + for _, h := range hs { + if h.Node != "" && !contains(out, h.Node) { + out = append(out, h.Node) + } + } + sort.Strings(out) + return out +} + +func orNobody(nodes []string) string { + if len(nodes) == 0 { + return "no machine the mesh knows of" + } + return strings.Join(nodes, ", ") +} + +// firstLine is a description's first line, for a list. +func firstLine(s string) string { + s = strings.TrimSpace(s) + if i := strings.IndexAny(s, "\n"); i >= 0 { + s = s[:i] + } + if len(s) > 160 { + s = s[:157] + "…" + } + return s +} + +// schemaWithoutNode is a tool's schema as an agent passes it: the machine is in the address. +func schemaWithoutNode(raw json.RawMessage) map[string]any { + schema := asSchema(raw) + out := map[string]any{} + for k, v := range schema { + out[k] = v + } + if p, ok := schema["properties"].(map[string]any); ok { + props := map[string]any{} + for k, v := range p { + if k != "node" { + props[k] = v + } + } + out["properties"] = props + } + if r, ok := schema["required"].([]any); ok { + var keep []any + for _, k := range r { + if k != "node" { + keep = append(keep, k) + } + } + if len(keep) == 0 { + delete(out, "required") + } else { + out["required"] = keep + } + } + return out +} + +// answerText is a discovery verb's answer as MCP content: JSON, indented. +func answerText(v any) map[string]any { + var b strings.Builder + enc := json.NewEncoder(&b) + enc.SetEscapeHTML(false) // `/…` is read by an agent, not embedded in a page + enc.SetIndent("", " ") + _ = enc.Encode(v) + return map[string]any{"content": []map[string]any{{"type": "text", "text": strings.TrimRight(b.String(), "\n")}}} +} + +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. +func (s *Surface) discover(name string, args map[string]any) map[string]any { + x, err := s.index() + if err != nil { + return failure(whyItFailed(catalogueModules, err)) + } + str := func(k string) string { v, _ := args[k].(string); return strings.TrimSpace(v) } + + switch name { + case verbOverview: + type verbLine struct { + Address string `json:"address"` + Description string `json:"description"` + } + var mesh, node []map[string]any + for _, st := range x.Seats { + var verbs []verbLine + for _, v := range st.Verbs { + addr := st.Seat + "." + v.Name + if st.Scope == "node" { + addr = "/" + addr + } + verbs = append(verbs, verbLine{addr, firstLine(v.Description)}) + } + entry := map[string]any{"seat": st.Seat, "verbs": verbs} + if st.Scope == "node" { + entry["held on"] = nodesOf(st.Holders) + node = append(node, entry) + } else { + if h := nodesOf(st.Holders); len(h) > 0 { + entry["held on"] = h + } + mesh = append(mesh, entry) + } + } + modules := 0 + tools := 0 + for _, m := range x.Modules { + modules++ + tools += len(m.Tools) + } + return answerText(map[string]any{ + "seats of the mesh": mesh, + "seats every machine": node, + "machines": x.Machines, + "modules with tools": fmt.Sprintf("%d modules, %d tools — mesh_machine lists a machine's, mesh_search finds one", modules, tools), + "not answering": x.NotAnswering, + }) + + case verbMachine: + node := str("node") + if node == "" { + return failure("mesh_machine needs `node`: one of " + orNobody(x.Machines)) + } + if !contains(x.Machines, node) { + return failure(fmt.Sprintf("the mesh knows no machine %q; it has %s", node, orNobody(x.Machines))) + } + var seats []map[string]any + for _, st := range x.Seats { + if st.Scope != "node" || !contains(nodesOf(st.Holders), node) { + continue + } + var verbs []string + for _, v := range st.Verbs { + verbs = append(verbs, node+"/"+st.Seat+"."+v.Name) + } + var by string + for _, h := range st.Holders { + if h.Node == node { + by = h.Module + } + } + seats = append(seats, map[string]any{"seat": st.Seat, "held by": by, "verbs": verbs}) + } + var modules []map[string]any + names := make([]string, 0, len(x.Modules)) + for n := range x.Modules { + names = append(names, n) + } + sort.Strings(names) + for _, n := range names { + m := x.Modules[n] + if !contains(m.On, node) { + continue + } + var tools []map[string]string + for _, t := range m.Tools { + tools = append(tools, map[string]string{"address": node + "/" + m.Module + "." + t.Name, "does": firstLine(t.Description)}) + } + entry := map[string]any{"module": m.Module, "tools": tools} + if m.Interchangeable { + entry["interchangeable"] = "any instance answers . as well" + } + modules = append(modules, entry) + } + return answerText(map[string]any{"machine": node, "seats": seats, "modules": modules}) + + case verbSearch: + query := strings.ToLower(str("query")) + if query == "" { + return failure("mesh_search needs `query`: words to find, e.g. `postgres databases`") + } + words := strings.Fields(query) + type hit struct { + Address string `json:"address"` + Does string `json:"does"` + Also string `json:"also,omitempty"` + } + var hits []hit + matches := func(parts ...string) bool { + hay := strings.ToLower(strings.Join(parts, " ")) + for _, w := range words { + if !strings.Contains(hay, w) { + return false + } + } + return true + } + for _, st := range x.Seats { + for _, v := range st.Verbs { + if !matches(st.Seat, v.Name, v.Description) { + continue + } + if st.Scope == "node" { + on := nodesOf(st.Holders) + first := "" + also := "" + if len(on) > 0 { + first = on[0] + if len(on) > 1 { + also = "also on " + strings.Join(on[1:], ", ") + } + } + hits = append(hits, hit{first + "/" + st.Seat + "." + v.Name, firstLine(v.Description), also}) + } else { + hits = append(hits, hit{st.Seat + "." + v.Name, firstLine(v.Description), ""}) + } + } + } + names := make([]string, 0, len(x.Modules)) + for n := range x.Modules { + names = append(names, n) + } + sort.Strings(names) + for _, n := range names { + m := x.Modules[n] + for _, t := range m.Tools { + if !matches(m.Module, t.Name, t.Description) { + continue + } + switch { + case m.Interchangeable: + hits = append(hits, hit{m.Module + "." + t.Name, firstLine(t.Description), "any instance; on " + orNobody(m.On)}) + case len(m.On) == 0: + hits = append(hits, hit{"/" + m.Module + "." + t.Name, firstLine(t.Description), "the mesh places it on no machine"}) + default: + also := "" + if len(m.On) > 1 { + also = "also on " + strings.Join(m.On[1:], ", ") + } + hits = append(hits, hit{m.On[0] + "/" + m.Module + "." + t.Name, firstLine(t.Description), also}) + } + } + } + out := map[string]any{"matches": hits} + if len(hits) > searchCap { + out["matches"] = hits[:searchCap] + out["more"] = fmt.Sprintf("%d more; narrow the words", len(hits)-searchCap) + } + if len(hits) == 0 { + out["matches"] = []hit{} + out["hint"] = "nothing matched every word; try fewer words, or mesh_overview and mesh_machine to browse" + } + return answerText(out) + + case verbDescribe: + t, err := resolve(x, str("address")) + if err != nil { + return failure(err.Error()) + } + description := t.Tool.Description + if description == "" { + description = t.Name + } + return answerText(map[string]any{"address": t.Address, "description": description, + "arguments": schemaWithoutNode(t.Tool.Input)}) + + case verbCall: + t, err := resolve(x, str("address")) + if err != nil { + return failure(err.Error()) + } + callArgs := map[string]any{} + if a, ok := args["arguments"].(map[string]any); ok { + for k, v := range a { + if k != "node" { + callArgs[k] = v + } + } + } + var got bus.Answered + if t.Seat { + got, err = s.conn.Ask(t.Key, callArgs, "") + } else { + got, err = callTool(s.conn, t.Key, callArgs, nil, x.Listing) + } + if err != nil { + return failure(whyItFailed(t.Address, err)) + } + content := []map[string]any{{"type": "text", "text": pretty(got.Result)}} + if got.Node != "" { + content = append(content, map[string]any{"type": "text", "text": "answered by " + got.Node}) + } + return map[string]any{"content": content} + } + return failure("no discovery verb " + name) +} diff --git a/node-tools/internal/console/address_test.go b/node-tools/internal/console/address_test.go new file mode 100644 index 0000000..25d897b --- /dev/null +++ b/node-tools/internal/console/address_test.go @@ -0,0 +1,222 @@ +package console + +import ( + "encoding/json" + "fmt" + "strings" + "testing" + "time" + + mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" + "github.com/novox/mesh-tools/node-tools/internal/runtime" +) + +// text is a tool result's first text, and whether it was an error. +func text(t *testing.T, reply map[string]any) (string, bool) { + t.Helper() + result, ok := reply["result"].(map[string]any) + if !ok { + t.Fatalf("no result: %v", reply) + } + content := result["content"].([]any) + isErr, _ := result["isError"].(bool) + var parts []string + for _, c := range content { + parts = append(parts, c.(map[string]any)["text"].(string)) + } + return strings.Join(parts, "\n"), isErr +} + +func call(t *testing.T, endpoint, tool string, args map[string]any) (string, bool) { + t.Helper() + body, _ := json.Marshal(map[string]any{"jsonrpc": "2.0", "id": 9, "method": "tools/call", + "params": map[string]any{"name": tool, "arguments": args}}) + return text(t, post(t, endpoint, string(body))) +} + +// novox/hq ADR 0195: the console announces five tools, and everything the mesh answers is reached +// through them by one address per layer. +func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { + was := IndexKept + IndexKept = 0 // every discovery asks the mesh, so a module arriving mid-test is found + t.Cleanup(func() { IndexKept = was }) + + mesh := mt.New(t) + // alpha: interchangeable (the mesh issued it a plain subject); beta: state on its machine, holds + // the node-shelf seat there. + mesh.Issue(t, mt.MembershipOf("alpha", "desk", true, nil)) + mesh.Issue(t, mt.MembershipOf("beta", "desk", false, map[string][]string{"node-shelf": {"list", "clear"}})) + nodeTools := connect(t, "node-tools", "desk") + stop, err := runtime.Run(nodeTools, []runtime.Served{ + {Module: "alpha", Entrypoints: []string{mt.Fixture("many-alpha.serve.mjs")}}, + {Module: "beta", Entrypoints: []string{mt.Fixture("many-beta.serve.mjs")}}, + }, nil, (&mt.Logs{}).Logf) + if err != nil { + t.Fatal(err) + } + defer stop() + + catalogue := connect(t, "mesh-catalog", "") + stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) { + return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "gamma"}}}, nil + }) + defer stopCat() + + // The controller, answering as its seat verbs do: the command's printed output. + controller := connect(t, "mesh-controller", "") + out := func(s string) map[string]any { return map[string]any{"output": s, "ok": true} } + serve := func(verb string, answer func() any) { + stop, err := controller.HandleSubject("mesh.seat.mesh-controller.tool."+verb, func(json.RawMessage) (any, error) { + return answer(), nil + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stop) + } + serve("tools", func() any { + return map[string]any{"seats": []map[string]any{ + {"seat": "mesh-controller", "scope": "mesh", "tools": []map[string]any{ + {"name": "nodes", "description": "Every machine the mesh knows.", "input": map[string]any{}}}}, + {"seat": "node-shelf", "scope": "node", "tools": []map[string]any{ + {"name": "list", "description": "what is on the shelf", "input": map[string]any{}}, + {"name": "clear", "description": "take it all off", "input": map[string]any{}}}}, + }} + }) + serve("nodes", func() any { + return out("bench 3m ago converged 1f2e\ndesk here converged 9a8b\n") + }) + serve("seats", func() any { + return out("a warning the command printed first\n" + `{ + "seats": [ + {"seat": "mesh-controller", "scope": "mesh", "decision": "x", "holders": [{"module": "mesh-controller", "node": "bench"}]}, + {"seat": "node-shelf", "scope": "node", "decision": "y", "holders": [{"module": "beta", "node": "desk"}]} + ] +}`) + }) + gammaOn := "nothing" + serve("modules", func() any { + return out(fmt.Sprintf("alpha 1 built 1a2b3c4d on desk\n needs container-runtime\n"+ + "beta 1 built 1a2b3c4d on desk\n"+ + "gamma 1 built 1a2b3c4d on %s\n", gammaOn)) + }) + catalogue.Flush() + controller.Flush() + + up, err := Serve(NewSurface(nodeTools, "desk.node-tools"), "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer up.Close() + endpoint := "http://" + up.Address + "/mcp" + + // Five 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) { + name := x.(map[string]any)["name"].(string) + names = append(names, name) + for _, r := range name { + if !(r >= 'a' && r <= 'z' || r >= 'A' && r <= 'Z' || r >= '0' && r <= '9' || r == '_' || r == '-') || len(name) > 64 { + t.Errorf("%q is not a name the API takes", name) + } + } + } + if got := strings.Join(names, ","); got != "mesh_overview,mesh_machine,mesh_search,mesh_describe,mesh_call" { + t.Errorf("announced %s", got) + } + + // The overview names the mesh's seats, the machines' seats and the machines. + overview, isErr := call(t, endpoint, "mesh_overview", nil) + if isErr || !strings.Contains(overview, "mesh-controller.nodes") || !strings.Contains(overview, "/node-shelf.list") || + !strings.Contains(overview, `"bench"`) || !strings.Contains(overview, `"desk"`) { + t.Errorf("overview: %s", overview) + } + machine, isErr := call(t, endpoint, "mesh_machine", map[string]any{"node": "desk"}) + if isErr || !strings.Contains(machine, "desk/node-shelf.list") || !strings.Contains(machine, "desk/beta.three") || + !strings.Contains(machine, "desk/alpha.one") { + t.Errorf("machine: %s", machine) + } + + // A mesh seat, a node seat, an assignment and an interchangeable module, each by address. + if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "mesh-controller.nodes"}); isErr || !strings.Contains(got, "converged") { + t.Errorf("mesh seat: %s", got) + } + if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "desk/node-shelf.list"}); isErr || !strings.Contains(got, `"a"`) { + t.Errorf("node seat: %s", got) + } + if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "desk/beta.three"}); isErr || + !strings.Contains(got, `"beta": 3`) || !strings.Contains(got, "answered by desk") { + t.Errorf("assignment: %s", got) + } + if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "alpha.one"}); isErr || !strings.Contains(got, `"alpha": 1`) { + t.Errorf("interchangeable module: %s", got) + } + + // Refused, by name, where the address does not say enough or says the wrong thing. + for address, want := range map[string]string{ + "beta.three": "keeps state on each machine it runs on, so a call names the machine: write /beta.three — it runs on desk", + "node-shelf.list": "held once per machine: write /node-shelf.list — it is held on desk", + "desk/mesh-controller.nodes": "held once for the whole mesh", + "bench/beta.three": "beta does not run on bench; it runs on desk", + "nonsense": "is not an address", + } { + got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": address}) + if !isErr || !strings.Contains(got, want) { + t.Errorf("%s: %s (want %q)", address, got, want) + } + } + + // Described without `node`: the address carries the machine. + described, isErr := call(t, endpoint, "mesh_describe", map[string]any{"address": "desk/beta.three"}) + if isErr || strings.Contains(described, `"node"`) || !strings.Contains(described, `"address": "desk/beta.three"`) { + t.Errorf("describe: %s", described) + } + + // Search finds across the layers; a module that starts serving after the first answer is found. + if got, _ := call(t, endpoint, "mesh_search", map[string]any{"query": "shelf"}); !strings.Contains(got, "desk/node-shelf.list") { + t.Errorf("search a node seat: %s", got) + } + if got, _ := call(t, endpoint, "mesh_search", map[string]any{"query": "gamma"}); strings.Contains(got, "gamma.given") { + t.Fatalf("gamma was found before it served: %s", got) + } + mesh.Issue(t, mt.MembershipOf("gamma", "desk", false, nil)) + gammaOn = "desk" + late := connect(t, "node-tools", "desk") + stopLate, err := runtime.Run(late, []runtime.Served{{Module: "gamma", Entrypoints: []string{mt.Fixture("env-gamma.serve.mjs")}}}, + nil, (&mt.Logs{}).Logf) + if err != nil { + t.Fatal(err) + } + defer stopLate() + var found string + for i := 0; i < 30; i++ { + found, _ = call(t, endpoint, "mesh_search", map[string]any{"query": "gamma given"}) + if strings.Contains(found, "desk/gamma.given") { + break + } + time.Sleep(100 * time.Millisecond) + } + if !strings.Contains(found, "desk/gamma.given") { + t.Errorf("a module that arrived later was not found: %s", found) + } + + // The old names still answer, unannounced. + if got, isErr := call(t, endpoint, "alpha.one", nil); isErr || !strings.Contains(got, `"alpha": 1`) { + t.Errorf("an old name: %s", got) + } +} + +func TestTheControllersPrintedListsAreRead(t *testing.T) { + machines := machinesIn("the bus's user list leaves out 2 user(s)\nace 2m ago converged 0c1d\n" + + "g14 here converged 77aa\nnovox 5s ago adopted 3e4f\n") + if got := strings.Join(machines, ","); got != "ace,g14,novox" { + t.Errorf("machines: %s", got) + } + on := assignmentsIn("baserow 1 built 7c800705 on ace\n requires postgres-database\n" + + "confluence 1 built 7c800705 on nothing\n" + + "mesh-wireguard 1 with the control plane on ace, g14, novox, shanks\n") + if strings.Join(on["baserow"], ",") != "ace" || len(on["confluence"]) != 0 || strings.Join(on["mesh-wireguard"], ",") != "ace,g14,novox,shanks" { + t.Errorf("assignments: %v", on) + } +} diff --git a/node-tools/internal/console/console_test.go b/node-tools/internal/console/console_test.go index 0645ae9..8052ea9 100644 --- a/node-tools/internal/console/console_test.go +++ b/node-tools/internal/console/console_test.go @@ -71,7 +71,9 @@ func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) { catalogue.Flush() controller.Flush() - up, err := Serve(NewSurface(nodeTools, "desk.node-tools"), "127.0.0.1:0") + flat := NewSurface(nodeTools, "desk.node-tools") + flat.Flat = true // the whole catalogue, as before ADR 0195: still reachable, no longer announced + up, err := Serve(flat, "127.0.0.1:0") if err != nil { t.Fatal(err) } diff --git a/node-tools/internal/console/mcp.go b/node-tools/internal/console/mcp.go index 07de2b1..d26b543 100644 --- a/node-tools/internal/console/mcp.go +++ b/node-tools/internal/console/mcp.go @@ -5,6 +5,7 @@ package console import ( "encoding/json" + "os" "strings" "sync" "time" @@ -47,10 +48,17 @@ type Surface struct { mu sync.Mutex known *Listing at time.Time + idx *index + idxAt time.Time + // 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 } // NewSurface is the surface over a connection, as `who`. -func NewSurface(conn *bus.Conn, who string) *Surface { return &Surface{conn: conn, who: who} } +func NewSurface(conn *bus.Conn, who string) *Surface { + return &Surface{conn: conn, who: who, Flat: os.Getenv("MESH_CONSOLE_FLAT") == "1"} +} func (s *Surface) listing() (*Listing, error) { s.mu.Lock() @@ -99,12 +107,7 @@ func (s *Surface) Handle(r Request) *Reply { "protocolVersion": Protocol, "capabilities": map[string]any{"tools": map[string]any{}}, "serverInfo": map[string]any{"name": "mesh", "version": "1"}, - "instructions": "These are the tools of a Novox mesh, reached as " + s.who + ". Every call goes to the module " + - "that serves it; what may be called was fixed when this account was issued, so a " + - "refusal means the account, not the tool. The list is what the running modules " + - "answered, plus every role's tools from the mesh's records — the mesh's own verbs " + - "(mesh-controller.status, .push, .assign …) among them; a module that did not answer " + - "is named in the list's _meta and can still be called by ..", + "instructions": s.instructions(), }) case "notifications/initialized": return nil @@ -114,6 +117,9 @@ func (s *Surface) Handle(r Request) *Reply { } return answer(r.ID, map[string]any{}) case "tools/list": + if !s.Flat { + return answer(r.ID, map[string]any{"tools": discovery()}) + } l, err := s.listing() if err != nil { return refuse(r.ID, -32603, whyItFailed(catalogueModules, err)) @@ -142,6 +148,12 @@ func (s *Surface) Handle(r Request) *Reply { Arguments map[string]any `json:"arguments"` } _ = json.Unmarshal(r.Params, &p) + if isDiscovery(p.Name) { + if p.Arguments == nil { + p.Arguments = map[string]any{} + } + return answer(r.ID, s.discover(p.Name, p.Arguments)) + } args := map[string]any{} for k, v := range p.Arguments { args[k] = v @@ -266,3 +278,20 @@ func withNode(schema map[string]any, description string, required bool) map[stri } return out } + +// instructions is what an agent host is told about this surface when it connects. +func (s *Surface) instructions() string { + if s.Flat { + return "These are the tools of a Novox mesh, reached as " + s.who + ". Every call goes to the module " + + "that serves it; what may be called was fixed when this account was issued, so a " + + "refusal means the account, not the tool. The list is what the running modules " + + "answered, plus every role's tools from the mesh's records — the mesh's own verbs " + + "(mesh-controller.status, .push, .assign …) among them; a module that did not answer " + + "is named in the list's _meta and can still be called by .." + } + 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 " + + "was issued, so a refusal means the account, not the tool." +} -- 2.54.0 From 7722668220992614cde418be73e67c2874c58cb4 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 22:14:16 +0200 Subject: [PATCH 2/3] Every runtime announces what it serves in the NATS services protocol; the console discovers by asking the bus (hq ADR 0197) The Go runtime answers $SRV.PING, $SRV.INFO and $SRV.STATS (and per name and id) in the io.nats.micro.v1 format with what it serves at the moment it is asked: one service per runtime process, since the bus admits one reply per request from each responder, and one endpoint per tool per subject, its metadata saying module, seat, scope, machine, description, schema and whether the module is interchangeable. Serving is unchanged. The console gathers one $SRV.INFO request's answers instead of asking the catalogue's roster and each module's tools, and reads the controller's records as JSON for what should have answered: an assignment with tools that did not announce is named, a module without tools never is. The text parsers of node list and module list are gone. Packages share the test bus: go test -p 1. --- node-tools/internal/announce/announce.go | 209 ++++++++++++++ node-tools/internal/bus/bus.go | 53 ++++ node-tools/internal/console/address.go | 281 +++++++++---------- node-tools/internal/console/address_test.go | 88 +++--- node-tools/internal/console/client.go | 109 +------ node-tools/internal/console/console_test.go | 17 +- node-tools/internal/console/mcp.go | 2 +- node-tools/internal/meshtest/meshtest.go | 4 + node-tools/internal/runtime/announce_test.go | 113 ++++++++ node-tools/internal/runtime/runtime.go | 74 +++++ 10 files changed, 647 insertions(+), 303 deletions(-) create mode 100644 node-tools/internal/announce/announce.go create mode 100644 node-tools/internal/runtime/announce_test.go diff --git a/node-tools/internal/announce/announce.go b/node-tools/internal/announce/announce.go new file mode 100644 index 0000000..c9f7f97 --- /dev/null +++ b/node-tools/internal/announce/announce.go @@ -0,0 +1,209 @@ +// Package announce answers the NATS services protocol's discovery for what a runtime serves, and +// gathers the answers (novox/hq ADR 0197). +// +// A runtime does not re-serve its tools through a services library: serving is unchanged. It answers +// `$SRV.PING`, `$SRV.INFO` and `$SRV.STATS` — and the same followed by its service name, and by its +// name and id — in the format NATS's own tools read, with what it is serving at the moment it is +// asked. One service per runtime process: the bus admits one reply per request from each responder, +// so a runtime serving many modules and seats answers once, one endpoint per tool per subject, and +// says in each endpoint's metadata which module, seat, scope and machine it is. +package announce + +import ( + "encoding/json" + "strings" + "time" + + "github.com/nats-io/nats.go/micro" + + "github.com/novox/mesh-tools/node-tools/internal/bus" +) + +// 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. +var Window = 750 * time.Millisecond + +// Kinds of endpoint. +const ( + KindTool = "tool" // a module's own tool + KindSeat = "seat" // a seat's verb, served by the module holding it +) + +// Endpoint is one tool served on one subject, as it is announced. +type Endpoint struct { + Kind string + Module string // the module whose code answers + Tool string // the tool's or the verb's name + Seat string // for a seat's verb + Scope string // "mesh" or "node", for a seat's verb + Node string + Description string + Schema json.RawMessage + Interchangeable bool + Subject string + Queue string +} + +// Name is the endpoint's name as the protocol allows it — letters, digits, `-` and `_` — the +// prefix and the tool joined by `__`; the metadata, not the name, is what identifies it. +func (e Endpoint) Name() string { + prefix := e.Module + if e.Kind == KindSeat { + prefix = e.Seat + } + return clean(prefix) + "__" + clean(e.Tool) +} + +func clean(s string) string { + var b strings.Builder + for _, r := range s { + if r == '-' || r == '_' || (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') { + b.WriteRune(r) + } else { + b.WriteRune('_') + } + } + return b.String() +} + +func (e Endpoint) info() micro.EndpointInfo { + md := map[string]string{ + "kind": e.Kind, "module": e.Module, "tool": e.Tool, "node": e.Node, + "description": e.Description, "interchangeable": boolWord(e.Interchangeable), + } + schema := strings.TrimSpace(string(e.Schema)) + if schema == "" || schema == "null" { + schema = "{}" + } + md["schema"] = schema + if e.Kind == KindSeat { + md["seat"] = e.Seat + md["scope"] = e.Scope + } + return micro.EndpointInfo{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue, Metadata: md} +} + +func boolWord(b bool) string { + if b { + return "true" + } + return "false" +} + +// Service is who answers: the runtime's name, its instance, and a word about it. +type Service struct { + Name string + ID string + Description string + Metadata map[string]string +} + +func (s Service) identity() micro.ServiceIdentity { + md := s.Metadata + if md == nil { + md = map[string]string{} + } + return micro.ServiceIdentity{Name: s.Name, ID: s.ID, Version: Version, Metadata: md} +} + +// Info is the info_response for these endpoints. +func Info(s Service, endpoints []Endpoint) micro.Info { + out := micro.Info{ServiceIdentity: s.identity(), Type: micro.InfoResponseType, + Description: s.Description, Endpoints: []micro.EndpointInfo{}} + for _, e := range endpoints { + out.Endpoints = append(out.Endpoints, e.info()) + } + return out +} + +// Serve answers discovery for one service until stopped, asking `current` for its endpoints each +// time — so what is announced is what is served now, re-served memberships included. +func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error) { + started := time.Now().UTC() + answer := func(subject string, _ []byte) []byte { + parts := strings.Split(subject, ".") + if len(parts) < 2 || parts[0] != "$SRV" { + return nil + } + if len(parts) >= 3 && parts[2] != s.Name { + return nil // another service's + } + if len(parts) >= 4 && parts[3] != s.ID { + return nil // another instance's + } + var v any + switch parts[1] { + case "PING": + v = micro.Ping{ServiceIdentity: s.identity(), Type: micro.PingResponseType} + case "INFO": + v = Info(s, current()) + case "STATS": + st := micro.Stats{ServiceIdentity: s.identity(), Type: micro.StatsResponseType, Started: started, + Endpoints: []*micro.EndpointStats{}} + for _, e := range current() { + st.Endpoints = append(st.Endpoints, µ.EndpointStats{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue}) + } + v = st + default: + return nil + } + body, err := json.Marshal(v) + if err != nil { + return nil + } + return body + } + var stops []func() + for _, verb := range []string{"PING", "INFO", "STATS"} { + for _, subject := range []string{"$SRV." + verb, "$SRV." + verb + ".>"} { + stop, err := conn.Raw(subject, answer) + if err != nil { + for _, st := range stops { + st() + } + return func() {}, err + } + stops = append(stops, stop) + } + } + return func() { + for _, st := range stops { + st() + } + }, 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 + } + var out []micro.Info + for _, b := range raw { + var i micro.Info + if json.Unmarshal(b, &i) != nil || i.Type != micro.InfoResponseType { + continue + } + out = append(out, i) + } + return out, nil +} + +// Endpoints reads an info_response's endpoints back into what they announce. +func Endpoints(i micro.Info) []Endpoint { + var out []Endpoint + for _, e := range i.Endpoints { + md := e.Metadata + out = append(out, Endpoint{ + Kind: md["kind"], Module: md["module"], Tool: md["tool"], Seat: md["seat"], Scope: md["scope"], + Node: md["node"], Description: md["description"], Schema: json.RawMessage(md["schema"]), + Interchangeable: md["interchangeable"] == "true", Subject: e.Subject, Queue: e.QueueGroup, + }) + } + return out +} diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index 8fddf95..608dadc 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -512,6 +512,59 @@ func (c *Conn) PublishAs(module string, env Envelope) error { return err } +// ServedOn is where a served module's tool is answered right now: the membership's subjects when +// issued, the derived shape otherwise — what Handle subscribes, for what announces it (ADR 0197). +func (c *Conn) ServedOn(module, tool string) []Served { return c.servedOn(module, tool) } + +// Raw answers one subject with a function of the request, not a tool's reply envelope: the NATS +// services protocol's discovery subjects answer in their own format (novox/hq ADR 0197). A nil +// answer is no reply — the request was for another service. +func (c *Conn) Raw(subject string, answer func(subject string, data []byte) []byte) (func(), error) { + sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) { + if body := answer(msg.Subject, msg.Data); body != nil { + _ = msg.Respond(body) + } + }) + if err != nil { + return func() {}, err + } + c.track(sub) + 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) { + inbox := c.nc.NewRespInbox() + sub, err := c.nc.SubscribeSync(inbox) + if err != nil { + return nil, err + } + 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) + } + } +} + // Flush waits until the bus has every subscription made so far, so what is served is answerable // when this returns. func (c *Conn) Flush() { _ = c.nc.Flush() } diff --git a/node-tools/internal/console/address.go b/node-tools/internal/console/address.go index 3909b59..f8b2cb4 100644 --- a/node-tools/internal/console/address.go +++ b/node-tools/internal/console/address.go @@ -19,12 +19,12 @@ package console import ( "encoding/json" "fmt" - "regexp" "sort" "strings" "sync" "time" + "github.com/novox/mesh-tools/node-tools/internal/announce" "github.com/novox/mesh-tools/node-tools/internal/bus" ) @@ -159,188 +159,181 @@ func controllerOutput(conn *bus.Conn, verb string) (string, error) { // jsonIn is the JSON document a command printed, after any lines it said first: a seat verb runs // the controller's command, and a command may warn before it answers. func jsonIn(output string) string { - if strings.HasPrefix(strings.TrimSpace(output), "{") { - return strings.TrimSpace(output) + t := strings.TrimSpace(output) + if strings.HasPrefix(t, "{") || strings.HasPrefix(t, "[") { + return t } - if i := strings.Index(output, "\n{"); i >= 0 { - return strings.TrimSpace(output[i+1:]) + for _, open := range []string{"\n{", "\n["} { + if i := strings.Index(output, open); i >= 0 { + return strings.TrimSpace(output[i+1:]) + } } return "" } -// A machine as `node list` prints it: its name, when it was last heard from, its mode — converged -// or adopted, the only two the command prints — and its id. Lines the command says around them -// (a warning about the bus's users, "no node records yet") are not machines and are skipped. -var nodeLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+.*\s(converged|adopted)\s+\S+\s*$`) - -// machinesIn reads the machines from `node list`'s output. -func machinesIn(output string) []string { - var out []string - for _, line := range strings.Split(output, "\n") { - if m := nodeLine.FindStringSubmatch(line); m != nil { - out = append(out, m[1]) - } - } - sort.Strings(out) - return out +// recordedModule is a module as the controller's records hold it (`module list --json`): where the +// mesh assigned it, and whether it declares tools — what should announce itself, and where. +type recordedModule struct { + Module string `json:"module"` + On []string `json:"on"` + Tools bool `json:"tools"` } -// A module as `module list` prints it: name, version, how it was built, and where it runs. -var moduleLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+\S+\s+.*?\s+on (.+)$`) - -// assignmentsIn reads, from `module list`'s output, the machines each module runs on. -func assignmentsIn(output string) map[string][]string { - out := map[string][]string{} - for _, line := range strings.Split(output, "\n") { - if line == "" || line[0] == ' ' || line[0] == '\t' { - continue - } - m := moduleLine.FindStringSubmatch(strings.TrimRight(line, " ")) - if m == nil { - continue - } - on := strings.TrimSpace(m[2]) - if on == "nothing" { - out[m[1]] = []string{} - continue - } - var nodes []string - for _, n := range strings.Split(on, ",") { - if n = strings.TrimSpace(n); n != "" { - nodes = append(nodes, n) - } - } - sort.Strings(nodes) - out[m[1]] = nodes - } - return out +// recordedMachine is a machine as the controller's records hold it (`node list --json`). +type recordedMachine struct { + Name string `json:"name"` } -// interchangeable is whether the mesh issued the module a plain subject for this tool — one any of -// its instances answers (ADR 0160): the module's own subject with no machine after it. -func interchangeable(module string, t Tool) bool { - plain := "mesh.mod." + module + ".tool." + t.Name - for _, s := range t.Subjects { - if s == plain { - return true - } - } - return false -} - -// indexOn asks the mesh what it holds: the flat listing (catalogue, modules, the seats' tools), -// and from the controller the seats' holders, the machines and the assignments — at once. +// indexOn asks the mesh what it holds (novox/hq ADR 0197): what answers, from every runtime's own +// announcement on the bus — one `$SRV.INFO` request — and what should, from the controller's records +// read as JSON. Nothing is inferred from a roster and nothing is parsed from print. func indexOn(conn *bus.Conn) (*index, error) { var wg sync.WaitGroup - var seatsOut, nodesOut, modulesOut string - wg.Add(3) - go func() { defer wg.Done(); seatsOut, _ = controllerOutput(conn, "seats") }() + var nodesOut, modulesOut string + wg.Add(2) go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }() go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }() - l, err := toolsOn(conn) + infos, err := announce.Gather(conn) wg.Wait() if err != nil { return nil, err } - x := &index{Modules: map[string]*moduleInfo{}, NotAnswering: l.NotAnswering, Listing: l} - - // The seats: their verbs from the mesh's records, their holders from the controller. - var held struct { - Seats []struct { - Seat string `json:"seat"` - Scope string `json:"scope"` - Holders []holder `json:"holders"` - } `json:"seats"` - } - _ = json.Unmarshal([]byte(jsonIn(seatsOut)), &held) - holders := map[string][]holder{} - scopes := map[string]string{} - for _, s := range held.Seats { - holders[s.Seat] = s.Holders - scopes[s.Seat] = s.Scope - } - bySeat := map[string]*seatInfo{} - for _, t := range l.Tools { - if !t.Seat { - continue - } - s := bySeat[t.Module] - if s == nil { - scope := t.Scope - if scope == "" { - scope = scopes[t.Module] + l := &Listing{Tools: []Tool{}, NotAnswering: []string{}} + x := &index{Modules: map[string]*moduleInfo{}, Listing: l} + 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) { + if e.Node != "" { + machines[e.Node] = true } - if scope != "node" { - scope = "mesh" + if announced[e.Module] == nil { + announced[e.Module] = map[string]bool{} } - s = &seatInfo{Seat: t.Module, Scope: scope, Holders: holders[t.Module]} - bySeat[t.Module] = s - } - s.Verbs = append(s.Verbs, t) - } - for _, s := range bySeat { - x.Seats = append(x.Seats, *s) - } - sort.Slice(x.Seats, func(i, j int) bool { return x.Seats[i].Seat < x.Seats[j].Seat }) - - // The machines: what the controller knows, and any that hold a seat. - seen := map[string]bool{} - for _, n := range machinesIn(nodesOut) { - seen[n] = true - } - for _, s := range x.Seats { - for _, h := range s.Holders { - if h.Node != "" { - seen[h.Node] = true + announced[e.Module][e.Node] = true + switch e.Kind { + case announce.KindSeat: + st := seats[e.Seat] + if st == nil { + st = &seatInfo{Seat: e.Seat, Scope: e.Scope} + seats[e.Seat] = st + } + if st.Scope != "node" && e.Scope == "node" { + st.Scope = "node" + } + h := holder{Module: e.Module, Node: e.Node} + if !containsHolder(st.Holders, h) { + st.Holders = append(st.Holders, h) + } + key := e.Seat + "." + e.Tool + if _, have := toolAt[key]; !have { + toolAt[key] = len(l.Tools) + t := Tool{Module: e.Seat, Name: e.Tool, Description: e.Description, Input: e.Schema, Seat: true, Scope: st.Scope} + l.Tools = append(l.Tools, t) + st.Verbs = append(st.Verbs, t) + } + case announce.KindTool: + m := x.Modules[e.Module] + if m == nil { + m = &moduleInfo{Module: e.Module} + x.Modules[e.Module] = m + } + if e.Node != "" && !contains(m.On, e.Node) { + m.On = append(m.On, e.Node) + } + m.Interchangeable = m.Interchangeable || e.Interchangeable + key := e.Module + "." + e.Tool + i, have := toolAt[key] + if !have { + i = len(l.Tools) + toolAt[key] = i + l.Tools = append(l.Tools, Tool{Module: e.Module, Name: e.Tool, Description: e.Description, Input: e.Schema}) + } + if !contains(l.Tools[i].Subjects, e.Subject) { + l.Tools[i].Subjects = append(l.Tools[i].Subjects, e.Subject) + } } } } - - // The modules that answer tools, where they run, and whether any instance will do. - on := assignmentsIn(modulesOut) - for _, t := range l.Tools { + // A tool's subjects as a call looks them up: the plain one any instance answers first, then each + // machine's. + for i := range l.Tools { + t := &l.Tools[i] if t.Seat { continue } - m := x.Modules[t.Module] - if m == nil { - m = &moduleInfo{Module: t.Module, On: on[t.Module]} - x.Modules[t.Module] = m - } - m.Tools = append(m.Tools, t) - if interchangeable(t.Module, t) { - m.Interchangeable = true + plain := "mesh.mod." + t.Module + ".tool." + t.Name + sort.SliceStable(t.Subjects, func(a, b int) bool { + if (t.Subjects[a] == plain) != (t.Subjects[b] == plain) { + return t.Subjects[a] == plain + } + return t.Subjects[a] < t.Subjects[b] + }) + } + for _, t := range l.Tools { + if !t.Seat { + x.Modules[t.Module].Tools = append(x.Modules[t.Module].Tools, t) } } for _, m := range x.Modules { - // A module the controller could not place is placed where its own answer says it runs. - if len(m.On) == 0 { - nodes := map[string]bool{} - for _, t := range m.Tools { - for _, s := range t.Subjects { - base := "mesh.mod." + m.Module + ".tool." + t.Name + "." - if strings.HasPrefix(s, base) { - nodes[strings.TrimPrefix(s, base)] = true - } + sort.Strings(m.On) + } + for _, st := range seats { + sort.Slice(st.Holders, func(a, b int) bool { return st.Holders[a].Node < st.Holders[b].Node }) + x.Seats = append(x.Seats, *st) + } + sort.Slice(x.Seats, func(i, j int) bool { return x.Seats[i].Seat < x.Seats[j].Seat }) + sort.SliceStable(l.Tools, func(i, j int) bool { + return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name + }) + + // What should have answered: every assignment of a module that declares tools. Silence is named; + // a module with no tools is never a name here. + var recorded []recordedModule + if json.Unmarshal([]byte(jsonIn(modulesOut)), &recorded) == nil { + for _, m := range recorded { + if !m.Tools { + continue + } + for _, n := range m.On { + if !announced[m.Module][n] { + l.NotAnswering = append(l.NotAnswering, m.Module+" on "+n) } } - for n := range nodes { - m.On = append(m.On, n) - } - sort.Strings(m.On) } - for _, n := range m.On { - seen[n] = true + } else { + l.NotAnswering = append(l.NotAnswering, "mesh-controller (its records of the modules did not answer, so what is missing cannot be said)") + } + sort.Strings(l.NotAnswering) + x.NotAnswering = l.NotAnswering + + var known []recordedMachine + if json.Unmarshal([]byte(jsonIn(nodesOut)), &known) == nil { + for _, n := range known { + if n.Name != "" { + machines[n.Name] = true + } } } - for n := range seen { + for n := range machines { x.Machines = append(x.Machines, n) } sort.Strings(x.Machines) return x, nil } +func containsHolder(hs []holder, h holder) bool { + for _, x := range hs { + if x == h { + return true + } + } + return false +} + func (s *Surface) index() (*index, error) { s.mu.Lock() if s.idx != nil && time.Since(s.idxAt) <= IndexKept { @@ -520,7 +513,7 @@ func failure(text string) map[string]any { func (s *Surface) discover(name string, args map[string]any) map[string]any { x, err := s.index() if err != nil { - return failure(whyItFailed(catalogueModules, err)) + return failure("the mesh's discovery failed: " + err.Error()) } str := func(k string) string { v, _ := args[k].(string); return strings.TrimSpace(v) } diff --git a/node-tools/internal/console/address_test.go b/node-tools/internal/console/address_test.go index 25d897b..f2a3208 100644 --- a/node-tools/internal/console/address_test.go +++ b/node-tools/internal/console/address_test.go @@ -4,9 +4,11 @@ import ( "encoding/json" "fmt" "strings" + "sync/atomic" "testing" "time" + "github.com/novox/mesh-tools/node-tools/internal/announce" mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" "github.com/novox/mesh-tools/node-tools/internal/runtime" ) @@ -56,14 +58,20 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { } defer stop() - catalogue := connect(t, "mesh-catalog", "") - stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) { - return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "gamma"}}}, nil - }) - defer stopCat() + // Nothing may be asked for a roster or a module's `tools` any more (ADR 0197): counted. + var asked atomic.Int32 + watcher := connect(t, "watcher", "") + for _, subject := range []string{"mesh.mod.*.tool.tools", "mesh.mod.*.tool.tools.*", "mesh.mod.mesh-catalog.>"} { + stop, err := watcher.Raw(subject, func(string, []byte) []byte { asked.Add(1); return nil }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stop) + } - // The controller, answering as its seat verbs do: the command's printed output. - controller := connect(t, "mesh-controller", "") + // The controller: its records as JSON, as its seat verbs answer them, and its own seat's verbs + // announced on the bus like every runtime's. + controller := connect(t, "mesh-controller", "bench") out := func(s string) map[string]any { return map[string]any{"output": s, "ok": true} } serve := func(verb string, answer func() any) { stop, err := controller.HandleSubject("mesh.seat.mesh-controller.tool."+verb, func(json.RawMessage) (any, error) { @@ -74,34 +82,28 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { } t.Cleanup(stop) } - serve("tools", func() any { - return map[string]any{"seats": []map[string]any{ - {"seat": "mesh-controller", "scope": "mesh", "tools": []map[string]any{ - {"name": "nodes", "description": "Every machine the mesh knows.", "input": map[string]any{}}}}, - {"seat": "node-shelf", "scope": "node", "tools": []map[string]any{ - {"name": "list", "description": "what is on the shelf", "input": map[string]any{}}, - {"name": "clear", "description": "take it all off", "input": map[string]any{}}}}, - }} - }) serve("nodes", func() any { - return out("bench 3m ago converged 1f2e\ndesk here converged 9a8b\n") + return out("the bus's user list leaves out 2 user(s)\n" + + `[{"name":"bench","heard":"3m ago","mode":"converged","id":"1f2e"},{"name":"desk","heard":"here","mode":"converged","id":"9a8b"}]`) }) - serve("seats", func() any { - return out("a warning the command printed first\n" + `{ - "seats": [ - {"seat": "mesh-controller", "scope": "mesh", "decision": "x", "holders": [{"module": "mesh-controller", "node": "bench"}]}, - {"seat": "node-shelf", "scope": "node", "decision": "y", "holders": [{"module": "beta", "node": "desk"}]} - ] -}`) - }) - gammaOn := "nothing" + var gammaOn atomic.Value + gammaOn.Store("[]") serve("modules", func() any { - return out(fmt.Sprintf("alpha 1 built 1a2b3c4d on desk\n needs container-runtime\n"+ - "beta 1 built 1a2b3c4d on desk\n"+ - "gamma 1 built 1a2b3c4d on %s\n", gammaOn)) + return out(fmt.Sprintf(`[{"module":"alpha","on":["desk"],"tools":true},{"module":"beta","on":["desk"],"tools":true},`+ + `{"module":"gamma","on":%s,"tools":true},{"module":"delta","on":["desk"],"tools":false},`+ + `{"module":"epsilon","on":["bench"],"tools":true}]`, gammaOn.Load().(string))) }) - catalogue.Flush() + stopAnn, err := announce.Serve(controller, announce.Service{Name: "mesh-controller", ID: "bench"}, func() []announce.Endpoint { + return []announce.Endpoint{{Kind: announce.KindSeat, Module: "mesh-controller", Seat: "mesh-controller", Scope: "mesh", + Tool: "nodes", Node: "bench", Description: "Every machine the mesh knows.", Schema: json.RawMessage(`{}`), + Subject: "mesh.seat.mesh-controller.tool.nodes"}} + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stopAnn) controller.Flush() + watcher.Flush() up, err := Serve(NewSurface(nodeTools, "desk.node-tools"), "127.0.0.1:0") if err != nil { @@ -132,6 +134,11 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { !strings.Contains(overview, `"bench"`) || !strings.Contains(overview, `"desk"`) { t.Errorf("overview: %s", overview) } + // Silence is named only where tools should have answered: epsilon declares tools on bench and + // nothing there announced it; delta declares none and is never a name. + if !strings.Contains(overview, "epsilon on bench") || strings.Contains(overview, "delta") { + t.Errorf("not answering: %s", overview) + } machine, isErr := call(t, endpoint, "mesh_machine", map[string]any{"node": "desk"}) if isErr || !strings.Contains(machine, "desk/node-shelf.list") || !strings.Contains(machine, "desk/beta.three") || !strings.Contains(machine, "desk/alpha.one") { @@ -181,7 +188,7 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { t.Fatalf("gamma was found before it served: %s", got) } mesh.Issue(t, mt.MembershipOf("gamma", "desk", false, nil)) - gammaOn = "desk" + gammaOn.Store(`["desk"]`) late := connect(t, "node-tools", "desk") stopLate, err := runtime.Run(late, []runtime.Served{{Module: "gamma", Entrypoints: []string{mt.Fixture("env-gamma.serve.mjs")}}}, nil, (&mt.Logs{}).Logf) @@ -201,22 +208,21 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) { t.Errorf("a module that arrived later was not found: %s", found) } + if n := asked.Load(); n != 0 { + t.Errorf("discovery asked a roster or a module's tools %d time(s); it asks the bus", n) + } + // The old names still answer, unannounced. if got, isErr := call(t, endpoint, "alpha.one", nil); isErr || !strings.Contains(got, `"alpha": 1`) { t.Errorf("an old name: %s", got) } } -func TestTheControllersPrintedListsAreRead(t *testing.T) { - machines := machinesIn("the bus's user list leaves out 2 user(s)\nace 2m ago converged 0c1d\n" + - "g14 here converged 77aa\nnovox 5s ago adopted 3e4f\n") - if got := strings.Join(machines, ","); got != "ace,g14,novox" { - t.Errorf("machines: %s", got) +func TestTheControllersRecordsAreReadAsJSON(t *testing.T) { + if got := jsonIn("a warning printed first\n[{\"name\":\"ace\"}]"); got != `[{"name":"ace"}]` { + t.Errorf("an array after a warning: %q", got) } - on := assignmentsIn("baserow 1 built 7c800705 on ace\n requires postgres-database\n" + - "confluence 1 built 7c800705 on nothing\n" + - "mesh-wireguard 1 with the control plane on ace, g14, novox, shanks\n") - if strings.Join(on["baserow"], ",") != "ace" || len(on["confluence"]) != 0 || strings.Join(on["mesh-wireguard"], ",") != "ace,g14,novox,shanks" { - t.Errorf("assignments: %v", on) + if got := jsonIn(`{"seats":[]}`); got != `{"seats":[]}` { + t.Errorf("an object: %q", got) } } diff --git a/node-tools/internal/console/client.go b/node-tools/internal/console/client.go index 9cd8c0d..f47a805 100644 --- a/node-tools/internal/console/client.go +++ b/node-tools/internal/console/client.go @@ -5,19 +5,11 @@ import ( "errors" "fmt" "regexp" - "sort" "strings" - "sync" "github.com/novox/mesh-tools/node-tools/internal/bus" ) -const ( - catalogueModules = "mesh-catalog.catalog_modules" - seatTools = "seat:mesh-controller.tools" - toolsVerb = "tools" -) - // Tool is a tool as its module — or, for a role's tool, the mesh's records — describes it. type Tool struct { Module string @@ -67,107 +59,14 @@ func toolKey(name string, seats Seats) string { return name } -// toolsOn asks the mesh what tools it has: the catalogue which modules it holds, each module what it -// serves, the controller's seat every role's tools — at once, so a restarting control plane hides -// nothing else. +// toolsOn is the flat catalogue (MESH_CONSOLE_FLAT=1): what announced itself on the bus, every +// tool and seat verb, and the assignments with tools that did not (novox/hq ADR 0197). func toolsOn(conn *bus.Conn) (*Listing, error) { - type rolesAnswer struct { - Seats []struct { - Seat string `json:"seat"` - Scope string `json:"scope"` - Tools []struct { - Name string `json:"name"` - Description string `json:"description"` - Input json.RawMessage `json:"input"` - } `json:"tools"` - } `json:"seats"` - } - var wg sync.WaitGroup - var roles *rolesAnswer - wg.Add(1) - go func() { - defer wg.Done() - if got, err := conn.Ask(seatTools, map[string]any{}, ""); err == nil { - var r rolesAnswer - if json.Unmarshal(got.Result, &r) == nil { - roles = &r - } - } - }() - answered, err := conn.Ask(catalogueModules, map[string]any{}, "") + x, err := indexOn(conn) if err != nil { - wg.Wait() return nil, err } - var held struct { - Modules []struct { - Module string `json:"module"` - } `json:"modules"` - } - _ = json.Unmarshal(answered.Result, &held) - names := make([]string, 0, len(held.Modules)) - for _, m := range held.Modules { - if m.Module != "" { - names = append(names, m.Module) - } - } - type outcome struct { - ok bool - answer struct { - Tools *[]struct { - Name string `json:"name"` - Description string `json:"description"` - Input json.RawMessage `json:"input"` - Subjects []string `json:"subjects"` - } `json:"tools"` - Failed *string `json:"failed"` - } - } - outcomes := make([]outcome, len(names)) - for i, module := range names { - wg.Add(1) - go func(i int, module string) { - defer wg.Done() - got, err := conn.Ask(module+"."+toolsVerb, map[string]any{}, "") - if err != nil { - return - } - if json.Unmarshal(got.Result, &outcomes[i].answer) == nil { - outcomes[i].ok = true - } - }(i, module) - } - wg.Wait() - l := &Listing{Tools: []Tool{}, NotAnswering: []string{}} - if roles != nil { - for _, s := range roles.Seats { - for _, t := range s.Tools { - l.Tools = append(l.Tools, Tool{Module: s.Seat, Name: t.Name, Description: t.Description, - Input: t.Input, Seat: true, Scope: s.Scope}) - } - } - } else { - l.NotAnswering = append(l.NotAnswering, "mesh-controller (seat)") - } - for i, module := range names { - o := outcomes[i] - switch { - case o.ok && o.answer.Failed != nil: - l.NotAnswering = append(l.NotAnswering, fmt.Sprintf("%s (its tools bundle failed to load: %s)", module, *o.answer.Failed)) - case o.ok && o.answer.Tools != nil: - for _, t := range *o.answer.Tools { - l.Tools = append(l.Tools, Tool{Module: module, Name: t.Name, Description: t.Description, - Input: t.Input, Subjects: t.Subjects}) - } - default: - l.NotAnswering = append(l.NotAnswering, module) - } - } - sort.SliceStable(l.Tools, func(i, j int) bool { - return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name - }) - sort.Strings(l.NotAnswering) - return l, nil + return x.Listing, nil } // callTool calls `.[@]`, on the subject the listing names for it when it names one. diff --git a/node-tools/internal/console/console_test.go b/node-tools/internal/console/console_test.go index 8052ea9..d41b55f 100644 --- a/node-tools/internal/console/console_test.go +++ b/node-tools/internal/console/console_test.go @@ -55,20 +55,13 @@ func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) { } defer stop() - catalogue := connect(t, "mesh-catalog", "") - stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) { - return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "ghost"}}}, nil - }) - defer stopCat() + // The controller's records (ADR 0197): ghost declares tools on desk and nothing announces it. controller := connect(t, "mesh-controller", "") - stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.tools", func(json.RawMessage) (any, error) { - return map[string]any{"seats": []map[string]any{{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{ - {"name": "list", "description": "what is on the shelf", "input": map[string]any{}}, - {"name": "clear", "description": "take it all off", "input": map[string]any{}}, - }}}}, nil + stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.modules", func(json.RawMessage) (any, error) { + return map[string]any{"ok": true, "output": `[{"module":"alpha","on":["desk"],"tools":true},` + + `{"module":"beta","on":["desk"],"tools":true},{"module":"ghost","on":["desk"],"tools":true}]`}, nil }) defer stopSeat() - catalogue.Flush() controller.Flush() flat := NewSurface(nodeTools, "desk.node-tools") @@ -92,7 +85,7 @@ func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) { if got := strings.Join(names, ","); got != "alpha.one,alpha.two,beta.five,beta.four,beta.three,node-shelf.clear,node-shelf.list" { t.Errorf("listed %s", got) } - if got := listed["_meta"].(map[string]any)["notAnswering"]; len(got.([]any)) != 1 || got.([]any)[0] != "ghost" { + if got := listed["_meta"].(map[string]any)["notAnswering"]; len(got.([]any)) != 1 || got.([]any)[0] != "ghost on desk" { t.Errorf("not answering: %v", got) } for _, x := range listed["tools"].([]any) { diff --git a/node-tools/internal/console/mcp.go b/node-tools/internal/console/mcp.go index d26b543..330e2f1 100644 --- a/node-tools/internal/console/mcp.go +++ b/node-tools/internal/console/mcp.go @@ -122,7 +122,7 @@ func (s *Surface) Handle(r Request) *Reply { } l, err := s.listing() if err != nil { - return refuse(r.ID, -32603, whyItFailed(catalogueModules, err)) + return refuse(r.ID, -32603, "the mesh's discovery failed: "+err.Error()) } tools := make([]map[string]any, 0, len(l.Tools)) for _, t := range l.Tools { diff --git a/node-tools/internal/meshtest/meshtest.go b/node-tools/internal/meshtest/meshtest.go index 185eb0c..7342c0d 100644 --- a/node-tools/internal/meshtest/meshtest.go +++ b/node-tools/internal/meshtest/meshtest.go @@ -1,5 +1,9 @@ // Package meshtest raises what the controller would, for tests against a real bus: the ASSIGNMENTS // and EVENTS streams, memberships issued by hand, and a fixture's path. +// +// **The packages share one bus, so run them one at a time: `go test -p 1 ./...`.** Each test raises +// the streams afresh, and the console discovers every runtime that announces itself on the bus +// (novox/hq ADR 0197) — a runtime from another package's test is, correctly, found. package meshtest import ( diff --git a/node-tools/internal/runtime/announce_test.go b/node-tools/internal/runtime/announce_test.go new file mode 100644 index 0000000..69408d7 --- /dev/null +++ b/node-tools/internal/runtime/announce_test.go @@ -0,0 +1,113 @@ +package runtime + +import ( + "encoding/json" + "sort" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" + "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" +) + +// novox/hq ADR 0197: what a runtime serves, it announces — the NATS services protocol's discovery, +// read here with NATS's own types, each endpoint a subject actually served — and the announcement +// follows a re-issued membership. +func TestTheRuntimeAnnouncesWhatItServesInTheServicesProtocol(t *testing.T) { + mesh := mt.New(t) + mesh.Issue(t, mt.MembershipOf("alpha", "anchor", false, nil)) + mesh.Issue(t, mt.MembershipOf("beta", "anchor", false, map[string][]string{"node-shelf": {"list", "clear"}})) + nodeTools := connect(t, "node-tools", "anchor") + stop, err := Run(nodeTools, []Served{ + {"alpha", []string{mt.Fixture("many-alpha.serve.mjs")}}, + {"beta", []string{mt.Fixture("many-beta.serve.mjs")}}, + }, nil, (&mt.Logs{}).Logf) + if err != nil { + t.Fatal(err) + } + defer stop() + + nc, err := nats.Connect(mt.URL(t)) + if err != nil { + t.Fatal(err) + } + defer nc.Close() + info := func() micro.Info { + t.Helper() + msg, err := nc.Request("$SRV.INFO.node-tools.anchor", nil, 2*time.Second) + if err != nil { + t.Fatal(err) + } + var i micro.Info + if err := json.Unmarshal(msg.Data, &i); err != nil { + t.Fatal(err) + } + return i + } + subjects := func(i micro.Info) string { + var out []string + for _, e := range i.Endpoints { + out = append(out, e.Subject+"|"+e.QueueGroup+"|"+e.Metadata["kind"]+"|"+e.Metadata["module"]+"|"+e.Metadata["seat"]+"|"+e.Metadata["scope"]+"|"+e.Metadata["interchangeable"]) + } + sort.Strings(out) + return strings.Join(out, "\n") + } + + got := info() + if got.Type != micro.InfoResponseType || got.Name != "node-tools" || got.ID != "anchor" || got.Version == "" { + t.Errorf("identity: %+v", got.ServiceIdentity) + } + want := strings.Join([]string{ + "mesh.mod.alpha.tool.one.anchor||tool|alpha|||false", + "mesh.mod.alpha.tool.two.anchor||tool|alpha|||false", + "mesh.mod.beta.tool.five.anchor||tool|beta|||false", + "mesh.mod.beta.tool.four.anchor||tool|beta|||false", + "mesh.mod.beta.tool.three.anchor||tool|beta|||false", + "mesh.seat.node-shelf.tool.clear.anchor||seat|beta|node-shelf|node|false", + "mesh.seat.node-shelf.tool.list.anchor||seat|beta|node-shelf|node|false", + }, "\n") + if s := subjects(got); s != want { + t.Errorf("announced:\n%s\nwant:\n%s", s, want) + } + for _, e := range got.Endpoints { + if e.Metadata["node"] != "anchor" || e.Metadata["description"] == "" || !json.Valid([]byte(e.Metadata["schema"])) { + t.Errorf("endpoint metadata: %+v", e) + } + } + + // Every subject announced is answered. + asker := connect(t, "console", "workstation") + for _, e := range got.Endpoints { + if _, err := asker.Ask("", map[string]any{}, e.Subject); err != nil { + t.Errorf("announced %s and does not answer it: %v", e.Subject, err) + } + } + + // PING answers with the same identity, and a request for another service is not answered. + if msg, err := nc.Request("$SRV.PING", nil, 2*time.Second); err != nil || !strings.Contains(string(msg.Data), micro.PingResponseType) { + t.Errorf("ping: %v %v", msg, err) + } + if _, err := nc.Request("$SRV.INFO.somebody-else", nil, 300*time.Millisecond); err == nil { + t.Error("answered a request for another service") + } + + // alpha is re-issued a plain subject: the announcement says so, and says it is interchangeable. + mesh.Issue(t, mt.MembershipOf("alpha", "anchor", true, nil)) + var after string + for i := 0; i < 40; i++ { + after = subjects(info()) + if strings.Contains(after, "mesh.mod.alpha.tool.one|serve.alpha|tool|alpha|||true") { + break + } + time.Sleep(100 * time.Millisecond) + } + if !strings.Contains(after, "mesh.mod.alpha.tool.one|serve.alpha|tool|alpha|||true") || + !strings.Contains(after, "mesh.mod.alpha.tool.one.anchor||tool|alpha|||true") { + t.Errorf("after a re-issued membership:\n%s", after) + } + _ = bus.Served{} +} diff --git a/node-tools/internal/runtime/runtime.go b/node-tools/internal/runtime/runtime.go index 674d2df..9b616e9 100644 --- a/node-tools/internal/runtime/runtime.go +++ b/node-tools/internal/runtime/runtime.go @@ -14,6 +14,7 @@ import ( "strings" "sync" + "github.com/novox/mesh-tools/node-tools/internal/announce" "github.com/novox/mesh-tools/node-tools/internal/bus" "github.com/novox/mesh-tools/node-tools/internal/launch" ) @@ -276,10 +277,83 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log } logf("%s", line) stops = append(stops, serveSeats(conn, modules, registrations, logf)) + + // **What it serves, it announces** (novox/hq ADR 0197): the NATS services protocol's discovery, + // answered with what is served at the moment it is asked — re-served memberships included. + announced, err := announce.Serve(conn, announce.Service{ + Name: conn.Module(), ID: instanceOf(conn), + Description: "the mesh's tool runtime on " + node + ": every assigned module's tools and the seats they hold", + Metadata: map[string]string{"node": node}, + }, func() []announce.Endpoint { return endpointsOf(conn, own, registrations, modules) }) + if err != nil { + logf("[mesh-tools] cannot announce what it serves: %v", err) + } else { + stops = append(stops, announced) + } conn.Flush() return stopAll, nil } +// instanceOf is this runtime's instance on the bus: its machine, which is what tells two instances of +// one service apart; the connection's module where it has no machine. +func instanceOf(conn *bus.Conn) string { + if n := conn.Node(); n != "" { + return n + } + return conn.Module() +} + +// endpointsOf is everything this runtime serves now: each served module's tools on every subject the +// mesh issued for them, and each held seat's verbs on the seat's subject (ADR 0197). +func endpointsOf(conn *bus.Conn, own []registration, registrations []registration, modules []string) []announce.Endpoint { + node := conn.Node() + var out []announce.Endpoint + for _, r := range own { + for _, t := range r.tools { + served := conn.ServedOn(r.module, t.Name) + interchangeable := false + for _, s := range served { + interchangeable = interchangeable || s.Subject == "mesh.mod."+r.module+".tool."+t.Name + } + for _, s := range served { + out = append(out, announce.Endpoint{Kind: announce.KindTool, Module: r.module, Tool: t.Name, + Node: node, Description: t.Description, Schema: t.Input, Interchangeable: interchangeable, + Subject: s.Subject, Queue: s.Queue}) + } + } + } + impl := map[string]map[string]launch.Tool{} + for _, r := range registrations { + if impl[r.module] == nil { + impl[r.module] = map[string]launch.Tool{} + } + for _, t := range r.tools { + impl[r.module][t.Name] = t + } + } + have := map[string]bool{} + for _, module := range modules { + m := conn.Membership(module) + if m == nil { + continue + } + for _, v := range m.Seats { + t, ok := impl[v.Seat][v.Verb] + if !ok || have[v.Subject] { + continue + } + have[v.Subject] = true + scope := "mesh" + if node != "" && strings.HasSuffix(v.Subject, "."+node) { + scope = "node" + } + out = append(out, announce.Endpoint{Kind: announce.KindSeat, Module: module, Tool: v.Verb, Seat: v.Seat, + Scope: scope, Node: node, Description: t.Description, Schema: t.Input, Subject: v.Subject}) + } + } + return out +} + func orNone(names []string) string { if len(names) == 0 { return "(none)" -- 2.54.0 From 66e8be0e31e6aa4e51338a29f7e23dc4d76aac01 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 3 Oct 2026 22:15:48 +0200 Subject: [PATCH 3/3] The TypeScript runtime announces what it serves too, in the same services format (hq ADR 0197) The per-module containers still run this runtime; their tools and the seats they hold (the store's, the catalogue's) must be found by the console the same way as the node runtime's. It answers $SRV.PING, $SRV.INFO and $SRV.STATS with one service per process, one endpoint per tool per subject and per seat verb served, the metadata as the Go runtime writes it. --- node-tools/src/announce.ts | 134 +++++++++++++++++++++++++++++++ node-tools/src/broker-nats.ts | 20 +++++ node-tools/src/runtime.ts | 28 ++++++- node-tools/test/announce.test.ts | 69 ++++++++++++++++ 4 files changed, 247 insertions(+), 4 deletions(-) create mode 100644 node-tools/src/announce.ts create mode 100644 node-tools/test/announce.test.ts diff --git a/node-tools/src/announce.ts b/node-tools/src/announce.ts new file mode 100644 index 0000000..a152a51 --- /dev/null +++ b/node-tools/src/announce.ts @@ -0,0 +1,134 @@ +// What a runtime serves, it announces (novox/hq ADR 0197): the NATS services protocol's discovery — +// `$SRV.PING`, `$SRV.INFO`, `$SRV.STATS`, and the same followed by the service's name and its id — +// answered in the io.nats.micro.v1 format with what is served at the moment of the request. Serving is +// unchanged; this only says what is served. One service per runtime process, because the bus admits +// one reply per request from each responder: one endpoint per tool per subject, its metadata saying +// which module, seat, scope and machine it is. The same shape the Go runtime answers. + +import { StringCodec } from "nats"; +import { asSchema } from "@novox/mesh-sdk/stdio"; +import type { ToolDefinition } from "@novox/mesh-sdk/tools"; +import type { RuntimeBroker } from "./broker-nats.js"; + +const sc = StringCodec(); + +export const VERSION = "0.1.0"; +export const INFO_RESPONSE = "io.nats.micro.v1.info_response"; +export const PING_RESPONSE = "io.nats.micro.v1.ping_response"; +export const STATS_RESPONSE = "io.nats.micro.v1.stats_response"; + +/** One tool served on one subject, as announced. */ +export interface Endpoint { + kind: "tool" | "seat"; + module: string; + tool: string; + seat?: string; + scope?: "mesh" | "node"; + node: string; + description: string; + schema: unknown; + interchangeable: boolean; + subject: string; + queue?: string; +} + +/** A seat's verb as the runtime serves it, with the definition that answers it. */ +export interface ServedSeatVerb { + seat: string; + verb: string; + subject: string; + holder: string; + tool: ToolDefinition; +} + +export interface Service { + name: string; + id: string; + description: string; + metadata: Record; +} + +/** The endpoint's name as the protocol allows it; the metadata, not the name, identifies it. */ +function nameOf(e: Endpoint): string { + const clean = (s: string) => s.replace(/[^A-Za-z0-9_-]/g, "_"); + return `${clean(e.kind === "seat" ? e.seat ?? "" : e.module)}__${clean(e.tool)}`; +} + +/** The info_response for these endpoints. */ +export function info(s: Service, endpoints: Endpoint[]): Record { + return { + name: s.name, id: s.id, version: VERSION, metadata: s.metadata, type: INFO_RESPONSE, description: s.description, + endpoints: endpoints.map((e) => { + const metadata: Record = { + kind: e.kind, module: e.module, tool: e.tool, node: e.node, description: e.description, + schema: JSON.stringify(e.schema ?? {}), interchangeable: e.interchangeable ? "true" : "false", + }; + if (e.kind === "seat") { + metadata.seat = e.seat ?? ""; + metadata.scope = e.scope ?? "mesh"; + } + return { name: nameOf(e), subject: e.subject, queue_group: e.queue ?? "", metadata }; + }), + }; +} + +/** Everything served now: each module's tools on every subject issued for them, and each held + * seat's verbs on the seat's subject. */ +export function endpointsOf( + broker: RuntimeBroker, + own: { module: string; tools: ToolDefinition[] }[], + seats: ServedSeatVerb[], +): Endpoint[] { + const node = broker.node ?? ""; + const out: Endpoint[] = []; + for (const { module, tools } of own) { + for (const t of tools) { + const served = broker.servedOn ? broker.servedOn(module, t.name) : []; + const interchangeable = served.some((s) => s.subject === `mesh.mod.${module}.tool.${t.name}`); + for (const s of served) { + out.push({ kind: "tool", module, tool: t.name, node, description: t.description, schema: asSchema(t.input), + interchangeable, subject: s.subject, queue: s.queue }); + } + } + } + for (const v of seats) { + out.push({ kind: "seat", module: v.holder, tool: v.verb, seat: v.seat, + scope: node && v.subject.endsWith(`.${node}`) ? "node" : "mesh", node, description: v.tool.description, + schema: asSchema(v.tool.input), interchangeable: false, subject: v.subject }); + } + return out; +} + +/** Answer discovery for one service until stopped. */ +export function announce(broker: RuntimeBroker, s: Service, current: () => Endpoint[]): () => void { + const started = new Date().toISOString(); + const identity = { name: s.name, id: s.id, version: VERSION, metadata: s.metadata }; + const answer = (subject: string): Uint8Array | undefined => { + const parts = subject.split("."); + if (parts[0] !== "$SRV" || parts.length < 2) return undefined; + if (parts.length >= 3 && parts[2] !== s.name) return undefined; // another service's + if (parts.length >= 4 && parts[3] !== s.id) return undefined; // another instance's + let v: unknown; + switch (parts[1]) { + case "PING": + v = { ...identity, type: PING_RESPONSE }; + break; + case "INFO": + v = info(s, current()); + break; + case "STATS": + v = { ...identity, type: STATS_RESPONSE, started, endpoints: current().map((e) => ({ + name: nameOf(e), subject: e.subject, queue_group: e.queue ?? "", num_requests: 0, num_errors: 0, + last_error: "", processing_time: 0, average_processing_time: 0 })) }; + break; + default: + return undefined; + } + return sc.encode(JSON.stringify(v)); + }; + const stops: (() => void)[] = []; + for (const verb of ["PING", "INFO", "STATS"]) { + for (const subject of [`$SRV.${verb}`, `$SRV.${verb}.>`]) stops.push(broker.raw!(subject, (subj) => answer(subj))); + } + return () => stops.forEach((stop) => stop()); +} diff --git a/node-tools/src/broker-nats.ts b/node-tools/src/broker-nats.ts index cbf53c9..e7546eb 100644 --- a/node-tools/src/broker-nats.ts +++ b/node-tools/src/broker-nats.ts @@ -97,6 +97,13 @@ export interface RuntimeBroker extends Broker { serving(): string[]; /** The module this connection is: what its credential named, and what a bare key serves as. */ readonly module: string; + /** Where a served module's tool is answered right now (ADR 0197: what it serves, it announces). */ + servedOn?(module: string, tool: string): { subject: string; queue?: string }[]; + /** Answer a subject in a format of its own, not a tool's reply envelope — the NATS services + * protocol's discovery (novox/hq ADR 0197). An undefined answer is no reply. */ + raw?(subject: string, answer: (subject: string, data: Uint8Array) => Uint8Array | undefined): () => void; + /** The machine this connection serves on, when its credential names one. */ + readonly node?: string; } /** @@ -348,6 +355,19 @@ export async function connectNats( }, module: self, + node, + servedOn, + raw(subject: string, answer: (subject: string, data: Uint8Array) => Uint8Array | undefined): () => void { + const sub = conn.subscribe(subject); + subs.push(sub); + void (async () => { + for await (const msg of sub) { + const body = answer(msg.subject, msg.data); + if (body) msg.respond(body); + } + })(); + return () => sub.unsubscribe(); + }, follow, serving: () => [...issued.keys()], membership: (module?: string) => issued.get(module ?? self), diff --git a/node-tools/src/runtime.ts b/node-tools/src/runtime.ts index e18b1c9..578d32b 100644 --- a/node-tools/src/runtime.ts +++ b/node-tools/src/runtime.ts @@ -16,6 +16,7 @@ import { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tool import type { Broker } from "@novox/mesh-sdk/messaging"; import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js"; import { launch, launches } from "./launch.js"; +import { announce, endpointsOf, type ServedSeatVerb } from "./announce.js"; /** * The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the @@ -259,7 +260,16 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> { console.log(`[mesh-tools] serving ${names.length} tool(s) for ${served.size} module(s): ${names.join(", ") || "(none)"}` + (failed.size ? `; not serving ${[...failed.keys()].join(", ")}, whose bundle(s) failed to load` : "")); - stops.push(await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations)); + const seats = await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations); + stops.push(seats.stop); + // **What it serves, it announces** (novox/hq ADR 0197), asked at the moment of the request. + if (typeof runtime.raw === "function") { + stops.push(announce(runtime, { + name: self ?? "runtime", id: runtime.node ?? self ?? "runtime", + description: `the tool runtime of ${self ?? "a module"}${runtime.node ? ` on ${runtime.node}` : ""}`, + metadata: runtime.node ? { node: runtime.node } : {}, + }, () => endpointsOf(runtime, ownRegistrations, seats.serving()))); + } return () => stop(); } @@ -303,11 +313,17 @@ async function serveClaimedSeats( self: string | undefined, credential: Credential | undefined, registrations: { module: string; owner: string; tools: ToolDefinition[] }[], -): Promise<() => void> { - if (typeof broker.handleSubject !== "function") return () => {}; +): Promise<{ stop: () => void; serving: () => ServedSeatVerb[] }> { + if (typeof broker.handleSubject !== "function") return { stop: () => {}, serving: () => [] }; // A seat's verbs are the role's, not the software's (ADR 0159): implemented under the seat's // name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools. const implementations = new Map) => Promise>>(); + const definitions = new Map>(); + for (const { module, tools } of registrations) { + const defs = definitions.get(module) ?? new Map(); + for (const t of tools) defs.set(t.name, t); + definitions.set(module, defs); + } for (const { module, owner, tools } of registrations) { const verbs = implementations.get(module) ?? new Map) => Promise>(); for (const t of tools) verbs.set(t.name, (args) => atWork.run({ module: owner }, () => t.run(args))); @@ -341,21 +357,25 @@ async function serveClaimedSeats( }; let stops: (() => void)[] = []; + let servingNow: ServedSeatVerb[] = []; const serve = async (): Promise => { stops.forEach((s) => s()); stops = []; + servingNow = []; for (const v of wanted()) { const run = implementations.get(v.seat)?.get(v.verb); + const tool = definitions.get(v.seat)?.get(v.verb); if (!run) { console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`); continue; } stops.push(await broker.handleSubject(v.subject, run)); + if (tool) servingNow.push({ ...v, tool }); console.log(`[mesh-tools] serving ${v.seat}'s ${v.verb} on ${v.subject}, admitted where ${v.holder} holds the seat`); } }; await serve(); // A membership issued to any served module may add, move or withdraw a seat's verbs. if (typeof broker.onMembership === "function") broker.onMembership(() => void serve()); - return () => stops.forEach((s) => s()); + return { stop: () => stops.forEach((s) => s()), serving: () => servingNow }; } diff --git a/node-tools/test/announce.test.ts b/node-tools/test/announce.test.ts new file mode 100644 index 0000000..7df625c --- /dev/null +++ b/node-tools/test/announce.test.ts @@ -0,0 +1,69 @@ +/** + * What a runtime serves, it announces (novox/hq ADR 0197): the TypeScript runtime the per-module + * containers still run answers the NATS services protocol's discovery in the same shape as the Go + * tool runtime — its module's tools on every subject issued, and the seat verbs it serves. + * + * MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/announce.test.ts + */ +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { fileURLToPath } from "node:url"; +import { connect, StringCodec } from "nats"; + +import { resetTools } from "@novox/mesh-sdk/tools"; + +import { connectNats, membershipSubject } from "../dist/broker-nats.js"; +import { runTools } from "../dist/runtime.js"; + +const url = process.env.MESH_TEST_NATS; +const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url)); +const sc = StringCodec(); + +test("the runtime answers $SRV.INFO with what it serves, in the services protocol's format", async (t) => { + if (!url) return t.skip("MESH_TEST_NATS unset"); + resetTools(); + const nc = await connect({ servers: url }); + const jsm = await nc.jetstreamManager(); + try { + await jsm.streams.delete("ASSIGNMENTS"); + } catch { + // none yet + } + await jsm.streams.add({ name: "ASSIGNMENTS", subjects: ["mesh.assignment.>"], max_msgs_per_subject: 1, allow_direct: true } as never); + await nc.jetstream().publish(membershipSubject("anchor", "shop"), sc.encode(JSON.stringify({ + node: "anchor", module: "shop", + serves: [{ subject: "mesh.mod.shop.tool.{tool}.anchor" }, { subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }], + emits: "mesh.mod.shop.event.{event}", tools: "mesh.mod.shop.tool.tools", + }))); + const shop = await connectNats({ url, node: "anchor", module: "shop" }); + let stop = () => {}; + try { + stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] }); + const msg = await nc.request("$SRV.INFO.shop.anchor", sc.encode(""), { timeout: 2000 }); + const info = JSON.parse(sc.decode(msg.data)) as { + type: string; name: string; id: string; version: string; + endpoints: { name: string; subject: string; queue_group: string; metadata: Record }[]; + }; + assert.equal(info.type, "io.nats.micro.v1.info_response"); + assert.equal(info.name, "shop"); + assert.equal(info.id, "anchor"); + assert.ok(info.version); + const price = info.endpoints.filter((e) => e.metadata.tool === "price").map((e) => `${e.subject}|${e.queue_group}`).sort(); + assert.deepEqual(price, ["mesh.mod.shop.tool.price.anchor|", "mesh.mod.shop.tool.price|serve.shop"]); + const one = info.endpoints.find((e) => e.metadata.tool === "price")!; + assert.equal(one.metadata.kind, "tool"); + assert.equal(one.metadata.module, "shop"); + assert.equal(one.metadata.node, "anchor"); + assert.equal(one.metadata.interchangeable, "true"); + assert.ok(JSON.parse(one.metadata.schema).type === "object"); + // Ping answers with the same identity; another service's request is not answered. + const ping = JSON.parse(sc.decode((await nc.request("$SRV.PING", sc.encode(""), { timeout: 2000 })).data)); + assert.equal(ping.type, "io.nats.micro.v1.ping_response"); + await assert.rejects(nc.request("$SRV.INFO.somebody-else", sc.encode(""), { timeout: 300 })); + } finally { + stop(); + await shop.close(); + await nc.close(); + resetTools(); + } +}); -- 2.54.0