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.
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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=<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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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.",
|
||||
|
||||
@@ -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{}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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.<machine>` the
|
||||
// node-engine, `<machine>.node-tools` the tool runner, `controller` the controller, `<machine>.<module>` 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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user