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.
This commit is contained in:
jochen
2026-10-06 18:40:33 +02:00
parent 7aa98e64ce
commit 8e8712e352
16 changed files with 1215 additions and 85 deletions
+1 -1
View File
@@ -146,7 +146,7 @@ func TestBindingFindingsSayAKeptMoveAndAMoveNothingAskedFor(t *testing.T) {
"game": {"postgres-database": home},
"network": {"wildcard-resolution": {Node: "home", Module: "resolver"}},
}
found := bindingFindings(plan, bound, nil)
found := linted(bindingFindings(plan, bound, nil))
if len(found) != 2 {
t.Fatalf("found %+v", found)
}
+67
View File
@@ -0,0 +1,67 @@
package main
import (
"context"
"sync"
"github.com/novox/mesh-controller/internal/conditions"
)
// **A finding one look can be wrong about is raised on the second look in a row** (novox/hq issue 277).
//
// D2 raised its resolver urgent on one unanswered question, asked once, while the resolver's machine
// was loaded by a push and a build starting; thirty questions right after were answered at once. A
// single sample of something that is answered over the network, or timed on a machine under load,
// is not the invariant failing. So a source marks such a finding `Confirm`, and this holds it back:
//
// - raised when the source's previous look saw it too — two runs of the self-check in a row (five
// minutes apart), or two ticks of the watchdogs (half a minute) — at the severity the source says;
// - kept while it is already open, however it is seen, so a condition the next look still sees is
// never cleared and raised again (flapping is not news; a reopening within ReopenWithin still is);
// - cleared, as everything is, by the look that no longer sees it.
//
// A finding held back is not a pass: the self-check's verdict names it as unconfirmed. A finding that is
// a definite answer — a resolver that answered wrongly, a stream that is not there — is not marked, and
// is raised at once.
type confirming struct {
mu sync.Mutex
// last is, by source, the keys of the Confirm findings its previous look saw.
last map[string]map[string]bool
}
// pass splits one source's findings into what is raised now and what is held for the next look, and
// remembers what it saw. An open condition that cannot be read is kept: unknown is not a reason to
// hold a finding back.
func (c *confirming) pass(ctx context.Context, keeper *conditions.Keeper, source string,
found []conditions.Observation) (raise, held []conditions.Observation) {
c.mu.Lock()
defer c.mu.Unlock()
if c.last == nil {
c.last = map[string]map[string]bool{}
}
before, seen := c.last[source], map[string]bool{}
for _, o := range found {
if !o.Confirm {
raise = append(raise, o)
continue
}
key := o.Key()
seen[key] = true
if before[key] || isOpen(ctx, keeper, key) {
raise = append(raise, o)
continue
}
held = append(held, o)
}
c.last[source] = seen
return raise, held
}
// isOpen says whether a condition is open; one that cannot be read is taken as open.
func isOpen(ctx context.Context, keeper *conditions.Keeper, key string) bool {
if keeper == nil {
return false
}
_, open, err := keeper.Get(ctx, key)
return open || err != nil
}
+19 -14
View File
@@ -220,9 +220,9 @@ func probeData(ctx context.Context, d *doctor) ([]conditions.Observation, error)
for _, m := range machines {
if m.askErr != nil {
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: m.name, Token: kindDataUnmeasured,
Kind: kindDataUnmeasured, Machine: m.name, Severity: conditions.Warning,
Kind: kindDataUnmeasured, Machine: m.name, Severity: conditions.Warning, Confirm: true,
Summary: fmt.Sprintf("%s's backup holder did not say what it measured of the data declared there, so "+
"nothing about that data is known this run: %s", m.name, firstLine(m.askErr.Error())),
"nothing about that data is known", m.name),
Said: firstLine(m.askErr.Error())})
}
why := fmt.Sprintf("no longer declared on %s: its module was unassigned there, or is no longer pulled in", m.name)
@@ -407,10 +407,11 @@ func dataFindings(records []inventory.DataRecord, peaks map[string]int64, shelf
if now.Sub(*r.RetiredAt) > cleanupAfter {
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: id, Token: "cleanup",
Kind: kindCleanupWaiting, Machine: r.Machine, Severity: conditions.Warning, Resolver: conditions.ResolverOperator,
Summary: fmt.Sprintf("%s of %s on %s (%s, %s) has been retired %d days — kept at %s since %s; `cleanup "+
"delete %s %s %s --why …` once a person has decided, or assign %s there again",
Summary: fmt.Sprintf("%s of %s on %s (%s, %s) has been retired %d days — kept where it was since %s; "+
"`cleanup delete %s %s %s --why …` once a person has decided, or assign %s there again",
r.Item, r.Module, r.Machine, r.Class, sizeWords(r.Size), int(now.Sub(*r.RetiredAt).Hours()/24),
orUnknownPath(r.Path), r.RetiredWhy, r.Machine, r.Module, r.Item, r.Module)})
r.RetiredWhy, r.Machine, r.Module, r.Item, r.Module),
Said: "kept at " + orUnknownPath(r.Path)})
}
continue
}
@@ -425,14 +426,16 @@ func dataFindings(records []inventory.DataRecord, peaks map[string]int64, shelf
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: id, Token: kindProtectionMissing,
Kind: kindProtectionMissing, Machine: r.Machine, Severity: severity, Resolver: conditions.ResolverOperator,
Summary: fmt.Sprintf("%s of %s on %s (%s) is said to be protected by the redundancy of the storage it is on, "+
"and %s is on nothing the backup holder can read as redundant: it has no protection the mesh can see",
r.Item, r.Module, r.Machine, class, orUnknownPath(r.Path))})
"and it is on nothing the backup holder can read as redundant: it has no protection the mesh can see",
r.Item, r.Module, r.Machine, class),
Said: orUnknownPath(r.Path) + " is on no redundant storage the backup holder can read"})
}
if r.MeasureError != "" && strings.Contains(r.MeasureError, "does not exist") {
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: id, Token: kindDataMissing,
Kind: kindDataMissing, Machine: r.Machine, Severity: severity, Resolver: conditions.ResolverOperator,
Summary: fmt.Sprintf("%s of %s on %s (%s) is gone: %s does not exist any more", r.Item, r.Module, r.Machine,
class, orUnknownPath(r.Path))})
Summary: fmt.Sprintf("%s of %s on %s (%s) is gone: where it was declared, nothing exists any more", r.Item,
r.Module, r.Machine, class),
Said: orUnknownPath(r.Path) + " does not exist: " + firstLine(r.MeasureError)})
} else if peak, ok := peaks[r.Key()]; ok && r.Size != nil && inventory.Comparable(r.Precision) &&
*r.Size*2 < peak && peak-*r.Size >= shrinkFloor && inventory.Dataset(r.Precision) != "" {
// Several items on one dataset share its size: one condition for the dataset, as loud as the
@@ -486,9 +489,10 @@ func dataFindings(records []inventory.DataRecord, peaks map[string]int64, shelf
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine,
ID: ds.machine + ".dataset." + strings.ReplaceAll(ds.dataset, "/", "-"), Token: kindDataShrank,
Kind: kindDataShrank, Machine: ds.machine, Severity: severityOf(ds.class), Resolver: conditions.ResolverOperator,
Summary: fmt.Sprintf("the dataset %s on %s shrank to %s from %s within %d days — more than half of what it held "+
"is gone; it holds %s", ds.dataset, ds.machine, sizeWords(&ds.size), sizeWords(&ds.peak),
int(shrinkWindow.Hours()/24), strings.Join(ds.items, ", "))})
Summary: fmt.Sprintf("the dataset on %s that holds %s shrank to %s from %s within %d days — more than half of "+
"what it held is gone", ds.machine, strings.Join(ds.items, ", "), sizeWords(&ds.size), sizeWords(&ds.peak),
int(shrinkWindow.Hours()/24)),
Said: "the dataset " + ds.dataset})
}
// The redundant storage watched data is on: one condition per array, as loud as the most precious
// item on it — the array, not each item, is what degrades.
@@ -512,8 +516,9 @@ func dataFindings(records []inventory.DataRecord, peaks map[string]int64, shelf
ID: machine + ".array." + strings.NewReplacer("/", "-", ":", "-").Replace(red.Kind+"-"+red.Where),
Token: kindArrayDegraded, Kind: kindArrayDegraded, Machine: machine, Severity: severityOf(class),
Resolver: conditions.ResolverOperator,
Summary: fmt.Sprintf("the %s storage %s on %s %s: %s — and it is what protects %s", red.Kind, red.Where, machine,
state, firstLine(red.Said), strings.Join(names, ", "))})
Summary: fmt.Sprintf("the %s storage on %s that protects %s %s", red.Kind, machine, strings.Join(names, ", "),
state),
Said: fmt.Sprintf("the %s storage %s %s: %s", red.Kind, red.Where, state, firstLine(red.Said))})
}
// The same item on several machines: a copy that is in use and far smaller than one kept elsewhere is
// an empty replacement. Only against a retired copy — a module running on two machines on purpose
+6 -6
View File
@@ -42,7 +42,7 @@ const houseManifest = `{"module":"house","version":"1",
func findingsByKind(obs []conditions.Observation) map[string]conditions.Observation {
out := map[string]conditions.Observation{}
for _, o := range obs {
for _, o := range linted(obs) {
out[o.Kind] = o
}
return out
@@ -62,7 +62,7 @@ func TestAnEmptyReplacementOfAConsumersDataIsSaid(t *testing.T) {
consumerCopy{Node: "anchor", Module: "postgres", Consumer: app, Size: bytesOf(9 << 20), Class: "valuable"})
}
copies[1].Class = "irreplaceable" // the photo site's own database says so
got := dataFindings(nil, nil, nil, copies, nil, time.Now())
got := linted(dataFindings(nil, nil, nil, copies, nil, time.Now()))
if len(got) != 5 {
t.Fatalf("%d findings for five empty replacements: %+v", len(got), got)
}
@@ -110,7 +110,7 @@ func TestConsumerDataActiveTwiceWithoutSizesIsAWarning(t *testing.T) {
{Node: "home", Module: "minio", Consumer: "mesh_home_photos"},
{Node: "anchor", Module: "minio", Consumer: "mesh_home_photos"},
}
got := dataFindings(nil, nil, nil, copies, nil, time.Now())
got := linted(dataFindings(nil, nil, nil, copies, nil, time.Now()))
if len(got) != 1 || got[0].Kind != kindDataHeldTwice || got[0].Severity != conditions.Warning {
t.Fatalf("%+v", got)
}
@@ -206,7 +206,7 @@ func TestRetiredDataWaitingThirtyDaysIsSaid(t *testing.T) {
{Machine: "home", Module: "house", Item: "config", Class: "irreplaceable", Size: bytesOf(1 << 30), RetiredAt: &old},
{Machine: "home", Module: "attic", Item: "boxes", Class: "irreplaceable", Size: bytesOf(1 << 30), RetiredAt: &recent},
}
got := dataFindings(records, nil, shelfFor(t, houseManifest), nil, nil, now)
got := linted(dataFindings(records, nil, shelfFor(t, houseManifest), nil, nil, now))
if len(got) != 1 || got[0].Kind != kindCleanupWaiting || !strings.Contains(got[0].Summary, "cleanup delete home house config") {
t.Fatalf("%+v", got)
}
@@ -387,7 +387,7 @@ func TestTheArrayUnderDataIsWatched(t *testing.T) {
Path: "/tank/" + item, Size: bytesOf(40 << 40), FirstSeen: now.Add(-24 * time.Hour), MeasuredAt: when(now),
Redundancy: &inventory.Redundancy{Kind: "zfs", Where: "tank", Healthy: healthy, Said: "pool 'tank' is DEGRADED"}}
}
got := dataFindings([]inventory.DataRecord{on("films", &sick), on("shows", &sick)}, nil, shelf, nil, nil, now)
got := linted(dataFindings([]inventory.DataRecord{on("films", &sick), on("shows", &sick)}, nil, shelf, nil, nil, now))
if len(got) != 1 || got[0].Kind != kindArrayDegraded || got[0].Severity != conditions.Urgent ||
!strings.Contains(got[0].Summary, "media/films, media/shows") {
t.Fatalf("%+v", got)
@@ -450,7 +450,7 @@ func TestADatasetShrinksOnceAndAPartialSizeIsNeverCompared(t *testing.T) {
}
films, shows := on("films", 30<<40), on("shows", 30<<40)
peaks := map[string]int64{films.Key(): 90 << 40, shows.Key(): 90 << 40}
got := dataFindings([]inventory.DataRecord{films, shows}, peaks, shelf, nil, nil, now)
got := linted(dataFindings([]inventory.DataRecord{films, shows}, peaks, shelf, nil, nil, now))
if len(got) != 1 || got[0].Kind != kindDataShrank || got[0].Severity != conditions.Urgent ||
!strings.Contains(got[0].Summary, "media/films, media/shows") {
t.Fatalf("%+v", got)
+24 -7
View File
@@ -120,8 +120,11 @@ type probeVerdict struct {
ID string `json:"id"`
Verdict string `json:"verdict"`
Found []string `json:"found,omitempty"`
Error string `json:"error,omitempty"`
Took string `json:"took,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"`
}
// doctorCounts are a run's verdicts, counted.
@@ -156,6 +159,8 @@ type doctor struct {
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
@@ -267,21 +272,27 @@ func (d *doctor) runOnce(ctx context.Context, why string) doctorRun {
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,
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:
if err := d.keeper.Reconcile(ctx, p.ID, kindedAs(r.obs, p.Kind)); err != nil {
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()
}
if len(r.obs) == 0 {
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 r.obs {
for _, o := range raise {
v.Found = append(v.Found, o.Summary)
}
}
@@ -291,6 +302,7 @@ func (d *doctor) runOnce(ctx context.Context, why string) doctorRun {
}
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)
}
@@ -406,7 +418,9 @@ func doctorAnswer(ctx context.Context, sub string) (any, error) {
// verdictAnswer is a run as the verb answers it, with its age.
func verdictAnswer(run doctorRun, now time.Time) map[string]any {
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"}
"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"}
}
// probesAnswer is the registry.
@@ -472,6 +486,9 @@ func doctorText(answer any) string {
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.Error != "" {
fmt.Fprintf(&b, " %s\n", p.Error)
}
+11 -1
View File
@@ -74,6 +74,15 @@ func TestARunKeepsWhatEachProbeFoundAndSaysItsHeartbeat(t *testing.T) {
k := conditions.NewKeeper(t.Context(), conditions.Options{Store: store, History: store})
defer k.Close(context.Background())
d := &doctor{keeper: k, teller: told, host: "anchor"}
// A probe that cannot run is said in the verdict at once, and as a condition when the next run cannot
// run it either (novox/hq issue 277): one timed-out question is not a probe gone blind.
first := d.runOnce(t.Context(), "a test")
if first.Counts != (doctorCounts{Passed: 1, Failed: 1, FailedToRun: 2, Deferred: 1}) {
t.Fatalf("counted %+v", first.Counts)
}
if open, _ := k.Open(t.Context()); len(open) != 1 || open[0].Key != "machine.anchor.refused" {
t.Fatalf("after one run, open: %+v", open)
}
run := d.runOnce(t.Context(), "a test")
if run.Counts != (doctorCounts{Passed: 1, Failed: 1, FailedToRun: 2, Deferred: 1}) {
t.Fatalf("counted %+v", run.Counts)
@@ -93,7 +102,7 @@ func TestARunKeepsWhatEachProbeFoundAndSaysItsHeartbeat(t *testing.T) {
}
}
// The heartbeat, in the shape mesh-watcher reads (the contract with the operator's channel).
if len(told.Names) != 1 || told.Names[0] != conditions.HeartbeatEvent {
if len(told.Names) != 2 || told.Names[1] != conditions.HeartbeatEvent {
t.Fatalf("said %v", told.Names)
}
if d.lastRunEnded().IsZero() || d.lastRun().Run != run.Run {
@@ -178,6 +187,7 @@ func TestAMachineWhoseDeclarationDoesNotComposeIsSaid(t *testing.T) {
}
}
got, err := probeDeclarations(ctx, &doctor{open: open})
linted(got)
if err != nil {
t.Fatal(err)
}
+292 -48
View File
@@ -140,6 +140,19 @@ func awaitingPush(node string, foreseen []string, since time.Time, waited time.D
// 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)
@@ -158,46 +171,206 @@ func probeResolvers(ctx context.Context, d *doctor) ([]conditions.Observation, e
if err != nil {
return nil, err
}
suffix := overlay.Suffix()
var out []conditions.Observation
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 {
var wrong []string
for _, p := range places {
if p.Address == "" {
continue
}
name := p.Name + "." + suffix
v4, rcode, err := askResolver(ctx, at, name, dnsmessage.TypeA)
switch {
case err != nil:
wrong = append(wrong, fmt.Sprintf("%s for %s: %v", "A", name, err))
continue
case rcode != dnsmessage.RCodeSuccess:
wrong = append(wrong, fmt.Sprintf("%s answered %s for %s (A)", holder, rcode, name))
case !slices.Contains(v4, p.Address):
wrong = append(wrong, fmt.Sprintf("%s answered %v for %s (A), not %s", holder, v4, name, p.Address))
for _, kind := range []dnsmessage.Type{dnsmessage.TypeA, dnsmessage.TypeAAAA} {
questions = append(questions, &resolverQuestion{holder: holder, at: at, place: p, kind: kind})
}
v6, rcode, err := askResolver(ctx, at, name, dnsmessage.TypeAAAA)
switch {
case err != nil:
wrong = append(wrong, fmt.Sprintf("%s for %s: %v", "AAAA", name, err))
case rcode != dnsmessage.RCodeSuccess:
wrong = append(wrong, fmt.Sprintf("%s answered %s for %s (AAAA), not NODATA: a musl machine "+
"takes that as no such name", holder, rcode, name))
case len(v6) > 0:
wrong = append(wrong, fmt.Sprintf("%s answered %v for %s (AAAA); the mesh has no IPv6 addresses", holder, v6, name))
}
}
if len(wrong) > 0 {
node := strings.TrimSuffix(holder, "."+suffix)
out = append(out, conditions.Observation{Scope: conditions.ScopeSeat, ID: "mesh-dns-resolver." + node,
Token: "wrong", Machine: node, Severity: conditions.Urgent,
Summary: fmt.Sprintf("the mesh's resolver on %s does not answer machine names as it must: %s",
node, wrong[0]),
Said: strings.Join(wrong, "; ")})
}
}
return sortedFound(out), nil
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.
@@ -215,13 +388,13 @@ func askResolver(ctx context.Context, at, name string, kind dnsmessage.Type) ([]
if err != nil {
return nil, 0, err
}
dialer := net.Dialer{Timeout: 3 * time.Second}
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(3 * time.Second))
_ = conn.SetDeadline(time.Now().Add(resolverWithin))
if _, err := conn.Write(packed); err != nil {
return nil, 0, err
}
@@ -308,6 +481,34 @@ func probeHolders(ctx context.Context, d *doctor) ([]conditions.Observation, err
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 {
@@ -315,7 +516,7 @@ func probeHolders(ctx context.Context, d *doctor) ([]conditions.Observation, err
continue
}
out = append(out, conditions.Observation{Scope: conditions.ScopeSeat, ID: seat + "." + node,
Token: "silent", Machine: node, Severity: conditions.Warning,
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)})
}
@@ -529,8 +730,10 @@ func probeConsumers(ctx context.Context, d *doctor) ([]conditions.Observation, e
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,
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)})
}
@@ -654,7 +857,9 @@ func probeBans(ctx context.Context, d *doctor) ([]conditions.Observation, error)
if err != nil {
return nil, err
}
owned := map[string]string{} // address → whose
// 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"
@@ -664,7 +869,7 @@ func probeBans(ctx context.Context, d *doctor) ([]conditions.Observation, error)
owned[ip.String()] = p.Name + "'s endpoint"
} else if addrs, err := net.DefaultResolver.LookupHost(ctx, host); err == nil {
for _, a := range addrs {
owned[a] = p.Name + "'s endpoint (" + host + ")"
owned[a], via[a] = p.Name+"'s endpoint", host
}
}
}
@@ -708,17 +913,26 @@ func probeBans(ctx context.Context, d *doctor) ([]conditions.Observation, error)
unasked = append(unasked, err.Error())
continue
}
var banned []string
var banned, whose []string
for _, ip := range ipLike.FindAllString(string(answer), -1) {
if whose, ours := owned[ip]; ours && !slices.Contains(banned, ip+" ("+whose+")") {
banned = append(banned, ip+" ("+whose+")")
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(banned, ", ")),
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, ", ")})
}
}
@@ -755,9 +969,12 @@ func probeStatus(ctx context.Context, d *doctor) ([]conditions.Observation, erro
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,
Summary: "status does not answer in full within ten seconds: " + why, Said: why})
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
}
@@ -910,6 +1127,14 @@ func probeWatchdogs(_ context.Context, d *doctor) ([]conditions.Observation, err
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) {
@@ -917,7 +1142,26 @@ func askSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string)
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)
}
answer, err := link.AskSeatTool(ctx, conn, seat, verb, node, map[string]any{}, 10*time.Second)
// **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
}
@@ -944,8 +1188,8 @@ func probeLease(ctx context.Context, d *doctor) ([]conditions.Observation, error
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, 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())
+287
View File
@@ -0,0 +1,287 @@
package main
import (
"context"
"fmt"
"net"
"os"
"sort"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
"golang.org/x/net/dns/dnsmessage"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/outward"
)
// **Every condition this suite raises says itself in machine names and words** (novox/hq issue 277,
// ADR 0234 §6). Every observation a keeper takes in any test of this package — every signals row
// suppressed past its bound, every probe's finding, every event's condition — is held to the operator
// channel's content rule; one whose summary or key carries an address, a domain, a path or a secret's
// shape fails the suite, naming its source. The keeper would say such a summary in words at run time;
// this is where the producer is made to say it rightly in the first place.
func TestMain(m *testing.M) {
conditions.Unsayable = func(o conditions.Observation, field string, r outward.Refusal) {
unsaid.note(o, field, r)
}
code := m.Run()
if said := unsaid.all(); len(said) > 0 {
fmt.Fprintf(os.Stderr, "FAIL: %d condition(s) raised in these tests say what the operator's channel withholds "+
"(an address, a domain, a path or a secret's shape belongs in the evidence, not the summary):\n %s\n",
len(said), strings.Join(said, "\n "))
if code == 0 {
code = 1
}
}
os.Exit(code)
}
// unsaid collects, across the suite, every finding whose words the operator's channel would withhold.
var unsaid unsayable
type unsayable struct {
mu sync.Mutex
seen map[string]bool
}
func (u *unsayable) note(o conditions.Observation, field string, r outward.Refusal) {
text := o.Summary
if field == "key" {
text = o.Key()
}
u.mu.Lock()
defer u.mu.Unlock()
if u.seen == nil {
u.seen = map[string]bool{}
}
u.seen[fmt.Sprintf("%s (source %s, kind %s): its %s carries %s — %q", o.Key(), o.Source, o.Kind, field, r.What,
text)] = true
}
func (u *unsayable) all() []string {
u.mu.Lock()
defer u.mu.Unlock()
var out []string
for s := range u.seen {
out = append(out, s)
}
sort.Strings(out)
return out
}
// linted is a producer's findings, held to the content rule as a keeper would hold them: for a test
// that reads a producer's findings without raising them.
func linted(obs []conditions.Observation) []conditions.Observation {
for _, o := range obs {
machines := append([]string{o.Machine}, o.Also...)
if r, ok := outward.Check(o.Key(), machines...); !ok {
unsaid.note(o, "key", r)
}
if r, ok := outward.Check(o.Summary, machines...); !ok {
unsaid.note(o, "summary", r)
}
}
return obs
}
// resolverStandIn is a resolver on loopback that answers as answer says; nil answers nothing.
func resolverStandIn(t *testing.T, answer func(q dnsmessage.Question, n int64) *dnsmessage.Message) (port string, asked *atomic.Int64) {
t.Helper()
conn, err := net.ListenPacket("udp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = conn.Close() })
asked = &atomic.Int64{}
go func() {
buf := make([]byte, 1500)
for {
n, from, err := conn.ReadFrom(buf)
if err != nil {
return
}
var q dnsmessage.Message
if q.Unpack(buf[:n]) != nil || len(q.Questions) == 0 {
continue
}
reply := answer(q.Questions[0], asked.Add(1))
if reply == nil {
continue
}
reply.Header.ID, reply.Header.Response, reply.Questions = q.ID, true, q.Questions
packed, _ := reply.Pack()
_, _ = conn.WriteTo(packed, from)
}
}()
_, port, _ = net.SplitHostPort(conn.LocalAddr().String())
return port, asked
}
// answersRightly answers A with the address the test's machine has, and NODATA for AAAA.
func answersRightly(q dnsmessage.Question) *dnsmessage.Message {
reply := &dnsmessage.Message{}
if q.Type == dnsmessage.TypeA {
reply.Answers = []dnsmessage.Resource{{Header: dnsmessage.ResourceHeader{Name: q.Name, Type: dnsmessage.TypeA,
Class: dnsmessage.ClassINET}, Body: &dnsmessage.AResource{A: [4]byte{10, 77, 0, 1}}}}
}
return reply
}
// quickResolvers makes D2's patience short for a test.
func quickResolvers(t *testing.T, port string) {
t.Helper()
before, within, pause := resolverPort, resolverWithin, resolverPause
resolverPort, resolverWithin, resolverPause = port, 150*time.Millisecond, 10*time.Millisecond
t.Cleanup(func() { resolverPort, resolverWithin, resolverPause = before, within, pause })
}
var anchorOnTheNetwork = []inventory.Overlay{{Name: "anchor", Address: "10.77.0.1"}}
// **D2 asks again before it says anything** (novox/hq issue 277): a resolver that misses a question on
// a loaded machine and answers the next try is a resolver that answers.
func TestAResolverThatMissesOneTryAndAnswersTheNextIsWell(t *testing.T) {
port, asked := resolverStandIn(t, func(q dnsmessage.Question, n int64) *dnsmessage.Message {
if n <= 2 {
return nil // the first question of each kind lost, as under a push and a build starting
}
return answersRightly(q)
})
quickResolvers(t, port)
got := askEveryResolver(t.Context(), map[string]string{"anchor.internal": "127.0.0.1"}, anchorOnTheNetwork, "internal")
if len(got) != 0 {
t.Fatalf("a resolver that answered its second try was said: %+v", got)
}
if asked.Load() < 3 {
t.Fatalf("asked %d times", asked.Load())
}
}
// **A resolver that answers nothing is held for the next run, and said in machine names**: the address
// it was asked at and the socket's words are the evidence, never the summary the operator reads.
func TestAResolverThatAnswersNothingIsHeldAndSaidInMachineNames(t *testing.T) {
port, asked := resolverStandIn(t, func(dnsmessage.Question, int64) *dnsmessage.Message { return nil })
quickResolvers(t, port)
got := askEveryResolver(t.Context(), map[string]string{"anchor.internal": "127.0.0.1"}, anchorOnTheNetwork, "internal")
if len(got) != 1 {
t.Fatalf("%+v", got)
}
o := got[0]
if !o.Confirm || o.Key() != "seat.mesh-dns-resolver.anchor.wrong" || o.Severity != conditions.Urgent {
t.Fatalf("not held for a second look, or not the resolver's condition: %+v", o)
}
if r, ok := outward.Check(o.Summary, "anchor"); !ok {
t.Fatalf("the summary carries %s: %q", r, o.Summary)
}
if !strings.Contains(o.Summary, "anchor's") || !strings.Contains(o.Said, "127.0.0.1") ||
!strings.Contains(o.Said, "anchor.internal") {
t.Fatalf("summary %q, evidence %q", o.Summary, o.Said)
}
if want := int64(2 * resolverTries); asked.Load() != want {
t.Fatalf("asked %d times, want %d: each question %d times", asked.Load(), want, resolverTries)
}
}
// **A resolver that answers wrongly is said at once**, in machine names: an answer is not the absence
// of one (issue 262's NXDOMAIN for IPv6).
func TestAResolverAnsweringWronglyIsSaidAtOnceInMachineNames(t *testing.T) {
port, _ := resolverStandIn(t, func(q dnsmessage.Question, _ int64) *dnsmessage.Message {
if q.Type == dnsmessage.TypeAAAA {
return &dnsmessage.Message{Header: dnsmessage.Header{RCode: dnsmessage.RCodeNameError}}
}
return answersRightly(q)
})
quickResolvers(t, port)
got := askEveryResolver(t.Context(), map[string]string{"anchor.internal": "127.0.0.1"}, anchorOnTheNetwork, "internal")
if len(got) != 1 || got[0].Confirm {
t.Fatalf("%+v", got)
}
if r, ok := outward.Check(got[0].Summary, "anchor"); !ok || !strings.Contains(got[0].Summary, "NXDOMAIN") ||
!strings.Contains(got[0].Summary, "anchor's IPv6 address") {
t.Fatalf("summary %q (%v)", got[0].Summary, r)
}
}
// **The self-check raises a finding one look can be wrong about on the second run in a row**, keeps it
// open while it is seen, and clears it when it is not — and a probe that cannot run is said as a
// condition only when the next run cannot run it either. Neither is ever a pass in the verdict.
func TestASingleLookIsHeldAndTheSecondInARowRaises(t *testing.T) {
var unanswered, broken atomic.Bool
held := conditions.Observation{Scope: conditions.ScopeSeat, ID: "mesh-dns-resolver.anchor", Token: "wrong",
Machine: "anchor", Severity: conditions.Urgent, Confirm: true,
Summary: "the mesh's resolver on anchor does not answer", Said: "no answer from 10.77.0.1"}
withProbes(t,
probe{ID: "P1", Asserts: "asks over the network", Kind: "resolver-wrong", Phase: 1,
run: func(context.Context, *doctor) ([]conditions.Observation, error) {
if unanswered.Load() {
return []conditions.Observation{held}, nil
}
return nil, nil
}},
probe{ID: "P2", Asserts: "sometimes cannot run", Kind: "x", Phase: 1,
run: func(context.Context, *doctor) ([]conditions.Observation, error) {
if broken.Load() {
return nil, fmt.Errorf("a question timed out")
}
return nil, nil
}},
)
store := conditions.NewInMemory()
k := conditions.NewKeeper(t.Context(), conditions.Options{Store: store, History: store})
defer k.Close(context.Background())
d := &doctor{keeper: k, host: "anchor"}
keys := func() []string {
open, err := k.Open(t.Context())
if err != nil {
t.Fatal(err)
}
var out []string
for _, c := range open {
out = append(out, c.Key)
}
sort.Strings(out)
return out
}
unanswered.Store(true)
broken.Store(true)
run := d.runOnce(t.Context(), "a test")
if open := keys(); len(open) != 0 {
t.Fatalf("one look raised %v", open)
}
if run.Probes[0].Verdict != verdictPass || len(run.Probes[0].Unconfirmed) != 1 || run.Probes[1].Verdict != verdictFailedToRun {
t.Fatalf("the first run's verdict hides what it saw: %+v", run.Probes)
}
// Seen twice in a row: raised, and so is the probe that could not run twice.
run = d.runOnce(t.Context(), "a test")
if open := keys(); strings.Join(open, " ") != "probe.P2.failed seat.mesh-dns-resolver.anchor.wrong" {
t.Fatalf("two looks in a row left open %v", open)
}
if run.Probes[0].Verdict != verdictFail {
t.Fatalf("%+v", run.Probes[0])
}
// Open, and seen again: kept, never cleared and raised again.
d.runOnce(t.Context(), "a test")
if c, open, _ := k.Get(t.Context(), "seat.mesh-dns-resolver.anchor.wrong"); !open || c.Count != 1 || c.Observations != 2 {
t.Fatalf("the open condition was not kept as it was: %+v", c)
}
// Answered again: cleared. Then one look alone raises nothing again.
unanswered.Store(false)
broken.Store(false)
d.runOnce(t.Context(), "a test")
if open := keys(); len(open) != 0 {
t.Fatalf("a passing run left open %v", open)
}
unanswered.Store(true)
d.runOnce(t.Context(), "a test")
if open := keys(); len(open) != 0 {
t.Fatalf("one look after a pass raised %v", open)
}
}
+4 -3
View File
@@ -527,7 +527,7 @@ func watchLease(f *signalFacts) []conditions.Observation {
Summary: "the controller serves WITHOUT the lease: nothing keeps a second controller from acting " +
"beside it, and its declarations carry no epoch. A bus whose user list is older than this " +
"controller does not grant it the lease's bucket: a push of the machine holding the bus sends " +
"the list that does, and the controller takes the lease within five seconds — " + l.unleased,
"the list that does, and the controller takes the lease within five seconds",
Said: l.unleased})
} else if l.held && !l.renewed.IsZero() && f.now.Sub(l.renewed) > leaseBound {
out = append(out, conditions.Observation{Scope: conditions.ScopeCore, ID: "controller.lease", Token: "late",
@@ -539,8 +539,9 @@ func watchLease(f *signalFacts) []conditions.Observation {
if !l.reset.IsZero() && f.now.Sub(l.reset) <= advisoryQuiet {
out = append(out, conditions.Observation{Scope: conditions.ScopeCore, ID: "controller.lease", Token: "reset",
Kind: "lease-lost", Machine: f.host, Severity: conditions.Urgent,
Summary: "the controller lease's bucket was raised again from nothing: " + l.resetSaid,
Said: l.resetSaid})
Summary: "the controller lease's bucket was raised again from nothing: whatever held the lease before " +
"does not know it lost it",
Said: l.resetSaid})
}
var lost []string
newest := inventory.Epoch{}
+6
View File
@@ -300,6 +300,12 @@ func TestABlindWatchdogSaysSoAndClearsNothing(t *testing.T) {
w.see(t.Context(), silent)
blind := calm(now)
blind.machines, blind.machinesErr = nil, errors.New("the store is away")
// Blind for one tick is a read that did not answer in time; for two in a row, a watchdog that cannot
// see (novox/hq issue 277).
w.see(t.Context(), blind)
if open, _ := k.Open(t.Context()); len(open) != 1 {
t.Fatalf("one blind tick raised %+v", open)
}
w.see(t.Context(), blind)
open, err := k.Open(t.Context())
if err != nil {
+10 -3
View File
@@ -164,6 +164,9 @@ type watchdogs struct {
// standing by hears no heartbeat and would call every machine silent.
acting func() bool
// confirm holds back a row blind for one tick: a read that timed out once on a loaded store.
confirm confirming
mu sync.Mutex
ticked time.Time
last *signalFacts
@@ -226,7 +229,9 @@ func (w *watchdogs) see(running context.Context, f *signalFacts) {
problems = append(problems, row.Row+": "+err.Error())
}
}
// The rows that could not see, said; the ones that see again, cleared.
// The rows that could not see, said once they could not two ticks in a row; the ones that see
// again, cleared.
blind, _ = w.confirm.pass(running, w.keeper, sourceWatchdogs, blind)
if err := w.keeper.Reconcile(running, sourceWatchdogs, blind); err != nil {
problems = append(problems, err.Error())
}
@@ -259,10 +264,12 @@ func (w *watchdogs) see(running context.Context, f *signalFacts) {
// sourceWatchdogs is what raises a blind row's condition.
const sourceWatchdogs = "watchdogs"
// blindRow is a row whose facts could not be gathered, as a condition of its own.
// blindRow is a row whose facts could not be gathered, as a condition of its own: raised when the next
// tick cannot gather them either (confirm.go), since one read that did not answer in time is not a
// watchdog gone blind.
func blindRow(row signalRow, err error) conditions.Observation {
return conditions.Observation{Scope: conditions.ScopeProbe, ID: row.Row, Kind: "probe-failed", Token: "failed",
Severity: conditions.Warning,
Severity: conditions.Warning, Confirm: true,
Summary: fmt.Sprintf("the watchdog of %s (%s) cannot see: what it reads could not be read, so nothing "+
"it would raise can be — and nothing it raised before is cleared", row.Row, row.Signal),
Said: firstLine(err.Error())}