diff --git a/node-tools/internal/announce/announce.go b/node-tools/internal/announce/announce.go new file mode 100644 index 0000000..c9f7f97 --- /dev/null +++ b/node-tools/internal/announce/announce.go @@ -0,0 +1,209 @@ +// Package announce answers the NATS services protocol's discovery for what a runtime serves, and +// gathers the answers (novox/hq ADR 0197). +// +// A runtime does not re-serve its tools through a services library: serving is unchanged. It answers +// `$SRV.PING`, `$SRV.INFO` and `$SRV.STATS` — and the same followed by its service name, and by its +// name and id — in the format NATS's own tools read, with what it is serving at the moment it is +// asked. One service per runtime process: the bus admits one reply per request from each responder, +// so a runtime serving many modules and seats answers once, one endpoint per tool per subject, and +// says in each endpoint's metadata which module, seat, scope and machine it is. +package announce + +import ( + "encoding/json" + "strings" + "time" + + "github.com/nats-io/nats.go/micro" + + "github.com/novox/mesh-tools/node-tools/internal/bus" +) + +// Version is the announced service version (semver, as the protocol requires). +const Version = "0.1.0" + +// Window is how long the console gathers discovery answers: every instance answers one request, and +// how many will is what is being found out. +var Window = 750 * time.Millisecond + +// Kinds of endpoint. +const ( + KindTool = "tool" // a module's own tool + KindSeat = "seat" // a seat's verb, served by the module holding it +) + +// Endpoint is one tool served on one subject, as it is announced. +type Endpoint struct { + Kind string + Module string // the module whose code answers + Tool string // the tool's or the verb's name + Seat string // for a seat's verb + Scope string // "mesh" or "node", for a seat's verb + Node string + Description string + Schema json.RawMessage + Interchangeable bool + Subject string + Queue string +} + +// Name is the endpoint's name as the protocol allows it — letters, digits, `-` and `_` — the +// prefix and the tool joined by `__`; the metadata, not the name, is what identifies it. +func (e Endpoint) Name() string { + prefix := e.Module + if e.Kind == KindSeat { + prefix = e.Seat + } + return clean(prefix) + "__" + clean(e.Tool) +} + +func clean(s string) string { + var b strings.Builder + for _, r := range s { + if r == '-' || r == '_' || (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') { + b.WriteRune(r) + } else { + b.WriteRune('_') + } + } + return b.String() +} + +func (e Endpoint) info() micro.EndpointInfo { + md := map[string]string{ + "kind": e.Kind, "module": e.Module, "tool": e.Tool, "node": e.Node, + "description": e.Description, "interchangeable": boolWord(e.Interchangeable), + } + schema := strings.TrimSpace(string(e.Schema)) + if schema == "" || schema == "null" { + schema = "{}" + } + md["schema"] = schema + if e.Kind == KindSeat { + md["seat"] = e.Seat + md["scope"] = e.Scope + } + return micro.EndpointInfo{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue, Metadata: md} +} + +func boolWord(b bool) string { + if b { + return "true" + } + return "false" +} + +// Service is who answers: the runtime's name, its instance, and a word about it. +type Service struct { + Name string + ID string + Description string + Metadata map[string]string +} + +func (s Service) identity() micro.ServiceIdentity { + md := s.Metadata + if md == nil { + md = map[string]string{} + } + return micro.ServiceIdentity{Name: s.Name, ID: s.ID, Version: Version, Metadata: md} +} + +// Info is the info_response for these endpoints. +func Info(s Service, endpoints []Endpoint) micro.Info { + out := micro.Info{ServiceIdentity: s.identity(), Type: micro.InfoResponseType, + Description: s.Description, Endpoints: []micro.EndpointInfo{}} + for _, e := range endpoints { + out.Endpoints = append(out.Endpoints, e.info()) + } + return out +} + +// Serve answers discovery for one service until stopped, asking `current` for its endpoints each +// time — so what is announced is what is served now, re-served memberships included. +func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error) { + started := time.Now().UTC() + answer := func(subject string, _ []byte) []byte { + parts := strings.Split(subject, ".") + if len(parts) < 2 || parts[0] != "$SRV" { + return nil + } + if len(parts) >= 3 && parts[2] != s.Name { + return nil // another service's + } + if len(parts) >= 4 && parts[3] != s.ID { + return nil // another instance's + } + var v any + switch parts[1] { + case "PING": + v = micro.Ping{ServiceIdentity: s.identity(), Type: micro.PingResponseType} + case "INFO": + v = Info(s, current()) + case "STATS": + st := micro.Stats{ServiceIdentity: s.identity(), Type: micro.StatsResponseType, Started: started, + Endpoints: []*micro.EndpointStats{}} + for _, e := range current() { + st.Endpoints = append(st.Endpoints, µ.EndpointStats{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue}) + } + v = st + default: + return nil + } + body, err := json.Marshal(v) + if err != nil { + return nil + } + return body + } + var stops []func() + for _, verb := range []string{"PING", "INFO", "STATS"} { + for _, subject := range []string{"$SRV." + verb, "$SRV." + verb + ".>"} { + stop, err := conn.Raw(subject, answer) + if err != nil { + for _, st := range stops { + st() + } + return func() {}, err + } + stops = append(stops, stop) + } + } + return func() { + for _, st := range stops { + st() + } + }, nil +} + +// Gather asks every service on the bus what it serves and answers what came back within the window. +// An answer that is not an info_response is skipped. +func Gather(conn *bus.Conn) ([]micro.Info, error) { + raw, err := conn.Gather("$SRV.INFO", nil, Window) + if err != nil { + return nil, err + } + var out []micro.Info + for _, b := range raw { + var i micro.Info + if json.Unmarshal(b, &i) != nil || i.Type != micro.InfoResponseType { + continue + } + out = append(out, i) + } + return out, nil +} + +// Endpoints reads an info_response's endpoints back into what they announce. +func Endpoints(i micro.Info) []Endpoint { + var out []Endpoint + for _, e := range i.Endpoints { + md := e.Metadata + out = append(out, Endpoint{ + Kind: md["kind"], Module: md["module"], Tool: md["tool"], Seat: md["seat"], Scope: md["scope"], + Node: md["node"], Description: md["description"], Schema: json.RawMessage(md["schema"]), + Interchangeable: md["interchangeable"] == "true", Subject: e.Subject, Queue: e.QueueGroup, + }) + } + return out +} diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index 8fddf95..608dadc 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -512,6 +512,59 @@ func (c *Conn) PublishAs(module string, env Envelope) error { return err } +// ServedOn is where a served module's tool is answered right now: the membership's subjects when +// issued, the derived shape otherwise — what Handle subscribes, for what announces it (ADR 0197). +func (c *Conn) ServedOn(module, tool string) []Served { return c.servedOn(module, tool) } + +// Raw answers one subject with a function of the request, not a tool's reply envelope: the NATS +// services protocol's discovery subjects answer in their own format (novox/hq ADR 0197). A nil +// answer is no reply — the request was for another service. +func (c *Conn) Raw(subject string, answer func(subject string, data []byte) []byte) (func(), error) { + sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) { + if body := answer(msg.Subject, msg.Data); body != nil { + _ = msg.Respond(body) + } + }) + if err != nil { + return func() {}, err + } + c.track(sub) + return func() { _ = sub.Unsubscribe() }, nil +} + +// Gather publishes one request and collects every answer that arrives within the window: a +// discovery request every service instance answers (ADR 0197). It never stops early — how many will +// answer is what it is finding out. +func (c *Conn) Gather(subject string, body []byte, window time.Duration) ([][]byte, error) { + inbox := c.nc.NewRespInbox() + sub, err := c.nc.SubscribeSync(inbox) + if err != nil { + return nil, err + } + defer func() { _ = sub.Unsubscribe() }() + if err := c.nc.PublishRequest(subject, inbox, body); err != nil { + return nil, err + } + var out [][]byte + deadline := time.Now().Add(window) + for { + left := time.Until(deadline) + if left <= 0 { + return out, nil + } + msg, err := sub.NextMsg(left) + if err != nil { + if errors.Is(err, nats.ErrTimeout) { + return out, nil + } + return out, err + } + if len(msg.Data) > 0 { + out = append(out, msg.Data) + } + } +} + // Flush waits until the bus has every subscription made so far, so what is served is answerable // when this returns. func (c *Conn) Flush() { _ = c.nc.Flush() } diff --git a/node-tools/internal/console/address.go b/node-tools/internal/console/address.go index 3909b59..f8b2cb4 100644 --- a/node-tools/internal/console/address.go +++ b/node-tools/internal/console/address.go @@ -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{} // . or . → index in l.Tools + machines := map[string]bool{} + for _, info := range infos { + for _, e := range announce.Endpoints(info) { + if e.Node != "" { + machines[e.Node] = true } - if 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) } diff --git a/node-tools/internal/console/address_test.go b/node-tools/internal/console/address_test.go index 25d897b..f2a3208 100644 --- a/node-tools/internal/console/address_test.go +++ b/node-tools/internal/console/address_test.go @@ -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) } } diff --git a/node-tools/internal/console/client.go b/node-tools/internal/console/client.go index 9cd8c0d..f47a805 100644 --- a/node-tools/internal/console/client.go +++ b/node-tools/internal/console/client.go @@ -5,19 +5,11 @@ import ( "errors" "fmt" "regexp" - "sort" "strings" - "sync" "github.com/novox/mesh-tools/node-tools/internal/bus" ) -const ( - catalogueModules = "mesh-catalog.catalog_modules" - seatTools = "seat:mesh-controller.tools" - toolsVerb = "tools" -) - // Tool is a tool as its module — or, for a role's tool, the mesh's records — describes it. type Tool struct { Module string @@ -67,107 +59,14 @@ func toolKey(name string, seats Seats) string { return name } -// toolsOn asks the mesh what tools it has: the catalogue which modules it holds, each module what it -// serves, the controller's seat every role's tools — at once, so a restarting control plane hides -// nothing else. +// toolsOn is the flat catalogue (MESH_CONSOLE_FLAT=1): what announced itself on the bus, every +// tool and seat verb, and the assignments with tools that did not (novox/hq ADR 0197). func toolsOn(conn *bus.Conn) (*Listing, error) { - type rolesAnswer struct { - Seats []struct { - Seat string `json:"seat"` - Scope string `json:"scope"` - Tools []struct { - Name string `json:"name"` - Description string `json:"description"` - Input json.RawMessage `json:"input"` - } `json:"tools"` - } `json:"seats"` - } - var wg sync.WaitGroup - var roles *rolesAnswer - wg.Add(1) - go func() { - defer wg.Done() - if got, err := conn.Ask(seatTools, map[string]any{}, ""); err == nil { - var r rolesAnswer - if json.Unmarshal(got.Result, &r) == nil { - roles = &r - } - } - }() - answered, err := conn.Ask(catalogueModules, map[string]any{}, "") + x, err := indexOn(conn) if err != nil { - wg.Wait() return nil, err } - var held struct { - Modules []struct { - Module string `json:"module"` - } `json:"modules"` - } - _ = json.Unmarshal(answered.Result, &held) - names := make([]string, 0, len(held.Modules)) - for _, m := range held.Modules { - if m.Module != "" { - names = append(names, m.Module) - } - } - type outcome struct { - ok bool - answer struct { - Tools *[]struct { - Name string `json:"name"` - Description string `json:"description"` - Input json.RawMessage `json:"input"` - Subjects []string `json:"subjects"` - } `json:"tools"` - Failed *string `json:"failed"` - } - } - outcomes := make([]outcome, len(names)) - for i, module := range names { - wg.Add(1) - go func(i int, module string) { - defer wg.Done() - got, err := conn.Ask(module+"."+toolsVerb, map[string]any{}, "") - if err != nil { - return - } - if json.Unmarshal(got.Result, &outcomes[i].answer) == nil { - outcomes[i].ok = true - } - }(i, module) - } - wg.Wait() - l := &Listing{Tools: []Tool{}, NotAnswering: []string{}} - if roles != nil { - for _, s := range roles.Seats { - for _, t := range s.Tools { - l.Tools = append(l.Tools, Tool{Module: s.Seat, Name: t.Name, Description: t.Description, - Input: t.Input, Seat: true, Scope: s.Scope}) - } - } - } else { - l.NotAnswering = append(l.NotAnswering, "mesh-controller (seat)") - } - for i, module := range names { - o := outcomes[i] - switch { - case o.ok && o.answer.Failed != nil: - l.NotAnswering = append(l.NotAnswering, fmt.Sprintf("%s (its tools bundle failed to load: %s)", module, *o.answer.Failed)) - case o.ok && o.answer.Tools != nil: - for _, t := range *o.answer.Tools { - l.Tools = append(l.Tools, Tool{Module: module, Name: t.Name, Description: t.Description, - Input: t.Input, Subjects: t.Subjects}) - } - default: - l.NotAnswering = append(l.NotAnswering, module) - } - } - sort.SliceStable(l.Tools, func(i, j int) bool { - return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name - }) - sort.Strings(l.NotAnswering) - return l, nil + return x.Listing, nil } // callTool calls `.[@]`, on the subject the listing names for it when it names one. diff --git a/node-tools/internal/console/console_test.go b/node-tools/internal/console/console_test.go index 8052ea9..d41b55f 100644 --- a/node-tools/internal/console/console_test.go +++ b/node-tools/internal/console/console_test.go @@ -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) { diff --git a/node-tools/internal/console/mcp.go b/node-tools/internal/console/mcp.go index d26b543..330e2f1 100644 --- a/node-tools/internal/console/mcp.go +++ b/node-tools/internal/console/mcp.go @@ -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 { diff --git a/node-tools/internal/meshtest/meshtest.go b/node-tools/internal/meshtest/meshtest.go index 185eb0c..7342c0d 100644 --- a/node-tools/internal/meshtest/meshtest.go +++ b/node-tools/internal/meshtest/meshtest.go @@ -1,5 +1,9 @@ // Package meshtest raises what the controller would, for tests against a real bus: the ASSIGNMENTS // and EVENTS streams, memberships issued by hand, and a fixture's path. +// +// **The packages share one bus, so run them one at a time: `go test -p 1 ./...`.** Each test raises +// the streams afresh, and the console discovers every runtime that announces itself on the bus +// (novox/hq ADR 0197) — a runtime from another package's test is, correctly, found. package meshtest import ( diff --git a/node-tools/internal/runtime/announce_test.go b/node-tools/internal/runtime/announce_test.go new file mode 100644 index 0000000..69408d7 --- /dev/null +++ b/node-tools/internal/runtime/announce_test.go @@ -0,0 +1,113 @@ +package runtime + +import ( + "encoding/json" + "sort" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/micro" + + "github.com/novox/mesh-tools/node-tools/internal/bus" + mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" +) + +// novox/hq ADR 0197: what a runtime serves, it announces — the NATS services protocol's discovery, +// read here with NATS's own types, each endpoint a subject actually served — and the announcement +// follows a re-issued membership. +func TestTheRuntimeAnnouncesWhatItServesInTheServicesProtocol(t *testing.T) { + mesh := mt.New(t) + mesh.Issue(t, mt.MembershipOf("alpha", "anchor", false, nil)) + mesh.Issue(t, mt.MembershipOf("beta", "anchor", false, map[string][]string{"node-shelf": {"list", "clear"}})) + nodeTools := connect(t, "node-tools", "anchor") + stop, err := Run(nodeTools, []Served{ + {"alpha", []string{mt.Fixture("many-alpha.serve.mjs")}}, + {"beta", []string{mt.Fixture("many-beta.serve.mjs")}}, + }, nil, (&mt.Logs{}).Logf) + if err != nil { + t.Fatal(err) + } + defer stop() + + nc, err := nats.Connect(mt.URL(t)) + if err != nil { + t.Fatal(err) + } + defer nc.Close() + info := func() micro.Info { + t.Helper() + msg, err := nc.Request("$SRV.INFO.node-tools.anchor", nil, 2*time.Second) + if err != nil { + t.Fatal(err) + } + var i micro.Info + if err := json.Unmarshal(msg.Data, &i); err != nil { + t.Fatal(err) + } + return i + } + subjects := func(i micro.Info) string { + var out []string + for _, e := range i.Endpoints { + out = append(out, e.Subject+"|"+e.QueueGroup+"|"+e.Metadata["kind"]+"|"+e.Metadata["module"]+"|"+e.Metadata["seat"]+"|"+e.Metadata["scope"]+"|"+e.Metadata["interchangeable"]) + } + sort.Strings(out) + return strings.Join(out, "\n") + } + + got := info() + if got.Type != micro.InfoResponseType || got.Name != "node-tools" || got.ID != "anchor" || got.Version == "" { + t.Errorf("identity: %+v", got.ServiceIdentity) + } + want := strings.Join([]string{ + "mesh.mod.alpha.tool.one.anchor||tool|alpha|||false", + "mesh.mod.alpha.tool.two.anchor||tool|alpha|||false", + "mesh.mod.beta.tool.five.anchor||tool|beta|||false", + "mesh.mod.beta.tool.four.anchor||tool|beta|||false", + "mesh.mod.beta.tool.three.anchor||tool|beta|||false", + "mesh.seat.node-shelf.tool.clear.anchor||seat|beta|node-shelf|node|false", + "mesh.seat.node-shelf.tool.list.anchor||seat|beta|node-shelf|node|false", + }, "\n") + if s := subjects(got); s != want { + t.Errorf("announced:\n%s\nwant:\n%s", s, want) + } + for _, e := range got.Endpoints { + if e.Metadata["node"] != "anchor" || e.Metadata["description"] == "" || !json.Valid([]byte(e.Metadata["schema"])) { + t.Errorf("endpoint metadata: %+v", e) + } + } + + // Every subject announced is answered. + asker := connect(t, "console", "workstation") + for _, e := range got.Endpoints { + if _, err := asker.Ask("", map[string]any{}, e.Subject); err != nil { + t.Errorf("announced %s and does not answer it: %v", e.Subject, err) + } + } + + // PING answers with the same identity, and a request for another service is not answered. + if msg, err := nc.Request("$SRV.PING", nil, 2*time.Second); err != nil || !strings.Contains(string(msg.Data), micro.PingResponseType) { + t.Errorf("ping: %v %v", msg, err) + } + if _, err := nc.Request("$SRV.INFO.somebody-else", nil, 300*time.Millisecond); err == nil { + t.Error("answered a request for another service") + } + + // alpha is re-issued a plain subject: the announcement says so, and says it is interchangeable. + mesh.Issue(t, mt.MembershipOf("alpha", "anchor", true, nil)) + var after string + for i := 0; i < 40; i++ { + after = subjects(info()) + if strings.Contains(after, "mesh.mod.alpha.tool.one|serve.alpha|tool|alpha|||true") { + break + } + time.Sleep(100 * time.Millisecond) + } + if !strings.Contains(after, "mesh.mod.alpha.tool.one|serve.alpha|tool|alpha|||true") || + !strings.Contains(after, "mesh.mod.alpha.tool.one.anchor||tool|alpha|||true") { + t.Errorf("after a re-issued membership:\n%s", after) + } + _ = bus.Served{} +} diff --git a/node-tools/internal/runtime/runtime.go b/node-tools/internal/runtime/runtime.go index 674d2df..9b616e9 100644 --- a/node-tools/internal/runtime/runtime.go +++ b/node-tools/internal/runtime/runtime.go @@ -14,6 +14,7 @@ import ( "strings" "sync" + "github.com/novox/mesh-tools/node-tools/internal/announce" "github.com/novox/mesh-tools/node-tools/internal/bus" "github.com/novox/mesh-tools/node-tools/internal/launch" ) @@ -276,10 +277,83 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log } logf("%s", line) stops = append(stops, serveSeats(conn, modules, registrations, logf)) + + // **What it serves, it announces** (novox/hq ADR 0197): the NATS services protocol's discovery, + // answered with what is served at the moment it is asked — re-served memberships included. + announced, err := announce.Serve(conn, announce.Service{ + Name: conn.Module(), ID: instanceOf(conn), + Description: "the mesh's tool runtime on " + node + ": every assigned module's tools and the seats they hold", + Metadata: map[string]string{"node": node}, + }, func() []announce.Endpoint { return endpointsOf(conn, own, registrations, modules) }) + if err != nil { + logf("[mesh-tools] cannot announce what it serves: %v", err) + } else { + stops = append(stops, announced) + } conn.Flush() return stopAll, nil } +// instanceOf is this runtime's instance on the bus: its machine, which is what tells two instances of +// one service apart; the connection's module where it has no machine. +func instanceOf(conn *bus.Conn) string { + if n := conn.Node(); n != "" { + return n + } + return conn.Module() +} + +// endpointsOf is everything this runtime serves now: each served module's tools on every subject the +// mesh issued for them, and each held seat's verbs on the seat's subject (ADR 0197). +func endpointsOf(conn *bus.Conn, own []registration, registrations []registration, modules []string) []announce.Endpoint { + node := conn.Node() + var out []announce.Endpoint + for _, r := range own { + for _, t := range r.tools { + served := conn.ServedOn(r.module, t.Name) + interchangeable := false + for _, s := range served { + interchangeable = interchangeable || s.Subject == "mesh.mod."+r.module+".tool."+t.Name + } + for _, s := range served { + out = append(out, announce.Endpoint{Kind: announce.KindTool, Module: r.module, Tool: t.Name, + Node: node, Description: t.Description, Schema: t.Input, Interchangeable: interchangeable, + Subject: s.Subject, Queue: s.Queue}) + } + } + } + impl := map[string]map[string]launch.Tool{} + for _, r := range registrations { + if impl[r.module] == nil { + impl[r.module] = map[string]launch.Tool{} + } + for _, t := range r.tools { + impl[r.module][t.Name] = t + } + } + have := map[string]bool{} + for _, module := range modules { + m := conn.Membership(module) + if m == nil { + continue + } + for _, v := range m.Seats { + t, ok := impl[v.Seat][v.Verb] + if !ok || have[v.Subject] { + continue + } + have[v.Subject] = true + scope := "mesh" + if node != "" && strings.HasSuffix(v.Subject, "."+node) { + scope = "node" + } + out = append(out, announce.Endpoint{Kind: announce.KindSeat, Module: module, Tool: v.Verb, Seat: v.Seat, + Scope: scope, Node: node, Description: t.Description, Schema: t.Input, Subject: v.Subject}) + } + } + return out +} + func orNone(names []string) string { if len(names) == 0 { return "(none)"