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 new file mode 100644 index 0000000..f8b2cb4 --- /dev/null +++ b/node-tools/internal/console/address.go @@ -0,0 +1,728 @@ +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" + "sort" + "strings" + "sync" + "time" + + "github.com/novox/mesh-tools/node-tools/internal/announce" + "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 { + t := strings.TrimSpace(output) + if strings.HasPrefix(t, "{") || strings.HasPrefix(t, "[") { + return t + } + for _, open := range []string{"\n{", "\n["} { + if i := strings.Index(output, open); i >= 0 { + return strings.TrimSpace(output[i+1:]) + } + } + return "" +} + +// 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"` +} + +// recordedMachine is a machine as the controller's records hold it (`node list --json`). +type recordedMachine struct { + Name string `json:"name"` +} + +// 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 nodesOut, modulesOut string + wg.Add(2) + go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }() + go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }() + infos, err := announce.Gather(conn) + wg.Wait() + if err != nil { + return nil, err + } + + 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 announced[e.Module] == nil { + announced[e.Module] = map[string]bool{} + } + 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) + } + } + } + } + // 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 + } + 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 { + 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) + } + } + } + } 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 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 { + 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("the mesh's discovery failed: " + err.Error()) + } + 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..f2a3208 --- /dev/null +++ b/node-tools/internal/console/address_test.go @@ -0,0 +1,228 @@ +package console + +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" +) + +// 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() + + // 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: 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) { + return answer(), nil + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stop) + } + serve("nodes", func() any { + 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"}]`) + }) + var gammaOn atomic.Value + gammaOn.Store("[]") + serve("modules", func() any { + 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))) + }) + 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 { + 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) + } + // 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") { + 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.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) + 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) + } + + 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 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) + } + 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 0645ae9..d41b55f 100644 --- a/node-tools/internal/console/console_test.go +++ b/node-tools/internal/console/console_test.go @@ -55,23 +55,18 @@ 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() - 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) } @@ -90,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 07de2b1..330e2f1 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,9 +117,12 @@ 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)) + 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 { @@ -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." +} 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)" 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(); + } +});