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-group group feat/asks-answered-on-any-channel rejected: a member's own check failed
mesh/delivery superseded: a newer head of the same pull request
669 lines
26 KiB
Go
669 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},
|
|
// Root where the trusted parties run (novox/hq ADR 0259 §8): while an agent can become root there without a
|
|
// person, an answer proven there proves nothing.
|
|
{ID: "D-root", Asserts: "no agent can become root without a person on a machine where the router or a channel " +
|
|
"proving its sender runs: not by its own account, and not through a tool that runs its command as an account " +
|
|
"that can", From: "ADR 0259 §8", Kind: kindAgentRoot, Phase: 2, run: probeAgentRoot},
|
|
{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
|
|
}
|