From cc7fb99f29344246103da64d4ce6ccb425d587fd Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 18:45:04 +0200 Subject: [PATCH] Answer a panicking verb with an error, say flag errors in the answer, and leave refused logins out of D15 Read verbs now run in the serving process, where a panic would end every call; refused logins are nobody's reconnect loop (review of hq issue 327). --- cmd/mesh-controller/conditions.go | 3 ++ cmd/mesh-controller/handacts.go | 1 + cmd/mesh-controller/queue.go | 1 + cmd/mesh-controller/reconnects.go | 12 +++++++- cmd/mesh-controller/reconnects_test.go | 33 ++++++++++++++++++++++ cmd/mesh-controller/replay327_test.go | 38 ++++---------------------- cmd/mesh-controller/servingbus.go | 11 ++++++++ internal/link/calls.go | 14 ++++++++++ internal/link/calls_test.go | 30 ++++++++++++++++++++ 9 files changed, 110 insertions(+), 33 deletions(-) diff --git a/cmd/mesh-controller/conditions.go b/cmd/mesh-controller/conditions.go index c26b181b..eed26b61 100644 --- a/cmd/mesh-controller/conditions.go +++ b/cmd/mesh-controller/conditions.go @@ -121,6 +121,7 @@ func conditionsCommand(ctx context.Context, args []string) error { func listConditions(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions", flag.ContinueOnError) + usageTo(set, w) scope := set.String("scope", "", "only this scope: "+strings.Join(conditions.Scopes, ", ")) severity := set.String("severity", "", "only urgent, or only warning") machine := set.String("machine", "", "only those about this machine") @@ -236,6 +237,7 @@ func conditionLines(list []conditions.Condition, now time.Time) []string { func showCondition(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions show", flag.ContinueOnError) + usageTo(set, w) asJSON := set.Bool("json", false, "as data") rest, err := parseAround(set, args) if err != nil { @@ -381,6 +383,7 @@ func parseFor(s string) (time.Duration, error) { func conditionHistory(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("conditions history", flag.ContinueOnError) + usageTo(set, w) days := set.Int("days", 7, "how many days back, at most 90") key := set.String("key", "", "only this condition") asJSON := set.Bool("json", false, "as data") diff --git a/cmd/mesh-controller/handacts.go b/cmd/mesh-controller/handacts.go index 0cc4e3c0..38aea23b 100644 --- a/cmd/mesh-controller/handacts.go +++ b/cmd/mesh-controller/handacts.go @@ -254,6 +254,7 @@ func handActCommand(ctx context.Context, args []string) error { // listHandActs is `hand-acts`: what was done by hand lately, and the causes done more than once. func listHandActs(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("hand-acts", flag.ContinueOnError) + usageTo(set, w) days := set.Int("days", 14, "how many days back") asJSON := set.Bool("json", false, "as data") if _, err := parseAround(set, args); err != nil { diff --git a/cmd/mesh-controller/queue.go b/cmd/mesh-controller/queue.go index e4c97e24..d36a6814 100644 --- a/cmd/mesh-controller/queue.go +++ b/cmd/mesh-controller/queue.go @@ -52,6 +52,7 @@ func queueCommand(ctx context.Context, args []string) error { // listQueue is `queue`: the build queue, as a person reads it or as JSON. func listQueue(ctx context.Context, args []string, w io.Writer) error { set := flag.NewFlagSet("queue", flag.ContinueOnError) + usageTo(set, w) asJSON := set.Bool("json", false, "the queue as JSON") if _, err := parseAround(set, args); err != nil { return err diff --git a/cmd/mesh-controller/reconnects.go b/cmd/mesh-controller/reconnects.go index 1e50e5c3..ed12cfe7 100644 --- a/cmd/mesh-controller/reconnects.go +++ b/cmd/mesh-controller/reconnects.go @@ -64,7 +64,12 @@ func probeReconnects(ctx context.Context, d *doctor) ([]conditions.Observation, if len(on) == 0 { return nil, nil // no bus module assigned: a mesh whose bus is not the mesh's module } - read, err := askClosed(ctx, d.js.Conn(), on[0]) + return reconnectsOn(ctx, d.js.Conn(), on[0]) +} + +// reconnectsOn asks the bus module on its machine and says who reconnects in a loop. +func reconnectsOn(ctx context.Context, conn *nats.Conn, node string) ([]conditions.Observation, error) { + read, err := askClosed(ctx, conn, node) if isNothingServes(err) { // A bus module older than its tool has nothing it can say, which is not a failure of the probe. return nil, nil @@ -100,6 +105,11 @@ func reconnecting(read closedConnections) []conditions.Observation { } var out []conditions.Observation for _, u := range read.Users { + if u.User == "" || strings.HasPrefix(u.User, "(") { + // Refused before it logged in: no client of the mesh's, so nobody's reconnect loop. A login + // refused again and again is a question of its own, not this probe's. + continue + } perHour := float64(u.Dropped) / hours if perHour <= reconnectBound { continue diff --git a/cmd/mesh-controller/reconnects_test.go b/cmd/mesh-controller/reconnects_test.go index 60334df5..d6b26b95 100644 --- a/cmd/mesh-controller/reconnects_test.go +++ b/cmd/mesh-controller/reconnects_test.go @@ -112,6 +112,8 @@ func TestAClientReconnectingInALoopIsSaid(t *testing.T) { {"name":"ace.node-tools","closed":40,"dropped":40,"reasons":{"Stale Connection":30,"Read Error":10}}]}, {"user":"controller","closed":300,"dropped":0,"names":[ {"name":"mesh-controller verb delivery-check","closed":300,"dropped":0,"reasons":{"Client Closed":300}}]}, + {"user":"(no user: refused before it logged in)","closed":90,"dropped":90,"names":[{"name":"","closed":90, + "dropped":90,"reasons":{"Authentication Failure":90}}]}, {"user":"node.anchor","closed":5,"dropped":5,"names":[{"name":"mesh-host/anchor","closed":5,"dropped":5, "reasons":{"Read Error":5}}]}]}`), &read); err != nil { t.Fatal(err) @@ -154,3 +156,34 @@ func TestTheBusModuleIsAskedWhoClosedConnections(t *testing.T) { t.Fatalf("%+v", read) } } + +// A bus module older than the tool answers nothing to the question: the probe passes over it quietly. +func TestAnOlderBusModuleIsPassedOverQuietly(t *testing.T) { + bus := testbus.Start(t) + conn, err := nats.Connect(bus.ClientURL()) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + found, err := reconnectsOn(context.Background(), conn, "anchor") + if err != nil || len(found) != 0 { + t.Fatalf("a bus module without the tool: %v %v", found, err) + } +} + +// A verb answered in the serving controller says its flag errors in its answer, not in the controller's log. +func TestAFlagErrorIsSaidInTheAnswer(t *testing.T) { + bus := testbus.Start(t) + serving, err := broker.Dial(bus.ClientURL()) + if err != nil { + t.Fatal(err) + } + defer serving.Close() + was := servingBus.Load() + servingBus.Store(serving) + defer servingBus.Store(was) + answer, read := readHere(context.Background(), []string{"hand-acts", "--bogus"}) + if !read || answer.OK || !strings.Contains(answer.Output, "flag provided but not defined") { + t.Fatalf("answered %+v", answer) + } +} diff --git a/cmd/mesh-controller/replay327_test.go b/cmd/mesh-controller/replay327_test.go index 20e02cca..559debe9 100644 --- a/cmd/mesh-controller/replay327_test.go +++ b/cmd/mesh-controller/replay327_test.go @@ -3,10 +3,8 @@ package main import ( "context" "encoding/json" - "os" "os/exec" "path/filepath" - "sync" "testing" "time" @@ -70,45 +68,21 @@ func TestReplay327(t *testing.T) { } } -// The binary a verb run as a process of its own is, built once for the tests that run one. -var verbBinary struct { - once sync.Once - path string - err error -} - -// asAProcess makes a verb that runs as a process of its own run this package's binary, on the bus at url. +// asAProcess makes a verb that runs as a process of its own run this package's binary, on the bus at url: +// built into the test's own directory, which goes with the test. func asAProcess(t *testing.T, url string) { t.Helper() - verbBinary.once.Do(func() { - dir, err := os.MkdirTemp("", "mesh-controller-verb-") - if err != nil { - verbBinary.err = err - return - } - verbBinary.path = filepath.Join(dir, "mesh-controller") - out, err := exec.Command("go", "build", "-o", verbBinary.path, ".").CombinedOutput() - if err != nil { - verbBinary.err = &buildError{out: string(out), err: err} - } - }) - if verbBinary.err != nil { - t.Fatalf("the controller could not be built to run a verb as its own process: %v", verbBinary.err) + path := filepath.Join(t.TempDir(), "mesh-controller") + if out, err := exec.Command("go", "build", "-o", path, ".").CombinedOutput(); err != nil { + t.Fatalf("the controller could not be built to run a verb as its own process: %v: %s", err, out) } was := ownImage - ownImage = func() string { return verbBinary.path } + ownImage = func() string { return path } t.Cleanup(func() { ownImage = was }) t.Setenv(broker.NATSVar, url) t.Setenv(broker.CertificateVar, "") } -type buildError struct { - out string - err error -} - -func (e *buildError) Error() string { return e.err.Error() + ": " + e.out } - // closedNames are the names of the connections the bus saw closed. func closedNames(t *testing.T, bus *server.Server) []string { t.Helper() diff --git a/cmd/mesh-controller/servingbus.go b/cmd/mesh-controller/servingbus.go index 48d3cfcb..727f4fb9 100644 --- a/cmd/mesh-controller/servingbus.go +++ b/cmd/mesh-controller/servingbus.go @@ -1,7 +1,10 @@ package main import ( + "flag" "fmt" + "io" + "os" "sync/atomic" "github.com/novox/mesh-controller/internal/broker" @@ -36,3 +39,11 @@ func aBus() (*broker.JetStream, error) { } return js, nil } + +// usageTo sends a command's flag errors and usage to where its answer goes when that is not this process's +// output: a verb answered in the serving controller says them in its answer, not in the controller's log. +func usageTo(set *flag.FlagSet, w io.Writer) { + if w != os.Stdout { + set.SetOutput(w) + } +} diff --git a/internal/link/calls.go b/internal/link/calls.go index 49d13e83..daaaf6e1 100644 --- a/internal/link/calls.go +++ b/internal/link/calls.go @@ -7,6 +7,7 @@ import ( "fmt" "log" "regexp" + "runtime/debug" "sort" "strconv" "strings" @@ -575,6 +576,19 @@ func (l *CallLog) serveCallWithin(seat, verb string, args json.RawMessage, reply done := make(chan outcome, 1) go func() { var body []byte + // **A handler that panics is an answer, not a dead controller.** A verb answered in the serving + // process (novox/hq issue 327: conditions, hand-acts, queue, dead-letters) runs on this goroutine, + // where a panic would take every other call and the controller with it. + defer func() { + if p := recover(); p != nil { + if logger != nil { + logger.Printf("%s.%s: call %s panicked: %v\n%s", seat, verb, c.ID, p, debug.Stack()) + } + body, _ = json.Marshal(map[string]any{"error": fmt.Sprintf("%s failed inside the controller and "+ + "answered nothing: %v. It is in the controller's log", verb, p)}) + done <- outcome{body, true} + } + }() result, err := handle(ctx, args) failed := err != nil if errors.Is(err, ErrHandingOver) { diff --git a/internal/link/calls_test.go b/internal/link/calls_test.go index b06ee147..9c2d83fd 100644 --- a/internal/link/calls_test.go +++ b/internal/link/calls_test.go @@ -11,6 +11,10 @@ import ( "sync" "testing" "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/testbus" ) // answers collects what a call answered, and fails a second answer: the bus permits one. @@ -223,3 +227,29 @@ func TestAJoinTokenIsShownToItsCallerAndNotKept(t *testing.T) { } } } + +// A verb whose handler panics answers an error, and the controller serves the next call (review of +// novox/hq issue 327: read verbs are answered in the serving process now). +func TestAHandlerThatPanicsAnswersAnError(t *testing.T) { + conn, err := nats.Connect(testbus.URL(t)) + if err != nil { + t.Fatal(err) + } + defer conn.Close() + stop, err := OverNATS{Conn: conn}.ServeSeatTools("panicky", map[string]ToolHandler{ + "boom": func(context.Context, json.RawMessage) (any, error) { panic("nil map") }, + "fine": func(context.Context, json.RawMessage) (any, error) { return "ok", nil }, + }, nil) + if err != nil { + t.Fatal(err) + } + defer stop() + answer, err := AskMeshSeatTool(context.Background(), conn, "panicky", "boom", map[string]any{}, 5*time.Second) + if err != nil || !strings.Contains(answer.Error, "failed inside the controller") { + t.Fatalf("a panic answered %+v (%v)", answer, err) + } + answer, err = AskMeshSeatTool(context.Background(), conn, "panicky", "fine", map[string]any{}, 5*time.Second) + if err != nil || answer.Error != "" { + t.Fatalf("the call after a panic answered %+v (%v)", answer, err) + } +}