Page an answer larger than one message of the bus, and keep overviews brief (hq issue 314)
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
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

The client library refuses to send a reply over the bus's max_payload, and the
controller only logged it: conditions and status answered nobody for hours on
2026-10-08 while calls said each was answered in 130 ms, and the operator's
channel read nothing. An answer too large is now held under its call and paged
to the caller that asks, on the same subject; a caller that does not page is
told in words, and calls says it. The overviews no longer carry every finding:
conditions and status list each condition with its newest evidence, doctor at
most twenty findings a probe (probe= gives one whole), and the JSON overviews
are sent once, as data, instead of twice.
This commit is contained in:
jochen
2026-10-08 12:19:16 +02:00
parent 4d337ff897
commit 175b28ee42
12 changed files with 793 additions and 27 deletions
+1
View File
@@ -351,6 +351,7 @@ var ControllerVerbs = []Verb{
"run": "\"true\": run every probe now and answer the verdict",
"probes": "\"true\": the registry — what each probe asserts, and the condition it raises",
"signals": "\"true\": the signals table, each row with the age of its newest signal",
"probe": "one probe's id (D14): its verdict alone, with every finding — the verdict shows at most twenty a probe",
}, nil, "run", "probes", "signals")},
// How a module's new builds reach its machines, and the bus's planned step (novox/hq ADR 0236).
{Name: "upgrade", Description: "How each module's new builds reach its machines (novox/hq ADR 0236): rolled " +
+47 -2
View File
@@ -81,6 +81,9 @@ type Call struct {
// Refused is the bus refusing the answer this holder sent: the caller got nothing, and this is
// the only place that says what it would have.
Refused string `json:"answer refused by the bus,omitempty"`
// Paged is an answer too large for one message, sent in parts to a caller that asked for them
// (novox/hq issue 314): how large, and until when it is held.
Paged string `json:"answer paged,omitempty"`
// Caller is the bus principal that asked, read from the inbox its answer went to — every principal
// is granted only its own (novox/hq to-be 45 §7: a hand act says who).
Caller string `json:"caller,omitempty"`
@@ -125,6 +128,9 @@ type CallLog struct {
// epoch is the lease this process serves under (UnderLease): a call carries its epoch, and its
// record is written only while the lease is held. Nil writes every record with no epoch.
epoch func() (uint64, error)
// pages are the answers too large for one message, held for their callers (pages.go).
pages pageBook
}
// UnderLease makes every call carry the controller lease's epoch, and every write of a call's record
@@ -316,9 +322,21 @@ func (l *CallLog) finish(c *Call, answer []byte, failed, answeredAlready bool) {
l.keep(onTheBus)
return
}
l.keep(*c)
onTheBus := *c
if len(onTheBus.Answer) > keptOnTheBusAtMost {
// A record on the bus is one message too, and one larger than the bus carries is refused whole —
// the call's state with it (novox/hq issue 314). The answer stays in this process's memory, and
// `calls` with its id pages it from there.
onTheBus.Answer, _ = json.Marshal(map[string]any{"not kept": fmt.Sprintf("an answer of %d bytes, more than "+
"a record on the bus holds: kept in the memory of the controller that answered it", len(c.Answer))})
}
l.keep(onTheBus)
}
// keptOnTheBusAtMost is the largest answer a call's record on the bus carries: half of the bus's default
// message limit, leaving the rest of the record room.
const keptOnTheBusAtMost = 512 << 10
// shownOnce are the verbs whose answer is a secret shown once to its caller, as `<seat>.<verb>`.
var shownOnce = map[string]bool{"mesh-controller.token": true}
@@ -529,6 +547,15 @@ func running(c *Call, within time.Duration, acknowledged bool, follow string) []
// and has not acknowledged, and otherwise that it is running — its answer then kept, and logged.
func (l *CallLog) serveCall(seat, verb string, args json.RawMessage, reply string, handle ToolHandler,
respond func([]byte) error, logger *log.Logger) {
l.serveCallWithin(seat, verb, args, reply, handle, respond, 0, logger)
}
// serveCallWithin is serveCall on a bus that carries at most limit bytes in one message: an answer
// larger than that is paged (pages.go), and one the bus would not carry anyway is said to its caller
// and to `calls` — never dropped with only a line in the journal (novox/hq issue 314). A limit of 0 is
// one not known.
func (l *CallLog) serveCallWithin(seat, verb string, args json.RawMessage, reply string, handle ToolHandler,
respond func([]byte) error, limit int64, logger *log.Logger) {
ctx, cancel := context.WithTimeout(context.Background(), HandlerTimeout)
defer cancel()
if len(args) == 0 {
@@ -563,9 +590,27 @@ func (l *CallLog) serveCall(seat, verb string, args json.RawMessage, reply strin
}()
say := func(body []byte) {
if err := respond(body); err != nil && logger != nil {
if limit > 0 && int64(len(body)) > limit {
if logger != nil {
logger.Printf("%s.%s: the answer to %s is %d bytes, more than the bus carries in one message (%d): "+
"paged", seat, verb, c.ID, len(body), limit)
}
body = l.page(c, body, limit)
}
err := respond(body)
if err == nil {
return
}
if logger != nil {
logger.Printf("%s.%s: could not answer %s: %v", seat, verb, c.ID, err)
}
l.unsent(c, err)
if errors.Is(err, nats.ErrMaxPayload) {
// The limit was not known, or was not the bus's: its caller is told, in a message that fits.
if again := respond(cannotCarry(seat, verb, c.ID, len(body), err)); again != nil && logger != nil {
logger.Printf("%s.%s: could not say to the caller of %s why it has no answer: %v", seat, verb, c.ID, again)
}
}
}
timer := time.NewTimer(AnswerWithin)
defer timer.Stop()
+302
View File
@@ -0,0 +1,302 @@
package link
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"log"
"strconv"
"strings"
"sync"
"time"
"github.com/nats-io/nats.go"
)
// An answer larger than the bus carries in one message, paged (novox/hq issue 314).
//
// **The bus carries one message of at most its max_payload** (a mebibyte unless the server says
// otherwise), and the client library refuses to send a larger one: `nats: maximum payload exceeded`,
// returned to the one who sent it and to nobody else. A holder that answered a call with more than
// that had its answer refused before it left the process — the caller waited out its whole timeout and
// read "did not answer in time", and `calls` said the call was answered, because it was, as far as the
// holder's own record went. On 2026-10-08 that was every `conditions` and `status` for hours, once 689
// stalled deliveries were open conditions each, and the operator's channel, which reads `conditions`
// every minute, among them.
//
// So an answer too large for one message is never sent as one. Its holder keeps it for a while, under
// the call's id, and answers with an error that says so — how large, in how many parts, and how to ask
// for them — so a caller that does not page is told, loudly, instead of timing out. A caller that does
// page (the mesh's console, and this package's own askers) asks the same subject again, once per part,
// with the PageHeader naming the call and the part: the same subject because the caller's grant
// already names it, and a part is no more than the caller was entitled to ask for. The bus's
// `allow_responses` permits one answer per request, so a part is a request of its own, never a second
// answer to the first. A part is raw bytes with PartHeader saying which, so no part grows past the
// limit for being quoted.
const (
// PageHeader asks for one part of a paged answer: `<call id> <part>`, parts counted from zero.
PageHeader = "Mesh-Page"
// PartHeader is on a part: `<part>/<parts>`.
PartHeader = "Mesh-Part"
)
// PagesKeptFor is how long a paged answer is held for its caller to ask for its parts; a caller pages
// at once, so this is for a slow bus, not for coming back later.
var PagesKeptFor = 5 * time.Minute
// pagesHeldAtMost is the most a holder keeps in paged answers at once; the oldest goes first.
const pagesHeldAtMost = 64 << 20
// partHeadroom is what a part leaves of the limit for its headers.
const partHeadroom = 1024
// Paged is what an answer too large for one message says instead of itself.
type Paged struct {
Call string `json:"call"`
Bytes int `json:"bytes"`
Parts int `json:"parts"`
Limit int64 `json:"limit"`
Until time.Time `json:"until"`
}
// heldAnswer is one paged answer, kept for its caller.
type heldAnswer struct {
body []byte
seat, verb string
caller string
partSize, parts int
until, heldSince time.Time
}
// pageBook holds the paged answers of one process.
type pageBook struct {
mu sync.Mutex
held map[string]*heldAnswer
order []string
bytes int
}
func (p *pageBook) put(id string, h *heldAnswer, now time.Time) {
p.mu.Lock()
defer p.mu.Unlock()
if p.held == nil {
p.held = map[string]*heldAnswer{}
}
p.dropExpired(now)
for p.bytes+len(h.body) > pagesHeldAtMost && len(p.order) > 0 {
p.drop(p.order[0])
}
p.held[id] = h
p.order = append(p.order, id)
p.bytes += len(h.body)
}
func (p *pageBook) drop(id string) {
if h, ok := p.held[id]; ok {
p.bytes -= len(h.body)
delete(p.held, id)
}
for i, o := range p.order {
if o == id {
p.order = append(p.order[:i], p.order[i+1:]...)
break
}
}
}
func (p *pageBook) dropExpired(now time.Time) {
for _, id := range append([]string(nil), p.order...) {
if h := p.held[id]; h != nil && now.After(h.until) {
p.drop(id)
}
}
}
func (p *pageBook) get(id string, now time.Time) (*heldAnswer, bool) {
p.mu.Lock()
defer p.mu.Unlock()
p.dropExpired(now)
h, ok := p.held[id]
return h, ok
}
// tooLarge is the answer a caller reads in place of one the bus cannot carry: an error, so a caller
// that does not page is told rather than left waiting, and `paged` for one that does.
func tooLarge(seat, verb string, p Paged) []byte {
said := fmt.Sprintf("the answer to %s.%s is %d bytes, more than the bus carries in one message (%d): it "+
"was not sent whole. It is held as %s until %s, in %d parts — a caller that pages asks %s.%s again "+
"with the header %q set to \"%s <part>\", parts 0 to %d, and joins them; the mesh's console does this "+
"itself. Or ask for less: one condition by its key, a scope, one machine.", seat, verb, p.Bytes, p.Limit,
p.Call, p.Until.UTC().Format(time.RFC3339), p.Parts, seat, verb, PageHeader, p.Call, p.Parts-1)
body, _ := json.Marshal(map[string]any{"error": said, "paged": p})
return body
}
// cannotCarry is the answer when even paging is not possible — the limit unknown, the bus having
// refused the send anyway: never silence.
func cannotCarry(seat, verb, id string, size int, err error) []byte {
body, _ := json.Marshal(map[string]any{"error": fmt.Sprintf("the answer to %s.%s (%s) is %d bytes and the bus "+
"would not carry it: %v. Ask for less — one condition by its key, a scope, one machine", seat, verb, id,
size, err)})
return body
}
// page holds an answer too large for limit and returns what its caller is sent instead.
func (l *CallLog) page(c *Call, body []byte, limit int64) []byte {
partSize := int(limit) - partHeadroom
if partSize < 1024 {
partSize = int(limit) / 2
}
parts := (len(body) + partSize - 1) / partSize
now := l.now()
p := Paged{Call: c.ID, Bytes: len(body), Parts: parts, Limit: limit, Until: now.Add(PagesKeptFor)}
l.pages.put(c.ID, &heldAnswer{body: body, seat: c.Seat, verb: c.Verb, caller: c.Caller, partSize: partSize,
parts: parts, until: p.Until, heldSince: now}, now)
l.mu.Lock()
c.Paged = fmt.Sprintf("%d bytes, more than the bus carries in one message (%d): held until %s and sent in "+
"%d parts to a caller that asks for them", len(body), limit, p.Until.UTC().Format(time.RFC3339), parts)
if c.State != CallFinishedAfter {
// Its caller is sent it in parts; a copy on the bus would be refused for the same size.
c.Answer, _ = json.Marshal(map[string]any{"not kept": fmt.Sprintf("an answer of %d bytes, paged to its caller",
len(body))})
}
kept := *c
l.mu.Unlock()
l.keep(kept)
return tooLarge(c.Seat, c.Verb, p)
}
// unsent records the bus refusing to carry an answer, against its call: `calls` says it.
func (l *CallLog) unsent(c *Call, err error) {
l.mu.Lock()
c.Refused = fmt.Sprintf("%s: not sent: %v", l.now().Format(time.RFC3339), err)
kept := *c
l.mu.Unlock()
l.keep(kept)
}
// servePage answers one part of a paged answer: only to the caller it was paged to, and only on the
// verb it answered.
func (l *CallLog) servePage(seat, verb string, msg *nats.Msg, logger *log.Logger) {
refuse := func(why string) {
body, _ := json.Marshal(map[string]any{"error": why})
if err := msg.Respond(body); err != nil && logger != nil {
logger.Printf("%s.%s: could not refuse a part: %v", seat, verb, err)
}
}
id, n, err := parsePage(msg.Header.Get(PageHeader))
if err != nil {
refuse(err.Error())
return
}
h, ok := l.pages.get(id, l.now())
if !ok {
refuse(fmt.Sprintf("%s is not held here: a paged answer is held for %s by the controller that answered "+
"it — ask %s.%s again", id, PagesKeptFor, seat, verb))
return
}
if h.seat != seat || h.verb != verb {
refuse(fmt.Sprintf("%s answered %s.%s, not %s.%s", id, h.seat, h.verb, seat, verb))
return
}
if caller := callerOf(msg.Reply); h.caller != "" && caller != h.caller {
refuse(fmt.Sprintf("%s was paged to another caller", id))
return
}
if n < 0 || n >= h.parts {
refuse(fmt.Sprintf("%s has parts 0 to %d, not %d", id, h.parts-1, n))
return
}
end := (n + 1) * h.partSize
if end > len(h.body) {
end = len(h.body)
}
part := nats.NewMsg(msg.Reply)
part.Header.Set(PartHeader, fmt.Sprintf("%d/%d", n, h.parts))
part.Data = h.body[n*h.partSize : end]
if err := msg.RespondMsg(part); err != nil && logger != nil {
logger.Printf("%s.%s: could not send part %d of %s: %v", seat, verb, n, id, err)
}
}
func parsePage(v string) (string, int, error) {
id, n, ok := strings.Cut(strings.TrimSpace(v), " ")
if !ok {
return "", 0, fmt.Errorf("%s is \"<call> <part>\", not %q", PageHeader, v)
}
part, err := strconv.Atoi(strings.TrimSpace(n))
if err != nil {
return "", 0, fmt.Errorf("%s is \"<call> <part>\", not %q", PageHeader, v)
}
return id, part, nil
}
// pagedIn reads a reply's `paged`, when it is the answer of one too large for a message.
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
}
// partTries is how often one part is asked before the paging is given up: during a handover a part can
// reach the controller that did not hold it.
const partTries = 3
// Whole is a reply's whole answer: the reply itself, or — when it says its answer was paged — every
// part asked of the same subject and joined. A part that cannot be had is an error naming it, never a
// shorter answer.
func Whole(ctx context.Context, conn *nats.Conn, subject string, data []byte) ([]byte, error) {
p, ok := pagedIn(data)
if !ok {
return data, nil
}
whole := 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 = askPart(ctx, conn, subject, p, n); err == nil {
break
}
}
if err != nil {
return nil, fmt.Errorf("the answer was %d bytes, paged as %s in %d parts, and part %d could not be had: %w",
p.Bytes, p.Call, p.Parts, n, err)
}
whole = append(whole, part...)
}
if len(whole) != 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(whole))
}
return whole, nil
}
func askPart(ctx context.Context, conn *nats.Conn, 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(`{}`)
asking, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
reply, err := conn.RequestMsgWithContext(asking, ask)
if err != nil {
return nil, err
}
if got := reply.Header.Get(PartHeader); got != fmt.Sprintf("%d/%d", n, p.Parts) {
var r Answer
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
}
+101
View File
@@ -0,0 +1,101 @@
package link
import (
"context"
"encoding/json"
"strings"
"testing"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/testbus"
)
// An answer too large for one message is recorded as paged, and a call's record says so: `calls` no
// longer reads "answered" for an answer nobody received (novox/hq issue 314).
func TestAPagedAnswerIsSaidInCalls(t *testing.T) {
l, a := NewCallLog(), newAnswers(t)
big := strings.Repeat("x", 5000)
l.serveCallWithin("mesh-controller", "conditions", nil, "_INBOX.operator.ABCDEFGHIJKLMNOPQRSTUV",
func(context.Context, json.RawMessage) (any, error) { return big, nil }, a.respond, 2048, nil)
got := a.only()
said, _ := got["error"].(string)
paged, _ := got["paged"].(map[string]any)
if !strings.Contains(said, "more than the bus carries") || paged == nil || paged["parts"].(float64) < 3 {
t.Fatalf("answered %v", got)
}
recent, _ := l.Recent()
if len(recent) != 1 || recent[0].Paged == "" || strings.Contains(string(recent[0].Answer), "in full") {
t.Fatalf("kept %+v", recent[0])
}
}
// An answer the bus refuses anyway is said to its caller and to `calls`, never only to the journal.
func TestAnAnswerTheBusRefusesIsSaid(t *testing.T) {
l := NewCallLog()
var sent [][]byte
respond := func(body []byte) error {
if len(body) > 1000 {
return nats.ErrMaxPayload
}
sent = append(sent, body)
return nil
}
l.serveCallWithin("mesh-controller", "status", nil, "_INBOX.x.1",
func(context.Context, json.RawMessage) (any, error) { return strings.Repeat("y", 4000), nil }, respond, 0, nil)
if len(sent) != 1 || !strings.Contains(string(sent[0]), "would not carry it") {
t.Fatalf("the caller was told %q", sent)
}
recent, _ := l.Recent()
if recent[0].Refused == "" {
t.Fatalf("calls does not say the answer was not sent: %+v", recent[0])
}
}
// A part goes only to the caller its answer was paged to, on the verb that answered it.
func TestAPartGoesOnlyToItsCaller(t *testing.T) {
conn, err := nats.Connect(testbus.URL(t))
if err != nil {
t.Fatal(err)
}
defer conn.Close()
l := NewCallLog()
c := l.begin("mesh-controller", "status", nil, "_INBOX.operator.ABCDEFGHIJKLMNOPQRSTUV")
l.page(c, []byte(strings.Repeat("z", 5000)), 2048)
sub, err := conn.Subscribe("pages.test", func(m *nats.Msg) { l.servePage("mesh-controller", "status", m, nil) })
if err != nil {
t.Fatal(err)
}
defer sub.Unsubscribe()
_, err = askPart(context.Background(), conn, "pages.test", Paged{Call: c.ID, Parts: 3}, 0)
if err == nil || !strings.Contains(err.Error(), "another caller") {
t.Fatalf("a part was handed to a caller it was not paged to: %v", err)
}
_, err = askPart(context.Background(), conn, "pages.test", Paged{Call: "call-0-0", Parts: 3}, 0)
if err == nil || !strings.Contains(err.Error(), "not held here") {
t.Fatalf("a part of nothing held: %v", err)
}
}
// A part that cannot be had fails the whole answer, naming it — never a shorter answer.
func TestAMissingPartFailsTheWhole(t *testing.T) {
conn, err := nats.Connect(testbus.URL(t))
if err != nil {
t.Fatal(err)
}
defer conn.Close()
sub, err := conn.Subscribe("pages.none", func(m *nats.Msg) { _ = m.Respond([]byte(`{"error":"gone"}`)) })
if err != nil {
t.Fatal(err)
}
defer sub.Unsubscribe()
first, _ := json.Marshal(map[string]any{"error": "too large", "paged": Paged{Call: "call-1-1", Bytes: 10, Parts: 2}})
_, err = Whole(context.Background(), conn, "pages.none", first)
if err == nil || !strings.Contains(err.Error(), "part 0") || !strings.Contains(err.Error(), "gone") {
t.Fatalf("a missing part read as %v", err)
}
if plain, err := Whole(context.Background(), conn, "pages.none", []byte(`{"result":1}`)); err != nil ||
string(plain) != `{"result":1}` {
t.Fatalf("an answer that fits was changed: %s %v", plain, err)
}
}
+10 -2
View File
@@ -446,8 +446,12 @@ func AskSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string,
case err != nil:
return Answer{}, err
}
data, err := Whole(ctx, conn, subject, reply.Data)
if err != nil {
return Answer{}, fmt.Errorf("%s answered %s.%s too large for one message: %w", node, seat, verb, err)
}
var answer Answer
if err := json.Unmarshal(reply.Data, &answer); err != nil {
if err := json.Unmarshal(data, &answer); err != nil {
return Answer{}, fmt.Errorf("%s answered %s.%s with something unreadable: %w", node, seat, verb, err)
}
return answer, nil
@@ -527,8 +531,12 @@ func AskMeshSeatTool(ctx context.Context, conn *nats.Conn, seat, verb string, ar
case err != nil:
return Answer{}, err
}
data, err := Whole(ctx, conn, subject, reply.Data)
if err != nil {
return Answer{}, fmt.Errorf("%s answered %s too large for one message: %w", seat, verb, err)
}
var answer Answer
if err := json.Unmarshal(reply.Data, &answer); err != nil {
if err := json.Unmarshal(data, &answer); err != nil {
return Answer{}, fmt.Errorf("%s answered %s with something unreadable: %w", seat, verb, err)
}
return answer, nil
+80
View File
@@ -0,0 +1,80 @@
package link
import (
"bytes"
"context"
"encoding/json"
"strings"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/testbus"
)
// novox/hq issue 314, replayed with only what the link had before its fix, so it can be laid over the
// older commit. On 2026-10-08 the controller's `conditions` answered 689 stalled deliveries, each an
// open condition with its evidence, and `status` led with them: more than the bus carries in one message
// (its max_payload, a mebibyte). The client library refused the answer before it left the controller —
// `nats: maximum payload exceeded`, in the controller's journal only — and every caller, the console and
// the operator's channel among them, waited out its timeout and read that the controller did not answer,
// while `calls` said each call answered in 130 ms. An answer larger than one message arrives whole to a
// caller that pages, and is said, in one message that fits, to one that does not — never silence.
func TestReplay314(t *testing.T) {
conn, err := nats.Connect(testbus.URL(t))
if err != nil {
t.Fatal(err)
}
defer conn.Close()
limit := conn.MaxPayload()
// Three times what one message carries: the size of the answers of that morning.
type finding struct {
Key, Summary string
}
var large []finding
for len(large)*200 < int(3*limit) {
large = append(large, finding{Key: "delivery.d-" + strings.Repeat("7", 8) + ".stalled",
Summary: strings.Repeat("held past its bound ", 9)})
}
stop, err := OverNATS{Conn: conn}.ServeSeatTools("replay-314", map[string]ToolHandler{
"conditions": func(context.Context, json.RawMessage) (any, error) { return large, nil },
}, nil)
if err != nil {
t.Fatal(err)
}
defer stop()
t.Run("a caller of the controller's own reads it whole", func(t *testing.T) {
answer, err := AskMeshSeatTool(context.Background(), conn, "replay-314", "conditions", map[string]any{},
5*time.Second)
if err != nil {
t.Fatalf("an answer of about %d bytes did not arrive: %v", 3*limit, err)
}
if answer.Error != "" {
t.Fatalf("answered an error: %.300s", answer.Error)
}
var got []finding
if err := json.Unmarshal(answer.Result, &got); err != nil || len(got) != len(large) {
t.Fatalf("answered %d findings of %d (%v)", len(got), len(large), err)
}
})
t.Run("a caller that does not page is told, not left waiting", func(t *testing.T) {
reply, err := conn.Request(SeatToolSubject("replay-314", "conditions"), []byte(`{}`), 5*time.Second)
if err != nil {
t.Fatalf("no answer at all: %v — the caller waits out its timeout and reads silence", err)
}
if int64(len(reply.Data)) > limit {
t.Fatalf("an answer of %d bytes, more than the bus carries", len(reply.Data))
}
var r struct {
Error string `json:"error"`
}
if json.Unmarshal(reply.Data, &r) != nil || !strings.Contains(r.Error, "more than the bus carries") ||
!bytes.Contains(reply.Data, []byte("ask for less")) && !bytes.Contains(reply.Data, []byte("Ask for less")) {
t.Fatalf("the answer does not say it was too large and what to do: %.400s", reply.Data)
}
})
}
+10 -2
View File
@@ -69,10 +69,18 @@ func (b OverNATS) serveTools(seat string, subjectOf func(string) string, handler
subject := subjectOf(verb)
bind := func() (*nats.Subscription, error) {
return b.Conn.QueueSubscribe(subject, "seat."+seat, func(msg *nats.Msg) {
if msg.Header != nil && msg.Header.Get(PageHeader) != "" {
// A part of an answer too large for one message, asked by the caller it was paged
// to (novox/hq issue 314): not a call of its own.
go Calls.servePage(seat, verb, msg, logger)
return
}
// Its own goroutine per call: a slow `push` must not hold up a `status` asked beside it,
// and the library would otherwise run handlers one after another. Answered once, within
// AnswerWithin, and kept (novox/hq issue 265).
go Calls.serveCall(seat, verb, json.RawMessage(msg.Data), msg.Reply, handle, msg.Respond, logger)
// AnswerWithin, and kept (novox/hq issue 265). Within what the bus carries in one
// message, or paged (issue 314).
go Calls.serveCallWithin(seat, verb, json.RawMessage(msg.Data), msg.Reply, handle, msg.Respond,
b.Conn.MaxPayload(), logger)
})
}
sub, err := bind()