From 175b28ee421f68e15a999fd6a8b6d763ae79520c Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 11:26:00 +0200 Subject: [PATCH] Page an answer larger than one message of the bus, and keep overviews brief (hq issue 314) 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. --- cmd/mesh-controller/conditions.go | 48 ++- cmd/mesh-controller/conditions_brief_test.go | 108 +++++++ cmd/mesh-controller/doctor.go | 79 ++++- cmd/mesh-controller/readable.go | 11 +- cmd/mesh-controller/seatverbs.go | 17 +- internal/catalogue/verbs.go | 1 + internal/link/calls.go | 49 ++- internal/link/pages.go | 302 +++++++++++++++++++ internal/link/pages_test.go | 101 +++++++ internal/link/queue.go | 12 +- internal/link/replay314_test.go | 80 +++++ internal/link/seattools.go | 12 +- 12 files changed, 793 insertions(+), 27 deletions(-) create mode 100644 cmd/mesh-controller/conditions_brief_test.go create mode 100644 internal/link/pages.go create mode 100644 internal/link/pages_test.go create mode 100644 internal/link/replay314_test.go diff --git a/cmd/mesh-controller/conditions.go b/cmd/mesh-controller/conditions.go index 8795174c..458382fd 100644 --- a/cmd/mesh-controller/conditions.go +++ b/cmd/mesh-controller/conditions.go @@ -144,11 +144,9 @@ func listConditions(ctx context.Context, args []string) error { } } if *asJSON { - if out == nil { - out = []conditions.Condition{} - } - return printJSON(map[string]any{"conditions": out, "open": len(open), - "note": "urgent first, then oldest first; a condition clears when observation says so, never by hand"}) + return printJSON(map[string]any{"conditions": inBrief(out), "open": len(open), "counted": counted(out), + "note": "urgent first, then oldest first; a condition clears when observation says so, never by hand; " + + "each with its newest evidence — `conditions key=` gives one whole"}) } if len(out) == 0 { if len(open) == 0 { @@ -164,6 +162,46 @@ func listConditions(ctx context.Context, args []string) error { return nil } +// briefEvidence is how many observations a condition carries in a list of them: its newest. The store +// keeps ten, and a list of hundreds of conditions with ten each was the larger half of an answer more +// than the bus carries in one message (novox/hq issue 314). One condition whole is `conditions key=`. +const briefEvidence = 1 + +// conditionInBrief is a condition as a list carries it: its newest evidence, with how much more there is +// and where to read it. Everything else of it, so a reader of the list — the operator's channel among +// them — reads the same condition it always did. +type conditionInBrief struct { + conditions.Condition + // EvidenceKept is how many observations the condition keeps, when the list shows fewer. + EvidenceKept int `json:"evidence-kept,omitempty"` + More string `json:"more,omitempty"` +} + +// inBrief is a list of conditions in brief; never null. +func inBrief(list []conditions.Condition) []conditionInBrief { + out := make([]conditionInBrief, 0, len(list)) + for _, c := range list { + b := conditionInBrief{Condition: c} + if len(c.Evidence) > briefEvidence { + b.EvidenceKept = len(c.Evidence) + b.Evidence = c.Evidence[:briefEvidence:briefEvidence] + b.More = c.Show() + } + out = append(out, b) + } + return out +} + +// counted is how many of a list each source raised of each kind — `D14 stalled: 689` — so a list of +// hundreds says its shape before its lines. +func counted(list []conditions.Condition) map[string]int { + out := map[string]int{} + for _, c := range list { + out[c.Source+" "+c.Kind]++ + } + return out +} + // concerns says whether a condition is about a machine: it names it, or its key does. func concerns(c conditions.Condition, machine string) bool { if c.Subject.Machine == machine || slices.Contains(c.Subject.Also, machine) { diff --git a/cmd/mesh-controller/conditions_brief_test.go b/cmd/mesh-controller/conditions_brief_test.go new file mode 100644 index 00000000..8fae4898 --- /dev/null +++ b/cmd/mesh-controller/conditions_brief_test.go @@ -0,0 +1,108 @@ +package main + +import ( + "encoding/json" + "fmt" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/conditions" +) + +// stalledConditions are n open conditions like those of 2026-10-08: one per delivery held past its bound, +// each keeping all its evidence. +func stalledConditions(n int) []conditions.Condition { + at := time.Date(2026, 10, 8, 9, 0, 0, 0, time.UTC) + var out []conditions.Condition + for i := 0; i < n; i++ { + id := fmt.Sprintf("d-%012x", 0x51a11ed+i) + lines := stalledObservations([]stalledLine{{ID: id, State: "checking", For: "13h2m", Bound: "45m", + H2: "close: recheck", Says: "the check was asked and nothing has answered"}}) + o := lines[0] + c := conditions.Condition{Key: o.Key(), Kind: o.Kind, Subject: conditions.Subject{Scope: o.Scope, ID: o.ID}, + Severity: o.Severity, Summary: o.Summary, Source: probeDeliveriesID, Raised: at, LastObserved: at, + Observations: 160, Count: 1, Resolver: conditions.ResolverSelf} + for e := 0; e < conditions.KeptEvidence; e++ { + c.Evidence = append(c.Evidence, conditions.Evidence{At: at.Add(time.Duration(e) * 5 * time.Minute), + Said: "checking for " + (13*time.Hour + time.Duration(e)*time.Minute).String()}) + } + out = append(out, c) + } + return out +} + +// novox/hq issue 314: a list of hundreds of open conditions carries each one's newest evidence, with +// how much more it keeps and where to read it, and the verb's answer says it once — as data, not again +// as the text the command printed. On 2026-10-08, 689 stalled deliveries made that answer more than the +// bus carries in one message; now it is less than half of one. +func TestAListOfConditionsCarriesEachInBrief(t *testing.T) { + list := stalledConditions(689) + printed, _ := json.MarshalIndent(map[string]any{"conditions": inBrief(list), "open": len(list), + "counted": counted(list)}, "", " ") + var parsed any + if err := json.Unmarshal(printed, &parsed); err != nil { + t.Fatal(err) + } + wire, _ := json.Marshal(map[string]any{"result": verbAnswer{OK: true, Answer: parsed, + Output: "its answer, as data, is `answer`\n"}}) + before, _ := json.MarshalIndent(map[string]any{"conditions": list}, "", " ") + t.Logf("689 stalled deliveries: %d bytes on the wire, where the whole list printed was %d and was sent twice", + len(wire), len(before)) + if len(wire) >= 512<<10 { + t.Fatalf("the answer of 689 conditions is %d bytes, not under half of one message of the bus", len(wire)) + } + var doc struct { + Conditions []conditionInBrief `json:"conditions"` + Counted map[string]int `json:"counted"` + } + if err := json.Unmarshal(printed, &doc); err != nil { + t.Fatal(err) + } + c := doc.Conditions[0] + if len(doc.Conditions) != 689 || len(c.Evidence) != 1 || c.EvidenceKept != conditions.KeptEvidence || + c.Key != list[0].Key || c.Summary != list[0].Summary || !strings.Contains(c.More, "key="+c.Key) { + t.Fatalf("a condition in brief reads %+v", c) + } + if doc.Counted["D14 stalled"] != 689 { + t.Fatalf("counted %v", doc.Counted) + } + if one := inBrief(list[:1])[0]; len(list[0].Evidence) != conditions.KeptEvidence || one.Evidence[0] != list[0].Evidence[0] { + t.Fatal("the brief changed the condition it was made from, or kept other than the newest") + } +} + +// The self-check's verdict carries at most overviewFindings of one probe's findings, with their count and +// how to read them all; `doctor probe=` gives that probe whole (novox/hq issue 314). +func TestAVerdictCarriesAProbesFindingsInBrief(t *testing.T) { + var found []string + for i := 0; i < 689; i++ { + found = append(found, fmt.Sprintf("delivery d-%012x has been checking for 13h, past its bound of 45m", i)) + } + run := doctorRun{Run: "run-1", Probes: []probeVerdict{{ID: "D14", Verdict: "failed", Found: found}, + {ID: "D1", Verdict: "passed"}}} + answer, err := verdictAnswer(run, time.Now(), "") + if err != nil { + t.Fatal(err) + } + brief := answer["run"].(doctorRun).Probes[0] + if len(brief.Found) != overviewFindings || brief.FoundCount != 689 || !strings.Contains(brief.More, "probe=D14") { + t.Fatalf("the verdict carries %d findings of D14, count %d, more %q", len(brief.Found), brief.FoundCount, brief.More) + } + if len(run.Probes[0].Found) != 689 { + t.Fatal("the overview cut the run it was made from") + } + whole, err := verdictAnswer(run, time.Now(), "d14") + if err != nil { + t.Fatal(err) + } + if p := whole["run"].(doctorRun).Probes; len(p) != 1 || len(p[0].Found) != 689 || p[0].More != "" { + t.Fatalf("doctor probe=D14 answered %+v", p) + } + if _, err := verdictAnswer(run, time.Now(), "D99"); err == nil { + t.Fatal("a probe the run does not have was answered") + } + if !strings.Contains(doctorText(answer), "probe=D14` gives them all") { + t.Fatalf("the text does not say how to read them all:\n%s", doctorText(answer)) + } +} diff --git a/cmd/mesh-controller/doctor.go b/cmd/mesh-controller/doctor.go index 66158437..ac5da269 100644 --- a/cmd/mesh-controller/doctor.go +++ b/cmd/mesh-controller/doctor.go @@ -151,6 +151,50 @@ type probeVerdict struct { Unconfirmed []string `json:"unconfirmed,omitempty"` Error string `json:"error,omitempty"` Took string `json:"took,omitempty"` + // FoundCount and More are an overview's: how many findings the probe had when the overview shows + // fewer, and how to read them all (novox/hq issue 314). Never in a heartbeat. + FoundCount int `json:"found-count,omitempty"` + More string `json:"more,omitempty"` +} + +// overviewFindings is how many findings of one probe an overview carries. On 2026-10-08 D14 found 689 +// stalled deliveries, and an overview carrying every one of them, beside the conditions they raised, +// was more than the bus carries in one message (novox/hq issue 314); one probe's verdict whole is +// `doctor probe=`. +const overviewFindings = 20 + +// inOverview is a run as an overview carries it: each probe's findings at most overviewFindings, with +// their count and how to read them all. The run itself is not changed. +func inOverview(run doctorRun) doctorRun { + probes := make([]probeVerdict, len(run.Probes)) + for i, p := range run.Probes { + total := len(p.Found) + len(p.Unconfirmed) + if total > overviewFindings { + p.FoundCount = len(p.Found) + if len(p.Found) > overviewFindings { + p.Found = p.Found[:overviewFindings:overviewFindings] + } + if room := overviewFindings - len(p.Found); len(p.Unconfirmed) > room { + p.Unconfirmed = p.Unconfirmed[:room:room] + } + p.More = fmt.Sprintf("%d found and %d unconfirmed, %d of them shown — `doctor probe=%s` gives them all", + p.FoundCount, total-p.FoundCount, len(p.Found)+len(p.Unconfirmed), p.ID) + } + probes[i] = p + } + run.Probes = probes + return run +} + +// oneProbe is a run with only one probe's verdict, whole. +func oneProbe(run doctorRun, id string) (doctorRun, error) { + for _, p := range run.Probes { + if strings.EqualFold(p.ID, id) { + run.Probes = []probeVerdict{p} + return run, nil + } + } + return doctorRun{}, fmt.Errorf("the run %s has no probe %q — `doctor probes=true` lists them", run.Run, id) } // doctorCounts are a run's verdicts, counted. @@ -386,12 +430,13 @@ func doctorCommand(ctx context.Context, args []string) error { } set := flag.NewFlagSet("doctor", flag.ContinueOnError) asJSON := set.Bool("json", false, "as data") + probe := set.String("probe", "", "one probe's verdict, with every finding") if rest, err := parseAround(set, args); err != nil { return err } else if len(rest) > 0 { - return errors.New("doctor [run|probes|signals] [--json]") + return errors.New("doctor [run|probes|signals] [--probe ] [--json]") } - answer, err := doctorAnswer(ctx, sub) + answer, err := doctorAnswer(ctx, sub, *probe) if err != nil { return err } @@ -403,12 +448,18 @@ func doctorCommand(ctx context.Context, args []string) error { } // doctorAnswer is what the verb answers, as data. -func doctorAnswer(ctx context.Context, sub string) (any, error) { +// +// A verdict is an overview — each probe's findings at most overviewFindings — unless probe names one, +// which is then answered alone and whole (novox/hq issue 314). +func doctorAnswer(ctx context.Context, sub, probe string) (any, error) { + if probe != "" && sub != "" && sub != "run" { + return nil, fmt.Errorf("probe names one probe of a verdict; %s has none", sub) + } switch sub { case "": if doctorFrom != nil { if run := doctorFrom.lastRun(); run != nil { - return verdictAnswer(*run, time.Now()), nil + return verdictAnswer(*run, time.Now(), probe) } return nil, fmt.Errorf("the self-check has not finished its first run yet: it runs %s after the "+ "controller starts, then every %s — `doctor run` runs it now", doctorFirstAfter, doctorEvery) @@ -417,7 +468,7 @@ func doctorAnswer(ctx context.Context, sub string) (any, error) { if err != nil { return nil, err } - return verdictAnswer(run, time.Now()), nil + return verdictAnswer(run, time.Now(), probe) case "run": d := doctorFrom if d == nil { @@ -428,7 +479,7 @@ func doctorAnswer(ctx context.Context, sub string) (any, error) { defer closeIt() d = local } - return verdictAnswer(d.runOnce(ctx, "asked by "+link.Caller()), time.Now()), nil + return verdictAnswer(d.runOnce(ctx, "asked by "+link.Caller()), time.Now(), probe) case "probes": return probesAnswer(), nil case "signals": @@ -442,11 +493,20 @@ func doctorAnswer(ctx context.Context, sub string) (any, error) { } // verdictAnswer is a run as the verb answers it, with its age. -func verdictAnswer(run doctorRun, now time.Time) map[string]any { +func verdictAnswer(run doctorRun, now time.Time, probe string) (map[string]any, error) { + if probe != "" { + one, err := oneProbe(run, probe) + if err != nil { + return nil, err + } + run = one + } else { + run = inOverview(run) + } return map[string]any{"run": run, "age": now.Sub(run.At).Round(time.Second).String(), "note": "a probe that could not run is never a pass; each failure is an open condition until a run passes it, " + "and one a single look can be wrong about — an unanswered question, a slow answer — is raised when two " + - "runs in a row see it"} + "runs in a row see it"}, nil } // probesAnswer is the registry. @@ -515,6 +575,9 @@ func doctorText(answer any) string { for _, f := range p.Unconfirmed { fmt.Fprintf(&b, " unconfirmed, raised if the next run sees it too: %s\n", f) } + if p.More != "" { + fmt.Fprintf(&b, " … %s\n", p.More) + } if p.Error != "" { fmt.Fprintf(&b, " %s\n", p.Error) } diff --git a/cmd/mesh-controller/readable.go b/cmd/mesh-controller/readable.go index e7f8d8a7..a84a1760 100644 --- a/cmd/mesh-controller/readable.go +++ b/cmd/mesh-controller/readable.go @@ -91,8 +91,8 @@ type meshStatus struct { // Conditions is every open condition, urgent first and then oldest first (novox/hq to-be 45 §2): // what is wrong, as the watchdogs, the self-check and the providers say it. Always present — an // empty list is "none open" — unless they could not be read, which ConditionsUnread says. - Conditions []conditions.Condition `json:"conditions"` - ConditionsUnread string `json:"conditionsUnread,omitempty"` + Conditions []conditionInBrief `json:"conditions"` + ConditionsUnread string `json:"conditionsUnread,omitempty"` // Failing is every consumer a provider says it keeps failing (novox/hq ADR 0224): the open // conditions of that kind, carried here as well because ADR 0224 names this field. Absent when no // provider says so. A document without it called the mesh well while the identity provider @@ -235,10 +235,9 @@ func statusAsJSON(asked answers) ([]byte, error) { out.Unheld = asked.unheld out.HandActsThisWeek, out.HandActsUnread = asked.handActs, asked.handActsUnread out.HealsThisWeek, out.HealsUnread = asked.heals, asked.healsUnread - out.Conditions, out.ConditionsUnread = asked.conditions, asked.conditionsUnread - if out.Conditions == nil { - out.Conditions = []conditions.Condition{} - } + // In brief, as `conditions` lists them: status leads with every open condition, and their whole + // evidence is each one's own (novox/hq issue 314). + out.Conditions, out.ConditionsUnread = inBrief(asked.conditions), asked.conditionsUnread out.Failing = providerStandings(asked.conditions) out.Overflowing = asked.overflowing for name := range asked.refused { diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 1eb0ccfa..b34cd634 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -641,6 +641,9 @@ func (a *verbArguments) commandLine() ([]string, error) { if which > 1 { return nil, errors.New("doctor answers one of run, probes or signals at a time") } + if p := str("probe"); p != "" { + argv = append(argv, "--probe", p) + } return append(argv, "--json"), nil case "rotate": if p := str("provision"); p != "" { @@ -771,6 +774,11 @@ var jsonVerbs = map[string]bool{"status": true, "seats": true, "plan": true, "co // The delivery's owner's verbs answer JSON where they read (plan, order, check, walks) — novox/hq ADR 0239. "delivery": true} +// overviewVerbs are the JSON verbs whose answer is said once, as data, and not again as the text the +// command printed: the overviews, which grow with the mesh (novox/hq issue 314). Every reader of theirs +// reads `answer` first. +var overviewVerbs = map[string]bool{"status": true, "conditions": true, "doctor": true} + // repairingCommand names a command line that repairs by hand, and so says why: a push, a plan stopped // or closed, a consumer re-made (novox/hq to-be 45 §7). Empty for any other. func repairingCommand(argv []string) string { @@ -834,6 +842,11 @@ func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { var parsed any if json.Unmarshal(bytes.TrimSpace(stdout.Bytes()), &parsed) == nil { answer.Answer = parsed + if overviewVerbs[argv[0]] { + // Once, as data: the same document again as text doubled an answer that already + // outgrew one message of the bus (novox/hq issue 314). + answer.Output = stderr.String() + "its answer, as data, is `answer`\n" + } } } var exit *exec.ExitError @@ -887,10 +900,10 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { return nil, err } sub := "" - if len(argv) > 2 { + if len(argv) > 2 && !strings.HasPrefix(argv[1], "-") { sub = argv[1] } - return doctorAnswer(ctx, sub) + return doctorAnswer(ctx, sub, a.given["probe"]) } return seatTools(), nil } diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index a6601ff6..d5e07ec5 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -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 " + diff --git a/internal/link/calls.go b/internal/link/calls.go index 56101e31..49d13e83 100644 --- a/internal/link/calls.go +++ b/internal/link/calls.go @@ -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 `.`. 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() diff --git a/internal/link/pages.go b/internal/link/pages.go new file mode 100644 index 00000000..0592533d --- /dev/null +++ b/internal/link/pages.go @@ -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: ` `, 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 +} diff --git a/internal/link/pages_test.go b/internal/link/pages_test.go new file mode 100644 index 00000000..844fb099 --- /dev/null +++ b/internal/link/pages_test.go @@ -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) + } +} diff --git a/internal/link/queue.go b/internal/link/queue.go index 55190e18..86c9ab07 100644 --- a/internal/link/queue.go +++ b/internal/link/queue.go @@ -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 diff --git a/internal/link/replay314_test.go b/internal/link/replay314_test.go new file mode 100644 index 00000000..f7d57965 --- /dev/null +++ b/internal/link/replay314_test.go @@ -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) + } + }) +} diff --git a/internal/link/seattools.go b/internal/link/seattools.go index 10dc8dd6..092398d9 100644 --- a/internal/link/seattools.go +++ b/internal/link/seattools.go @@ -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()