From 1e04670052ac6e1cd97ac249fabef3f8f58fb497 Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 17:52:36 +0200 Subject: [PATCH 1/3] Serve the read verbs on the serving controller's own connection, and name every connection Each verb ran as a process that dialled the bus, so hundreds of short connections an hour, all named mesh-controller, hid any client reconnecting in a loop (hq issue 327). D15 now says a user whose connections keep dropping. --- cmd/mesh-builder/main.go | 6 +- cmd/mesh-controller/acting.go | 4 +- cmd/mesh-controller/build.go | 10 +- cmd/mesh-controller/conditions.go | 51 +++---- cmd/mesh-controller/doctor.go | 5 + cmd/mesh-controller/handacts.go | 27 ++-- cmd/mesh-controller/main.go | 19 +++ cmd/mesh-controller/nodes.go | 9 +- cmd/mesh-controller/plain_words.go | 6 + cmd/mesh-controller/plan_retry.go | 9 +- cmd/mesh-controller/push.go | 2 + cmd/mesh-controller/queue.go | 23 ++-- cmd/mesh-controller/reconnects.go | 175 +++++++++++++++++++++++++ cmd/mesh-controller/reconnects_test.go | 156 ++++++++++++++++++++++ cmd/mesh-controller/replay327_test.go | 132 +++++++++++++++++++ cmd/mesh-controller/rollout.go | 5 +- cmd/mesh-controller/seatverbs.go | 87 ++++++++++-- cmd/mesh-controller/servingbus.go | 38 ++++++ internal/broker/jetstream.go | 23 +++- internal/broker/jetstream_test.go | 31 +++++ internal/link/builds_nats.go | 6 + 21 files changed, 744 insertions(+), 80 deletions(-) create mode 100644 cmd/mesh-controller/reconnects.go create mode 100644 cmd/mesh-controller/reconnects_test.go create mode 100644 cmd/mesh-controller/replay327_test.go create mode 100644 cmd/mesh-controller/servingbus.go diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index 5de94ed5..672e6d8e 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -27,6 +27,8 @@ import ( "syscall" "time" + "github.com/nats-io/nats.go" + "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/builder" "github.com/novox/mesh-controller/internal/link" @@ -175,7 +177,9 @@ func dialFor(credential Credential) (*broker.JetStream, string, error) { if !credential.onTheNewBus() { return nil, "", fmt.Errorf("the credential at hand names %q, which is not the mesh's bus", credential.URL) } - js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint) + // Named for what it is, not the controller whose code dials it (novox/hq issue 327). + host, _ := os.Hostname() + js, err := broker.DialPinned(credential.natsURL(), credential.Fingerprint, nats.Name("build agent on "+host)) if err != nil { return nil, "", err } diff --git a/cmd/mesh-controller/acting.go b/cmd/mesh-controller/acting.go index e796a6a7..63f7aea2 100644 --- a/cmd/mesh-controller/acting.go +++ b/cmd/mesh-controller/acting.go @@ -8,6 +8,7 @@ import ( "sync" "time" + "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" "github.com/novox/mesh-controller/internal/broker" @@ -187,7 +188,8 @@ func (a *actor) serveUnderTheLease(ctx context.Context, inv *inventory.Inventory a.mu.Lock() a.serving = true a.mu.Unlock() - js, err := broker.Dial(address) + // Its own connection, held as long as the lease, and named so (novox/hq issue 327). + js, err := broker.Dial(address, nats.Name(broker.ConnectionName+" lease")) if err != nil { return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it to take the "+ "lease: %w", broker.BareAddress(address), err) diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index ec5d2b2f..083139ea 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -879,6 +879,10 @@ func heldBy(ctx context.Context) map[string]string { // the mesh runs on today this needs the controller's own connection, so it is handed one; on the bus // being built it dials, because a build request is a one-shot and holds nothing else. func askOverOn(seat string) (link.Builders, error) { + // The serving controller asks on its own connection (novox/hq issue 327). + if serving := servingBus.Load(); serving != nil { + return link.BuildsOn(serving, seat), nil + } address, err := broker.BusAddress() if err != nil { return nil, err @@ -937,11 +941,7 @@ func buildSeatAmong(entries []inventory.Entry) string { // with a consumer of its own that is gone when this returns, so nothing accumulates in the server // for the reading, and filtered by subject, so one build's lines are all that travel. func buildLog(ctx context.Context, id string) error { - address, err := broker.BusAddress() - if err != nil { - return err - } - js, err := broker.Dial(address) + js, err := aBus() if err != nil { return fmt.Errorf("cannot reach the bus to read a build's log: %w", err) } diff --git a/cmd/mesh-controller/conditions.go b/cmd/mesh-controller/conditions.go index 458382fd..c26b181b 100644 --- a/cmd/mesh-controller/conditions.go +++ b/cmd/mesh-controller/conditions.go @@ -5,6 +5,7 @@ import ( "errors" "flag" "fmt" + "io" "os" "slices" "strconv" @@ -107,18 +108,18 @@ func conditionsCommand(ctx context.Context, args []string) error { } switch sub { case "list": - return listConditions(ctx, args) + return listConditions(ctx, args, os.Stdout) case "show": - return showCondition(ctx, args) + return showCondition(ctx, args, os.Stdout) case "silence": return silenceCondition(ctx, args) case "history": - return conditionHistory(ctx, args) + return conditionHistory(ctx, args, os.Stdout) } return errors.New(conditionsUsage) } -func listConditions(ctx context.Context, args []string) error { +func listConditions(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions", flag.ContinueOnError) scope := set.String("scope", "", "only this scope: "+strings.Join(conditions.Scopes, ", ")) severity := set.String("severity", "", "only urgent, or only warning") @@ -144,20 +145,20 @@ func listConditions(ctx context.Context, args []string) error { } } if *asJSON { - return printJSON(map[string]any{"conditions": inBrief(out), "open": len(open), "counted": counted(out), + return printJSONTo(w, 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 { - fmt.Println("no open conditions") + fmt.Fprintln(w, "no open conditions") } else { - fmt.Printf("none of the %d open condition(s) is about that\n", len(open)) + fmt.Fprintf(w, "none of the %d open condition(s) is about that\n", len(open)) } return nil } for _, line := range conditionLines(out, time.Now()) { - fmt.Println(line) + fmt.Fprintln(w, line) } return nil } @@ -233,7 +234,7 @@ func conditionLines(list []conditions.Condition, now time.Time) []string { return out } -func showCondition(ctx context.Context, args []string) error { +func showCondition(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions show", flag.ContinueOnError) asJSON := set.Bool("json", false, "as data") rest, err := parseAround(set, args) @@ -266,34 +267,34 @@ func showCondition(ctx context.Context, args []string) error { "history --key %s` what became of it", key, key) } if *asJSON { - return printJSON(c) + return printJSONTo(w, c) } now := time.Now() - fmt.Printf("%s %s\n %s\n\n", strings.ToUpper(string(c.Severity)), c.Key, c.Summary) - fmt.Printf(" kind %s\n about %s %s", c.Kind, c.Subject.Scope, c.Subject.ID) + fmt.Fprintf(w, "%s %s\n %s\n\n", strings.ToUpper(string(c.Severity)), c.Key, c.Summary) + fmt.Fprintf(w, " kind %s\n about %s %s", c.Kind, c.Subject.Scope, c.Subject.ID) if c.Subject.Machine != "" { - fmt.Printf(", on %s", c.Subject.Machine) + fmt.Fprintf(w, ", on %s", c.Subject.Machine) } - fmt.Printf("\n raised by %s\n since %s (%s ago), observed %d time(s), last %s ago\n", + fmt.Fprintf(w, "\n raised by %s\n since %s (%s ago), observed %d time(s), last %s ago\n", c.Source, c.Raised.Local().Format("2006-01-02 15:04:05"), roughly(now.Sub(c.Raised)), c.Observations, now.Sub(c.LastObserved).Round(time.Second)) if c.Count > 1 { - fmt.Printf(" raised %d times, each within ten minutes of clearing\n", c.Count) + fmt.Fprintf(w, " raised %d times, each within ten minutes of clearing\n", c.Count) } - fmt.Printf(" resolved by %s\n", resolverWords(c.Resolver)) + fmt.Fprintf(w, " resolved by %s\n", resolverWords(c.Resolver)) if c.Silenced != nil { - fmt.Printf(" silenced until %s by %s: %s\n", c.Silenced.Until.Local().Format("2006-01-02 15:04"), + fmt.Fprintf(w, " silenced until %s by %s: %s\n", c.Silenced.Until.Local().Format("2006-01-02 15:04"), c.Silenced.By, c.Silenced.Why) } if len(c.Tried) > 0 { - fmt.Println("\n tried:") + fmt.Fprintln(w, "\n tried:") for _, t := range c.Tried { - fmt.Printf(" %s %s — %s: %s\n", t.At.Local().Format("2006-01-02 15:04"), orHealer(t.By), t.What, t.Outcome) + fmt.Fprintf(w, " %s %s — %s: %s\n", t.At.Local().Format("2006-01-02 15:04"), orHealer(t.By), t.What, t.Outcome) } } - fmt.Println("\n evidence, newest first:") + fmt.Fprintln(w, "\n evidence, newest first:") for _, e := range c.Evidence { - fmt.Printf(" %s %s\n", e.At.Local().Format("2006-01-02 15:04:05"), e.Said) + fmt.Fprintf(w, " %s %s\n", e.At.Local().Format("2006-01-02 15:04:05"), e.Said) } return nil } @@ -378,7 +379,7 @@ func parseFor(s string) (time.Duration, error) { return d, nil } -func conditionHistory(ctx context.Context, args []string) error { +func conditionHistory(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions history", flag.ContinueOnError) days := set.Int("days", 7, "how many days back, at most 90") key := set.String("key", "", "only this condition") @@ -420,10 +421,10 @@ func conditionHistory(ctx context.Context, args []string) error { if out == nil { out = []conditions.Event{} } - return printJSON(map[string]any{"history": out, "days": *days}) + return printJSONTo(w, map[string]any{"history": out, "days": *days}) } if len(out) == 0 { - fmt.Printf("nothing was raised, changed or cleared in the last %d day(s)\n", *days) + fmt.Fprintf(w, "nothing was raised, changed or cleared in the last %d day(s)\n", *days) return nil } for _, e := range out { @@ -440,7 +441,7 @@ func conditionHistory(ctx context.Context, args []string) error { line += " — " + e.Why } } - fmt.Println(line) + fmt.Fprintln(w, line) } return nil } diff --git a/cmd/mesh-controller/doctor.go b/cmd/mesh-controller/doctor.go index ac5da269..54f533ba 100644 --- a/cmd/mesh-controller/doctor.go +++ b/cmd/mesh-controller/doctor.go @@ -116,6 +116,11 @@ var probeRegistry = []probe{ {ID: probeDeliveriesID, Asserts: "no delivery is held past its state's bound unsaid: mesh-delivery's " + "`stalled`, each with the transition its table lets healer H2 take", From: "ADR 0239", Kind: kindDeliveryStalled, Phase: 3, run: probeDeliveries}, + // A client of the bus reconnecting in a loop (novox/hq issue 327), from the server's record of closed + // connections, which the bus's own module reads. + {ID: probeReconnectsID, Asserts: "no user of the bus had its connection dropped more than twelve times in the " + + "last hour: the bus module's nats_closed_connections", From: "issue 327", Kind: kindBusReconnects, + Phase: 1, run: probeReconnects}, {ID: "DW", Asserts: "the watchdogs of the signals table ran within three of their intervals", From: "ADR 0227 rule 6: the watchers are watched", Kind: "watchdogs-silent", Phase: 1, run: probeWatchdogs}, // The core's health definitions (novox/hq to-be 45 §8, ADR 0236): what a core component's new build is diff --git a/cmd/mesh-controller/handacts.go b/cmd/mesh-controller/handacts.go index 5d5ad5d7..0cc4e3c0 100644 --- a/cmd/mesh-controller/handacts.go +++ b/cmd/mesh-controller/handacts.go @@ -6,6 +6,7 @@ import ( "errors" "flag" "fmt" + "io" "os" "slices" "sort" @@ -14,7 +15,6 @@ import ( "github.com/nats-io/nats.go" - "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/conditions" "github.com/novox/mesh-controller/internal/link" ) @@ -154,14 +154,10 @@ func onTheBus(f func(*nats.Conn) error) error { if handActConn != nil { return f(handActConn) } - address, err := broker.BusAddress() + js, err := aBus() if err != nil { return err } - js, err := broker.Dial(address) - if err != nil { - return fmt.Errorf("cannot reach the bus: %w", err) - } defer js.Close() return f(js.Conn()) } @@ -252,6 +248,11 @@ func handActCommand(ctx context.Context, args []string) error { if len(args) > 0 && args[0] == "list" { args = args[1:] } + return listHandActs(ctx, args, os.Stdout) +} + +// listHandActs is `hand-acts`: what was done by hand lately, and the causes done more than once. +func listHandActs(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("hand-acts", flag.ContinueOnError) days := set.Int("days", 14, "how many days back") asJSON := set.Bool("json", false, "as data") @@ -270,27 +271,27 @@ func handActCommand(ctx context.Context, args []string) error { if err != nil { return err } - fmt.Println(string(body)) + fmt.Fprintln(w, string(body)) return nil } if len(acts) == 0 { - fmt.Printf("nothing was done by hand in the last %d day(s)\n", *days) + fmt.Fprintf(w, "nothing was done by hand in the last %d day(s)\n", *days) return nil } for i := len(acts) - 1; i >= 0; i-- { a := acts[i] - fmt.Printf("%s %s %s %s\n by %s — %s (cause: %s", a.At.Local().Format("2006-01-02 15:04"), a.ID, + fmt.Fprintf(w, "%s %s %s %s\n by %s — %s (cause: %s", a.At.Local().Format("2006-01-02 15:04"), a.ID, a.Verb, strings.Join(a.Args, " "), a.By, a.Why, a.Cause) if a.Condition != "" { - fmt.Printf(", condition %s", a.Condition) + fmt.Fprintf(w, ", condition %s", a.Condition) } - fmt.Println(")") + fmt.Fprintln(w, ")") if pushedRecorded(a) { carried := strings.Join(a.Carried, "; ") if carried == "" { carried = recordedBefore[a.ID] } - fmt.Printf(" a push of recorded builds, no repair: %s\n", carried) + fmt.Fprintf(w, " a push of recorded builds, no repair: %s\n", carried) } } if len(repeated) > 0 { @@ -299,7 +300,7 @@ func handActCommand(ctx context.Context, args []string) error { causes = append(causes, fmt.Sprintf("%s ×%d", c, n)) } sort.Strings(causes) - fmt.Printf("\ndone by hand more than once in a fortnight — a healer is wanted (to-be 45 S15): %s\n", + fmt.Fprintf(w, "\ndone by hand more than once in a fortnight — a healer is wanted (to-be 45 S15): %s\n", strings.Join(causes, ", ")) } return nil diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index ac5ab1a3..b480107b 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -15,6 +15,7 @@ import ( "os/signal" "syscall" + "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/identity" "github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/licences" @@ -53,6 +54,9 @@ func run() error { return fmt.Errorf("no command given") } + // Every connection this process dials says what it is (novox/hq issue 327). + broker.ConnectionName = connectionName(args[0], os.Getenv(verbVar)) + ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer stop() // Whatever this process holds of the controller's lease is given back as it ends (novox/hq to-be @@ -416,3 +420,18 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error { statusFrom.nudge() return nil } + +// verbVar carries the seat verb a command runs for, from the serving controller to the process it starts. +const verbVar = "MESH_VERB" + +// connectionName is what this process's connections say they are in the bus's list (novox/hq issue 327): +// the serving controller, a verb's own process and which verb, or a command run at a shell and which. +func connectionName(command, verb string) string { + switch { + case command == "serve": + return "mesh-controller serving" + case verb != "": + return "mesh-controller verb " + verb + } + return "mesh-controller command " + command +} diff --git a/cmd/mesh-controller/nodes.go b/cmd/mesh-controller/nodes.go index 6a50adc4..579c01de 100644 --- a/cmd/mesh-controller/nodes.go +++ b/cmd/mesh-controller/nodes.go @@ -7,7 +7,9 @@ import ( "errors" "flag" "fmt" + "io" "net" + "os" "strings" "time" @@ -684,11 +686,14 @@ func orNotReported(s string) string { } // printJSON prints a value as indented JSON, the shape every `--json` answers in. -func printJSON(v any) error { +func printJSON(v any) error { return printJSONTo(os.Stdout, v) } + +// printJSONTo is printJSON to a writer of the caller's. +func printJSONTo(w io.Writer, v any) error { body, err := json.MarshalIndent(v, "", " ") if err != nil { return err } - fmt.Println(string(body)) + fmt.Fprintln(w, string(body)) return nil } diff --git a/cmd/mesh-controller/plain_words.go b/cmd/mesh-controller/plain_words.go index 0b333b26..81a568eb 100644 --- a/cmd/mesh-controller/plain_words.go +++ b/cmd/mesh-controller/plain_words.go @@ -543,6 +543,12 @@ var plainWordings = map[string]func(conditions.Observation) words{ Explanation: "The bus reports a listener too slow to keep up, so messages to it are late.", Resolved: "The listener keeps up again"} }), + kindBusReconnects: worded(func(o conditions.Observation) words { + return words{Headline: "A client keeps losing the bus", + Explanation: "One of the mesh's clients lost its connection to the bus again and again in the last hour. " + + "While it reconnects, what it says and what it is asked waits.", + Resolved: "Resolved: the client stays connected"} + }), "max-deliveries": worded(func(o conditions.Observation) words { return words{Headline: "A listener gave up on messages", Needs: "deliver them again or drop them, from the mesh MCP server.", diff --git a/cmd/mesh-controller/plan_retry.go b/cmd/mesh-controller/plan_retry.go index dc8a2b49..7a2b16af 100644 --- a/cmd/mesh-controller/plan_retry.go +++ b/cmd/mesh-controller/plan_retry.go @@ -9,7 +9,6 @@ import ( "strings" "time" - "github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/inventory" "github.com/novox/mesh-controller/internal/link" ) @@ -98,11 +97,9 @@ func buildSeatPause(ctx context.Context, inv *inventory.Inventory, plans []inven if len(holders) == 0 { return pauseView{} } - address, err := broker.BusAddress() - if err != nil { - return pauseView{} - } - js, err := broker.Dial(address) + // On the serving controller's own connection when this is it: a watchdog tick while a walk waits for a + // build dialled one every 30 seconds (novox/hq issue 327). + js, err := aBus() if err != nil { fmt.Fprintf(os.Stderr, "could not reach the bus to read whether the build seat is paused: %v\n", err) return pauseView{} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index a84dfd56..b4dabd7e 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -219,6 +219,8 @@ func serve(ctx context.Context) (err error) { } // The hand-act log is counted for `status` on this connection rather than a new one a minute. handActConn = bus.Conn + // And everything else this controller does on the bus for a moment (novox/hq issue 327). + servingBus.Store(server.JetStream()) // And says when it replaced a value given by hand (novox/hq ADR 0228). givenEvents = bus // And a pull request's merge check, asked when the forge announces its head and said when judged diff --git a/cmd/mesh-controller/queue.go b/cmd/mesh-controller/queue.go index fbe64894..e4c97e24 100644 --- a/cmd/mesh-controller/queue.go +++ b/cmd/mesh-controller/queue.go @@ -6,6 +6,7 @@ import ( "errors" "flag" "fmt" + "io" "os" "sort" "strings" @@ -37,21 +38,19 @@ const ( killAnswer = 75 * time.Second ) -// dialTheBus opens the controller's own connection, for a command that reads or changes the queue. +// dialTheBus is the controller's connection, for a command that reads or changes the queue: the serving +// controller's own, lent, when this process is it (novox/hq issue 327). func dialTheBus() (*broker.JetStream, error) { - address, err := broker.BusAddress() - if err != nil { - return nil, err - } - js, err := broker.Dial(address) - if err != nil { - return nil, fmt.Errorf("cannot reach the bus: %w", err) - } - return js, nil + return aBus() } // queueCommand prints every ask in the build seat's work queue. func queueCommand(ctx context.Context, args []string) error { + return listQueue(ctx, args, os.Stdout) +} + +// listQueue is `queue`: the build queue, as a person reads it or as JSON. +func listQueue(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("queue", flag.ContinueOnError) asJSON := set.Bool("json", false, "the queue as JSON") if _, err := parseAround(set, args); err != nil { @@ -72,10 +71,10 @@ func queueCommand(ctx context.Context, args []string) error { if err != nil { return err } - fmt.Println(string(body)) + fmt.Fprintln(w, string(body)) return nil } - fmt.Print(queueText(q, time.Now())) + fmt.Fprint(w, queueText(q, time.Now())) return nil } diff --git a/cmd/mesh-controller/reconnects.go b/cmd/mesh-controller/reconnects.go new file mode 100644 index 00000000..1e50e5c3 --- /dev/null +++ b/cmd/mesh-controller/reconnects.go @@ -0,0 +1,175 @@ +package main + +import ( + "context" + "fmt" + "sort" + "strings" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/conditions" + "github.com/novox/mesh-controller/internal/link" +) + +// D15: no client of the bus reconnects in a loop (novox/hq issue 327). +// +// The server's connection total was the one number that would show a client reconnecting, and every verb +// the controller served opened and closed a connection of its own, hundreds an hour, so a loop was +// invisible in it. The bus's own module reads the server's record of closed connections +// (`nats_closed_connections`); this asks it every run, and says each user whose connections were dropped — +// closed by anything but the client itself: a read or write error, a stale connection, a slow consumer, a +// refused login — more often than reconnectBound in the last hour. + +const ( + probeReconnectsID = "D15" + kindBusReconnects = "bus-reconnects" + // reconnectBound is how many dropped connections in an hour one user may have before it is said: + // a client that loses its connection every five minutes. Provisional. + reconnectBound = 12 + // busModule is the module that is the bus, and closedTool its tool that reads closed connections. + busModule = "nats" + closedTool = "nats_closed_connections" + closedAsk = 20 * time.Second +) + +// closedConnections is what the bus's module answers. +type closedConnections struct { + Hours float64 `json:"hours"` + Reaches bool `json:"reaches"` + Users []struct { + User string `json:"user"` + Closed int `json:"closed"` + Dropped int `json:"dropped"` + DroppedPerHour float64 `json:"dropped_per_hour"` + Names []struct { + Name string `json:"name"` + Closed int `json:"closed"` + Dropped int `json:"dropped"` + Reasons map[string]int `json:"reasons"` + } `json:"names"` + } `json:"users"` +} + +// probeReconnects is D15. +func probeReconnects(ctx context.Context, d *doctor) ([]conditions.Observation, error) { + if d.js == nil { + return nil, fmt.Errorf("no bus to ask the bus's module over") + } + on, err := d.open.inventory.Running(ctx, busModule) + if err != nil { + return nil, err + } + if len(on) == 0 { + return nil, nil // no bus module assigned: a mesh whose bus is not the mesh's module + } + read, err := askClosed(ctx, d.js.Conn(), on[0]) + if isNothingServes(err) { + // A bus module older than its tool has nothing it can say, which is not a failure of the probe. + return nil, nil + } + if err != nil { + return nil, err + } + return reconnecting(read), nil +} + +// askClosed asks the bus's module, on its machine, who closed connections in the last hour. +func askClosed(ctx context.Context, conn *nats.Conn, node string) (closedConnections, error) { + var read closedConnections + answer, err := link.AskModuleToolOn(ctx, conn, busModule, closedTool, node, map[string]any{"hours": 1}, closedAsk) + if err != nil { + return read, err + } + if answer.Error != "" { + return read, fmt.Errorf("%s on %s answered %s with an error: %s", busModule, node, closedTool, answer.Error) + } + if err := unmarshalAnswer(answer, &read); err != nil { + return read, fmt.Errorf("%s on %s answered %s with something unreadable: %w", busModule, node, closedTool, err) + } + return read, nil +} + +// reconnecting is one observation per user whose connections were dropped more than reconnectBound times +// an hour. +func reconnecting(read closedConnections) []conditions.Observation { + hours := read.Hours + if hours <= 0 { + hours = 1 + } + var out []conditions.Observation + for _, u := range read.Users { + perHour := float64(u.Dropped) / hours + if perHour <= reconnectBound { + continue + } + reasons := map[string]int{} + var names []string + for _, n := range u.Names { + if n.Dropped == 0 { + continue + } + names = append(names, fmt.Sprintf("%q ×%d", n.Name, n.Dropped)) + for r, c := range n.Reasons { + if r != "Client Closed" { + reasons[r] += c + } + } + } + var why []string + for r, c := range reasons { + why = append(why, fmt.Sprintf("%s ×%d", r, c)) + } + sort.Strings(why) + partial := "" + if !read.Reaches { + partial = " (at least: the server's record of closed connections does not reach back the whole hour)" + } + who, machine := busUserWords(u.User) + out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: u.User, Kind: kindBusReconnects, + Machine: machine, Severity: conditions.Warning, + Summary: fmt.Sprintf("the bus dropped %s's connection %d times in the last hour%s, more than %d: a "+ + "client reconnecting in a loop", u.User, u.Dropped, partial, reconnectBound), + Said: fmt.Sprintf("%d of %d closed connections dropped in %.0f h; by name %s; why %s", u.Dropped, + u.Closed, hours, strings.Join(names, ", "), strings.Join(why, ", ")), + Headline: clipWords(conditions.Capital(who)+" keeps losing the bus", 60), + Explanation: conditions.Capital(fmt.Sprintf("%s lost its connection to the bus %d times in the last hour "+ + "and connected again each time. While it reconnects, what it says and what it is asked waits.", who, + u.Dropped)), + Resolved: "Resolved: " + who + " stays connected"}) + } + return out +} + +// busUserWords is a bus user as the operator says it, and the machine it is on: `node.` the +// node-engine, `.node-tools` the tool runner, `controller` the controller, `.` a +// module. +func busUserWords(user string) (string, string) { + switch { + case user == "controller": + return "the controller", "" + case strings.HasPrefix(user, "node."): + m := strings.TrimPrefix(user, "node.") + return "the node-engine on " + m, m + } + if m, module, ok := strings.Cut(user, "."); ok { + if module == "node-tools" { + return "the tool runner on " + m, m + } + return module + " on " + m, m + } + return "a client of the bus", "" +} + +// clipWords is words at most n characters long, cut at a word. +func clipWords(s string, n int) string { + if len(s) <= n { + return s + } + cut := s[:n] + if i := strings.LastIndex(cut, " "); i > 0 { + cut = cut[:i] + } + return cut +} diff --git a/cmd/mesh-controller/reconnects_test.go b/cmd/mesh-controller/reconnects_test.go new file mode 100644 index 00000000..60334df5 --- /dev/null +++ b/cmd/mesh-controller/reconnects_test.go @@ -0,0 +1,156 @@ +package main + +import ( + "context" + "encoding/json" + "strings" + "testing" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/conditions" + "github.com/novox/mesh-controller/internal/link" + "github.com/novox/mesh-controller/internal/testbus" +) + +// The serving controller's own connection serves what it does for a moment: no connection is opened, +// and closing what it was lent leaves the serving one open (novox/hq issue 327). +func TestTheServingControllerLendsItsOwnConnection(t *testing.T) { + bus := testbus.Start(t) + serving, err := broker.Dial(bus.ClientURL()) + if err != nil { + t.Fatal(err) + } + defer serving.Close() + if err := serving.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + was := servingBus.Load() + servingBus.Store(serving) + defer servingBus.Store(was) + t.Setenv(broker.NATSVar, bus.ClientURL()) + + before, _ := bus.Varz(nil) + lent, err := aBus() + if err != nil { + t.Fatal(err) + } + lent.Close() + if !serving.Conn().IsConnected() { + t.Fatal("closing a lent connection closed the serving controller's") + } + if err := onTheBus(func(*nats.Conn) error { return nil }); err != nil { + t.Fatal(err) + } + handlers, _, err := seatToolHandlers() + if err != nil { + t.Fatal(err) + } + answer, err := handlers["hand-acts"](context.Background(), json.RawMessage(`{}`)) + if err != nil { + t.Fatal(err) + } + if a := answer.(verbAnswer); !a.OK || a.Answer == nil { + t.Fatalf("hand-acts answered %+v", a) + } + _, _ = handlers["queue"](context.Background(), json.RawMessage(`{}`)) + after, _ := bus.Varz(nil) + if opened := after.TotalConnections - before.TotalConnections; opened != 0 { + t.Fatalf("the serving controller opened %d connection(s) of its own", opened) + } +} + +// A verb that still runs as a process of its own says which in the bus's list of connections. +func TestAVerbsOwnProcessNamesItsConnection(t *testing.T) { + bus := testbus.Start(t) + setup, err := broker.Dial(bus.ClientURL()) + if err != nil { + t.Fatal(err) + } + if err := setup.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + setup.Close() + asAProcess(t, bus.ClientURL()) + handlers, _, err := seatToolHandlers() + if err != nil { + t.Fatal(err) + } + // A silence acts, so it is still its own process; it dials, and finds no such condition. + _, _ = handlers["conditions"](context.Background(), + json.RawMessage(`{"silence":"bus.nothing.here","for":"1h","why":"a test"}`)) + names := closedNames(t, bus) + found := false + for _, n := range names { + found = found || n == "mesh-controller verb conditions" + } + if !found { + t.Fatalf("the verb's connection was named %q", names) + } +} + +func TestEachProcessNamesItsConnectionsForWhatItIs(t *testing.T) { + for _, c := range []struct{ command, verb, want string }{ + {"serve", "", "mesh-controller serving"}, + {"conditions", "conditions", "mesh-controller verb conditions"}, + {"delivery", "delivery-check", "mesh-controller verb delivery-check"}, + {"push", "", "mesh-controller command push"}, + } { + if got := connectionName(c.command, c.verb); got != c.want { + t.Errorf("%s/%s: %q", c.command, c.verb, got) + } + } +} + +// D15: a user whose connections are dropped more than twelve times an hour is said, in plain words; one +// whose client closes its own short connections is not. +func TestAClientReconnectingInALoopIsSaid(t *testing.T) { + var read closedConnections + if err := json.Unmarshal([]byte(`{"hours":1,"reaches":true,"users":[ + {"user":"ace.node-tools","closed":40,"dropped":40,"dropped_per_hour":40,"names":[ + {"name":"ace.node-tools","closed":40,"dropped":40,"reasons":{"Stale Connection":30,"Read Error":10}}]}, + {"user":"controller","closed":300,"dropped":0,"names":[ + {"name":"mesh-controller verb delivery-check","closed":300,"dropped":0,"reasons":{"Client Closed":300}}]}, + {"user":"node.anchor","closed":5,"dropped":5,"names":[{"name":"mesh-host/anchor","closed":5,"dropped":5, + "reasons":{"Read Error":5}}]}]}`), &read); err != nil { + t.Fatal(err) + } + found := reconnecting(read) + if len(found) != 1 { + t.Fatalf("%+v", found) + } + o := found[0] + if o.ID != "ace.node-tools" || o.Machine != "ace" || o.Kind != kindBusReconnects || + o.Headline != "The tool runner on ace keeps losing the bus" || !strings.Contains(o.Said, "Stale Connection ×30") { + t.Fatalf("%+v", o) + } + if why, ok := conditions.PlainWords(conditions.Words{Headline: o.Headline, Explanation: o.Explanation, + Resolved: o.Resolved}, "ace"); !ok { + t.Fatalf("not plain: %s", why) + } +} + +// The bus's module is asked on its machine, over the bus. +func TestTheBusModuleIsAskedWhoClosedConnections(t *testing.T) { + bus := testbus.Start(t) + conn, err := nats.Connect(bus.ClientURL()) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + sub, err := conn.Subscribe(link.ModuleToolOn("nats", "nats_closed_connections", "anchor"), func(m *nats.Msg) { + _ = m.Respond([]byte(`{"result":{"hours":1,"reaches":true,"users":[{"user":"controller","closed":2,"dropped":0}]}}`)) + }) + if err != nil { + t.Fatal(err) + } + defer func() { _ = sub.Unsubscribe() }() + read, err := askClosed(context.Background(), conn, "anchor") + if err != nil { + t.Fatal(err) + } + if len(read.Users) != 1 || read.Users[0].Closed != 2 { + t.Fatalf("%+v", read) + } +} diff --git a/cmd/mesh-controller/replay327_test.go b/cmd/mesh-controller/replay327_test.go new file mode 100644 index 00000000..20e02cca --- /dev/null +++ b/cmd/mesh-controller/replay327_test.go @@ -0,0 +1,132 @@ +package main + +import ( + "context" + "encoding/json" + "os" + "os/exec" + "path/filepath" + "sync" + "testing" + "time" + + "github.com/nats-io/nats-server/v2/server" + + "github.com/novox/mesh-controller/internal/broker" + "github.com/novox/mesh-controller/internal/testbus" +) + +// novox/hq issue 327, replayed with only what the controller had before its fix, so it can be laid over +// the older commit. On 2026-10-08 the bus's connection total rose by 5.2 a minute while the same 16 +// connections stayed open, and three calls of `conditions` in one second added three: each verb ran as a +// process of the controller's own binary, which dialled the bus, logged in and left. The operator's +// channel reads `conditions` at least once a minute. A verb that only reads, served by the serving +// controller, opens no connection of its own. +func TestReplay327(t *testing.T) { + bus := testbus.Start(t) + serving, err := broker.Dial(bus.ClientURL()) + if err != nil { + t.Fatal(err) + } + defer serving.Close() + if err := serving.EnsureControllerBuckets(); err != nil { + t.Fatal(err) + } + keeper, err := keeperOn(context.Background(), serving.Conn()) + if err != nil { + t.Fatal(err) + } + // The serving controller: its keeper and its connection, as serve sets them. + keptBefore, connBefore := conditionsFrom, handActConn + conditionsFrom, handActConn = keeper, serving.Conn() + defer func() { conditionsFrom, handActConn = keptBefore, connBefore }() + // A verb that runs as a process of its own runs this controller's binary, on this bus. + asAProcess(t, bus.ClientURL()) + + handlers, _, err := seatToolHandlers() + if err != nil { + t.Fatal(err) + } + accepted := func() uint64 { + v, err := bus.Varz(nil) + if err != nil { + t.Fatal(err) + } + return v.TotalConnections + } + before := accepted() + for i := 0; i < 3; i++ { + answer, err := handlers["conditions"](context.Background(), json.RawMessage(`{}`)) + if err != nil { + t.Fatal(err) + } + if a, ok := answer.(verbAnswer); !ok || !a.OK { + t.Fatalf("conditions answered %+v", answer) + } + } + if opened := accepted() - before; opened != 0 { + t.Fatalf("three conditions calls opened %d connection(s) to the bus; the serving controller is on it already", + opened) + } +} + +// The binary a verb run as a process of its own is, built once for the tests that run one. +var verbBinary struct { + once sync.Once + path string + err error +} + +// asAProcess makes a verb that runs as a process of its own run this package's binary, on the bus at url. +func asAProcess(t *testing.T, url string) { + t.Helper() + verbBinary.once.Do(func() { + dir, err := os.MkdirTemp("", "mesh-controller-verb-") + if err != nil { + verbBinary.err = err + return + } + verbBinary.path = filepath.Join(dir, "mesh-controller") + out, err := exec.Command("go", "build", "-o", verbBinary.path, ".").CombinedOutput() + if err != nil { + verbBinary.err = &buildError{out: string(out), err: err} + } + }) + if verbBinary.err != nil { + t.Fatalf("the controller could not be built to run a verb as its own process: %v", verbBinary.err) + } + was := ownImage + ownImage = func() string { return verbBinary.path } + t.Cleanup(func() { ownImage = was }) + t.Setenv(broker.NATSVar, url) + t.Setenv(broker.CertificateVar, "") +} + +type buildError struct { + out string + err error +} + +func (e *buildError) Error() string { return e.err.Error() + ": " + e.out } + +// closedNames are the names of the connections the bus saw closed. +func closedNames(t *testing.T, bus *server.Server) []string { + t.Helper() + deadline := time.Now().Add(5 * time.Second) + var names []string + for time.Now().Before(deadline) { + connz, err := bus.Connz(&server.ConnzOptions{State: server.ConnClosed}) + if err != nil { + t.Fatal(err) + } + names = names[:0] + for _, c := range connz.Conns { + names = append(names, c.Name) + } + if len(names) > 0 { + return names + } + time.Sleep(50 * time.Millisecond) + } + return names +} diff --git a/cmd/mesh-controller/rollout.go b/cmd/mesh-controller/rollout.go index 4e459c51..81f0c2a6 100644 --- a/cmd/mesh-controller/rollout.go +++ b/cmd/mesh-controller/rollout.go @@ -131,7 +131,10 @@ func readinessOf(ctx context.Context, inv *inventory.Inventory) (broker.Readines // **Dialled the way the mesh dials it** — credential and pin — because a bare connect to a // bus that requires TLS and a user fails at the handshake, and the check then reported a // standing server as absent (seen live, 2026-09-28). - if js, err := broker.Dial(address, nats.Timeout(5*time.Second)); err == nil { + // The serving controller's own connection answers it without a second (novox/hq issue 327). + if serving := servingBus.Load(); serving != nil && serving.Conn().IsConnected() { + state.ServerStanding = true + } else if js, err := broker.Dial(address, nats.Timeout(5*time.Second)); err == nil { state.ServerStanding = true js.Close() } diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 58aed4a9..f3b5a4d1 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -7,6 +7,7 @@ import ( "errors" "fmt" "github.com/nats-io/nats.go/micro" + "io" "os" "os/exec" "slices" @@ -28,6 +29,9 @@ import ( // answer. It also means a refusal is the same refusal in the same words, because it is the same // output. +// verbKey carries the verb a call is for, to the process it runs. +type verbKey struct{} + // verbAnswer is what a verb answers: what the command printed, whether it succeeded, and — where the // command speaks JSON — the same as data. type verbAnswer struct { @@ -921,6 +925,12 @@ func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { caller = "a seat call whose caller the bus did not name" } cmd.Env = append(cmd.Env, link.CallerVar+"="+caller+", through the "+catalogue.ControllerSeatName+" seat") + // And which verb, so the connection it dials says so in the bus's list (novox/hq issue 327). + verb, _ := ctx.Value(verbKey{}).(string) + if verb == "" { + verb = argv[0] + } + cmd.Env = append(cmd.Env, verbVar+"="+verb) // Two buffers, one answer. What the command *says* is both streams, in the order a person at // a shell would read them; what it *answers as data* is standard output alone — `status --json` // prints its warnings beside the document, and a JSON parsed from the two together parsed @@ -929,18 +939,7 @@ func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { cmd.Stdout = &stdout cmd.Stderr = &stderr runErr := cmd.Run() - answer := verbAnswer{Output: stdout.String() + stderr.String(), OK: runErr == nil} - if jsonVerbs[argv[0]] && runErr == nil { - 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" - } - } - } + answer := answerOf(argv, stdout.Bytes(), stderr.String(), runErr == nil) var exit *exec.ExitError if runErr != nil && !errors.As(runErr, &exit) { // Not the command refusing — the command not running at all, which is this process's fault. @@ -955,6 +954,66 @@ func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { return answer, nil } +// answerOf is what a command said, as a verb answers it: both streams as text, and standard output as data +// where the command speaks JSON. +func answerOf(argv []string, stdout []byte, stderr string, ok bool) verbAnswer { + answer := verbAnswer{Output: string(stdout) + stderr, OK: ok} + if jsonVerbs[argv[0]] && ok { + var parsed any + if json.Unmarshal(bytes.TrimSpace(stdout), &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 + "its answer, as data, is `answer`\n" + } + } + } + return answer +} + +// readHere answers a verb that only reads the bus in the serving controller itself, on its own connection +// and keeper (novox/hq issue 327): the same command, writing to the answer rather than to a process's +// output, so the answer is the one the command prints. False for any other command line, which runs as a +// command of its own. Every `conditions` call — the operator's channel reads it at least once a minute — +// was a process that dialled the bus, logged in and left. +func readHere(ctx context.Context, argv []string) (verbAnswer, bool) { + var read func(context.Context, []string, io.Writer) error + args := argv[1:] + // Each where this process holds what it reads: the serving keeper, the hand-act log's connection, the + // serving connection. + switch argv[0] { + case "conditions": + sub := "list" + if len(args) > 0 && !strings.HasPrefix(args[0], "-") { + sub, args = args[0], args[1:] + } + if conditionsFrom != nil { + read = map[string]func(context.Context, []string, io.Writer) error{ + "list": listConditions, "show": showCondition, "history": conditionHistory}[sub] + } + case "hand-acts": + if handActConn != nil || servingBus.Load() != nil { + read = listHandActs + } + case "queue": + if servingBus.Load() != nil { + read = listQueue + } + } + if read == nil { + return verbAnswer{}, false + } + var out bytes.Buffer + stderr := "" + err := read(ctx, args, &out) + if err != nil { + // As the command says it when it fails (main). + stderr = "mesh-controller: " + err.Error() + "\n" + } + return answerOf(argv, out.Bytes(), stderr, err == nil), true +} + // seatToolHandlers are the handlers for every verb the mesh-controller seat declares, from the // store's row, so a verb the row does not carry is not served. A verb it carries that this binary // cannot run is named at start and answers the reason when called — never a refusal to serve, which @@ -1038,6 +1097,7 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { if err != nil { return nil, err } + ctx = context.WithValue(ctx, verbKey{}, verb) if verb == "status" && statusFrom != nil { // At once, from the summary the serving controller keeps (novox/hq to-be 45 Phase 0). return statusFrom.answer(ctx) @@ -1046,6 +1106,9 @@ func seatToolHandlers() (map[string]link.ToolHandler, []string, error) { // Whatever it did, `status` is composed again once it has. defer statusFrom.nudge() } + if answer, read := readHere(ctx, argv); read { + return answer, nil + } if answersFirst(argv) { // Before anything is sent: a push sends the bus's own machine first, and a broker // reloading its user list forgets the answer it was about to permit (novox/hq issue 265). diff --git a/cmd/mesh-controller/servingbus.go b/cmd/mesh-controller/servingbus.go new file mode 100644 index 00000000..48d3cfcb --- /dev/null +++ b/cmd/mesh-controller/servingbus.go @@ -0,0 +1,38 @@ +package main + +import ( + "fmt" + "sync/atomic" + + "github.com/novox/mesh-controller/internal/broker" +) + +// The serving controller's own connection, lent to whatever it does for a moment (novox/hq issue 327). +// +// Every place that needed the bus for a moment dialled it: a verb's own process, and in the serving +// controller a watchdog tick reading whether the build seat paused, a walk's step reading readiness, a +// queue read. Each paid a connection, a TLS handshake and a login on the control node, and hundreds an +// hour hid in the server's connection total the one thing it would show: a client reconnecting in a loop. +// The serving controller is on the bus already; what it does is done on that connection. + +// servingBus is the serving controller's connection, set when it starts serving; nil in every other +// process, which dials its own. +var servingBus atomic.Pointer[broker.JetStream] + +// aBus is a connection for something done for a moment: the serving controller's own, lent — so its +// Close closes nothing — when this process is it, and otherwise one dialled for it, named for this +// process (broker.ConnectionName), which its Close closes. +func aBus() (*broker.JetStream, error) { + if serving := servingBus.Load(); serving != nil { + return broker.Borrow(serving), nil + } + address, err := broker.BusAddress() + if err != nil { + return nil, err + } + js, err := broker.Dial(address) + if err != nil { + return nil, fmt.Errorf("cannot reach the bus: %w", err) + } + return js, nil +} diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 7480f578..70ae35b8 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -28,6 +28,8 @@ import ( type JetStream struct { conn *nats.Conn js nats.JetStreamContext + // borrowed is a connection lent by its owner (Borrow): closing it is the owner's. + borrowed bool // Note is how this says something it decided not to fail over. Nil is silent, which is only // right for a caller that has no way to report; the controller sets it. Note func(string, ...any) @@ -40,9 +42,17 @@ func (j *JetStream) note(format string, args ...any) { } } +// ConnectionName is what a connection this process dials says it is, in the server's list of +// connections: the controller's role and what it is doing (novox/hq issue 327). Every connection was +// named `mesh-controller` — the serving controller's, each verb's own process, the build agents' — so the +// server's list could not say which was which. The process sets it once, at its start; an option a +// caller passes to Dial names one connection otherwise. +var ConnectionName = "mesh-controller" + // Dial connects and returns the controller's JetStream handle. func Dial(url string, opts ...nats.Option) (*JetStream, error) { - opts = append(opts, nats.Name("mesh-controller"), nats.Timeout(10*time.Second)) + // The name first, so a name the caller gives is the one that stands. + opts = append([]nats.Option{nats.Name(ConnectionName), nats.Timeout(10 * time.Second)}, opts...) // **Pinned, not named.** The bus presents the mesh's own certificate, which names nothing a // public verifier would accept (design 25 §4: a host pins the server's exact certificate and // checks nothing else, and so does this). Without this, the first connection failed with @@ -130,11 +140,20 @@ func (j *JetStream) Conn() *nats.Conn { return j.conn } func (j *JetStream) Context() nats.JetStreamContext { return j.js } func (j *JetStream) Close() { - if j.conn != nil { + if j.conn != nil && !j.borrowed { j.conn.Close() } } +// Borrow is the same connection for a caller that will close what it was handed when it is done: its +// Close closes nothing, and the connection stays its owner's (novox/hq issue 327). How the serving +// controller lends its own connection to work that would otherwise dial one of its own. +func Borrow(j *JetStream) *JetStream { + lent := *j + lent.borrowed = true + return &lent +} + // EnsureStream creates the stream if it is absent and brings it to match if it is present. // // **Idempotent, because the controller asserts on every start** rather than creating once at diff --git a/internal/broker/jetstream_test.go b/internal/broker/jetstream_test.go index e0c237a4..5124696d 100644 --- a/internal/broker/jetstream_test.go +++ b/internal/broker/jetstream_test.go @@ -3,6 +3,8 @@ package broker import ( "testing" + "github.com/nats-io/nats.go" + "github.com/novox/mesh-controller/internal/testbus" ) @@ -66,3 +68,32 @@ func TestAgainstARealServer(t *testing.T) { } }) } + +// A connection says what it is in the server's list (novox/hq issue 327): the process's name, unless the +// caller names this one; and a lent connection's Close leaves its owner's open. +func TestAConnectionIsNamedAndALentOneIsNotClosed(t *testing.T) { + url := testbus.URL(t) + was := ConnectionName + ConnectionName = "mesh-controller verb conditions" + defer func() { ConnectionName = was }() + named, err := Dial(url) + if err != nil { + t.Fatal(err) + } + defer named.Close() + if got := named.Conn().Opts.Name; got != "mesh-controller verb conditions" { + t.Errorf("named %q", got) + } + lease, err := Dial(url, nats.Name("mesh-controller serving lease")) + if err != nil { + t.Fatal(err) + } + defer lease.Close() + if got := lease.Conn().Opts.Name; got != "mesh-controller serving lease" { + t.Errorf("a name the caller gave became %q", got) + } + Borrow(named).Close() + if !named.Conn().IsConnected() { + t.Error("closing a lent connection closed its owner's") + } +} diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index a2cbd1f6..1358dd92 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -47,6 +47,12 @@ func BuildsOverNATSOn(address, seat string) (Builders, error) { return &natsBuilds{js: js, owned: true, seat: seat}, nil } +// BuildsOn is the asking side on a connection the caller holds, which Close leaves open: the serving +// controller asks on its own rather than dialling one per build (novox/hq issue 327). +func BuildsOn(js *broker.JetStream, seat string) Builders { + return &natsBuilds{js: js, seat: seat} +} + // role is the seat asked: what the asker was made for, or the current build role for one made // without saying (a test building the struct by hand). func (b *natsBuilds) role() string { From cc7fb99f29344246103da64d4ce6ccb425d587fd Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 18:45:04 +0200 Subject: [PATCH 2/3] Answer a panicking verb with an error, say flag errors in the answer, and leave refused logins out of D15 Read verbs now run in the serving process, where a panic would end every call; refused logins are nobody's reconnect loop (review of hq issue 327). --- cmd/mesh-controller/conditions.go | 3 ++ cmd/mesh-controller/handacts.go | 1 + cmd/mesh-controller/queue.go | 1 + cmd/mesh-controller/reconnects.go | 12 +++++++- cmd/mesh-controller/reconnects_test.go | 33 ++++++++++++++++++++++ cmd/mesh-controller/replay327_test.go | 38 ++++---------------------- cmd/mesh-controller/servingbus.go | 11 ++++++++ internal/link/calls.go | 14 ++++++++++ internal/link/calls_test.go | 30 ++++++++++++++++++++ 9 files changed, 110 insertions(+), 33 deletions(-) diff --git a/cmd/mesh-controller/conditions.go b/cmd/mesh-controller/conditions.go index c26b181b..eed26b61 100644 --- a/cmd/mesh-controller/conditions.go +++ b/cmd/mesh-controller/conditions.go @@ -121,6 +121,7 @@ func conditionsCommand(ctx context.Context, args []string) error { func listConditions(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions", flag.ContinueOnError) + usageTo(set, w) scope := set.String("scope", "", "only this scope: "+strings.Join(conditions.Scopes, ", ")) severity := set.String("severity", "", "only urgent, or only warning") machine := set.String("machine", "", "only those about this machine") @@ -236,6 +237,7 @@ func conditionLines(list []conditions.Condition, now time.Time) []string { func showCondition(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions show", flag.ContinueOnError) + usageTo(set, w) asJSON := set.Bool("json", false, "as data") rest, err := parseAround(set, args) if err != nil { @@ -381,6 +383,7 @@ func parseFor(s string) (time.Duration, error) { func conditionHistory(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions history", flag.ContinueOnError) + usageTo(set, w) days := set.Int("days", 7, "how many days back, at most 90") key := set.String("key", "", "only this condition") asJSON := set.Bool("json", false, "as data") diff --git a/cmd/mesh-controller/handacts.go b/cmd/mesh-controller/handacts.go index 0cc4e3c0..38aea23b 100644 --- a/cmd/mesh-controller/handacts.go +++ b/cmd/mesh-controller/handacts.go @@ -254,6 +254,7 @@ func handActCommand(ctx context.Context, args []string) error { // listHandActs is `hand-acts`: what was done by hand lately, and the causes done more than once. func listHandActs(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("hand-acts", flag.ContinueOnError) + usageTo(set, w) days := set.Int("days", 14, "how many days back") asJSON := set.Bool("json", false, "as data") if _, err := parseAround(set, args); err != nil { diff --git a/cmd/mesh-controller/queue.go b/cmd/mesh-controller/queue.go index e4c97e24..d36a6814 100644 --- a/cmd/mesh-controller/queue.go +++ b/cmd/mesh-controller/queue.go @@ -52,6 +52,7 @@ func queueCommand(ctx context.Context, args []string) error { // listQueue is `queue`: the build queue, as a person reads it or as JSON. func listQueue(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("queue", flag.ContinueOnError) + usageTo(set, w) asJSON := set.Bool("json", false, "the queue as JSON") if _, err := parseAround(set, args); err != nil { return err diff --git a/cmd/mesh-controller/reconnects.go b/cmd/mesh-controller/reconnects.go index 1e50e5c3..ed12cfe7 100644 --- a/cmd/mesh-controller/reconnects.go +++ b/cmd/mesh-controller/reconnects.go @@ -64,7 +64,12 @@ func probeReconnects(ctx context.Context, d *doctor) ([]conditions.Observation, if len(on) == 0 { return nil, nil // no bus module assigned: a mesh whose bus is not the mesh's module } - read, err := askClosed(ctx, d.js.Conn(), on[0]) + return reconnectsOn(ctx, d.js.Conn(), on[0]) +} + +// reconnectsOn asks the bus module on its machine and says who reconnects in a loop. +func reconnectsOn(ctx context.Context, conn *nats.Conn, node string) ([]conditions.Observation, error) { + read, err := askClosed(ctx, conn, node) if isNothingServes(err) { // A bus module older than its tool has nothing it can say, which is not a failure of the probe. return nil, nil @@ -100,6 +105,11 @@ func reconnecting(read closedConnections) []conditions.Observation { } var out []conditions.Observation for _, u := range read.Users { + if u.User == "" || strings.HasPrefix(u.User, "(") { + // Refused before it logged in: no client of the mesh's, so nobody's reconnect loop. A login + // refused again and again is a question of its own, not this probe's. + continue + } perHour := float64(u.Dropped) / hours if perHour <= reconnectBound { continue diff --git a/cmd/mesh-controller/reconnects_test.go b/cmd/mesh-controller/reconnects_test.go index 60334df5..d6b26b95 100644 --- a/cmd/mesh-controller/reconnects_test.go +++ b/cmd/mesh-controller/reconnects_test.go @@ -112,6 +112,8 @@ func TestAClientReconnectingInALoopIsSaid(t *testing.T) { {"name":"ace.node-tools","closed":40,"dropped":40,"reasons":{"Stale Connection":30,"Read Error":10}}]}, {"user":"controller","closed":300,"dropped":0,"names":[ {"name":"mesh-controller verb delivery-check","closed":300,"dropped":0,"reasons":{"Client Closed":300}}]}, + {"user":"(no user: refused before it logged in)","closed":90,"dropped":90,"names":[{"name":"","closed":90, + "dropped":90,"reasons":{"Authentication Failure":90}}]}, {"user":"node.anchor","closed":5,"dropped":5,"names":[{"name":"mesh-host/anchor","closed":5,"dropped":5, "reasons":{"Read Error":5}}]}]}`), &read); err != nil { t.Fatal(err) @@ -154,3 +156,34 @@ func TestTheBusModuleIsAskedWhoClosedConnections(t *testing.T) { t.Fatalf("%+v", read) } } + +// A bus module older than the tool answers nothing to the question: the probe passes over it quietly. +func TestAnOlderBusModuleIsPassedOverQuietly(t *testing.T) { + bus := testbus.Start(t) + conn, err := nats.Connect(bus.ClientURL()) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + found, err := reconnectsOn(context.Background(), conn, "anchor") + if err != nil || len(found) != 0 { + t.Fatalf("a bus module without the tool: %v %v", found, err) + } +} + +// A verb answered in the serving controller says its flag errors in its answer, not in the controller's log. +func TestAFlagErrorIsSaidInTheAnswer(t *testing.T) { + bus := testbus.Start(t) + serving, err := broker.Dial(bus.ClientURL()) + if err != nil { + t.Fatal(err) + } + defer serving.Close() + was := servingBus.Load() + servingBus.Store(serving) + defer servingBus.Store(was) + answer, read := readHere(context.Background(), []string{"hand-acts", "--bogus"}) + if !read || answer.OK || !strings.Contains(answer.Output, "flag provided but not defined") { + t.Fatalf("answered %+v", answer) + } +} diff --git a/cmd/mesh-controller/replay327_test.go b/cmd/mesh-controller/replay327_test.go index 20e02cca..559debe9 100644 --- a/cmd/mesh-controller/replay327_test.go +++ b/cmd/mesh-controller/replay327_test.go @@ -3,10 +3,8 @@ package main import ( "context" "encoding/json" - "os" "os/exec" "path/filepath" - "sync" "testing" "time" @@ -70,45 +68,21 @@ func TestReplay327(t *testing.T) { } } -// The binary a verb run as a process of its own is, built once for the tests that run one. -var verbBinary struct { - once sync.Once - path string - err error -} - -// asAProcess makes a verb that runs as a process of its own run this package's binary, on the bus at url. +// asAProcess makes a verb that runs as a process of its own run this package's binary, on the bus at url: +// built into the test's own directory, which goes with the test. func asAProcess(t *testing.T, url string) { t.Helper() - verbBinary.once.Do(func() { - dir, err := os.MkdirTemp("", "mesh-controller-verb-") - if err != nil { - verbBinary.err = err - return - } - verbBinary.path = filepath.Join(dir, "mesh-controller") - out, err := exec.Command("go", "build", "-o", verbBinary.path, ".").CombinedOutput() - if err != nil { - verbBinary.err = &buildError{out: string(out), err: err} - } - }) - if verbBinary.err != nil { - t.Fatalf("the controller could not be built to run a verb as its own process: %v", verbBinary.err) + path := filepath.Join(t.TempDir(), "mesh-controller") + if out, err := exec.Command("go", "build", "-o", path, ".").CombinedOutput(); err != nil { + t.Fatalf("the controller could not be built to run a verb as its own process: %v: %s", err, out) } was := ownImage - ownImage = func() string { return verbBinary.path } + ownImage = func() string { return path } t.Cleanup(func() { ownImage = was }) t.Setenv(broker.NATSVar, url) t.Setenv(broker.CertificateVar, "") } -type buildError struct { - out string - err error -} - -func (e *buildError) Error() string { return e.err.Error() + ": " + e.out } - // closedNames are the names of the connections the bus saw closed. func closedNames(t *testing.T, bus *server.Server) []string { t.Helper() diff --git a/cmd/mesh-controller/servingbus.go b/cmd/mesh-controller/servingbus.go index 48d3cfcb..727f4fb9 100644 --- a/cmd/mesh-controller/servingbus.go +++ b/cmd/mesh-controller/servingbus.go @@ -1,7 +1,10 @@ package main import ( + "flag" "fmt" + "io" + "os" "sync/atomic" "github.com/novox/mesh-controller/internal/broker" @@ -36,3 +39,11 @@ func aBus() (*broker.JetStream, error) { } return js, nil } + +// usageTo sends a command's flag errors and usage to where its answer goes when that is not this process's +// output: a verb answered in the serving controller says them in its answer, not in the controller's log. +func usageTo(set *flag.FlagSet, w io.Writer) { + if w != os.Stdout { + set.SetOutput(w) + } +} diff --git a/internal/link/calls.go b/internal/link/calls.go index 49d13e83..daaaf6e1 100644 --- a/internal/link/calls.go +++ b/internal/link/calls.go @@ -7,6 +7,7 @@ import ( "fmt" "log" "regexp" + "runtime/debug" "sort" "strconv" "strings" @@ -575,6 +576,19 @@ func (l *CallLog) serveCallWithin(seat, verb string, args json.RawMessage, reply done := make(chan outcome, 1) go func() { var body []byte + // **A handler that panics is an answer, not a dead controller.** A verb answered in the serving + // process (novox/hq issue 327: conditions, hand-acts, queue, dead-letters) runs on this goroutine, + // where a panic would take every other call and the controller with it. + defer func() { + if p := recover(); p != nil { + if logger != nil { + logger.Printf("%s.%s: call %s panicked: %v\n%s", seat, verb, c.ID, p, debug.Stack()) + } + body, _ = json.Marshal(map[string]any{"error": fmt.Sprintf("%s failed inside the controller and "+ + "answered nothing: %v. It is in the controller's log", verb, p)}) + done <- outcome{body, true} + } + }() result, err := handle(ctx, args) failed := err != nil if errors.Is(err, ErrHandingOver) { diff --git a/internal/link/calls_test.go b/internal/link/calls_test.go index b06ee147..9c2d83fd 100644 --- a/internal/link/calls_test.go +++ b/internal/link/calls_test.go @@ -11,6 +11,10 @@ import ( "sync" "testing" "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/testbus" ) // answers collects what a call answered, and fails a second answer: the bus permits one. @@ -223,3 +227,29 @@ func TestAJoinTokenIsShownToItsCallerAndNotKept(t *testing.T) { } } } + +// A verb whose handler panics answers an error, and the controller serves the next call (review of +// novox/hq issue 327: read verbs are answered in the serving process now). +func TestAHandlerThatPanicsAnswersAnError(t *testing.T) { + conn, err := nats.Connect(testbus.URL(t)) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + stop, err := OverNATS{Conn: conn}.ServeSeatTools("panicky", map[string]ToolHandler{ + "boom": func(context.Context, json.RawMessage) (any, error) { panic("nil map") }, + "fine": func(context.Context, json.RawMessage) (any, error) { return "ok", nil }, + }, nil) + if err != nil { + t.Fatal(err) + } + defer stop() + answer, err := AskMeshSeatTool(context.Background(), conn, "panicky", "boom", map[string]any{}, 5*time.Second) + if err != nil || !strings.Contains(answer.Error, "failed inside the controller") { + t.Fatalf("a panic answered %+v (%v)", answer, err) + } + answer, err = AskMeshSeatTool(context.Background(), conn, "panicky", "fine", map[string]any{}, 5*time.Second) + if err != nil || answer.Error != "" { + t.Fatalf("the call after a panic answered %+v (%v)", answer, err) + } +} From 5b7e6ff453eb991e21304d2fd31b4a13ae2c214b Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 21:09:08 +0200 Subject: [PATCH 3/3] Answer dead-letters on the lent serving connection, and keep one clip helper Two handles to the same serving connection, and two copies of one helper, would drift (review of hq issues 327 and 330). --- cmd/mesh-controller/deadletters.go | 10 +++------- cmd/mesh-controller/deadletters_test.go | 6 +++--- cmd/mesh-controller/push.go | 2 ++ cmd/mesh-controller/reconnects.go | 14 +------------- cmd/mesh-controller/watchdogs.go | 4 ---- 5 files changed, 9 insertions(+), 27 deletions(-) diff --git a/cmd/mesh-controller/deadletters.go b/cmd/mesh-controller/deadletters.go index d2c7b735..c01dd2c4 100644 --- a/cmd/mesh-controller/deadletters.go +++ b/cmd/mesh-controller/deadletters.go @@ -6,7 +6,6 @@ import ( "fmt" "strconv" "strings" - "sync/atomic" "github.com/nats-io/nats.go" @@ -22,10 +21,6 @@ import ( // serving controller is already on it. A verb run as a fresh process would open a connection of its own // for each call (novox/hq issue 327). -// deadLettersOn is the serving controller's bus, set when it starts watching the mesh; nil in a -// process that serves nothing. -var deadLettersOn atomic.Pointer[busHandles] - // busHandles are the serving controller's connection and JetStream handle. type busHandles struct { conn *nats.Conn @@ -41,11 +36,12 @@ const causeDeadLetter = "dead-letter" // deadLettersAnswer is what `dead-letters` answers: the list, one whole, or what came of delivering one // again or dropping it. func deadLettersAnswer(ctx context.Context, a *verbArguments) (any, error) { - on := deadLettersOn.Load() - if on == nil { + serving := servingBus.Load() + if serving == nil { return nil, errors.New("this controller is not serving, so it does not read DEAD_LETTERS: ask again, " + "and the serving controller answers") } + on := &busHandles{conn: serving.Conn(), js: serving.Context()} // The shape first, and every argument it reads; one given beside it is refused before anything is // done, as every verb refuses what it would pass over (novox/hq issue 244). var deliver, drop, why, cause, idText, consumer, limit string diff --git a/cmd/mesh-controller/deadletters_test.go b/cmd/mesh-controller/deadletters_test.go index 5b3aef7a..e714456f 100644 --- a/cmd/mesh-controller/deadletters_test.go +++ b/cmd/mesh-controller/deadletters_test.go @@ -26,9 +26,9 @@ func servingDeadLetters(t *testing.T) *broker.JetStream { if err := js.EnsureControllerBuckets(); err != nil { t.Fatal(err) } - before := deadLettersOn.Load() - deadLettersOn.Store(&busHandles{conn: js.Conn(), js: js.Context()}) - t.Cleanup(func() { deadLettersOn.Store(before) }) + before := servingBus.Load() + servingBus.Store(js) + t.Cleanup(func() { servingBus.Store(before) }) return js } diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index b4dabd7e..a4f54d63 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -221,6 +221,8 @@ func serve(ctx context.Context) (err error) { handActConn = bus.Conn // And everything else this controller does on the bus for a moment (novox/hq issue 327). servingBus.Store(server.JetStream()) + // No longer serving: nothing is lent, and dead-letters says it is not read here (novox/hq issue 330). + defer servingBus.Store(nil) // And says when it replaced a value given by hand (novox/hq ADR 0228). givenEvents = bus // And a pull request's merge check, asked when the forge announces its head and said when judged diff --git a/cmd/mesh-controller/reconnects.go b/cmd/mesh-controller/reconnects.go index ed12cfe7..42f14479 100644 --- a/cmd/mesh-controller/reconnects.go +++ b/cmd/mesh-controller/reconnects.go @@ -143,7 +143,7 @@ func reconnecting(read closedConnections) []conditions.Observation { "client reconnecting in a loop", u.User, u.Dropped, partial, reconnectBound), Said: fmt.Sprintf("%d of %d closed connections dropped in %.0f h; by name %s; why %s", u.Dropped, u.Closed, hours, strings.Join(names, ", "), strings.Join(why, ", ")), - Headline: clipWords(conditions.Capital(who)+" keeps losing the bus", 60), + Headline: clip(conditions.Capital(who)+" keeps losing the bus", 60), Explanation: conditions.Capital(fmt.Sprintf("%s lost its connection to the bus %d times in the last hour "+ "and connected again each time. While it reconnects, what it says and what it is asked waits.", who, u.Dropped)), @@ -171,15 +171,3 @@ func busUserWords(user string) (string, string) { } return "a client of the bus", "" } - -// clipWords is words at most n characters long, cut at a word. -func clipWords(s string, n int) string { - if len(s) <= n { - return s - } - cut := s[:n] - if i := strings.LastIndex(cut, " "); i > 0 { - cut = cut[:i] - } - return cut -} diff --git a/cmd/mesh-controller/watchdogs.go b/cmd/mesh-controller/watchdogs.go index a4c07cef..88bf4434 100644 --- a/cmd/mesh-controller/watchdogs.go +++ b/cmd/mesh-controller/watchdogs.go @@ -660,8 +660,6 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li return nil } conditionsFrom = keeper - // What a consumer gave up on is read and acted on by the verb in this process (novox/hq issue 330). - deadLettersOn.Store(&busHandles{conn: server.JetStream().Conn(), js: server.JetStream().Context()}) logf := func(format string, args ...any) { fmt.Printf(format+"\n", args...) } stopHearing, err := server.HearAdvisories(logf) if err != nil { @@ -688,8 +686,6 @@ func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus li return func() { stop() stopHearing() - // No longer serving: the verb answers that it does not read DEAD_LETTERS here (novox/hq issue 330). - deadLettersOn.Store(nil) flushing, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() keeper.Close(flushing)