From 809ab4cf29de1d3a214bd3f78baa1edc4b74ce97 Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 16:44:01 +0200 Subject: [PATCH 1/2] Ask a declared tool check on the NATS link the node-engine holds The queue asks a tool through a ToolBus assertion that only OverNATS satisfied; the link the engine actually holds (natsLink) had no Ask, so every tool check failed as "no link to the bus is open" while the link was up (hq issue 331). Give natsLink Ask, assert at build time every optional interface the queue expects of it, and tell a link that cannot ask apart from no link. --- internal/link/asktool_nats_test.go | 40 ++++++++++++++++++++++++++++++ internal/link/health_test.go | 31 +++++++++++++++++++++++ internal/link/hearing_nats.go | 27 ++++++++++++++++++++ internal/link/queue.go | 7 ++++-- 4 files changed, 103 insertions(+), 2 deletions(-) create mode 100644 internal/link/asktool_nats_test.go 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..b2c4885 100644 --- a/internal/link/hearing_nats.go +++ b/internal/link/hearing_nats.go @@ -220,6 +220,33 @@ 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) { + reply, err := OverNATS{Conn: l.conn, JS: l.js}.Ask(ctx, subject, body) + switch { + case errors.Is(err, nats.ErrNoResponders): + return nil, fmt.Errorf("nothing on this machine answers %s", subject) + case err != nil: + if refused := l.refusal(subject); refused != "" { + return nil, fmt.Errorf("asking %s: %v%s", subject, err, refused) + } + return nil, fmt.Errorf("asking %s: %w", subject, err) + } + return reply, nil +} + +// 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 From 5e189c2fcad250369894acea1eaf8688990baa7e Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 17:01:15 +0200 Subject: [PATCH 2/2] Say a tool check's errors once and as they stand, and test the link on a real bus The evidence read "asking it: asking ": the words are now the link's own. A deadline is said as the time the check gave it. A refusal is said only when the bus refused this question, not an earlier one under a grant since widened. merge-check.sh runs the tests that need a bus against a throwaway nats-server when the toolchain has one, and says so when it does not (hq issue 331). --- internal/link/hearing_nats.go | 17 ++++++++++------ internal/link/witnessing.go | 16 +++++++++------ internal/liveness/readiness.go | 4 +++- internal/liveness/readiness_test.go | 21 ++++++++++++++++++++ merge-check.sh | 30 +++++++++++++++++++++++++++++ 5 files changed, 75 insertions(+), 13 deletions(-) diff --git a/internal/link/hearing_nats.go b/internal/link/hearing_nats.go index b2c4885..cbc57e2 100644 --- a/internal/link/hearing_nats.go +++ b/internal/link/hearing_nats.go @@ -224,17 +224,22 @@ func (l *natsLink) Health(ctx context.Context, node string, body []byte) error { // 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) - case err != nil: - if refused := l.refusal(subject); refused != "" { - return nil, fmt.Errorf("asking %s: %v%s", subject, err, refused) - } - return nil, fmt.Errorf("asking %s: %w", subject, err) } - return reply, nil + 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 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