diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index e34ced0..e716625 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -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 } diff --git a/internal/apply/apply.go b/internal/apply/apply.go index 2f04871..a821353 100644 --- a/internal/apply/apply.go +++ b/internal/apply/apply.go @@ -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 diff --git a/internal/apply/health_check_test.go b/internal/apply/health_check_test.go new file mode 100644 index 0000000..ab0d561 --- /dev/null +++ b/internal/apply/health_check_test.go @@ -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") + } +} diff --git a/internal/declaration/declaration.go b/internal/declaration/declaration.go index de2460c..2472adf 100644 --- a/internal/declaration/declaration.go +++ b/internal/declaration/declaration.go @@ -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") } diff --git a/internal/declaration/health.go b/internal/declaration/health.go new file mode 100644 index 0000000..d6395bf --- /dev/null +++ b/internal/declaration/health.go @@ -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 +} diff --git a/internal/declaration/health_test.go b/internal/declaration/health_test.go new file mode 100644 index 0000000..000d53a --- /dev/null +++ b/internal/declaration/health_test.go @@ -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) + } + } +} diff --git a/internal/link/bus.go b/internal/link/bus.go index 6266ed7..da91d15 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -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 diff --git a/internal/link/health_test.go b/internal/link/health_test.go index 38e999a..3021e09 100644 --- a/internal/link/health_test.go +++ b/internal/link/health_test.go @@ -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")) + } +} diff --git a/internal/link/messages.go b/internal/link/messages.go index cfd0eaf..c37263d 100644 --- a/internal/link/messages.go +++ b/internal/link/messages.go @@ -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. diff --git a/internal/link/messages_test.go b/internal/link/messages_test.go index d3aa3a7..1c2578e 100644 --- a/internal/link/messages_test.go +++ b/internal/link/messages_test.go @@ -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) diff --git a/internal/link/queue.go b/internal/link/queue.go index 3efa30d..a21ff72 100644 --- a/internal/link/queue.go +++ b/internal/link/queue.go @@ -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 +} diff --git a/internal/liveness/liveness.go b/internal/liveness/liveness.go index 86f773f..506ccad 100644 --- a/internal/liveness/liveness.go +++ b/internal/liveness/liveness.go @@ -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) { diff --git a/internal/liveness/readiness.go b/internal/liveness/readiness.go new file mode 100644 index 0000000..585c08e --- /dev/null +++ b/internal/liveness/readiness.go @@ -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()) +} diff --git a/internal/liveness/readiness_test.go b/internal/liveness/readiness_test.go new file mode 100644 index 0000000..9203844 --- /dev/null +++ b/internal/liveness/readiness_test.go @@ -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") +} diff --git a/internal/liveness/replay_test.go b/internal/liveness/replay_test.go index 1805bf5..39c438f 100644 --- a/internal/liveness/replay_test.go +++ b/internal/liveness/replay_test.go @@ -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) + } + } +} diff --git a/internal/liveness/runtime.go b/internal/liveness/runtime.go index 376f352..eccdc5a 100644 --- a/internal/liveness/runtime.go +++ b/internal/liveness/runtime.go @@ -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 {