Merge pull request 'node-tools in Go: the node's runtime, launch-only, wire-compatible with the TypeScript (hq ADR 0193)' (#38) from feat/0193-node-tools-in-go into main
This commit was merged in pull request #38.
This commit is contained in:
@@ -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 <module>=<entrypoint>,… the bundles to launch
|
||||||
|
// MESH_TOOL_ENV {"<module>": {"<word>": "<value>"}} 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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
)
|
||||||
@@ -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=
|
||||||
@@ -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.<module>.event.<type> an event a module emits
|
||||||
|
// mesh.mod.<module>.tool.<tool> a tool a module serves
|
||||||
|
// mesh.seat.<seat>.tool.<verb> 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 `<module>.<tool>` 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; `<module>.<tool>` another's;
|
||||||
|
// `seat:<seat>.<verb>` a role's, with `@<node>` 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:<seat>.<verb>", 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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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 `<seat>.<verb>` 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 `<module>.<tool>[@<node>]`, 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 <module>.<tool>, as `mesh tools` lists "+
|
||||||
|
"them, or <module>.<tool>@<node> 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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -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 <module>.<tool>.",
|
||||||
|
})
|
||||||
|
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
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
@@ -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") }
|
||||||
@@ -0,0 +1,5 @@
|
|||||||
|
package meshtest
|
||||||
|
|
||||||
|
import "fmt"
|
||||||
|
|
||||||
|
func sprintf(format string, args ...any) string { return fmt.Sprintf(format, args...) }
|
||||||
@@ -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: `<module>=<entrypoint>` 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 <module>=<entrypoint> 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 {"<module>": {"<word>": "<value>"}}: %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
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -35,10 +35,10 @@
|
|||||||
{
|
{
|
||||||
"name": "runtime",
|
"name": "runtime",
|
||||||
"kind": "bundle",
|
"kind": "bundle",
|
||||||
"language": "typescript",
|
"language": "go",
|
||||||
"entrypoints": [
|
"system": "arch",
|
||||||
"src/main.js"
|
"from": "cmd/node-tools",
|
||||||
]
|
"binary": "node-tools"
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user