Merge pull request 'Keep a machine's discovery answer under the bus's message limit' (#52) from fix/discovery-fits-in-a-message into main
This commit is contained in:
@@ -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))
|
||||
}
|
||||
}
|
||||
@@ -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() }
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
+1
-1
@@ -3,7 +3,7 @@
|
||||
import { registerModuleTools } from "@novox/mesh-sdk/tools";
|
||||
|
||||
registerModuleTools("beta", () => [
|
||||
{ name: "three", description: "beta's", input: {}, run: async () => ({ beta: 3 }) },
|
||||
{ name: "three", description: "beta's", input: { verbose: { type: "boolean", description: "say more" } }, run: async () => ({ beta: 3 }) },
|
||||
{ name: "four", description: "beta's", input: {}, run: async () => ({ beta: 4 }) },
|
||||
{ name: "five", description: "beta's", input: {}, run: async () => ({ beta: 5 }) },
|
||||
]);
|
||||
|
||||
Reference in New Issue
Block a user