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.
303 lines
11 KiB
Go
303 lines
11 KiB
Go
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
|
|
}
|