Merge pull request 'Console: wait for every runtime that answered, name the ones it missed, add mesh_runtimes' (#13) from fix/console-sees-every-runtime into main
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
This commit was merged in pull request #13.
This commit is contained in:
@@ -61,10 +61,22 @@ The same package carries the client (novox/hq design 25 §7, design 34): `mesh t
|
||||
refuses to bind anything but loopback. With `--console <url>`, `tools` and `call` go through a console
|
||||
already on the machine and need no credential.
|
||||
|
||||
Discovery asks the modules: every runtime answers a `tools` verb for each module it serves, with names,
|
||||
descriptions and schemas from the code that answers them, and the console asks the catalogue which
|
||||
modules the mesh holds and each module what it serves. A module that does not answer is named, never
|
||||
dropped. A module may not name a tool of its own `tools`; the runtime refuses it at load.
|
||||
The console announces six tools and reaches everything else by address (novox/hq ADR 0195):
|
||||
`mesh_overview` (the mesh's seats and machines), `mesh_machine` (one machine's seats and modules),
|
||||
`mesh_search` (a tool by words), `mesh_describe` (one tool's arguments), `mesh_call` (call one by
|
||||
address) and `mesh_runtimes` (which runtimes answered discovery: per runtime its machine, how long its
|
||||
answer took, its size in bytes, how many modules and tools it announced, whether it was shortened to
|
||||
fit the bus, when it was last heard — and who was expected and not heard).
|
||||
|
||||
Discovery asks the bus (ADR 0197): every runtime answers the NATS services protocol's `$SRV.PING` and
|
||||
`$SRV.INFO` with what it serves at that moment, and the console reads where the controller's records
|
||||
place each module. The console waits at least 750 ms, and up to 5 s for every runtime that answered
|
||||
PING to send what it serves — a large runtime's answer crossing the bus to a broker on another machine
|
||||
can arrive well after the rest. A runtime that said it is there and did not say what it serves, or
|
||||
that answered earlier and not now, is named: the console never calls a module missing, or on another
|
||||
machine, while a runtime that might serve it was not heard. A runtime whose answer outgrows the bus
|
||||
announces first-line descriptions, then none, and says so. A module may not name a tool of its own
|
||||
`tools`; the runtime refuses it at load.
|
||||
|
||||
## Verified
|
||||
|
||||
|
||||
@@ -22,10 +22,22 @@ import (
|
||||
// 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.
|
||||
// Window is how long the console gathers discovery answers at the least: every instance answers one
|
||||
// request, and how many will is what is being found out.
|
||||
var Window = 750 * time.Millisecond
|
||||
|
||||
// Patience is how long the console waits, at the most, for the full answer of a service that has
|
||||
// said it is there. **A fixed window lost the largest runtime** (2026-10-05): the laptop's answer —
|
||||
// 48 modules, 341 endpoints, 164 kB — went up to the broker on another machine and back, arrived
|
||||
// last every time (a median of 365 ms, one in 25 after 813 ms on a quiet link, later still while the
|
||||
// runtime re-served), and when it missed the 750 ms window the console said its modules ran nowhere.
|
||||
// A PING answer is a hundred bytes and arrives at once, so who is there is known early; what each
|
||||
// serves is waited for until it arrives or this passes, and a service that never sends it is named.
|
||||
var Patience = 5 * time.Second
|
||||
|
||||
// Shortened is the service metadata key a runtime sets when its announcement was cut to fit the bus.
|
||||
const Shortened = "shortened"
|
||||
|
||||
// Kinds of endpoint.
|
||||
const (
|
||||
KindTool = "tool" // a module's own tool
|
||||
@@ -163,17 +175,10 @@ func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error)
|
||||
return nil
|
||||
}
|
||||
// Too large for the bus is no answer at all, and was silent: shorten every description to its
|
||||
// first line and say so, rather than vanish.
|
||||
// first line and say so, rather than vanish; still too large, leave descriptions out.
|
||||
if limit := conn.MaxPayload(); limit > 0 && int64(len(body)) > limit {
|
||||
if info, ok := v.(micro.Info); ok {
|
||||
for i := range info.Endpoints {
|
||||
info.Endpoints[i].Metadata["description"] = firstLine(info.Endpoints[i].Metadata["description"])
|
||||
}
|
||||
if shorter, err := json.Marshal(info); err == nil {
|
||||
conn.Logf("[mesh-tools] what this runtime serves is %d bytes, beyond the bus's %d; announced with each description cut to its first line (%d bytes)",
|
||||
len(body), limit, len(shorter))
|
||||
body = shorter
|
||||
}
|
||||
body = shorten(conn.Logf, info, body, limit)
|
||||
}
|
||||
}
|
||||
return body
|
||||
@@ -200,22 +205,149 @@ func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error)
|
||||
}, 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
|
||||
// shorten is an info_response cut to fit the bus: descriptions to their first line, then none, the
|
||||
// service's metadata saying which (so discovery can tell a person what it was given).
|
||||
func shorten(logf func(string, ...any), info micro.Info, body []byte, limit int64) []byte {
|
||||
md := map[string]string{}
|
||||
for k, v := range info.Metadata {
|
||||
md[k] = v
|
||||
}
|
||||
var out []micro.Info
|
||||
for _, b := range raw {
|
||||
var i micro.Info
|
||||
if json.Unmarshal(b, &i) != nil || i.Type != micro.InfoResponseType {
|
||||
continue
|
||||
info.Metadata = md
|
||||
for _, step := range []string{"descriptions cut to their first line", "descriptions left out"} {
|
||||
for i := range info.Endpoints {
|
||||
e := info.Endpoints[i].Metadata
|
||||
if step == "descriptions left out" {
|
||||
delete(e, "description")
|
||||
} else {
|
||||
e["description"] = firstLine(e["description"])
|
||||
}
|
||||
}
|
||||
md[Shortened] = step
|
||||
shorter, err := json.Marshal(info)
|
||||
if err != nil {
|
||||
return body
|
||||
}
|
||||
logf("[mesh-tools] what this runtime serves is %d bytes, beyond the bus's %d; announced with %s (%d bytes)",
|
||||
len(body), limit, step, len(shorter))
|
||||
if int64(len(shorter)) <= limit {
|
||||
return shorter
|
||||
}
|
||||
body = shorter
|
||||
}
|
||||
return body
|
||||
}
|
||||
|
||||
// Heard is one service instance's answer to discovery: what it serves, how long it took to arrive,
|
||||
// and how large it was.
|
||||
type Heard struct {
|
||||
Info micro.Info
|
||||
Took time.Duration
|
||||
Bytes int
|
||||
Shortened string // why its descriptions are not whole, or ""
|
||||
}
|
||||
|
||||
// Machine is the machine an instance runs on: its metadata's, or the one its endpoints name.
|
||||
func (h Heard) Machine() string {
|
||||
if n := h.Info.Metadata["node"]; n != "" {
|
||||
return n
|
||||
}
|
||||
for _, e := range h.Info.Endpoints {
|
||||
if n := e.Metadata["node"]; n != "" {
|
||||
return n
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// Instance names one service instance.
|
||||
type Instance struct {
|
||||
Name string
|
||||
ID string
|
||||
Machine string
|
||||
}
|
||||
|
||||
// Discovery is what one round of discovery heard, and who said it was there but did not say what it
|
||||
// serves in time.
|
||||
type Discovery struct {
|
||||
At time.Time
|
||||
Waited time.Duration
|
||||
Heard []Heard
|
||||
Silent []Instance
|
||||
}
|
||||
|
||||
// Gather asks every service on the bus who it is and what it serves, and waits at least the window
|
||||
// and at most the patience: until every instance that answered PING has answered INFO. An answer
|
||||
// that is not the protocol's is skipped.
|
||||
func Gather(conn *bus.Conn) (Discovery, error) {
|
||||
start := time.Now()
|
||||
d := Discovery{At: start.UTC()}
|
||||
pings, stopPings, err := conn.Collect("$SRV.PING", nil)
|
||||
if err != nil {
|
||||
return d, err
|
||||
}
|
||||
defer stopPings()
|
||||
infos, stopInfos, err := conn.Collect("$SRV.INFO", nil)
|
||||
if err != nil {
|
||||
return d, err
|
||||
}
|
||||
defer stopInfos()
|
||||
|
||||
type key struct{ name, id string }
|
||||
pinged := map[key]int{}
|
||||
answered := map[key]int{}
|
||||
where := map[key]string{}
|
||||
var order []key
|
||||
complete := func() bool {
|
||||
for k, n := range pinged {
|
||||
if answered[k] < n {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
least := time.NewTimer(Window)
|
||||
defer least.Stop()
|
||||
most := time.NewTimer(Patience)
|
||||
defer most.Stop()
|
||||
windowOver := false
|
||||
for {
|
||||
select {
|
||||
case a := <-pings:
|
||||
var p micro.Ping
|
||||
if json.Unmarshal(a.Data, &p) != nil || p.Type != micro.PingResponseType {
|
||||
continue
|
||||
}
|
||||
k := key{p.Name, p.ID}
|
||||
if _, seen := pinged[k]; !seen {
|
||||
order = append(order, k)
|
||||
}
|
||||
pinged[k]++
|
||||
if where[k] == "" {
|
||||
where[k] = p.Metadata["node"]
|
||||
}
|
||||
case a := <-infos:
|
||||
var i micro.Info
|
||||
if json.Unmarshal(a.Data, &i) != nil || i.Type != micro.InfoResponseType {
|
||||
continue
|
||||
}
|
||||
answered[key{i.Name, i.ID}]++
|
||||
d.Heard = append(d.Heard, Heard{Info: i, Took: a.At.Sub(start), Bytes: len(a.Data), Shortened: i.Metadata[Shortened]})
|
||||
case <-least.C:
|
||||
windowOver = true
|
||||
case <-most.C:
|
||||
d.Waited = time.Since(start)
|
||||
for _, k := range order {
|
||||
if answered[k] < pinged[k] {
|
||||
d.Silent = append(d.Silent, Instance{Name: k.name, ID: k.id, Machine: where[k]})
|
||||
}
|
||||
}
|
||||
return d, nil
|
||||
}
|
||||
if windowOver && complete() {
|
||||
d.Waited = time.Since(start)
|
||||
return d, nil
|
||||
}
|
||||
out = append(out, i)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// Endpoints reads an info_response's endpoints back into what they announce.
|
||||
|
||||
@@ -0,0 +1,156 @@
|
||||
package announce
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"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"
|
||||
)
|
||||
|
||||
func connect(t *testing.T, module, node string) *bus.Conn {
|
||||
t.Helper()
|
||||
c, err := bus.Connect(bus.Credential{URL: mt.URL(t), Module: module, Node: node})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
c.Logf = func(string, ...any) {}
|
||||
t.Cleanup(c.Close)
|
||||
return c
|
||||
}
|
||||
|
||||
// lagging answers PING at once and INFO after `lag` — or never, with a negative lag: a runtime whose
|
||||
// one large answer is slow to cross the bus, as the laptop's was (2026-10-05).
|
||||
func lagging(t *testing.T, conn *bus.Conn, s Service, lag time.Duration, endpoints []Endpoint) {
|
||||
t.Helper()
|
||||
ping, _ := json.Marshal(micro.Ping{ServiceIdentity: s.identity(), Type: micro.PingResponseType})
|
||||
stop, err := conn.Raw("$SRV.PING", func(string, []byte) []byte { return ping })
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
stop, err = conn.Raw("$SRV.INFO", func(string, []byte) []byte {
|
||||
if lag < 0 {
|
||||
return nil
|
||||
}
|
||||
time.Sleep(lag)
|
||||
body, _ := json.Marshal(Info(s, endpoints))
|
||||
return body
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
conn.Flush()
|
||||
}
|
||||
|
||||
// A runtime whose answer arrives after the window is waited for, because it said it was there; one
|
||||
// that never says what it serves is named, not dropped. With the fixed 750 ms window the laptop's
|
||||
// modules were said to run nowhere whenever its answer was late (2026-10-05).
|
||||
func TestDiscoveryWaitsForEveryRuntimeThatSaidItIsThere(t *testing.T) {
|
||||
was := Patience
|
||||
Patience = 2500 * time.Millisecond
|
||||
t.Cleanup(func() { Patience = was })
|
||||
|
||||
fast := connect(t, "node-tools", "desk")
|
||||
stopFast, err := Serve(fast, Service{Name: "node-tools", ID: "desk", Metadata: map[string]string{"node": "desk"}}, func() []Endpoint {
|
||||
return []Endpoint{{Kind: KindTool, Module: "alpha", Tool: "one", Node: "desk", Subject: "mesh.mod.alpha.tool.one.desk"}}
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stopFast)
|
||||
fast.Flush()
|
||||
slow := connect(t, "node-tools", "laptop")
|
||||
lagging(t, slow, Service{Name: "node-tools", ID: "laptop", Metadata: map[string]string{"node": "laptop"}}, Window+400*time.Millisecond,
|
||||
[]Endpoint{{Kind: KindTool, Module: "slack", Tool: "slack_check", Node: "laptop", Subject: "mesh.mod.slack.tool.slack_check.laptop"}})
|
||||
mute := connect(t, "node-tools", "sleeper")
|
||||
lagging(t, mute, Service{Name: "node-tools", ID: "sleeper", Metadata: map[string]string{"node": "sleeper"}}, -1, nil)
|
||||
|
||||
asker := connect(t, "console", "desk")
|
||||
start := time.Now()
|
||||
d, err := Gather(asker)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
took := time.Since(start)
|
||||
|
||||
heard := map[string]Heard{}
|
||||
for _, h := range d.Heard {
|
||||
heard[h.Info.ID] = h
|
||||
}
|
||||
if h, ok := heard["laptop"]; !ok {
|
||||
t.Fatalf("the slow runtime was not heard: %+v", d.Heard)
|
||||
} else if h.Took <= Window || h.Machine() != "laptop" || h.Bytes == 0 {
|
||||
t.Errorf("the slow runtime's answer: took %s (window %s), machine %q, %d bytes", h.Took, Window, h.Machine(), h.Bytes)
|
||||
}
|
||||
if _, ok := heard["desk"]; !ok {
|
||||
t.Errorf("the fast runtime was not heard: %+v", d.Heard)
|
||||
}
|
||||
if len(d.Silent) != 1 || d.Silent[0].ID != "sleeper" || d.Silent[0].Machine != "sleeper" {
|
||||
t.Errorf("the runtime that never said what it serves is not named: %+v", d.Silent)
|
||||
}
|
||||
if took < Patience || took > Patience+time.Second {
|
||||
t.Errorf("waited %s with a silent runtime; patience is %s", took, Patience)
|
||||
}
|
||||
}
|
||||
|
||||
// With nobody slow, discovery takes the window and no longer.
|
||||
func TestDiscoveryDoesNotWaitWhenEveryoneAnswered(t *testing.T) {
|
||||
fast := connect(t, "node-tools", "desk")
|
||||
stop, err := Serve(fast, Service{Name: "node-tools", ID: "desk"}, func() []Endpoint { return nil })
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
fast.Flush()
|
||||
start := time.Now()
|
||||
d, err := Gather(connect(t, "console", "desk"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if took := time.Since(start); took > Window+500*time.Millisecond || len(d.Silent) != 0 || len(d.Heard) == 0 {
|
||||
t.Errorf("took %s, heard %d, silent %+v", took, len(d.Heard), d.Silent)
|
||||
}
|
||||
}
|
||||
|
||||
// Too large even with first lines, descriptions are left out — and the answer says which, without
|
||||
// touching the service's own metadata.
|
||||
func TestAnAnswerTooLargeIsShortenedUntilItFits(t *testing.T) {
|
||||
long := strings.Repeat("A sentence that is long. ", 20)
|
||||
var endpoints []Endpoint
|
||||
for i := 0; i < 50; i++ {
|
||||
endpoints = append(endpoints, Endpoint{Kind: KindTool, Module: "m", Tool: "t" + strings.Repeat("x", i), Node: "laptop",
|
||||
Description: long, Subject: "mesh.mod.m.tool.t"})
|
||||
}
|
||||
own := map[string]string{"node": "laptop"}
|
||||
info := Info(Service{Name: "node-tools", ID: "laptop", Metadata: own}, endpoints)
|
||||
full, _ := json.Marshal(info)
|
||||
firstLines := 0
|
||||
for range endpoints {
|
||||
firstLines += len(firstLine(long))
|
||||
}
|
||||
limit := int64(len(full) - len(long)*len(endpoints) + firstLines/2) // first lines alone do not fit
|
||||
var said []string
|
||||
got := shorten(func(f string, a ...any) { said = append(said, f) }, info, full, limit)
|
||||
if int64(len(got)) > limit {
|
||||
t.Fatalf("still %d bytes, beyond %d", len(got), limit)
|
||||
}
|
||||
var back micro.Info
|
||||
if err := json.Unmarshal(got, &back); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if back.Metadata[Shortened] != "descriptions left out" || back.Metadata["node"] != "laptop" {
|
||||
t.Errorf("metadata: %v", back.Metadata)
|
||||
}
|
||||
if _, has := own[Shortened]; has {
|
||||
t.Error("the service's own metadata was changed")
|
||||
}
|
||||
if len(said) != 2 {
|
||||
t.Errorf("each step is said: %v", said)
|
||||
}
|
||||
}
|
||||
@@ -649,37 +649,51 @@ func (c *Conn) Raw(subject string, answer func(subject string, data []byte) []by
|
||||
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) {
|
||||
// Arrival is one answer to a request, and when it arrived.
|
||||
type Arrival struct {
|
||||
Data []byte
|
||||
At time.Time
|
||||
}
|
||||
|
||||
// Collect publishes one request and hands every answer to the channel as it arrives, until stopped:
|
||||
// for a caller that decides for itself when it has heard enough — discovery waits for every service
|
||||
// that said it is there, not for a fixed moment (2026-10-05). Answers beyond what the caller has
|
||||
// taken are kept up to the channel's room; an empty answer is no answer.
|
||||
func (c *Conn) Collect(subject string, body []byte) (<-chan Arrival, func(), error) {
|
||||
inbox := c.nc.NewRespInbox()
|
||||
sub, err := c.nc.SubscribeSync(inbox)
|
||||
out := make(chan Arrival, 1024)
|
||||
var mu sync.Mutex
|
||||
stopped := false
|
||||
sub, err := c.nc.Subscribe(inbox, func(msg *nats.Msg) {
|
||||
if len(msg.Data) == 0 {
|
||||
return
|
||||
}
|
||||
a := Arrival{Data: msg.Data, At: time.Now()}
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if stopped {
|
||||
return
|
||||
}
|
||||
select {
|
||||
case out <- a:
|
||||
default:
|
||||
c.Logf("[mesh-tools] more answers to %s than are kept; one (%d bytes) dropped", subject, len(msg.Data))
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return nil, func() {}, err
|
||||
}
|
||||
stop := func() {
|
||||
_ = sub.Unsubscribe()
|
||||
mu.Lock()
|
||||
stopped = true
|
||||
mu.Unlock()
|
||||
}
|
||||
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)
|
||||
}
|
||||
stop()
|
||||
return nil, func() {}, err
|
||||
}
|
||||
return out, stop, nil
|
||||
}
|
||||
|
||||
// MaxPayload is the largest message the bus carries, as the server told this connection.
|
||||
|
||||
@@ -2,8 +2,8 @@ package console
|
||||
|
||||
// The mesh's tools found by address, not announced whole (novox/hq ADR 0195, to-be 34 §3a).
|
||||
//
|
||||
// The console announces five tools. Everything the mesh answers is reached through them by one
|
||||
// address per layer:
|
||||
// The console announces six tools — five to find and call, and mesh_runtimes to say which runtimes
|
||||
// discovery heard. Everything the mesh answers is reached through them by one address per layer:
|
||||
//
|
||||
// <seat>.<verb> a seat held once for the mesh — its holder answers
|
||||
// <node>/<seat>.<verb> a seat held once per machine — that machine's holder answers
|
||||
@@ -32,13 +32,14 @@ import (
|
||||
// one agent turn's search, describe and call ask once, short enough that nothing goes stale.
|
||||
var IndexKept = 5 * time.Second
|
||||
|
||||
// The five tools the console announces. Names of the API's kind — letters, digits, `_`, `-`.
|
||||
// The six tools the console announces. Names of the API's kind — letters, digits, `_`, `-`.
|
||||
const (
|
||||
verbOverview = "mesh_overview"
|
||||
verbMachine = "mesh_machine"
|
||||
verbSearch = "mesh_search"
|
||||
verbDescribe = "mesh_describe"
|
||||
verbCall = "mesh_call"
|
||||
verbRuntimes = "mesh_runtimes"
|
||||
)
|
||||
|
||||
// searchCap is how many matches a search answers before it says how many more there were.
|
||||
@@ -49,7 +50,7 @@ const grammar = "Addresses: `<seat>.<verb>` for a seat held once for the mesh (e
|
||||
"`<node>/<module>.<tool>` for a module on one machine (e.g. `novox/postgres.postgres_list_databases`); " +
|
||||
"and `<module>.<tool>` also for a module whose instances are interchangeable."
|
||||
|
||||
// discovery is the five tools as tools/list announces them.
|
||||
// discovery is the six tools as tools/list announces them.
|
||||
func discovery() []map[string]any {
|
||||
str := func(desc string) map[string]any { return map[string]any{"type": "string", "description": desc} }
|
||||
obj := func(props map[string]any, required ...string) map[string]any {
|
||||
@@ -77,12 +78,17 @@ func discovery() []map[string]any {
|
||||
"arguments": map[string]any{"type": "object", "description": "the tool's arguments, as mesh_describe gives its schema"},
|
||||
}, "address"),
|
||||
"description": "Call one tool by its address with its arguments; the answer says which machine gave it. " + grammar},
|
||||
{"name": verbRuntimes, "inputSchema": obj(map[string]any{}),
|
||||
"description": "Which runtimes answered discovery, asked now: for each, its machine, how long its answer took to " +
|
||||
"arrive, its size in bytes, how many modules and tools it announced, whether it was shortened to fit the " +
|
||||
"bus, and when it was last heard — and every runtime or machine expected that was not heard. Read-only; " +
|
||||
"for when an address is said to be missing or on another machine."},
|
||||
}
|
||||
}
|
||||
|
||||
func isDiscovery(name string) bool {
|
||||
switch name {
|
||||
case verbOverview, verbMachine, verbSearch, verbDescribe, verbCall:
|
||||
case verbOverview, verbMachine, verbSearch, verbDescribe, verbCall, verbRuntimes:
|
||||
return true
|
||||
}
|
||||
return false
|
||||
@@ -117,6 +123,44 @@ type index struct {
|
||||
Modules map[string]*moduleInfo
|
||||
NotAnswering []string
|
||||
Listing *Listing // the flat catalogue the call path resolves subjects with
|
||||
Discovery announce.Discovery
|
||||
// Recorded is where the controller's records place each module that declares tools.
|
||||
Recorded map[string][]string
|
||||
// Unheard is every runtime whose answer this index lacks: it said it was there and did not say
|
||||
// what it serves in time, or it answered before and not now. An address it might answer is never
|
||||
// called missing while it is here (2026-10-05).
|
||||
Unheard []unheard
|
||||
}
|
||||
|
||||
// unheard is one runtime discovery did not hear in full, and why.
|
||||
type unheard struct {
|
||||
Runtime string `json:"runtime"`
|
||||
Machine string `json:"machine,omitempty"`
|
||||
Why string `json:"why"`
|
||||
}
|
||||
|
||||
// unheardOn is what is unheard on one machine; with no machine, everything unheard.
|
||||
func (x *index) unheardOn(node string) []unheard {
|
||||
var out []unheard
|
||||
for _, u := range x.Unheard {
|
||||
if node == "" || u.Machine == node || u.Machine == "" {
|
||||
out = append(out, u)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// sayUnheard is the unheard, as a sentence's end.
|
||||
func sayUnheard(us []unheard) string {
|
||||
var parts []string
|
||||
for _, u := range us {
|
||||
who := u.Runtime
|
||||
if u.Machine != "" {
|
||||
who += " on " + u.Machine
|
||||
}
|
||||
parts = append(parts, who+" ("+u.Why+")")
|
||||
}
|
||||
return strings.Join(parts, "; ")
|
||||
}
|
||||
|
||||
func (x *index) seat(name string) *seatInfo {
|
||||
@@ -193,20 +237,24 @@ func indexOn(conn *bus.Conn) (*index, error) {
|
||||
wg.Add(2)
|
||||
go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }()
|
||||
go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }()
|
||||
infos, err := announce.Gather(conn)
|
||||
d, err := announce.Gather(conn)
|
||||
wg.Wait()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
l := &Listing{Tools: []Tool{}, NotAnswering: []string{}}
|
||||
x := &index{Modules: map[string]*moduleInfo{}, Listing: l}
|
||||
x := &index{Modules: map[string]*moduleInfo{}, Listing: l, Discovery: d, Recorded: map[string][]string{}}
|
||||
for _, s := range d.Silent {
|
||||
x.Unheard = append(x.Unheard, unheard{Runtime: s.Name, Machine: s.Machine,
|
||||
Why: fmt.Sprintf("it answered PING; what it serves did not arrive within %s", announce.Patience)})
|
||||
}
|
||||
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) {
|
||||
for _, h := range d.Heard {
|
||||
for _, e := range announce.Endpoints(h.Info) {
|
||||
if e.Node != "" {
|
||||
machines[e.Node] = true
|
||||
}
|
||||
@@ -298,6 +346,7 @@ func indexOn(conn *bus.Conn) (*index, error) {
|
||||
if !m.Tools {
|
||||
continue
|
||||
}
|
||||
x.Recorded[m.Module] = append([]string{}, m.On...)
|
||||
for _, n := range m.On {
|
||||
// A holder of a seat held once for the mesh announces no machine — the seat is the mesh's,
|
||||
// not a machine's — so what it announced without one answers for wherever it is assigned.
|
||||
@@ -309,6 +358,9 @@ func indexOn(conn *bus.Conn) (*index, error) {
|
||||
} else {
|
||||
l.NotAnswering = append(l.NotAnswering, "mesh-controller (its records of the modules did not answer, so what is missing cannot be said)")
|
||||
}
|
||||
for _, u := range x.Unheard {
|
||||
l.NotAnswering = append(l.NotAnswering, "the runtime "+sayUnheard([]unheard{u}))
|
||||
}
|
||||
sort.Strings(l.NotAnswering)
|
||||
x.NotAnswering = l.NotAnswering
|
||||
|
||||
@@ -336,9 +388,12 @@ func containsHolder(hs []holder, h holder) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func (s *Surface) index() (*index, error) {
|
||||
func (s *Surface) index() (*index, error) { return s.indexAsked(false) }
|
||||
|
||||
// indexAsked is what the mesh answered, asked again when `fresh` or when the kept answer is old.
|
||||
func (s *Surface) indexAsked(fresh bool) (*index, error) {
|
||||
s.mu.Lock()
|
||||
if s.idx != nil && time.Since(s.idxAt) <= IndexKept {
|
||||
if !fresh && s.idx != nil && time.Since(s.idxAt) <= IndexKept {
|
||||
x := s.idx
|
||||
s.mu.Unlock()
|
||||
return x, nil
|
||||
@@ -349,6 +404,7 @@ func (s *Surface) index() (*index, error) {
|
||||
return nil, err
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.remember(x)
|
||||
s.idx, s.idxAt = x, time.Now()
|
||||
s.mu.Unlock()
|
||||
return x, nil
|
||||
@@ -397,6 +453,18 @@ func resolve(x *index, address string) (target, error) {
|
||||
|
||||
m := x.Modules[prefix]
|
||||
if m == nil {
|
||||
if on := x.Recorded[prefix]; len(on) > 0 {
|
||||
why := "it did not announce itself on the bus"
|
||||
if us := x.unheardOnAny(on); len(us) > 0 {
|
||||
why = "discovery did not hear in full from " + sayUnheard(us)
|
||||
}
|
||||
return target{}, fmt.Errorf("%s is assigned to %s in the mesh's records, but nothing that answered discovery serves it: %s. "+
|
||||
"Ask again in a moment; mesh_runtimes shows which runtimes answered", prefix, strings.Join(on, ", "), why)
|
||||
}
|
||||
if len(x.Unheard) > 0 {
|
||||
return target{}, fmt.Errorf("nothing that answered discovery is called %s, but not every runtime answered: %s — "+
|
||||
"%s may be theirs. Ask again in a moment; mesh_runtimes shows which runtimes answered", prefix, sayUnheard(x.Unheard), prefix)
|
||||
}
|
||||
return target{}, fmt.Errorf("nothing in the mesh is called %s: no seat, and no module that answers tools. "+
|
||||
"mesh_search finds a tool by words", prefix)
|
||||
}
|
||||
@@ -412,11 +480,30 @@ func resolve(x *index, address string) (target, error) {
|
||||
return target{Address: rest, Key: rest, Name: rest, Tool: *tool}, nil
|
||||
}
|
||||
if len(m.On) > 0 && !contains(m.On, node) {
|
||||
if us := x.unheardOn(node); len(us) > 0 || contains(x.Recorded[prefix], node) {
|
||||
why := "it did not announce itself there"
|
||||
if len(us) > 0 {
|
||||
why = "discovery did not hear in full from " + sayUnheard(us)
|
||||
}
|
||||
return target{}, fmt.Errorf("%s on %s did not answer discovery: %s. It was heard on %s. "+
|
||||
"Ask again in a moment; mesh_runtimes shows which runtimes answered", prefix, node, why, orNobody(m.On))
|
||||
}
|
||||
return target{}, fmt.Errorf("%s does not run on %s; it runs on %s", prefix, node, orNobody(m.On))
|
||||
}
|
||||
return target{Address: node + "/" + rest, Key: rest + "@" + node, Name: rest, Node: node, Tool: *tool}, nil
|
||||
}
|
||||
|
||||
// unheardOnAny is what is unheard on any of these machines.
|
||||
func (x *index) unheardOnAny(nodes []string) []unheard {
|
||||
var out []unheard
|
||||
for _, u := range x.Unheard {
|
||||
if u.Machine == "" || contains(nodes, u.Machine) {
|
||||
out = append(out, u)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func contains(xs []string, s string) bool {
|
||||
for _, x := range xs {
|
||||
if x == s {
|
||||
@@ -511,8 +598,11 @@ func failure(text string) map[string]any {
|
||||
return map[string]any{"content": []map[string]any{{"type": "text", "text": text}}, "isError": true}
|
||||
}
|
||||
|
||||
// discover answers one of the five.
|
||||
// discover answers one of the discovery verbs.
|
||||
func (s *Surface) discover(name string, args map[string]any) map[string]any {
|
||||
if name == verbRuntimes {
|
||||
return s.runtimesAnswer()
|
||||
}
|
||||
x, err := s.index()
|
||||
if err != nil {
|
||||
return failure("the mesh's discovery failed: " + err.Error())
|
||||
@@ -606,7 +696,18 @@ func (s *Surface) discover(name string, args map[string]any) map[string]any {
|
||||
}
|
||||
modules = append(modules, entry)
|
||||
}
|
||||
return answerText(map[string]any{"machine": node, "seats": seats, "modules": modules})
|
||||
out := map[string]any{"machine": node, "seats": seats, "modules": modules}
|
||||
if seats == nil {
|
||||
out["seats"] = []map[string]any{}
|
||||
}
|
||||
if modules == nil {
|
||||
out["modules"] = []map[string]any{}
|
||||
}
|
||||
if us := x.unheardOn(node); len(us) > 0 {
|
||||
out["incomplete"] = "discovery did not hear in full from " + sayUnheard(us) +
|
||||
"; what is listed may be missing what it serves. Ask again in a moment; mesh_runtimes shows which runtimes answered"
|
||||
}
|
||||
return answerText(out)
|
||||
|
||||
case verbSearch:
|
||||
query := strings.ToLower(str("query"))
|
||||
@@ -684,6 +785,10 @@ func (s *Surface) discover(name string, args map[string]any) map[string]any {
|
||||
out["matches"] = []hit{}
|
||||
out["hint"] = "nothing matched every word; try fewer words, or mesh_overview and mesh_machine to browse"
|
||||
}
|
||||
if len(x.Unheard) > 0 {
|
||||
out["incomplete"] = "discovery did not hear in full from " + sayUnheard(x.Unheard) +
|
||||
"; what they serve is not searched. mesh_runtimes shows which runtimes answered"
|
||||
}
|
||||
return answerText(out)
|
||||
|
||||
case verbDescribe:
|
||||
|
||||
@@ -36,7 +36,7 @@ func call(t *testing.T, endpoint, tool string, args map[string]any) (string, boo
|
||||
return text(t, post(t, endpoint, string(body)))
|
||||
}
|
||||
|
||||
// novox/hq ADR 0195: the console announces five tools, and everything the mesh answers is reached
|
||||
// novox/hq ADR 0195: the console announces six tools, and everything the mesh answers is reached
|
||||
// through them by one address per layer.
|
||||
func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
|
||||
was := IndexKept
|
||||
@@ -113,7 +113,7 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
|
||||
defer up.Close()
|
||||
endpoint := "http://" + up.Address + "/mcp"
|
||||
|
||||
// Five tools, nothing else.
|
||||
// Six tools, nothing else.
|
||||
listed := post(t, endpoint, `{"jsonrpc":"2.0","id":2,"method":"tools/list"}`)["result"].(map[string]any)
|
||||
var names []string
|
||||
for _, x := range listed["tools"].([]any) {
|
||||
@@ -125,7 +125,7 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
if got := strings.Join(names, ","); got != "mesh_overview,mesh_machine,mesh_search,mesh_describe,mesh_call" {
|
||||
if got := strings.Join(names, ","); got != "mesh_overview,mesh_machine,mesh_search,mesh_describe,mesh_call,mesh_runtimes" {
|
||||
t.Errorf("announced %s", got)
|
||||
}
|
||||
|
||||
|
||||
@@ -50,6 +50,8 @@ type Surface struct {
|
||||
at time.Time
|
||||
idx *index
|
||||
idxAt time.Time
|
||||
// runtimes is every runtime discovery has heard, as it last heard it (mesh_runtimes).
|
||||
runtimes map[string]*runtimeSeen
|
||||
// Flat announces the whole catalogue, as the console did before ADR 0195: for a person reading
|
||||
// it or a client that wants it. Off by default; MESH_CONSOLE_FLAT=1 turns it on.
|
||||
Flat bool
|
||||
@@ -300,7 +302,8 @@ func (s *Surface) instructions() string {
|
||||
return "The tools of a Novox mesh, reached as " + s.who + ", found by address rather than listed " +
|
||||
"whole (novox/hq ADR 0195). mesh_overview shows the mesh's seats and machines; mesh_machine one " +
|
||||
"machine's seats and modules; mesh_search finds a tool by words; mesh_describe gives one tool's " +
|
||||
"arguments; mesh_call calls it. " + grammar + " What may be called was fixed when this account " +
|
||||
"arguments; mesh_call calls it; mesh_runtimes says which machines' runtimes answered discovery, how fast " +
|
||||
"and how much — ask it when an address is said to be missing or on another machine. " + grammar + " What may be called was fixed when this account " +
|
||||
"was issued, so a refusal means the account, not the tool."
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,200 @@
|
||||
package console
|
||||
|
||||
// What discovery heard from each runtime (2026-10-05): the console said a module ran nowhere when the
|
||||
// largest runtime's answer arrived after the window closed, and nothing showed that it had. Every
|
||||
// round of discovery is remembered per runtime, so a runtime heard before and not now is named, and
|
||||
// `mesh_runtimes` says who answered, how fast and how much.
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/announce"
|
||||
)
|
||||
|
||||
// RuntimeForgotten is how long a runtime heard once is expected again: within it, its silence is said
|
||||
// wherever an address it might answer is asked for; after it, it is listed only by mesh_runtimes.
|
||||
var RuntimeForgotten = 15 * time.Minute
|
||||
|
||||
// runtimeSeen is one runtime instance as discovery last heard it.
|
||||
type runtimeSeen struct {
|
||||
Name string
|
||||
ID string
|
||||
Machine string
|
||||
LastHeard time.Time
|
||||
Took time.Duration
|
||||
Bytes int
|
||||
Modules int
|
||||
Tools int
|
||||
SeatVerbs int
|
||||
Shortened string
|
||||
Answered bool // in the latest round
|
||||
Why string // why not, when it did not
|
||||
}
|
||||
|
||||
func runtimeKey(name, id string) string { return name + "/" + id }
|
||||
|
||||
// seenOf is what one answer announced, counted.
|
||||
func seenOf(h announce.Heard, at time.Time) *runtimeSeen {
|
||||
modules := map[string]bool{}
|
||||
tools := map[string]bool{}
|
||||
verbs := map[string]bool{}
|
||||
for _, e := range announce.Endpoints(h.Info) {
|
||||
switch e.Kind {
|
||||
case announce.KindTool:
|
||||
modules[e.Module] = true
|
||||
tools[e.Module+"."+e.Tool] = true
|
||||
case announce.KindSeat:
|
||||
verbs[e.Seat+"."+e.Tool] = true
|
||||
}
|
||||
}
|
||||
return &runtimeSeen{Name: h.Info.Name, ID: h.Info.ID, Machine: h.Machine(), LastHeard: at, Took: h.Took,
|
||||
Bytes: h.Bytes, Modules: len(modules), Tools: len(tools), SeatVerbs: len(verbs), Shortened: h.Shortened, Answered: true}
|
||||
}
|
||||
|
||||
// remember records a round of discovery and adds to the index every runtime heard lately that did not
|
||||
// answer this time. Called with s.mu held.
|
||||
func (s *Surface) remember(x *index) {
|
||||
if s.runtimes == nil {
|
||||
s.runtimes = map[string]*runtimeSeen{}
|
||||
}
|
||||
at := x.Discovery.At
|
||||
if at.IsZero() {
|
||||
at = time.Now().UTC()
|
||||
}
|
||||
heard := map[string]bool{}
|
||||
for _, h := range x.Discovery.Heard {
|
||||
k := runtimeKey(h.Info.Name, h.Info.ID)
|
||||
seen := seenOf(h, at)
|
||||
if heard[k] {
|
||||
// Two instances under one name and id: both counted, the slower answer's time kept.
|
||||
prev := s.runtimes[k]
|
||||
seen.Modules += prev.Modules
|
||||
seen.Tools += prev.Tools
|
||||
seen.SeatVerbs += prev.SeatVerbs
|
||||
seen.Bytes += prev.Bytes
|
||||
if prev.Took > seen.Took {
|
||||
seen.Took = prev.Took
|
||||
}
|
||||
}
|
||||
heard[k] = true
|
||||
s.runtimes[k] = seen
|
||||
}
|
||||
// A runtime restarted answers under a new instance id: the old one is gone, not silent, and saying
|
||||
// it was missed would hide nothing but cry wolf after every restart.
|
||||
now := map[string]bool{}
|
||||
for k := range heard {
|
||||
now[s.runtimes[k].Name+"@"+s.runtimes[k].Machine] = true
|
||||
}
|
||||
for k, r := range s.runtimes {
|
||||
if !heard[k] && r.Machine != "" && now[r.Name+"@"+r.Machine] {
|
||||
delete(s.runtimes, k)
|
||||
}
|
||||
}
|
||||
silent := map[string]string{}
|
||||
for _, u := range x.Discovery.Silent {
|
||||
silent[runtimeKey(u.Name, u.ID)] = fmt.Sprintf("it answered PING; what it serves did not arrive within %s", announce.Patience)
|
||||
}
|
||||
for k, r := range s.runtimes {
|
||||
if heard[k] {
|
||||
continue
|
||||
}
|
||||
r.Answered = false
|
||||
if why, ok := silent[k]; ok {
|
||||
r.Why = why
|
||||
continue // already in the index's unheard
|
||||
}
|
||||
r.Why = "it did not answer this round of discovery"
|
||||
if time.Since(r.LastHeard) <= RuntimeForgotten {
|
||||
x.Unheard = append(x.Unheard, unheard{Runtime: r.Name, Machine: r.Machine,
|
||||
Why: fmt.Sprintf("it answered discovery %s ago, not now", time.Since(r.LastHeard).Round(time.Second))})
|
||||
}
|
||||
}
|
||||
for _, u := range x.Discovery.Silent {
|
||||
k := runtimeKey(u.Name, u.ID)
|
||||
if _, known := s.runtimes[k]; !known {
|
||||
s.runtimes[k] = &runtimeSeen{Name: u.Name, ID: u.ID, Machine: u.Machine, Why: silent[k]}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// runtimesAnswer is mesh_runtimes: discovery asked now, and every runtime as it was heard.
|
||||
func (s *Surface) runtimesAnswer() map[string]any {
|
||||
x, err := s.indexAsked(true)
|
||||
if err != nil {
|
||||
return failure("the mesh's discovery failed: " + err.Error())
|
||||
}
|
||||
s.mu.Lock()
|
||||
var all []runtimeSeen
|
||||
for _, r := range s.runtimes {
|
||||
all = append(all, *r)
|
||||
}
|
||||
s.mu.Unlock()
|
||||
sort.Slice(all, func(i, j int) bool {
|
||||
if all[i].Machine != all[j].Machine {
|
||||
return all[i].Machine < all[j].Machine
|
||||
}
|
||||
return runtimeKey(all[i].Name, all[i].ID) < runtimeKey(all[j].Name, all[j].ID)
|
||||
})
|
||||
now := time.Now()
|
||||
var answered, notHeard []map[string]any
|
||||
onMachine := map[string]bool{}
|
||||
for _, r := range all {
|
||||
entry := map[string]any{"runtime": r.Name, "instance": r.ID}
|
||||
if r.Machine != "" {
|
||||
entry["machine"] = r.Machine
|
||||
}
|
||||
if !r.LastHeard.IsZero() {
|
||||
entry["last heard"] = r.LastHeard.Format(time.RFC3339)
|
||||
entry["last heard, ago"] = now.Sub(r.LastHeard).Round(time.Second).String()
|
||||
entry["took"] = r.Took.Round(time.Millisecond).String()
|
||||
entry["bytes"] = r.Bytes
|
||||
entry["modules"] = r.Modules
|
||||
entry["tools"] = r.Tools
|
||||
entry["seat verbs"] = r.SeatVerbs
|
||||
entry["shortened"] = r.Shortened != ""
|
||||
if r.Shortened != "" {
|
||||
entry["shortened, how"] = r.Shortened
|
||||
}
|
||||
}
|
||||
if r.Answered {
|
||||
onMachine[r.Machine] = true
|
||||
answered = append(answered, entry)
|
||||
} else {
|
||||
entry["why"] = r.Why
|
||||
if r.LastHeard.IsZero() {
|
||||
entry["last heard"] = "never in full, since this console started"
|
||||
}
|
||||
notHeard = append(notHeard, entry)
|
||||
}
|
||||
}
|
||||
for _, m := range x.Machines {
|
||||
if !onMachine[m] && !machineListed(notHeard, m) {
|
||||
notHeard = append(notHeard, map[string]any{"machine": m,
|
||||
"why": "the mesh knows this machine and no runtime on it has answered discovery since this console started"})
|
||||
}
|
||||
}
|
||||
if answered == nil {
|
||||
answered = []map[string]any{}
|
||||
}
|
||||
if notHeard == nil {
|
||||
notHeard = []map[string]any{}
|
||||
}
|
||||
return answerText(map[string]any{
|
||||
"asked": x.Discovery.At.Format(time.RFC3339),
|
||||
"waited": x.Discovery.Waited.Round(time.Millisecond).String(),
|
||||
"window": fmt.Sprintf("at least %s; up to %s for a runtime that answered PING", announce.Window, announce.Patience),
|
||||
"answered": answered,
|
||||
"not heard": notHeard,
|
||||
})
|
||||
}
|
||||
|
||||
func machineListed(entries []map[string]any, machine string) bool {
|
||||
for _, e := range entries {
|
||||
if e["machine"] == machine {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
@@ -0,0 +1,219 @@
|
||||
package console
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go/micro"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/announce"
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
|
||||
)
|
||||
|
||||
// slowRuntime answers PING at once and INFO after `lag`, or never with a negative lag.
|
||||
func slowRuntime(t *testing.T, conn *bus.Conn, node string, lag time.Duration, endpoints []announce.Endpoint) {
|
||||
t.Helper()
|
||||
id := micro.ServiceIdentity{Name: "node-tools", ID: node, Version: announce.Version, Metadata: map[string]string{"node": node}}
|
||||
ping, _ := json.Marshal(micro.Ping{ServiceIdentity: id, Type: micro.PingResponseType})
|
||||
stop, err := conn.Raw("$SRV.PING", func(string, []byte) []byte { return ping })
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
stop, err = conn.Raw("$SRV.INFO", func(string, []byte) []byte {
|
||||
if lag < 0 {
|
||||
return nil
|
||||
}
|
||||
time.Sleep(lag)
|
||||
body, _ := json.Marshal(announce.Info(announce.Service{Name: id.Name, ID: id.ID, Metadata: id.Metadata}, endpoints))
|
||||
return body
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
conn.Flush()
|
||||
}
|
||||
|
||||
// The laptop's symptom (2026-10-05), reproduced: its runtime's one large answer reached the console
|
||||
// after the 750 ms window, and the console said its modules ran nowhere — "nothing in the mesh is
|
||||
// called slack" — or ran only on another machine. Now a runtime that said it is there is waited for,
|
||||
// one that never says what it serves is named wherever an address it might answer is asked for, and
|
||||
// mesh_runtimes shows who answered, how fast and how much.
|
||||
func TestARuntimeThatAnswersLateIsFoundAndOneThatNeverDoesIsNamed(t *testing.T) {
|
||||
keptWas, patienceWas := IndexKept, announce.Patience
|
||||
// Kept for the test: one round of discovery answers every question but mesh_runtimes, which asks anew.
|
||||
IndexKept, announce.Patience = time.Minute, 2500*time.Millisecond
|
||||
t.Cleanup(func() { IndexKept, announce.Patience = keptWas, patienceWas })
|
||||
mt.New(t)
|
||||
|
||||
laptop := connect(t, "node-tools", "laptop")
|
||||
stopTool, err := laptop.HandleSubject("mesh.mod.slack.tool.slack_check.laptop", func(json.RawMessage) (any, error) {
|
||||
return map[string]any{"ok": true}, nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stopTool)
|
||||
slowRuntime(t, laptop, "laptop", announce.Window+400*time.Millisecond, []announce.Endpoint{
|
||||
{Kind: announce.KindTool, Module: "slack", Tool: "slack_check", Node: "laptop", Description: "Is Slack started once?",
|
||||
Subject: "mesh.mod.slack.tool.slack_check.laptop"},
|
||||
{Kind: announce.KindTool, Module: "nextcloud-client", Tool: "nextcloud_check", Node: "laptop",
|
||||
Subject: "mesh.mod.nextcloud-client.tool.nextcloud_check.laptop"},
|
||||
})
|
||||
sleeper := connect(t, "node-tools", "sleeper")
|
||||
slowRuntime(t, sleeper, "sleeper", -1, nil)
|
||||
desktop := connect(t, "node-tools", "desktop")
|
||||
slowRuntime(t, desktop, "desktop", 0, []announce.Endpoint{
|
||||
{Kind: announce.KindTool, Module: "nextcloud-client", Tool: "nextcloud_check", Node: "desktop",
|
||||
Subject: "mesh.mod.nextcloud-client.tool.nextcloud_check.desktop"},
|
||||
})
|
||||
|
||||
controller := connect(t, "mesh-controller", "bench")
|
||||
for verb, output := range map[string]string{
|
||||
"nodes": `[{"name":"laptop"},{"name":"sleeper"},{"name":"desktop"},{"name":"bench"}]`,
|
||||
"modules": `[{"module":"slack","on":["laptop"],"tools":true},{"module":"nap","on":["sleeper"],"tools":true},` +
|
||||
`{"module":"nextcloud-client","on":["laptop","desktop","sleeper"],"tools":true}]`,
|
||||
} {
|
||||
output := output
|
||||
stop, err := controller.HandleSubject("mesh.seat.mesh-controller.tool."+verb, func(json.RawMessage) (any, error) {
|
||||
return map[string]any{"output": output, "ok": true}, nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(stop)
|
||||
}
|
||||
controller.Flush()
|
||||
|
||||
console := connect(t, "node-tools", "laptop-console")
|
||||
up, err := Serve(NewSurface(console, "laptop.node-tools"), "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer up.Close()
|
||||
endpoint := "http://" + up.Address + "/mcp"
|
||||
|
||||
// The late answer is waited for: the laptop's module is found and called.
|
||||
if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "laptop/slack.slack_check"}); isErr || !strings.Contains(got, `"ok": true`) {
|
||||
t.Errorf("a late runtime's module: %s", got)
|
||||
}
|
||||
if got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": "laptop/nextcloud-client.nextcloud_check"}); isErr && strings.Contains(got, "does not run on laptop") {
|
||||
t.Errorf("a late runtime's module was said to run elsewhere: %s", got)
|
||||
}
|
||||
machine, _ := call(t, endpoint, "mesh_machine", map[string]any{"node": "laptop"})
|
||||
if !strings.Contains(machine, "laptop/slack.slack_check") {
|
||||
t.Errorf("the late runtime's machine lists nothing: %s", machine)
|
||||
}
|
||||
|
||||
// The silent one is named, never called missing.
|
||||
for address, want := range map[string]string{
|
||||
"sleeper/nap.doze": "nap is assigned to sleeper in the mesh's records, but nothing that answered discovery serves it: discovery did not hear in full from node-tools on sleeper (it answered PING",
|
||||
"sleeper/nextcloud-client.nextcloud_check": "nextcloud-client on sleeper did not answer discovery: discovery did not hear in full from node-tools on sleeper",
|
||||
} {
|
||||
got, isErr := call(t, endpoint, "mesh_call", map[string]any{"address": address})
|
||||
if !isErr || !strings.Contains(got, want) || strings.Contains(got, "nothing in the mesh is called") || strings.Contains(got, "does not run on") {
|
||||
t.Errorf("%s: %s\nwant %q", address, got, want)
|
||||
}
|
||||
}
|
||||
machine, _ = call(t, endpoint, "mesh_machine", map[string]any{"node": "sleeper"})
|
||||
if !strings.Contains(machine, `"incomplete"`) || !strings.Contains(machine, `"modules": []`) {
|
||||
t.Errorf("a machine whose runtime did not answer: %s", machine)
|
||||
}
|
||||
overview, _ := call(t, endpoint, "mesh_overview", nil)
|
||||
if !strings.Contains(overview, "the runtime node-tools on sleeper") {
|
||||
t.Errorf("overview does not name the silent runtime: %s", overview)
|
||||
}
|
||||
|
||||
// mesh_runtimes: who answered, how fast, how much; who did not.
|
||||
got, isErr := call(t, endpoint, "mesh_runtimes", nil)
|
||||
if isErr {
|
||||
t.Fatal(got)
|
||||
}
|
||||
var report struct {
|
||||
Answered []map[string]any `json:"answered"`
|
||||
NotHeard []map[string]any `json:"not heard"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(got), &report); err != nil {
|
||||
t.Fatalf("%v: %s", err, got)
|
||||
}
|
||||
var lap map[string]any
|
||||
for _, r := range report.Answered {
|
||||
if r["machine"] == "laptop" {
|
||||
lap = r
|
||||
}
|
||||
}
|
||||
if lap == nil {
|
||||
t.Fatalf("the laptop's runtime is not listed as answering: %s", got)
|
||||
}
|
||||
took, _ := time.ParseDuration(lap["took"].(string))
|
||||
if took < announce.Window || lap["bytes"].(float64) <= 0 || lap["modules"].(float64) != 2 || lap["tools"].(float64) != 2 ||
|
||||
lap["shortened"] != false || lap["last heard"] == nil {
|
||||
t.Errorf("the laptop's runtime: %v", lap)
|
||||
}
|
||||
if !strings.Contains(got, `"machine": "sleeper"`) || !strings.Contains(got, "answered PING") {
|
||||
t.Errorf("the silent runtime is not listed: %s", got)
|
||||
}
|
||||
silent := false
|
||||
for _, r := range report.NotHeard {
|
||||
silent = silent || r["machine"] == "sleeper"
|
||||
}
|
||||
if !silent {
|
||||
t.Errorf("sleeper is not among the not heard: %s", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A runtime heard once and silent now is named in the answers about its machine.
|
||||
func TestARuntimeHeardBeforeAndNotNowIsNamed(t *testing.T) {
|
||||
s := &Surface{}
|
||||
before := &index{Modules: map[string]*moduleInfo{}, Recorded: map[string][]string{}, Discovery: announce.Discovery{At: time.Now().UTC(),
|
||||
Heard: []announce.Heard{{Info: micro.Info{ServiceIdentity: micro.ServiceIdentity{Name: "node-tools", ID: "laptop",
|
||||
Metadata: map[string]string{"node": "laptop"}}}, Took: 300 * time.Millisecond, Bytes: 1000}}}}
|
||||
s.remember(before)
|
||||
now := &index{Modules: map[string]*moduleInfo{"nextcloud-client": {Module: "nextcloud-client", On: []string{"desktop"},
|
||||
Tools: []Tool{{Module: "nextcloud-client", Name: "nextcloud_check"}}}}, Recorded: map[string][]string{},
|
||||
Discovery: announce.Discovery{At: time.Now().UTC()}}
|
||||
s.remember(now)
|
||||
if len(now.Unheard) != 1 || now.Unheard[0].Machine != "laptop" {
|
||||
t.Fatalf("unheard: %+v", now.Unheard)
|
||||
}
|
||||
_, err := resolve(now, "laptop/nextcloud-client.nextcloud_check")
|
||||
if err == nil || !strings.Contains(err.Error(), "nextcloud-client on laptop did not answer discovery") ||
|
||||
!strings.Contains(err.Error(), "answered discovery") {
|
||||
t.Errorf("%v", err)
|
||||
}
|
||||
_, err = resolve(now, "laptop/slack.slack_check")
|
||||
if err == nil || strings.Contains(err.Error(), "nothing in the mesh is called") || !strings.Contains(err.Error(), "slack may be theirs") {
|
||||
t.Errorf("%v", err)
|
||||
}
|
||||
// Nobody unheard: the plain answers stand.
|
||||
quiet := &index{Modules: now.Modules, Recorded: map[string][]string{}}
|
||||
if _, err := resolve(quiet, "laptop/nextcloud-client.nextcloud_check"); err == nil || !strings.Contains(err.Error(), "does not run on laptop; it runs on desktop") {
|
||||
t.Errorf("%v", err)
|
||||
}
|
||||
if _, err := resolve(quiet, "laptop/slack.slack_check"); err == nil || !strings.Contains(err.Error(), "nothing in the mesh is called slack") {
|
||||
t.Errorf("%v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// A runtime restarted answers under a new instance id; its old instance is gone, not unheard.
|
||||
func TestARestartedRuntimeIsNotNamedUnheard(t *testing.T) {
|
||||
heard := func(id string) *index {
|
||||
return &index{Modules: map[string]*moduleInfo{}, Recorded: map[string][]string{}, Discovery: announce.Discovery{At: time.Now().UTC(),
|
||||
Heard: []announce.Heard{{Info: micro.Info{ServiceIdentity: micro.ServiceIdentity{Name: "node-tools", ID: id,
|
||||
Metadata: map[string]string{"node": "laptop"}}}}}}}
|
||||
}
|
||||
s := &Surface{}
|
||||
s.remember(heard("before-restart"))
|
||||
after := heard("after-restart")
|
||||
s.remember(after)
|
||||
if len(after.Unheard) != 0 {
|
||||
t.Errorf("a restarted runtime was named unheard: %+v", after.Unheard)
|
||||
}
|
||||
if len(s.runtimes) != 1 {
|
||||
t.Errorf("the old instance is still kept: %d runtimes", len(s.runtimes))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user