Files
mesh-controller/cmd/mesh-controller/doctor.go
T
jochen 175b28ee42
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
mesh/delivery-group group fix/314-a-large-answer-is-paged-not-lost delivered: every member is delivered
Page an answer larger than one message of the bus, and keep overviews brief (hq issue 314)
The client library refuses to send a reply over the bus's max_payload, and the
controller only logged it: conditions and status answered nobody for hours on
2026-10-08 while calls said each was answered in 130 ms, and the operator's
channel read nothing. An answer too large is now held under its call and paged
to the caller that asks, on the same subject; a caller that does not page is
told in words, and calls says it. The overviews no longer carry every finding:
conditions and status list each condition with its newest evidence, doctor at
most twenty findings a probe (probe= gives one whole), and the JSON overviews
are sent once, as data, instead of twice.
2026-10-08 12:19:16 +02:00

664 lines
26 KiB
Go

package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"sort"
"strings"
"sync"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/link"
)
// The self-check: `doctor` (novox/hq to-be 45 §4, ADR 0227 rule 6).
//
// **The design's invariants, run against the running mesh.** A probe is one live invariant of a
// design — every machine's declaration composes and validates, every resolver answers, every seat's
// holder answers, every stream and consumer is there as defined — with an id, a bound, and the
// condition it raises when the invariant does not hold. The serving controller runs the registry every
// five minutes, each probe given thirty seconds; a probe that errors or does not finish raises
// `probe-failed` for itself, because an unanswered probe is never a pass. Every run ends with a
// heartbeat on the bus (`doctor-heartbeat`, S10), which mesh-watcher listens for from a second
// machine: a controller that stops checking is itself said, through a channel that does not pass
// through it.
//
// `doctor` answers the last run's verdict at once; `doctor run` runs now; `doctor probes` lists the
// registry; `doctor signals` says, for every row of the signals table, the age of its newest signal.
// The self-check's clocks.
var (
// doctorEvery is how often the registry runs; doctorFirstAfter how long after the controller
// starts the first run waits, so what the controller hears at its start has arrived.
doctorEvery = 5 * time.Minute
doctorFirstAfter = time.Minute
// probeWithin is each probe's bound.
probeWithin = 30 * time.Second
)
// The verdicts a probe can have.
const (
verdictPass = "pass"
verdictFail = "fail"
verdictFailedToRun = "failed-to-run"
verdictDeferred = "deferred"
)
// probe is one live invariant.
type probe struct {
ID string
Asserts string
From string
// Kind is the condition kind raised when the invariant does not hold.
Kind string
// Raises are the other kinds its findings carry, each its own (a consumer missing, one far behind):
// what a healer may be registered against (healers.go).
Raises []string
Phase int
// Deferred says why it is not run yet; empty for one that is.
Deferred string
// Asks are the seat verbs it calls. A probe may call no other (askSeatTool refuses), and the
// controller's grant names every one (a test over this registry): a probe whose question the bus
// refuses checks nothing (D8, 2026-10-06).
Asks []broker.SeatVerb
run func(ctx context.Context, d *doctor) ([]conditions.Observation, error)
}
// probeRegistry is the registry, in to-be 45's order. **The registry is the design's live form**: a
// probe added to a design is a row added here.
var probeRegistry = []probe{
{ID: "D1", Asserts: "every machine's declaration composes, and passes the node-engine's validation",
From: "issues 236, 263, 275", Kind: "declaration-refused", Raises: []string{kindAwaitingPush}, Phase: 1,
run: probeDeclarations},
{ID: "D2", Asserts: "every holder of the mesh's resolver answers a machine name for IPv4, and NODATA for IPv6",
From: "issue 262", Kind: "resolver-wrong", Phase: 1, run: probeResolvers},
{ID: "D3", Asserts: "every seat on record that serves verbs has a live holder that answers, on every " +
"machine that is heard from", From: "issues 208, 218", Kind: "holder-silent", Phase: 1, run: probeHolders},
{ID: "D4", Asserts: "every kept archive is held by a manifest", From: "issue 253",
Kind: "archives-unheld", Phase: 1, run: probeArchives},
{ID: "D5", Asserts: "exactly one lease holder — the key names this controller at its epoch, and the record " +
"holds no other epoch open; no message from a stale epoch refused in the last interval",
From: "issue 204", Kind: "lease-split", Phase: 2, run: probeLease},
{ID: "D6", Asserts: "every durable consumer the mesh expects exists with its definition, and is near its " +
"stream's head", From: "issues 248, 266", Kind: "consumer-wrong", Raises: []string{"consumer-lost", kindConsumerBehind},
Phase: 1, run: probeConsumers},
{ID: "D7", Asserts: "every stream the controller defines exists with its definition, and its own buckets",
From: "issue 208", Kind: "stream-wrong", Phase: 1, run: probeStreams},
{ID: "D8", Asserts: "no address the mesh owns — a machine's private address or its endpoint — is in a ban list",
From: "issue 238", Kind: "own-address-banned", Phase: 1, run: probeBans,
Asks: []broker.SeatVerb{{Seat: "node-intrusion-prevention", Verb: "banned"}}},
{ID: "D9", Asserts: "status answers in full within ten seconds, from a summary composed lately",
From: "issue 265", Kind: "status-slow", Phase: 1, run: probeStatus},
{ID: "D10", Asserts: "every machine runs the node-engine and node tools builds the mesh holds, or is inside " +
"a plan's window", From: "the version split", Kind: "core-behind", Phase: 1, run: probeCoreBuilds},
{ID: "D11", Asserts: "no provider holds a consumer retired more than thirty days without a person deciding " +
"its cleanup", From: "ADR 0230", Kind: kindCleanupWaiting, Phase: 2, run: probeRetired},
{ID: probeBindingsID, Asserts: "every consumer of a provision that keeps its data is bound where it was last " +
"sent, or moves by a pin", From: "issue 273, ADR 0232", Kind: kindBindingMoved,
Raises: []string{kindBindingKept, kindBindingMoving}, Phase: 2, run: probeBindings},
{ID: probeDataID, Asserts: "every item of data a machine declares is measured, is there, holds what it held, is " +
"written where it should be, is backed up within its bound or sits on healthy redundant storage, and is no " +
"empty replacement of a copy kept elsewhere; what a machine no longer declares that is irreplaceable or " +
"valuable is retired, not forgotten", From: "issue 273, ADR 0233",
Kind: kindDataShrank, Raises: []string{kindEmptyReplacement, kindDataHeldTwice, kindDataQuiet, kindBackupStale,
kindDataUnmeasured, kindDataMissing, kindArrayDegraded, kindProtectionMissing, kindCleanupWaiting},
Phase: 2, run: probeData,
Asks: []broker.SeatVerb{{Seat: "node-backup", Verb: "backed-up"}}},
// The delivery's owner (novox/hq ADR 0239): what it holds past a bound of its own table is said here, by
// the controller, and H2 works it through the owner's `close`.
{ID: probeDeliveriesID, Asserts: "no delivery is held past its state's bound unsaid: mesh-delivery's " +
"`stalled`, each with the transition its table lets healer H2 take", From: "ADR 0239",
Kind: kindDeliveryStalled, Phase: 3, run: probeDeliveries},
{ID: "DW", Asserts: "the watchdogs of the signals table ran within three of their intervals",
From: "ADR 0227 rule 6: the watchers are watched", Kind: "watchdogs-silent", Phase: 1, run: probeWatchdogs},
// The core's health definitions (novox/hq to-be 45 §8, ADR 0236): what a core component's new build is
// judged by on its first machine, run against every machine between upgrades too.
{ID: "H-controller", Asserts: "the controller lease is held, renewed in time, by a controller that says it is " +
"ready: its self-check ran and status answered in full within ten seconds", From: "ADR 0236, to-be 45 §8",
Kind: kindCoreUnhealthy, Phase: 4, run: probeControllerHealth},
{ID: "H-engine", Asserts: "every machine heard from has reported its current declaration, under a node-engine " +
"build it names", From: "ADR 0236, to-be 45 §8", Kind: kindCoreUnhealthy, Phase: 4, run: probeEngineHealth},
{ID: "H-tools", Asserts: "every machine heard from that runs the node tools has them answering the bus",
From: "ADR 0236, to-be 45 §8", Kind: kindCoreUnhealthy, Phase: 4, run: probeToolsHealth},
{ID: "H-bus", Asserts: "every stream and durable consumer the mesh defines is on the bus, and a request crosses " +
"it to the machines' node tools and back", From: "ADR 0236, to-be 45 §8", Kind: kindCoreUnhealthy, Phase: 4,
run: probeBusHealth},
// The gate's verdicts and the witnesses' rollbacks (ADR 0236): each build that failed its gate keeps its
// condition until a newer build passes; each rollback a witness stands by is said.
{ID: gateProbe, Asserts: "no build that failed its gate, and no core component a witness put back, goes unsaid; " +
"a newer build that passes its gate clears it", From: "ADR 0236, to-be 45 §8", Kind: kindRolledBack,
Raises: []string{kindRollbackFailed}, Phase: 4, run: probeGates},
// The bus's planned step (ADR 0236): open while a person's bus upgrade runs, then checked by H-bus.
{ID: busStepProbe, Asserts: "a bus upgrade a person started is said while it runs, and is followed by the bus's " +
"health within its bound — or is said failed, with its snapshot as the way back", From: "ADR 0236, to-be 45 §8",
Kind: kindBusMaintenance, Raises: []string{kindBusUpgradeFailed}, Phase: 4, run: probeBusStep},
}
// probeVerdict is one probe's outcome in a run.
type probeVerdict struct {
ID string `json:"id"`
Verdict string `json:"verdict"`
Found []string `json:"found,omitempty"`
// Unconfirmed are findings one look can be wrong about, seen by this run and not the one before:
// raised if the next run sees them too (confirm.go). Not a pass, and not yet a condition.
Unconfirmed []string `json:"unconfirmed,omitempty"`
Error string `json:"error,omitempty"`
Took string `json:"took,omitempty"`
// FoundCount and More are an overview's: how many findings the probe had when the overview shows
// fewer, and how to read them all (novox/hq issue 314). Never in a heartbeat.
FoundCount int `json:"found-count,omitempty"`
More string `json:"more,omitempty"`
}
// overviewFindings is how many findings of one probe an overview carries. On 2026-10-08 D14 found 689
// stalled deliveries, and an overview carrying every one of them, beside the conditions they raised,
// was more than the bus carries in one message (novox/hq issue 314); one probe's verdict whole is
// `doctor probe=<id>`.
const overviewFindings = 20
// inOverview is a run as an overview carries it: each probe's findings at most overviewFindings, with
// their count and how to read them all. The run itself is not changed.
func inOverview(run doctorRun) doctorRun {
probes := make([]probeVerdict, len(run.Probes))
for i, p := range run.Probes {
total := len(p.Found) + len(p.Unconfirmed)
if total > overviewFindings {
p.FoundCount = len(p.Found)
if len(p.Found) > overviewFindings {
p.Found = p.Found[:overviewFindings:overviewFindings]
}
if room := overviewFindings - len(p.Found); len(p.Unconfirmed) > room {
p.Unconfirmed = p.Unconfirmed[:room:room]
}
p.More = fmt.Sprintf("%d found and %d unconfirmed, %d of them shown — `doctor probe=%s` gives them all",
p.FoundCount, total-p.FoundCount, len(p.Found)+len(p.Unconfirmed), p.ID)
}
probes[i] = p
}
run.Probes = probes
return run
}
// oneProbe is a run with only one probe's verdict, whole.
func oneProbe(run doctorRun, id string) (doctorRun, error) {
for _, p := range run.Probes {
if strings.EqualFold(p.ID, id) {
run.Probes = []probeVerdict{p}
return run, nil
}
}
return doctorRun{}, fmt.Errorf("the run %s has no probe %q — `doctor probes=true` lists them", run.Run, id)
}
// doctorCounts are a run's verdicts, counted.
type doctorCounts struct {
Passed int `json:"passed"`
Failed int `json:"failed"`
FailedToRun int `json:"failed-to-run"`
Deferred int `json:"deferred"`
}
// doctorRun is one run of the registry, and the body of its heartbeat.
type doctorRun struct {
Run string `json:"run"`
At time.Time `json:"at"`
Started time.Time `json:"started"`
Took string `json:"took"`
IntervalSeconds int `json:"interval-seconds"`
Counts doctorCounts `json:"counts"`
Probes []probeVerdict `json:"probes"`
// Controller is the machine that ran it, and Why what started it: the schedule, or a person.
Controller string `json:"controller"`
Why string `json:"why"`
// Unsaid is how many condition transitions this controller could not say, since it started.
Unsaid int `json:"unsaid,omitempty"`
}
// doctor is the registry and what its probes need.
type doctor struct {
open *stores
js *broker.JetStream
keeper *conditions.Keeper
teller conditions.Teller
watchdogs *watchdogs
host string
// confirm holds back what one run alone saw of a finding a single look can be wrong about.
confirm confirming
running sync.Mutex
mu sync.Mutex
last *doctorRun
ended time.Time
}
// lastRunEnded is when the last run ended; zero before the first.
func (d *doctor) lastRunEnded() time.Time {
d.mu.Lock()
defer d.mu.Unlock()
return d.ended
}
// lastRun is the last run's verdict; nil before the first.
func (d *doctor) lastRun() *doctorRun {
d.mu.Lock()
defer d.mu.Unlock()
return d.last
}
// keep runs the registry on its schedule until ctx ends.
func (d *doctor) keep(ctx context.Context) {
select {
case <-ctx.Done():
return
case <-time.After(doctorFirstAfter):
}
tick := time.NewTicker(doctorEvery)
defer tick.Stop()
for {
// Only the controller acting checks the mesh on a schedule: one standing by would say a
// heartbeat for a self-check that is not the mesh's.
if d.watchdogs == nil || d.watchdogs.acting == nil || d.watchdogs.acting() {
d.runOnce(ctx, "the schedule")
}
select {
case <-ctx.Done():
return
case <-tick.C:
}
}
}
var doctorRuns struct {
sync.Mutex
n uint64
}
// runOnce runs every probe the registry runs, keeps what each found, and says the heartbeat. One run
// at a time: a person's `doctor run` during a scheduled one waits for it.
func (d *doctor) runOnce(ctx context.Context, why string) doctorRun {
d.running.Lock()
defer d.running.Unlock()
doctorRuns.Lock()
doctorRuns.n++
n := doctorRuns.n
doctorRuns.Unlock()
started := time.Now()
run := doctorRun{Run: fmt.Sprintf("doctor-%d-%d", started.Unix(), n), Started: started.UTC(),
IntervalSeconds: int(doctorEvery / time.Second), Controller: d.host, Why: why}
type result struct {
obs []conditions.Observation
err error
took time.Duration
}
results := make([]result, len(probeRegistry))
var wg sync.WaitGroup
for i, p := range probeRegistry {
if p.run == nil {
continue
}
wg.Add(1)
go func(i int, p probe) {
defer wg.Done()
probing, cancel := context.WithTimeout(context.WithValue(ctx, probeAsksKey{}, p), probeWithin)
defer cancel()
began := time.Now()
done := make(chan result, 1)
go func() {
defer func() {
if r := recover(); r != nil {
done <- result{err: fmt.Errorf("the probe panicked: %v", r)}
}
}()
obs, err := p.run(probing, d)
done <- result{obs: obs, err: err}
}()
select {
case r := <-done:
r.took = time.Since(began)
results[i] = r
case <-probing.Done():
results[i] = result{err: fmt.Errorf("it did not finish within %s", probeWithin), took: time.Since(began)}
}
}(i, p)
}
wg.Wait()
var blind []conditions.Observation
for i, p := range probeRegistry {
v := probeVerdict{ID: p.ID}
r := results[i]
switch {
case p.run == nil:
v.Verdict = verdictDeferred
run.Counts.Deferred++
case r.err != nil:
v.Verdict, v.Error = verdictFailedToRun, r.err.Error()
run.Counts.FailedToRun++
// A probe that could not run once — a question timed out on a loaded machine — is said in
// the verdict at once, and raised as a condition when the next run cannot run it either.
blind = append(blind, conditions.Observation{Scope: conditions.ScopeProbe, ID: p.ID, Kind: "probe-failed",
Token: "failed", Severity: conditions.Warning, Confirm: true,
Summary: fmt.Sprintf("the probe %s (%s) could not run: what it checks is not known — never a pass", p.ID, p.Asserts),
Said: firstLine(r.err.Error())})
default:
raise, held := d.confirm.pass(ctx, d.keeper, p.ID, kindedAs(r.obs, p.Kind))
if err := d.keeper.Reconcile(ctx, p.ID, raise); err != nil {
v.Error = "what it found could not be kept: " + err.Error()
}
for _, o := range held {
v.Unconfirmed = append(v.Unconfirmed, o.Summary)
}
if len(raise) == 0 {
v.Verdict = verdictPass
run.Counts.Passed++
} else {
v.Verdict = verdictFail
run.Counts.Failed++
for _, o := range raise {
v.Found = append(v.Found, o.Summary)
}
}
}
if p.run != nil {
v.Took = r.took.Round(time.Millisecond).String()
}
run.Probes = append(run.Probes, v)
}
blind, _ = d.confirm.pass(ctx, d.keeper, sourceDoctor, blind)
if err := d.keeper.Reconcile(ctx, sourceDoctor, blind); err != nil {
fmt.Printf("the self-check's own failures could not be kept: %v\n", err)
}
ended := time.Now()
run.At, run.Took, run.Unsaid = ended.UTC(), ended.Sub(started).Round(time.Millisecond).String(), d.keeper.Unsaid()
d.mu.Lock()
d.last, d.ended = &run, ended
d.mu.Unlock()
d.sayHeartbeat(ctx, run)
return run
}
// sourceDoctor is what raises a probe's own failure to run.
const sourceDoctor = "doctor"
// kindedAs gives each observation of a probe the probe's kind where it named none.
func kindedAs(obs []conditions.Observation, kind string) []conditions.Observation {
out := make([]conditions.Observation, 0, len(obs))
for _, o := range obs {
if o.Kind == "" {
o.Kind = kind
}
out = append(out, o)
}
return out
}
// sayHeartbeat publishes the run's heartbeat. Not said is said here, and S10 on the second machine
// says it outward: the watcher hears nothing.
func (d *doctor) sayHeartbeat(ctx context.Context, run doctorRun) {
if d.teller == nil {
return
}
body, err := json.Marshal(run)
if err != nil {
fmt.Printf("the self-check's heartbeat could not be written: %v\n", err)
return
}
saying, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
if err := d.teller.PublishSeatEvent(saying, conditions.Seat, conditions.HeartbeatEvent, body); err != nil {
fmt.Printf("the self-check's heartbeat (%s) could NOT be said, so the watcher on the second machine "+
"will say the self-check is silent: %v\n", run.Run, err)
}
}
// doctorFrom is the serving controller's self-check; nil in any other process.
var doctorFrom *doctor
// doctorCommand is `doctor`, `doctor run`, `doctor probes` and `doctor signals`.
func doctorCommand(ctx context.Context, args []string) error {
sub := ""
if len(args) > 0 && !strings.HasPrefix(args[0], "-") {
sub, args = args[0], args[1:]
}
set := flag.NewFlagSet("doctor", flag.ContinueOnError)
asJSON := set.Bool("json", false, "as data")
probe := set.String("probe", "", "one probe's verdict, with every finding")
if rest, err := parseAround(set, args); err != nil {
return err
} else if len(rest) > 0 {
return errors.New("doctor [run|probes|signals] [--probe <id>] [--json]")
}
answer, err := doctorAnswer(ctx, sub, *probe)
if err != nil {
return err
}
if *asJSON {
return printJSON(answer)
}
fmt.Print(doctorText(answer))
return nil
}
// doctorAnswer is what the verb answers, as data.
//
// A verdict is an overview — each probe's findings at most overviewFindings — unless probe names one,
// which is then answered alone and whole (novox/hq issue 314).
func doctorAnswer(ctx context.Context, sub, probe string) (any, error) {
if probe != "" && sub != "" && sub != "run" {
return nil, fmt.Errorf("probe names one probe of a verdict; %s has none", sub)
}
switch sub {
case "":
if doctorFrom != nil {
if run := doctorFrom.lastRun(); run != nil {
return verdictAnswer(*run, time.Now(), probe)
}
return nil, fmt.Errorf("the self-check has not finished its first run yet: it runs %s after the "+
"controller starts, then every %s — `doctor run` runs it now", doctorFirstAfter, doctorEvery)
}
run, err := lastHeartbeat(ctx)
if err != nil {
return nil, err
}
return verdictAnswer(run, time.Now(), probe)
case "run":
d := doctorFrom
if d == nil {
local, closeIt, err := localDoctor(ctx)
if err != nil {
return nil, err
}
defer closeIt()
d = local
}
return verdictAnswer(d.runOnce(ctx, "asked by "+link.Caller()), time.Now(), probe)
case "probes":
return probesAnswer(), nil
case "signals":
if doctorFrom == nil || doctorFrom.watchdogs == nil {
return nil, errors.New("the age of each signal is known to the serving controller alone, which " +
"hears them: ask it through the mesh-controller seat's doctor verb")
}
return signalsAnswer(doctorFrom.watchdogs.lastFacts(), doctorFrom.watchdogs.lastTick()), nil
}
return nil, fmt.Errorf("doctor answers the last run, or `run`, `probes` or `signals` — not %q", sub)
}
// verdictAnswer is a run as the verb answers it, with its age.
func verdictAnswer(run doctorRun, now time.Time, probe string) (map[string]any, error) {
if probe != "" {
one, err := oneProbe(run, probe)
if err != nil {
return nil, err
}
run = one
} else {
run = inOverview(run)
}
return map[string]any{"run": run, "age": now.Sub(run.At).Round(time.Second).String(),
"note": "a probe that could not run is never a pass; each failure is an open condition until a run passes it, " +
"and one a single look can be wrong about — an unanswered question, a slow answer — is raised when two " +
"runs in a row see it"}, nil
}
// probesAnswer is the registry.
func probesAnswer() map[string]any {
var out []map[string]any
for _, p := range probeRegistry {
row := map[string]any{"id": p.ID, "asserts": p.Asserts, "from": p.From, "kind": p.Kind, "phase": p.Phase}
if len(p.Raises) > 0 {
row["raises"] = p.Raises
}
if p.Deferred != "" {
row["deferred"] = p.Deferred
}
out = append(out, row)
}
return map[string]any{"probes": out, "every": doctorEvery.String(), "each within": probeWithin.String()}
}
// signalsAnswer is every row of the signals table with the age of its newest signal.
func signalsAnswer(f *signalFacts, ticked time.Time) map[string]any {
var rows []map[string]any
for _, r := range signalsTable {
row := map[string]any{"row": r.Row, "signal": r.Signal, "emitter": r.Emitter, "bound": r.Bound,
"kind": r.Kind, "severity": r.Severity, "phase": r.Phase}
switch {
case r.Deferred != "":
row["deferred"] = r.Deferred
case f == nil:
row["newest"] = "not yet looked at"
default:
if err := r.needs(f); err != nil {
row["blind"] = err.Error()
} else if newest := r.newest(f); newest.IsZero() {
row["newest"] = "none heard"
} else {
row["newest"] = newest.UTC().Format(time.RFC3339)
row["age"] = f.now.Sub(newest).Round(time.Second).String()
}
}
rows = append(rows, row)
}
out := map[string]any{"signals": rows, "every": watchEvery.String()}
if !ticked.IsZero() {
out["looked"] = ticked.UTC().Format(time.RFC3339)
}
return out
}
// doctorText is an answer as a person reads it.
func doctorText(answer any) string {
body, _ := json.Marshal(answer)
var b strings.Builder
var verdict struct {
Run doctorRun `json:"run"`
Age string `json:"age"`
}
if json.Unmarshal(body, &verdict) == nil && verdict.Run.Run != "" {
r := verdict.Run
fmt.Fprintf(&b, "%s, %s ago (took %s, %s): %d passed, %d failed, %d could not run, %d not built yet\n\n",
r.Run, verdict.Age, r.Took, r.Why, r.Counts.Passed, r.Counts.Failed, r.Counts.FailedToRun, r.Counts.Deferred)
for _, p := range r.Probes {
fmt.Fprintf(&b, " %-4s %-14s %s\n", p.ID, p.Verdict, p.Took)
for _, f := range p.Found {
fmt.Fprintf(&b, " %s\n", f)
}
for _, f := range p.Unconfirmed {
fmt.Fprintf(&b, " unconfirmed, raised if the next run sees it too: %s\n", f)
}
if p.More != "" {
fmt.Fprintf(&b, " … %s\n", p.More)
}
if p.Error != "" {
fmt.Fprintf(&b, " %s\n", p.Error)
}
}
return b.String()
}
pretty, _ := json.MarshalIndent(answer, "", " ")
return string(pretty) + "\n"
}
// lastHeartbeat is the newest run's heartbeat, read from the events stream: what a process other than
// the serving controller answers `doctor` from.
func lastHeartbeat(ctx context.Context) (doctorRun, error) {
var run doctorRun
err := onTheBus(func(conn *nats.Conn) error {
js, err := conn.JetStream(nats.Context(ctx))
if err != nil {
return err
}
msg, err := js.GetLastMsg(broker.EventsStream, link.SeatEventSubject(conditions.Seat, conditions.HeartbeatEvent))
if errors.Is(err, nats.ErrMsgNotFound) {
return errors.New("the self-check has said no heartbeat on the bus in the last week: it is not running")
}
if err != nil {
return fmt.Errorf("the self-check's last heartbeat cannot be read: %w", err)
}
return json.Unmarshal(msg.Data, &run)
})
return run, err
}
// localDoctor is a self-check run by a process other than the serving controller: its own stores,
// its own connection, and its own keeper, closed after.
func localDoctor(ctx context.Context) (*doctor, func(), error) {
open, err := openStores(ctx)
if err != nil {
return nil, nil, err
}
js, err := dialTheBus()
if err != nil {
open.Close()
return nil, nil, err
}
k, err := keeperOn(ctx, js.Conn())
if err != nil {
js.Close()
open.Close()
return nil, nil, err
}
jsCtx := js.Context()
d := &doctor{open: open, js: js, keeper: k, teller: link.OverNATS{Conn: js.Conn(), JS: jsCtx},
host: controlHost(ctx, open.inventory)}
return d, func() {
flushing, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
k.Close(flushing)
js.Close()
open.Close()
}, nil
}
// sortedFound is a probe's findings in a stated order, so two runs over one mesh say the same.
func sortedFound(obs []conditions.Observation) []conditions.Observation {
sort.Slice(obs, func(i, j int) bool { return obs[i].Key() < obs[j].Key() })
return obs
}
// probeAsksKey carries the running probe, so a seat verb it calls is checked against what it declares.
type probeAsksKey struct{}
// declaredBy says whether the probe running in ctx declared a seat verb; outside a probe, false.
func declaredBy(ctx context.Context, seat, verb string) (string, bool) {
p, ok := ctx.Value(probeAsksKey{}).(probe)
if !ok {
return "a caller outside the self-check", false
}
for _, v := range p.Asks {
if v.Seat == seat && v.Verb == verb {
return p.ID, true
}
}
return p.ID, false
}