Merge pull request 'Look twice before saying a probe failed, and say conditions in machine names (hq issue 277)' (#91) from fix/probes-look-twice-and-say-no-addresses into main
mesh/delivery delivered
mesh/delivery delivered
This commit was merged in pull request #91.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
@@ -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())
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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{}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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())}
|
||||
|
||||
@@ -170,10 +170,19 @@ type Observation struct {
|
||||
Also []string
|
||||
Severity Severity
|
||||
Summary string
|
||||
// Said is this observation's evidence, in the mesh's words; Summary when empty.
|
||||
// Said is this observation's evidence, in the mesh's words; Summary when empty. **Detail goes
|
||||
// here, never in Summary**: an address, a socket's error, a path or a name with its domain is
|
||||
// kept in the condition's evidence, which stays inside the mesh. The summary leaves it — to the
|
||||
// operator's channel, whose content rule withholds a message that carries any of them (ADR 0234
|
||||
// §6), and names machines in words.
|
||||
Said string
|
||||
Source string
|
||||
Resolver string
|
||||
// Confirm says a single look can be wrong about this finding — a question over the network that
|
||||
// went unanswered, a time measured once on a loaded machine. The keeper does not read it: the
|
||||
// source that looks again does, and raises it only when the next look sees it too, or while it is
|
||||
// already open (novox/hq issue 277).
|
||||
Confirm bool
|
||||
}
|
||||
|
||||
// Key is where the observation's condition is kept: `<scope>.<id>.<kind>`, so the same fault said
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
package conditions
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/outward"
|
||||
)
|
||||
|
||||
// **A summary the operator's channel would withhold is said in words, and kept whole in the evidence**
|
||||
// (novox/hq issue 277, ADR 0234 §6): the condition that carried a resolver's address reached the
|
||||
// operator as "its words are withheld" instead of the alert.
|
||||
func TestASummaryCarryingAnAddressIsSaidInWordsAndKeptInTheEvidence(t *testing.T) {
|
||||
k, _, _, _ := keeper(t)
|
||||
before := Unsayable
|
||||
var told []string
|
||||
Unsayable = func(o Observation, field string, r outward.Refusal) { told = append(told, field+": "+r.What) }
|
||||
t.Cleanup(func() { Unsayable = before })
|
||||
|
||||
raw := "AAAA for anchor.internal: no answer from 10.77.0.1: read udp 10.77.0.3:41234->10.77.0.1:53: i/o timeout"
|
||||
c, err := k.Observe(t.Context(), Observation{Scope: ScopeSeat, ID: "mesh-dns-resolver.anchor", Token: "wrong",
|
||||
Kind: "resolver-wrong", Machine: "anchor", Severity: Urgent, Source: "D2",
|
||||
Summary: "the mesh's resolver on anchor does not answer: " + raw})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if r, ok := outward.Check(c.Summary, "anchor"); !ok {
|
||||
t.Fatalf("the summary kept is still withheld (%s): %q", r, c.Summary)
|
||||
}
|
||||
if !strings.HasPrefix(c.Summary, "the mesh's resolver on anchor does not answer: ") {
|
||||
t.Errorf("the summary's words were not kept: %q", c.Summary)
|
||||
}
|
||||
if !strings.Contains(c.Evidence[0].Said, raw) {
|
||||
t.Errorf("the evidence lost what the summary carried: %q", c.Evidence[0].Said)
|
||||
}
|
||||
if len(told) != 1 || !strings.HasPrefix(told[0], "summary") {
|
||||
t.Errorf("the producer was not named to the test: %v", told)
|
||||
}
|
||||
|
||||
// A summary that may leave is kept as it is said, machine names and all.
|
||||
ok, err := k.Observe(t.Context(), Observation{Scope: ScopeMachine, ID: "anchor", Kind: "silent", Machine: "anchor",
|
||||
Severity: Warning, Source: "S1", Summary: "anchor has not been heard from since 12:00 UTC (bound 3m0s)"})
|
||||
if err != nil || ok.Summary != "anchor has not been heard from since 12:00 UTC (bound 3m0s)" {
|
||||
t.Fatalf("%q %v", ok.Summary, err)
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,8 @@ import (
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/outward"
|
||||
)
|
||||
|
||||
// Backend is where the open conditions are kept: one value per key, written by compare-and-set.
|
||||
@@ -70,6 +72,8 @@ type Keeper struct {
|
||||
closing sync.Once
|
||||
// Unsaid counts the transitions given up on, for the self-check to say.
|
||||
unsaid int
|
||||
// reworded are the keys whose summary was said in words for the operator's channel, said once each.
|
||||
reworded map[string]bool
|
||||
}
|
||||
|
||||
type clearing struct {
|
||||
@@ -105,7 +109,7 @@ var TellFor = 10 * time.Minute
|
||||
func NewKeeper(ctx context.Context, o Options) *Keeper {
|
||||
k := &Keeper{store: o.Store, history: o.History, teller: o.Teller, now: o.Now, say: o.Say, changed: o.Changed,
|
||||
epoch: o.Epoch,
|
||||
cleared: map[string]clearing{}, out: make(chan Event, 1024), drained: make(chan struct{})}
|
||||
cleared: map[string]clearing{}, reworded: map[string]bool{}, out: make(chan Event, 1024), drained: make(chan struct{})}
|
||||
if k.now == nil {
|
||||
k.now = time.Now
|
||||
}
|
||||
@@ -170,6 +174,7 @@ func (k *Keeper) Observe(ctx context.Context, o Observation) (Condition, error)
|
||||
if err := o.check(); err != nil {
|
||||
return Condition{}, err
|
||||
}
|
||||
o = k.sayable(o)
|
||||
key := o.Key()
|
||||
for i := 0; i < tries; i++ {
|
||||
now := k.now().UTC()
|
||||
@@ -627,3 +632,49 @@ func orSelf(resolver string) string {
|
||||
}
|
||||
return resolver
|
||||
}
|
||||
|
||||
// Unsayable is told of every observation whose words the operator's channel would withhold (novox/hq
|
||||
// ADR 0234 §6): its summary or its key carries an address, a domain, a path or a secret's shape. The
|
||||
// keeper says such a summary in words itself, and keeps what it said whole in the evidence; a test
|
||||
// suite sets this to fail the producer, which is where the summary should have been said rightly.
|
||||
var Unsayable func(o Observation, field string, r outward.Refusal)
|
||||
|
||||
// sayable is an observation whose summary may leave the mesh: **a summary says things in machine names
|
||||
// and the mesh's words** (novox/hq issue 277). One that carries what the channel withholds — a raw error
|
||||
// with an address in it, a path — is said in words here, as the last stand before the operator would
|
||||
// read "its words are withheld" instead of the alert; what it carried goes to the evidence, which
|
||||
// stays inside the mesh.
|
||||
func (k *Keeper) sayable(o Observation) Observation {
|
||||
machines := append([]string{o.Machine}, o.Also...)
|
||||
if r, ok := outward.Check(o.Key(), machines...); !ok && Unsayable != nil {
|
||||
Unsayable(o, "key", r)
|
||||
}
|
||||
r, ok := outward.Check(o.Summary, machines...)
|
||||
if ok {
|
||||
return o
|
||||
}
|
||||
if Unsayable != nil {
|
||||
Unsayable(o, "summary", r)
|
||||
}
|
||||
whole := o.Summary
|
||||
o.Summary = outward.Scrub(whole, fmt.Sprintf("a %s condition about %s %s: what it says is kept in its evidence "+
|
||||
"(`conditions show`)", o.Kind, o.Scope, o.ID), machines...)
|
||||
switch {
|
||||
case o.Said == "":
|
||||
o.Said = whole
|
||||
case !strings.Contains(o.Said, whole):
|
||||
o.Said += " — as raised: " + whole
|
||||
}
|
||||
k.mu.Lock()
|
||||
if k.reworded == nil {
|
||||
k.reworded = map[string]bool{}
|
||||
}
|
||||
first := !k.reworded[o.Key()]
|
||||
k.reworded[o.Key()] = true
|
||||
k.mu.Unlock()
|
||||
if first {
|
||||
k.say("the condition %s's summary carried %s, which the operator's channel withholds: said in words, "+
|
||||
"and kept whole in its evidence — its source should say it so", o.Key(), r.What)
|
||||
}
|
||||
return o
|
||||
}
|
||||
|
||||
@@ -0,0 +1,292 @@
|
||||
// Package outward is the operator channel's content rule, as the controller holds its own words to it
|
||||
// (novox/hq ADR 0234 §6, to-be 45 §5).
|
||||
//
|
||||
// **What may leave the mesh is machine names and words.** The messenger refuses a message that carries
|
||||
// an address, a domain, a path or anything shaped like a secret, and withholds its words — so a
|
||||
// condition whose summary carried a resolver's address reached the operator as "this message carried an
|
||||
// IPv4 address, so its words are withheld" instead of the alert (2026-10-06). The rule is the
|
||||
// messenger's, and stays the messenger's: this package mirrors its patterns so that the controller can
|
||||
// hold a condition's summary to it where the summary is made, and a test can fail a summary that would
|
||||
// be withheld. Detail — an address, a socket's error, a path — belongs in a condition's evidence, which
|
||||
// stays inside the mesh.
|
||||
//
|
||||
// Mirrored, not imported: the messenger is a module of the catalogue with its own module path, and a
|
||||
// pattern changed there is changed here (the table in outward_test.go names the shapes both refuse).
|
||||
// Where the two differ, this one may only be the stricter: what passes here passes there.
|
||||
package outward
|
||||
|
||||
import (
|
||||
"math"
|
||||
"regexp"
|
||||
"strings"
|
||||
"unicode"
|
||||
)
|
||||
|
||||
// Refusal says why a text may not leave: the class of what it carried, never the text itself.
|
||||
type Refusal struct {
|
||||
Class string // address, path or secret
|
||||
What string // a few words: "an IPv4 address", "a URL", …
|
||||
}
|
||||
|
||||
func (r Refusal) String() string { return r.Class + " (" + r.What + ")" }
|
||||
|
||||
// The messenger's patterns (mesh-catalog modules/messenger content.go), one for one.
|
||||
var (
|
||||
reURL = regexp.MustCompile(`(?i)\b[a-z][a-z0-9+.-]*://`)
|
||||
reEmail = regexp.MustCompile(`[A-Za-z0-9._%+-]+@[A-Za-z0-9-]+(\.[A-Za-z0-9-]+)*\.[A-Za-z]{2,}`)
|
||||
reIPv4 = regexp.MustCompile(`\b\d{1,3}(\.\d{1,3}){3}\b`)
|
||||
reIPv6 = regexp.MustCompile(`(?i)(^|[^0-9a-z:])(([0-9a-f]{1,4}:){4,7}[0-9a-f]{1,4}|([0-9a-f]{1,4}:)*[0-9a-f]{0,4}::([0-9a-f]{1,4}:)*[0-9a-f]{0,4})([^0-9a-z:]|$)`)
|
||||
reMAC = regexp.MustCompile(`(?i)\b([0-9a-f]{2}[:-]){5}[0-9a-f]{2}\b`)
|
||||
rePEM = regexp.MustCompile(`-----BEGIN [A-Z ]+-----`)
|
||||
reJWT = regexp.MustCompile(`\beyJ[A-Za-z0-9_-]{8,}\.[A-Za-z0-9_-]{8,}`)
|
||||
reBotToken = regexp.MustCompile(`\b\d{6,}:[A-Za-z0-9_-]{30,}`)
|
||||
reKnown = regexp.MustCompile(`\b(gh[pousr]_[A-Za-z0-9]{20,}|glpat-[A-Za-z0-9_-]{16,}|sk-[A-Za-z0-9_-]{16,}|xox[abprs]-[A-Za-z0-9-]{10,}|AKIA[0-9A-Z]{16})`)
|
||||
reAssigned = regexp.MustCompile(`(?i)\b(password|passwd|passphrase|secret|token|api[_-]?key|apikey|credential|private[_-]?key)\s*[=:]\s*\S`)
|
||||
reHex = regexp.MustCompile(`(?i)\b[0-9a-f]{32,}\b`)
|
||||
reRun = regexp.MustCompile(`[A-Za-z0-9+/=_]{20,}`)
|
||||
reWinPath = regexp.MustCompile(`(?i)\b[a-z]:\\`)
|
||||
)
|
||||
|
||||
// topLevel are names that end a host name, as the messenger reads them.
|
||||
var topLevel = map[string]bool{}
|
||||
|
||||
func init() {
|
||||
for _, t := range strings.Fields(`com net org edu gov mil int io dev app cloud ai co me info biz xyz
|
||||
site online tech page link
|
||||
be nl de fr uk lu eu ch at it es pt se no dk fi pl cz us ca au nz jp cn ru in br ie
|
||||
internal lan home local localdomain corp intranet private arpa test example invalid localhost`) {
|
||||
topLevel[t] = true
|
||||
}
|
||||
}
|
||||
|
||||
// wordSeparators split a text into the words the messenger reads one by one.
|
||||
const wordSeparators = "\"'`()[]{}<>,;|"
|
||||
|
||||
func isSeparator(r rune) bool { return unicode.IsSpace(r) || strings.ContainsRune(wordSeparators, r) }
|
||||
|
||||
// Check says whether a text may leave the mesh, and when not, why. **The mesh's own machine names may
|
||||
// appear** (ADR 0234 §6): a word that is one of machines is read as a name, whatever its shape —
|
||||
// never as a host name, a path or a random string. A machine name joined to a domain is a domain.
|
||||
func Check(text string, machines ...string) (Refusal, bool) {
|
||||
text = withoutMachines(text, machines)
|
||||
switch {
|
||||
case reURL.MatchString(text):
|
||||
return Refusal{"address", "a URL"}, false
|
||||
case reEmail.MatchString(text):
|
||||
return Refusal{"address", "a mail address"}, false
|
||||
case reIPv4.MatchString(text):
|
||||
return Refusal{"address", "an IPv4 address"}, false
|
||||
case reMAC.MatchString(text):
|
||||
return Refusal{"address", "a hardware address"}, false
|
||||
case reIPv6.MatchString(text):
|
||||
return Refusal{"address", "an IPv6 address"}, false
|
||||
case rePEM.MatchString(text):
|
||||
return Refusal{"secret", "a key block"}, false
|
||||
case reJWT.MatchString(text):
|
||||
return Refusal{"secret", "a signed token"}, false
|
||||
case reBotToken.MatchString(text):
|
||||
return Refusal{"secret", "a bot token"}, false
|
||||
case reKnown.MatchString(text):
|
||||
return Refusal{"secret", "a known token shape"}, false
|
||||
case reAssigned.MatchString(text):
|
||||
return Refusal{"secret", "a value given to a secret's name"}, false
|
||||
case reHex.MatchString(text):
|
||||
return Refusal{"secret", "a long hexadecimal string"}, false
|
||||
case reWinPath.MatchString(text):
|
||||
return Refusal{"path", "a drive path"}, false
|
||||
}
|
||||
for _, run := range reRun.FindAllString(text, -1) {
|
||||
if looksRandom(run) {
|
||||
return Refusal{"secret", "a long random-looking string"}, false
|
||||
}
|
||||
}
|
||||
for _, word := range strings.FieldsFunc(text, isSeparator) {
|
||||
w := strings.TrimRight(word, ".:!?")
|
||||
if isPath(w) {
|
||||
return Refusal{"path", "a file path"}, false
|
||||
}
|
||||
if isHostName(w) {
|
||||
return Refusal{"address", "a host name"}, false
|
||||
}
|
||||
}
|
||||
return Refusal{}, true
|
||||
}
|
||||
|
||||
// What Scrub says in words, beyond the messenger's patterns: an address with its port and what a
|
||||
// socket says around it ("udp 192.0.2.1:53"), a bracketed IPv6 address, a key block whole, a secret's
|
||||
// value, a drive path whole.
|
||||
var (
|
||||
scrubURL = regexp.MustCompile(`(?i)\b[a-z][a-z0-9+.-]*://\S*`)
|
||||
scrubIPv4 = regexp.MustCompile(`\b\d{1,3}(\.\d{1,3}){3}(:\d+)?\b`)
|
||||
scrubIPv6 = regexp.MustCompile(`\[[0-9a-fA-F:.%]*:[0-9a-fA-F:.%]*\](:\d+)?`)
|
||||
scrubPEM = regexp.MustCompile(`(?s)-----BEGIN [A-Z ]+-----.*?(-----END [A-Z ]+-----|$)`)
|
||||
scrubAssigned = regexp.MustCompile(`(?i)\b(password|passwd|passphrase|secret|token|api[_-]?key|apikey|credential|private[_-]?key)\s*[=:]\s*\S+`)
|
||||
scrubWinPath = regexp.MustCompile(`(?i)\b[a-z]:\\\S*`)
|
||||
)
|
||||
|
||||
// Scrub is a text with everything Check refuses said in words instead: an address as "an address", a
|
||||
// path as "a path", a secret's shape as "(withheld)". The text's own words, and the machine names, are
|
||||
// kept. What Scrub cannot make pass is replaced whole by fallback — never sent as it was.
|
||||
func Scrub(text, fallback string, machines ...string) string {
|
||||
if _, ok := Check(text, machines...); ok {
|
||||
return text
|
||||
}
|
||||
out := scrubURL.ReplaceAllString(text, "an address")
|
||||
out = reEmail.ReplaceAllString(out, "an address")
|
||||
out = scrubIPv4.ReplaceAllString(out, "an address")
|
||||
out = reMAC.ReplaceAllString(out, "a hardware address")
|
||||
out = scrubIPv6.ReplaceAllString(out, "an address")
|
||||
for i := 0; i < 4 && reIPv6.MatchString(out); i++ {
|
||||
out = reIPv6.ReplaceAllString(out, "${1}an address${6}")
|
||||
}
|
||||
out = scrubPEM.ReplaceAllString(out, "(withheld)")
|
||||
for _, re := range []*regexp.Regexp{reJWT, reBotToken, reKnown, reHex} {
|
||||
out = re.ReplaceAllString(out, "(withheld)")
|
||||
}
|
||||
out = scrubAssigned.ReplaceAllString(out, "${1} (withheld)")
|
||||
out = scrubWinPath.ReplaceAllString(out, "a path")
|
||||
out = reRun.ReplaceAllStringFunc(out, func(run string) string {
|
||||
if looksRandom(run) {
|
||||
return "(withheld)"
|
||||
}
|
||||
return run
|
||||
})
|
||||
out = eachWord(out, func(w string) string {
|
||||
switch {
|
||||
case isMachine(w, machines):
|
||||
return w
|
||||
case isPath(w):
|
||||
return "a path"
|
||||
case isHostName(w):
|
||||
return "a host name"
|
||||
}
|
||||
return w
|
||||
})
|
||||
if _, ok := Check(out, machines...); ok {
|
||||
return out
|
||||
}
|
||||
return fallback
|
||||
}
|
||||
|
||||
// withoutMachines is a text with every machine name that stands as a word of its own read as a plain
|
||||
// word, so the patterns do not read a name as anything else.
|
||||
func withoutMachines(text string, machines []string) string {
|
||||
if len(machines) == 0 {
|
||||
return text
|
||||
}
|
||||
return eachWord(text, func(w string) string {
|
||||
if isMachine(w, machines) {
|
||||
return "machine"
|
||||
}
|
||||
return w
|
||||
})
|
||||
}
|
||||
|
||||
// eachWord is a text with each word, as the messenger splits and trims it, put through say; what
|
||||
// separates the words, and the punctuation a word ends in, are kept.
|
||||
func eachWord(text string, say func(w string) string) string {
|
||||
var b strings.Builder
|
||||
start := -1
|
||||
flush := func(end int) {
|
||||
if start < 0 {
|
||||
return
|
||||
}
|
||||
word := text[start:end]
|
||||
w := strings.TrimRight(word, ".:!?")
|
||||
if w == "" {
|
||||
b.WriteString(word)
|
||||
} else {
|
||||
b.WriteString(say(w) + word[len(w):])
|
||||
}
|
||||
start = -1
|
||||
}
|
||||
for i, r := range text {
|
||||
if isSeparator(r) {
|
||||
flush(i)
|
||||
b.WriteRune(r)
|
||||
continue
|
||||
}
|
||||
if start < 0 {
|
||||
start = i
|
||||
}
|
||||
}
|
||||
flush(len(text))
|
||||
return b.String()
|
||||
}
|
||||
|
||||
func isMachine(w string, machines []string) bool {
|
||||
for _, m := range machines {
|
||||
if m != "" && strings.EqualFold(w, m) {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// isPath: absolute, home-relative or dot-relative, or two separators deep. A mesh address names one
|
||||
// machine and one tool (`ace/postgres.query`) and has one; a ratio ("3/4") has digits only.
|
||||
func isPath(w string) bool {
|
||||
if w == "" {
|
||||
return false
|
||||
}
|
||||
if strings.HasPrefix(w, "/") && len(w) > 1 {
|
||||
return true
|
||||
}
|
||||
for _, p := range []string{"~/", "./", "../", "$HOME", "${"} {
|
||||
if strings.HasPrefix(w, p) {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return strings.Count(w, "/") >= 2 || strings.Contains(w, "\\")
|
||||
}
|
||||
|
||||
// isHostName: two names or more, the last a top-level one. A condition key's last name is its kind
|
||||
// (`machine.ace.silent`), which none of these is.
|
||||
func isHostName(w string) bool {
|
||||
w = strings.ToLower(w)
|
||||
if w == "localhost" {
|
||||
return true
|
||||
}
|
||||
parts := strings.Split(w, ".")
|
||||
if len(parts) < 2 {
|
||||
return false
|
||||
}
|
||||
for _, p := range parts {
|
||||
if p == "" {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return topLevel[parts[len(parts)-1]]
|
||||
}
|
||||
|
||||
// looksRandom: letters and digits mixed, and the characters spread as a random string's are.
|
||||
func looksRandom(s string) bool {
|
||||
var letters, digits int
|
||||
counts := map[rune]int{}
|
||||
for _, r := range s {
|
||||
counts[r]++
|
||||
switch {
|
||||
case unicode.IsLetter(r):
|
||||
letters++
|
||||
case unicode.IsDigit(r):
|
||||
digits++
|
||||
}
|
||||
}
|
||||
if letters == 0 || digits == 0 {
|
||||
return letters > 0 && hasUpperAndLower(s) && entropy(counts, len(s)) >= 4.0
|
||||
}
|
||||
return entropy(counts, len(s)) >= 3.3
|
||||
}
|
||||
|
||||
func hasUpperAndLower(s string) bool {
|
||||
return strings.IndexFunc(s, unicode.IsUpper) >= 0 && strings.IndexFunc(s, unicode.IsLower) >= 0
|
||||
}
|
||||
|
||||
func entropy(counts map[rune]int, n int) float64 {
|
||||
e := 0.0
|
||||
for _, c := range counts {
|
||||
p := float64(c) / float64(n)
|
||||
e -= p * math.Log2(p)
|
||||
}
|
||||
return e
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
package outward
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// machines are the names a mesh in these tests has.
|
||||
var machines = []string{"anchor", "laptop", "g14", "home-server"}
|
||||
|
||||
// **The content rule's table** (novox/hq ADR 0234, "machine names pass the content rule; domains,
|
||||
// addresses, paths and secrets do not"): the shapes the messenger refuses are refused here, and the
|
||||
// words a condition is made of pass — the mesh's machine names among them.
|
||||
func TestMachineNamesPassAndAddressesDomainsPathsAndSecretsDoNot(t *testing.T) {
|
||||
pass := []string{
|
||||
"the mesh's resolver on anchor has not answered for two runs of the self-check",
|
||||
"laptop has not been heard from since 2026-10-06 15:02 UTC (bound 3m0s)",
|
||||
"g14's node-engine would refuse its declaration whole",
|
||||
"home-server runs the node-engine 3f2a9c1b7d0e, the mesh holds 5e6f7a8b9c0d",
|
||||
"seat.mesh-dns-resolver.anchor.wrong",
|
||||
"anchor/postgres.query answers",
|
||||
"3/4 of the consumers are behind; 1200 message(s) behind its stream's head",
|
||||
"the plan for novox/app c0ffee00 has been at tier 1 of 2",
|
||||
"mesh-controller.conditions key=machine.anchor.silent",
|
||||
"backed-up 2 days ago: tank/data shrank to 1.2 GiB",
|
||||
}
|
||||
for _, text := range pass {
|
||||
if r, ok := Check(text, machines...); !ok {
|
||||
t.Errorf("refused as %s: %q", r, text)
|
||||
}
|
||||
}
|
||||
refuse := map[string]string{
|
||||
"AAAA for anchor: no answer from 10.77.0.1: read udp 10.77.0.3:41234->10.77.0.1:53: i/o timeout": "address",
|
||||
"the resolver at [fd00::1]:53 did not answer": "address",
|
||||
"fe80::1ff:fe23:4567:890a is banned": "address",
|
||||
"anchor.internal answered NXDOMAIN": "address",
|
||||
"the store at https://artifacts.example.org/v2 is away": "address",
|
||||
"write to jochen@example.org": "address",
|
||||
"aa:bb:cc:dd:ee:ff": "address",
|
||||
"data at /srv/postgres/data is gone": "path",
|
||||
"kept at ~/backups": "path",
|
||||
`kept at C:\backups`: "path",
|
||||
"password=hunter2": "secret",
|
||||
"the token glpat-abcdefghijklmnop12 leaked": "secret",
|
||||
"0123456789abcdef0123456789abcdef01": "secret",
|
||||
}
|
||||
for text, class := range refuse {
|
||||
r, ok := Check(text, machines...)
|
||||
if ok || r.Class != class {
|
||||
t.Errorf("%q: want refused as %s, got %v %v", text, class, r, ok)
|
||||
}
|
||||
}
|
||||
// A machine name joined to a domain is a domain, the machine's name notwithstanding.
|
||||
if _, ok := Check("anchor.lan answered", machines...); ok {
|
||||
t.Error("a machine's name with a domain passed")
|
||||
}
|
||||
// And a machine whose name has a shape the rule would otherwise refuse still passes as a name.
|
||||
if _, ok := Check("node.internal is silent", "node.internal"); !ok {
|
||||
t.Error("a machine name was refused")
|
||||
}
|
||||
}
|
||||
|
||||
// **Scrub says what Check refuses in words, and what it gives back always passes** — or is the
|
||||
// fallback, never the text as it was.
|
||||
func TestScrubAlwaysGivesWhatMayLeave(t *testing.T) {
|
||||
for _, text := range []string{
|
||||
"AAAA for anchor.internal: no answer from 10.77.0.1: read udp 10.77.0.3:41234->10.77.0.1:53: i/o timeout",
|
||||
"dial tcp [fd00::1]:4222: connect: connection refused",
|
||||
"open /var/lib/mesh/data/x.db: no such file or directory",
|
||||
"GET https://artifacts.example.org/v2/blobs/sha256:0123456789abcdef0123456789abcdef0123456789abcdef: 404",
|
||||
"password=hunter2 and a key -----BEGIN PRIVATE KEY-----\nAAAA\n-----END PRIVATE KEY-----",
|
||||
"laptop's ban list holds 192.0.2.7 (anchor's endpoint)",
|
||||
} {
|
||||
got := Scrub(text, "FALLBACK", machines...)
|
||||
if r, ok := Check(got, machines...); !ok {
|
||||
t.Errorf("Scrub(%q) = %q still refused as %s", text, got, r)
|
||||
}
|
||||
if got == "FALLBACK" {
|
||||
t.Errorf("Scrub(%q) fell back where words could be kept", text)
|
||||
}
|
||||
}
|
||||
if got := Scrub("laptop's ban list holds 192.0.2.7", "", machines...); !strings.HasPrefix(got, "laptop's ban list holds an address") {
|
||||
t.Errorf("the machine's name was not kept: %q", got)
|
||||
}
|
||||
if got := Scrub("nothing wrong with anchor", "", machines...); got != "nothing wrong with anchor" {
|
||||
t.Errorf("a text that passes was changed: %q", got)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user