diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index 0013c71..8c50880 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -449,7 +449,11 @@ func (c *Conn) answerOn(subject, queue string, h Handler) (func(), error) { } r.Node = c.node body, _ := wire.Marshal(r) - _ = msg.Respond(body) + // Within what the bus carries, or said to be too large — never a send the library + // refuses and the caller waits out (novox/hq issue 314). + if err := msg.Respond(c.fits(subject, body)); err != nil { + c.Logf("[mesh-tools] could not answer on %s: %v", subject, err) + } }() } var sub *nats.Subscription @@ -614,13 +618,18 @@ func (c *Conn) askOnce(subject string, data []byte) (Answered, bool, error) { } return Answered{}, false, err } + // An answer larger than one message comes in parts (pages.go, novox/hq issue 314). + body, err := c.whole(subject, msg.Data) + if err != nil { + return Answered{}, false, err + } var r struct { Result json.RawMessage `json:"result"` Error string `json:"error"` Node string `json:"node"` Retry string `json:"retry"` } - if err := json.Unmarshal(msg.Data, &r); err != nil { + if err := json.Unmarshal(body, &r); err != nil { return Answered{}, false, err } if r.Error != "" { diff --git a/node-tools/internal/bus/pages.go b/node-tools/internal/bus/pages.go new file mode 100644 index 0000000..6351f4d --- /dev/null +++ b/node-tools/internal/bus/pages.go @@ -0,0 +1,123 @@ +package bus + +import ( + "bytes" + "encoding/json" + "errors" + "fmt" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-tools/node-tools/internal/wire" +) + +// An answer larger than the bus carries in one message (novox/hq issue 314). +// +// The bus carries one message of at most its max_payload, and the client library refuses a larger one +// before it leaves the process that sent it. Until this, a holder whose answer was too large wrote one +// line to its own log and its caller waited out RequestTimeout and read "did not answer" — the +// controller's `conditions` and `status` for hours on 2026-10-08, and the operator's channel with them. +// +// The controller pages such an answer (mesh-controller internal/link, pages.go): it answers an error +// carrying `paged` — the call it holds the answer under, its size and its parts — and a caller asks the +// same subject again once per part with PageHeader, each part answered raw with PartHeader. Asking the +// same subject means no grant beyond the one the caller already had. A module's answer too large is +// refused here, loudly, in a message that fits: its instances may be several, and a part asked again +// could reach another. + +const ( + // PageHeader asks for one part of a paged answer: ` `, from zero. + PageHeader = "Mesh-Page" + // PartHeader is on a part: `/`. + PartHeader = "Mesh-Part" + // partTries is how often one part is asked: during a handover a part can reach the controller that + // did not hold it. + partTries = 3 +) + +// PartTimeout bounds the asking of one part. +var PartTimeout = 10 * time.Second + +// paged is what a too-large answer says instead of itself. +type paged struct { + Call string `json:"call"` + Bytes int `json:"bytes"` + Parts int `json:"parts"` +} + +func pagedIn(data []byte) (paged, bool) { + if !bytes.Contains(data, []byte(`"paged"`)) { + return paged{}, false + } + var r struct { + Paged *paged `json:"paged"` + } + if json.Unmarshal(data, &r) != nil || r.Paged == nil || r.Paged.Call == "" || r.Paged.Parts < 1 { + return paged{}, false + } + return *r.Paged, true +} + +// whole is a reply's whole answer: the reply itself, or every part of a paged one, joined. A part that +// cannot be had is an error naming it, never a shorter answer. +func (c *Conn) whole(subject string, data []byte) ([]byte, error) { + p, ok := pagedIn(data) + if !ok { + return data, nil + } + out := make([]byte, 0, p.Bytes) + for n := 0; n < p.Parts; n++ { + var part []byte + var err error + for try := 0; try < partTries; try++ { + if part, err = c.askPart(subject, p, n); err == nil { + break + } + } + if err != nil { + return nil, fmt.Errorf("the answer was %d bytes, more than the bus carries in one message, paged as %s "+ + "in %d parts, and part %d could not be had: %w", p.Bytes, p.Call, p.Parts, n, err) + } + out = append(out, part...) + } + if len(out) != p.Bytes { + return nil, fmt.Errorf("the answer paged as %s was %d bytes, and its parts joined are %d", p.Call, p.Bytes, len(out)) + } + return out, nil +} + +func (c *Conn) askPart(subject string, p paged, n int) ([]byte, error) { + ask := nats.NewMsg(subject) + ask.Header.Set(PageHeader, fmt.Sprintf("%s %d", p.Call, n)) + ask.Data = []byte(`{}`) + reply, err := c.nc.RequestMsg(ask, PartTimeout) + if err != nil { + return nil, err + } + if got := reply.Header.Get(PartHeader); got != fmt.Sprintf("%d/%d", n, p.Parts) { + var r struct { + Error string `json:"error"` + } + if json.Unmarshal(reply.Data, &r) == nil && r.Error != "" { + return nil, errors.New(r.Error) + } + return nil, fmt.Errorf("asked part %d/%d, answered %q", n, p.Parts, got) + } + return reply.Data, nil +} + +// fits is an answer the bus can carry: the answer itself, or — when it is larger than one message — +// an error that says so, which is. +func (c *Conn) fits(subject string, body []byte) []byte { + limit := c.nc.MaxPayload() + if limit <= 0 || int64(len(body)) <= limit { + return body + } + c.Logf("[mesh-tools] the answer on %s is %d bytes, more than the bus carries in one message (%d): its caller "+ + "is told so instead", subject, len(body), limit) + said, _ := wire.Marshal(reply{Node: c.node, Error: fmt.Sprintf("the answer is %d bytes, more than the bus "+ + "carries in one message (%d), so it was not sent: ask for less — a narrower question, a shorter window, "+ + "fewer lines", len(body), limit)}) + return said +} diff --git a/node-tools/internal/bus/pages_test.go b/node-tools/internal/bus/pages_test.go new file mode 100644 index 0000000..b06b1c1 --- /dev/null +++ b/node-tools/internal/bus/pages_test.go @@ -0,0 +1,80 @@ +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) + } +}