Files
mesh-host/internal/liveness/readiness.go
T
jochen 5e189c2fca
mesh/merge-gate pass: builds mesh-host → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without (4 of 4 compose)
mesh/repo-check pass: THE CHANGE ALTERS ITS OWN CHECK (merge-check.sh): main's version judged it; the change's judges the pull requests after it merges; it…
mesh/delivery delivered
Say a tool check's errors once and as they stand, and test the link on a real bus
The evidence read "asking it: asking <subject>": the words are now the link's own.
A deadline is said as the time the check gave it. A refusal is said only when the bus
refused this question, not an earlier one under a grant since widened. merge-check.sh
runs the tests that need a bus against a throwaway nats-server when the toolchain has
one, and says so when it does not (hq issue 331).
2026-10-08 17:01:15 +02:00

337 lines
11 KiB
Go

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 errors.Is(err, context.DeadlineExceeded):
return false, fmt.Sprintf("no answer within %s", timeout)
case err != nil:
return false, 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())
}