diff --git a/node-tools/cmd/node-tools/main.go b/node-tools/cmd/node-tools/main.go new file mode 100644 index 0000000..a52242d --- /dev/null +++ b/node-tools/cmd/node-tools/main.go @@ -0,0 +1,130 @@ +// node-tools — the node's tool runtime, in Go (novox/hq ADR 0175, ADR 0193; to-be 38 WP4d). +// +// One process per machine: it connects to the bus on the node's credential, launches every assigned +// module's tools bundle and serves its tools and its seats' verbs; as the node-tools module — or +// wherever MESH_CONSOLE_LISTEN says — it is also the console, MCP over HTTP on loopback. +// +// MESH_BROKER_FILE the credential the mesh sealed to this machine for the runtime +// MESH_TOOL_MODULES =,… the bundles to launch +// MESH_TOOL_ENV {"": {"": ""}} what each is given (ADR 0192) +// MESH_OPERATOR_ACCOUNT whose machine this is, and MESH_OPERATOR_HOME their home +// MESH_CONSOLE_LISTEN where the console listens, overriding 127.0.0.1:4270 +package main + +import ( + "encoding/json" + "fmt" + "log" + "math/rand/v2" + "os" + "os/signal" + "syscall" + "time" + + "github.com/novox/mesh-tools/node-tools/internal/bus" + "github.com/novox/mesh-tools/node-tools/internal/console" + "github.com/novox/mesh-tools/node-tools/internal/runtime" +) + +// runtimeModule is the module that is the node's tool runtime; on its credential, serving is also +// the console. +const runtimeModule = "node-tools" + +// consoleListen is where the console listens when nothing says otherwise. +const consoleListen = "127.0.0.1:4270" + +func main() { + log.SetFlags(0) + if len(os.Args) > 1 && os.Args[1] != "serve" { + fmt.Fprintf(os.Stderr, "node-tools: %q is not a mode; this runtime serves (novox/hq ADR 0193)\n", os.Args[1]) + os.Exit(2) + } + cred := credential() + conn := connectPatiently(cred) + + served, err := runtime.ServedModulesFrom(os.Getenv(runtime.ToolModules), cred.Module) + if err != nil { + log.Fatalf("node-tools: %v", err) + } + envs, err := runtime.TakeToolEnvs() + if err != nil { + log.Fatalf("node-tools: %v", err) + } + stop, err := runtime.Run(conn, served, envs, log.Printf) + if err != nil { + log.Fatalf("node-tools: %v", err) + } + + listen := os.Getenv("MESH_CONSOLE_LISTEN") + if listen == "" && cred.Module == runtimeModule { + listen = consoleListen + } + var up *console.Listening + if listen != "" { + node := cred.Node + if node == "" { + node = "?" + } + who := node + "." + cred.Module + up, err = console.Serve(console.NewSurface(conn, who), listen) + if err != nil { + log.Fatalf("node-tools: %v", err) + } + log.Printf("mesh console listening on http://%s/mcp as %s", up.Address, who) + } + + signals := make(chan os.Signal, 1) + signal.Notify(signals, syscall.SIGTERM, syscall.SIGINT) + <-signals + stop() + if up != nil { + _ = up.Close() + } + conn.Close() +} + +// credential reads the credential the mesh delivered; a missing or unreadable one is a fault of +// configuration, said and final. +func credential() bus.Credential { + file := os.Getenv("MESH_BROKER_FILE") + if file == "" { + if url := os.Getenv("MESH_BROKER_URL"); url != "" { + return bus.Credential{URL: url, Module: runtimeModule} + } + fmt.Fprintln(os.Stderr, "node-tools: set MESH_BROKER_FILE (a sealed credential) or MESH_BROKER_URL — there is no broker to reach") + os.Exit(1) + } + raw, err := os.ReadFile(file) + if err != nil { + fmt.Fprintf(os.Stderr, "node-tools: cannot read the broker credential at %s: %v\n", file, err) + os.Exit(1) + } + var cred bus.Credential + if err := json.Unmarshal(raw, &cred); err != nil { + fmt.Fprintf(os.Stderr, "node-tools: cannot read the broker credential at %s: %v\n", file, err) + os.Exit(1) + } + if cred.URL == "" { + fmt.Fprintf(os.Stderr, "node-tools: %s carries no url — it is not a broker credential\n", file) + os.Exit(1) + } + return cred +} + +// connectPatiently retries while the bus is merely not reachable yet — the normal case at startup — +// and gives up at once on what waiting cannot fix (issue 058). +func connectPatiently(cred bus.Credential) *bus.Conn { + for delay := 2 * time.Second; ; delay = min(delay*2, 30*time.Second) { + conn, err := bus.Connect(cred) + if err == nil { + return conn + } + if fatal := bus.Fatal(err); fatal != "" { + fmt.Fprintf(os.Stderr, "node-tools: %s — waiting will not fix this; giving up\n", fatal) + os.Exit(1) + } + wait := delay + time.Duration(rand.IntN(1000))*time.Millisecond + fmt.Fprintf(os.Stderr, "node-tools: the broker is not reachable yet (%v); retrying in %ds\n", err, int(wait.Round(time.Second)/time.Second)) + time.Sleep(wait) + } +} diff --git a/node-tools/go.mod b/node-tools/go.mod new file mode 100644 index 0000000..06b9b50 --- /dev/null +++ b/node-tools/go.mod @@ -0,0 +1,15 @@ +module github.com/novox/mesh-tools/node-tools + +go 1.26.0 + +require ( + github.com/nats-io/nats.go v1.54.0 + golang.org/x/sys v0.48.0 +) + +require ( + github.com/klauspost/compress v1.20.0 // indirect + github.com/nats-io/nkeys v0.4.16 // indirect + github.com/nats-io/nuid v1.0.1 // indirect + golang.org/x/crypto v0.57.0 // indirect +) diff --git a/node-tools/go.sum b/node-tools/go.sum new file mode 100644 index 0000000..65100b8 --- /dev/null +++ b/node-tools/go.sum @@ -0,0 +1,12 @@ +github.com/klauspost/compress v1.20.0 h1:a3C1ke2ohxFymNlb2HWAHjDeKCI90scRskErZkR0ezA= +github.com/klauspost/compress v1.20.0/go.mod h1:LUdAzn7YLVvxLpc7y3V1m40wESHTgc1422pwwBSKYuI= +github.com/nats-io/nats.go v1.54.0 h1:vsXoOxjHp/GmPUN+EcI7uOf/uB+iAP+kEsAFNQN0yzA= +github.com/nats-io/nats.go v1.54.0/go.mod h1:y+DZoD1oBOYfZTU681eTUiUjI0vbqYGixNVFHcjHJ0k= +github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg= +github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs= +github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= +github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= +golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M= +golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA= +golang.org/x/sys v0.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo= +golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og= diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go new file mode 100644 index 0000000..8fddf95 --- /dev/null +++ b/node-tools/internal/bus/bus.go @@ -0,0 +1,568 @@ +// Package bus is the node runtime's connection to the mesh bus, on NATS (novox/hq design 25, design +// 29, ADR 0160, ADR 0175). It is the Go port of node-tools' broker-nats.ts, and speaks the same +// wire: the same subjects, the same JSON request and reply bodies, the same event headers. +// +// mesh.mod..event. an event a module emits +// mesh.mod..tool. a tool a module serves +// mesh.seat..tool. a role's verb, answered by whoever holds the seat +package bus + +import ( + "crypto/sha256" + "crypto/tls" + "crypto/x509" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "log" + "regexp" + "strings" + "sync" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-tools/node-tools/internal/wire" +) + +// RequestTimeout is how long a call waits for its answer — what modules already expect. +const RequestTimeout = 30 * time.Second + +const assignmentsStream = "ASSIGNMENTS" + +// Credential is a broker credential as the mesh delivers it (novox/hq ADR 0120). +type Credential struct { + URL string `json:"url"` + Fingerprint string `json:"fingerprint,omitempty"` + Node string `json:"node,omitempty"` + Module string `json:"module,omitempty"` + User string `json:"user,omitempty"` + Password string `json:"password,omitempty"` + Claims []Claim `json:"claims,omitempty"` +} + +// Claim is a seat a module claims, with the verbs it promises (novox/hq ADR 0159). +type Claim struct { + Seat string `json:"seat"` + Scope string `json:"scope,omitempty"` + Serves []string `json:"serves,omitempty"` +} + +// Membership is what the mesh issued one assignment (novox/hq ADR 0160). +type Membership struct { + Node string `json:"node"` + Module string `json:"module"` + Serves []Served `json:"serves"` + Seats []SeatVerb `json:"seats,omitempty"` + Emits string `json:"emits"` + Reaches map[string][]string `json:"reaches,omitempty"` + Tools string `json:"tools"` +} + +// Served is an address a tool is answered on; `{tool}` stands for the tool's name. +type Served struct { + Subject string `json:"subject"` + Queue string `json:"queue,omitempty"` +} + +// SeatVerb is one verb of a seat a module holds, where it is answered. +type SeatVerb struct { + Seat string `json:"seat"` + Verb string `json:"verb"` + Subject string `json:"subject"` +} + +// Envelope is an event as a module emits it: the body is the payload, the metadata rides as headers +// (novox/hq ADR 0042). +type Envelope struct { + Key string `json:"key"` + Node string `json:"node,omitempty"` + Body json.RawMessage `json:"body"` + Headers map[string]string `json:"headers,omitempty"` +} + +// Answered is a call's result and the machine that gave it (novox/hq ADR 0159). +type Answered struct { + Result json.RawMessage + Node string +} + +// Handler answers one request body. +type Handler func(body json.RawMessage) (any, error) + +// ErrPin is a bus whose certificate is not the one the mesh pinned: final, never retried. +var ErrPin = errors.New("the bus's certificate does not match the pin") + +// Fatal is why a connection failure is final rather than "not yet", or "" when waiting may fix it — +// the same classification the TypeScript runtime makes. +func Fatal(err error) string { + if err == nil { + return "" + } + msg := err.Error() + if errors.Is(err, ErrPin) || strings.Contains(msg, ErrPin.Error()) { + return "the bus's certificate does not match the pin" + } + if regexp.MustCompile(`(?i)invalid url|no servers available for connection: .*url`).MatchString(msg) || + strings.Contains(msg, "nats: invalid url") { + return "the bus address is not a usable URL" + } + if regexp.MustCompile(`(?i)authorization violation|user authentication expired|permissions violation`).MatchString(msg) { + return "the bus refused this account" + } + return "" +} + +func normalizeFingerprint(f string) string { + f = strings.TrimSpace(f) + if len(f) > 7 && strings.EqualFold(f[:7], "sha256:") { + f = f[7:] + } + return strings.ToLower(strings.ReplaceAll(f, ":", "")) +} + +// pinned accepts exactly the certificate with this SHA-256 and no other. The pin is the only check: +// the bus's certificate names the seat, not the address a machine dials it by. +func pinned(want string) *tls.Config { + want = normalizeFingerprint(want) + return &tls.Config{ + InsecureSkipVerify: true, //nolint:gosec // replaced by the pin, which is stricter + MinVersion: tls.VersionTLS12, + VerifyPeerCertificate: func(raw [][]byte, _ [][]*x509.Certificate) error { + if len(raw) == 0 { + return fmt.Errorf("%w: it presented none", ErrPin) + } + sum := sha256.Sum256(raw[0]) + if got := hex.EncodeToString(sum[:]); got != want { + return fmt.Errorf("%w: it presented %s, not the pinned %s", ErrPin, got, want) + } + return nil + }, + } +} + +// Conn is the runtime's connection: the module it is, the memberships it follows, the subjects it +// answers. +type Conn struct { + nc *nats.Conn + js nats.JetStreamContext + self string + node string + cred Credential + mu sync.Mutex + issued map[string]*Membership // module → membership; present with nil = followed, none issued + onNew []func(Membership) + subs []*nats.Subscription + Logf func(format string, args ...any) +} + +// Connect dials the bus as the credential's module. A module's subjects come from its credential, +// never from its calls (ADR 0074). +func Connect(cred Credential) (*Conn, error) { + if cred.Module == "" { + return nil, errors.New("a broker credential with no module: the runtime derives its subjects " + + "from the account the mesh issued, and cannot guess which module it is") + } + node := cred.Node + if node == "" { + node = "?" + } + opts := []nats.Option{ + nats.Name(node + "." + cred.Module), + // Reconnect forever: the bus restarting is an upgrade, not a reason to exit. + nats.MaxReconnects(-1), + } + if cred.User != "" { + opts = append(opts, nats.UserInfo(cred.User, cred.Password), + // Its own inbox: every user's inbox is private to it (design 25 §4). + nats.CustomInboxPrefix("_INBOX."+cred.User)) + } + if strings.TrimSpace(cred.Fingerprint) != "" { + opts = append(opts, nats.Secure(pinned(cred.Fingerprint))) + } + nc, err := nats.Connect(cred.URL, opts...) + if err != nil { + return nil, err + } + js, err := nc.JetStream() + if err != nil { + nc.Close() + return nil, err + } + c := &Conn{nc: nc, js: js, self: cred.Module, node: cred.Node, cred: cred, + issued: map[string]*Membership{}, Logf: log.Printf} + c.Follow(cred.Module) + return c, nil +} + +// Module is what this connection is. +func (c *Conn) Module() string { return c.self } + +// Node is the machine this connection's account is scoped to. +func (c *Conn) Node() string { return c.node } + +// Credential is what this connection was opened with. +func (c *Conn) Credential() Credential { return c.cred } + +// MembershipSubject is the one address a runtime derives for an assignment (ADR 0160). +func MembershipSubject(node, module string) string { + return "mesh.assignment." + node + "." + module +} + +// Follow reads a module's membership on this machine once and follows it live, so its tools are +// served where the mesh issued them (ADR 0175). +func (c *Conn) Follow(module string) { + c.mu.Lock() + if c.node == "" { + c.mu.Unlock() + return + } + if _, has := c.issued[module]; has { + c.mu.Unlock() + return + } + c.issued[module] = nil + c.mu.Unlock() + subject := MembershipSubject(c.node, module) + // The subject-addressed direct get: the one address the mesh grants this account on the + // stream's API. + if got, err := c.nc.Request("$JS.API.DIRECT.GET."+assignmentsStream+"."+subject, nil, 5*time.Second); err == nil { + if got.Header.Get("Status") == "" && len(got.Data) > 0 { + var m Membership + if json.Unmarshal(got.Data, &m) == nil { + c.mu.Lock() + c.issued[module] = &m + c.mu.Unlock() + } + } + } + if c.Membership(module) == nil { + c.Logf("[mesh-tools] no membership issued for %s on %s yet; serving the derived shape until one arrives", module, c.node) + } + sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) { + var m Membership + if err := json.Unmarshal(msg.Data, &m); err != nil { + c.Logf("[mesh-tools] a membership arrived that is not one: %v", err) + return + } + c.mu.Lock() + c.issued[module] = &m + handlers := append([]func(Membership){}, c.onNew...) + c.mu.Unlock() + c.Logf("[mesh-tools] %s on %s was issued a new membership; re-serving on it", module, c.node) + for _, h := range handlers { + h(m) + } + _ = c.nc.Flush() + }) + if err == nil { + c.track(sub) + } +} + +// Membership is what the mesh issued a module here, or nil when nothing has been issued. +func (c *Conn) Membership(module string) *Membership { + c.mu.Lock() + defer c.mu.Unlock() + return c.issued[module] +} + +// Following says whether this connection follows a module's membership. +func (c *Conn) Following(module string) bool { + c.mu.Lock() + defer c.mu.Unlock() + _, has := c.issued[module] + return has +} + +// Serving is every module whose membership this connection follows. +func (c *Conn) Serving() []string { + c.mu.Lock() + defer c.mu.Unlock() + out := make([]string, 0, len(c.issued)) + for m := range c.issued { + out = append(out, m) + } + return out +} + +// OnMembership is called with every new membership any followed module is issued. +func (c *Conn) OnMembership(h func(Membership)) { + c.mu.Lock() + c.onNew = append(c.onNew, h) + c.mu.Unlock() +} + +func (c *Conn) track(s *nats.Subscription) { + c.mu.Lock() + c.subs = append(c.subs, s) + c.mu.Unlock() +} + +// servedOn is where a served module's tool is answered: its membership's subjects when issued, the +// derived shape otherwise (the shape the mesh issues on day one). +func (c *Conn) servedOn(module, tool string) []Served { + if m := c.Membership(module); m != nil { + out := make([]Served, 0, len(m.Serves)+1) + for _, s := range m.Serves { + out = append(out, Served{Subject: strings.ReplaceAll(s.Subject, "{tool}", tool), Queue: s.Queue}) + } + if tool == "tools" && m.Tools != "" { + found := false + for _, s := range out { + found = found || s.Subject == m.Tools + } + if !found { + out = append([]Served{{Subject: m.Tools, Queue: "serve." + module}}, out...) + } + } + return out + } + base := "mesh.mod." + module + ".tool." + tool + out := []Served{{Subject: base, Queue: "serve." + module}} + if c.node != "" { + out = append(out, Served{Subject: base + "." + c.node}) + } + return out +} + +type reply struct { + Result any `json:"result,omitempty"` + Error string `json:"error,omitempty"` + Node string `json:"node,omitempty"` +} + +// answerOn answers one subject with one handler, and says which machine answered (ADR 0159). +func (c *Conn) answerOn(subject, queue string, h Handler) (func(), error) { + cb := func(msg *nats.Msg) { + go func() { + var r reply + result, err := h(json.RawMessage(msg.Data)) + if err != nil { + r.Error = err.Error() + } else { + r.Result = nullable(result) + } + r.Node = c.node + body, _ := wire.Marshal(r) + _ = msg.Respond(body) + }() + } + var sub *nats.Subscription + var err error + if queue != "" { + sub, err = c.nc.QueueSubscribe(subject, queue, cb) + } else { + sub, err = c.nc.Subscribe(subject, cb) + } + if err != nil { + return func() {}, err + } + c.track(sub) + return func() { _ = sub.Unsubscribe() }, nil +} + +// nullable keeps a nil result as JSON null rather than dropping the key: the TypeScript reply always +// carries `result` when the handler did not throw. +func nullable(v any) any { + if v == nil { + return json.RawMessage("null") + } + return v +} + +// HandleSubject answers a subject outright: a seat's verb where the mesh issued it. +func (c *Conn) HandleSubject(subject string, h Handler) (func(), error) { + return c.answerOn(subject, "", h) +} + +// Handle serves `.` where the mesh issued that module, and follows its membership: +// when a new one arrives, it serves where it now says and stops where it no longer does. +func (c *Conn) Handle(key string, h Handler) (func(), error) { + if strings.HasPrefix(key, "seat:") { + subject, err := ToolSubject(key, c.self) + if err != nil { + return func() {}, err + } + return c.answerOn(subject, "", h) + } + module, tool := c.self, key + if dot := strings.Index(key, "."); dot >= 0 { + module, tool = key[:dot], key[dot+1:] + } + if module != c.self && !c.Following(module) { + return func() {}, fmt.Errorf("%s cannot serve %s: a module serves its own tools, and a runtime "+ + "those of the modules it follows", c.self, key) + } + var mu sync.Mutex + var stops []func() + serve := func() { + mu.Lock() + defer mu.Unlock() + for _, s := range stops { + s() + } + stops = nil + for _, s := range c.servedOn(module, tool) { + if stop, err := c.answerOn(s.Subject, s.Queue, h); err == nil { + stops = append(stops, stop) + } else { + c.Logf("[mesh-tools] cannot serve %s on %s: %v", key, s.Subject, err) + } + } + } + serve() + c.OnMembership(func(m Membership) { + if m.Module == module { + serve() + } + }) + return func() { + mu.Lock() + defer mu.Unlock() + for _, s := range stops { + s() + } + stops = nil + }, nil +} + +// reachedAt is where a call by key goes: a subject this connection's own membership says it +// reaches — the machine's when named — else the derived shape. +func (c *Conn) reachedAt(key string) (string, error) { + name, wanted, _ := strings.Cut(key, "@") + if m := c.Membership(c.self); m != nil { + if reach := m.Reaches[name]; len(reach) > 0 { + if wanted == "" { + return reach[0], nil + } + for _, s := range reach { + if strings.HasSuffix(s, "."+wanted) { + return s, nil + } + } + } + } + return ToolSubject(key, c.self) +} + +// Ask calls a tool by key — or on a subject the mesh listed for it — and learns which machine +// answered. Core request/reply: a tool call is never persisted (design 25 §3). +func (c *Conn) Ask(key string, body any, on string) (Answered, error) { + subject := on + if subject == "" { + s, err := c.reachedAt(key) + if err != nil { + return Answered{}, err + } + subject = s + } + data, err := wire.Marshal(body) + if err != nil { + return Answered{}, err + } + msg, err := c.nc.Request(subject, data, RequestTimeout) + if err != nil { + if errors.Is(err, nats.ErrNoResponders) { + return Answered{}, errors.New("503 no responders") + } + if errors.Is(err, nats.ErrTimeout) { + return Answered{}, errors.New("timeout") + } + return Answered{}, err + } + var r struct { + Result json.RawMessage `json:"result"` + Error string `json:"error"` + Node string `json:"node"` + } + if err := json.Unmarshal(msg.Data, &r); err != nil { + return Answered{}, err + } + if r.Error != "" { + return Answered{}, errors.New(r.Error) + } + return Answered{Result: r.Result, Node: r.Node}, nil +} + +// PublishAs emits an event as a module: published into JetStream and awaited, de-duplicated by its +// own id (ADR 0042). The body is the payload; the metadata rides as headers. +func (c *Conn) PublishAs(module string, env Envelope) error { + msg := nats.NewMsg("mesh.mod." + module + ".event." + env.Key) + for k, v := range env.Headers { + msg.Header.Set(k, v) + } + if env.Headers["content-type"] == "" { + msg.Header.Set("content-type", "application/json") + } + if env.Node != "" { + msg.Header.Set("x-node", env.Node) + } + body := env.Body + if len(body) == 0 { + body = json.RawMessage("null") + } + msg.Data = body + var opts []nats.PubOpt + if id := env.Headers["x-event-id"]; id != "" { + opts = append(opts, nats.MsgId(id)) + } + _, err := c.js.PublishMsg(msg, opts...) + return err +} + +// Flush waits until the bus has every subscription made so far, so what is served is answerable +// when this returns. +func (c *Conn) Flush() { _ = c.nc.Flush() } + +// Close unsubscribes everything and drains, so an in-flight reply is finished rather than dropped. +func (c *Conn) Close() { + c.mu.Lock() + subs := c.subs + c.subs = nil + c.mu.Unlock() + for _, s := range subs { + _ = s.Unsubscribe() + } + _ = c.nc.Drain() +} + +// ToolSubject is a tool's subject. A bare name is this module's own; `.` another's; +// `seat:.` a role's, with `@` for a node-scoped seat (design 33 §4). +func ToolSubject(key, self string) (string, error) { + if strings.HasPrefix(key, "seat:") { + rest := strings.TrimPrefix(key, "seat:") + dot := strings.Index(rest, ".") + if dot < 0 { + return "", fmt.Errorf("%q names a seat and no verb: seat:.", key) + } + seat := rest[:dot] + verb, node, _ := strings.Cut(rest[dot+1:], "@") + if node != "" { + return "mesh.seat." + seat + ".tool." + verb + "." + node, nil + } + return "mesh.seat." + seat + ".tool." + verb, nil + } + name, node, _ := strings.Cut(key, "@") + var base string + if dot := strings.Index(name, "."); dot < 0 { + base = "mesh.mod." + self + ".tool." + name + } else { + base = "mesh.mod." + name[:dot] + ".tool." + name[dot+1:] + } + if node != "" { + return base + "." + node, nil + } + return base, nil +} + +// SeatToolSubject is a seat's verb as its holder serves it: flat for a mesh seat, carrying the +// machine for a node-scoped one. +func SeatToolSubject(seat, verb, scope, node string) string { + base := "mesh.seat." + seat + ".tool." + verb + if scope == "node" && node != "" { + return base + "." + node + } + return base +} diff --git a/node-tools/internal/bus/bus_test.go b/node-tools/internal/bus/bus_test.go new file mode 100644 index 0000000..a74c066 --- /dev/null +++ b/node-tools/internal/bus/bus_test.go @@ -0,0 +1,48 @@ +package bus + +import ( + "errors" + "fmt" + "testing" +) + +func TestSubjectsAreTheOnesTheTypeScriptRuntimeUses(t *testing.T) { + for key, want := range map[string]string{ + "status": "mesh.mod.self.tool.status", + "alpha.one": "mesh.mod.alpha.tool.one", + "alpha.one@anchor": "mesh.mod.alpha.tool.one.anchor", + "seat:node-shelf.list": "mesh.seat.node-shelf.tool.list", + "seat:node-shelf.list@anchor": "mesh.seat.node-shelf.tool.list.anchor", + } { + if got, err := ToolSubject(key, "self"); err != nil || got != want { + t.Errorf("%s: %s %v, want %s", key, got, err, want) + } + } + if _, err := ToolSubject("seat:nothing", "self"); err == nil { + t.Error("a seat with no verb was accepted") + } + if SeatToolSubject("s", "v", "node", "n") != "mesh.seat.s.tool.v.n" || SeatToolSubject("s", "v", "mesh", "n") != "mesh.seat.s.tool.v" { + t.Error("seat subjects") + } +} + +func TestAFingerprintIsReadHoweverItIsWritten(t *testing.T) { + for _, f := range []string{"sha256:AB:CD:ef", "abcdef", "ABCDEF", "SHA256:abcdef"} { + if got := normalizeFingerprint(f); got != "abcdef" { + t.Errorf("%s → %s", f, got) + } + } +} + +func TestWhatWaitingCannotFixIsFinal(t *testing.T) { + for err, want := range map[error]string{ + fmt.Errorf("x509: %w: it presented aa", ErrPin): "the bus's certificate does not match the pin", + errors.New("nats: Authorization Violation"): "the bus refused this account", + errors.New("nats: invalid url"): "the bus address is not a usable URL", + errors.New("dial tcp 10.0.0.1:4222: connect: connection refused"): "", + } { + if got := Fatal(err); got != want { + t.Errorf("%v → %q, want %q", err, got, want) + } + } +} diff --git a/node-tools/internal/console/client.go b/node-tools/internal/console/client.go new file mode 100644 index 0000000..9cd8c0d --- /dev/null +++ b/node-tools/internal/console/client.go @@ -0,0 +1,240 @@ +package console + +import ( + "encoding/json" + "errors" + "fmt" + "regexp" + "sort" + "strings" + "sync" + + "github.com/novox/mesh-tools/node-tools/internal/bus" +) + +const ( + catalogueModules = "mesh-catalog.catalog_modules" + seatTools = "seat:mesh-controller.tools" + toolsVerb = "tools" +) + +// Tool is a tool as its module — or, for a role's tool, the mesh's records — describes it. +type Tool struct { + Module string + Name string + Description string + Input json.RawMessage + Seat bool + Scope string + Subjects []string +} + +// Listing is what the mesh could say about its tools; silence is named, never dropped (design 34 §3). +type Listing struct { + Tools []Tool + NotAnswering []string +} + +// Seats is the seats and their verbs, so `.` resolves to the role. +type Seats map[string]map[string]bool + +func seatsIn(l *Listing) Seats { + out := Seats{} + for _, t := range l.Tools { + if !t.Seat { + continue + } + if out[t.Module] == nil { + out[t.Module] = map[string]bool{} + } + out[t.Module][t.Name] = true + } + return out +} + +// toolKey is the key a call uses: a role's when the prefix is a seat declaring that verb. +func toolKey(name string, seats Seats) string { + if strings.HasPrefix(name, "seat:") { + return name + } + dot := strings.Index(name, ".") + if dot < 0 { + return name + } + if seats[name[:dot]][name[dot+1:]] { + return "seat:" + name + } + return name +} + +// toolsOn asks the mesh what tools it has: the catalogue which modules it holds, each module what it +// serves, the controller's seat every role's tools — at once, so a restarting control plane hides +// nothing else. +func toolsOn(conn *bus.Conn) (*Listing, error) { + type rolesAnswer struct { + Seats []struct { + Seat string `json:"seat"` + Scope string `json:"scope"` + Tools []struct { + Name string `json:"name"` + Description string `json:"description"` + Input json.RawMessage `json:"input"` + } `json:"tools"` + } `json:"seats"` + } + var wg sync.WaitGroup + var roles *rolesAnswer + wg.Add(1) + go func() { + defer wg.Done() + if got, err := conn.Ask(seatTools, map[string]any{}, ""); err == nil { + var r rolesAnswer + if json.Unmarshal(got.Result, &r) == nil { + roles = &r + } + } + }() + answered, err := conn.Ask(catalogueModules, map[string]any{}, "") + if err != nil { + wg.Wait() + return nil, err + } + var held struct { + Modules []struct { + Module string `json:"module"` + } `json:"modules"` + } + _ = json.Unmarshal(answered.Result, &held) + names := make([]string, 0, len(held.Modules)) + for _, m := range held.Modules { + if m.Module != "" { + names = append(names, m.Module) + } + } + type outcome struct { + ok bool + answer struct { + Tools *[]struct { + Name string `json:"name"` + Description string `json:"description"` + Input json.RawMessage `json:"input"` + Subjects []string `json:"subjects"` + } `json:"tools"` + Failed *string `json:"failed"` + } + } + outcomes := make([]outcome, len(names)) + for i, module := range names { + wg.Add(1) + go func(i int, module string) { + defer wg.Done() + got, err := conn.Ask(module+"."+toolsVerb, map[string]any{}, "") + if err != nil { + return + } + if json.Unmarshal(got.Result, &outcomes[i].answer) == nil { + outcomes[i].ok = true + } + }(i, module) + } + wg.Wait() + l := &Listing{Tools: []Tool{}, NotAnswering: []string{}} + if roles != nil { + for _, s := range roles.Seats { + for _, t := range s.Tools { + l.Tools = append(l.Tools, Tool{Module: s.Seat, Name: t.Name, Description: t.Description, + Input: t.Input, Seat: true, Scope: s.Scope}) + } + } + } else { + l.NotAnswering = append(l.NotAnswering, "mesh-controller (seat)") + } + for i, module := range names { + o := outcomes[i] + switch { + case o.ok && o.answer.Failed != nil: + l.NotAnswering = append(l.NotAnswering, fmt.Sprintf("%s (its tools bundle failed to load: %s)", module, *o.answer.Failed)) + case o.ok && o.answer.Tools != nil: + for _, t := range *o.answer.Tools { + l.Tools = append(l.Tools, Tool{Module: module, Name: t.Name, Description: t.Description, + Input: t.Input, Subjects: t.Subjects}) + } + default: + l.NotAnswering = append(l.NotAnswering, module) + } + } + sort.SliceStable(l.Tools, func(i, j int) bool { + return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name + }) + sort.Strings(l.NotAnswering) + return l, nil +} + +// callTool calls `.[@]`, on the subject the listing names for it when it names one. +func callTool(conn *bus.Conn, key string, args any, seats Seats, l *Listing) (bus.Answered, error) { + name, node, _ := strings.Cut(key, "@") + if !strings.Contains(name, ".") { + return bus.Answered{}, fmt.Errorf("%q does not name a tool: write ., as `mesh tools` lists "+ + "them, or .@ for the instance on one machine", key) + } + resolved := toolKey(name, seats) + if node != "" { + resolved += "@" + node + } + return conn.Ask(resolved, args, subjectListed(name, node, l)) +} + +func subjectListed(name, node string, l *Listing) string { + if l == nil { + return "" + } + dot := strings.Index(name, ".") + module, tool := name[:dot], name[dot+1:] + for _, t := range l.Tools { + if t.Module == module && t.Name == tool && !t.Seat { + if len(t.Subjects) == 0 { + return "" + } + if node == "" { + return t.Subjects[0] + } + for _, s := range t.Subjects { + if strings.HasSuffix(s, "."+node) { + return s + } + } + return "" + } + } + return "" +} + +var ( + noResponders = regexp.MustCompile(`(?i)no responders|503`) + refused = regexp.MustCompile(`(?i)permissions violation|authorization`) + timedOut = regexp.MustCompile(`(?i)timeout`) +) + +// whyItFailed says why a call failed, so the remedy is in the words. +func whyItFailed(key string, err error) string { + if err == nil { + err = errors.New("failed") + } + msg := err.Error() + switch { + case noResponders.MatchString(msg): + extra := "" + if strings.HasPrefix(key, "seat:") { + extra = ", or nothing holds that seat" + } + return "nothing serves " + key + ". The module may not be assigned to any machine, or it is down" + + extra + " — `mesh tools` lists what answered." + case refused.MatchString(msg): + return "this account may not call " + key + ". What it may call was fixed when it was issued — a " + + "person's by `operator issue`, the console's by its manifest." + case timedOut.MatchString(msg): + return key + " did not answer in time. Something is serving it, so this is the tool being slow " + + "rather than absent." + } + return key + " failed: " + msg +} diff --git a/node-tools/internal/console/console_test.go b/node-tools/internal/console/console_test.go new file mode 100644 index 0000000..0645ae9 --- /dev/null +++ b/node-tools/internal/console/console_test.go @@ -0,0 +1,136 @@ +package console + +import ( + "bytes" + "encoding/json" + "io" + "net/http" + "strings" + "testing" + + "github.com/novox/mesh-tools/node-tools/internal/bus" + mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" + "github.com/novox/mesh-tools/node-tools/internal/runtime" +) + +func connect(t *testing.T, module, node string) *bus.Conn { + t.Helper() + c, err := bus.Connect(bus.Credential{URL: mt.URL(t), Module: module, Node: node}) + if err != nil { + t.Fatal(err) + } + c.Logf = func(string, ...any) {} + t.Cleanup(c.Close) + return c +} + +func post(t *testing.T, endpoint string, body string) map[string]any { + t.Helper() + res, err := http.Post(endpoint, "application/json", bytes.NewBufferString(body)) + if err != nil { + t.Fatal(err) + } + defer res.Body.Close() + raw, _ := io.ReadAll(res.Body) + var out map[string]any + if err := json.Unmarshal(raw, &out); err != nil { + t.Fatalf("%d %s", res.StatusCode, raw) + } + return out +} + +// As node-tools, the runtime serves the bundles and is the console on loopback: the listing is what +// the modules and the mesh's records answered, and a call reaches the module on the machine named. +func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) { + mesh := mt.New(t) + mesh.Issue(t, mt.MembershipOf("alpha", "desk", true, nil)) + mesh.Issue(t, mt.MembershipOf("beta", "desk", false, map[string][]string{"node-shelf": {"list", "clear"}})) + nodeTools := connect(t, "node-tools", "desk") + stop, err := runtime.Run(nodeTools, []runtime.Served{ + {Module: "alpha", Entrypoints: []string{mt.Fixture("many-alpha.serve.mjs")}}, + {Module: "beta", Entrypoints: []string{mt.Fixture("many-beta.serve.mjs")}}, + }, nil, (&mt.Logs{}).Logf) + if err != nil { + t.Fatal(err) + } + defer stop() + + catalogue := connect(t, "mesh-catalog", "") + stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) { + return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "ghost"}}}, nil + }) + defer stopCat() + controller := connect(t, "mesh-controller", "") + stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.tools", func(json.RawMessage) (any, error) { + return map[string]any{"seats": []map[string]any{{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{ + {"name": "list", "description": "what is on the shelf", "input": map[string]any{}}, + {"name": "clear", "description": "take it all off", "input": map[string]any{}}, + }}}}, nil + }) + defer stopSeat() + catalogue.Flush() + controller.Flush() + + up, err := Serve(NewSurface(nodeTools, "desk.node-tools"), "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer up.Close() + endpoint := "http://" + up.Address + "/mcp" + + init := post(t, endpoint, `{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-03-26","capabilities":{}}}`) + if !strings.Contains(init["result"].(map[string]any)["instructions"].(string), "reached as desk.node-tools") { + t.Errorf("initialize: %v", init) + } + 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) { + names = append(names, x.(map[string]any)["name"].(string)) + } + if got := strings.Join(names, ","); got != "alpha.one,alpha.two,beta.five,beta.four,beta.three,node-shelf.clear,node-shelf.list" { + t.Errorf("listed %s", got) + } + if got := listed["_meta"].(map[string]any)["notAnswering"]; len(got.([]any)) != 1 || got.([]any)[0] != "ghost" { + t.Errorf("not answering: %v", got) + } + for _, x := range listed["tools"].([]any) { + tool := x.(map[string]any) + if tool["name"] == "node-shelf.list" { + schema := tool["inputSchema"].(map[string]any) + if req, _ := schema["required"].([]any); len(req) != 1 || req[0] != "node" { + t.Errorf("a node seat's verb does not require node: %v", schema) + } + } + } + + called := post(t, endpoint, `{"jsonrpc":"2.0","id":3,"method":"tools/call","params":{"name":"alpha.one","arguments":{"node":"desk"}}}`) + content := called["result"].(map[string]any)["content"].([]any) + var got map[string]any + _ = json.Unmarshal([]byte(content[0].(map[string]any)["text"].(string)), &got) + if got["alpha"] != float64(1) || content[1].(map[string]any)["text"] != "answered by desk" { + t.Errorf("called: %v", called) + } + seat := post(t, endpoint, `{"jsonrpc":"2.0","id":4,"method":"tools/call","params":{"name":"node-shelf.list","arguments":{"node":"desk"}}}`) + if text := seat["result"].(map[string]any)["content"].([]any)[0].(map[string]any)["text"].(string); !strings.Contains(text, `"a"`) { + t.Errorf("seat verb: %v", seat) + } + refused := post(t, endpoint, `{"jsonrpc":"2.0","id":5,"method":"tools/call","params":{"name":"node-shelf.list","arguments":{}}}`) + if refused["error"] == nil { + t.Errorf("a node seat's verb was called without its machine: %v", refused) + } + absent := post(t, endpoint, `{"jsonrpc":"2.0","id":6,"method":"tools/call","params":{"name":"ghost.boo","arguments":{}}}`) + result := absent["result"].(map[string]any) + if result["isError"] != true || !strings.Contains(result["content"].([]any)[0].(map[string]any)["text"].(string), "nothing serves ghost.boo") { + t.Errorf("an absent tool: %v", absent) + } + res, err := http.Post(endpoint, "application/json", strings.NewReader(`{"jsonrpc":"2.0","method":"notifications/initialized"}`)) + if err != nil || res.StatusCode != 202 { + t.Errorf("a notification: %v %v", res, err) + } +} + +func TestTheConsoleListensOnLoopbackAndNowhereElse(t *testing.T) { + if _, err := Serve(NewSurface(nil, "x"), "0.0.0.0:0"); err == nil || !strings.Contains(err.Error(), "loopback and nowhere else") { + t.Errorf("a non-loopback console was not refused: %v", err) + } +} diff --git a/node-tools/internal/console/http.go b/node-tools/internal/console/http.go new file mode 100644 index 0000000..40e19d0 --- /dev/null +++ b/node-tools/internal/console/http.go @@ -0,0 +1,125 @@ +package console + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net" + "net/http" + "strconv" + "strings" + + "github.com/novox/mesh-tools/node-tools/internal/wire" +) + +// bodyLimit is the most a request body may be. +const bodyLimit = 1 << 20 + +var loopback = map[string]bool{"127.0.0.1": true, "::1": true, "localhost": true, "[::1]": true} + +// Listening is a console that listens: where, with the port the machine gave, and how to stop it. +type Listening struct { + Address string + Close func() error +} + +// Serve listens on host:port, refused unless the host is loopback — said before binding, so a +// console that would open to a network is a startup failure (ADR 0152). +func Serve(s *Surface, listen string) (*Listening, error) { + at := strings.LastIndex(listen, ":") + if at < 0 { + return nil, fmt.Errorf("%q is not host:port", listen) + } + host, portText := listen[:at], listen[at+1:] + if !loopback[host] { + return nil, fmt.Errorf(`the console listens on loopback and nowhere else (novox/hq ADR 0152): %q is not this `+ + "machine's own address — whoever is on the machine owns the mesh there, and nobody else may reach this", host) + } + port, err := strconv.Atoi(portText) + if err != nil || port < 0 || port > 65535 { + return nil, fmt.Errorf("%q is not a port", portText) + } + ln, err := net.Listen("tcp", net.JoinHostPort(strings.Trim(host, "[]"), portText)) + if err != nil { + return nil, err + } + server := &http.Server{Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { route(w, r, s) })} + go func() { _ = server.Serve(ln) }() + bound := ln.Addr().(*net.TCPAddr).Port + return &Listening{Address: host + ":" + strconv.Itoa(bound), Close: func() error { return server.Shutdown(context.Background()) }}, nil +} + +func route(w http.ResponseWriter, r *http.Request, s *Surface) { + switch r.URL.Path { + case "/": + w.Header().Set("content-type", "text/plain; charset=utf-8") + _, _ = io.WriteString(w, "the mesh's console: MCP over HTTP at POST /mcp (novox/hq design 34)\n") + return + case "/mcp": + default: + writeJSON(w, 404, map[string]any{"error": "the console serves /mcp and nothing else"}) + return + } + switch r.Method { + case http.MethodPost: + case http.MethodDelete: + w.WriteHeader(204) // no session to end + return + default: + w.Header().Set("allow", "POST, DELETE") + w.WriteHeader(405) + return + } + body, err := io.ReadAll(io.LimitReader(r.Body, bodyLimit+1)) + if err != nil || len(body) > bodyLimit { + if err == nil { + err = fmt.Errorf("the request is larger than %d bytes", bodyLimit) + } + writeJSON(w, 413, map[string]any{"jsonrpc": "2.0", "id": nil, "error": map[string]any{"code": -32600, "message": err.Error()}}) + return + } + trimmed := strings.TrimSpace(string(body)) + if strings.HasPrefix(trimmed, "[") { + var batch []Request + if err := json.Unmarshal(body, &batch); err != nil { + writeJSON(w, 400, map[string]any{"jsonrpc": "2.0", "id": nil, "error": map[string]any{"code": -32700, "message": "the body is not JSON"}}) + return + } + replies := []*Reply{} + for _, req := range batch { + if reply := s.Handle(req); reply != nil { + replies = append(replies, reply) + } + } + if len(replies) == 0 { + w.WriteHeader(202) + return + } + writeJSON(w, 200, replies) + return + } + var req Request + if err := json.Unmarshal(body, &req); err != nil { + writeJSON(w, 400, map[string]any{"jsonrpc": "2.0", "id": nil, "error": map[string]any{"code": -32700, "message": "the body is not JSON"}}) + return + } + reply := s.Handle(req) + if reply == nil { + w.WriteHeader(202) + return + } + writeJSON(w, 200, reply) +} + +func writeJSON(w http.ResponseWriter, status int, body any) { + text, err := wire.Marshal(body) + if err != nil { + text, _ = json.Marshal(map[string]any{"error": errors.New("unencodable answer").Error()}) + } + w.Header().Set("content-type", "application/json; charset=utf-8") + w.Header().Set("content-length", strconv.Itoa(len(text))) + w.WriteHeader(status) + _, _ = w.Write(text) +} diff --git a/node-tools/internal/console/mcp.go b/node-tools/internal/console/mcp.go new file mode 100644 index 0000000..07de2b1 --- /dev/null +++ b/node-tools/internal/console/mcp.go @@ -0,0 +1,268 @@ +// Package console is the mesh's tools as an MCP server on a machine's loopback (novox/hq design 34, +// ADR 0152, ADR 0175 §6): the Go port of node-tools' http.ts, mcp.ts and client.ts. A thin adapter: +// every tool listed is one a module answered for, the schema is the module's, the answer the module's. +package console + +import ( + "encoding/json" + "strings" + "sync" + "time" + + "github.com/novox/mesh-tools/node-tools/internal/bus" +) + +// Protocol is the MCP version spoken to an agent host. +const Protocol = "2025-03-26" + +// ListingKept is how long a fetched tool list is kept before the modules are asked again. +var ListingKept = 30 * time.Second + +// Request is one JSON-RPC message from a host. +type Request struct { + JSONRPC string `json:"jsonrpc"` + ID json.RawMessage `json:"id,omitempty"` + Method string `json:"method"` + Params json.RawMessage `json:"params,omitempty"` +} + +// Reply is one JSON-RPC answer. +type Reply struct { + JSONRPC string `json:"jsonrpc"` + ID json.RawMessage `json:"id"` + Result any `json:"result,omitempty"` + Error *RPCError `json:"error,omitempty"` +} + +// RPCError is a protocol-level refusal. +type RPCError struct { + Code int `json:"code"` + Message string `json:"message"` +} + +// Surface answers MCP requests over one bus connection, as one account. +type Surface struct { + conn *bus.Conn + who string + mu sync.Mutex + known *Listing + at time.Time +} + +// NewSurface is the surface over a connection, as `who`. +func NewSurface(conn *bus.Conn, who string) *Surface { return &Surface{conn: conn, who: who} } + +func (s *Surface) listing() (*Listing, error) { + s.mu.Lock() + if s.known != nil && time.Since(s.at) <= ListingKept { + l := s.known + s.mu.Unlock() + return l, nil + } + s.mu.Unlock() + l, err := toolsOn(s.conn) + if err != nil { + return nil, err + } + s.mu.Lock() + s.known, s.at = l, time.Now() + s.mu.Unlock() + return l, nil +} + +func isNotification(id json.RawMessage) bool { + t := strings.TrimSpace(string(id)) + return t == "" || t == "null" +} + +func answer(id json.RawMessage, result any) *Reply { + return &Reply{JSONRPC: "2.0", ID: idOrNull(id), Result: result} +} + +func refuse(id json.RawMessage, code int, message string) *Reply { + return &Reply{JSONRPC: "2.0", ID: idOrNull(id), Error: &RPCError{Code: code, Message: message}} +} + +func idOrNull(id json.RawMessage) json.RawMessage { + if isNotification(id) { + return json.RawMessage("null") + } + return id +} + +// Handle answers one request; nil for a notification, which expects none. +func (s *Surface) Handle(r Request) *Reply { + notification := isNotification(r.ID) + switch r.Method { + case "initialize": + return answer(r.ID, map[string]any{ + "protocolVersion": Protocol, + "capabilities": map[string]any{"tools": map[string]any{}}, + "serverInfo": map[string]any{"name": "mesh", "version": "1"}, + "instructions": "These are the tools of a Novox mesh, reached as " + s.who + ". Every call goes to the module " + + "that serves it; what may be called was fixed when this account was issued, so a " + + "refusal means the account, not the tool. The list is what the running modules " + + "answered, plus every role's tools from the mesh's records — the mesh's own verbs " + + "(mesh-controller.status, .push, .assign …) among them; a module that did not answer " + + "is named in the list's _meta and can still be called by ..", + }) + case "notifications/initialized": + return nil + case "ping": + if notification { + return nil + } + return answer(r.ID, map[string]any{}) + case "tools/list": + l, err := s.listing() + if err != nil { + return refuse(r.ID, -32603, whyItFailed(catalogueModules, err)) + } + tools := make([]map[string]any, 0, len(l.Tools)) + for _, t := range l.Tools { + var schema map[string]any + switch { + case t.Seat && t.Scope != "node": + schema = asSchema(t.Input) + case t.Seat: + schema = withNode(asSchema(t.Input), "the machine whose seat answers; required, the seat is held once per machine", true) + default: + schema = withNode(asSchema(t.Input), "", false) + } + description := t.Description + if description == "" { + description = t.Name + ", served by " + t.Module + } + tools = append(tools, map[string]any{"name": t.Module + "." + t.Name, "description": description, "inputSchema": schema}) + } + return answer(r.ID, map[string]any{"tools": tools, "_meta": map[string]any{"notAnswering": l.NotAnswering}}) + case "tools/call": + var p struct { + Name string `json:"name"` + Arguments map[string]any `json:"arguments"` + } + _ = json.Unmarshal(r.Params, &p) + args := map[string]any{} + for k, v := range p.Arguments { + args[k] = v + } + l, _ := s.listing() + var roles Seats + if l != nil { + roles = seatsIn(l) + } + bare, _, _ := strings.Cut(p.Name, "@") + isSeatVerb := roles != nil && strings.HasPrefix(toolKey(bare, roles), "seat:") + nodeScoped := false + if isSeatVerb && l != nil { + for _, t := range l.Tools { + if t.Seat && t.Scope == "node" && t.Module+"."+t.Name == bare { + nodeScoped = true + } + } + } + takesNode := !isSeatVerb || nodeScoped + node := "" + if takesNode { + if n, ok := args["node"].(string); ok { + node = n + } + delete(args, "node") + } + if nodeScoped && node == "" && !strings.Contains(p.Name, "@") { + return refuse(r.ID, -32602, p.Name+" is a machine's seat's verb: name the machine with `node`") + } + name := p.Name + if node != "" && !strings.Contains(p.Name, "@") { + name = p.Name + "@" + node + } + got, err := callTool(s.conn, name, args, roles, l) + if err != nil { + return answer(r.ID, map[string]any{ + "content": []map[string]any{{"type": "text", "text": whyItFailed(name, err)}}, + "isError": true, + }) + } + 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 answer(r.ID, map[string]any{"content": content}) + } + if notification { + return nil + } + return refuse(r.ID, -32601, "mesh's MCP surface has no "+r.Method) +} + +// pretty is a module's answer as JSON text, indented as JSON.stringify(result, null, 2) writes it. +func pretty(raw json.RawMessage) string { + if len(raw) == 0 { + return "null" + } + var v any + if json.Unmarshal(raw, &v) != nil { + return string(raw) + } + b, err := json.MarshalIndent(v, "", " ") + if err != nil { + return string(raw) + } + return strings.NewReplacer(`<`, "<", `>`, ">", `&`, "&").Replace(string(b)) +} + +// asSchema is a module's declared input as a JSON schema: wrapped when it is a bare map of +// properties, passed through when it is a schema, empty when nothing was declared. +func asSchema(raw json.RawMessage) map[string]any { + var given map[string]any + if json.Unmarshal(raw, &given) != nil || given == nil { + return map[string]any{"type": "object", "properties": map[string]any{}} + } + if given["type"] == "object" { + return given + } + if _, has := given["properties"]; has { + return given + } + if len(given) == 0 { + return map[string]any{"type": "object", "properties": map[string]any{}} + } + return map[string]any{"type": "object", "properties": given} +} + +// withNode adds the optional — or, for a node seat, required — `node` argument (ADR 0159). +func withNode(schema map[string]any, description string, required bool) map[string]any { + properties := map[string]any{} + if p, ok := schema["properties"].(map[string]any); ok { + for k, v := range p { + properties[k] = v + } + } + if _, has := properties["node"]; !has { + if description == "" { + description = "the machine to ask, when this module runs on several; else whichever answers, and the answer says which" + } + properties["node"] = map[string]any{"type": "string", "description": description} + } + out := map[string]any{} + for k, v := range schema { + out[k] = v + } + out["type"] = "object" + out["properties"] = properties + if required { + have := []any{} + if r, ok := schema["required"].([]any); ok { + have = r + } + hasNode := false + for _, x := range have { + hasNode = hasNode || x == "node" + } + if !hasNode { + have = append(have, "node") + } + out["required"] = have + } + return out +} diff --git a/node-tools/internal/launch/launch.go b/node-tools/internal/launch/launch.go new file mode 100644 index 0000000..45caf5d --- /dev/null +++ b/node-tools/internal/launch/launch.go @@ -0,0 +1,401 @@ +// Package launch starts a bundle the runtime serves and speaks MCP over stdio to it (novox/hq ADR +// 0188, ADR 0193): `initialize`, `tools/list` once, `tools/call` per call. A Go binary, a Python +// script and a Node launcher are the same thing here: an executable that answers those. The runtime +// knows no language; it starts the path it is given. +package launch + +import ( + "bufio" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "os/exec" + "regexp" + "strings" + "sync" + "syscall" + "time" + + "golang.org/x/sys/unix" + + "github.com/novox/mesh-tools/node-tools/internal/wire" +) + +// Protocol is the MCP version spoken to a bundle. +const Protocol = "2025-03-26" + +// How long a child has for its handshake, and a call before the caller is told it is slow. +var ( + HandshakeTimeout = 10 * time.Second + CallTimeout = 30 * time.Second +) + +// Tool is one tool a launched bundle listed, and how to call it. +type Tool struct { + Name string + Description string + Input json.RawMessage + Run func(args json.RawMessage) (json.RawMessage, error) +} + +// Registration is the tools a bundle listed under one name: its module's, or a seat's. +type Registration struct { + Module string + Tools []Tool +} + +// Publisher publishes an event a bundle asked the runtime to emit, as the bundle's module. +type Publisher func(params json.RawMessage) error + +// Launched is a running bundle: what it registered, and how to stop it. +type Launched struct { + Registrations []Registration + Stop func() +} + +// Executable says whether an entrypoint can be started: a bundle the runtime serves is executable, +// and one that is not was not built to be served (ADR 0193). +func Executable(path string) bool { + info, err := os.Stat(path) + if err != nil || info.IsDir() { + return false + } + return info.Mode().Perm()&0o111 != 0 +} + +type message struct { + JSONRPC string `json:"jsonrpc,omitempty"` + ID json.RawMessage `json:"id,omitempty"` + Method string `json:"method,omitempty"` + Params json.RawMessage `json:"params,omitempty"` + Result json.RawMessage `json:"result,omitempty"` + Error *struct { + Code int `json:"code"` + Message string `json:"message"` + } `json:"error,omitempty"` +} + +type child struct { + cmd *exec.Cmd + stdin io.WriteCloser + writeMu sync.Mutex + mu sync.Mutex + pending map[int64]chan message + next int64 + dead chan struct{} + why error +} + +func (c *child) write(m any) error { + b, err := wire.Marshal(m) + if err != nil { + return err + } + c.writeMu.Lock() + defer c.writeMu.Unlock() + _, err = c.stdin.Write(append(b, '\n')) + return err +} + +var stackLine = regexp.MustCompile(`^\s+at\s`) + +// Start launches one bundle and learns its tools. It fails when the child cannot be started or does +// not complete the handshake. A child that exits later is started again on its next call. +func Start(module, entry string, env []string, publish Publisher, logf func(string, ...any)) (*Launched, error) { + var mu sync.Mutex + var current *child + stopped := false + + start := func() (*child, error) { + cmd := exec.Command(entry) + cmd.Env = env + stdin, err := cmd.StdinPipe() + if err != nil { + return nil, err + } + stdout, err := cmd.StdoutPipe() + if err != nil { + return nil, err + } + stderr, err := cmd.StderrPipe() + if err != nil { + return nil, err + } + if err := cmd.Start(); err != nil { + return nil, err + } + c := &child{cmd: cmd, stdin: stdin, pending: map[int64]chan message{}, next: 1, dead: make(chan struct{})} + var lastSaid string + var saidMu sync.Mutex + stderrDone := make(chan struct{}) + go func() { + defer close(stderrDone) + scan := bufio.NewScanner(stderr) + scan.Buffer(make([]byte, 64*1024), 1<<20) + for scan.Scan() { + line := scan.Text() + if strings.TrimSpace(line) == "" { + continue + } + logf("[%s] %s", module, line) + if !stackLine.MatchString(line) && !strings.HasPrefix(line, "Node.js v") { + saidMu.Lock() + lastSaid = strings.TrimSpace(line) + saidMu.Unlock() + } + } + }() + go func() { + scan := bufio.NewScanner(stdout) + scan.Buffer(make([]byte, 64*1024), 16<<20) + for scan.Scan() { + line := strings.TrimSpace(scan.Text()) + if line == "" { + continue + } + var m message + if err := json.Unmarshal([]byte(line), &m); err != nil { + if len(line) > 120 { + line = line[:120] + } + logf("[mesh-tools] %s's bundle said something that is not a reply: %s", module, line) + continue + } + // The bundle asks the runtime to emit (ADR 0193): published as this module, answered + // once the bus has accepted it. Nothing else a bundle may ask. + if m.Method != "" { + go func(m message) { + if m.Method != "mesh/publish" { + if len(m.ID) > 0 { + _ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, + "error": map[string]any{"code": -32601, "message": "the runtime answers no " + m.Method + " from a bundle"}}) + } + return + } + err := publish(m.Params) + if len(m.ID) == 0 { + return + } + if err != nil { + _ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, + "error": map[string]any{"code": -32000, "message": err.Error()}}) + return + } + _ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, "result": map[string]any{}}) + }(m) + continue + } + var id int64 + if json.Unmarshal(m.ID, &id) != nil { + continue + } + c.mu.Lock() + ch := c.pending[id] + delete(c.pending, id) + c.mu.Unlock() + if ch != nil { + ch <- m + } + } + <-stderrDone + err := cmd.Wait() + code := "0" + if exit := (*exec.ExitError)(nil); errors.As(err, &exit) { + if status, ok := exit.Sys().(syscall.WaitStatus); ok && status.Signaled() { + code = unix.SignalName(status.Signal()) + } else { + code = fmt.Sprint(exit.ExitCode()) + } + } + saidMu.Lock() + why := fmt.Sprintf("%s's bundle exited (%s)", module, code) + if lastSaid != "" { + why += ": " + lastSaid + } + saidMu.Unlock() + c.mu.Lock() + c.why = errors.New(why) + c.mu.Unlock() + close(c.dead) + mu.Lock() + if current == c { + current = nil + } + wasStopped := stopped + mu.Unlock() + if !wasStopped { + logf("[mesh-tools] %s; started again on its next call", why) + } + }() + if _, err := c.ask(module, "initialize", map[string]any{"protocolVersion": Protocol, "capabilities": map[string]any{}, + "clientInfo": map[string]any{"name": "node-tools", "version": "1"}}, HandshakeTimeout); err != nil { + _ = cmd.Process.Kill() + return nil, err + } + _ = c.write(map[string]any{"jsonrpc": "2.0", "method": "notifications/initialized"}) + return c, nil + } + + asking := func(method string, params any, timeout time.Duration) (json.RawMessage, error) { + mu.Lock() + c := current + mu.Unlock() + if c == nil { + fresh, err := start() + if err != nil { + return nil, err + } + mu.Lock() + current = fresh + c = fresh + mu.Unlock() + } + return c.ask(module, method, params, timeout) + } + + first, err := start() + if err != nil { + return nil, err + } + mu.Lock() + current = first + mu.Unlock() + + raw, err := asking("tools/list", map[string]any{}, HandshakeTimeout) + if err != nil { + return nil, err + } + var listed struct { + Tools []struct { + Name string `json:"name"` + Description string `json:"description"` + InputSchema json.RawMessage `json:"inputSchema"` + } `json:"tools"` + } + if err := json.Unmarshal(raw, &listed); err != nil { + return nil, fmt.Errorf("%s's bundle listed its tools in a shape that is not MCP's: %w", module, err) + } + order := []string{} + groups := map[string][]Tool{} + for _, t := range listed.Tools { + under, name := module, t.Name + if dot := strings.Index(t.Name, "."); dot >= 0 { + under, name = t.Name[:dot], t.Name[dot+1:] + } + full := t.Name + input := t.InputSchema + if len(input) == 0 || string(input) == "null" { + input = json.RawMessage("{}") + } + tool := Tool{Name: name, Description: t.Description, Input: input, + Run: func(args json.RawMessage) (json.RawMessage, error) { + if len(args) == 0 || string(args) == "null" { + args = json.RawMessage("{}") + } + res, err := asking("tools/call", map[string]any{"name": full, "arguments": args}, CallTimeout) + if err != nil { + return nil, err + } + var called struct { + Content []struct { + Type string `json:"type"` + Text string `json:"text"` + } `json:"content"` + IsError bool `json:"isError"` + } + _ = json.Unmarshal(res, &called) + text := "" + for _, c := range called.Content { + if c.Type == "text" { + text = c.Text + break + } + } + if called.IsError { + if text == "" { + text = module + "." + full + " failed" + } + return nil, errors.New(text) + } + // The bundle's answer is JSON as text; handed back as the value it encodes. + if json.Valid([]byte(text)) && text != "" { + return json.RawMessage(text), nil + } + b, _ := json.Marshal(text) + return b, nil + }} + if _, seen := groups[under]; !seen { + order = append(order, under) + } + groups[under] = append(groups[under], tool) + } + out := &Launched{Stop: func() { + mu.Lock() + stopped = true + c := current + current = nil + mu.Unlock() + if c != nil && c.cmd.Process != nil { + _ = c.cmd.Process.Signal(syscall.SIGTERM) + } + }} + for _, under := range order { + out.Registrations = append(out.Registrations, Registration{Module: under, Tools: groups[under]}) + } + return out, nil +} + +func (c *child) ask(module, method string, params any, timeout time.Duration) (json.RawMessage, error) { + c.mu.Lock() + if c.why != nil { + c.mu.Unlock() + return nil, c.why + } + id := c.next + c.next++ + ch := make(chan message, 1) + c.pending[id] = ch + c.mu.Unlock() + if err := c.write(map[string]any{"jsonrpc": "2.0", "id": id, "method": method, "params": params}); err != nil { + c.mu.Lock() + delete(c.pending, id) + c.mu.Unlock() + select { + case <-c.dead: + return nil, c.deathReason() + case <-time.After(100 * time.Millisecond): + return nil, err + } + } + timer := time.NewTimer(timeout) + defer timer.Stop() + select { + case m := <-ch: + if m.Error != nil { + msg := m.Error.Message + if msg == "" { + msg = "the bundle refused the request" + } + return nil, errors.New(msg) + } + return m.Result, nil + case <-c.dead: + return nil, c.deathReason() + case <-timer.C: + c.mu.Lock() + delete(c.pending, id) + c.mu.Unlock() + return nil, fmt.Errorf("%s's bundle did not answer %s in %ds", module, method, int(timeout/time.Second)) + } +} + +func (c *child) deathReason() error { + c.mu.Lock() + defer c.mu.Unlock() + if c.why != nil { + return c.why + } + return errors.New("the bundle exited") +} diff --git a/node-tools/internal/meshtest/meshtest.go b/node-tools/internal/meshtest/meshtest.go new file mode 100644 index 0000000..185eb0c --- /dev/null +++ b/node-tools/internal/meshtest/meshtest.go @@ -0,0 +1,143 @@ +// 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. +package meshtest + +import ( + "encoding/json" + "os" + "path/filepath" + "runtime" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-tools/node-tools/internal/bus" +) + +// URL is the test bus, or the test is skipped. +func URL(t *testing.T) string { + t.Helper() + url := os.Getenv("MESH_TEST_NATS") + if url == "" { + t.Skip("MESH_TEST_NATS unset") + } + return url +} + +// Mesh is the controller's job, done by hand. +type Mesh struct { + nc *nats.Conn + js nats.JetStreamContext +} + +// New raises the streams afresh. +func New(t *testing.T) *Mesh { + t.Helper() + nc, err := nats.Connect(URL(t)) + if err != nil { + t.Fatal(err) + } + js, _ := nc.JetStream() + for _, s := range []string{"ASSIGNMENTS", "EVENTS"} { + _ = js.DeleteStream(s) + } + if _, err := js.AddStream(&nats.StreamConfig{Name: "ASSIGNMENTS", Subjects: []string{"mesh.assignment.>"}, + MaxMsgsPerSubject: 1, AllowDirect: true}); err != nil { + t.Fatal(err) + } + if _, err := js.AddStream(&nats.StreamConfig{Name: "EVENTS", Subjects: []string{"mesh.mod.*.event.>"}}); err != nil { + t.Fatal(err) + } + t.Cleanup(nc.Close) + return &Mesh{nc: nc, js: js} +} + +// Issue publishes a membership. +func (m *Mesh) Issue(t *testing.T, mem bus.Membership) { + t.Helper() + body, _ := json.Marshal(mem) + if _, err := m.js.Publish(bus.MembershipSubject(mem.Node, mem.Module), body); err != nil { + t.Fatal(err) + } +} + +// NextEvent is the subject the next event under a pattern lands on. +func (m *Mesh) NextEvent(t *testing.T, pattern string) <-chan string { + t.Helper() + ch := make(chan string, 1) + sub, err := m.nc.Subscribe(pattern, func(msg *nats.Msg) { + select { + case ch <- msg.Subject: + default: + } + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = sub.Unsubscribe() }) + _ = m.nc.Flush() + return ch +} + +// MembershipOf is a membership as the controller issues one on a machine. +func MembershipOf(module, node string, plain bool, seats map[string][]string) bus.Membership { + own := "mesh.mod." + module + serves := []bus.Served{{Subject: own + ".tool.{tool}." + node}} + if plain { + serves = append(serves, bus.Served{Subject: own + ".tool.{tool}", Queue: "serve." + module}) + } + var verbs []bus.SeatVerb + for seat, vs := range seats { + for _, v := range vs { + verbs = append(verbs, bus.SeatVerb{Seat: seat, Verb: v, Subject: "mesh.seat." + seat + ".tool." + v + "." + node}) + } + } + return bus.Membership{Node: node, Module: module, Serves: serves, Seats: verbs, + Emits: own + ".event.{event}", Tools: own + ".tool.tools"} +} + +// Fixture is a path in node-tools/test/fixtures, which these tests share with the TypeScript ones. +func Fixture(name string) string { + _, here, _, _ := runtime.Caller(0) + return filepath.Join(filepath.Dir(here), "..", "..", "test", "fixtures", name) +} + +// Until retries while the bus answers "no responders" — something not yet served. +func Until(t *testing.T, try func() error) { + t.Helper() + var err error + for i := 0; i < 50; i++ { + if err = try(); err == nil { + return + } + time.Sleep(100 * time.Millisecond) + } + t.Fatal(err) +} + +// Logs collects what the runtime says. +type Logs struct{ lines []string } + +// Logf is a logger that keeps the lines. +func (l *Logs) Logf(format string, args ...any) { + l.lines = append(l.lines, sprintf(format, args...)) +} + +// Has says whether a line contains every fragment. +func (l *Logs) Has(fragments ...string) bool { + for _, line := range l.lines { + all := true + for _, f := range fragments { + all = all && strings.Contains(line, f) + } + if all { + return true + } + } + return false +} + +// All is every line. +func (l *Logs) All() string { return strings.Join(l.lines, "\n") } diff --git a/node-tools/internal/meshtest/sprintf.go b/node-tools/internal/meshtest/sprintf.go new file mode 100644 index 0000000..19a5819 --- /dev/null +++ b/node-tools/internal/meshtest/sprintf.go @@ -0,0 +1,5 @@ +package meshtest + +import "fmt" + +func sprintf(format string, args ...any) string { return fmt.Sprintf(format, args...) } diff --git a/node-tools/internal/runtime/runtime.go b/node-tools/internal/runtime/runtime.go new file mode 100644 index 0000000..674d2df --- /dev/null +++ b/node-tools/internal/runtime/runtime.go @@ -0,0 +1,378 @@ +// Package runtime is the node's tool runtime (novox/hq ADR 0175, ADR 0193), the Go port of +// node-tools' runtime.ts in its launch-only form: it launches every assigned module's bundle, +// serves each module's tools on that module's subjects and each held seat's verbs on the seat's, +// and answers for every module the verb that says what it serves. It imports nothing and knows no +// language. +package runtime + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "sort" + "strings" + "sync" + + "github.com/novox/mesh-tools/node-tools/internal/bus" + "github.com/novox/mesh-tools/node-tools/internal/launch" +) + +// ToolsVerb is the verb every module's runtime answers for it (ADR 0152): its tools, from the code +// that answers them. +const ToolsVerb = "tools" + +// Words the mesh sets for the runtime (ADR 0175, ADR 0192). +const ( + ToolModules = "MESH_TOOL_MODULES" + ToolEnv = "MESH_TOOL_ENV" + OperatorAccount = "MESH_OPERATOR_ACCOUNT" + OperatorHome = "MESH_OPERATOR_HOME" +) + +// Served is one module this runtime serves and its entrypoints. +type Served struct { + Module string + Entrypoints []string +} + +// ToolsAnswer is what `tools` answers for one module. +type ToolsAnswer struct { + Module string `json:"module"` + Tools []ListedTool `json:"tools"` + Failed string `json:"failed,omitempty"` +} + +// ListedTool is one tool as a module's `tools` answer lists it. +type ListedTool struct { + Name string `json:"name"` + Description string `json:"description"` + Input json.RawMessage `json:"input"` + Subjects []string `json:"subjects,omitempty"` +} + +// ServedModulesFrom reads MESH_TOOL_MODULES: `=` entries, comma-separated, several +// per module. The one-module form — a bare path, or the runtime's own module — is the per-module +// containers' (to-be 38 WP4c) and refused here: the node's runtime imports nothing. +func ServedModulesFrom(spec, own string) ([]Served, error) { + order := []string{} + by := map[string][]string{} + for _, raw := range strings.Split(spec, ",") { + entry := strings.TrimSpace(raw) + if entry == "" { + continue + } + module, path, ok := strings.Cut(entry, "=") + module, path = strings.TrimSpace(module), strings.TrimSpace(path) + if !ok || module == "" || path == "" || module == own { + return nil, fmt.Errorf("%s: %q is not = of another module; the node's "+ + "runtime launches the bundles it is given and imports nothing (novox/hq ADR 0193)", ToolModules, entry) + } + if _, seen := by[module]; !seen { + order = append(order, module) + } + by[module] = append(by[module], path) + } + out := make([]Served, 0, len(order)) + for _, m := range order { + out = append(out, Served{Module: m, Entrypoints: by[m]}) + } + return out, nil +} + +// TakeToolEnvs reads the composed environments (ADR 0192) and removes them from the process's, so no +// bundle finds another's there. +func TakeToolEnvs() (map[string]map[string]string, error) { + raw := os.Getenv(ToolEnv) + os.Unsetenv(ToolEnv) + out := map[string]map[string]string{} + if raw == "" { + return out, nil + } + var parsed map[string]map[string]any + if err := json.Unmarshal([]byte(raw), &parsed); err != nil { + return nil, fmt.Errorf(`%s is not JSON of the shape {"": {"": ""}}: %w`, ToolEnv, err) + } + for module, words := range parsed { + own := map[string]string{} + for k, v := range words { + if s, ok := v.(string); ok { + own[k] = s + } else { + b, _ := json.Marshal(v) + own[k] = string(b) + } + } + out[module] = own + } + return out, nil +} + +type registration struct { + module string // the module's own name, or a seat's + owner string // the module whose bundle made it + tools []launch.Tool +} + +// Run launches, binds and serves. It answers a stop function. +func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, logf func(string, ...any)) (func(), error) { + if account := os.Getenv(OperatorAccount); account != "" { + home := "" + if h := os.Getenv(OperatorHome); h != "" { + home = " (home " + h + ")" + } + logf("[mesh-tools] the operator's account here is %s%s", account, home) + } + modules := make([]string, 0, len(served)) + isServed := map[string]bool{} + for _, s := range served { + modules = append(modules, s.Module) + isServed[s.Module] = true + conn.Follow(s.Module) + } + node := conn.Node() + + base := os.Environ() + envFor := func(module string) []string { + words := map[string]string{} + for _, kv := range base { + if k, v, ok := strings.Cut(kv, "="); ok && k != ToolEnv { + words[k] = v + } + } + for k, v := range envs[module] { + words[k] = v + } + words["MESH_SERVED_MODULE"] = module + words["MESH_MODULE"] = module + if node != "" { + words["MESH_NODE"] = node + } + out := make([]string, 0, len(words)) + for k, v := range words { + out = append(out, k+"="+v) + } + sort.Strings(out) + return out + } + + failed := map[string]string{} + var registrations []registration + var stops []func() + for _, s := range served { + for _, entry := range s.Entrypoints { + path, _ := filepath.Abs(entry) + module := s.Module + fail := func(why string) { + failed[module] = why + logf("[mesh-tools] %s's bundle %s failed to load: %s; its tools are not served here", module, entry, why) + } + if !launch.Executable(path) { + fail(path + " is not executable; a bundle the runtime serves is started, never imported, and its build makes it executable (novox/hq ADR 0193)") + continue + } + child, err := launch.Start(module, path, envFor(module), func(params json.RawMessage) error { + var env bus.Envelope + if err := json.Unmarshal(params, &env); err != nil { + return fmt.Errorf("not an event envelope: %w", err) + } + return conn.PublishAs(module, env) + }, logf) + if err != nil { + fail(err.Error()) + continue + } + stops = append(stops, child.Stop) + for _, r := range child.Registrations { + registrations = append(registrations, registration{module: r.Module, owner: module, tools: r.Tools}) + } + } + } + + claimed := map[string]bool{} + for _, m := range modules { + if mem := conn.Membership(m); mem != nil { + for _, s := range mem.Seats { + claimed[s.Seat] = true + } + } + } + var own []registration + for _, r := range registrations { + switch { + case isServed[r.module]: + own = append(own, r) + case claimed[r.module]: + default: + logf(`[mesh-tools] %s registers tools under "%s", which is neither a module served here nor a seat one of them claims; not served until the mesh issues the claim`, r.owner, r.module) + } + } + stopAll := func() { + for i := len(stops) - 1; i >= 0; i-- { + stops[i]() + } + } + for _, r := range own { + seen := map[string]bool{} + for _, t := range r.tools { + if t.Name == ToolsVerb { + stopAll() + return nil, fmt.Errorf(`%s names a tool "%s", which is the verb the runtime answers for every module with what it serves (novox/hq ADR 0152) — refused, rename it`, r.module, ToolsVerb) + } + if seen[t.Name] { + stopAll() + return nil, fmt.Errorf("%s exposes two tools named %s — refused", r.module, t.Name) + } + seen[t.Name] = true + } + } + + var names []string + byModule := map[string][]launch.Tool{} + for _, r := range own { + for _, t := range r.tools { + t := t + stop, err := conn.Handle(r.module+"."+t.Name, func(body json.RawMessage) (any, error) { + return t.Run(argsOf(body)) + }) + if err != nil { + logf("[mesh-tools] cannot serve %s.%s: %v", r.module, t.Name, err) + continue + } + names = append(names, r.module+"."+t.Name) + stops = append(stops, stop) + } + byModule[r.module] = append(byModule[r.module], r.tools...) + } + for _, module := range modules { + module := module + tools := byModule[module] + why := failed[module] + if len(tools) == 0 && why == "" { + continue // a pure-events module: silent, as it always was + } + stop, err := conn.Handle(module+"."+ToolsVerb, func(json.RawMessage) (any, error) { + answer := ToolsAnswer{Module: module, Tools: []ListedTool{}, Failed: why} + for _, t := range tools { + answer.Tools = append(answer.Tools, ListedTool{Name: t.Name, Description: t.Description, + Input: t.Input, Subjects: subjectsOf(conn, module, t.Name)}) + } + return answer, nil + }) + if err == nil { + stops = append(stops, stop) + } + } + + failedNames := make([]string, 0, len(failed)) + for _, m := range modules { + if _, f := failed[m]; f { + failedNames = append(failedNames, m) + } + } + line := fmt.Sprintf("[mesh-tools] serving %d tool(s) for %d module(s): %s", len(names), len(modules), orNone(names)) + if len(failedNames) > 0 { + line += fmt.Sprintf("; not serving %s, whose bundle(s) failed to load", strings.Join(failedNames, ", ")) + } + logf("%s", line) + stops = append(stops, serveSeats(conn, modules, registrations, logf)) + conn.Flush() + return stopAll, nil +} + +func orNone(names []string) string { + if len(names) == 0 { + return "(none)" + } + return strings.Join(names, ", ") +} + +// argsOf is a call's arguments as the tool receives them: an object, `{}` for none. +func argsOf(body json.RawMessage) json.RawMessage { + trimmed := strings.TrimSpace(string(body)) + if trimmed == "" || trimmed == "null" { + return json.RawMessage("{}") + } + return body +} + +// subjectsOf is where a tool is answered as the mesh issued it: the plain subject first, then this +// machine's; nothing before a membership is issued. +func subjectsOf(conn *bus.Conn, module, tool string) []string { + m := conn.Membership(module) + if m == nil { + return nil + } + var plain, mine []string + for _, s := range m.Serves { + subject := strings.ReplaceAll(s.Subject, "{tool}", tool) + if s.Queue != "" { + plain = append(plain, subject) + } else { + mine = append(mine, subject) + } + } + return append(plain, mine...) +} + +// serveSeats serves every verb of every seat a served module holds, where the mesh issued it, by +// the tool of the same name registered under the seat's name (ADR 0159, 0160) — and serves again +// whenever a membership changes. Whether this machine holds the seat is the bus's to decide. +func serveSeats(conn *bus.Conn, modules []string, registrations []registration, logf func(string, ...any)) func() { + 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 + } + } + var mu sync.Mutex + var stops []func() + serve := func() { + mu.Lock() + defer mu.Unlock() + for _, s := range stops { + s() + } + stops = nil + have := map[string]bool{} + for _, module := range modules { + m := conn.Membership(module) + if m == nil { + continue + } + for _, v := range m.Seats { + if have[v.Subject] { + continue + } + have[v.Subject] = true + t, ok := impl[v.Seat][v.Verb] + if !ok { + logf("[mesh-tools] %s claims %s and implements no %s, which that seat promises; not served", module, v.Seat, v.Verb) + continue + } + stop, err := conn.HandleSubject(v.Subject, func(body json.RawMessage) (any, error) { + return t.Run(argsOf(body)) + }) + if err != nil { + logf("[mesh-tools] cannot serve %s's %s on %s: %v", v.Seat, v.Verb, v.Subject, err) + continue + } + stops = append(stops, stop) + logf("[mesh-tools] serving %s's %s on %s, admitted where %s holds the seat", v.Seat, v.Verb, v.Subject, module) + } + } + } + serve() + conn.OnMembership(func(bus.Membership) { go func() { serve(); conn.Flush() }() }) + return func() { + mu.Lock() + defer mu.Unlock() + for _, s := range stops { + s() + } + stops = nil + } +} diff --git a/node-tools/internal/runtime/runtime_test.go b/node-tools/internal/runtime/runtime_test.go new file mode 100644 index 0000000..0cb991b --- /dev/null +++ b/node-tools/internal/runtime/runtime_test.go @@ -0,0 +1,251 @@ +package runtime + +import ( + "encoding/json" + "os" + "strings" + "testing" + "time" + + "github.com/novox/mesh-tools/node-tools/internal/bus" + mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" +) + +func connect(t *testing.T, module, node string) *bus.Conn { + t.Helper() + c, err := bus.Connect(bus.Credential{URL: mt.URL(t), Module: module, Node: node}) + if err != nil { + t.Fatal(err) + } + c.Logf = func(string, ...any) {} + t.Cleanup(c.Close) + return c +} + +func call(t *testing.T, asker *bus.Conn, key string, args any) (string, error) { + t.Helper() + got, err := asker.Ask(key, args, "") + return string(got.Result), err +} + +func same(t *testing.T, got, want string) { + t.Helper() + var a, b any + if json.Unmarshal([]byte(got), &a) != nil || json.Unmarshal([]byte(want), &b) != nil { + t.Fatalf("not JSON: got %s want %s", got, want) + } + ga, _ := json.Marshal(a) + gb, _ := json.Marshal(b) + if string(ga) != string(gb) { + t.Errorf("got %s, want %s", got, want) + } +} + +// The node's runtime serves five modules' bundles on one credential — two TypeScript, one broken, +// one Python, one written against the protocol — and follows a re-issued membership (ADR 0175, 0193). +func TestTheNodesRuntimeServesFiveModulesAndFollowsAReissuedMembership(t *testing.T) { + mesh := mt.New(t) + mesh.Issue(t, mt.MembershipOf("alpha", "anchor", true, nil)) + mesh.Issue(t, mt.MembershipOf("beta", "anchor", false, map[string][]string{"node-shelf": {"list", "clear"}})) + mesh.Issue(t, mt.MembershipOf("gamma", "anchor", false, nil)) + mesh.Issue(t, mt.MembershipOf("delta", "anchor", false, map[string][]string{"node-lamp": {"on"}})) + mesh.Issue(t, mt.MembershipOf("epsilon", "anchor", false, nil)) + nodeTools := connect(t, "node-tools", "anchor") + asker := connect(t, "console", "workstation") + logs := &mt.Logs{} + t.Setenv("MESH_OPERATOR_ACCOUNT", "somebody") + t.Setenv("MESH_OPERATOR_HOME", "/home/somebody") + stop, err := Run(nodeTools, []Served{ + {"alpha", []string{mt.Fixture("many-alpha.serve.mjs")}}, + {"beta", []string{mt.Fixture("many-beta.serve.mjs")}}, + {"gamma", []string{mt.Fixture("many-broken.serve.mjs")}}, + {"delta", []string{mt.Fixture("many-delta.py")}}, + {"epsilon", []string{mt.Fixture("many-epsilon.mjs")}}, + }, nil, logs.Logf) + if err != nil { + t.Fatal(err) + } + defer stop() + if !logs.Has("the operator's account here is somebody (home /home/somebody)") { + t.Errorf("the operator was not said:\n%s", logs.All()) + } + if !logs.Has("gamma's bundle", "many-broken.serve.mjs failed to load: gamma's bundle exited (1): Error: gamma's bundle cannot find its client; its tools are not served here") { + t.Errorf("the broken bundle was not named with its own words:\n%s", logs.All()) + } + if !logs.Has("serving 8 tool(s) for 5 module(s): alpha.one, alpha.two, beta.three, beta.four, beta.five, delta.greet, delta.die, epsilon.seven; not serving gamma") { + t.Errorf("not serving what it should:\n%s", logs.All()) + } + + for key, want := range map[string]string{ + "alpha.one": `{"alpha":1}`, "alpha.one@anchor": `{"alpha":1}`, "beta.three@anchor": `{"beta":3}`, + "beta.four@anchor": `{"beta":4}`, "beta.five@anchor": `{"beta":5}`, + "seat:node-shelf.list@anchor": `{"shelf":["a","b"]}`, "seat:node-shelf.clear@anchor": `{"cleared":true}`, + "seat:node-lamp.on@anchor": `{"on":true,"language":"python"}`, "epsilon.seven@anchor": `{"epsilon":7,"via":"stdio"}`, + } { + got, err := call(t, asker, key, map[string]any{}) + if err != nil { + t.Errorf("%s: %v", key, err) + continue + } + same(t, got, want) + } + if _, err := call(t, asker, "beta.three", map[string]any{}); err == nil || !strings.Contains(err.Error(), "no responders") { + t.Errorf("beta was not issued the plain subject, and answered on it: %v", err) + } + got, err := call(t, asker, "delta.greet@anchor", map[string]any{"who": "mesh"}) + if err != nil { + t.Fatal(err) + } + same(t, got, `{"greeting":"hello mesh","language":"python"}`) + + // A launched tool that emits does so as its module, through the runtime. + landed := mesh.NextEvent(t, "mesh.mod.*.event.>") + got, err = call(t, asker, "alpha.two", map[string]any{}) + if err != nil { + t.Fatal(err) + } + same(t, got, `{"alpha":2}`) + select { + case subject := <-landed: + if subject != "mesh.mod.alpha.event.happened" { + t.Errorf("the event landed on %s", subject) + } + case <-time.After(5 * time.Second): + t.Error("the event never landed") + } + + // A bundle that dies mid-call is said, and started again on its next call. + if _, err := call(t, asker, "delta.die@anchor", map[string]any{}); err == nil { + t.Error("a bundle that died answered") + } + got, err = call(t, asker, "delta.greet@anchor", map[string]any{"who": "again"}) + if err != nil { + t.Fatalf("not started again: %v", err) + } + same(t, got, `{"greeting":"hello again","language":"python"}`) + + // `tools` answers for each, and why gamma serves nothing. + gamma, err := call(t, asker, "gamma.tools@anchor", map[string]any{}) + if err != nil { + t.Fatal(err) + } + same(t, gamma, `{"module":"gamma","tools":[],"failed":"gamma's bundle exited (1): Error: gamma's bundle cannot find its client"}`) + beta, _ := call(t, asker, "beta.tools@anchor", map[string]any{}) + var answer ToolsAnswer + _ = json.Unmarshal([]byte(beta), &answer) + if len(answer.Tools) != 3 || answer.Tools[0].Name != "three" || strings.Join(answer.Tools[0].Subjects, ",") != "mesh.mod.beta.tool.three.anchor" { + t.Errorf("beta's tools answer: %s", beta) + } + + // Re-issued mid-run, now answering for the module anywhere: served without a restart. + mesh.Issue(t, mt.MembershipOf("beta", "anchor", true, map[string][]string{"node-shelf": {"list", "clear"}})) + mt.Until(t, func() error { _, err := call(t, asker, "beta.three", map[string]any{}); return err }) + got, err = call(t, asker, "seat:node-shelf.list@anchor", map[string]any{}) + if err != nil { + t.Fatal(err) + } + same(t, got, `{"shelf":["a","b"]}`) +} + +// Each bundle is given its own environment and none of another's (ADR 0192), and the composed +// environments are not left in the runtime's. +func TestEachBundleIsGivenItsOwnEnvironment(t *testing.T) { + mesh := mt.New(t) + for _, m := range []string{"gamma", "delta", "zeta"} { + mesh.Issue(t, mt.MembershipOf(m, "anchor", false, nil)) + } + nodeTools := connect(t, "node-tools", "anchor") + asker := connect(t, "console", "workstation") + t.Setenv("MESH_OPERATOR_ACCOUNT", "somebody") + t.Setenv(ToolEnv, `{"gamma":{"GAMMA_CONFIG_FILE":"/var/lib/mesh/gamma/config.json"},"delta":{"DELTA_TOKEN_FILE":"/var/lib/mesh/delta/token"},"zeta":{"ZETA_URL":"http://127.0.0.1:3000"}}`) + envs, err := TakeToolEnvs() + if err != nil { + t.Fatal(err) + } + if _, left := os.LookupEnv(ToolEnv); left { + t.Error("the composed environments were left in the runtime's") + } + logs := &mt.Logs{} + stop, err := Run(nodeTools, []Served{ + {"gamma", []string{mt.Fixture("env-gamma.serve.mjs")}}, + {"delta", []string{mt.Fixture("env-delta.serve.mjs")}}, + {"zeta", []string{mt.Fixture("env-zeta.mjs")}}, + }, envs, logs.Logf) + if err != nil { + t.Fatal(err) + } + defer stop() + for key, want := range map[string]string{ + "gamma.given@anchor": `{"mine":"/var/lib/mesh/gamma/config.json","theirs":null,"runtime":"somebody","composed":null}`, + "delta.given@anchor": `{"mine":"/var/lib/mesh/delta/token","theirs":null}`, + "zeta.given@anchor": `{"mine":"http://127.0.0.1:3000","theirs":null,"composed":null}`, + } { + got, err := call(t, asker, key, map[string]any{}) + if err != nil { + t.Fatalf("%s: %v\n%s", key, err, logs.All()) + } + same(t, got, want) + } +} + +// An entrypoint that is not executable is refused by name, and the others serve (ADR 0193). +func TestAnEntrypointThatIsNotExecutableIsRefused(t *testing.T) { + mesh := mt.New(t) + mesh.Issue(t, mt.MembershipOf("alpha", "anchor", false, nil)) + mesh.Issue(t, mt.MembershipOf("plain", "anchor", false, nil)) + nodeTools := connect(t, "node-tools", "anchor") + asker := connect(t, "console", "workstation") + logs := &mt.Logs{} + stop, err := Run(nodeTools, []Served{ + {"alpha", []string{mt.Fixture("many-alpha.serve.mjs")}}, + {"plain", []string{mt.Fixture("many-alpha.mjs")}}, + }, nil, logs.Logf) + if err != nil { + t.Fatal(err) + } + defer stop() + if !logs.Has("plain's bundle", "many-alpha.mjs failed to load:", "is not executable; a bundle the runtime serves is started, never imported") { + t.Errorf("the non-executable entrypoint was not refused by name:\n%s", logs.All()) + } + got, err := call(t, asker, "alpha.one@anchor", map[string]any{}) + if err != nil { + t.Fatal(err) + } + same(t, got, `{"alpha":1}`) +} + +// A launched bundle is told the module it serves, so its seat's verbs stay the seat's (ADR 0193). +func TestALaunchedBundleRegisteringItsSeatFirstServesTheSeat(t *testing.T) { + mesh := mt.New(t) + mesh.Issue(t, mt.MembershipOf("theta", "anchor", false, map[string][]string{"node-shelf": {"list"}})) + nodeTools := connect(t, "node-tools", "anchor") + asker := connect(t, "console", "workstation") + stop, err := Run(nodeTools, []Served{{"theta", []string{mt.Fixture("served-seat-first.mjs")}}}, + map[string]map[string]string{"theta": {"THETA_WORD": "given"}}, (&mt.Logs{}).Logf) + if err != nil { + t.Fatal(err) + } + defer stop() + got, err := call(t, asker, "theta.own@anchor", map[string]any{}) + if err != nil { + t.Fatal(err) + } + same(t, got, `{"theta":"given"}`) + got, err = call(t, asker, "seat:node-shelf.list@anchor", map[string]any{}) + if err != nil { + t.Fatal(err) + } + same(t, got, `{"shelf":["x"]}`) +} + +func TestToolModulesNamesOtherModulesOnly(t *testing.T) { + got, err := ServedModulesFrom(" alpha=/a/tools/index.serve.mjs, beta=/b/one, beta=/b/two ", "node-tools") + if err != nil || len(got) != 2 || got[1].Module != "beta" || len(got[1].Entrypoints) != 2 { + t.Fatalf("%v %v", got, err) + } + for _, bad := range []string{"/mine/index.js", "node-tools=/own.js", "=/x"} { + if _, err := ServedModulesFrom(bad, "node-tools"); err == nil { + t.Errorf("%q was accepted", bad) + } + } +} diff --git a/node-tools/internal/wire/wire.go b/node-tools/internal/wire/wire.go new file mode 100644 index 0000000..a65a6c3 --- /dev/null +++ b/node-tools/internal/wire/wire.go @@ -0,0 +1,18 @@ +// Package wire is JSON as the TypeScript runtime writes it: no HTML escaping of <, > and &. +package wire + +import ( + "bytes" + "encoding/json" +) + +// Marshal encodes v the way JSON.stringify does, without a trailing newline. +func Marshal(v any) ([]byte, error) { + var b bytes.Buffer + enc := json.NewEncoder(&b) + enc.SetEscapeHTML(false) + if err := enc.Encode(v); err != nil { + return nil, err + } + return bytes.TrimRight(b.Bytes(), "\n"), nil +} diff --git a/node-tools/module.json b/node-tools/module.json index e0e8d85..e5630b0 100644 --- a/node-tools/module.json +++ b/node-tools/module.json @@ -35,10 +35,10 @@ { "name": "runtime", "kind": "bundle", - "language": "typescript", - "entrypoints": [ - "src/main.js" - ] + "language": "go", + "system": "arch", + "from": "cmd/node-tools", + "binary": "node-tools" } ] }