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: ` `, parts counted from zero. PageHeader = "Mesh-Page" // PartHeader is on a part: `/`. 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 \", 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 \" \", not %q", PageHeader, v) } part, err := strconv.Atoi(strings.TrimSpace(n)) if err != nil { return "", 0, fmt.Errorf("%s is \" \", 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 }