mesh/merge-gate pass: builds mesh-tools, node-tools → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without (4 of …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
mesh/delivery-group group fix/314-a-large-answer-is-paged-not-lost delivered: every member is delivered
An answer larger than the bus's max_payload was refused by the client library inside its holder, and the caller waited out its 30 s and read silence: the controller's conditions and status for hours on 2026-10-08. The console now asks for the parts of an answer the controller paged and joins them; a module's answer too large for one message is answered as an error saying so, at once.
81 lines
2.7 KiB
Go
81 lines
2.7 KiB
Go
package bus_test
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
|
"github.com/novox/mesh-tools/node-tools/internal/meshtest"
|
|
)
|
|
|
|
// novox/hq issue 314: an answer larger than the bus carries in one message arrives whole when the
|
|
// controller pages it, and a module's is said to be too large — neither is a 30-second silence.
|
|
func TestAnAnswerLargerThanOneMessage(t *testing.T) {
|
|
url := meshtest.URL(t)
|
|
seat, err := nats.Connect(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer seat.Close()
|
|
c, err := bus.Connect(bus.Credential{URL: url, Module: "console", Node: "anchor"})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer c.Close()
|
|
limit := int(seat.MaxPayload())
|
|
|
|
// The controller as it pages (mesh-controller internal/link, pages.go): the first answer an error
|
|
// with `paged`, each part asked again on the same subject with the page header.
|
|
large := strings.Repeat("a condition held past its bound; ", 3*limit/32)
|
|
whole, _ := json.Marshal(map[string]any{"result": large})
|
|
partSize := limit - 1024
|
|
parts := (len(whole) + partSize - 1) / partSize
|
|
sub, err := seat.Subscribe(bus.SeatToolSubject("mesh-controller", "conditions", "", ""), func(m *nats.Msg) {
|
|
if page := m.Header.Get(bus.PageHeader); page != "" {
|
|
var id string
|
|
var n int
|
|
fmt.Sscanf(page, "%s %d", &id, &n)
|
|
end := min((n+1)*partSize, len(whole))
|
|
part := nats.NewMsg(m.Reply)
|
|
part.Header.Set(bus.PartHeader, fmt.Sprintf("%d/%d", n, parts))
|
|
part.Data = whole[n*partSize : end]
|
|
_ = m.RespondMsg(part)
|
|
return
|
|
}
|
|
body, _ := json.Marshal(map[string]any{"error": "too large for one message",
|
|
"paged": map[string]any{"call": "call-1-1", "bytes": len(whole), "parts": parts}})
|
|
_ = m.Respond(body)
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer sub.Unsubscribe()
|
|
_ = seat.Flush()
|
|
|
|
got, err := c.Ask("seat:mesh-controller.conditions", map[string]any{}, "")
|
|
var said string
|
|
if err != nil || json.Unmarshal(got.Result, &said) != nil || said != large {
|
|
t.Fatalf("an answer of %d bytes in %d parts arrived as %d bytes: %v", len(whole), parts, len(got.Result), err)
|
|
}
|
|
|
|
// A module's answer too large is refused here in words, at once.
|
|
stop, err := c.HandleSubject("mesh.mod.journal.tool.lines", func(json.RawMessage) (any, error) { return large, nil })
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer stop()
|
|
asked := time.Now()
|
|
_, err = c.Ask("journal.lines", map[string]any{}, "mesh.mod.journal.tool.lines")
|
|
if err == nil || !strings.Contains(err.Error(), "more than the bus carries") {
|
|
t.Fatalf("a module's answer too large for one message read as %v", err)
|
|
}
|
|
if took := time.Since(asked); took > 5*time.Second {
|
|
t.Fatalf("it took %s to say", took)
|
|
}
|
|
}
|