Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
df4f492a72 | ||
|
|
66e8be0e31 | ||
|
|
7722668220 | ||
|
|
e8989f3cf5 | ||
|
|
915f372a85 | ||
|
|
486dad99a5 | ||
|
|
65f3b68076 | ||
|
|
42e0987c64 | ||
|
|
ae6bdc9b06 | ||
|
|
b182943c24 |
@@ -0,0 +1,209 @@
|
||||
// Package announce answers the NATS services protocol's discovery for what a runtime serves, and
|
||||
// gathers the answers (novox/hq ADR 0197).
|
||||
//
|
||||
// A runtime does not re-serve its tools through a services library: serving is unchanged. It answers
|
||||
// `$SRV.PING`, `$SRV.INFO` and `$SRV.STATS` — and the same followed by its service name, and by its
|
||||
// name and id — in the format NATS's own tools read, with what it is serving at the moment it is
|
||||
// asked. One service per runtime process: the bus admits one reply per request from each responder,
|
||||
// so a runtime serving many modules and seats answers once, one endpoint per tool per subject, and
|
||||
// says in each endpoint's metadata which module, seat, scope and machine it is.
|
||||
package announce
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go/micro"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
)
|
||||
|
||||
// Version is the announced service version (semver, as the protocol requires).
|
||||
const Version = "0.1.0"
|
||||
|
||||
// Window is how long the console gathers discovery answers: every instance answers one request, and
|
||||
// how many will is what is being found out.
|
||||
var Window = 750 * time.Millisecond
|
||||
|
||||
// Kinds of endpoint.
|
||||
const (
|
||||
KindTool = "tool" // a module's own tool
|
||||
KindSeat = "seat" // a seat's verb, served by the module holding it
|
||||
)
|
||||
|
||||
// Endpoint is one tool served on one subject, as it is announced.
|
||||
type Endpoint struct {
|
||||
Kind string
|
||||
Module string // the module whose code answers
|
||||
Tool string // the tool's or the verb's name
|
||||
Seat string // for a seat's verb
|
||||
Scope string // "mesh" or "node", for a seat's verb
|
||||
Node string
|
||||
Description string
|
||||
Schema json.RawMessage
|
||||
Interchangeable bool
|
||||
Subject string
|
||||
Queue string
|
||||
}
|
||||
|
||||
// Name is the endpoint's name as the protocol allows it — letters, digits, `-` and `_` — the
|
||||
// prefix and the tool joined by `__`; the metadata, not the name, is what identifies it.
|
||||
func (e Endpoint) Name() string {
|
||||
prefix := e.Module
|
||||
if e.Kind == KindSeat {
|
||||
prefix = e.Seat
|
||||
}
|
||||
return clean(prefix) + "__" + clean(e.Tool)
|
||||
}
|
||||
|
||||
func clean(s string) string {
|
||||
var b strings.Builder
|
||||
for _, r := range s {
|
||||
if r == '-' || r == '_' || (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') {
|
||||
b.WriteRune(r)
|
||||
} else {
|
||||
b.WriteRune('_')
|
||||
}
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
|
||||
func (e Endpoint) info() micro.EndpointInfo {
|
||||
md := map[string]string{
|
||||
"kind": e.Kind, "module": e.Module, "tool": e.Tool, "node": e.Node,
|
||||
"description": e.Description, "interchangeable": boolWord(e.Interchangeable),
|
||||
}
|
||||
schema := strings.TrimSpace(string(e.Schema))
|
||||
if schema == "" || schema == "null" {
|
||||
schema = "{}"
|
||||
}
|
||||
md["schema"] = schema
|
||||
if e.Kind == KindSeat {
|
||||
md["seat"] = e.Seat
|
||||
md["scope"] = e.Scope
|
||||
}
|
||||
return micro.EndpointInfo{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue, Metadata: md}
|
||||
}
|
||||
|
||||
func boolWord(b bool) string {
|
||||
if b {
|
||||
return "true"
|
||||
}
|
||||
return "false"
|
||||
}
|
||||
|
||||
// Service is who answers: the runtime's name, its instance, and a word about it.
|
||||
type Service struct {
|
||||
Name string
|
||||
ID string
|
||||
Description string
|
||||
Metadata map[string]string
|
||||
}
|
||||
|
||||
func (s Service) identity() micro.ServiceIdentity {
|
||||
md := s.Metadata
|
||||
if md == nil {
|
||||
md = map[string]string{}
|
||||
}
|
||||
return micro.ServiceIdentity{Name: s.Name, ID: s.ID, Version: Version, Metadata: md}
|
||||
}
|
||||
|
||||
// Info is the info_response for these endpoints.
|
||||
func Info(s Service, endpoints []Endpoint) micro.Info {
|
||||
out := micro.Info{ServiceIdentity: s.identity(), Type: micro.InfoResponseType,
|
||||
Description: s.Description, Endpoints: []micro.EndpointInfo{}}
|
||||
for _, e := range endpoints {
|
||||
out.Endpoints = append(out.Endpoints, e.info())
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// Serve answers discovery for one service until stopped, asking `current` for its endpoints each
|
||||
// time — so what is announced is what is served now, re-served memberships included.
|
||||
func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error) {
|
||||
started := time.Now().UTC()
|
||||
answer := func(subject string, _ []byte) []byte {
|
||||
parts := strings.Split(subject, ".")
|
||||
if len(parts) < 2 || parts[0] != "$SRV" {
|
||||
return nil
|
||||
}
|
||||
if len(parts) >= 3 && parts[2] != s.Name {
|
||||
return nil // another service's
|
||||
}
|
||||
if len(parts) >= 4 && parts[3] != s.ID {
|
||||
return nil // another instance's
|
||||
}
|
||||
var v any
|
||||
switch parts[1] {
|
||||
case "PING":
|
||||
v = micro.Ping{ServiceIdentity: s.identity(), Type: micro.PingResponseType}
|
||||
case "INFO":
|
||||
v = Info(s, current())
|
||||
case "STATS":
|
||||
st := micro.Stats{ServiceIdentity: s.identity(), Type: micro.StatsResponseType, Started: started,
|
||||
Endpoints: []*micro.EndpointStats{}}
|
||||
for _, e := range current() {
|
||||
st.Endpoints = append(st.Endpoints, µ.EndpointStats{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue})
|
||||
}
|
||||
v = st
|
||||
default:
|
||||
return nil
|
||||
}
|
||||
body, err := json.Marshal(v)
|
||||
if err != nil {
|
||||
return nil
|
||||
}
|
||||
return body
|
||||
}
|
||||
var stops []func()
|
||||
for _, verb := range []string{"PING", "INFO", "STATS"} {
|
||||
for _, subject := range []string{"$SRV." + verb, "$SRV." + verb + ".>"} {
|
||||
stop, err := conn.Raw(subject, answer)
|
||||
if err != nil {
|
||||
for _, st := range stops {
|
||||
st()
|
||||
}
|
||||
return func() {}, err
|
||||
}
|
||||
stops = append(stops, stop)
|
||||
}
|
||||
}
|
||||
return func() {
|
||||
for _, st := range stops {
|
||||
st()
|
||||
}
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Gather asks every service on the bus what it serves and answers what came back within the window.
|
||||
// An answer that is not an info_response is skipped.
|
||||
func Gather(conn *bus.Conn) ([]micro.Info, error) {
|
||||
raw, err := conn.Gather("$SRV.INFO", nil, Window)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out []micro.Info
|
||||
for _, b := range raw {
|
||||
var i micro.Info
|
||||
if json.Unmarshal(b, &i) != nil || i.Type != micro.InfoResponseType {
|
||||
continue
|
||||
}
|
||||
out = append(out, i)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// Endpoints reads an info_response's endpoints back into what they announce.
|
||||
func Endpoints(i micro.Info) []Endpoint {
|
||||
var out []Endpoint
|
||||
for _, e := range i.Endpoints {
|
||||
md := e.Metadata
|
||||
out = append(out, Endpoint{
|
||||
Kind: md["kind"], Module: md["module"], Tool: md["tool"], Seat: md["seat"], Scope: md["scope"],
|
||||
Node: md["node"], Description: md["description"], Schema: json.RawMessage(md["schema"]),
|
||||
Interchangeable: md["interchangeable"] == "true", Subject: e.Subject, Queue: e.QueueGroup,
|
||||
})
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -512,6 +512,59 @@ func (c *Conn) PublishAs(module string, env Envelope) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// ServedOn is where a served module's tool is answered right now: the membership's subjects when
|
||||
// issued, the derived shape otherwise — what Handle subscribes, for what announces it (ADR 0197).
|
||||
func (c *Conn) ServedOn(module, tool string) []Served { return c.servedOn(module, tool) }
|
||||
|
||||
// Raw answers one subject with a function of the request, not a tool's reply envelope: the NATS
|
||||
// services protocol's discovery subjects answer in their own format (novox/hq ADR 0197). A nil
|
||||
// answer is no reply — the request was for another service.
|
||||
func (c *Conn) Raw(subject string, answer func(subject string, data []byte) []byte) (func(), error) {
|
||||
sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) {
|
||||
if body := answer(msg.Subject, msg.Data); body != nil {
|
||||
_ = msg.Respond(body)
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
return func() {}, err
|
||||
}
|
||||
c.track(sub)
|
||||
return func() { _ = sub.Unsubscribe() }, nil
|
||||
}
|
||||
|
||||
// Gather publishes one request and collects every answer that arrives within the window: a
|
||||
// discovery request every service instance answers (ADR 0197). It never stops early — how many will
|
||||
// answer is what it is finding out.
|
||||
func (c *Conn) Gather(subject string, body []byte, window time.Duration) ([][]byte, error) {
|
||||
inbox := c.nc.NewRespInbox()
|
||||
sub, err := c.nc.SubscribeSync(inbox)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
if err := c.nc.PublishRequest(subject, inbox, body); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out [][]byte
|
||||
deadline := time.Now().Add(window)
|
||||
for {
|
||||
left := time.Until(deadline)
|
||||
if left <= 0 {
|
||||
return out, nil
|
||||
}
|
||||
msg, err := sub.NextMsg(left)
|
||||
if err != nil {
|
||||
if errors.Is(err, nats.ErrTimeout) {
|
||||
return out, nil
|
||||
}
|
||||
return out, err
|
||||
}
|
||||
if len(msg.Data) > 0 {
|
||||
out = append(out, msg.Data)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Flush waits until the bus has every subscription made so far, so what is served is answerable
|
||||
// when this returns.
|
||||
func (c *Conn) Flush() { _ = c.nc.Flush() }
|
||||
|
||||
@@ -19,12 +19,12 @@ package console
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/announce"
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
)
|
||||
|
||||
@@ -159,188 +159,181 @@ func controllerOutput(conn *bus.Conn, verb string) (string, error) {
|
||||
// jsonIn is the JSON document a command printed, after any lines it said first: a seat verb runs
|
||||
// the controller's command, and a command may warn before it answers.
|
||||
func jsonIn(output string) string {
|
||||
if strings.HasPrefix(strings.TrimSpace(output), "{") {
|
||||
return strings.TrimSpace(output)
|
||||
t := strings.TrimSpace(output)
|
||||
if strings.HasPrefix(t, "{") || strings.HasPrefix(t, "[") {
|
||||
return t
|
||||
}
|
||||
if i := strings.Index(output, "\n{"); i >= 0 {
|
||||
return strings.TrimSpace(output[i+1:])
|
||||
for _, open := range []string{"\n{", "\n["} {
|
||||
if i := strings.Index(output, open); i >= 0 {
|
||||
return strings.TrimSpace(output[i+1:])
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// A machine as `node list` prints it: its name, when it was last heard from, its mode — converged
|
||||
// or adopted, the only two the command prints — and its id. Lines the command says around them
|
||||
// (a warning about the bus's users, "no node records yet") are not machines and are skipped.
|
||||
var nodeLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+.*\s(converged|adopted)\s+\S+\s*$`)
|
||||
|
||||
// machinesIn reads the machines from `node list`'s output.
|
||||
func machinesIn(output string) []string {
|
||||
var out []string
|
||||
for _, line := range strings.Split(output, "\n") {
|
||||
if m := nodeLine.FindStringSubmatch(line); m != nil {
|
||||
out = append(out, m[1])
|
||||
}
|
||||
}
|
||||
sort.Strings(out)
|
||||
return out
|
||||
// recordedModule is a module as the controller's records hold it (`module list --json`): where the
|
||||
// mesh assigned it, and whether it declares tools — what should announce itself, and where.
|
||||
type recordedModule struct {
|
||||
Module string `json:"module"`
|
||||
On []string `json:"on"`
|
||||
Tools bool `json:"tools"`
|
||||
}
|
||||
|
||||
// A module as `module list` prints it: name, version, how it was built, and where it runs.
|
||||
var moduleLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+\S+\s+.*?\s+on (.+)$`)
|
||||
|
||||
// assignmentsIn reads, from `module list`'s output, the machines each module runs on.
|
||||
func assignmentsIn(output string) map[string][]string {
|
||||
out := map[string][]string{}
|
||||
for _, line := range strings.Split(output, "\n") {
|
||||
if line == "" || line[0] == ' ' || line[0] == '\t' {
|
||||
continue
|
||||
}
|
||||
m := moduleLine.FindStringSubmatch(strings.TrimRight(line, " "))
|
||||
if m == nil {
|
||||
continue
|
||||
}
|
||||
on := strings.TrimSpace(m[2])
|
||||
if on == "nothing" {
|
||||
out[m[1]] = []string{}
|
||||
continue
|
||||
}
|
||||
var nodes []string
|
||||
for _, n := range strings.Split(on, ",") {
|
||||
if n = strings.TrimSpace(n); n != "" {
|
||||
nodes = append(nodes, n)
|
||||
}
|
||||
}
|
||||
sort.Strings(nodes)
|
||||
out[m[1]] = nodes
|
||||
}
|
||||
return out
|
||||
// recordedMachine is a machine as the controller's records hold it (`node list --json`).
|
||||
type recordedMachine struct {
|
||||
Name string `json:"name"`
|
||||
}
|
||||
|
||||
// interchangeable is whether the mesh issued the module a plain subject for this tool — one any of
|
||||
// its instances answers (ADR 0160): the module's own subject with no machine after it.
|
||||
func interchangeable(module string, t Tool) bool {
|
||||
plain := "mesh.mod." + module + ".tool." + t.Name
|
||||
for _, s := range t.Subjects {
|
||||
if s == plain {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// indexOn asks the mesh what it holds: the flat listing (catalogue, modules, the seats' tools),
|
||||
// and from the controller the seats' holders, the machines and the assignments — at once.
|
||||
// indexOn asks the mesh what it holds (novox/hq ADR 0197): what answers, from every runtime's own
|
||||
// announcement on the bus — one `$SRV.INFO` request — and what should, from the controller's records
|
||||
// read as JSON. Nothing is inferred from a roster and nothing is parsed from print.
|
||||
func indexOn(conn *bus.Conn) (*index, error) {
|
||||
var wg sync.WaitGroup
|
||||
var seatsOut, nodesOut, modulesOut string
|
||||
wg.Add(3)
|
||||
go func() { defer wg.Done(); seatsOut, _ = controllerOutput(conn, "seats") }()
|
||||
var nodesOut, modulesOut string
|
||||
wg.Add(2)
|
||||
go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }()
|
||||
go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }()
|
||||
l, err := toolsOn(conn)
|
||||
infos, err := announce.Gather(conn)
|
||||
wg.Wait()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
x := &index{Modules: map[string]*moduleInfo{}, NotAnswering: l.NotAnswering, Listing: l}
|
||||
|
||||
// The seats: their verbs from the mesh's records, their holders from the controller.
|
||||
var held struct {
|
||||
Seats []struct {
|
||||
Seat string `json:"seat"`
|
||||
Scope string `json:"scope"`
|
||||
Holders []holder `json:"holders"`
|
||||
} `json:"seats"`
|
||||
}
|
||||
_ = json.Unmarshal([]byte(jsonIn(seatsOut)), &held)
|
||||
holders := map[string][]holder{}
|
||||
scopes := map[string]string{}
|
||||
for _, s := range held.Seats {
|
||||
holders[s.Seat] = s.Holders
|
||||
scopes[s.Seat] = s.Scope
|
||||
}
|
||||
bySeat := map[string]*seatInfo{}
|
||||
for _, t := range l.Tools {
|
||||
if !t.Seat {
|
||||
continue
|
||||
}
|
||||
s := bySeat[t.Module]
|
||||
if s == nil {
|
||||
scope := t.Scope
|
||||
if scope == "" {
|
||||
scope = scopes[t.Module]
|
||||
l := &Listing{Tools: []Tool{}, NotAnswering: []string{}}
|
||||
x := &index{Modules: map[string]*moduleInfo{}, Listing: l}
|
||||
announced := map[string]map[string]bool{} // module → node → announced something
|
||||
seats := map[string]*seatInfo{}
|
||||
toolAt := map[string]int{} // <module>.<tool> or <seat>.<verb> → index in l.Tools
|
||||
machines := map[string]bool{}
|
||||
for _, info := range infos {
|
||||
for _, e := range announce.Endpoints(info) {
|
||||
if e.Node != "" {
|
||||
machines[e.Node] = true
|
||||
}
|
||||
if scope != "node" {
|
||||
scope = "mesh"
|
||||
if announced[e.Module] == nil {
|
||||
announced[e.Module] = map[string]bool{}
|
||||
}
|
||||
s = &seatInfo{Seat: t.Module, Scope: scope, Holders: holders[t.Module]}
|
||||
bySeat[t.Module] = s
|
||||
}
|
||||
s.Verbs = append(s.Verbs, t)
|
||||
}
|
||||
for _, s := range bySeat {
|
||||
x.Seats = append(x.Seats, *s)
|
||||
}
|
||||
sort.Slice(x.Seats, func(i, j int) bool { return x.Seats[i].Seat < x.Seats[j].Seat })
|
||||
|
||||
// The machines: what the controller knows, and any that hold a seat.
|
||||
seen := map[string]bool{}
|
||||
for _, n := range machinesIn(nodesOut) {
|
||||
seen[n] = true
|
||||
}
|
||||
for _, s := range x.Seats {
|
||||
for _, h := range s.Holders {
|
||||
if h.Node != "" {
|
||||
seen[h.Node] = true
|
||||
announced[e.Module][e.Node] = true
|
||||
switch e.Kind {
|
||||
case announce.KindSeat:
|
||||
st := seats[e.Seat]
|
||||
if st == nil {
|
||||
st = &seatInfo{Seat: e.Seat, Scope: e.Scope}
|
||||
seats[e.Seat] = st
|
||||
}
|
||||
if st.Scope != "node" && e.Scope == "node" {
|
||||
st.Scope = "node"
|
||||
}
|
||||
h := holder{Module: e.Module, Node: e.Node}
|
||||
if !containsHolder(st.Holders, h) {
|
||||
st.Holders = append(st.Holders, h)
|
||||
}
|
||||
key := e.Seat + "." + e.Tool
|
||||
if _, have := toolAt[key]; !have {
|
||||
toolAt[key] = len(l.Tools)
|
||||
t := Tool{Module: e.Seat, Name: e.Tool, Description: e.Description, Input: e.Schema, Seat: true, Scope: st.Scope}
|
||||
l.Tools = append(l.Tools, t)
|
||||
st.Verbs = append(st.Verbs, t)
|
||||
}
|
||||
case announce.KindTool:
|
||||
m := x.Modules[e.Module]
|
||||
if m == nil {
|
||||
m = &moduleInfo{Module: e.Module}
|
||||
x.Modules[e.Module] = m
|
||||
}
|
||||
if e.Node != "" && !contains(m.On, e.Node) {
|
||||
m.On = append(m.On, e.Node)
|
||||
}
|
||||
m.Interchangeable = m.Interchangeable || e.Interchangeable
|
||||
key := e.Module + "." + e.Tool
|
||||
i, have := toolAt[key]
|
||||
if !have {
|
||||
i = len(l.Tools)
|
||||
toolAt[key] = i
|
||||
l.Tools = append(l.Tools, Tool{Module: e.Module, Name: e.Tool, Description: e.Description, Input: e.Schema})
|
||||
}
|
||||
if !contains(l.Tools[i].Subjects, e.Subject) {
|
||||
l.Tools[i].Subjects = append(l.Tools[i].Subjects, e.Subject)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// The modules that answer tools, where they run, and whether any instance will do.
|
||||
on := assignmentsIn(modulesOut)
|
||||
for _, t := range l.Tools {
|
||||
// A tool's subjects as a call looks them up: the plain one any instance answers first, then each
|
||||
// machine's.
|
||||
for i := range l.Tools {
|
||||
t := &l.Tools[i]
|
||||
if t.Seat {
|
||||
continue
|
||||
}
|
||||
m := x.Modules[t.Module]
|
||||
if m == nil {
|
||||
m = &moduleInfo{Module: t.Module, On: on[t.Module]}
|
||||
x.Modules[t.Module] = m
|
||||
}
|
||||
m.Tools = append(m.Tools, t)
|
||||
if interchangeable(t.Module, t) {
|
||||
m.Interchangeable = true
|
||||
plain := "mesh.mod." + t.Module + ".tool." + t.Name
|
||||
sort.SliceStable(t.Subjects, func(a, b int) bool {
|
||||
if (t.Subjects[a] == plain) != (t.Subjects[b] == plain) {
|
||||
return t.Subjects[a] == plain
|
||||
}
|
||||
return t.Subjects[a] < t.Subjects[b]
|
||||
})
|
||||
}
|
||||
for _, t := range l.Tools {
|
||||
if !t.Seat {
|
||||
x.Modules[t.Module].Tools = append(x.Modules[t.Module].Tools, t)
|
||||
}
|
||||
}
|
||||
for _, m := range x.Modules {
|
||||
// A module the controller could not place is placed where its own answer says it runs.
|
||||
if len(m.On) == 0 {
|
||||
nodes := map[string]bool{}
|
||||
for _, t := range m.Tools {
|
||||
for _, s := range t.Subjects {
|
||||
base := "mesh.mod." + m.Module + ".tool." + t.Name + "."
|
||||
if strings.HasPrefix(s, base) {
|
||||
nodes[strings.TrimPrefix(s, base)] = true
|
||||
}
|
||||
sort.Strings(m.On)
|
||||
}
|
||||
for _, st := range seats {
|
||||
sort.Slice(st.Holders, func(a, b int) bool { return st.Holders[a].Node < st.Holders[b].Node })
|
||||
x.Seats = append(x.Seats, *st)
|
||||
}
|
||||
sort.Slice(x.Seats, func(i, j int) bool { return x.Seats[i].Seat < x.Seats[j].Seat })
|
||||
sort.SliceStable(l.Tools, func(i, j int) bool {
|
||||
return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name
|
||||
})
|
||||
|
||||
// What should have answered: every assignment of a module that declares tools. Silence is named;
|
||||
// a module with no tools is never a name here.
|
||||
var recorded []recordedModule
|
||||
if json.Unmarshal([]byte(jsonIn(modulesOut)), &recorded) == nil {
|
||||
for _, m := range recorded {
|
||||
if !m.Tools {
|
||||
continue
|
||||
}
|
||||
for _, n := range m.On {
|
||||
if !announced[m.Module][n] {
|
||||
l.NotAnswering = append(l.NotAnswering, m.Module+" on "+n)
|
||||
}
|
||||
}
|
||||
for n := range nodes {
|
||||
m.On = append(m.On, n)
|
||||
}
|
||||
sort.Strings(m.On)
|
||||
}
|
||||
for _, n := range m.On {
|
||||
seen[n] = true
|
||||
} else {
|
||||
l.NotAnswering = append(l.NotAnswering, "mesh-controller (its records of the modules did not answer, so what is missing cannot be said)")
|
||||
}
|
||||
sort.Strings(l.NotAnswering)
|
||||
x.NotAnswering = l.NotAnswering
|
||||
|
||||
var known []recordedMachine
|
||||
if json.Unmarshal([]byte(jsonIn(nodesOut)), &known) == nil {
|
||||
for _, n := range known {
|
||||
if n.Name != "" {
|
||||
machines[n.Name] = true
|
||||
}
|
||||
}
|
||||
}
|
||||
for n := range seen {
|
||||
for n := range machines {
|
||||
x.Machines = append(x.Machines, n)
|
||||
}
|
||||
sort.Strings(x.Machines)
|
||||
return x, nil
|
||||
}
|
||||
|
||||
func containsHolder(hs []holder, h holder) bool {
|
||||
for _, x := range hs {
|
||||
if x == h {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func (s *Surface) index() (*index, error) {
|
||||
s.mu.Lock()
|
||||
if s.idx != nil && time.Since(s.idxAt) <= IndexKept {
|
||||
@@ -520,7 +513,7 @@ func failure(text string) map[string]any {
|
||||
func (s *Surface) discover(name string, args map[string]any) map[string]any {
|
||||
x, err := s.index()
|
||||
if err != nil {
|
||||
return failure(whyItFailed(catalogueModules, err))
|
||||
return failure("the mesh's discovery failed: " + err.Error())
|
||||
}
|
||||
str := func(k string) string { v, _ := args[k].(string); return strings.TrimSpace(v) }
|
||||
|
||||
|
||||
@@ -4,9 +4,11 @@ import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/announce"
|
||||
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
|
||||
"github.com/novox/mesh-tools/node-tools/internal/runtime"
|
||||
)
|
||||
@@ -56,14 +58,20 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
|
||||
}
|
||||
defer stop()
|
||||
|
||||
catalogue := connect(t, "mesh-catalog", "")
|
||||
stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) {
|
||||
return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "gamma"}}}, nil
|
||||
})
|
||||
defer stopCat()
|
||||
// Nothing may be asked for a roster or a module's `tools` any more (ADR 0197): counted.
|
||||
var asked atomic.Int32
|
||||
watcher := connect(t, "watcher", "")
|
||||
for _, subject := range []string{"mesh.mod.*.tool.tools", "mesh.mod.*.tool.tools.*", "mesh.mod.mesh-catalog.>"} {
|
||||
stop, err := watcher.Raw(subject, func(string, []byte) []byte { asked.Add(1); return nil })
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
}
|
||||
|
||||
// The controller, answering as its seat verbs do: the command's printed output.
|
||||
controller := connect(t, "mesh-controller", "")
|
||||
// The controller: its records as JSON, as its seat verbs answer them, and its own seat's verbs
|
||||
// announced on the bus like every runtime's.
|
||||
controller := connect(t, "mesh-controller", "bench")
|
||||
out := func(s string) map[string]any { return map[string]any{"output": s, "ok": true} }
|
||||
serve := func(verb string, answer func() any) {
|
||||
stop, err := controller.HandleSubject("mesh.seat.mesh-controller.tool."+verb, func(json.RawMessage) (any, error) {
|
||||
@@ -74,34 +82,28 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
}
|
||||
serve("tools", func() any {
|
||||
return map[string]any{"seats": []map[string]any{
|
||||
{"seat": "mesh-controller", "scope": "mesh", "tools": []map[string]any{
|
||||
{"name": "nodes", "description": "Every machine the mesh knows.", "input": map[string]any{}}}},
|
||||
{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{
|
||||
{"name": "list", "description": "what is on the shelf", "input": map[string]any{}},
|
||||
{"name": "clear", "description": "take it all off", "input": map[string]any{}}}},
|
||||
}}
|
||||
})
|
||||
serve("nodes", func() any {
|
||||
return out("bench 3m ago converged 1f2e\ndesk here converged 9a8b\n")
|
||||
return out("the bus's user list leaves out 2 user(s)\n" +
|
||||
`[{"name":"bench","heard":"3m ago","mode":"converged","id":"1f2e"},{"name":"desk","heard":"here","mode":"converged","id":"9a8b"}]`)
|
||||
})
|
||||
serve("seats", func() any {
|
||||
return out("a warning the command printed first\n" + `{
|
||||
"seats": [
|
||||
{"seat": "mesh-controller", "scope": "mesh", "decision": "x", "holders": [{"module": "mesh-controller", "node": "bench"}]},
|
||||
{"seat": "node-shelf", "scope": "node", "decision": "y", "holders": [{"module": "beta", "node": "desk"}]}
|
||||
]
|
||||
}`)
|
||||
})
|
||||
gammaOn := "nothing"
|
||||
var gammaOn atomic.Value
|
||||
gammaOn.Store("[]")
|
||||
serve("modules", func() any {
|
||||
return out(fmt.Sprintf("alpha 1 built 1a2b3c4d on desk\n needs container-runtime\n"+
|
||||
"beta 1 built 1a2b3c4d on desk\n"+
|
||||
"gamma 1 built 1a2b3c4d on %s\n", gammaOn))
|
||||
return out(fmt.Sprintf(`[{"module":"alpha","on":["desk"],"tools":true},{"module":"beta","on":["desk"],"tools":true},`+
|
||||
`{"module":"gamma","on":%s,"tools":true},{"module":"delta","on":["desk"],"tools":false},`+
|
||||
`{"module":"epsilon","on":["bench"],"tools":true}]`, gammaOn.Load().(string)))
|
||||
})
|
||||
catalogue.Flush()
|
||||
stopAnn, err := announce.Serve(controller, announce.Service{Name: "mesh-controller", ID: "bench"}, func() []announce.Endpoint {
|
||||
return []announce.Endpoint{{Kind: announce.KindSeat, Module: "mesh-controller", Seat: "mesh-controller", Scope: "mesh",
|
||||
Tool: "nodes", Node: "bench", Description: "Every machine the mesh knows.", Schema: json.RawMessage(`{}`),
|
||||
Subject: "mesh.seat.mesh-controller.tool.nodes"}}
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stopAnn)
|
||||
controller.Flush()
|
||||
watcher.Flush()
|
||||
|
||||
up, err := Serve(NewSurface(nodeTools, "desk.node-tools"), "127.0.0.1:0")
|
||||
if err != nil {
|
||||
@@ -132,6 +134,11 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
|
||||
!strings.Contains(overview, `"bench"`) || !strings.Contains(overview, `"desk"`) {
|
||||
t.Errorf("overview: %s", overview)
|
||||
}
|
||||
// Silence is named only where tools should have answered: epsilon declares tools on bench and
|
||||
// nothing there announced it; delta declares none and is never a name.
|
||||
if !strings.Contains(overview, "epsilon on bench") || strings.Contains(overview, "delta") {
|
||||
t.Errorf("not answering: %s", overview)
|
||||
}
|
||||
machine, isErr := call(t, endpoint, "mesh_machine", map[string]any{"node": "desk"})
|
||||
if isErr || !strings.Contains(machine, "desk/node-shelf.list") || !strings.Contains(machine, "desk/beta.three") ||
|
||||
!strings.Contains(machine, "desk/alpha.one") {
|
||||
@@ -181,7 +188,7 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
|
||||
t.Fatalf("gamma was found before it served: %s", got)
|
||||
}
|
||||
mesh.Issue(t, mt.MembershipOf("gamma", "desk", false, nil))
|
||||
gammaOn = "desk"
|
||||
gammaOn.Store(`["desk"]`)
|
||||
late := connect(t, "node-tools", "desk")
|
||||
stopLate, err := runtime.Run(late, []runtime.Served{{Module: "gamma", Entrypoints: []string{mt.Fixture("env-gamma.serve.mjs")}}},
|
||||
nil, (&mt.Logs{}).Logf)
|
||||
@@ -201,22 +208,21 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
|
||||
t.Errorf("a module that arrived later was not found: %s", found)
|
||||
}
|
||||
|
||||
if n := asked.Load(); n != 0 {
|
||||
t.Errorf("discovery asked a roster or a module's tools %d time(s); it asks the bus", n)
|
||||
}
|
||||
|
||||
// The old names still answer, unannounced.
|
||||
if got, isErr := call(t, endpoint, "alpha.one", nil); isErr || !strings.Contains(got, `"alpha": 1`) {
|
||||
t.Errorf("an old name: %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheControllersPrintedListsAreRead(t *testing.T) {
|
||||
machines := machinesIn("the bus's user list leaves out 2 user(s)\nace 2m ago converged 0c1d\n" +
|
||||
"g14 here converged 77aa\nnovox 5s ago adopted 3e4f\n")
|
||||
if got := strings.Join(machines, ","); got != "ace,g14,novox" {
|
||||
t.Errorf("machines: %s", got)
|
||||
func TestTheControllersRecordsAreReadAsJSON(t *testing.T) {
|
||||
if got := jsonIn("a warning printed first\n[{\"name\":\"ace\"}]"); got != `[{"name":"ace"}]` {
|
||||
t.Errorf("an array after a warning: %q", got)
|
||||
}
|
||||
on := assignmentsIn("baserow 1 built 7c800705 on ace\n requires postgres-database\n" +
|
||||
"confluence 1 built 7c800705 on nothing\n" +
|
||||
"mesh-wireguard 1 with the control plane on ace, g14, novox, shanks\n")
|
||||
if strings.Join(on["baserow"], ",") != "ace" || len(on["confluence"]) != 0 || strings.Join(on["mesh-wireguard"], ",") != "ace,g14,novox,shanks" {
|
||||
t.Errorf("assignments: %v", on)
|
||||
if got := jsonIn(`{"seats":[]}`); got != `{"seats":[]}` {
|
||||
t.Errorf("an object: %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,19 +5,11 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
)
|
||||
|
||||
const (
|
||||
catalogueModules = "mesh-catalog.catalog_modules"
|
||||
seatTools = "seat:mesh-controller.tools"
|
||||
toolsVerb = "tools"
|
||||
)
|
||||
|
||||
// Tool is a tool as its module — or, for a role's tool, the mesh's records — describes it.
|
||||
type Tool struct {
|
||||
Module string
|
||||
@@ -67,107 +59,14 @@ func toolKey(name string, seats Seats) string {
|
||||
return name
|
||||
}
|
||||
|
||||
// toolsOn asks the mesh what tools it has: the catalogue which modules it holds, each module what it
|
||||
// serves, the controller's seat every role's tools — at once, so a restarting control plane hides
|
||||
// nothing else.
|
||||
// toolsOn is the flat catalogue (MESH_CONSOLE_FLAT=1): what announced itself on the bus, every
|
||||
// tool and seat verb, and the assignments with tools that did not (novox/hq ADR 0197).
|
||||
func toolsOn(conn *bus.Conn) (*Listing, error) {
|
||||
type rolesAnswer struct {
|
||||
Seats []struct {
|
||||
Seat string `json:"seat"`
|
||||
Scope string `json:"scope"`
|
||||
Tools []struct {
|
||||
Name string `json:"name"`
|
||||
Description string `json:"description"`
|
||||
Input json.RawMessage `json:"input"`
|
||||
} `json:"tools"`
|
||||
} `json:"seats"`
|
||||
}
|
||||
var wg sync.WaitGroup
|
||||
var roles *rolesAnswer
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
if got, err := conn.Ask(seatTools, map[string]any{}, ""); err == nil {
|
||||
var r rolesAnswer
|
||||
if json.Unmarshal(got.Result, &r) == nil {
|
||||
roles = &r
|
||||
}
|
||||
}
|
||||
}()
|
||||
answered, err := conn.Ask(catalogueModules, map[string]any{}, "")
|
||||
x, err := indexOn(conn)
|
||||
if err != nil {
|
||||
wg.Wait()
|
||||
return nil, err
|
||||
}
|
||||
var held struct {
|
||||
Modules []struct {
|
||||
Module string `json:"module"`
|
||||
} `json:"modules"`
|
||||
}
|
||||
_ = json.Unmarshal(answered.Result, &held)
|
||||
names := make([]string, 0, len(held.Modules))
|
||||
for _, m := range held.Modules {
|
||||
if m.Module != "" {
|
||||
names = append(names, m.Module)
|
||||
}
|
||||
}
|
||||
type outcome struct {
|
||||
ok bool
|
||||
answer struct {
|
||||
Tools *[]struct {
|
||||
Name string `json:"name"`
|
||||
Description string `json:"description"`
|
||||
Input json.RawMessage `json:"input"`
|
||||
Subjects []string `json:"subjects"`
|
||||
} `json:"tools"`
|
||||
Failed *string `json:"failed"`
|
||||
}
|
||||
}
|
||||
outcomes := make([]outcome, len(names))
|
||||
for i, module := range names {
|
||||
wg.Add(1)
|
||||
go func(i int, module string) {
|
||||
defer wg.Done()
|
||||
got, err := conn.Ask(module+"."+toolsVerb, map[string]any{}, "")
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
if json.Unmarshal(got.Result, &outcomes[i].answer) == nil {
|
||||
outcomes[i].ok = true
|
||||
}
|
||||
}(i, module)
|
||||
}
|
||||
wg.Wait()
|
||||
l := &Listing{Tools: []Tool{}, NotAnswering: []string{}}
|
||||
if roles != nil {
|
||||
for _, s := range roles.Seats {
|
||||
for _, t := range s.Tools {
|
||||
l.Tools = append(l.Tools, Tool{Module: s.Seat, Name: t.Name, Description: t.Description,
|
||||
Input: t.Input, Seat: true, Scope: s.Scope})
|
||||
}
|
||||
}
|
||||
} else {
|
||||
l.NotAnswering = append(l.NotAnswering, "mesh-controller (seat)")
|
||||
}
|
||||
for i, module := range names {
|
||||
o := outcomes[i]
|
||||
switch {
|
||||
case o.ok && o.answer.Failed != nil:
|
||||
l.NotAnswering = append(l.NotAnswering, fmt.Sprintf("%s (its tools bundle failed to load: %s)", module, *o.answer.Failed))
|
||||
case o.ok && o.answer.Tools != nil:
|
||||
for _, t := range *o.answer.Tools {
|
||||
l.Tools = append(l.Tools, Tool{Module: module, Name: t.Name, Description: t.Description,
|
||||
Input: t.Input, Subjects: t.Subjects})
|
||||
}
|
||||
default:
|
||||
l.NotAnswering = append(l.NotAnswering, module)
|
||||
}
|
||||
}
|
||||
sort.SliceStable(l.Tools, func(i, j int) bool {
|
||||
return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name
|
||||
})
|
||||
sort.Strings(l.NotAnswering)
|
||||
return l, nil
|
||||
return x.Listing, nil
|
||||
}
|
||||
|
||||
// callTool calls `<module>.<tool>[@<node>]`, on the subject the listing names for it when it names one.
|
||||
|
||||
@@ -55,20 +55,13 @@ func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) {
|
||||
}
|
||||
defer stop()
|
||||
|
||||
catalogue := connect(t, "mesh-catalog", "")
|
||||
stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) {
|
||||
return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "ghost"}}}, nil
|
||||
})
|
||||
defer stopCat()
|
||||
// The controller's records (ADR 0197): ghost declares tools on desk and nothing announces it.
|
||||
controller := connect(t, "mesh-controller", "")
|
||||
stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.tools", func(json.RawMessage) (any, error) {
|
||||
return map[string]any{"seats": []map[string]any{{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{
|
||||
{"name": "list", "description": "what is on the shelf", "input": map[string]any{}},
|
||||
{"name": "clear", "description": "take it all off", "input": map[string]any{}},
|
||||
}}}}, nil
|
||||
stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.modules", func(json.RawMessage) (any, error) {
|
||||
return map[string]any{"ok": true, "output": `[{"module":"alpha","on":["desk"],"tools":true},` +
|
||||
`{"module":"beta","on":["desk"],"tools":true},{"module":"ghost","on":["desk"],"tools":true}]`}, nil
|
||||
})
|
||||
defer stopSeat()
|
||||
catalogue.Flush()
|
||||
controller.Flush()
|
||||
|
||||
flat := NewSurface(nodeTools, "desk.node-tools")
|
||||
@@ -92,7 +85,7 @@ func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) {
|
||||
if got := strings.Join(names, ","); got != "alpha.one,alpha.two,beta.five,beta.four,beta.three,node-shelf.clear,node-shelf.list" {
|
||||
t.Errorf("listed %s", got)
|
||||
}
|
||||
if got := listed["_meta"].(map[string]any)["notAnswering"]; len(got.([]any)) != 1 || got.([]any)[0] != "ghost" {
|
||||
if got := listed["_meta"].(map[string]any)["notAnswering"]; len(got.([]any)) != 1 || got.([]any)[0] != "ghost on desk" {
|
||||
t.Errorf("not answering: %v", got)
|
||||
}
|
||||
for _, x := range listed["tools"].([]any) {
|
||||
|
||||
@@ -122,7 +122,7 @@ func (s *Surface) Handle(r Request) *Reply {
|
||||
}
|
||||
l, err := s.listing()
|
||||
if err != nil {
|
||||
return refuse(r.ID, -32603, whyItFailed(catalogueModules, err))
|
||||
return refuse(r.ID, -32603, "the mesh's discovery failed: "+err.Error())
|
||||
}
|
||||
tools := make([]map[string]any, 0, len(l.Tools))
|
||||
for _, t := range l.Tools {
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
// Package meshtest raises what the controller would, for tests against a real bus: the ASSIGNMENTS
|
||||
// and EVENTS streams, memberships issued by hand, and a fixture's path.
|
||||
//
|
||||
// **The packages share one bus, so run them one at a time: `go test -p 1 ./...`.** Each test raises
|
||||
// the streams afresh, and the console discovers every runtime that announces itself on the bus
|
||||
// (novox/hq ADR 0197) — a runtime from another package's test is, correctly, found.
|
||||
package meshtest
|
||||
|
||||
import (
|
||||
|
||||
@@ -0,0 +1,113 @@
|
||||
package runtime
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"sort"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
"github.com/nats-io/nats.go/micro"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
|
||||
)
|
||||
|
||||
// novox/hq ADR 0197: what a runtime serves, it announces — the NATS services protocol's discovery,
|
||||
// read here with NATS's own types, each endpoint a subject actually served — and the announcement
|
||||
// follows a re-issued membership.
|
||||
func TestTheRuntimeAnnouncesWhatItServesInTheServicesProtocol(t *testing.T) {
|
||||
mesh := mt.New(t)
|
||||
mesh.Issue(t, mt.MembershipOf("alpha", "anchor", false, nil))
|
||||
mesh.Issue(t, mt.MembershipOf("beta", "anchor", false, map[string][]string{"node-shelf": {"list", "clear"}}))
|
||||
nodeTools := connect(t, "node-tools", "anchor")
|
||||
stop, err := Run(nodeTools, []Served{
|
||||
{"alpha", []string{mt.Fixture("many-alpha.serve.mjs")}},
|
||||
{"beta", []string{mt.Fixture("many-beta.serve.mjs")}},
|
||||
}, nil, (&mt.Logs{}).Logf)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer stop()
|
||||
|
||||
nc, err := nats.Connect(mt.URL(t))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer nc.Close()
|
||||
info := func() micro.Info {
|
||||
t.Helper()
|
||||
msg, err := nc.Request("$SRV.INFO.node-tools.anchor", nil, 2*time.Second)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var i micro.Info
|
||||
if err := json.Unmarshal(msg.Data, &i); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return i
|
||||
}
|
||||
subjects := func(i micro.Info) string {
|
||||
var out []string
|
||||
for _, e := range i.Endpoints {
|
||||
out = append(out, e.Subject+"|"+e.QueueGroup+"|"+e.Metadata["kind"]+"|"+e.Metadata["module"]+"|"+e.Metadata["seat"]+"|"+e.Metadata["scope"]+"|"+e.Metadata["interchangeable"])
|
||||
}
|
||||
sort.Strings(out)
|
||||
return strings.Join(out, "\n")
|
||||
}
|
||||
|
||||
got := info()
|
||||
if got.Type != micro.InfoResponseType || got.Name != "node-tools" || got.ID != "anchor" || got.Version == "" {
|
||||
t.Errorf("identity: %+v", got.ServiceIdentity)
|
||||
}
|
||||
want := strings.Join([]string{
|
||||
"mesh.mod.alpha.tool.one.anchor||tool|alpha|||false",
|
||||
"mesh.mod.alpha.tool.two.anchor||tool|alpha|||false",
|
||||
"mesh.mod.beta.tool.five.anchor||tool|beta|||false",
|
||||
"mesh.mod.beta.tool.four.anchor||tool|beta|||false",
|
||||
"mesh.mod.beta.tool.three.anchor||tool|beta|||false",
|
||||
"mesh.seat.node-shelf.tool.clear.anchor||seat|beta|node-shelf|node|false",
|
||||
"mesh.seat.node-shelf.tool.list.anchor||seat|beta|node-shelf|node|false",
|
||||
}, "\n")
|
||||
if s := subjects(got); s != want {
|
||||
t.Errorf("announced:\n%s\nwant:\n%s", s, want)
|
||||
}
|
||||
for _, e := range got.Endpoints {
|
||||
if e.Metadata["node"] != "anchor" || e.Metadata["description"] == "" || !json.Valid([]byte(e.Metadata["schema"])) {
|
||||
t.Errorf("endpoint metadata: %+v", e)
|
||||
}
|
||||
}
|
||||
|
||||
// Every subject announced is answered.
|
||||
asker := connect(t, "console", "workstation")
|
||||
for _, e := range got.Endpoints {
|
||||
if _, err := asker.Ask("", map[string]any{}, e.Subject); err != nil {
|
||||
t.Errorf("announced %s and does not answer it: %v", e.Subject, err)
|
||||
}
|
||||
}
|
||||
|
||||
// PING answers with the same identity, and a request for another service is not answered.
|
||||
if msg, err := nc.Request("$SRV.PING", nil, 2*time.Second); err != nil || !strings.Contains(string(msg.Data), micro.PingResponseType) {
|
||||
t.Errorf("ping: %v %v", msg, err)
|
||||
}
|
||||
if _, err := nc.Request("$SRV.INFO.somebody-else", nil, 300*time.Millisecond); err == nil {
|
||||
t.Error("answered a request for another service")
|
||||
}
|
||||
|
||||
// alpha is re-issued a plain subject: the announcement says so, and says it is interchangeable.
|
||||
mesh.Issue(t, mt.MembershipOf("alpha", "anchor", true, nil))
|
||||
var after string
|
||||
for i := 0; i < 40; i++ {
|
||||
after = subjects(info())
|
||||
if strings.Contains(after, "mesh.mod.alpha.tool.one|serve.alpha|tool|alpha|||true") {
|
||||
break
|
||||
}
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
}
|
||||
if !strings.Contains(after, "mesh.mod.alpha.tool.one|serve.alpha|tool|alpha|||true") ||
|
||||
!strings.Contains(after, "mesh.mod.alpha.tool.one.anchor||tool|alpha|||true") {
|
||||
t.Errorf("after a re-issued membership:\n%s", after)
|
||||
}
|
||||
_ = bus.Served{}
|
||||
}
|
||||
@@ -14,6 +14,7 @@ import (
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/announce"
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
"github.com/novox/mesh-tools/node-tools/internal/launch"
|
||||
)
|
||||
@@ -276,10 +277,83 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
|
||||
}
|
||||
logf("%s", line)
|
||||
stops = append(stops, serveSeats(conn, modules, registrations, logf))
|
||||
|
||||
// **What it serves, it announces** (novox/hq ADR 0197): the NATS services protocol's discovery,
|
||||
// answered with what is served at the moment it is asked — re-served memberships included.
|
||||
announced, err := announce.Serve(conn, announce.Service{
|
||||
Name: conn.Module(), ID: instanceOf(conn),
|
||||
Description: "the mesh's tool runtime on " + node + ": every assigned module's tools and the seats they hold",
|
||||
Metadata: map[string]string{"node": node},
|
||||
}, func() []announce.Endpoint { return endpointsOf(conn, own, registrations, modules) })
|
||||
if err != nil {
|
||||
logf("[mesh-tools] cannot announce what it serves: %v", err)
|
||||
} else {
|
||||
stops = append(stops, announced)
|
||||
}
|
||||
conn.Flush()
|
||||
return stopAll, nil
|
||||
}
|
||||
|
||||
// instanceOf is this runtime's instance on the bus: its machine, which is what tells two instances of
|
||||
// one service apart; the connection's module where it has no machine.
|
||||
func instanceOf(conn *bus.Conn) string {
|
||||
if n := conn.Node(); n != "" {
|
||||
return n
|
||||
}
|
||||
return conn.Module()
|
||||
}
|
||||
|
||||
// endpointsOf is everything this runtime serves now: each served module's tools on every subject the
|
||||
// mesh issued for them, and each held seat's verbs on the seat's subject (ADR 0197).
|
||||
func endpointsOf(conn *bus.Conn, own []registration, registrations []registration, modules []string) []announce.Endpoint {
|
||||
node := conn.Node()
|
||||
var out []announce.Endpoint
|
||||
for _, r := range own {
|
||||
for _, t := range r.tools {
|
||||
served := conn.ServedOn(r.module, t.Name)
|
||||
interchangeable := false
|
||||
for _, s := range served {
|
||||
interchangeable = interchangeable || s.Subject == "mesh.mod."+r.module+".tool."+t.Name
|
||||
}
|
||||
for _, s := range served {
|
||||
out = append(out, announce.Endpoint{Kind: announce.KindTool, Module: r.module, Tool: t.Name,
|
||||
Node: node, Description: t.Description, Schema: t.Input, Interchangeable: interchangeable,
|
||||
Subject: s.Subject, Queue: s.Queue})
|
||||
}
|
||||
}
|
||||
}
|
||||
impl := map[string]map[string]launch.Tool{}
|
||||
for _, r := range registrations {
|
||||
if impl[r.module] == nil {
|
||||
impl[r.module] = map[string]launch.Tool{}
|
||||
}
|
||||
for _, t := range r.tools {
|
||||
impl[r.module][t.Name] = t
|
||||
}
|
||||
}
|
||||
have := map[string]bool{}
|
||||
for _, module := range modules {
|
||||
m := conn.Membership(module)
|
||||
if m == nil {
|
||||
continue
|
||||
}
|
||||
for _, v := range m.Seats {
|
||||
t, ok := impl[v.Seat][v.Verb]
|
||||
if !ok || have[v.Subject] {
|
||||
continue
|
||||
}
|
||||
have[v.Subject] = true
|
||||
scope := "mesh"
|
||||
if node != "" && strings.HasSuffix(v.Subject, "."+node) {
|
||||
scope = "node"
|
||||
}
|
||||
out = append(out, announce.Endpoint{Kind: announce.KindSeat, Module: module, Tool: v.Verb, Seat: v.Seat,
|
||||
Scope: scope, Node: node, Description: t.Description, Schema: t.Input, Subject: v.Subject})
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func orNone(names []string) string {
|
||||
if len(names) == 0 {
|
||||
return "(none)"
|
||||
|
||||
@@ -13,7 +13,7 @@
|
||||
"test": "node --test --test-concurrency=1 --experimental-strip-types 'test/*.test.ts'"
|
||||
},
|
||||
"dependencies": {
|
||||
"@novox/mesh-sdk": "^0.1.0",
|
||||
"@novox/mesh-sdk": "^0.1.5",
|
||||
"nats": "^2.29.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
// What a runtime serves, it announces (novox/hq ADR 0197): the NATS services protocol's discovery —
|
||||
// `$SRV.PING`, `$SRV.INFO`, `$SRV.STATS`, and the same followed by the service's name and its id —
|
||||
// answered in the io.nats.micro.v1 format with what is served at the moment of the request. Serving is
|
||||
// unchanged; this only says what is served. One service per runtime process, because the bus admits
|
||||
// one reply per request from each responder: one endpoint per tool per subject, its metadata saying
|
||||
// which module, seat, scope and machine it is. The same shape the Go runtime answers.
|
||||
|
||||
import { StringCodec } from "nats";
|
||||
import { asSchema } from "@novox/mesh-sdk/stdio";
|
||||
import type { ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import type { RuntimeBroker } from "./broker-nats.js";
|
||||
|
||||
const sc = StringCodec();
|
||||
|
||||
export const VERSION = "0.1.0";
|
||||
export const INFO_RESPONSE = "io.nats.micro.v1.info_response";
|
||||
export const PING_RESPONSE = "io.nats.micro.v1.ping_response";
|
||||
export const STATS_RESPONSE = "io.nats.micro.v1.stats_response";
|
||||
|
||||
/** One tool served on one subject, as announced. */
|
||||
export interface Endpoint {
|
||||
kind: "tool" | "seat";
|
||||
module: string;
|
||||
tool: string;
|
||||
seat?: string;
|
||||
scope?: "mesh" | "node";
|
||||
node: string;
|
||||
description: string;
|
||||
schema: unknown;
|
||||
interchangeable: boolean;
|
||||
subject: string;
|
||||
queue?: string;
|
||||
}
|
||||
|
||||
/** A seat's verb as the runtime serves it, with the definition that answers it. */
|
||||
export interface ServedSeatVerb {
|
||||
seat: string;
|
||||
verb: string;
|
||||
subject: string;
|
||||
holder: string;
|
||||
tool: ToolDefinition;
|
||||
}
|
||||
|
||||
export interface Service {
|
||||
name: string;
|
||||
id: string;
|
||||
description: string;
|
||||
metadata: Record<string, string>;
|
||||
}
|
||||
|
||||
/** The endpoint's name as the protocol allows it; the metadata, not the name, identifies it. */
|
||||
function nameOf(e: Endpoint): string {
|
||||
const clean = (s: string) => s.replace(/[^A-Za-z0-9_-]/g, "_");
|
||||
return `${clean(e.kind === "seat" ? e.seat ?? "" : e.module)}__${clean(e.tool)}`;
|
||||
}
|
||||
|
||||
/** The info_response for these endpoints. */
|
||||
export function info(s: Service, endpoints: Endpoint[]): Record<string, unknown> {
|
||||
return {
|
||||
name: s.name, id: s.id, version: VERSION, metadata: s.metadata, type: INFO_RESPONSE, description: s.description,
|
||||
endpoints: endpoints.map((e) => {
|
||||
const metadata: Record<string, string> = {
|
||||
kind: e.kind, module: e.module, tool: e.tool, node: e.node, description: e.description,
|
||||
schema: JSON.stringify(e.schema ?? {}), interchangeable: e.interchangeable ? "true" : "false",
|
||||
};
|
||||
if (e.kind === "seat") {
|
||||
metadata.seat = e.seat ?? "";
|
||||
metadata.scope = e.scope ?? "mesh";
|
||||
}
|
||||
return { name: nameOf(e), subject: e.subject, queue_group: e.queue ?? "", metadata };
|
||||
}),
|
||||
};
|
||||
}
|
||||
|
||||
/** Everything served now: each module's tools on every subject issued for them, and each held
|
||||
* seat's verbs on the seat's subject. */
|
||||
export function endpointsOf(
|
||||
broker: RuntimeBroker,
|
||||
own: { module: string; tools: ToolDefinition[] }[],
|
||||
seats: ServedSeatVerb[],
|
||||
): Endpoint[] {
|
||||
const node = broker.node ?? "";
|
||||
const out: Endpoint[] = [];
|
||||
for (const { module, tools } of own) {
|
||||
for (const t of tools) {
|
||||
const served = broker.servedOn ? broker.servedOn(module, t.name) : [];
|
||||
const interchangeable = served.some((s) => s.subject === `mesh.mod.${module}.tool.${t.name}`);
|
||||
for (const s of served) {
|
||||
out.push({ kind: "tool", module, tool: t.name, node, description: t.description, schema: asSchema(t.input),
|
||||
interchangeable, subject: s.subject, queue: s.queue });
|
||||
}
|
||||
}
|
||||
}
|
||||
for (const v of seats) {
|
||||
out.push({ kind: "seat", module: v.holder, tool: v.verb, seat: v.seat,
|
||||
scope: node && v.subject.endsWith(`.${node}`) ? "node" : "mesh", node, description: v.tool.description,
|
||||
schema: asSchema(v.tool.input), interchangeable: false, subject: v.subject });
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Answer discovery for one service until stopped. */
|
||||
export function announce(broker: RuntimeBroker, s: Service, current: () => Endpoint[]): () => void {
|
||||
const started = new Date().toISOString();
|
||||
const identity = { name: s.name, id: s.id, version: VERSION, metadata: s.metadata };
|
||||
const answer = (subject: string): Uint8Array | undefined => {
|
||||
const parts = subject.split(".");
|
||||
if (parts[0] !== "$SRV" || parts.length < 2) return undefined;
|
||||
if (parts.length >= 3 && parts[2] !== s.name) return undefined; // another service's
|
||||
if (parts.length >= 4 && parts[3] !== s.id) return undefined; // another instance's
|
||||
let v: unknown;
|
||||
switch (parts[1]) {
|
||||
case "PING":
|
||||
v = { ...identity, type: PING_RESPONSE };
|
||||
break;
|
||||
case "INFO":
|
||||
v = info(s, current());
|
||||
break;
|
||||
case "STATS":
|
||||
v = { ...identity, type: STATS_RESPONSE, started, endpoints: current().map((e) => ({
|
||||
name: nameOf(e), subject: e.subject, queue_group: e.queue ?? "", num_requests: 0, num_errors: 0,
|
||||
last_error: "", processing_time: 0, average_processing_time: 0 })) };
|
||||
break;
|
||||
default:
|
||||
return undefined;
|
||||
}
|
||||
return sc.encode(JSON.stringify(v));
|
||||
};
|
||||
const stops: (() => void)[] = [];
|
||||
for (const verb of ["PING", "INFO", "STATS"]) {
|
||||
for (const subject of [`$SRV.${verb}`, `$SRV.${verb}.>`]) stops.push(broker.raw!(subject, (subj) => answer(subj)));
|
||||
}
|
||||
return () => stops.forEach((stop) => stop());
|
||||
}
|
||||
@@ -97,6 +97,13 @@ export interface RuntimeBroker extends Broker {
|
||||
serving(): string[];
|
||||
/** The module this connection is: what its credential named, and what a bare key serves as. */
|
||||
readonly module: string;
|
||||
/** Where a served module's tool is answered right now (ADR 0197: what it serves, it announces). */
|
||||
servedOn?(module: string, tool: string): { subject: string; queue?: string }[];
|
||||
/** Answer a subject in a format of its own, not a tool's reply envelope — the NATS services
|
||||
* protocol's discovery (novox/hq ADR 0197). An undefined answer is no reply. */
|
||||
raw?(subject: string, answer: (subject: string, data: Uint8Array) => Uint8Array | undefined): () => void;
|
||||
/** The machine this connection serves on, when its credential names one. */
|
||||
readonly node?: string;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -348,6 +355,19 @@ export async function connectNats(
|
||||
},
|
||||
|
||||
module: self,
|
||||
node,
|
||||
servedOn,
|
||||
raw(subject: string, answer: (subject: string, data: Uint8Array) => Uint8Array | undefined): () => void {
|
||||
const sub = conn.subscribe(subject);
|
||||
subs.push(sub);
|
||||
void (async () => {
|
||||
for await (const msg of sub) {
|
||||
const body = answer(msg.subject, msg.data);
|
||||
if (body) msg.respond(body);
|
||||
}
|
||||
})();
|
||||
return () => sub.unsubscribe();
|
||||
},
|
||||
follow,
|
||||
serving: () => [...issued.keys()],
|
||||
membership: (module?: string) => issued.get(module ?? self),
|
||||
|
||||
@@ -13,6 +13,8 @@
|
||||
import { spawn, type ChildProcess } from "node:child_process";
|
||||
import { accessSync, constants } from "node:fs";
|
||||
import type { ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import { broker, type Envelope } from "@novox/mesh-sdk/messaging";
|
||||
import { atWork } from "./broker-nats.js";
|
||||
|
||||
/** The protocol version this speaks; a bundle says the same. */
|
||||
export const PROTOCOL = "2025-03-26";
|
||||
@@ -71,29 +73,50 @@ export async function launch(module: string, entry: string, env: NodeJS.ProcessE
|
||||
const line = buffered.slice(0, at).trim();
|
||||
buffered = buffered.slice(at + 1);
|
||||
if (!line) continue;
|
||||
let reply: { id?: number; result?: unknown; error?: { message?: string } };
|
||||
let reply: { id?: number | string; method?: string; params?: unknown; result?: unknown; error?: { message?: string } };
|
||||
try {
|
||||
reply = JSON.parse(line);
|
||||
} catch {
|
||||
console.log(`[mesh-tools] ${module}'s bundle said something that is not a reply: ${line.slice(0, 120)}`);
|
||||
continue;
|
||||
}
|
||||
// **The bundle asks the runtime to emit** (novox/hq ADR 0193): published on the bus as this
|
||||
// module, and answered once the bus has accepted it, so the tool's emit means what it means
|
||||
// in-process. Nothing else a bundle may ask.
|
||||
if (typeof reply.method === "string") {
|
||||
const id = reply.id;
|
||||
const answer = (m: Record<string, unknown>) => proc.stdin!.write(JSON.stringify({ jsonrpc: "2.0", id, ...m }) + "\n");
|
||||
if (reply.method !== "mesh/publish") {
|
||||
if (id !== undefined) answer({ error: { code: -32601, message: `the runtime answers no ${reply.method} from a bundle` } });
|
||||
continue;
|
||||
}
|
||||
atWork.run({ module }, () => broker().publish(reply.params as Envelope<unknown>))
|
||||
.then(() => { if (id !== undefined) answer({ result: {} }); })
|
||||
.catch((err: unknown) => { if (id !== undefined) answer({ error: { code: -32000, message: err instanceof Error ? err.message : String(err) } }); });
|
||||
continue;
|
||||
}
|
||||
const waiting = typeof reply.id === "number" ? pending.get(reply.id) : undefined;
|
||||
if (!waiting) continue;
|
||||
pending.delete(reply.id!);
|
||||
pending.delete(reply.id as number);
|
||||
clearTimeout(waiting.timer);
|
||||
if (reply.error) waiting.reject(new Error(reply.error.message ?? "the bundle refused the request"));
|
||||
else waiting.resolve(reply.result);
|
||||
}
|
||||
});
|
||||
// stderr is the bundle's log; kept under the module's name so a fault reads where it belongs.
|
||||
// The last thing it said is kept, so a bundle that dies says why in its own words, not by code.
|
||||
let lastSaid = "";
|
||||
proc.stderr!.on("data", (chunk: Buffer) => {
|
||||
for (const line of chunk.toString("utf8").split("\n")) if (line.trim()) console.log(`[${module}] ${line}`);
|
||||
for (const line of chunk.toString("utf8").split("\n")) {
|
||||
if (!line.trim()) continue;
|
||||
console.log(`[${module}] ${line}`);
|
||||
if (/\S/.test(line) && !/^\s+at\s/.test(line) && !/^Node\.js v/.test(line)) lastSaid = line.trim();
|
||||
}
|
||||
});
|
||||
const exited = new Promise<never>((_, reject) => {
|
||||
proc.once("error", (err) => reject(err));
|
||||
proc.once("exit", (code, signal) => {
|
||||
const why = `${module}'s bundle exited (${signal ?? code})`;
|
||||
const why = `${module}'s bundle exited (${signal ?? code})` + (lastSaid ? `: ${lastSaid}` : "");
|
||||
for (const [id, p] of pending) {
|
||||
pending.delete(id);
|
||||
clearTimeout(p.timer);
|
||||
|
||||
+51
-97
@@ -9,16 +9,14 @@
|
||||
// and the per-module shape is the list with one entry. A bundle that fails to import is named —
|
||||
// in the log and in what `tools` answers for it — and the others serve.
|
||||
|
||||
import { fileURLToPath, pathToFileURL } from "node:url";
|
||||
import { dirname, join, resolve } from "node:path";
|
||||
import { existsSync, readFileSync } from "node:fs";
|
||||
import { registerHooks } from "node:module";
|
||||
import { pathToFileURL } from "node:url";
|
||||
import { resolve } from "node:path";
|
||||
import { useBroker } from "@novox/mesh-sdk/messaging";
|
||||
import * as sdkTools from "@novox/mesh-sdk/tools";
|
||||
import { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import type { Broker } from "@novox/mesh-sdk/messaging";
|
||||
import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
|
||||
import { launch, launches } from "./launch.js";
|
||||
import { announce, endpointsOf, type ServedSeatVerb } from "./announce.js";
|
||||
|
||||
/**
|
||||
* The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the
|
||||
@@ -115,6 +113,7 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
||||
for (const s of opts.serves ?? []) {
|
||||
served.set(s.module, [...(served.get(s.module) ?? []), ...s.entrypoints]);
|
||||
}
|
||||
const ownEntrypoints = opts.moduleEntrypoints ?? [];
|
||||
if (opts.moduleEntrypoints?.length) {
|
||||
if (!self) {
|
||||
throw new Error(
|
||||
@@ -143,34 +142,38 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
||||
const envs = opts.envs ?? new Map<string, Record<string, string>>();
|
||||
const envFor = (module: string): NodeJS.ProcessEnv => ({ ...process.env, ...(envs.get(module) ?? {}) });
|
||||
|
||||
// Import each bundle, guarded (ADR 0175: one faulty bundle must not take the node's tools down).
|
||||
// Importing the entrypoint runs its registerModuleTools(...) — that is the whole handshake — and
|
||||
// the registrations it adds are the ones that appear after it, which is how each is attributed
|
||||
// to the module whose bundle made it.
|
||||
// A bundle that is not plain JavaScript — or is marked executable — is launched as a process
|
||||
// and spoken to over MCP on stdio instead (ADR 0188); what it lists is registered the same way.
|
||||
// Every bundle's import of the SDK resolves to this runtime's copy (04-ISSUES/209): one registry
|
||||
// of tools, one broker. Installed before the first bundle is imported.
|
||||
oneSdk();
|
||||
// Every bundle this runtime serves is launched as a process and spoken to over MCP on stdio (ADR
|
||||
// 0188, ADR 0193): given the runtime's words and its module's own, and told the module it serves
|
||||
// it as. The runtime knows no language; an entrypoint that is not executable was not built to be
|
||||
// served, and is refused by name. A bundle that fails to start is named, and the others serve
|
||||
// (ADR 0175: one faulty bundle must not take the node's tools down).
|
||||
//
|
||||
// The one-module form — the credential's own module's entrypoints, which the per-module containers
|
||||
// still use for their event handlers and provisioners until they move (to-be 38 WP4c) — is imported
|
||||
// into this process as before: one module, one SDK, its container's own environment.
|
||||
// The machine this runtime serves, which an event a launched tool emits is stamped with.
|
||||
const node = opts.credential?.node ?? (typeof (runtime as unknown as { node?: unknown }).node === "string" ? (runtime as unknown as { node: string }).node : undefined);
|
||||
const failed = new Map<string, string>();
|
||||
const owner: string[] = []; // registration index → the module whose bundle registered it
|
||||
const launched: { module: string; owner: string; tools: ToolDefinition[] }[] = [];
|
||||
const children: Array<() => void> = [];
|
||||
for (const [module, entrypoints] of served) {
|
||||
const imported = module === self && ownEntrypoints.length > 0;
|
||||
for (const entry of entrypoints) {
|
||||
const path = resolve(entry);
|
||||
try {
|
||||
if (launches(path)) {
|
||||
// Told the module it serves it as, so a seat's verbs are the seat's (ADR 0193).
|
||||
const child = await launch(module, path, { ...envFor(module), MESH_SERVED_MODULE: module });
|
||||
children.push(child.stop);
|
||||
for (const r of child.registrations) launched.push({ ...r, owner: module });
|
||||
if (imported) {
|
||||
await import(pathToFileURL(path).href);
|
||||
continue;
|
||||
}
|
||||
const before = collectTools().length;
|
||||
await import(pathToFileURL(path).href);
|
||||
const after = collectTools().length;
|
||||
for (let i = before; i < after; i++) owner[i] = module;
|
||||
if (!launches(path)) {
|
||||
throw new Error(`${path} is not executable; a bundle the runtime serves is started, never imported, and its build makes it executable (novox/hq ADR 0193)`);
|
||||
}
|
||||
const child = await launch(module, path, {
|
||||
...envFor(module), MESH_SERVED_MODULE: module, MESH_MODULE: module,
|
||||
...(node ? { MESH_NODE: node } : {}),
|
||||
});
|
||||
children.push(child.stop);
|
||||
for (const r of child.registrations) launched.push({ ...r, owner: module });
|
||||
} catch (err) {
|
||||
const why = err instanceof Error ? err.message : String(err);
|
||||
failed.set(module, why);
|
||||
@@ -187,16 +190,8 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
||||
// out rather than fatal — on 2026-10-01 the credential of a module that had just learned to
|
||||
// implement a seat did not yet name the claim, and the whole runtime restarted for it.
|
||||
const claimed = seatsClaimed(served.keys(), self, opts.credential, runtime);
|
||||
// Each registration's contributor is given its own module's environment and no other's (ADR 0192):
|
||||
// the runtime's own words, and over them what the mesh composed for the module whose bundle made
|
||||
// the registration. An SDK too old to ask per registration cannot do that; said, not hidden.
|
||||
const each = (sdkTools as { collectToolsEach?: (f: (module: string, i: number) => NodeJS.ProcessEnv) => { module: string; tools: ToolDefinition[] }[] }).collectToolsEach;
|
||||
if (!each && envs.size > 0) {
|
||||
console.log(`[mesh-tools] this runtime's SDK cannot give each bundle its own environment; ${[...envs.keys()].join(", ")} serve with the runtime's words only (novox/hq ADR 0192)`);
|
||||
}
|
||||
const collected = each ? each((module, i) => envFor(owner[i] ?? self ?? module)) : collectTools();
|
||||
const registrations = [
|
||||
...collected.map((r, i) => ({ ...r, owner: owner[i] ?? self ?? r.module })),
|
||||
...collectTools().map((r) => ({ ...r, owner: self ?? r.module })),
|
||||
...launched,
|
||||
];
|
||||
const ownRegistrations = registrations.filter(({ module, owner: by }) => {
|
||||
@@ -265,7 +260,16 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
|
||||
|
||||
console.log(`[mesh-tools] serving ${names.length} tool(s) for ${served.size} module(s): ${names.join(", ") || "(none)"}` +
|
||||
(failed.size ? `; not serving ${[...failed.keys()].join(", ")}, whose bundle(s) failed to load` : ""));
|
||||
stops.push(await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations));
|
||||
const seats = await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations);
|
||||
stops.push(seats.stop);
|
||||
// **What it serves, it announces** (novox/hq ADR 0197), asked at the moment of the request.
|
||||
if (typeof runtime.raw === "function") {
|
||||
stops.push(announce(runtime, {
|
||||
name: self ?? "runtime", id: runtime.node ?? self ?? "runtime",
|
||||
description: `the tool runtime of ${self ?? "a module"}${runtime.node ? ` on ${runtime.node}` : ""}`,
|
||||
metadata: runtime.node ? { node: runtime.node } : {},
|
||||
}, () => endpointsOf(runtime, ownRegistrations, seats.serving())));
|
||||
}
|
||||
return () => stop();
|
||||
}
|
||||
|
||||
@@ -309,11 +313,17 @@ async function serveClaimedSeats(
|
||||
self: string | undefined,
|
||||
credential: Credential | undefined,
|
||||
registrations: { module: string; owner: string; tools: ToolDefinition[] }[],
|
||||
): Promise<() => void> {
|
||||
if (typeof broker.handleSubject !== "function") return () => {};
|
||||
): Promise<{ stop: () => void; serving: () => ServedSeatVerb[] }> {
|
||||
if (typeof broker.handleSubject !== "function") return { stop: () => {}, serving: () => [] };
|
||||
// A seat's verbs are the role's, not the software's (ADR 0159): implemented under the seat's
|
||||
// name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools.
|
||||
const implementations = new Map<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>();
|
||||
const definitions = new Map<string, Map<string, ToolDefinition>>();
|
||||
for (const { module, tools } of registrations) {
|
||||
const defs = definitions.get(module) ?? new Map<string, ToolDefinition>();
|
||||
for (const t of tools) defs.set(t.name, t);
|
||||
definitions.set(module, defs);
|
||||
}
|
||||
for (const { module, owner, tools } of registrations) {
|
||||
const verbs = implementations.get(module) ?? new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
|
||||
for (const t of tools) verbs.set(t.name, (args) => atWork.run({ module: owner }, () => t.run(args)));
|
||||
@@ -347,81 +357,25 @@ async function serveClaimedSeats(
|
||||
};
|
||||
|
||||
let stops: (() => void)[] = [];
|
||||
let servingNow: ServedSeatVerb[] = [];
|
||||
const serve = async (): Promise<void> => {
|
||||
stops.forEach((s) => s());
|
||||
stops = [];
|
||||
servingNow = [];
|
||||
for (const v of wanted()) {
|
||||
const run = implementations.get(v.seat)?.get(v.verb);
|
||||
const tool = definitions.get(v.seat)?.get(v.verb);
|
||||
if (!run) {
|
||||
console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`);
|
||||
continue;
|
||||
}
|
||||
stops.push(await broker.handleSubject(v.subject, run));
|
||||
if (tool) servingNow.push({ ...v, tool });
|
||||
console.log(`[mesh-tools] serving ${v.seat}'s ${v.verb} on ${v.subject}, admitted where ${v.holder} holds the seat`);
|
||||
}
|
||||
};
|
||||
await serve();
|
||||
// A membership issued to any served module may add, move or withdraw a seat's verbs.
|
||||
if (typeof broker.onMembership === "function") broker.onMembership(() => void serve());
|
||||
return () => stops.forEach((s) => s());
|
||||
}
|
||||
|
||||
let sdkHooked = false;
|
||||
const SDK = "@novox/mesh-sdk";
|
||||
/**
|
||||
* The one SDK in a node's runtime (novox/hq ADR 0175, 04-ISSUES/209).
|
||||
*
|
||||
* A bundle carries its own dependencies — the toolchain copies them in so a bundle starts anywhere
|
||||
* (to-be 38 WP3) — and among them is a copy of the SDK. Imported in this process, that copy would be
|
||||
* a second SDK: its own registry of tools, its own broker handle. A bundle calling registerModuleTools
|
||||
* through it registers into a list this runtime never reads, and its tools are silently not served.
|
||||
* So every import of the SDK, from whichever bundle, is resolved as if this runtime had written it:
|
||||
* one registry, one broker — the runtime's. Everything else a bundle carries resolves from the
|
||||
* bundle's own tree, as before. A launched bundle (ADR 0188) is another process and is untouched.
|
||||
*
|
||||
* Installed once, in-thread, before the first bundle is imported; the hook sees every import after,
|
||||
* `require` included. A bundle whose own copy is another version than the runtime's is said once,
|
||||
* so a tool failing against the runtime's SDK points at the bundle rather than at the runtime.
|
||||
*/
|
||||
function oneSdk(): void {
|
||||
if (sdkHooked) return;
|
||||
registerHooks({
|
||||
resolve(specifier, context, next) {
|
||||
if (specifier === SDK || specifier.startsWith(SDK + "/")) {
|
||||
if (context.parentURL) sayOtherSdk(context.parentURL);
|
||||
return next(specifier, { ...context, parentURL: import.meta.url });
|
||||
}
|
||||
return next(specifier, context);
|
||||
},
|
||||
});
|
||||
sdkHooked = true;
|
||||
}
|
||||
|
||||
const sdkSaid = new Set<string>();
|
||||
/** The version of the SDK copy nearest a file, by its package.json, or nothing when the file has none above it. */
|
||||
function sdkVersionNear(fileURL: string): { dir: string; version: string } | undefined {
|
||||
let dir = dirname(fileURLToPath(fileURL));
|
||||
for (;;) {
|
||||
const pkg = join(dir, "node_modules", SDK, "package.json");
|
||||
if (existsSync(pkg)) {
|
||||
try {
|
||||
return { dir, version: String((JSON.parse(readFileSync(pkg, "utf8")) as { version?: string }).version ?? "?") };
|
||||
} catch {
|
||||
return { dir, version: "?" };
|
||||
}
|
||||
}
|
||||
const up = dirname(dir);
|
||||
if (up === dir) return undefined;
|
||||
dir = up;
|
||||
}
|
||||
}
|
||||
function sayOtherSdk(parentURL: string): void {
|
||||
if (!parentURL.startsWith("file:")) return;
|
||||
const own = sdkVersionNear(import.meta.url);
|
||||
const theirs = sdkVersionNear(parentURL);
|
||||
if (!theirs || !own || theirs.dir === own.dir || sdkSaid.has(theirs.dir)) return;
|
||||
sdkSaid.add(theirs.dir);
|
||||
if (theirs.version !== own.version) {
|
||||
console.log(`[mesh-tools] ${theirs.dir} carries ${SDK} ${theirs.version}; this runtime's is ${own.version}, and the bundle speaks to the runtime's`);
|
||||
}
|
||||
return { stop: () => stops.forEach((s) => s()), serving: () => servingNow };
|
||||
}
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
/**
|
||||
* What a runtime serves, it announces (novox/hq ADR 0197): the TypeScript runtime the per-module
|
||||
* containers still run answers the NATS services protocol's discovery in the same shape as the Go
|
||||
* tool runtime — its module's tools on every subject issued, and the seat verbs it serves.
|
||||
*
|
||||
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/announce.test.ts
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "node:test";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { connect, StringCodec } from "nats";
|
||||
|
||||
import { resetTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
import { connectNats, membershipSubject } from "../dist/broker-nats.js";
|
||||
import { runTools } from "../dist/runtime.js";
|
||||
|
||||
const url = process.env.MESH_TEST_NATS;
|
||||
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
|
||||
const sc = StringCodec();
|
||||
|
||||
test("the runtime answers $SRV.INFO with what it serves, in the services protocol's format", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const nc = await connect({ servers: url });
|
||||
const jsm = await nc.jetstreamManager();
|
||||
try {
|
||||
await jsm.streams.delete("ASSIGNMENTS");
|
||||
} catch {
|
||||
// none yet
|
||||
}
|
||||
await jsm.streams.add({ name: "ASSIGNMENTS", subjects: ["mesh.assignment.>"], max_msgs_per_subject: 1, allow_direct: true } as never);
|
||||
await nc.jetstream().publish(membershipSubject("anchor", "shop"), sc.encode(JSON.stringify({
|
||||
node: "anchor", module: "shop",
|
||||
serves: [{ subject: "mesh.mod.shop.tool.{tool}.anchor" }, { subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }],
|
||||
emits: "mesh.mod.shop.event.{event}", tools: "mesh.mod.shop.tool.tools",
|
||||
})));
|
||||
const shop = await connectNats({ url, node: "anchor", module: "shop" });
|
||||
let stop = () => {};
|
||||
try {
|
||||
stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] });
|
||||
const msg = await nc.request("$SRV.INFO.shop.anchor", sc.encode(""), { timeout: 2000 });
|
||||
const info = JSON.parse(sc.decode(msg.data)) as {
|
||||
type: string; name: string; id: string; version: string;
|
||||
endpoints: { name: string; subject: string; queue_group: string; metadata: Record<string, string> }[];
|
||||
};
|
||||
assert.equal(info.type, "io.nats.micro.v1.info_response");
|
||||
assert.equal(info.name, "shop");
|
||||
assert.equal(info.id, "anchor");
|
||||
assert.ok(info.version);
|
||||
const price = info.endpoints.filter((e) => e.metadata.tool === "price").map((e) => `${e.subject}|${e.queue_group}`).sort();
|
||||
assert.deepEqual(price, ["mesh.mod.shop.tool.price.anchor|", "mesh.mod.shop.tool.price|serve.shop"]);
|
||||
const one = info.endpoints.find((e) => e.metadata.tool === "price")!;
|
||||
assert.equal(one.metadata.kind, "tool");
|
||||
assert.equal(one.metadata.module, "shop");
|
||||
assert.equal(one.metadata.node, "anchor");
|
||||
assert.equal(one.metadata.interchangeable, "true");
|
||||
assert.ok(JSON.parse(one.metadata.schema).type === "object");
|
||||
// Ping answers with the same identity; another service's request is not answered.
|
||||
const ping = JSON.parse(sc.decode((await nc.request("$SRV.PING", sc.encode(""), { timeout: 2000 })).data));
|
||||
assert.equal(ping.type, "io.nats.micro.v1.ping_response");
|
||||
await assert.rejects(nc.request("$SRV.INFO.somebody-else", sc.encode(""), { timeout: 300 }));
|
||||
} finally {
|
||||
stop();
|
||||
await shop.close();
|
||||
await nc.close();
|
||||
resetTools();
|
||||
}
|
||||
});
|
||||
@@ -122,9 +122,9 @@ test("the node's runtime serves five modules' bundles on one credential — two
|
||||
broker: nodeTools,
|
||||
credential,
|
||||
serves: [
|
||||
{ module: "alpha", entrypoints: [fixture("many-alpha.mjs")] },
|
||||
{ module: "beta", entrypoints: [fixture("many-beta.mjs")] },
|
||||
{ module: "gamma", entrypoints: [fixture("many-broken.mjs")] },
|
||||
{ module: "alpha", entrypoints: [fixture("many-alpha.serve.mjs")] },
|
||||
{ module: "beta", entrypoints: [fixture("many-beta.serve.mjs")] },
|
||||
{ module: "gamma", entrypoints: [fixture("many-broken.serve.mjs")] },
|
||||
{ module: "delta", entrypoints: [fixture("many-delta.py")] },
|
||||
{ module: "epsilon", entrypoints: [fixture("many-epsilon.mjs")] },
|
||||
],
|
||||
@@ -132,7 +132,7 @@ test("the node's runtime serves five modules' bundles on one credential — two
|
||||
console.log = log;
|
||||
assert.deepEqual(nodeTools.serving().sort(), ["alpha", "beta", "delta", "epsilon", "gamma", "node-tools"]);
|
||||
assert.ok(said.some((s) => /the operator's account here is somebody \(home \/home\/somebody\)/.test(s)), said.join("\n"));
|
||||
assert.ok(said.some((s) => /gamma's bundle .*many-broken\.mjs failed to load: gamma's bundle cannot find its client; its tools are not served here/.test(s)), said.join("\n"));
|
||||
assert.ok(said.some((s) => /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/.test(s)), said.join("\n"));
|
||||
assert.ok(said.some((s) => /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/.test(s)), said.join("\n"));
|
||||
|
||||
// Five tools answer, each where its module's membership says: alpha anywhere and here, beta here only.
|
||||
@@ -163,7 +163,7 @@ test("the node's runtime serves five modules' bundles on one credential — two
|
||||
|
||||
// `tools` answers for each: what alpha and beta serve, and why gamma serves nothing.
|
||||
const gamma = await asker.request<Record<string, never>, { module: string; tools: unknown[]; failed?: string }>("gamma.tools@anchor", {});
|
||||
assert.deepEqual(gamma, { module: "gamma", tools: [], failed: "gamma's bundle cannot find its client" });
|
||||
assert.deepEqual(gamma, { module: "gamma", tools: [], failed: "gamma's bundle exited (1): Error: gamma's bundle cannot find its client" });
|
||||
const beta = await asker.request<Record<string, never>, { tools: { name: string; subjects?: string[] }[] }>("beta.tools@anchor", {});
|
||||
assert.deepEqual(beta.tools.map((x) => x.name), ["three", "four", "five"]);
|
||||
assert.deepEqual(beta.tools[0]!.subjects, ["mesh.mod.beta.tool.three.anchor"]);
|
||||
@@ -174,7 +174,7 @@ test("the node's runtime serves five modules' bundles on one credential — two
|
||||
try {
|
||||
const have = await toolsOn(asker);
|
||||
assert.deepEqual(have.tools.map((x) => `${x.module}.${x.name}`), ["alpha.one", "alpha.two", "beta.five", "beta.four", "beta.three", "delta.die", "delta.greet", "epsilon.seven"]);
|
||||
assert.deepEqual(have.notAnswering, ["gamma (its tools bundle failed to load: gamma's bundle cannot find its client)", "mesh-controller (seat)"]);
|
||||
assert.deepEqual(have.notAnswering, ["gamma (its tools bundle failed to load: gamma's bundle exited (1): Error: gamma's bundle cannot find its client)", "mesh-controller (seat)"]);
|
||||
} finally {
|
||||
await catalogue.close();
|
||||
}
|
||||
@@ -215,6 +215,8 @@ test("a bundle carrying its own copy of the SDK registers into the runtime's reg
|
||||
writeFileSync(join(dir, "node_modules", "zeta-flavour", "package.json"), '{"name":"zeta-flavour","type":"module","main":"index.js"}\n');
|
||||
writeFileSync(join(dir, "node_modules", "zeta-flavour", "index.js"), 'export const flavour = "the bundle\'s own";\n');
|
||||
writeFileSync(join(dir, "package.json"), '{"type":"module","private":true}\n');
|
||||
writeFileSync(join(dir, "index.serve.mjs"),
|
||||
'#!/usr/bin/env node\nimport { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";\nawait import("./index.js");\nawait serveRegisteredOverStdio();\n', { mode: 0o755 });
|
||||
writeFileSync(join(dir, "index.js"),
|
||||
'import { registerModuleTools } from "@novox/mesh-sdk/tools";\n' +
|
||||
'import { flavour } from "zeta-flavour";\n' +
|
||||
@@ -228,7 +230,7 @@ test("a bundle carrying its own copy of the SDK registers into the runtime's reg
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
closing.push(() => asker.close());
|
||||
console.log = (...a: unknown[]) => said.push(a.join(" "));
|
||||
stop = await runTools({ broker: nodeTools, credential, serves: [{ module: "zeta", entrypoints: [join(dir, "index.js")] }] });
|
||||
stop = await runTools({ broker: nodeTools, credential, serves: [{ module: "zeta", entrypoints: [join(dir, "index.serve.mjs")] }] });
|
||||
console.log = log;
|
||||
assert.ok(said.some((s) => /serving 1 tool\(s\) for 1 module\(s\): zeta\.probe/.test(s)), said.join("\n"));
|
||||
// The SDK is the runtime's (the registration arrived); the bundle's other dependency is its own.
|
||||
@@ -266,8 +268,8 @@ test("each bundle is given its own environment and none of another's, imported o
|
||||
stop = await runTools({
|
||||
broker: nodeTools, credential, envs,
|
||||
serves: [
|
||||
{ module: "gamma", entrypoints: [fixture("env-gamma.mjs")] },
|
||||
{ module: "delta", entrypoints: [fixture("env-delta.mjs")] },
|
||||
{ module: "gamma", entrypoints: [fixture("env-gamma.serve.mjs")] },
|
||||
{ module: "delta", entrypoints: [fixture("env-delta.serve.mjs")] },
|
||||
{ module: "zeta", entrypoints: [fixture("env-zeta.mjs")] },
|
||||
],
|
||||
});
|
||||
@@ -316,3 +318,33 @@ test("a launched bundle is told the module it serves, so its seat's verbs stay t
|
||||
resetTools();
|
||||
}
|
||||
});
|
||||
|
||||
test("an entrypoint that is not executable is refused by name, and the others serve (ADR 0193)", async (t) => {
|
||||
if (!url) return t.skip("MESH_TEST_NATS unset");
|
||||
resetTools();
|
||||
const mesh = await aMesh();
|
||||
for (const m of ["alpha", "plain"]) await mesh.issue(membershipOf(m, "anchor"));
|
||||
const credential = { url, node: "anchor", module: "node-tools" };
|
||||
const nodeTools = await connectNats(credential);
|
||||
const asker = await connectNats({ url, module: "console", node: "workstation" });
|
||||
const said: string[] = [];
|
||||
const log = console.log;
|
||||
let stop = () => {};
|
||||
try {
|
||||
console.log = (...a: unknown[]) => said.push(a.join(" "));
|
||||
stop = await runTools({ broker: nodeTools, credential, serves: [
|
||||
{ module: "alpha", entrypoints: [fixture("many-alpha.serve.mjs")] },
|
||||
{ module: "plain", entrypoints: [fixture("many-alpha.mjs")] },
|
||||
] });
|
||||
console.log = log;
|
||||
assert.ok(said.some((s) => /plain's bundle .*many-alpha\.mjs failed to load: .* is not executable; a bundle the runtime serves is started, never imported/.test(s)), said.join("\n"));
|
||||
assert.deepEqual((await callTool(asker, "alpha.one@anchor", {})).result, { alpha: 1 });
|
||||
} finally {
|
||||
console.log = log;
|
||||
stop();
|
||||
await asker.close();
|
||||
await nodeTools.close();
|
||||
await mesh.close();
|
||||
resetTools();
|
||||
}
|
||||
});
|
||||
|
||||
@@ -36,7 +36,7 @@ test("as node-tools, serve loads the bundles and is the console on loopback", as
|
||||
env: {
|
||||
...process.env,
|
||||
MESH_BROKER_FILE: credential,
|
||||
MESH_TOOL_MODULES: `alpha=${fixture("many-alpha.mjs")}`,
|
||||
MESH_TOOL_MODULES: `alpha=${fixture("many-alpha.serve.mjs")}`,
|
||||
MESH_CONSOLE_LISTEN: "127.0.0.1:0",
|
||||
},
|
||||
stdio: ["ignore", "pipe", "pipe"],
|
||||
|
||||
Reference in New Issue
Block a user