Say a tool check's errors once and as they stand, and test the link on a real bus
mesh/merge-gate pass: builds mesh-host → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without (4 of 4 compose)
mesh/repo-check pass: THE CHANGE ALTERS ITS OWN CHECK (merge-check.sh): main's version judged it; the change's judges the pull requests after it merges; it…
mesh/delivery delivered

The evidence read "asking it: asking <subject>": 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).
This commit is contained in:
jochen
2026-10-08 17:01:15 +02:00
parent 809ab4cf29
commit 5e189c2fca
5 changed files with 75 additions and 13 deletions
+11 -6
View File
@@ -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
+10 -6
View File
@@ -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())
+3 -1
View File
@@ -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")
}
+21
View File
@@ -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) {
+30
View File
@@ -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