The mesh's tools are found by address, from what announces itself on the bus (hq ADR 0195, 0197) #39
@@ -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
|
||||||
|
}
|
||||||
@@ -512,6 +512,59 @@ func (c *Conn) PublishAs(module string, env Envelope) error {
|
|||||||
return err
|
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
|
// Flush waits until the bus has every subscription made so far, so what is served is answerable
|
||||||
// when this returns.
|
// when this returns.
|
||||||
func (c *Conn) Flush() { _ = c.nc.Flush() }
|
func (c *Conn) Flush() { _ = c.nc.Flush() }
|
||||||
|
|||||||
@@ -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:
|
||||||
|
//
|
||||||
|
// <seat>.<verb> a seat held once for the mesh — its holder answers
|
||||||
|
// <node>/<seat>.<verb> a seat held once per machine — that machine's holder answers
|
||||||
|
// <node>/<module>.<tool> a module assigned to a machine — that assignment answers
|
||||||
|
// <module>.<tool> 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: `<seat>.<verb>` for a seat held once for the mesh (e.g. `mesh-controller.nodes`); " +
|
||||||
|
"`<node>/<seat>.<verb>` for a seat every machine holds (e.g. `ace/node-packet-filter.rules`); " +
|
||||||
|
"`<node>/<module>.<tool>` for a module on one machine (e.g. `novox/postgres.postgres_list_databases`); " +
|
||||||
|
"and `<module>.<tool>` 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{} // <module>.<tool> or <seat>.<verb> → 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:<s>.<v>[@node] or <m>.<t>[@node]
|
||||||
|
Name string // <seat>.<verb> or <module>.<tool>, 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 <node>/%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 <node>/%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) // `<node>/…` 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 = "<node>/" + 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 <module>.<tool> 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 := "<node>"
|
||||||
|
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{"<node>/" + 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)
|
||||||
|
}
|
||||||
@@ -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>/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 <node>/beta.three — it runs on desk",
|
||||||
|
"node-shelf.list": "held once per machine: write <node>/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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -5,19 +5,11 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"regexp"
|
"regexp"
|
||||||
"sort"
|
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
|
||||||
|
|
||||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
"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.
|
// Tool is a tool as its module — or, for a role's tool, the mesh's records — describes it.
|
||||||
type Tool struct {
|
type Tool struct {
|
||||||
Module string
|
Module string
|
||||||
@@ -67,107 +59,14 @@ func toolKey(name string, seats Seats) string {
|
|||||||
return name
|
return name
|
||||||
}
|
}
|
||||||
|
|
||||||
// toolsOn asks the mesh what tools it has: the catalogue which modules it holds, each module what it
|
// toolsOn is the flat catalogue (MESH_CONSOLE_FLAT=1): what announced itself on the bus, every
|
||||||
// serves, the controller's seat every role's tools — at once, so a restarting control plane hides
|
// tool and seat verb, and the assignments with tools that did not (novox/hq ADR 0197).
|
||||||
// nothing else.
|
|
||||||
func toolsOn(conn *bus.Conn) (*Listing, error) {
|
func toolsOn(conn *bus.Conn) (*Listing, error) {
|
||||||
type rolesAnswer struct {
|
x, err := indexOn(conn)
|
||||||
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{}, "")
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
wg.Wait()
|
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
var held struct {
|
return x.Listing, nil
|
||||||
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
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// callTool calls `<module>.<tool>[@<node>]`, on the subject the listing names for it when it names one.
|
// callTool calls `<module>.<tool>[@<node>]`, on the subject the listing names for it when it names one.
|
||||||
|
|||||||
@@ -55,23 +55,18 @@ func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) {
|
|||||||
}
|
}
|
||||||
defer stop()
|
defer stop()
|
||||||
|
|
||||||
catalogue := connect(t, "mesh-catalog", "")
|
// The controller's records (ADR 0197): ghost declares tools on desk and nothing announces it.
|
||||||
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()
|
|
||||||
controller := connect(t, "mesh-controller", "")
|
controller := connect(t, "mesh-controller", "")
|
||||||
stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.tools", func(json.RawMessage) (any, error) {
|
stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.modules", func(json.RawMessage) (any, error) {
|
||||||
return map[string]any{"seats": []map[string]any{{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{
|
return map[string]any{"ok": true, "output": `[{"module":"alpha","on":["desk"],"tools":true},` +
|
||||||
{"name": "list", "description": "what is on the shelf", "input": map[string]any{}},
|
`{"module":"beta","on":["desk"],"tools":true},{"module":"ghost","on":["desk"],"tools":true}]`}, nil
|
||||||
{"name": "clear", "description": "take it all off", "input": map[string]any{}},
|
|
||||||
}}}}, nil
|
|
||||||
})
|
})
|
||||||
defer stopSeat()
|
defer stopSeat()
|
||||||
catalogue.Flush()
|
|
||||||
controller.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 {
|
if err != nil {
|
||||||
t.Fatal(err)
|
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" {
|
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)
|
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)
|
t.Errorf("not answering: %v", got)
|
||||||
}
|
}
|
||||||
for _, x := range listed["tools"].([]any) {
|
for _, x := range listed["tools"].([]any) {
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ package console
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"os"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
@@ -47,10 +48,17 @@ type Surface struct {
|
|||||||
mu sync.Mutex
|
mu sync.Mutex
|
||||||
known *Listing
|
known *Listing
|
||||||
at time.Time
|
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`.
|
// 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) {
|
func (s *Surface) listing() (*Listing, error) {
|
||||||
s.mu.Lock()
|
s.mu.Lock()
|
||||||
@@ -99,12 +107,7 @@ func (s *Surface) Handle(r Request) *Reply {
|
|||||||
"protocolVersion": Protocol,
|
"protocolVersion": Protocol,
|
||||||
"capabilities": map[string]any{"tools": map[string]any{}},
|
"capabilities": map[string]any{"tools": map[string]any{}},
|
||||||
"serverInfo": map[string]any{"name": "mesh", "version": "1"},
|
"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 " +
|
"instructions": s.instructions(),
|
||||||
"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 <module>.<tool>.",
|
|
||||||
})
|
})
|
||||||
case "notifications/initialized":
|
case "notifications/initialized":
|
||||||
return nil
|
return nil
|
||||||
@@ -114,9 +117,12 @@ func (s *Surface) Handle(r Request) *Reply {
|
|||||||
}
|
}
|
||||||
return answer(r.ID, map[string]any{})
|
return answer(r.ID, map[string]any{})
|
||||||
case "tools/list":
|
case "tools/list":
|
||||||
|
if !s.Flat {
|
||||||
|
return answer(r.ID, map[string]any{"tools": discovery()})
|
||||||
|
}
|
||||||
l, err := s.listing()
|
l, err := s.listing()
|
||||||
if err != nil {
|
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))
|
tools := make([]map[string]any, 0, len(l.Tools))
|
||||||
for _, t := range l.Tools {
|
for _, t := range l.Tools {
|
||||||
@@ -142,6 +148,12 @@ func (s *Surface) Handle(r Request) *Reply {
|
|||||||
Arguments map[string]any `json:"arguments"`
|
Arguments map[string]any `json:"arguments"`
|
||||||
}
|
}
|
||||||
_ = json.Unmarshal(r.Params, &p)
|
_ = 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{}
|
args := map[string]any{}
|
||||||
for k, v := range p.Arguments {
|
for k, v := range p.Arguments {
|
||||||
args[k] = v
|
args[k] = v
|
||||||
@@ -266,3 +278,20 @@ func withNode(schema map[string]any, description string, required bool) map[stri
|
|||||||
}
|
}
|
||||||
return out
|
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 <module>.<tool>."
|
||||||
|
}
|
||||||
|
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."
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,5 +1,9 @@
|
|||||||
// Package meshtest raises what the controller would, for tests against a real bus: the ASSIGNMENTS
|
// 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.
|
// 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
|
package meshtest
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
|||||||
@@ -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{}
|
||||||
|
}
|
||||||
@@ -14,6 +14,7 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"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/bus"
|
||||||
"github.com/novox/mesh-tools/node-tools/internal/launch"
|
"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)
|
logf("%s", line)
|
||||||
stops = append(stops, serveSeats(conn, modules, registrations, logf))
|
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()
|
conn.Flush()
|
||||||
return stopAll, nil
|
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 {
|
func orNone(names []string) string {
|
||||||
if len(names) == 0 {
|
if len(names) == 0 {
|
||||||
return "(none)"
|
return "(none)"
|
||||||
|
|||||||
@@ -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<string, string>;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 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<string, unknown> {
|
||||||
|
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<string, string> = {
|
||||||
|
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());
|
||||||
|
}
|
||||||
@@ -97,6 +97,13 @@ export interface RuntimeBroker extends Broker {
|
|||||||
serving(): string[];
|
serving(): string[];
|
||||||
/** The module this connection is: what its credential named, and what a bare key serves as. */
|
/** The module this connection is: what its credential named, and what a bare key serves as. */
|
||||||
readonly module: string;
|
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,
|
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,
|
follow,
|
||||||
serving: () => [...issued.keys()],
|
serving: () => [...issued.keys()],
|
||||||
membership: (module?: string) => issued.get(module ?? self),
|
membership: (module?: string) => issued.get(module ?? self),
|
||||||
|
|||||||
@@ -16,6 +16,7 @@ import { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tool
|
|||||||
import type { Broker } from "@novox/mesh-sdk/messaging";
|
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||||
import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
|
import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
|
||||||
import { launch, launches } from "./launch.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
|
* 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)"}` +
|
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` : ""));
|
(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();
|
return () => stop();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -303,11 +313,17 @@ async function serveClaimedSeats(
|
|||||||
self: string | undefined,
|
self: string | undefined,
|
||||||
credential: Credential | undefined,
|
credential: Credential | undefined,
|
||||||
registrations: { module: string; owner: string; tools: ToolDefinition[] }[],
|
registrations: { module: string; owner: string; tools: ToolDefinition[] }[],
|
||||||
): Promise<() => void> {
|
): Promise<{ stop: () => void; serving: () => ServedSeatVerb[] }> {
|
||||||
if (typeof broker.handleSubject !== "function") return () => {};
|
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
|
// 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.
|
// name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools.
|
||||||
const implementations = new Map<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>();
|
const implementations = new Map<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>();
|
||||||
|
const definitions = new Map<string, Map<string, ToolDefinition>>();
|
||||||
|
for (const { module, tools } of registrations) {
|
||||||
|
const defs = definitions.get(module) ?? new Map<string, ToolDefinition>();
|
||||||
|
for (const t of tools) defs.set(t.name, t);
|
||||||
|
definitions.set(module, defs);
|
||||||
|
}
|
||||||
for (const { module, owner, tools } of registrations) {
|
for (const { module, owner, tools } of registrations) {
|
||||||
const verbs = implementations.get(module) ?? new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
|
const verbs = implementations.get(module) ?? new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
|
||||||
for (const t of tools) verbs.set(t.name, (args) => atWork.run({ module: owner }, () => t.run(args)));
|
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 stops: (() => void)[] = [];
|
||||||
|
let servingNow: ServedSeatVerb[] = [];
|
||||||
const serve = async (): Promise<void> => {
|
const serve = async (): Promise<void> => {
|
||||||
stops.forEach((s) => s());
|
stops.forEach((s) => s());
|
||||||
stops = [];
|
stops = [];
|
||||||
|
servingNow = [];
|
||||||
for (const v of wanted()) {
|
for (const v of wanted()) {
|
||||||
const run = implementations.get(v.seat)?.get(v.verb);
|
const run = implementations.get(v.seat)?.get(v.verb);
|
||||||
|
const tool = definitions.get(v.seat)?.get(v.verb);
|
||||||
if (!run) {
|
if (!run) {
|
||||||
console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`);
|
console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
stops.push(await broker.handleSubject(v.subject, run));
|
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`);
|
console.log(`[mesh-tools] serving ${v.seat}'s ${v.verb} on ${v.subject}, admitted where ${v.holder} holds the seat`);
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
await serve();
|
await serve();
|
||||||
// A membership issued to any served module may add, move or withdraw a seat's verbs.
|
// 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());
|
if (typeof broker.onMembership === "function") broker.onMembership(() => void serve());
|
||||||
return () => stops.forEach((s) => s());
|
return { stop: () => stops.forEach((s) => s()), serving: () => servingNow };
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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<string, string> }[];
|
||||||
|
};
|
||||||
|
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();
|
||||||
|
}
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user