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 d299a805..63ebf6d8 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -673,6 +673,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 != "" { @@ -805,6 +808,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 { @@ -870,6 +878,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 @@ -923,10 +936,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 1e637ace..1dd52358 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()