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).
This commit is contained in:
jochen
2026-10-08 21:08:34 +02:00
parent 1e04670052
commit cc7fb99f29
9 changed files with 110 additions and 33 deletions
+3
View File
@@ -121,6 +121,7 @@ func conditionsCommand(ctx context.Context, args []string) error {
func listConditions(ctx context.Context, args []string, w io.Writer) error { func listConditions(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("conditions", flag.ContinueOnError) set := flag.NewFlagSet("conditions", flag.ContinueOnError)
usageTo(set, w)
scope := set.String("scope", "", "only this scope: "+strings.Join(conditions.Scopes, ", ")) scope := set.String("scope", "", "only this scope: "+strings.Join(conditions.Scopes, ", "))
severity := set.String("severity", "", "only urgent, or only warning") severity := set.String("severity", "", "only urgent, or only warning")
machine := set.String("machine", "", "only those about this machine") 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 { func showCondition(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("conditions show", flag.ContinueOnError) set := flag.NewFlagSet("conditions show", flag.ContinueOnError)
usageTo(set, w)
asJSON := set.Bool("json", false, "as data") asJSON := set.Bool("json", false, "as data")
rest, err := parseAround(set, args) rest, err := parseAround(set, args)
if err != nil { 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 { func conditionHistory(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("conditions history", flag.ContinueOnError) set := flag.NewFlagSet("conditions history", flag.ContinueOnError)
usageTo(set, w)
days := set.Int("days", 7, "how many days back, at most 90") days := set.Int("days", 7, "how many days back, at most 90")
key := set.String("key", "", "only this condition") key := set.String("key", "", "only this condition")
asJSON := set.Bool("json", false, "as data") asJSON := set.Bool("json", false, "as data")
+1
View File
@@ -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. // 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 { func listHandActs(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("hand-acts", flag.ContinueOnError) set := flag.NewFlagSet("hand-acts", flag.ContinueOnError)
usageTo(set, w)
days := set.Int("days", 14, "how many days back") days := set.Int("days", 14, "how many days back")
asJSON := set.Bool("json", false, "as data") asJSON := set.Bool("json", false, "as data")
if _, err := parseAround(set, args); err != nil { if _, err := parseAround(set, args); err != nil {
+1
View File
@@ -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. // 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 { func listQueue(ctx context.Context, args []string, w io.Writer) error {
set := flag.NewFlagSet("queue", flag.ContinueOnError) set := flag.NewFlagSet("queue", flag.ContinueOnError)
usageTo(set, w)
asJSON := set.Bool("json", false, "the queue as JSON") asJSON := set.Bool("json", false, "the queue as JSON")
if _, err := parseAround(set, args); err != nil { if _, err := parseAround(set, args); err != nil {
return err return err
+11 -1
View File
@@ -64,7 +64,12 @@ func probeReconnects(ctx context.Context, d *doctor) ([]conditions.Observation,
if len(on) == 0 { if len(on) == 0 {
return nil, nil // no bus module assigned: a mesh whose bus is not the mesh's module 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) { if isNothingServes(err) {
// A bus module older than its tool has nothing it can say, which is not a failure of the probe. // A bus module older than its tool has nothing it can say, which is not a failure of the probe.
return nil, nil return nil, nil
@@ -100,6 +105,11 @@ func reconnecting(read closedConnections) []conditions.Observation {
} }
var out []conditions.Observation var out []conditions.Observation
for _, u := range read.Users { 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 perHour := float64(u.Dropped) / hours
if perHour <= reconnectBound { if perHour <= reconnectBound {
continue continue
+33
View File
@@ -112,6 +112,8 @@ func TestAClientReconnectingInALoopIsSaid(t *testing.T) {
{"name":"ace.node-tools","closed":40,"dropped":40,"reasons":{"Stale Connection":30,"Read Error":10}}]}, {"name":"ace.node-tools","closed":40,"dropped":40,"reasons":{"Stale Connection":30,"Read Error":10}}]},
{"user":"controller","closed":300,"dropped":0,"names":[ {"user":"controller","closed":300,"dropped":0,"names":[
{"name":"mesh-controller verb delivery-check","closed":300,"dropped":0,"reasons":{"Client Closed":300}}]}, {"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, {"user":"node.anchor","closed":5,"dropped":5,"names":[{"name":"mesh-host/anchor","closed":5,"dropped":5,
"reasons":{"Read Error":5}}]}]}`), &read); err != nil { "reasons":{"Read Error":5}}]}]}`), &read); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -154,3 +156,34 @@ func TestTheBusModuleIsAskedWhoClosedConnections(t *testing.T) {
t.Fatalf("%+v", read) 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)
}
}
+6 -32
View File
@@ -3,10 +3,8 @@ package main
import ( import (
"context" "context"
"encoding/json" "encoding/json"
"os"
"os/exec" "os/exec"
"path/filepath" "path/filepath"
"sync"
"testing" "testing"
"time" "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. // asAProcess makes a verb that runs as a process of its own run this package's binary, on the bus at url:
var verbBinary struct { // built into the test's own directory, which goes with the test.
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) { func asAProcess(t *testing.T, url string) {
t.Helper() t.Helper()
verbBinary.once.Do(func() { path := filepath.Join(t.TempDir(), "mesh-controller")
dir, err := os.MkdirTemp("", "mesh-controller-verb-") if out, err := exec.Command("go", "build", "-o", path, ".").CombinedOutput(); err != nil {
if err != nil { t.Fatalf("the controller could not be built to run a verb as its own process: %v: %s", err, out)
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 was := ownImage
ownImage = func() string { return verbBinary.path } ownImage = func() string { return path }
t.Cleanup(func() { ownImage = was }) t.Cleanup(func() { ownImage = was })
t.Setenv(broker.NATSVar, url) t.Setenv(broker.NATSVar, url)
t.Setenv(broker.CertificateVar, "") 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. // closedNames are the names of the connections the bus saw closed.
func closedNames(t *testing.T, bus *server.Server) []string { func closedNames(t *testing.T, bus *server.Server) []string {
t.Helper() t.Helper()
+11
View File
@@ -1,7 +1,10 @@
package main package main
import ( import (
"flag"
"fmt" "fmt"
"io"
"os"
"sync/atomic" "sync/atomic"
"github.com/novox/mesh-controller/internal/broker" "github.com/novox/mesh-controller/internal/broker"
@@ -36,3 +39,11 @@ func aBus() (*broker.JetStream, error) {
} }
return js, nil 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)
}
}
+14
View File
@@ -7,6 +7,7 @@ import (
"fmt" "fmt"
"log" "log"
"regexp" "regexp"
"runtime/debug"
"sort" "sort"
"strconv" "strconv"
"strings" "strings"
@@ -575,6 +576,19 @@ func (l *CallLog) serveCallWithin(seat, verb string, args json.RawMessage, reply
done := make(chan outcome, 1) done := make(chan outcome, 1)
go func() { go func() {
var body []byte 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) result, err := handle(ctx, args)
failed := err != nil failed := err != nil
if errors.Is(err, ErrHandingOver) { if errors.Is(err, ErrHandingOver) {
+30
View File
@@ -11,6 +11,10 @@ import (
"sync" "sync"
"testing" "testing"
"time" "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. // 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)
}
}