Keep a machine's discovery answer under the bus's message limit

Every tool's schema in one $SRV.INFO reply outgrew max_payload on machines
serving 137-206 tools, and the refused reply was dropped silently, so search
and a machine's view went empty. Module tools now announce without schema;
describe and tools/list ask the module for it. An answer still too large has
its descriptions cut to a line, and a failed reply is logged.
This commit is contained in:
jochen
2026-10-04 17:06:20 +02:00
parent 51c79d8461
commit b3ebdd5edd
8 changed files with 184 additions and 10 deletions
+40 -4
View File
@@ -74,11 +74,19 @@ func (e Endpoint) info() micro.EndpointInfo {
"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 = "{}"
// **A seat's verb carries its schema; a module's tool does not.** Every tool's schema, once per
// subject it answers on, made a runtime's one answer outgrow the bus's largest message (1 MiB) once
// machines served a couple of hundred tools — and the answer that cannot be sent is silence: the
// machine vanished from discovery (2026-10-04). A module's tool's schema is asked of the runtime that
// serves it when a person describes it (its `tools` verb); a seat's verbs are few, and a seat has no
// `tools` verb of its own to ask.
if e.Kind == KindSeat {
schema := strings.TrimSpace(string(e.Schema))
if schema == "" || schema == "null" {
schema = "{}"
}
md["schema"] = schema
}
md["schema"] = schema
if e.Kind == KindSeat {
md["seat"] = e.Seat
md["scope"] = e.Scope
@@ -154,6 +162,20 @@ func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error)
if err != nil {
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.
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
}
}
}
return body
}
var stops []func()
@@ -209,3 +231,17 @@ func Endpoints(i micro.Info) []Endpoint {
}
return out
}
// firstLine is a description's first sentence or line, at most 160 characters.
func firstLine(s string) string {
if i := strings.IndexAny(s, "\n"); i >= 0 {
s = s[:i]
}
if i := strings.Index(s, ". "); i >= 0 {
s = s[:i+1]
}
if len(s) > 160 {
s = s[:157] + "..."
}
return s
}
@@ -0,0 +1,54 @@
package announce
import (
"encoding/json"
"fmt"
"strings"
"testing"
)
// What a machine serves fits the bus's largest message (1 MiB) with room to spare at several hundred
// tools: a module's tool is announced without its schema, which made a machine with ~200 tools outgrow it
// and vanish from discovery (2026-10-04). A seat's verb keeps its schema.
func TestAMachinesAnnouncementFitsTheBusAtSeveralHundredTools(t *testing.T) {
schema := json.RawMessage(`{"type":"object","properties":{` + strings.Repeat(`"field_with_a_long_name":{"type":"string","description":"a description of what this argument means, long enough to matter"},`, 12) + `"last":{"type":"string"}}}`)
description := strings.Repeat("A tool that does something useful, described at the length the catalogue's tools are. ", 3)
var endpoints []Endpoint
for i := 0; i < 400; i++ {
for _, subject := range []string{fmt.Sprintf("mesh.mod.m%d.tool.t", i), fmt.Sprintf("mesh.mod.m%d.tool.t.laptop", i)} {
endpoints = append(endpoints, Endpoint{Kind: KindTool, Module: fmt.Sprintf("m%d", i), Tool: "t", Node: "laptop",
Description: description, Schema: schema, Subject: subject})
}
}
endpoints = append(endpoints, Endpoint{Kind: KindSeat, Module: "m0", Tool: "verb", Seat: "a-seat", Scope: "mesh",
Description: "a verb", Schema: schema, Subject: "mesh.seat.a-seat.tool.verb"})
body, err := json.Marshal(Info(Service{Name: "node-tools", ID: "laptop"}, endpoints))
if err != nil {
t.Fatal(err)
}
if len(body) > 1<<20 {
t.Fatalf("400 tools announce %d bytes, beyond the bus's 1 MiB", len(body))
}
t.Logf("400 tools on two subjects each announce %d bytes", len(body))
back := Endpoints(Info(Service{Name: "node-tools", ID: "laptop"}, endpoints))
if len(back[0].Schema) != 0 {
t.Fatalf("a module's tool still carries its schema: %s", back[0].Schema)
}
if seat := back[len(back)-1]; !strings.Contains(string(seat.Schema), "field_with_a_long_name") {
t.Fatalf("a seat's verb lost its schema: %s", seat.Schema)
}
withSchemas := 0
for range endpoints {
withSchemas += len(schema) * 2 // escaped inside a string
}
t.Logf("with every schema it would have been over %d bytes", len(body)+withSchemas)
}
func TestFirstLineShortensADescription(t *testing.T) {
if got := firstLine("One thing. And more after it."); got != "One thing." {
t.Fatalf("%q", got)
}
if got := firstLine(strings.Repeat("x", 300)); len(got) != 160 {
t.Fatalf("%d", len(got))
}
}
+7 -1
View File
@@ -636,7 +636,10 @@ func (c *Conn) ServedOn(module, tool string) []Served { return c.servedOn(module
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)
// Said, never dropped: an answer the bus will not carry was silence before (2026-10-04).
if err := msg.Respond(body); err != nil {
c.Logf("[mesh-tools] could not answer %s (%d bytes): %v", msg.Subject, len(body), err)
}
}
})
if err != nil {
@@ -679,6 +682,9 @@ func (c *Conn) Gather(subject string, body []byte, window time.Duration) ([][]by
}
}
// MaxPayload is the largest message the bus carries, as the server told this connection.
func (c *Conn) MaxPayload() int64 { return c.nc.MaxPayload() }
// 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() }
+30
View File
@@ -691,6 +691,10 @@ func (s *Surface) discover(name string, args map[string]any) map[string]any {
if err != nil {
return failure(err.Error())
}
if !t.Seat && len(t.Tool.Input) == 0 {
// Discovery carries a seat's schema, not a module tool's: asked of the runtime serving it.
t.Tool.Input = inputOf(s.conn, t.Tool.Module, t.Tool.Name)
}
description := t.Tool.Description
if description == "" {
description = t.Name
@@ -728,3 +732,29 @@ func (s *Surface) discover(name string, args map[string]any) map[string]any {
}
return failure("no discovery verb " + name)
}
// inputOf is a module tool's input schema, asked of a runtime that serves the module — its `tools` verb,
// answered on the module's plain subject by any instance — because discovery no longer carries it: every
// tool's schema made a machine's one discovery answer outgrow the bus (2026-10-04). Nil when nothing
// answers, which describes the tool as taking nothing rather than failing the description.
func inputOf(conn *bus.Conn, module, tool string) json.RawMessage {
got, err := conn.Ask(module+".tools", map[string]any{}, "")
if err != nil {
return nil
}
var answer struct {
Tools []struct {
Name string `json:"name"`
Input json.RawMessage `json:"input"`
} `json:"tools"`
}
if json.Unmarshal(got.Result, &answer) != nil {
return nil
}
for _, t := range answer.Tools {
if t.Name == tool {
return t.Input
}
}
return nil
}
+13 -2
View File
@@ -178,11 +178,20 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
}
}
// Discovery itself asked nothing but the bus.
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)
}
// Described without `node`: the address carries the machine.
described, isErr := call(t, endpoint, "mesh_describe", map[string]any{"address": "desk/beta.three"})
if isErr || strings.Contains(described, `"node"`) || !strings.Contains(described, `"address": "desk/beta.three"`) {
t.Errorf("describe: %s", described)
}
// Its schema, which discovery no longer carries for a module's tool, asked of the runtime serving it.
if !strings.Contains(described, `"verbose"`) {
t.Errorf("describe lost the tool's arguments: %s", described)
}
// Search finds across the layers; a module that starts serving after the first answer is found.
if got, _ := call(t, endpoint, "mesh_search", map[string]any{"query": "shelf"}); !strings.Contains(got, "desk/node-shelf.list") {
@@ -212,8 +221,10 @@ 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)
// Only describing a module's tool asks its runtime — once, for the schema discovery no longer carries
// (2026-10-04: every tool's schema made a machine's discovery answer outgrow the bus).
if n := asked.Load(); n != 1 {
t.Errorf("asked a module's tools %d time(s); describing one tool asks once, and discovery never", n)
}
// The old names still answer, unannounced.
+30 -1
View File
@@ -125,7 +125,15 @@ func (s *Surface) Handle(r Request) *Reply {
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 {
inputs := map[string]map[string]json.RawMessage{} // module → tool → input, asked once per module
for i, t := range l.Tools {
if !t.Seat && len(t.Input) == 0 {
if inputs[t.Module] == nil {
inputs[t.Module] = toolInputs(s.conn, t.Module)
}
l.Tools[i].Input = inputs[t.Module][t.Name]
t = l.Tools[i]
}
var schema map[string]any
switch {
case t.Seat && t.Scope != "node":
@@ -295,3 +303,24 @@ func (s *Surface) instructions() string {
"arguments; mesh_call calls it. " + grammar + " What may be called was fixed when this account " +
"was issued, so a refusal means the account, not the tool."
}
// toolInputs is every tool's input schema of one module, asked of a runtime that serves it.
func toolInputs(conn *bus.Conn, module string) map[string]json.RawMessage {
out := map[string]json.RawMessage{}
got, err := conn.Ask(module+".tools", map[string]any{}, "")
if err != nil {
return out
}
var answer struct {
Tools []struct {
Name string `json:"name"`
Input json.RawMessage `json:"input"`
} `json:"tools"`
}
if json.Unmarshal(got.Result, &answer) == nil {
for _, t := range answer.Tools {
out[t.Name] = t.Input
}
}
return out
}
+9 -1
View File
@@ -74,9 +74,17 @@ func TestTheRuntimeAnnouncesWhatItServesInTheServicesProtocol(t *testing.T) {
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"])) {
if e.Metadata["node"] != "anchor" || e.Metadata["description"] == "" {
t.Errorf("endpoint metadata: %+v", e)
}
// A seat's verb carries its schema; a module's tool does not — it outgrew the bus (2026-10-04).
_, hasSchema := e.Metadata["schema"]
if e.Metadata["kind"] == "seat" && !json.Valid([]byte(e.Metadata["schema"])) {
t.Errorf("a seat's verb without a schema: %+v", e)
}
if e.Metadata["kind"] == "tool" && hasSchema {
t.Errorf("a module's tool still announces its schema: %+v", e)
}
}
// Every subject announced is answered.