Judge how a module says it is ready, beside whether it stays up (hq ADR 0240, to-be 48 Phase B)

Liveness alone could not see a web application whose port was open and whose
program ran while every request hung for eleven hours (issue 145). A resource
now carries the `health` its module declared: the engine makes http and tcp
looks itself from the machine to the endpoint's published port, reads a unit's
readiness from the show it already makes, hands an exec command or the image's
own check to the runtime as the container's check with the declared timing and
reads its state from the inspect it already makes, and asks a module's tool on
its own node tools. Starting until the check passed, unhealthy once its looks
after the grace fail the declared number of times; never more looks than the
measured budget; nothing restarted. The statement says contract 2, which tells
the controller this engine may be sent the field.
This commit is contained in:
jochen
2026-10-07 14:08:32 +02:00
parent 41f908b803
commit bdd44154cc
16 changed files with 1260 additions and 19 deletions
+18 -2
View File
@@ -1164,6 +1164,11 @@ func runLink(ctx context.Context, opts options) error {
// And whether what it runs stays up, looked at on its own clock and said when it changes (novox/hq
// ADR 0240) — between applies, which is when a container crash-loops.
if judging != nil {
// The looks it makes itself — http, tcp, a module's own tool — each at its declared interval,
// spaced to the budget (ADR 0240 Phase B); a tool is asked of this machine's node tools over the
// link open at the time.
judging.Probes = &liveness.Probes{AskTool: queue.AskTool}
go judging.Probe(aside)
go judgeWhatRuns(aside, judging, queue, say)
}
// And the core builds this host placed are judged, whichever host placed them (to-be 45 §8).
@@ -1373,6 +1378,7 @@ func judgeWhatRuns(ctx context.Context, j *liveness.Judge, queue *link.Queue, sa
defer ticker.Stop()
var lastSaid time.Time
owed := true
spaced := 1.0
for {
select {
case <-ctx.Done():
@@ -1381,6 +1387,14 @@ func judgeWhatRuns(ctx context.Context, j *liveness.Judge, queue *link.Queue, sa
}
st, changed := j.Look(ctx)
owed = owed || changed
// Never more looks than the budget (ADR 0240): said when the engine has to space its own out.
if sp := j.Spacing(); sp != spaced {
if sp > 1 {
say(fmt.Sprintf("the declared health checks here would cost more than %d looks a minute; the engine's "+
"own looks are spaced %.1f times their declared interval", liveness.Budget, sp))
}
spaced = sp
}
since := time.Since(lastSaid)
if !owed && !(!st.Healthy() && since >= sayUnhealthyAgain) && since < sayAnyway {
continue
@@ -1401,11 +1415,13 @@ func judgeWhatRuns(ctx context.Context, j *liveness.Judge, queue *link.Queue, sa
// healthAsReported is a statement as the report and the event carry it.
func healthAsReported(st liveness.Statement) *link.Health {
h := &link.Health{Contract: link.LivenessContract, At: st.At.UTC(), Resources: []link.ResourceHealth{}}
// ReadinessContract: this engine reads a resource's declared `health` and judges it (ADR 0240 Phase
// B), which is what tells the controller it may be sent the field.
h := &link.Health{Contract: link.ReadinessContract, At: st.At.UTC(), Resources: []link.ResourceHealth{}}
for _, r := range st.Resources {
h.Resources = append(h.Resources, link.ResourceHealth{Module: r.Module, Resource: r.ID, Kind: r.Kind,
Target: r.Target, State: r.State, Reason: r.Reason, Since: r.Since.UTC(), Streak: r.Streak,
Restarts: r.Restarts})
Restarts: r.Restarts, Check: r.CheckOf(), Needs: r.NeedsOf()})
}
return h
}
+26
View File
@@ -1987,6 +1987,12 @@ func containerSpecReading(r *declaration.Container, declares, reads map[string]s
if r.Schedule != "" {
b.WriteString("schedule " + r.Schedule + "\n")
}
// And the check the runtime runs as its own (novox/hq ADR 0240 rule 3): an exec command or the
// image's own, with the declared timing, is fixed when the container is created. Only those two:
// an http, tcp, unit or tool check the engine makes itself, so declaring one recreates nothing.
if h := r.Health; h.RunByRuntime() {
b.WriteString("health " + strings.Join(runtimeCheckArgs(h), " ") + "\n")
}
// **What this container reads is part of what it is.**
//
// A container takes its environment and its mounted files once, at start, and never looks
@@ -2230,6 +2236,12 @@ func applyContainer(ctx context.Context, r *declaration.Container, run Runner,
for _, v := range r.Volumes {
args = append(args, "--volume", v)
}
// **The one check the container carries is the one the module declared** (ADR 0240 rule 3): an
// exec command, or the image's own adopted by name, with the declared timing. Nothing else on the
// machine sets a container's check.
if r.Health.RunByRuntime() {
args = append(args, runtimeCheckArgs(r.Health)...)
}
for _, h := range r.Hosts {
// Written into the container's own hosts file by the runtime. Per container rather than
// by editing the machine's resolver configuration: that file belongs to something else on
@@ -2287,6 +2299,20 @@ func applyContainer(ctx context.Context, r *declaration.Container, run Runner,
return out, nil
}
// runtimeCheckArgs is a declared check as the runtime takes it (novox/hq ADR 0240 rule 3, to-be 48 §3):
// its timing, and for an exec check its command, run by the container's shell. The runtime's own retries
// are the declared failing looks, and its start period the grace — so its retries and start period do the
// timing, with no execution per look from outside. For the image's own check no command: the image's
// stays, under the declared timing.
func runtimeCheckArgs(h *declaration.Health) []string {
args := []string{"--health-interval", h.Interval, "--health-timeout", h.Timeout,
"--health-retries", strconv.Itoa(h.Looks), "--health-start-period", h.Grace}
if h.Kind == declaration.HealthExec {
args = append([]string{"--health-cmd", h.Command}, args...)
}
return args
}
// applyRunOnce runs a container to completion, once, and requires it to exit 0 (novox/hq ADR 0052).
//
// It is a step, not a service: the module's own code seeding a store, migrating a schema or
+84
View File
@@ -0,0 +1,84 @@
package apply
import (
"context"
"slices"
"strings"
"testing"
"github.com/novox/mesh-host/internal/declaration"
"github.com/novox/mesh-host/internal/store"
)
// The node-engine owns every verdict, and nothing else sets a container's check (novox/hq ADR 0240 rule 3,
// "how it is checked"): a declared command becomes the container's check with the declared timing, an
// adopted image check keeps the image's command under the declared timing, and a check the engine makes
// itself — http, tcp — sets nothing on the container and recreates nothing.
func TestADeclaredCommandBecomesTheContainersCheckAndNothingElseSetsOne(t *testing.T) {
pinned := "postgres@sha256:" + strings.Repeat("a", 64)
ranWith := func(health string) []string {
t.Helper()
var ran []string
run := func(_ context.Context, cmd string, args ...string) (string, error) {
if cmd == "docker" && len(args) > 0 && args[0] == "run" {
ran = args
return "deadbeef\n", nil
}
return "", nil
}
field := ""
if health != "" {
field = `,"health":` + health
}
d := parseTrusted(t, `{"declaration":1,"resources":[
{"id":"db","type":"container","name":"db","image":"`+pinned+`","ports":["31001:5432"]`+field+`}
]}`)
_, _, _ = Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil, nil)
if ran == nil {
t.Fatal("the container was not run")
}
return ran
}
healthFlags := func(args []string) []string {
var out []string
for i, a := range args {
if strings.HasPrefix(a, "--health") || a == "--no-healthcheck" {
out = append(out, a)
if i+1 < len(args) {
out = append(out, args[i+1])
}
}
}
return out
}
exec := healthFlags(ranWith(`{"kind":"exec","command":"pg_isready -q","interval":"30s","timeout":"5s","looks":3,"grace":"1m0s"}`))
want := []string{"--health-cmd", "pg_isready -q", "--health-interval", "30s", "--health-timeout", "5s",
"--health-retries", "3", "--health-start-period", "1m0s"}
if !slices.Equal(exec, want) {
t.Errorf("a declared command was run as %v, want %v", exec, want)
}
adopted := healthFlags(ranWith(`{"kind":"runtime","interval":"20s","timeout":"3s","looks":2,"grace":"30s"}`))
if slices.Contains(adopted, "--health-cmd") || !slices.Contains(adopted, "20s") || !slices.Contains(adopted, "2") {
t.Errorf("an adopted image check was run as %v: the image's command under the declared timing", adopted)
}
for _, h := range []string{"", `{"kind":"http","endpoint":"web","port":31001,"path":"/","interval":"30s","timeout":"5s","looks":3,"grace":"1m0s"}`,
`{"kind":"tcp","port":31001,"interval":"30s","timeout":"5s","looks":3,"grace":"1m0s"}`} {
if got := healthFlags(ranWith(h)); len(got) > 0 {
t.Errorf("a container whose check the engine makes itself (%q) was given the runtime's: %v", h, got)
}
}
// Declaring an http check recreates nothing; declaring a command does.
plain := &declaration.Container{ID: "db", Name: "db", Image: pinned}
withHTTP := *plain
withHTTP.Health = &declaration.Health{Kind: "http", Port: 31001, Path: "/", Interval: "30s", Timeout: "5s", Looks: 3, Grace: "1m0s"}
withExec := *plain
withExec.Health = &declaration.Health{Kind: "exec", Command: "true", Interval: "30s", Timeout: "5s", Looks: 3, Grace: "1m0s"}
if containerSpec(plain, inputs{}) != containerSpec(&withHTTP, inputs{}) {
t.Error("declaring an http check would recreate the container")
}
if containerSpec(plain, inputs{}) == containerSpec(&withExec, inputs{}) {
t.Error("declaring a command would not reach a running container: it is fixed when the container is made")
}
}
+15 -3
View File
@@ -603,6 +603,10 @@ type Process struct {
// (to-be 45 §8, rule 8). A build so declared that is not healthy in bound is left running and said
// as urgent; the build before it is never started against the newer data.
NotReversible string `json:"not-reversible,omitempty"`
// Health is how this resource is ready (novox/hq ADR 0240 rule 2, Phase B): one kind and its
// timing, judged by the node-engine beside liveness. Absent: judged alive or not, and nothing more.
Health *Health `json:"health,omitempty"`
}
func (d *Process) Identity() string { return d.ID }
@@ -634,7 +638,7 @@ func ProcessNameProblem(name string) string {
}
func (d *Process) validate(where string, _ bool) []string {
var problems []string
problems := d.Health.problems(where, false, !d.RunOnce && d.Schedule == "")
if problem := ProcessNameProblem(d.Name); problem != "" {
problems = append(problems, where+": "+problem)
}
@@ -782,6 +786,10 @@ type Service struct {
// said by the controller, which knows the found tunnel's key is this node's own: without that,
// starting this unit on the found one's port would drop every peer's packets.
TakesOver *TakeOver `json:"takes-over,omitempty"`
// Health is how this resource is ready (novox/hq ADR 0240 rule 2, Phase B): one kind and its
// timing, judged by the node-engine beside liveness. Absent: judged alive or not, and nothing more.
Health *Health `json:"health,omitempty"`
}
// TakeOver is a found tunnel a service replaces: its interface, the unit that raised it, and its
@@ -810,7 +818,7 @@ const (
func (s *Service) UserScoped() bool { return s.Scope == ScopeUser }
func (s *Service) validate(where string, _ bool) []string {
var problems []string
problems := s.Health.problems(where, false, s.State == "running")
if s.Unit == "" {
problems = append(problems, where+": a service needs a unit")
}
@@ -1122,6 +1130,10 @@ type Container struct {
// offline job says *before*, not *instead of*; a recurring window is the case order cannot
// express, and the only one this serves.
WhileStopped []string `json:"while-stopped,omitempty"`
// Health is how this resource is ready (novox/hq ADR 0240 rule 2, Phase B): one kind and its
// timing, judged by the node-engine beside liveness. Absent: judged alive or not, and nothing more.
Health *Health `json:"health,omitempty"`
}
func (c *Container) Identity() string { return c.ID }
@@ -1129,7 +1141,7 @@ func (c *Container) Kind() Type { return TypeContainer }
func (c *Container) Target() string { return c.Name }
func (c *Container) validate(where string, _ bool) []string {
var problems []string
problems := c.Health.problems(where, true, !c.RunOnce && c.Schedule == "")
if c.Name == "" {
problems = append(problems, where+": a container needs a name")
}
+188
View File
@@ -0,0 +1,188 @@
package declaration
import (
"bytes"
"encoding/json"
"fmt"
"strings"
"time"
)
// Health is how a long-running resource is ready, as the controller composed it from the module's
// `health` (novox/hq ADR 0240 rule 2, to-be 48 §2–§3, Phase B): one kind and its timing, the endpoint
// already the port this machine published it on.
//
// **The node-engine runs every kind and owns every verdict.** http and tcp it makes itself, from the
// machine to the port; unit it reads from the service manager it already reads; exec and runtime it hands
// to the container runtime as the container's own check, with this timing, and reads the state; tool it
// asks of its own node tools. Nothing else on the machine sets a container's check.
//
// Refused here as the controller refuses it near the author, in the same bounds: an engine that took a
// check it could not judge would say a module ready that nothing looked at.
type Health struct {
Kind string `json:"kind"`
// Endpoint is the module's name for what Port is: for the words a verdict is said in.
Endpoint string `json:"endpoint,omitempty"`
Port int `json:"port,omitempty"`
Path string `json:"path,omitempty"`
Status int `json:"status,omitempty"`
Body string `json:"body,omitempty"`
Scheme string `json:"scheme,omitempty"`
Command string `json:"command,omitempty"`
Tool string `json:"tool,omitempty"`
Interval string `json:"interval"`
Timeout string `json:"timeout"`
Looks int `json:"looks"`
Grace string `json:"grace"`
// Needs is the provision the check exercises (to-be 48 §6): said with every verdict, so the
// controller can hold what it finds under the provider's own condition.
Needs string `json:"needs,omitempty"`
}
// UnmarshalJSON reads a health strictly, as everything a declaration carries is read: a field this host
// does not know is a part of the check the controller believes it asked for, and nothing would look at it.
func (h *Health) UnmarshalJSON(raw []byte) error {
type plain Health
var p plain
dec := json.NewDecoder(bytes.NewReader(raw))
dec.DisallowUnknownFields()
if err := dec.Decode(&p); err != nil {
return fmt.Errorf("health: %w", err)
}
*h = Health(p)
return nil
}
// The kinds.
const (
HealthRuntime = "runtime"
HealthHTTP = "http"
HealthTCP = "tcp"
HealthExec = "exec"
HealthUnit = "unit"
HealthTool = "tool"
)
// The bounds (ADR 0240 rule 2) — the controller's, held again here.
const (
HealthIntervalFloor = 10 * time.Second
HealthLooksFloor = 2
HealthWithin = 5 * time.Minute
)
// Every, Within and GraceOf are the timing, read. Validated on arrival, so a parse error here is
// impossible on a declaration that was accepted; it reads as zero.
func (h *Health) Every() time.Duration { d, _ := time.ParseDuration(h.Interval); return d }
func (h *Health) Within() time.Duration { d, _ := time.ParseDuration(h.Timeout); return d }
func (h *Health) GraceOf() time.Duration { d, _ := time.ParseDuration(h.Grace); return d }
// RunByRuntime says the container runtime runs this check as the container's own: exec and runtime.
func (h *Health) RunByRuntime() bool {
return h != nil && (h.Kind == HealthExec || h.Kind == HealthRuntime)
}
// Words is the check in a few words, as a verdict is said: "http /healthz on web".
func (h *Health) Words() string {
switch h.Kind {
case HealthHTTP:
return "http " + h.Path + " on " + orPort(h.Endpoint, h.Port)
case HealthTCP:
return "tcp on " + orPort(h.Endpoint, h.Port)
case HealthTool:
return "its tool " + h.Tool
case HealthRuntime:
return "its image's own check"
case HealthExec:
return "its command"
case HealthUnit:
return "its unit"
}
return h.Kind
}
func orPort(endpoint string, port int) string {
if endpoint != "" {
return endpoint
}
return fmt.Sprint(port)
}
// problems holds a resource's health to its kind and bounds. container says whether the resource is a
// container; longRunning whether it stays up.
func (h *Health) problems(where string, container, longRunning bool) []string {
if h == nil {
return nil
}
var problems []string
say := func(format string, args ...any) {
problems = append(problems, where+": "+fmt.Sprintf(format, args...))
}
if !longRunning {
say("health is judged on what stays up; a step or a scheduled run is judged by its own outcome")
}
switch h.Kind {
case HealthHTTP, HealthTCP:
if h.Port < 1 || h.Port > 65535 {
say("a %s check needs the port it looks at", h.Kind)
}
case HealthExec:
if !container {
say("an exec check runs inside a container")
}
if strings.TrimSpace(h.Command) == "" {
say("an exec check needs a command")
}
case HealthRuntime:
if !container {
say("a runtime check is a container image's own")
}
case HealthUnit:
if container {
say("a unit check is a service's or a process's own")
}
case HealthTool:
if strings.TrimSpace(h.Tool) == "" {
say("a tool check names the tool")
}
default:
say("health of kind %q; it is runtime, http, tcp, exec, unit or tool", h.Kind)
}
if h.Kind == HealthHTTP {
if !strings.HasPrefix(h.Path, "/") {
say("an http check asks a path starting with /")
}
if h.Status != 0 && (h.Status < 100 || h.Status > 599) {
say("an http check expects status %d, which is not one", h.Status)
}
if h.Scheme != "" && h.Scheme != "http" && h.Scheme != "https" {
say("an http check is over http or https, not %q", h.Scheme)
}
}
if strings.ContainsAny(h.Command, "\n\r") {
say("an exec check's command is one line")
}
every, everyErr := time.ParseDuration(h.Interval)
within, withinErr := time.ParseDuration(h.Timeout)
grace, graceErr := time.ParseDuration(h.Grace)
switch {
case everyErr != nil || withinErr != nil || graceErr != nil:
say("health's interval, timeout and grace are durations")
default:
if every < HealthIntervalFloor {
say("a check looks no more often than every %s, not every %s", HealthIntervalFloor, every)
}
if within <= 0 || within >= every {
say("a look takes more than nothing and less than its interval")
}
if grace < 0 {
say("a grace is not negative")
}
if h.Looks >= HealthLooksFloor && grace+time.Duration(h.Looks)*every > HealthWithin {
say("a grace and the failing looks take at most %s", HealthWithin)
}
}
if h.Looks < HealthLooksFloor {
say("a check is unhealthy after at least %d failing looks, not %d", HealthLooksFloor, h.Looks)
}
return problems
}
+49
View File
@@ -0,0 +1,49 @@
package declaration
import (
"strings"
"testing"
)
// A resource's health is held to its kind and bounds on arrival, as the controller holds it near the
// author (novox/hq ADR 0240 rule 2): an engine that took a check it could not judge would say a module ready
// that nothing looked at.
func TestAHealthIsHeldToItsKindAndBounds(t *testing.T) {
pinned := "postgres@sha256:" + strings.Repeat("a", 64)
parse := func(resource string) error {
_, err := Parse([]byte(`{"declaration":1,"resources":[` + resource + `]}`))
return err
}
container := func(health string, more string) string {
return `{"id":"m.db","type":"container","name":"db","image":"` + pinned + `"` + more + `,"health":` + health + `}`
}
timing := `"interval":"30s","timeout":"5s","looks":3,"grace":"1m0s"`
for _, ok := range []string{
container(`{"kind":"http","endpoint":"web","port":31001,"path":"/",`+timing+`}`, ""),
container(`{"kind":"exec","command":"pg_isready",`+timing+`}`, ""),
container(`{"kind":"runtime",`+timing+`}`, ""),
`{"id":"m.d","type":"service","unit":"d.service","state":"running","health":{"kind":"unit",` + timing + `}}`,
} {
if err := parse(ok); err != nil {
t.Errorf("refused %s: %v", ok, err)
}
}
for _, c := range []struct{ resource, says string }{
{container(`{"kind":"http","path":"/",`+timing+`}`, ""), "needs the port"},
{container(`{"kind":"exec",`+timing+`}`, ""), "needs a command"},
{container(`{"kind":"unit",`+timing+`}`, ""), "a service's or a process's own"},
{container(`{"kind":"ping",`+timing+`}`, ""), `of kind "ping"`},
{container(`{"kind":"tcp","port":1,"interval":"5s","timeout":"1s","looks":3,"grace":"0s"}`, ""), "no more often than every 10s"},
{container(`{"kind":"tcp","port":1,"interval":"30s","timeout":"30s","looks":3,"grace":"0s"}`, ""), "less than its interval"},
{container(`{"kind":"tcp","port":1,"interval":"30s","timeout":"5s","looks":1,"grace":"0s"}`, ""), "at least 2 failing looks"},
{container(`{"kind":"tcp","port":1,"interval":"60s","timeout":"5s","looks":3,"grace":"3m"}`, ""), "at most 5m0s"},
{container(`{"kind":"runtime",`+timing+`}`, `,"run-once":true`), "judged on what stays up"},
{`{"id":"m.d","type":"service","unit":"d.service","state":"running","health":{"kind":"runtime",` + timing + `}}`, "a container image's own"},
{container(`{"kind":"tcp","port":1,`+timing+`,"retries":2}`, ""), "retries"},
} {
err := parse(c.resource)
if err == nil || !strings.Contains(err.Error(), c.says) {
t.Errorf("%s: want a refusal saying %q, got %v", c.resource, c.says, err)
}
}
}
+20
View File
@@ -67,6 +67,26 @@ func (b OverNATS) Health(ctx context.Context, node string, body []byte) error {
return b.Conn.FlushWithContext(flush)
}
// ToolSubject is a module's tool on one machine's node tools — the subject a node-engine asks a declared
// tool check on (novox/hq ADR 0240, to-be 48 §3), granted it for its own machine's instance and no other.
func ToolSubject(module, tool, node string) string {
return "mesh.mod." + module + ".tool." + tool + "." + node
}
// ToolBus is a link that can ask a module's tool and wait for its answer.
type ToolBus interface {
Ask(ctx context.Context, subject string, body []byte) ([]byte, error)
}
// Ask asks on core NATS and waits for the answer, within ctx.
func (b OverNATS) Ask(ctx context.Context, subject string, body []byte) ([]byte, error) {
msg, err := b.Conn.RequestWithContext(ctx, subject, body)
if err != nil {
return nil, err
}
return msg.Data, nil
}
func (b OverNATS) Report(ctx context.Context, node string, body []byte) error {
// Into the CONTROL stream and awaited: this is the message the store-window guarantee is
// about (novox/hq ADR 0083). The controller naks with a delay while its store is away and
+20
View File
@@ -50,3 +50,23 @@ func TestAHealthStatementIsSaidOnTheLinkOpenNow(t *testing.T) {
t.Fatalf("the subject %s is not inside the host's own grant", HealthSubject("anchor"))
}
}
// A declared tool check reads the node tools' reply envelope: healthy only when the tool says so (novox/hq
// ADR 0240, to-be 48 §2), and asked on this machine's own instance (to-be 48 §3).
func TestAToolsAnswerIsHealthyOnlyWhenItSaysSo(t *testing.T) {
for reply, want := range map[string]bool{
`{"result":{"healthy":true},"node":"anchor"}`: true,
`{"result":{"healthy":false,"why":"the admin refused the secret"}}`: false,
`{"result":{"status":"fine"}}`: false,
`{"error":"503 no responders"}`: false,
`{"result":"yes"}`: false,
} {
got, _, err := ReadToolAnswer([]byte(reply))
if err != nil || got != want {
t.Errorf("%s read as %v (%v)", reply, got, err)
}
}
if ToolSubject("keycloak", "keycloak_admin_health", "anchor") != "mesh.mod.keycloak.tool.keycloak_admin_health.anchor" {
t.Errorf("asked on %s", ToolSubject("keycloak", "keycloak_admin_health", "anchor"))
}
}
+9
View File
@@ -222,6 +222,11 @@ type Report struct {
// judged with no declaration).
const LivenessContract = 1
// ReadinessContract is the statement of a host that also reads a resource's declared `health` and judges
// it (ADR 0240 Phase B). Its presence is what tells the controller this host may be sent the field: an
// older host parses strictly and refuses the whole declaration for it.
const ReadinessContract = 2
// Health is one statement of every long-running resource's state on this machine (to-be 48 §4). Said in
// every report, as an event on each change, and again every minute while anything is not healthy — so a
// lost statement is not a lost fault.
@@ -270,6 +275,10 @@ type ResourceHealth struct {
// Restarts is how many restarts the engine counted after a grace, kept across recreates and across
// its own restarts: the runtime's own count is lost on every recreate.
Restarts int `json:"restarts,omitempty"`
// Check is the declared check's kind, empty for a resource judged by liveness alone; Needs the
// provision it exercises, for the controller to hold what it finds under the provider (to-be 48 §6).
Check string `json:"check,omitempty"`
Needs string `json:"needs,omitempty"`
}
// HealthSaid is the health event: a machine's statement between its reports, on HealthSubject.
+3 -2
View File
@@ -138,8 +138,9 @@ func TestTheWireFormatIsExactlyTheseFieldNames(t *testing.T) {
{Report{Node: "n", Health: &Health{Contract: LivenessContract}}, []string{"node", "health"}},
{Health{Contract: LivenessContract, Resources: []ResourceHealth{}}, []string{"contract", "at", "resources"}},
{ResourceHealth{Module: "m", Resource: "m.r", Kind: "container", Target: "t", State: StateUnhealthy,
Reason: ReasonRestarting, Streak: 2, Restarts: 3},
[]string{"module", "resource", "kind", "target", "state", "reason", "since", "streak", "restarts"}},
Reason: ReasonRestarting, Streak: 2, Restarts: 3, Check: "http", Needs: "postgres-database"},
[]string{"module", "resource", "kind", "target", "state", "reason", "since", "streak", "restarts",
"check", "needs"}},
{HealthSaid{Node: "n"}, []string{"node", "health"}},
} {
raw, err := json.Marshal(c.value)
+43
View File
@@ -5,6 +5,7 @@ import (
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"sync"
"time"
@@ -511,3 +512,45 @@ func (q *Queue) timeout() time.Duration {
}
return q.Timeout
}
// AskTool asks a module's tool on this machine's own node tools, as a declared health check (novox/hq ADR
// 0240, to-be 48 §3), and reads its answer: healthy or not, with why. An error when there is no link open
// or nothing answered — said as the check's failure, never as healthy.
func (q *Queue) AskTool(ctx context.Context, module, tool string) (bool, string, error) {
q.init()
q.mu.Lock()
bus := q.bus
q.mu.Unlock()
asker, ok := bus.(ToolBus)
if bus == nil || !ok {
return false, "", errors.New("no link to the bus is open")
}
reply, err := asker.Ask(ctx, ToolSubject(module, tool, q.Membership.Node), []byte("{}"))
if err != nil {
return false, "", err
}
return ReadToolAnswer(reply)
}
// ReadToolAnswer reads a tool's reply envelope — `{result, error}`, as the node tools answer every call —
// as a health answer: an error is not healthy, and so is a result that does not say `healthy: true`.
func ReadToolAnswer(reply []byte) (bool, string, error) {
var env struct {
Result json.RawMessage `json:"result"`
Error string `json:"error"`
}
if err := json.Unmarshal(reply, &env); err != nil {
return false, "", fmt.Errorf("its answer is not one: %w", err)
}
if env.Error != "" {
return false, env.Error, nil
}
var a struct {
Healthy bool `json:"healthy"`
Why string `json:"why"`
}
if err := json.Unmarshal(env.Result, &a); err != nil {
return false, "it answered something other than healthy or not", nil
}
return a.Healthy, a.Why, nil
}
+65 -9
View File
@@ -81,6 +81,9 @@ type Resource struct {
// Scope and User are a service's manager: "user" and the account for a unit in an account's own.
Scope string `json:"scope,omitempty"`
User string `json:"user,omitempty"`
// Check is how it is ready, as its module declared (ADR 0240 rule 2, Phase B); nil judges it alive
// or not, and nothing more.
Check *declaration.Health `json:"check,omitempty"`
}
// LongRunning is every long-running resource a declaration asks this machine to run for a module: its
@@ -102,12 +105,12 @@ func LongRunning(d *declaration.Declaration, held map[string]bool) []Resource {
if v.RunOnce || v.Schedule != "" {
continue
}
out = append(out, Resource{Module: module, ID: v.ID, Kind: KindContainer, Target: v.Name})
out = append(out, Resource{Module: module, ID: v.ID, Kind: KindContainer, Target: v.Name, Check: v.Health})
case *declaration.Service:
if v.State != "running" {
continue
}
res := Resource{Module: module, ID: v.ID, Kind: KindService, Target: v.Unit}
res := Resource{Module: module, ID: v.ID, Kind: KindService, Target: v.Unit, Check: v.Health}
if v.UserScoped() {
res.Scope, res.User = declaration.ScopeUser, v.User
}
@@ -116,7 +119,8 @@ func LongRunning(d *declaration.Declaration, held map[string]bool) []Resource {
if v.RunOnce || v.Schedule != "" {
continue
}
out = append(out, Resource{Module: module, ID: v.ID, Kind: KindProcess, Target: v.Name + ".service"})
out = append(out, Resource{Module: module, ID: v.ID, Kind: KindProcess, Target: v.Name + ".service",
Check: v.Health})
}
}
return out
@@ -150,6 +154,10 @@ type Observed struct {
// Started is when the runtime says the current run started, as it says it: a change with no restart
// counted is a restart somebody made, which is a new start.
Started string
// Health is the runtime's own word on the check it runs as the container's — healthy, unhealthy,
// starting — empty when the container carries none; HealthSaid what its last look printed.
Health string
HealthSaid string
}
// Runtime is what one look reads: every container's state in one read, and every unit's per manager.
@@ -169,6 +177,21 @@ type State struct {
Restarts int `json:"restarts,omitempty"`
}
// CheckOf and NeedsOf are a stated resource's declared check, in a word, and the provision it exercises.
func (s State) CheckOf() string {
if s.Check == nil {
return ""
}
return s.Check.Kind
}
func (s State) NeedsOf() string {
if s.Check == nil {
return ""
}
return s.Check.Needs
}
// Statement is one look at every long-running resource: when, and each resource's state.
type Statement struct {
At time.Time
@@ -206,6 +229,15 @@ type kept struct {
Reason string `json:"reason,omitempty"`
Since time.Time `json:"since"`
Streak int `json:"streak,omitempty"`
// Readiness since the current start (Phase B): whether the declared check has passed, how many of
// its looks after the grace failed in a row, and what the last failing one said.
Passed bool `json:"passed,omitempty"`
Failing int `json:"failing,omitempty"`
Why string `json:"why,omitempty"`
// probing says a look of the engine's own is under way; due when the next is.
probing bool
due time.Time
}
// file is the judge's file beside the node's state.
@@ -233,13 +265,20 @@ type Judge struct {
dirty bool
// said is the statement said last, to know a change by.
said map[string]string
// Probes make the looks the engine makes itself — http, tcp and a module's tool (Phase B). Nil
// makes them from this machine.
Probes *Probes
// Budget is the most looks a minute every declared check together may cost (ADR 0240: never more
// than measured on the busiest machine); over it the engine's own looks are spaced out.
Budget int
}
// Open is the judge whose file is at path, reading what it kept. A file that cannot be read is said and
// started afresh: a count lost is a crash loop judged from now, never a machine left unjudged.
func Open(path string, rt Runtime) (*Judge, error) {
j := &Judge{path: path, runtime: rt, Now: time.Now, Grace: DefaultGrace, Settle: SettleWindow,
f: file{Kept: map[string]*kept{}}, said: map[string]string{}}
f: file{Kept: map[string]*kept{}}, said: map[string]string{}, Budget: Budget}
raw, err := os.ReadFile(path)
switch {
case os.IsNotExist(err):
@@ -277,6 +316,10 @@ func (j *Judge) Set(resources []Resource) {
// The same id now names another thing: judged as a new one, its count kept.
counted := k.Counted
j.f.Kept[r.ID] = &kept{Resource: r, Counted: counted}
} else if ok && !sameCheck(k.Check, r.Check) {
// Its check changed — another kind, another port, another timing: what the old one found
// says nothing about the new, which looks again at once.
k.Check, k.Passed, k.Failing, k.Why, k.due = r.Check, false, 0, "", time.Time{}
}
}
for id := range j.f.Kept {
@@ -314,7 +357,7 @@ func (j *Judge) Look(ctx context.Context) (Statement, bool) {
o := observed[keyOf(r)]
j.judge(k, o, now, r.Kind == KindContainer && held[r.Target], blind, why)
seen[r.ID] = true
st.Resources = append(st.Resources, State{Resource: r, State: k.State, Reason: k.Reason, Since: k.Since,
st.Resources = append(st.Resources, State{Resource: k.Resource, State: k.State, Reason: k.Reason, Since: k.Since,
Streak: k.Streak, Restarts: k.Counted})
word := k.State + "/" + k.Reason
if j.said[r.ID] != word {
@@ -348,6 +391,8 @@ func (j *Judge) judge(k *kept, o Observed, now time.Time, held, blind bool, why
fresh := func(o Observed, started time.Time) {
k.Seen, k.Identity, k.RuntimeRestarts, k.RuntimeStarted = o.Found, o.Identity, o.Restarts, o.Started
k.Started, k.Recent = started, nil
// Every start is judged ready afresh (to-be 48 §4).
k.Passed, k.Failing, k.Why, k.due = false, 0, "", time.Time{}
}
// **Held is neither alive nor dead**, and the window ending is a start: judged from a fresh grace.
@@ -407,23 +452,34 @@ func (j *Judge) judge(k *kept, o Observed, now time.Time, held, blind bool, why
}
k.Recent = recent
inGrace := now.Before(k.Started.Add(j.Grace))
inGrace := now.Before(k.Started.Add(j.graceOf(k.Resource)))
switch {
case len(k.Recent) >= 2:
set(Unhealthy, ReasonRestarting)
case inGrace:
set(Starting, "")
case !o.Found || !o.Running:
case !inGrace && (!o.Found || !o.Running):
if o.Restarting {
set(Unhealthy, ReasonRestarting)
} else {
set(Unhealthy, ReasonDown)
}
case k.Check != nil:
// Alive, or still in its grace: how ready it is is the declared check's to say.
set(ready(k, o, inGrace))
case inGrace:
set(Starting, "")
default:
set(Healthy, "")
}
}
// graceOf is a resource's grace: its declared one, or the default (to-be 48 §1).
func (j *Judge) graceOf(r Resource) time.Duration {
if r.Check != nil {
return r.Check.GraceOf()
}
return j.Grace
}
// read is one read of everything: every container at once, every unit per manager. unread names each
// group that could not be read, with why.
func (j *Judge) read(ctx context.Context) (map[string]Observed, map[string]string) {
+334
View File
@@ -0,0 +1,334 @@
package liveness
import (
"context"
"crypto/tls"
"errors"
"fmt"
"io"
"net"
"net/http"
"net/url"
"os"
"strconv"
"strings"
"time"
"github.com/novox/mesh-host/internal/declaration"
)
// Readiness, declared (novox/hq ADR 0240 rules 2 and 3, to-be 48 §2–§3, Phase B).
//
// A long-running resource that declares `health` is judged alive as before, and then ready by its check:
// **starting** while in its grace and not yet passed, **healthy** once it passed, **unhealthy** once its
// declared number of looks after the grace failed in a row — with what the check found, in words.
//
// **The engine runs every kind where it is cheapest, and owns every verdict:**
//
// - http and tcp it makes itself, from the machine to the port the endpoint is published on — the path
// a caller takes (issue 145), and no execution inside the container;
// - unit it reads from the show of the service manager every look already makes;
// - exec and runtime the runtime runs as the container's own check, with the declared timing (its
// retries are the failing looks, its start period the grace), and the engine reads the state from the
// inspect every look already makes;
// - tool it asks of the module's tool on this machine's own node tools.
//
// **Never more looks than the measured budget.** The busiest machine runs about ninety looks a minute at
// the default interval (research 032 §7); when the declared checks together would cost more, the engine
// spaces its own looks out until they do not, and says so.
// Budget is the most looks a minute the declared checks on one machine may cost together: one check per
// long-running resource at the default interval on the busiest machine (research 032 §7).
const Budget = 90
// ProbesAtOnce is how many of the engine's own looks run at the same time.
const ProbesAtOnce = 8
// ready is a resource's readiness, on a look where it is alive or still in its grace: starting until its
// check passed, healthy once it has, unhealthy once its failing looks after the grace reach the declared
// number. The runtime's and the unit's checks are read from what the look itself read; the engine's own
// looks were recorded by Probe.
func ready(k *kept, o Observed, inGrace bool) (string, string) {
c := k.Check
switch c.Kind {
case declaration.HealthRuntime, declaration.HealthExec:
switch o.Health {
case Healthy:
k.Passed, k.Failing, k.Why = true, 0, ""
case Unhealthy:
// The runtime's retries are the declared failing looks: unhealthy is already that many.
if !inGrace {
k.Failing, k.Why = c.Looks, orSaid(o.HealthSaid, "its check failed")
}
case "":
if !inGrace {
why := "the container carries no check: it was made before its check was declared"
if c.Kind == declaration.HealthRuntime {
why = "its image ships no check to adopt"
}
k.Failing, k.Why = c.Looks, why
}
}
case declaration.HealthUnit:
// Active and not failed — and for a unit that notifies, notified, which the manager says by
// being active only once it was.
if o.Running {
k.Passed, k.Failing, k.Why = true, 0, ""
}
}
switch {
case !inGrace && k.Failing >= c.Looks:
return Unhealthy, c.Words() + ": " + orSaid(k.Why, "it failed")
case k.Passed:
return Healthy, ""
case k.Why != "":
return Starting, "not ready yet — " + c.Words() + ": " + k.Why
}
return Starting, ""
}
func orSaid(s, otherwise string) string {
if strings.TrimSpace(s) == "" {
return otherwise
}
return s
}
// sameCheck says two declarations of a check are the same check.
func sameCheck(a, b *declaration.Health) bool {
if a == nil || b == nil {
return a == b
}
return *a == *b
}
// ownLook says whether a check is one the engine makes itself.
func ownLook(c *declaration.Health) bool {
return c != nil && (c.Kind == declaration.HealthHTTP || c.Kind == declaration.HealthTCP || c.Kind == declaration.HealthTool)
}
// Spacing is how much the engine's own looks are spaced out so every declared check together costs no
// more than the budget: 1 when they fit, more when they would not.
func (j *Judge) Spacing() float64 {
j.mu.Lock()
defer j.mu.Unlock()
return j.spacing()
}
func (j *Judge) spacing() float64 {
perMinute := 0.0
for _, r := range j.f.Resources {
if r.Check == nil || r.Check.Kind == declaration.HealthUnit {
continue // a unit's readiness rides on the show every look already makes
}
if every := r.Check.Every(); every > 0 {
perMinute += float64(time.Minute) / float64(every)
}
}
budget := j.Budget
if budget <= 0 {
budget = Budget
}
if perMinute <= float64(budget) {
return 1
}
return perMinute / float64(budget)
}
// Probe makes the engine's own looks — http, tcp, a module's tool — each at its interval, spaced out to
// the budget, never two of one resource at once and never more than ProbesAtOnce together, until ctx
// ends. A held resource is not looked at. It reads; it never acts (ADR 0240 rule 6).
func (j *Judge) Probe(ctx context.Context) {
tick := time.NewTicker(time.Second)
defer tick.Stop()
slots := make(chan struct{}, ProbesAtOnce)
for {
select {
case <-ctx.Done():
return
case <-tick.C:
}
for _, d := range j.due() {
select {
case slots <- struct{}{}:
case <-ctx.Done():
return
}
go func(d dueLook) {
defer func() { <-slots }()
ok, why := j.probes().Look(ctx, d.module, d.check)
j.Record(d.id, d.started, ok, why)
}(d)
}
}
}
// dueLook is one of the engine's own looks to make now.
type dueLook struct {
id, module string
check *declaration.Health
started time.Time
}
// due is every look of the engine's own whose time has come, marked under way.
func (j *Judge) due() []dueLook {
j.mu.Lock()
defer j.mu.Unlock()
now := j.Now()
spacing := j.spacing()
var out []dueLook
for _, r := range j.f.Resources {
k := j.f.Kept[r.ID]
if k == nil || !ownLook(r.Check) || k.probing || k.State == Held || k.Started.IsZero() || now.Before(k.due) {
continue
}
k.probing = true
k.due = now.Add(time.Duration(float64(r.Check.Every()) * spacing))
out = append(out, dueLook{id: r.ID, module: r.Module, check: r.Check, started: k.Started})
}
return out
}
// Record keeps what one of the engine's own looks found, for the next look to fold in: a pass is a pass
// whenever it came; a failure counts only after the grace. A look begun before the resource's current
// start says nothing about it.
func (j *Judge) Record(id string, started time.Time, ok bool, why string) {
j.mu.Lock()
defer j.mu.Unlock()
k := j.f.Kept[id]
if k == nil {
return
}
k.probing = false
if !k.Started.Equal(started) || k.Check == nil {
return
}
switch {
case ok:
k.Passed, k.Failing, k.Why = true, 0, ""
case j.Now().Before(k.Started.Add(j.graceOf(k.Resource))):
k.Why = why
default:
k.Failing++
k.Why = why
}
j.dirty = true
}
func (j *Judge) probes() *Probes {
if j.Probes != nil {
return j.Probes
}
return &Probes{}
}
// Probes make the looks the engine makes itself (to-be 48 §3).
type Probes struct {
// Host is where an endpoint's port is dialled: this machine, by default its loopback.
Host string
// AskTool asks a module's tool on this machine's own node tools; nil says it cannot be asked.
AskTool func(ctx context.Context, module, tool string) (healthy bool, why string, err error)
}
// Look makes one look of a check, and says whether it passed and, when not, what it found.
func (p *Probes) Look(ctx context.Context, module string, c *declaration.Health) (bool, string) {
timeout := c.Within()
if timeout <= 0 {
timeout = 5 * time.Second
}
ctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
switch c.Kind {
case declaration.HealthHTTP:
return p.http(ctx, c, timeout)
case declaration.HealthTCP:
conn, err := (&net.Dialer{}).DialContext(ctx, "tcp", p.address(c.Port))
if err != nil {
return false, dialWords(err, timeout)
}
_ = conn.Close()
return true, ""
case declaration.HealthTool:
if p.AskTool == nil {
return false, "its tool cannot be asked: this engine has no link to the node tools"
}
healthy, why, err := p.AskTool(ctx, module, c.Tool)
switch {
case err != nil:
return false, "asking it: " + firstLine(err.Error())
case !healthy:
return false, orSaid(why, "it answered not healthy, and not why")
}
return true, ""
}
return false, "the engine does not make a " + c.Kind + " look itself"
}
func (p *Probes) address(port int) string {
host := p.Host
if host == "" {
host = "127.0.0.1"
}
return net.JoinHostPort(host, strconv.Itoa(port))
}
// http is one request: the path, on the port, expecting the declared status (any under 400 when none
// was declared) and, when declared, a text in the answer. A redirect is an answer, not followed: a login
// page elsewhere says nothing about this program.
func (p *Probes) http(ctx context.Context, c *declaration.Health, timeout time.Duration) (bool, string) {
scheme := c.Scheme
if scheme == "" {
scheme = "http"
}
url := scheme + "://" + p.address(c.Port) + c.Path
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
if err != nil {
return false, err.Error()
}
req.Header.Set("User-Agent", "mesh-node-engine health")
client := &http.Client{
CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse },
// Whether it answers, not whom to trust: a program's own certificate is its own business.
Transport: &http.Transport{TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, DisableKeepAlives: true},
}
res, err := client.Do(req)
if err != nil {
return false, dialWords(err, timeout)
}
defer res.Body.Close()
switch {
case c.Status != 0 && res.StatusCode != c.Status:
return false, fmt.Sprintf("answered %d, expected %d", res.StatusCode, c.Status)
case c.Status == 0 && res.StatusCode >= 400:
return false, fmt.Sprintf("answered %d", res.StatusCode)
}
if c.Body != "" {
body, err := io.ReadAll(io.LimitReader(res.Body, 64<<10))
if err != nil {
return false, "its answer could not be read: " + firstLine(err.Error())
}
if !strings.Contains(string(body), c.Body) {
return false, fmt.Sprintf("answered %d without %q", res.StatusCode, c.Body)
}
}
return true, ""
}
// dialWords is a failed look in words: no answer in time, refused, or what the system said.
func dialWords(err error, timeout time.Duration) string {
var asked *url.Error
if errors.As(err, &asked) {
err = asked.Err // what happened, without the address it happened at
}
switch {
case errors.Is(err, io.EOF) || errors.Is(err, io.ErrUnexpectedEOF):
return "it closed the connection without answering"
case errors.Is(err, context.DeadlineExceeded) || errors.Is(err, os.ErrDeadlineExceeded):
return fmt.Sprintf("no answer within %s", timeout)
case strings.Contains(err.Error(), "connection refused"):
return "connection refused"
case strings.Contains(err.Error(), "Client.Timeout") || strings.Contains(err.Error(), "timeout"):
return fmt.Sprintf("no answer within %s", timeout)
}
return firstLine(err.Error())
}
+231
View File
@@ -0,0 +1,231 @@
package liveness
import (
"context"
"net"
"net/http"
"net/http/httptest"
"path/filepath"
"strconv"
"strings"
"testing"
"time"
"github.com/novox/mesh-host/internal/declaration"
)
// Readiness, declared (novox/hq ADR 0240 rules 2 and 3, Phase B): starting until its check passed, healthy
// once it has, unhealthy once its failing looks after the grace reach the declared number; an http check
// dials the endpoint's current port after a port change; the runtime's check is read, never run by the
// engine; and never more looks than the budget.
func httpCheck(port int) *declaration.Health {
return &declaration.Health{Kind: declaration.HealthHTTP, Endpoint: "web", Port: port, Path: "/healthz",
Interval: "10s", Timeout: "2s", Looks: 2, Grace: "30s"}
}
func portOf(t *testing.T, url string) int {
t.Helper()
_, p, err := net.SplitHostPort(strings.TrimPrefix(url, "http://"))
if err != nil {
t.Fatal(err)
}
n, _ := strconv.Atoi(p)
return n
}
// lookOnce makes the engine's own due looks at once and waits for them, as Probe would on its tick.
func lookOnce(t *testing.T, j *Judge) {
t.Helper()
for _, d := range j.due() {
ok, why := j.probes().Look(t.Context(), d.module, d.check)
j.Record(d.id, d.started, ok, why)
}
}
func TestAnHTTPCheckIsStartingUntilItPassesAndUnhealthyAfterItsLooksFail(t *testing.T) {
answering := true
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if !answering || r.URL.Path != "/healthz" {
w.WriteHeader(http.StatusServiceUnavailable)
return
}
w.Write([]byte("ok"))
}))
defer srv.Close()
c := &clock{now: t0}
rt := &fakeRuntime{containers: map[string]Observed{"web": running("c1", 0, t0.Format(time.RFC3339Nano))}}
j := aJudge(t, rt, c, filepath.Join(t.TempDir(), FileName))
web := Resource{Module: "app", ID: "app.web", Kind: KindContainer, Target: "web", Check: httpCheck(portOf(t, srv.URL))}
j.Set([]Resource{web})
// In its grace, before any look: starting.
st, _ := j.Look(t.Context())
if s := stateOf(t, st, "app.web"); s.State != Starting || s.CheckOf() != "http" {
t.Fatalf("before any look: %+v", s)
}
// It answers: healthy, even inside its grace.
lookOnce(t, j)
st, _ = j.Look(t.Context())
if s := stateOf(t, st, "app.web"); s.State != Healthy {
t.Fatalf("after a passing look: %+v", s)
}
// It stops answering, after its grace: one failing look is not yet unhealthy; the second is.
answering = false
c.now = t0.Add(time.Minute)
lookOnce(t, j)
st, _ = j.Look(t.Context())
if s := stateOf(t, st, "app.web"); s.State != Healthy {
t.Fatalf("one failing look made it %s", s.State)
}
c.now = c.now.Add(10 * time.Second)
lookOnce(t, j)
st, _ = j.Look(t.Context())
s := stateOf(t, st, "app.web")
if s.State != Unhealthy || !strings.Contains(s.Reason, "http /healthz on web: answered 503") {
t.Fatalf("two failing looks: %+v", s)
}
// And nothing was restarted, recreated or stopped: the fake runtime was only read.
answering = true
c.now = c.now.Add(10 * time.Second)
lookOnce(t, j)
st, _ = j.Look(t.Context())
if s := stateOf(t, st, "app.web"); s.State != Healthy {
t.Fatalf("answering again: %+v", s)
}
}
func TestAnHTTPCheckDialsTheEndpointsCurrentPortAfterAPortChange(t *testing.T) {
hit := map[string]int{}
handler := func(name string) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { hit[name]++ })
}
before := httptest.NewServer(handler("before"))
defer before.Close()
after := httptest.NewServer(handler("after"))
defer after.Close()
c := &clock{now: t0}
rt := &fakeRuntime{containers: map[string]Observed{"web": running("c1", 0, t0.Format(time.RFC3339Nano))}}
j := aJudge(t, rt, c, "")
j.Set([]Resource{{Module: "app", ID: "app.web", Kind: KindContainer, Target: "web", Check: httpCheck(portOf(t, before.URL))}})
j.Look(t.Context())
lookOnce(t, j)
// The machine gave the endpoint another port: the next declaration says so, and the check follows.
j.Set([]Resource{{Module: "app", ID: "app.web", Kind: KindContainer, Target: "web", Check: httpCheck(portOf(t, after.URL))}})
j.Look(t.Context())
lookOnce(t, j)
if hit["before"] != 1 || hit["after"] != 1 {
t.Fatalf("the check dialled %v; after the port moved it must dial the new one at once", hit)
}
}
func TestTheRuntimesCheckIsReadAndAnImageWithoutOneIsSaid(t *testing.T) {
c := &clock{now: t0}
o := running("c1", 0, t0.Format(time.RFC3339Nano))
rt := &fakeRuntime{containers: map[string]Observed{"db": o}}
j := aJudge(t, rt, c, "")
check := &declaration.Health{Kind: declaration.HealthRuntime, Interval: "30s", Timeout: "5s", Looks: 3, Grace: "1m0s"}
j.Set([]Resource{{Module: "app", ID: "app.db", Kind: KindContainer, Target: "db", Check: check}})
look := func(health, said string) State {
o.Health, o.HealthSaid = health, said
rt.containers["db"] = o
st, _ := j.Look(t.Context())
return stateOf(t, st, "app.db")
}
if s := look("starting", ""); s.State != Starting {
t.Fatalf("the runtime says starting: %+v", s)
}
if s := look(Healthy, ""); s.State != Healthy {
t.Fatalf("the runtime says healthy: %+v", s)
}
c.now = t0.Add(2 * time.Minute)
if s := look(Unhealthy, "curl: (7) Failed to connect to localhost port 3000"); s.State != Unhealthy ||
!strings.Contains(s.Reason, "Failed to connect to localhost") {
t.Fatalf("the runtime says unhealthy: %+v", s)
}
if s := look("", ""); s.State != Unhealthy || !strings.Contains(s.Reason, "ships no check to adopt") {
t.Fatalf("an image with no check adopted by name: %+v", s)
}
}
func TestAUnitIsReadyWhenItsManagerSaysItIsActive(t *testing.T) {
c := &clock{now: t0}
rt := &fakeRuntime{units: map[string]Observed{"d.service": {Found: true, Identity: "i1", Running: true}}}
j := aJudge(t, rt, c, "")
check := &declaration.Health{Kind: declaration.HealthUnit, Interval: "30s", Timeout: "5s", Looks: 2, Grace: "0s"}
j.Set([]Resource{{Module: "app", ID: "app.d", Kind: KindService, Target: "d.service", Check: check}})
st, _ := j.Look(t.Context())
if s := stateOf(t, st, "app.d"); s.State != Healthy {
t.Fatalf("an active unit with no grace: %+v", s)
}
}
func TestNeverMoreLooksThanTheBudget(t *testing.T) {
j := aJudge(t, &fakeRuntime{}, &clock{now: t0}, "")
var rs []Resource
for i := 0; i < 30; i++ {
rs = append(rs, Resource{Module: "app", ID: "app.c" + strconv.Itoa(i), Kind: KindContainer, Target: "c" + strconv.Itoa(i),
Check: &declaration.Health{Kind: declaration.HealthTCP, Port: 1, Interval: "10s", Timeout: "1s", Looks: 2, Grace: "0s"}})
}
j.Set(rs)
// 30 checks every 10 s are 180 looks a minute: spaced twice their interval to stay within 90.
if got := j.Spacing(); got != 2 {
t.Fatalf("spaced %v", got)
}
j.Set(rs[:9])
if got := j.Spacing(); got != 1 {
t.Fatalf("54 looks a minute spaced %v", got)
}
}
func TestAToolCheckIsAskedOfTheNodeTools(t *testing.T) {
asked := ""
p := &Probes{AskTool: func(_ context.Context, module, tool string) (bool, string, error) {
asked = module + "." + tool
return false, "the admin refused the minted secret", nil
}}
ok, why := p.Look(t.Context(), "keycloak", &declaration.Health{Kind: declaration.HealthTool, Tool: "keycloak_admin_health",
Interval: "30s", Timeout: "5s", Looks: 2, Grace: "0s"})
if ok || asked != "keycloak.keycloak_admin_health" || !strings.Contains(why, "refused the minted secret") {
t.Fatalf("asked %q: %v %q", asked, 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) {
status, said := runtimeHealth(`{"Status":"running","Health":{"Status":"unhealthy","FailingStreak":3,"Log":[` +
`{"ExitCode":1,"Output":"curl: (7) Failed to connect\n"}]}}`)
if status != Unhealthy || said != "curl: (7) Failed to connect" {
t.Errorf("read %q %q", status, said)
}
if status, _ := runtimeHealth(`{"Status":"running","Running":true}`); status != "" {
t.Errorf("a container without a check read as %q", status)
}
if !strings.Contains(inspectFormat, "{{json .State}}") || strings.Contains(inspectFormat, ".State.Health") {
t.Errorf("the inspect names a key a container without a check does not have: %s", inspectFormat)
}
}
// The engine's own looks run on their own clock, beside the looks of liveness, and are folded in.
func TestProbeLooksOnItsOwnAndTheNextLookFoldsItIn(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}))
defer srv.Close()
rt := &fakeRuntime{containers: map[string]Observed{"web": running("c1", 0, time.Now().Format(time.RFC3339Nano))}}
j, err := Open("", rt)
if err != nil {
t.Fatal(err)
}
j.Set([]Resource{{Module: "app", ID: "app.web", Kind: KindContainer, Target: "web", Check: httpCheck(portOf(t, srv.URL))}})
ctx, cancel := context.WithCancel(t.Context())
defer cancel()
go j.Probe(ctx)
deadline := time.Now().Add(5 * time.Second)
for time.Now().Before(deadline) {
if st, _ := j.Look(t.Context()); stateOf(t, st, "app.web").State == Healthy {
return
}
time.Sleep(100 * time.Millisecond)
}
t.Fatal("the engine's own look was never folded in")
}
+113
View File
@@ -7,10 +7,12 @@ import (
"fmt"
"os"
"os/exec"
"strconv"
"strings"
"testing"
"time"
"github.com/novox/mesh-host/internal/declaration"
"github.com/novox/mesh-host/internal/link"
)
@@ -89,3 +91,114 @@ func TestReplayCrashLoopIsSaidUnhealthy(t *testing.T) {
}
}
}
// SilentWebImage and SilentWebProgram are the silent web application: a web server whose application
// never answers — it accepts every request and holds it, as the application waiting on a database it could
// not reach did. The lab raises the same (mesh-lab replays/silentweb_test.go).
const (
SilentWebImage = "busybox:1.36"
SilentWebProgram = "mkdir -p /www/cgi-bin && printf '#!/bin/sh\\nsleep 3600\\n' > /www/cgi-bin/app && " +
"chmod +x /www/cgi-bin/app && exec httpd -f -p 8080 -h /www"
)
// **R145, the engine's half — a web application that accepts TCP and answers nothing is said unhealthy by
// its HTTP check within two looks** (novox/hq ADR 0240 Phase B, issue 145). The application's port was
// open and its program ran while every request hung, for eleven hours. Liveness says it alive and a TCP
// check says it reachable; only the HTTP check its module declares sees it.
//
// A container that accepts every connection on its port and never answers is raised here as the
// node-engine raises a module's, its port published on the machine; the judge looks at it with an HTTP check
// of two looks and a TCP check beside it. MESH_REPLAY_STATEMENT, when set, is where the statement is written
// for the controller's half (mesh-controller TestReplaySilentWebAppIsRaisedWithinTwoLooks). The container is
// removed after, whatever happened.
func TestReplaySilentWebAppIsSaidUnhealthy(t *testing.T) {
if os.Getenv("MESH_REPLAY_RUNTIME") != "1" && os.Getenv("MESH_REPLAY_CONTAINER") == "" {
t.Skip("no MESH_REPLAY_CONTAINER, and MESH_REPLAY_RUNTIME is not 1: this replay raises a container")
}
run := func(ctx context.Context, cmd string, args ...string) (string, error) {
c := exec.CommandContext(ctx, cmd, args...)
out, err := c.Output()
var exit *exec.ExitError
if errors.As(err, &exit) {
return string(out), fmt.Errorf("%s exited %d: %s", cmd, exit.ExitCode(), strings.TrimSpace(string(exit.Stderr)))
}
return string(out), err
}
if _, err := run(t.Context(), "docker", "version", "--format", "{{.Server.Version}}"); err != nil {
t.Skipf("no container runtime answers here: %v", err)
}
// The container is raised by the lab (MESH_REPLAY_CONTAINER, reached at MESH_REPLAY_ADDRESS on its own
// port), or here, published on this machine's loopback.
name, host, published := os.Getenv("MESH_REPLAY_CONTAINER"), os.Getenv("MESH_REPLAY_ADDRESS"), 8080
if name == "" {
name = fmt.Sprintf("mesh-replay-silent-web-%d", time.Now().UnixNano())
if out, err := run(t.Context(), "docker", "run", "--detach", "--name", name, "--restart", "unless-stopped",
"--publish", "127.0.0.1::8080", "--label", "mesh-host.id=app.server", "--label", "mesh.replay=1", SilentWebImage,
"sh", "-c", SilentWebProgram); err != nil {
t.Fatalf("starting the silent web application: %v %s", err, out)
}
t.Cleanup(func() { _, _ = run(context.Background(), "docker", "rm", "-f", name) })
said, err := run(t.Context(), "docker", "port", name, "8080/tcp")
if err != nil {
t.Fatal(err)
}
_, port, _ := strings.Cut(strings.TrimSpace(strings.Split(said, "\n")[0]), "127.0.0.1:")
if published, err = strconv.Atoi(port); err != nil {
t.Fatalf("the port it was published on: %q", said)
}
host = "127.0.0.1"
}
timing := func(h *declaration.Health) *declaration.Health {
h.Interval, h.Timeout, h.Looks, h.Grace = "10s", "2s", 2, "0s"
return h
}
j, err := Open("", &Exec{Run: run})
if err != nil {
t.Fatal(err)
}
j.Set([]Resource{
{Module: "app", ID: "app.server", Kind: KindContainer, Target: name,
Check: timing(&declaration.Health{Kind: declaration.HealthHTTP, Endpoint: "web", Port: published, Path: "/cgi-bin/app"})},
})
j.Probes = &Probes{Host: host}
tcp := j.Probes
began := time.Now()
var st Statement
looks := 0
for time.Since(began) < 2*time.Minute {
for _, d := range j.due() {
ok, why := j.probes().Look(t.Context(), d.module, d.check)
j.Record(d.id, d.started, ok, why)
looks++
}
st, _ = j.Look(t.Context())
if st.Resources[0].State == Unhealthy {
break
}
time.Sleep(time.Second)
}
s := st.Resources[0]
if s.State != Unhealthy || looks > 2 {
t.Fatalf("a web application answering nothing was not said unhealthy within two looks (%d looks): %+v", looks, s)
}
// And a TCP check would have said it reachable: the port is open.
if ok, why := tcp.Look(t.Context(), "app", timing(&declaration.Health{Kind: declaration.HealthTCP, Port: published})); !ok {
t.Fatalf("a tcp check could not connect to the silent application: %s", why)
}
t.Logf("said %s (%s) after %d looks, %s after the judging began; a TCP check connects", s.State, s.Reason, looks,
time.Since(began).Round(time.Second))
if out := os.Getenv("MESH_REPLAY_STATEMENT"); out != "" {
h := link.Health{Contract: link.ReadinessContract, At: st.At.UTC(), Resources: []link.ResourceHealth{{
Module: s.Module, Resource: s.ID, Kind: s.Kind, Target: s.Target, State: s.State, Reason: s.Reason,
Since: s.Since.UTC(), Streak: s.Streak, Restarts: s.Restarts, Check: s.CheckOf(), Needs: s.NeedsOf()}}}
raw, err := json.Marshal(h)
if err != nil {
t.Fatal(err)
}
if err := os.WriteFile(out, raw, 0o644); err != nil {
t.Fatal(err)
}
}
}
+42 -3
View File
@@ -2,6 +2,7 @@ package liveness
import (
"context"
"encoding/json"
"errors"
"fmt"
"strconv"
@@ -52,8 +53,11 @@ func (e *Exec) runtime(ctx context.Context) (string, error) {
}
// inspectFormat is what one inspect says of each container: its name, id, status, the runtime's restart
// count and when its current run started.
const inspectFormat = "{{.Name}}\t{{.Id}}\t{{.State.Status}}\t{{.RestartCount}}\t{{.State.StartedAt}}"
// count, when its current run started — and the state of the check the runtime runs as the container's
// own, read out of the whole state as one line of JSON, for a check declared exec or runtime (Phase B).
// The whole state and not its Health: a container that carries no check has no such key, and the runtime
// refuses a template naming a key that is not there — for every container in the read.
const inspectFormat = "{{.Name}}\t{{.Id}}\t{{.State.Status}}\t{{.RestartCount}}\t{{.State.StartedAt}}\t{{json .State}}"
// Containers is one inspect of every container named. A name the runtime does not have is not found;
// a runtime that does not answer is an error — unknown, never down.
@@ -83,12 +87,47 @@ func (e *Exec) Containers(ctx context.Context, names []string) (map[string]Obser
name := strings.TrimPrefix(strings.TrimSpace(parts[0]), "/")
restarts, _ := strconv.ParseInt(strings.TrimSpace(parts[3]), 10, 64)
status := strings.TrimSpace(parts[2])
out[name] = Observed{Found: true, Identity: strings.TrimSpace(parts[1]), Running: status == "running",
o := Observed{Found: true, Identity: strings.TrimSpace(parts[1]), Running: status == "running",
Restarting: status == "restarting", Restarts: restarts, Started: strings.TrimSpace(parts[4])}
if len(parts) > 5 {
o.Health, o.HealthSaid = runtimeHealth(strings.Join(parts[5:], "\t"))
}
out[name] = o
}
return out, nil
}
// runtimeHealth reads, from a container's state, the runtime's state of its own check: its status and what its last look
// printed, in a line. Nothing when the container carries no check.
func runtimeHealth(said string) (string, string) {
var state struct {
Health *struct {
Status string
Log []struct {
ExitCode int
Output string
}
}
}
if err := json.Unmarshal([]byte(strings.TrimSpace(said)), &state); err != nil || state.Health == nil ||
state.Health.Status == "" {
return "", ""
}
h := state.Health
last := ""
if n := len(h.Log); n > 0 {
l := h.Log[n-1]
last = firstLine(l.Output)
if last == "" && l.ExitCode != 0 {
last = fmt.Sprintf("its check exited %d", l.ExitCode)
}
if len(last) > 200 {
last = last[:200] + "…"
}
}
return strings.ToLower(h.Status), last
}
// missing is an inspect that failed only because a name is not there: what it printed for the others is
// still the answer.
func missing(err error) bool {