diff --git a/internal/link/asktool_nats_test.go b/internal/link/asktool_nats_test.go new file mode 100644 index 0000000..e9caa58 --- /dev/null +++ b/internal/link/asktool_nats_test.go @@ -0,0 +1,40 @@ +package link + +import ( + "context" + "strings" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// A declared tool check asked over a real bus, through the queue, on the link the node-engine holds: +// the tool's answer is read, and a tool nothing serves is said as such (novox/hq ADR 0240, to-be 48 §3). +func TestNatsAToolCheckIsAskedOnTheLink(t *testing.T) { + conn, js := aBus(t) + const node = "toolcheck" + served, err := conn.Subscribe(ToolSubject("sensors", "sensors_health", node), func(m *nats.Msg) { + _ = m.Respond([]byte(`{"result":{"healthy":true}}`)) + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = served.Unsubscribe() }) + if err := conn.Flush(); err != nil { + t.Fatal(err) + } + + q := &Queue{Membership: Membership{Node: node}} + q.attach(t.Context(), &natsLink{conn: conn, js: js, node: node}) + ctx, cancel := context.WithTimeout(t.Context(), 2*time.Second) + defer cancel() + healthy, why, err := q.AskTool(ctx, "sensors", "sensors_health") + if err != nil || !healthy { + t.Fatalf("a served healthy tool read as %v %q (%v)", healthy, why, err) + } + _, _, err = q.AskTool(ctx, "sensors", "nothing_serves_this") + if err == nil || !strings.Contains(err.Error(), "nothing on this machine answers") { + t.Fatalf("a tool nothing serves read as %v", err) + } +} diff --git a/internal/link/health_test.go b/internal/link/health_test.go index 3021e09..dd701ae 100644 --- a/internal/link/health_test.go +++ b/internal/link/health_test.go @@ -3,6 +3,7 @@ package link import ( "context" "encoding/json" + "strings" "testing" "time" ) @@ -70,3 +71,33 @@ func TestAToolsAnswerIsHealthyOnlyWhenItSaysSo(t *testing.T) { t.Errorf("asked on %s", ToolSubject("keycloak", "keycloak_admin_health", "anchor")) } } + +// A declared tool check is asked on the link the node-engine actually holds — the NATS link — and not +// refused as if none were open (2026-10-08: every tool check read "no link to the bus is open" while +// the link was up, because the NATS link could not Ask). With no link open, and with a link that cannot +// ask, the check fails and says which. +func TestAToolCheckIsAskedOnTheLinkTheEngineHolds(t *testing.T) { + q := &Queue{Membership: Membership{Node: "anchor"}} + if _, _, err := q.AskTool(t.Context(), "sensors", "sensors_health"); err == nil || err.Error() != "no link to the bus is open" { + t.Fatalf("with no link open: %v", err) + } + + quiet := newQuietLink() + q.attach(t.Context(), quiet) + if _, _, err := q.AskTool(t.Context(), "sensors", "sensors_health"); err == nil || err.Error() != "the link to the bus open now cannot ask a tool" { + t.Fatalf("on a link that cannot ask: %v", err) + } + q.detach(quiet) + + // The link Open returns, with no connection under it: the question must reach its Ask (and fail + // there), never stop at the queue's assertion. + held := &natsLink{node: "anchor"} + q.attach(t.Context(), held) + _, _, err := q.AskTool(t.Context(), "sensors", "sensors_health") + if err == nil { + t.Fatal("asked on no connection and was answered") + } + if strings.Contains(err.Error(), "no link") || strings.Contains(err.Error(), "cannot ask a tool") { + t.Fatalf("the NATS link was not asked: %v", err) + } +} diff --git a/internal/link/hearing_nats.go b/internal/link/hearing_nats.go index 32e8623..cbc57e2 100644 --- a/internal/link/hearing_nats.go +++ b/internal/link/hearing_nats.go @@ -220,6 +220,38 @@ func (l *natsLink) Health(ctx context.Context, node string, body []byte) error { return OverNATS{Conn: l.conn, JS: l.js}.Health(ctx, node, body) } +// Ask asks a module's tool for a declared tool check (novox/hq ADR 0240, to-be 48 §3). The link the +// node-engine holds is this one, not OverNATS: without it here the queue's ToolBus assertion failed on +// every live link and each tool check read "no link to the bus is open" while the link was up. +func (l *natsLink) Ask(ctx context.Context, subject string, body []byte) ([]byte, error) { + before := l.conn.LastError() + reply, err := OverNATS{Conn: l.conn, JS: l.js}.Ask(ctx, subject, body) + switch { + case err == nil: + return reply, nil + case errors.Is(err, nats.ErrNoResponders): + return nil, fmt.Errorf("nothing on this machine answers %s", subject) + } + if refused := l.refusal(subject, before); refused != "" { + return nil, fmt.Errorf("%s: %v%s", subject, err, refused) + } + if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, nats.ErrTimeout) { + // Said as the deadline, which the caller words with the time it gave the question. + return nil, fmt.Errorf("no answer from %s: %w", subject, context.DeadlineExceeded) + } + return nil, fmt.Errorf("%s: %w", subject, err) +} + +// What the queue asks of a live link beside Bus, each by a type assertion that fails quietly: said here +// so a link that drops one does not build. +var ( + _ Bus = (*natsLink)(nil) + _ HealthBus = (*natsLink)(nil) + _ ToolBus = (*natsLink)(nil) + _ Asker = (*natsLink)(nil) + _ Asked = (*natsLink)(nil) +) + // natsDeclaration is one declaration off the NODES stream. type natsDeclaration struct{ msg *nats.Msg } diff --git a/internal/link/queue.go b/internal/link/queue.go index 6021f5f..1ae5f5c 100644 --- a/internal/link/queue.go +++ b/internal/link/queue.go @@ -530,10 +530,13 @@ func (q *Queue) AskTool(ctx context.Context, module, tool string) (bool, string, q.mu.Lock() bus := q.bus q.mu.Unlock() - asker, ok := bus.(ToolBus) - if bus == nil || !ok { + if bus == nil { return false, "", errors.New("no link to the bus is open") } + asker, ok := bus.(ToolBus) + if !ok { + return false, "", errors.New("the link to the bus open now cannot ask a tool") + } reply, err := asker.Ask(ctx, ToolSubject(module, tool, q.Membership.Node), []byte("{}")) if err != nil { return false, "", err diff --git a/internal/link/witnessing.go b/internal/link/witnessing.go index 6b8c2e1..340eaec 100644 --- a/internal/link/witnessing.go +++ b/internal/link/witnessing.go @@ -60,12 +60,13 @@ func (q *Queue) Asker() Asker { func (l *natsLink) ReadKey(ctx context.Context, bucket, key string) ([]byte, bool, error) { subject := DirectGetSubject(bucket, key) + before := l.conn.LastError() msg, err := l.conn.RequestWithContext(ctx, subject, nil) switch { case errors.Is(err, nats.ErrNoResponders): return nil, false, fmt.Errorf("%w: the bus has no bucket %s that answers a direct get", ErrCannotAsk, bucket) case err != nil: - return nil, false, fmt.Errorf("%w: reading %s from %s: %v%s", ErrCannotAsk, key, bucket, err, l.refusal(subject)) + return nil, false, fmt.Errorf("%w: reading %s from %s: %v%s", ErrCannotAsk, key, bucket, err, l.refusal(subject, before)) } if msg.Header != nil { switch status := msg.Header.Get("Status"); status { @@ -89,13 +90,14 @@ func (l *natsLink) ReadKey(ctx context.Context, bucket, key string) ([]byte, boo func (l *natsLink) Ping(ctx context.Context, service, id string) error { subject := PingSubject(service, id) + before := l.conn.LastError() msg, err := l.conn.RequestWithContext(ctx, subject, nil) switch { case errors.Is(err, nats.ErrNoResponders): // The bus delivered the question and nothing subscribes to it: the runtime is not on the bus. return fmt.Errorf("%w: no %s on this machine is on the bus", ErrNoAnswer, service) case err != nil: - if refused := l.refusal(subject); refused != "" { + if refused := l.refusal(subject, before); refused != "" { return fmt.Errorf("%w: asking %s: %v%s", ErrCannotAsk, subject, err, refused) } return fmt.Errorf("%w: %s did not answer: %v", ErrNoAnswer, subject, err) @@ -110,11 +112,13 @@ func (l *natsLink) Ping(ctx context.Context, service, id string) error { return nil } -// refusal is the bus's last word on this connection when it refused a publish to subject, said as -// the end of an error; empty when it did not. -func (l *natsLink) refusal(subject string) string { +// refusal is the bus's word on this connection when it refused a publish to subject while this question +// was asked, said as the end of an error; empty when it did not. before is the connection's last error +// as it stood before asking: the client keeps its last error for the life of the connection, so a +// refusal from an earlier question — under a grant since widened — is not said again as this one's. +func (l *natsLink) refusal(subject string, before error) string { last := l.conn.LastError() - if last == nil { + if last == nil || last == before { return "" } words := strings.ToLower(last.Error()) diff --git a/internal/liveness/readiness.go b/internal/liveness/readiness.go index 585c08e..a5e1c08 100644 --- a/internal/liveness/readiness.go +++ b/internal/liveness/readiness.go @@ -254,8 +254,10 @@ func (p *Probes) Look(ctx context.Context, module string, c *declaration.Health) } healthy, why, err := p.AskTool(ctx, module, c.Tool) switch { + case errors.Is(err, context.DeadlineExceeded): + return false, fmt.Sprintf("no answer within %s", timeout) case err != nil: - return false, "asking it: " + firstLine(err.Error()) + return false, firstLine(err.Error()) case !healthy: return false, orSaid(why, "it answered not healthy, and not why") } diff --git a/internal/liveness/readiness_test.go b/internal/liveness/readiness_test.go index 9203844..0c39419 100644 --- a/internal/liveness/readiness_test.go +++ b/internal/liveness/readiness_test.go @@ -2,6 +2,8 @@ package liveness import ( "context" + "errors" + "fmt" "net" "net/http" "net/http/httptest" @@ -191,6 +193,25 @@ func TestAToolCheckIsAskedOfTheNodeTools(t *testing.T) { } } +// A tool that does not answer in time is said as the time it was given, once, and an error from the link is +// said as it is, without a second "asking" in front of what the link already says. +func TestAToolCheckThatGetsNoAnswerSaysHowLongItWaited(t *testing.T) { + check := &declaration.Health{Kind: declaration.HealthTool, Tool: "sensors_health", + Interval: "30s", Timeout: "5s", Looks: 2, Grace: "0s"} + p := &Probes{AskTool: func(context.Context, string, string) (bool, string, error) { + return false, "", fmt.Errorf("no answer from mesh.mod.sensors.tool.sensors_health.anchor: %w", context.DeadlineExceeded) + }} + if ok, why := p.Look(t.Context(), "sensors", check); ok || why != "no answer within 5s" { + t.Fatalf("%v %q", ok, why) + } + p.AskTool = func(context.Context, string, string) (bool, string, error) { + return false, "", errors.New("no link to the bus is open") + } + if ok, why := p.Look(t.Context(), "sensors", check); ok || why != "no link to the bus is open" { + t.Fatalf("%v %q", ok, why) + } +} + // The runtime's state of a container's own check is read out of its whole state: a container that carries // none has no Health in it, and that is not an error (a template naming it would refuse every container). func TestTheRuntimesCheckIsReadOutOfTheContainersState(t *testing.T) { diff --git a/merge-check.sh b/merge-check.sh index 4fc67c7..dabdee8 100755 --- a/merge-check.sh +++ b/merge-check.sh @@ -21,6 +21,36 @@ if [ -n "$unformatted" ]; then fi CGO_ENABLED=0 go vet ./... +# **The tests that need a real bus** (MESH_TEST_NATS) run against a throwaway nats-server with JetStream when +# the toolchain holds one, so the link the node-engine dials is tested as it is used (hq issue 331: every +# unit test used fake links, and the real one could not ask a tool). Said when it cannot: a skipped test +# is not a passed one. A caller that sets MESH_TEST_NATS keeps its own bus. +if [ -z "${MESH_TEST_NATS:-}" ]; then + if command -v nats-server >/dev/null 2>&1; then + bus_dir=$(mktemp -d) + bus_port=$((20000 + $$ % 20000)) + nats-server -js -a 127.0.0.1 -p "$bus_port" -sd "$bus_dir" >"$bus_dir/log" 2>&1 & + bus_pid=$! + trap 'kill "$bus_pid" 2>/dev/null || true; rm -rf "$bus_dir"' EXIT + up=0 + for _ in 1 2 3 4 5 6 7 8 9 10; do + if grep -q "Server is ready" "$bus_dir/log"; then + up=1 + break + fi + sleep 0.5 + done + if [ "$up" -ne 1 ]; then + echo "the throwaway nats-server did not start:" + cat "$bus_dir/log" + exit 1 + fi + export MESH_TEST_NATS="nats://127.0.0.1:$bus_port" + else + echo "NOT RUN AGAINST A BUS: the toolchain holds no nats-server; the tests that need one are skipped" + fi +fi + if command -v gcc >/dev/null 2>&1; then CGO_ENABLED=1 go test -race -count=1 ./... else