Merge pull request 'Ask a declared tool check on the NATS link the node-engine holds (hq issue 331)' (#57) from fix/tool-check-asks-on-the-link into main
This commit was merged in pull request #57.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 }
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user