Merge pull request 'Join a paged answer, and say one too large instead of timing out (hq issue 314)' (#22) from fix/314-a-large-answer-is-paged-not-lost into main
This commit was merged in pull request #22.
This commit is contained in:
@@ -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 != "" {
|
||||
|
||||
@@ -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: `<call> <part>`, from zero.
|
||||
PageHeader = "Mesh-Page"
|
||||
// PartHeader is on a part: `<part>/<parts>`.
|
||||
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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user