From 1e04670052ac6e1cd97ac249fabef3f8f58fb497 Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 17:52:36 +0200 Subject: [PATCH] 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 {