Files
mesh-controller/cmd/mesh-controller/probes.go
T
jochen 8e8712e352 Look twice before saying a probe failed, and say conditions in machine names (hq issue 277)
D2 raised a resolver urgent on one query that timed out while its machine was
loaded, and its summary carried the resolver's address and socket text, so the
operator channel withheld the whole alert.

- D2 asks every question up to three times, all at once; a resolver that
  answers nothing is held for the next run and raised urgent when two runs
  in a row find it silent. A wrong answer is still raised at once.
- Findings a single look can be wrong about carry Confirm: raised on the
  second look in a row, kept while open, never cleared-and-reraised. Used by
  D2 silence, D3 (also asks discovery twice), D6 behind, D9, D13 unmeasured,
  probe-failed of the doctor, and blind watchdog rows.
- Probe seat asks (D8, D13) are asked again when the bus brought no answer.
- Summaries name machines and say things in words; addresses, paths,
  domains and raw errors move to the evidence (D2, D5, D8, D9, D13, S12).
- internal/outward mirrors the messenger's content rule, allowing the mesh's
  machine names; the keeper rewords a summary that would be withheld and keeps
  it whole in the evidence; a TestMain lint fails the suite on any raised or
  linted finding that would be withheld.
2026-10-06 18:40:33 +02:00

1246 lines
47 KiB
Go

package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"net"
"regexp"
"slices"
"sort"
"strings"
"sync"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/nats-io/nats.go/micro"
"github.com/novox/mesh-host/validate"
"golang.org/x/net/dns/dnsmessage"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/lease"
"github.com/novox/mesh-controller/internal/link"
"github.com/novox/mesh-controller/internal/overlay"
)
// The probes of the self-check's registry (doctor.go), one function each. A probe answers what is
// wrong now as observations, or an error when it could not tell — never "nothing" for "could not
// look" (ADR 0227 rule 4).
// awaitingPushBound is how long a machine may have something only its next push makes — an own
// secret of a module assigned since its last push — before that is said (novox/hq issue 275).
var awaitingPushBound = 30 * time.Minute
// kindAwaitingPush is D1's kind for a machine waiting for a push past awaitingPushBound.
const kindAwaitingPush = "awaiting-push"
// probeDeclarations is D1: every machine's declaration composes, and the node-engine's own validator
// (mesh-host's `validate`, the package the host runs) takes it.
//
// **Composed as the next push would compose it, without making anything** (Foreseeing, novox/hq issue
// 275). Between `assign` and `push` a module's own secrets are not made yet — the push makes them —
// and a composition that only reads them failed, which raised an urgent "nothing can be sent" about a
// machine the next push sent to without a word. So a secret the push WILL make is composed with a
// stand-in and named; one the push would be refused on is refused here, with the push's words. A
// machine whose declaration composes and validates with only such stand-ins is waiting for a push,
// not broken: nothing is said until awaitingPushBound passes from when it began waiting, and then a
// warning (`awaiting-push`) naming what the push will make — never urgent, because nothing is wrong
// that a push does not fix.
func probeDeclarations(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
open := d.open
nodes, err := open.inventory.Nodes(ctx)
if err != nil {
return nil, err
}
var out []conditions.Observation
// The private network every declaration is composed with. Not computable is not "could not
// look": it is the finding — no machine's declaration composes without it — and the machines that
// do not resolve, which are usually why, are named below.
gens, gensErr := generators(ctx, open)
if gensErr != nil {
if ctx.Err() != nil {
return nil, ctx.Err()
}
out = append(out, conditions.Observation{Scope: conditions.ScopeMesh, ID: "private-network",
Token: "uncomposable", Severity: conditions.Urgent,
Summary: "the private network cannot be computed, so no machine's declaration composes",
Said: firstLine(gensErr.Error())})
}
for _, n := range nodes {
plan, settings, err := planFor(ctx, open, n.Name)
if err != nil && !unresolvable(err) {
return nil, fmt.Errorf("%s cannot be worked out: %w", n.Name, err)
}
var problems, foreseen []string
if err == nil && gensErr == nil {
var declared sendable
if declared, err = declarationWith(ctx, open, n.Name, plan, settings, gens, Foreseeing); err == nil {
foreseen = declared.foreseen
// With the order it was last sent, so the validator reads the envelope a machine is sent
// — its epoch included (novox/hq to-be 45 §6).
var serr error
if declared.Sequence, serr = open.inventory.Sequence(ctx, n.ID); serr == nil {
declared.Epoch, serr = open.inventory.SentEpoch(ctx, n.ID)
}
if serr != nil {
return nil, serr // the store, not the machine: the probe could not run
}
var body []byte
if body, err = declared.Body(); err == nil {
problems = validate.Declaration(body)
}
}
}
if ctx.Err() != nil {
return nil, ctx.Err()
}
switch {
case err != nil:
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: n.Name, Token: "uncomposable",
Machine: n.Name, Severity: conditions.Urgent,
Summary: fmt.Sprintf("nothing can be sent to %s: its declaration does not compose — %s", n.Name,
oneLine(err.Error())),
Said: oneLine(err.Error())})
case len(problems) > 0:
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: n.Name, Token: "refused",
Machine: n.Name, Severity: conditions.Urgent,
Summary: fmt.Sprintf("%s's node-engine would refuse its declaration whole: %d problem(s), the first: %s",
n.Name, len(problems), problems[0]),
Said: strings.Join(problems, "; ")})
case len(foreseen) > 0:
since, err := open.inventory.AwaitingSince(ctx, n.Name)
if err != nil {
return nil, err // the store, not the machine: the probe could not run
}
if waited := time.Since(since); waited > awaitingPushBound {
out = append(out, awaitingPush(n.Name, foreseen, since, waited))
}
}
}
return out, nil
}
// awaitingPush is D1's warning for a machine whose declaration composes only once its next push has
// made what it names, past awaitingPushBound (novox/hq issue 275).
func awaitingPush(node string, foreseen []string, since time.Time, waited time.Duration) conditions.Observation {
made := strings.Join(foreseen, ", ")
return conditions.Observation{Scope: conditions.ScopeMachine, ID: node, Token: kindAwaitingPush,
Kind: kindAwaitingPush, Machine: node, Severity: conditions.Warning,
Summary: fmt.Sprintf("%s has waited %s for a push: its declaration composes once the next push makes %s — "+
"`push %s` sends it", node, waited.Round(time.Minute), made, node),
Said: fmt.Sprintf("waiting since %s; the next push makes %s", since.UTC().Format(time.RFC3339), made)}
}
// probeResolvers is D2: every holder of the mesh's resolver answers each machine's name with its
// address for IPv4, and with no address and no error for IPv6 — NODATA, not NXDOMAIN, which musl
// takes as final (issue 262).
//
// **One unanswered question is not a resolver that does not answer** (novox/hq issue 277). D2 raised a
// resolver urgent on one query that timed out while its machine was loaded — a push applying, a build
// starting twenty containers — and thirty asked right after were answered at once. So every question is
// asked up to resolverTries times, a little apart, and all of them at once rather than one after
// another; a resolver that answered none of a question's tries is held back (Confirm) and raised, urgent,
// when the next run finds it not answering too. A resolver that *answered wrongly* — another address,
// NXDOMAIN for IPv6, an IPv6 address — said so on every try that answered, and is raised at once: that is
// an answer, not the absence of one.
//
// **Said in machine names** (ADR 0234 §6): the summary names the resolver's machine and the machine
// asked about; the names with their domain, the resolver's address and the socket's words are the
// evidence.
func probeResolvers(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
inv := d.open.inventory
shelf, err := inv.Catalogue(ctx)
if err != nil {
return nil, err
}
holders, err := replicatedHolders(ctx, inv, shelf)
if err != nil {
return nil, err
}
resolvers := holders["mesh-dns-resolver"]
if len(resolvers) == 0 {
return nil, nil // the mesh holds no resolver of its own: nothing is asked of one
}
places, err := onTheNetwork(ctx, inv, shelf)
if err != nil {
return nil, err
}
return askEveryResolver(ctx, resolvers, places, overlay.Suffix()), nil
}
// resolverQuestion is one question D2 asks one resolver, and what came of it.
type resolverQuestion struct {
holder, at string
place inventory.Overlay
kind dnsmessage.Type
asked resolverAsked
}
// askEveryResolver asks every resolver about every machine, every question at once, and says what each
// resolver got wrong or did not answer.
func askEveryResolver(ctx context.Context, resolvers map[string]string, places []inventory.Overlay,
suffix string) []conditions.Observation {
var questions []*resolverQuestion
for holder, at := range resolvers {
for _, p := range places {
if p.Address == "" {
continue
}
for _, kind := range []dnsmessage.Type{dnsmessage.TypeA, dnsmessage.TypeAAAA} {
questions = append(questions, &resolverQuestion{holder: holder, at: at, place: p, kind: kind})
}
}
}
var wg sync.WaitGroup
for _, q := range questions {
wg.Add(1)
go func(q *resolverQuestion) {
defer wg.Done()
name := q.place.Name + "." + suffix
q.asked = askResolverPatiently(ctx, q.at, name, q.kind, func(addresses []string, rcode dnsmessage.RCode) bool {
if q.kind == dnsmessage.TypeA {
return rcode == dnsmessage.RCodeSuccess && slices.Contains(addresses, q.place.Address)
}
return rcode == dnsmessage.RCodeSuccess && len(addresses) == 0
})
}(q)
}
wg.Wait()
type verdict struct {
wrong, unanswered, said []string
asked int
}
byHolder := map[string]*verdict{}
for _, q := range questions {
v := byHolder[q.holder]
if v == nil {
v = &verdict{}
byHolder[q.holder] = v
}
v.asked++
words, said := q.judged(suffix)
switch {
case said == "":
case q.asked.answered:
v.wrong = append(v.wrong, words)
v.said = append(v.said, said)
default:
v.unanswered = append(v.unanswered, words)
v.said = append(v.said, said)
}
}
var out []conditions.Observation
for holder, v := range byHolder {
if len(v.said) == 0 {
continue
}
sort.Strings(v.wrong)
sort.Strings(v.unanswered)
sort.Strings(v.said)
node := strings.TrimSuffix(holder, "."+suffix)
o := conditions.Observation{Scope: conditions.ScopeSeat, ID: "mesh-dns-resolver." + node,
Token: "wrong", Machine: node, Severity: conditions.Urgent, Said: strings.Join(v.said, "; ")}
if len(v.wrong) > 0 {
o.Summary = fmt.Sprintf("the mesh's resolver on %s does not answer machine names as it must: %s%s", node,
v.wrong[0], andMore(len(v.wrong)-1))
} else {
// Nothing but silence: held for the next run, which raises it if the resolver is still silent.
o.Confirm = true
o.Summary = fmt.Sprintf("the mesh's resolver on %s does not answer: %d of the %d question(s) about the "+
"machines' names went unanswered, each asked %d times — the first, %s", node, len(v.unanswered), v.asked,
resolverTries, v.unanswered[0])
}
out = append(out, o)
}
return sortedFound(out)
}
// andMore is ", and n more" for n other things; nothing for none.
func andMore(n int) string {
if n <= 0 {
return ""
}
return fmt.Sprintf(", and %d more", n)
}
// family is a question's kind as the mesh says it.
func family(kind dnsmessage.Type) string {
if kind == dnsmessage.TypeAAAA {
return "IPv6"
}
return "IPv4"
}
// rcodeWords is a resolver's answer code as a person reads it, with the name a resolver's log uses.
func rcodeWords(rcode dnsmessage.RCode) string {
switch rcode {
case dnsmessage.RCodeNameError:
return "no such name (NXDOMAIN)"
case dnsmessage.RCodeServerFailure:
return "a failure (SERVFAIL)"
case dnsmessage.RCodeRefused:
return "a refusal (REFUSED)"
case dnsmessage.RCodeFormatError:
return "a format error (FORMERR)"
case dnsmessage.RCodeNotImplemented:
return "not implemented (NOTIMP)"
}
return strings.TrimPrefix(rcode.String(), "RCode")
}
// judged is a question's outcome: in words, naming machines only, for a summary; and in full — the
// name with its domain, the addresses, the socket's words — for the evidence. Both empty when it was
// answered as it must be.
func (q *resolverQuestion) judged(suffix string) (words, said string) {
a, node, name := q.asked, q.place.Name, q.place.Name+"."+suffix
about := fmt.Sprintf("%s's %s address", node, family(q.kind))
switch {
case a.right:
return "", ""
case !a.answered:
return about, fmt.Sprintf("%s for %s: no answer from %s in %d tries: %s", q.kind, name, q.at, a.tries,
oneLine(errors.Join(a.errs...).Error()))
case a.rcode != dnsmessage.RCodeSuccess && q.kind == dnsmessage.TypeAAAA:
return fmt.Sprintf("it answered %s for %s, not NODATA: a musl machine takes that as no such name",
rcodeWords(a.rcode), about),
fmt.Sprintf("%s answered %s for %s (AAAA), not NODATA", q.at, a.rcode, name)
case a.rcode != dnsmessage.RCodeSuccess:
return fmt.Sprintf("it answered %s for %s", rcodeWords(a.rcode), about),
fmt.Sprintf("%s answered %s for %s (A)", q.at, a.rcode, name)
case q.kind == dnsmessage.TypeAAAA:
return fmt.Sprintf("it answered an IPv6 address for %s; the mesh has none", node),
fmt.Sprintf("%s answered %v for %s (AAAA); the mesh has no IPv6 addresses", q.at, a.addresses, name)
}
return fmt.Sprintf("it answered %s with another address than the mesh gave %s", about, node),
fmt.Sprintf("%s answered %v for %s (A), not %s", q.at, a.addresses, name, q.place.Address)
}
// How patiently D2 asks: each question up to resolverTries times, each try waiting resolverWithin for
// its answer, the tries resolverPause apart and that pause doubled each time. Every question is asked at
// once, so a resolver that answers nothing costs the run about seven seconds, within probeWithin.
var (
resolverTries = 3
resolverWithin = 2 * time.Second
resolverPause = 250 * time.Millisecond
)
// resolverAsked is one question asked patiently: the answer that settled it, or every try's error.
type resolverAsked struct {
// right says a try was answered as it must be; answered that some try was answered at all.
right, answered bool
addresses []string
rcode dnsmessage.RCode
tries int
errs []error
}
// askResolverPatiently asks one question until it is answered as right says, at most resolverTries
// times. A wrong answer is asked again too — a resolver restarting may say NXDOMAIN for a moment — and
// stands if no try answers rightly.
func askResolverPatiently(ctx context.Context, at, name string, kind dnsmessage.Type,
right func([]string, dnsmessage.RCode) bool) resolverAsked {
var a resolverAsked
pause := resolverPause
for a.tries < resolverTries {
if a.tries > 0 {
select {
case <-ctx.Done():
a.errs = append(a.errs, ctx.Err())
return a
case <-time.After(pause):
}
pause *= 2
}
a.tries++
addresses, rcode, err := askResolver(ctx, at, name, kind)
if err != nil {
a.errs = append(a.errs, err)
continue
}
a.answered, a.addresses, a.rcode = true, addresses, rcode
if right(addresses, rcode) {
a.right = true
return a
}
}
return a
}
// resolverPort is where a resolver answers; a variable so a test can stand one up.
var resolverPort = "53"
// askResolver asks one resolver one question over UDP, and answers the addresses and the code.
func askResolver(ctx context.Context, at, name string, kind dnsmessage.Type) ([]string, dnsmessage.RCode, error) {
q, err := dnsmessage.NewName(strings.TrimSuffix(name, ".") + ".")
if err != nil {
return nil, 0, err
}
msg := dnsmessage.Message{Header: dnsmessage.Header{ID: uint16(time.Now().UnixNano()), RecursionDesired: true},
Questions: []dnsmessage.Question{{Name: q, Type: kind, Class: dnsmessage.ClassINET}}}
packed, err := msg.Pack()
if err != nil {
return nil, 0, err
}
dialer := net.Dialer{Timeout: resolverWithin}
conn, err := dialer.DialContext(ctx, "udp", net.JoinHostPort(at, resolverPort))
if err != nil {
return nil, 0, err
}
defer conn.Close()
_ = conn.SetDeadline(time.Now().Add(resolverWithin))
if _, err := conn.Write(packed); err != nil {
return nil, 0, err
}
buf := make([]byte, 1500)
n, err := conn.Read(buf)
if err != nil {
return nil, 0, fmt.Errorf("no answer from %s: %w", at, err)
}
var answer dnsmessage.Message
if err := answer.Unpack(buf[:n]); err != nil {
return nil, 0, fmt.Errorf("an answer from %s that cannot be read: %w", at, err)
}
var addresses []string
for _, a := range answer.Answers {
switch r := a.Body.(type) {
case *dnsmessage.AResource:
addresses = append(addresses, net.IP(r.A[:]).String())
case *dnsmessage.AAAAResource:
addresses = append(addresses, net.IP(r.AAAA[:]).String())
}
}
return addresses, answer.RCode, nil
}
// probeHolders is D3: every seat that serves verbs has, on every machine that holds it and is heard
// from, a holder answering the bus's discovery for that seat. A machine past its heartbeat's bound is
// S1's, and is not asked about here.
func probeHolders(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
entries, err := d.open.inventory.Catalogued(ctx)
if err != nil {
return nil, err
}
recorded, err := d.open.inventory.Holdings(ctx)
if err != nil {
return nil, err
}
heard := heardMachines(d)
delivered, err := readDeliveries(ctx, d.open.inventory)
if err != nil {
return nil, err
}
now := time.Now()
expected := map[string]map[string]bool{} // seat → machine
add := func(seat, node string) {
if expected[seat] == nil {
expected[seat] = map[string]bool{}
}
expected[seat][node] = true
}
for _, s := range catalogue.SeatsWithAProtocol() {
if len(s.Serves) == 0 || s.Name == catalogue.ControllerSeatName {
continue // a seat with no verb has nothing to answer with; this controller is answering now
}
onRecord := false
for _, h := range recorded {
if h.Claim == s.Name && heard[h.Node] {
add(s.Name, h.Node)
onRecord = true
}
}
if onRecord && s.Scope == catalogue.ScopeMesh {
continue
}
for _, e := range entries {
for _, c := range e.Manifest.Claims {
if c.Name != s.Name {
continue
}
for _, node := range e.On {
// Assigned is not there yet: a holder is expected once its machine was sent it and
// had time to report, never between `assign` and `push` (novox/hq issue 275).
if heard[node] && (s.Scope == catalogue.ScopeNode || !onRecord) &&
delivered[node].settled(e.Manifest.Module, now) {
add(s.Name, node)
}
}
}
}
}
if len(expected) == 0 {
return nil, nil
}
answering, err := discoverHolders(ctx, d.js.Conn())
if err != nil {
return nil, err
}
// **A holder whose answer to one discovery came late is asked again** (novox/hq issue 277): a
// runtime on a loaded machine answers last, and once past discoveryPatience. What a second discovery
// still does not hear is held for the next run (Confirm), and said when that run does not hear it
// either.
missing := func() bool {
for seat, nodes := range expected {
for node := range nodes {
if !answering[seat][node] {
return true
}
}
}
return false
}
if missing() {
again, err := discoverHolders(ctx, d.js.Conn())
if err != nil {
return nil, err
}
for seat, nodes := range again {
for node := range nodes {
if answering[seat] == nil {
answering[seat] = map[string]bool{}
}
answering[seat][node] = true
}
}
}
var out []conditions.Observation
for seat, nodes := range expected {
for node := range nodes {
if answering[seat][node] {
continue
}
out = append(out, conditions.Observation{Scope: conditions.ScopeSeat, ID: seat + "." + node,
Token: "silent", Machine: node, Severity: conditions.Warning, Confirm: true,
Summary: fmt.Sprintf("%s's holder on %s does not answer the bus: its verbs reach nothing there", seat, node),
Said: fmt.Sprintf("no answer for %s from %s to the bus's discovery", seat, node)})
}
}
return sortedFound(out), nil
}
// reportGrace is how long a machine is given, after it was sent a declaration, to apply it and report
// before what the declaration carries is expected to answer (novox/hq issue 275).
var reportGrace = 10 * time.Minute
// lastSend is what a machine was last sent, as far as the self-check needs it: which modules the send
// carried, when, and whether the machine has reported acting on it.
type lastSend struct {
// carried is the modules its last send carried; nil when that was not kept (a send from before
// it was, or one by hand), and then every module is taken as carried — as before this existed.
carried map[string]string
// sent is when its current declaration went to it; nil if nothing ever did.
sent *time.Time
// current says its last report names the declaration last sent.
current bool
}
// settled says whether a module's holder on this machine can be asked to answer now: the machine was
// sent a declaration carrying it, and has reported acting on that declaration or had reportGrace to.
// A machine never sent anything holds nothing the mesh put there.
func (d lastSend) settled(module string, now time.Time) bool {
if d.sent == nil {
return false
}
if d.carried != nil {
if _, in := d.carried[module]; !in {
return false
}
}
return d.current || now.Sub(*d.sent) > reportGrace
}
// readDeliveries is every machine's last send, by name, from the records a send and a report keep.
func readDeliveries(ctx context.Context, inv *inventory.Inventory) (map[string]lastSend, error) {
reports, err := inv.LastReports(ctx)
if err != nil {
return nil, err
}
out := map[string]lastSend{}
for _, r := range reports {
builds, known, err := inv.SentBuilds(ctx, r.Node)
if err != nil {
return nil, err
}
d := lastSend{sent: r.Sent, current: r.Current}
if known {
d.carried = builds
}
out[r.Node] = d
}
return out, nil
}
// heardMachines is every machine within its heartbeat's bound, as the watchdogs last saw.
func heardMachines(d *doctor) map[string]bool {
out := map[string]bool{}
if d.watchdogs == nil {
return out
}
f := d.watchdogs.lastFacts()
if f == nil {
return out
}
for _, m := range f.machines {
if !m.lastHeard.IsZero() && f.now.Sub(m.lastHeard) <= heartbeatBound(m.every) && !m.asleep() {
out[m.name] = true
}
}
return out
}
// discoveryPatience is how long the discovery's answers are waited for after the last arrived, and at
// the most (ADR 0197: a large runtime's answer arrives last).
const (
discoveryQuiet = 1500 * time.Millisecond
discoveryPatience = 8 * time.Second
)
// discoverHolders asks the bus's discovery who serves what, and answers seat → machine for every
// endpoint a seat's verb is served on.
func discoverHolders(ctx context.Context, conn *nats.Conn) (map[string]map[string]bool, error) {
inbox := conn.NewRespInbox()
sub, err := conn.SubscribeSync(inbox)
if err != nil {
return nil, err
}
defer func() { _ = sub.Unsubscribe() }()
if err := conn.PublishRequest("$SRV.INFO", inbox, nil); err != nil {
return nil, fmt.Errorf("asking the bus who serves what: %w", err)
}
out := map[string]map[string]bool{}
deadline := time.Now().Add(discoveryPatience)
for time.Now().Before(deadline) {
wait, cancel := context.WithTimeout(ctx, discoveryQuiet)
msg, err := sub.NextMsgWithContext(wait)
cancel()
if err != nil {
if ctx.Err() != nil {
return nil, ctx.Err()
}
break
}
var info micro.Info
if json.Unmarshal(msg.Data, &info) != nil {
continue
}
for _, e := range info.Endpoints {
seat, node := e.Metadata["seat"], e.Metadata["node"]
if seat == "" {
continue
}
if node == "" {
node = info.ID
}
if out[seat] == nil {
out[seat] = map[string]bool{}
}
out[seat][node] = true
}
// A holder that announces itself as its seat, by name and machine.
if out[info.Name] == nil {
out[info.Name] = map[string]bool{}
}
out[info.Name][info.ID] = true
}
return out, nil
}
// probeArchives is D4: every archive the mesh keeps is held by its manifest in the artifact store.
func probeArchives(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
inv := d.open.inventory
kept, err := inv.KeptArchives(ctx)
if err != nil {
return nil, err
}
if len(kept) == 0 {
return nil, nil
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
return nil, err
}
store, err := artifactStoreAddress(ctx, inv, shelf, "")
if err != nil {
return nil, err
}
if store == "" {
return nil, errors.New("the mesh keeps archives and has no artifact store on its network to ask")
}
var report collectionReport
if unasked, stopped := askHeld(ctx, artifacts.Store{Address: store}, kept, &report); stopped != "" {
return nil, fmt.Errorf("%d kept archive(s) could not be asked about: %s", unasked, stopped)
}
var out []conditions.Observation
if n := len(report.Unheld); n > 0 {
out = append(out, conditions.Observation{Scope: conditions.ScopeMesh, ID: "artifact-store", Token: "archives-unheld",
Severity: conditions.Warning,
Summary: fmt.Sprintf("%d of %d kept archive(s) are not held by a manifest: the store's collector would "+
"delete them", n, len(kept)),
Said: "first: " + report.Unheld[0]})
}
if n := len(report.Missing); n > 0 {
out = append(out, conditions.Observation{Scope: conditions.ScopeMesh, ID: "artifact-store", Token: "archives-missing",
Kind: "archives-missing", Severity: conditions.Warning,
Summary: fmt.Sprintf("%d kept archive(s) are not in the artifact store at all", n),
Said: "first: " + report.Missing[0]})
}
return out, nil
}
// consumerFarBehind is how far a durable consumer may be from its stream's head (D6).
const consumerFarBehind = 1000
// probeConsumers is D6: every durable consumer the mesh expects exists, as the controller defines it,
// and is near its stream's head.
func probeConsumers(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
_, consumers, err := expectedBusObjects(ctx, d.open.inventory)
if err != nil {
return nil, err
}
js := d.js.Context()
var out []conditions.Observation
for _, c := range consumers {
if ctx.Err() != nil {
return nil, ctx.Err()
}
info, err := js.ConsumerInfo(c.Stream, c.Name, nats.Context(ctx))
who := c.Stream + "." + c.Name
switch {
case errors.Is(err, nats.ErrConsumerNotFound) || errors.Is(err, nats.ErrStreamNotFound):
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: who, Token: "missing",
Kind: "consumer-lost", Severity: conditions.Warning,
Summary: fmt.Sprintf("%s is not on the bus: what it delivers reaches nobody", consumerWords(c)),
Said: "no consumer " + c.Name + " on " + c.Stream})
continue
case err != nil:
return nil, fmt.Errorf("%s cannot be read: %w", consumerWords(c), err)
}
if differs := consumerDiffers(c, info.Config); differs != "" {
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: who, Token: "redefined",
Severity: conditions.Warning,
Summary: fmt.Sprintf("%s is not as the controller defines it: %s", consumerWords(c), differs),
Said: differs})
}
if c.Stream != "NODES" && info.NumPending > consumerFarBehind {
// Its own kind (novox/hq to-be 45 §7): healer H4 answers a consumer far behind, and only
// this — a redefined consumer is not one a reset repairs.
// One look at a consumer that a burst has put behind is not one that stays behind: held for
// the next run, so H4 resets nothing on a single sample (novox/hq issue 277).
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: who, Token: "behind",
Kind: kindConsumerBehind, Severity: conditions.Warning, Confirm: true,
Summary: fmt.Sprintf("%s is %d message(s) behind its stream's head", consumerWords(c), info.NumPending),
Said: fmt.Sprintf("%d pending, %d handed out and not settled", info.NumPending, info.NumAckPending)})
}
}
return sortedFound(out), nil
}
// consumerWords is a consumer as the mesh says it, with its name.
func consumerWords(c broker.Consumer) string {
return fmt.Sprintf("%s (%s on %s)", link.ConsumerInWords(c.Stream, c.Name), c.Name, c.Stream)
}
// consumerDiffers is what about a consumer on the bus is not as defined; empty when nothing is. The
// delivery subject is not compared: one kept as it was while a holder is bound is the assertion's
// stated choice (novox/hq issue 156).
func consumerDiffers(want broker.Consumer, have nats.ConsumerConfig) string {
var differs []string
haveFilters := append([]string(nil), have.FilterSubjects...)
if have.FilterSubject != "" {
haveFilters = append(haveFilters, have.FilterSubject)
}
wantFilters := append([]string(nil), want.Filters...)
sort.Strings(haveFilters)
sort.Strings(wantFilters)
if !slices.Equal(haveFilters, wantFilters) {
differs = append(differs, fmt.Sprintf("filters %v, defined %v", haveFilters, wantFilters))
}
if have.AckPolicy != nats.AckExplicitPolicy {
differs = append(differs, "acknowledges "+have.AckPolicy.String()+", defined explicit")
}
if wantMax := want.MaxDeliver; wantMax != 0 && have.MaxDeliver != wantMax {
differs = append(differs, fmt.Sprintf("hands a message over %d times, defined %d", have.MaxDeliver, wantMax))
}
if wantWait := time.Duration(want.AckWaitSeconds) * time.Second; wantWait != 0 && have.AckWait != wantWait {
differs = append(differs, fmt.Sprintf("waits %s for an acknowledgement, defined %s", have.AckWait, wantWait))
}
if want.MaxAckPending != 0 && have.MaxAckPending != want.MaxAckPending {
differs = append(differs, fmt.Sprintf("hands out %d at once, defined %d", have.MaxAckPending, want.MaxAckPending))
}
if wantPush, havePush := want.Push || want.Queue != "", have.DeliverSubject != ""; wantPush != havePush {
differs = append(differs, "delivers by the other shape (push or pull) than defined")
}
return strings.Join(differs, "; ")
}
// probeStreams is D7: every stream the controller defines exists as defined, and its own buckets.
func probeStreams(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
streams, _, err := expectedBusObjects(ctx, d.open.inventory)
if err != nil {
return nil, err
}
js := d.js.Context()
var out []conditions.Observation
for _, s := range streams {
info, err := js.StreamInfo(s.Name, nats.Context(ctx))
switch {
case errors.Is(err, nats.ErrStreamNotFound):
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: s.Name, Token: "stream-missing",
Kind: "stream-wrong", Severity: conditions.Urgent,
Summary: fmt.Sprintf("the stream %s is not on the bus: %s", s.Name, firstLine(s.Why)),
Said: "no stream " + s.Name})
continue
case err != nil:
return nil, fmt.Errorf("the stream %s cannot be read: %w", s.Name, err)
}
if differs := streamDiffers(s, info.Config); differs != "" {
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: s.Name, Token: "stream-redefined",
Kind: "stream-wrong", Severity: conditions.Warning,
Summary: fmt.Sprintf("the stream %s is not as the controller defines it: %s", s.Name, differs),
Said: differs})
}
}
for _, bucket := range broker.ControllerBuckets() {
if _, err := js.StreamInfo("KV_"+bucket, nats.Context(ctx)); errors.Is(err, nats.ErrStreamNotFound) {
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: bucket, Token: "bucket-missing",
Kind: "stream-wrong", Severity: conditions.Urgent,
Summary: fmt.Sprintf("the controller's bucket %s is not on the bus: what it keeps there is not kept", bucket),
Said: "no bucket " + bucket})
} else if err != nil {
return nil, fmt.Errorf("the bucket %s cannot be read: %w", bucket, err)
}
}
return sortedFound(out), nil
}
// streamDiffers is what about a stream on the bus is not as defined; empty when nothing is.
func streamDiffers(want broker.Stream, have nats.StreamConfig) string {
var differs []string
haveSubjects := append([]string(nil), have.Subjects...)
wantSubjects := append([]string(nil), want.Subjects...)
sort.Strings(haveSubjects)
sort.Strings(wantSubjects)
if !slices.Equal(haveSubjects, wantSubjects) {
differs = append(differs, fmt.Sprintf("subjects %v, defined %v", haveSubjects, wantSubjects))
}
retention, perSubject := nats.LimitsPolicy, int64(want.MaxMsgsPerSubject)
switch want.Retention {
case broker.RetentionWorkQueue:
retention = nats.WorkQueuePolicy
case broker.RetentionLastPerSubject:
perSubject = 1
}
if have.Retention != retention {
differs = append(differs, fmt.Sprintf("keeps by %s, defined %s", have.Retention, retention))
}
if perSubject != 0 && have.MaxMsgsPerSubject != perSubject {
differs = append(differs, fmt.Sprintf("keeps %d per subject, defined %d", have.MaxMsgsPerSubject, perSubject))
}
return strings.Join(differs, "; ")
}
// ipLike finds the addresses in a ban list, whatever shape its holder answers in.
var ipLike = regexp.MustCompile(`\b(?:\d{1,3}\.){3}\d{1,3}\b`)
// probeBans is D8: no address the mesh owns is in any machine's ban list. The mesh's own addresses are
// every machine's private address and the address its endpoint names; the ban lists are asked of each
// machine's holder of `node-intrusion-prevention` that is heard from.
func probeBans(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
inv := d.open.inventory
places, err := inv.Overlays(ctx)
if err != nil {
return nil, err
}
// address → whose, in machine names: what the summary says. The address, and the endpoint's name
// it was found by, are the evidence (ADR 0234 §6).
owned, via := map[string]string{}, map[string]string{}
for _, p := range places {
if p.Address != "" {
owned[p.Address] = p.Name + "'s private address"
}
if host, _, err := net.SplitHostPort(p.Endpoint); err == nil && host != "" {
if ip := net.ParseIP(host); ip != nil {
owned[ip.String()] = p.Name + "'s endpoint"
} else if addrs, err := net.DefaultResolver.LookupHost(ctx, host); err == nil {
for _, a := range addrs {
owned[a], via[a] = p.Name+"'s endpoint", host
}
}
}
}
entries, err := inv.Catalogued(ctx)
if err != nil {
return nil, err
}
heard := heardMachines(d)
var holders []string
for _, e := range entries {
for _, c := range e.Manifest.Claims {
if c.Name == "node-intrusion-prevention" {
for _, node := range e.On {
if heard[node] && !slices.Contains(holders, node) {
holders = append(holders, node)
}
}
}
}
}
sort.Strings(holders)
var out []conditions.Observation
// Every machine asked at once: one after another, four holders that each wait their bound
// outlast the probe's thirty seconds, and the probe says nothing about any of them.
answers := make([]json.RawMessage, len(holders))
errs := make([]error, len(holders))
var wg sync.WaitGroup
for i, node := range holders {
wg.Add(1)
go func(i int, node string) {
defer wg.Done()
answers[i], errs[i] = askSeatTool(ctx, d.js.Conn(), "node-intrusion-prevention", "banned", node)
}(i, node)
}
wg.Wait()
var unasked []string
for i, node := range holders {
answer, err := answers[i], errs[i]
if err != nil {
unasked = append(unasked, err.Error())
continue
}
var banned, whose []string
for _, ip := range ipLike.FindAllString(string(answer), -1) {
who, ours := owned[ip]
said := ip + " (" + who + ")"
if h := via[ip]; h != "" {
said = ip + " (" + who + ", " + h + ")"
}
if !ours || slices.Contains(banned, said) {
continue
}
banned = append(banned, said)
if !slices.Contains(whose, who) {
whose = append(whose, who)
}
}
if len(banned) > 0 {
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: node, Token: "bans-the-mesh",
Machine: node, Severity: conditions.Urgent,
Summary: fmt.Sprintf("%s's ban list holds %d of the mesh's own address(es) — %s: the mesh is locked out "+
"of itself there (ADR 0186)", node, len(banned), strings.Join(whose, ", ")),
Said: strings.Join(banned, ", ")})
}
}
if len(unasked) > 0 {
return nil, fmt.Errorf("a ban list could not be read: %s", strings.Join(unasked, "; "))
}
return out, nil
}
// probeStatus is D9: `status` answers in full within ten seconds, from a summary composed lately.
func probeStatus(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
started := time.Now()
stale := ""
if statusFrom != nil {
answer, err := statusFrom.answer(ctx)
if err != nil {
stale = err.Error()
} else if m, ok := answer.(map[string]any); ok {
if failed, _ := m["lastAttemptFailed"].(string); failed != "" {
stale = "its last composition failed: " + failed
}
if composed, err := time.Parse(time.RFC3339, fmt.Sprint(m["composed"])); err == nil &&
time.Since(composed) > 3*statusEvery+statusComposeWithin {
stale = fmt.Sprintf("it answers a summary composed %s ago", time.Since(composed).Round(time.Second))
}
}
} else if _, err := theThreeQuestions(ctx, d.open); err != nil {
stale = err.Error()
}
took := time.Since(started)
var out []conditions.Observation
if took > 10*time.Second || stale != "" {
why := stale
if why == "" {
why = fmt.Sprintf("it took %s", took.Round(time.Millisecond))
}
// Timed once, on a controller that may be busy for a moment: held for the next run (novox/hq
// issue 277). Why it was slow is the evidence: a composition's failure carries a store's words.
out = append(out, conditions.Observation{Scope: conditions.ScopeCore, ID: "controller", Token: "status-slow",
Machine: d.host, Severity: conditions.Warning, Confirm: true,
Summary: "status does not answer in full within ten seconds, from a summary composed lately",
Said: why})
}
return out, nil
}
// probeCoreBuilds is D10: every machine runs the node-engine and the node tools the mesh holds, or a
// plan is rolling one of them out. The node-engine says its build in each report; the node tools' is
// the build the machine was last sent, once it reported applying that send.
func probeCoreBuilds(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
inv := d.open.inventory
current, err := inv.CurrentBuilds(ctx)
if err != nil {
return nil, err
}
plans, err := inv.OpenPlans(ctx)
if err != nil {
return nil, err
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
return nil, err
}
hostVersions := deliveredVersions(shelf[hostModule])
rolling := map[string]bool{}
for _, p := range plans {
for m := range p.Modules {
rolling[m] = true
}
}
nodes, err := inv.Nodes(ctx)
if err != nil {
return nil, err
}
reports, err := inv.LastReports(ctx)
if err != nil {
return nil, err
}
applied := map[string]bool{}
for _, r := range reports {
applied[r.Node] = r.Current
}
heard := heardMachines(d)
var out []conditions.Observation
for _, n := range nodes {
if !heard[n.Name] {
continue // a machine not heard from is S1's
}
var behind []string
// **A node-engine says its version as the directory it was delivered into** — its archive's
// digest, twelve characters (catalogue.versionOf, ADR 0141) — not the commit it was built
// from. Compared as a commit, every machine read as behind right after a push sent it the
// current one (2026-10-06). An engine placed by hand reports its link-time stamp instead,
// which a commit can match.
if !rolling[hostModule] && engineBehind(n.HostVersion, hostVersions, current[hostModule].Commit) {
behind = append(behind, fmt.Sprintf("the node-engine %s, the mesh holds %s (built from %s)",
n.HostVersion, strings.Join(hostVersions, " or "), short(current[hostModule].Commit)))
}
assigned, err := inv.Assigned(ctx, n.Name)
if err != nil {
return nil, err
}
tools := current[broker.RuntimeModule].Commit
if tools != "" && slices.Contains(assigned, broker.RuntimeModule) && !rolling[broker.RuntimeModule] {
sent, known, err := inv.SentBuilds(ctx, n.Name)
if err != nil {
return nil, err
}
switch {
case !known:
case !sameCommit(sent[broker.RuntimeModule], tools):
behind = append(behind, fmt.Sprintf("the node tools %s, the mesh holds %s",
short(orNotKnown(sent[broker.RuntimeModule])), short(tools)))
case !applied[n.Name]:
behind = append(behind, "the node tools it was last sent, not yet reported applied")
}
}
if len(behind) > 0 {
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: n.Name, Token: "core-behind",
Machine: n.Name, Severity: conditions.Warning,
Summary: fmt.Sprintf("%s runs %s, and no plan is rolling them out", n.Name, strings.Join(behind, "; ")),
Said: strings.Join(behind, "; ")})
}
}
return out, nil
}
// deliveredVersions are the versions a module's registered build is delivered as: the last element
// of every resource path under a `versions/` directory, which registration filled from the artifact's
// digest (catalogue `${version}`). The node-engine names itself by that directory.
func deliveredVersions(m catalogue.Manifest) []string {
var out []string
for _, r := range m.Resources {
path, _ := r["path"].(string)
before, version, found := strings.Cut(path, "/versions/")
if !found || before == "" || version == "" || strings.Contains(version, "/") || strings.Contains(version, "$") {
continue
}
if !slices.Contains(out, version) {
out = append(out, version)
}
}
return out
}
// engineBehind says a node-engine's reported version is not the build the mesh holds: neither the
// directory that build is delivered as, nor (for an engine placed by hand) its commit. A machine that
// has not said, or a mesh that holds no delivered build, is not behind anything.
func engineBehind(reported string, delivered []string, commit string) bool {
if reported == "" || len(delivered) == 0 {
return false
}
return !slices.Contains(delivered, reported) && !sameCommit(reported, commit)
}
// hostModule is the node-engine's module.
const hostModule = "mesh-host"
// sameCommit says two commits are one, either written short.
func sameCommit(a, b string) bool {
if a == "" || b == "" {
return false
}
return strings.HasPrefix(a, b) || strings.HasPrefix(b, a)
}
func orNotKnown(s string) string {
if s == "" {
return "(not known)"
}
return s
}
// probeWatchdogs is DW: the watchdogs ran within three of their intervals. The doctor and the
// watchdogs watch each other: S10 is the other half.
func probeWatchdogs(_ context.Context, d *doctor) ([]conditions.Observation, error) {
if d.watchdogs == nil {
return nil, nil // a self-check run outside the serving controller has none to watch
}
ticked := d.watchdogs.lastTick()
since := ticked
if since.IsZero() {
since = d.watchdogs.started
}
if time.Since(since) <= 3*watchEvery {
return nil, nil
}
return []conditions.Observation{{Scope: conditions.ScopeCore, ID: "watchdogs", Token: "silent",
Machine: d.host, Severity: conditions.Urgent,
Summary: fmt.Sprintf("the watchdogs of the signals table have not run since %s: no late signal is being said",
since.UTC().Format("2006-01-02 15:04 MST")),
Said: fmt.Sprintf("no tick for %s", time.Since(since).Round(time.Second))}}, nil
}
// How a probe asks a seat's holder: each try waits seatAskWithin, at most seatAskTries tries,
// seatAskPause apart.
var (
seatAskTries = 2
seatAskWithin = 10 * time.Second
seatAskPause = 500 * time.Millisecond
)
// askSeatTool asks one machine's holder of a node seat a verb and answers its result; the holder's
// own refusal is an error.
func askSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string) (json.RawMessage, error) {
if who, declared := declaredBy(ctx, seat, verb); !declared {
return nil, fmt.Errorf("%s asks %s.%s, which it does not declare in the probe registry — and so the "+
"controller is not granted it", who, seat, verb)
}
// **Asked again when the bus brought no answer** (novox/hq issue 277), while the probe's bound
// leaves room for a whole try: a holder on a loaded machine that misses one question is not a
// holder that cannot answer. The holder's own refusal is an answer, and is not asked again.
var answer link.Answer
var err error
for try := 0; try < seatAskTries; try++ {
if try > 0 {
if deadline, ok := ctx.Deadline(); ok && time.Until(deadline) < seatAskWithin+seatAskPause {
break
}
select {
case <-ctx.Done():
return nil, err
case <-time.After(seatAskPause):
}
}
if answer, err = link.AskSeatTool(ctx, conn, seat, verb, node, map[string]any{}, seatAskWithin); err == nil {
break
}
}
if err != nil {
return nil, err
}
if answer.Error != "" {
return nil, fmt.Errorf("%s's %s.%s refused: %s", node, seat, verb, answer.Error)
}
return answer.Result, nil
}
// oneLine is a message of several lines said on one, its runs of space made one.
func oneLine(s string) string { return strings.Join(strings.Fields(s), " ") }
// probeLease is D5 (novox/hq to-be 45 §4, §6): exactly one lease holder — the key on the bus names this
// controller at the epoch it acts under, and the mesh's record has that epoch and no other open — and no
// message from a stale epoch refused in the last interval.
func probeLease(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
if d.js == nil {
return nil, errors.New("this controller is not on the bus")
}
st := theLease.standing()
var out []conditions.Observation
split := func(token, summary, said string) {
out = append(out, conditions.Observation{Scope: conditions.ScopeCore, ID: "controller.lease", Token: token,
Machine: d.host, Severity: conditions.Urgent, Summary: summary, Said: said})
}
if st.Unleased != "" {
split("unheld", "this controller acts without the lease, so nothing keeps another from acting beside it",
st.Unleased)
return out, nil
}
api, err := jetstream.New(d.js.Conn())
if err != nil {
return nil, err
}
kv, err := api.KeyValue(ctx, broker.LeaseBucket)
if err != nil {
return nil, fmt.Errorf("the lease bucket cannot be read: %w", err)
}
holder, found, err := lease.Current(ctx, kv)
if err != nil {
return nil, fmt.Errorf("the lease cannot be read: %w", err)
}
switch {
case !found:
split("split", fmt.Sprintf("nobody holds the lease on the bus, and this controller acts as epoch %d", st.Epoch),
"the lease's key is absent")
case holder.Instance != instance || holder.Epoch != st.Epoch:
split("split", fmt.Sprintf("the lease on the bus names %s at epoch %d, and this controller (%s) acts as "+
"epoch %d: two controllers believe they may act", holder.Instance, holder.Epoch, instance, st.Epoch),
fmt.Sprintf("held by %s, epoch %d", holder.Instance, holder.Epoch))
}
// The record: one epoch open, this one.
epochs, err := d.open.inventory.EpochsSince(ctx, time.Now().Add(-time.Hour))
if err != nil {
return nil, fmt.Errorf("the epochs the mesh issued cannot be read: %w", err)
}
var open []string
for _, e := range epochs {
if e.Ended == nil && e.Epoch != st.Epoch {
open = append(open, fmt.Sprintf("epoch %d (%s)", e.Epoch, e.Instance))
}
}
if len(open) > 0 {
split("open-epochs", fmt.Sprintf("the mesh's record holds %s open beside this controller's epoch %d: a "+
"holder that neither gave the lease back nor was found expired", strings.Join(open, ", "), st.Epoch),
strings.Join(open, ", "))
}
// And no message from a stale epoch in the last interval.
for _, w := range link.StaleRefusals.Within(time.Now().Add(-doctorEvery)) {
if w.Epoch <= 0 || uint64(w.Epoch) >= st.Epoch {
continue
}
out = append(out, conditions.Observation{Scope: conditions.ScopeCore, ID: fmt.Sprintf("controller.epoch-%d", w.Epoch),
Token: "stale-epoch", Machine: firstOf(w.Receivers), Severity: conditions.Urgent,
Summary: fmt.Sprintf("%d declaration(s) from epoch %d — older than this controller's %d — reached %s in the "+
"last %s and were refused: a controller that lost the lease is still sending", w.Count, w.Epoch, st.Epoch,
strings.Join(w.Receivers, ", "), doctorEvery),
Said: fmt.Sprintf("%d refused, the last at %s", w.Count, w.Last.UTC().Format(time.RFC3339))})
}
return sortedFound(out), nil
}