Files
mesh-controller/cmd/mesh-controller/signals.go
T
jochen 1a4305213d
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
mesh/delivery-group group feat/plain-notifications delivered: every member is delivered
Open every explanation with what the operator needs to do, and offer the answers
The operator could not tell from a notification whether to act, and was told to
have an agent do it. Each condition now says "Nothing for you to do." or
"Needs you:" with one thing they can do themselves, and carries the actions the
operator channel performs when chosen (hq ADR 0253).
2026-10-08 13:58:53 +02:00

722 lines
32 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package main
import (
"fmt"
"sort"
"strings"
"time"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/conditions"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// The signals table (novox/hq to-be 45 §3, ADR 0227 rule 5), compiled in.
//
// **Every signal the core expects has a watchdog; its absence is a condition.** A row says what is
// expected, from whom, after what, within what bound, and what is raised when it does not come. One
// table, read three ways: by the watchdog loop (watchdogs.go), which runs every row's watch over the
// facts it gathered; by `doctor signals`, which says the age of every row's newest signal; and by the
// test generated from it (signals_test.go), which suppresses each signal in turn and asserts its
// condition — so a row added without a watch, or a watch that does not fire, fails the build.
//
// A row not watched yet says why, and in which phase it will be: a watchdog of a signal nothing emits
// would be a condition that can only cry wolf. Bounds marked provisional are set from what Phase 0
// measured (`durations`) and corrected in Phase 1's first live week; a corrected bound is a change to
// this table, reviewed like code.
// signalRow is one row of the table.
type signalRow struct {
Row string
Signal string
Emitter string
Trigger string
Bound string
Kind string
Severity conditions.Severity
// Phase is the phase of to-be 45 its watchdog is built in.
Phase int
// Deferred says why it is not watched yet; empty for a row that is.
Deferred string
// needs is the part of the facts the row reads: an error there is the row blind, which is
// itself said (probe-failed) rather than read as nothing wrong.
needs func(f *signalFacts) error
// watch is what is wrong now, from the facts.
watch func(f *signalFacts) []conditions.Observation
// newest is the time of the newest signal of this row, for `doctor signals`; zero when none.
newest func(f *signalFacts) time.Time
}
// The bounds, provisional where to-be 45 says so.
const (
// heartbeatEvery is the interval a node-engine or node tools that say none have always used.
heartbeatEvery = 60 * time.Second
// heartbeatsMissed is how many intervals may pass in silence (S1, S11).
heartbeatsMissed = 3
// controlNodeUrgentAfter is how long the control node may be silent before it is urgent (S1).
controlNodeUrgentAfter = 30 * time.Minute
// reportAtLeast is the least a machine is given to report a send (S2).
reportAtLeast = 2 * time.Minute
// tierAtLeast is the least a plan's tier is given (S3), the bound `status` calls a plan late at.
tierAtLeast = planWaitBound
// loopDeafAfter is how long the event loop may take nothing while its consumers hold some (S4).
loopDeafAfter = 2 * time.Minute
// askAtLeast is the least a build ask is given, and askDefault the bound while nothing is
// measured (S6): the build seat declares no timeout of its own.
askAtLeast = 20 * time.Minute
askDefault = time.Hour
// callDefault is the bound of a verb that declares none (S7); push's and build's are longer.
callDefault = 10 * time.Minute
// advisoryQuiet is how long the bus must be quiet about a thing before its advisory clears (S9).
advisoryQuiet = time.Hour
// staleRefusalsAllowed in staleRefusalsWithin are what S13 lets pass from one writer.
staleRefusalsAllowed = 5
staleRefusalsWithin = 5 * time.Minute
// leaseBound is how long the lease may go unrenewed (S12): the key's age.
leaseBound = broker.LeaseTTL
// waitBound is how long a walk may wait for its delivery's word before it is said, and waitUrgentAfter
// before it is urgent (S16, novox/hq ADR 0239): a group member may wait behind another for a while; four
// hours without a word is the delivery's owner broken, whatever it says of itself.
waitBound = 30 * time.Minute
waitUrgentAfter = 4 * time.Hour
)
// callBounds are the verbs that may run longer than callDefault, and how long (S7).
var callBounds = map[string]time.Duration{
"push": 30 * time.Minute, "rotate": 30 * time.Minute, "assign": 15 * time.Minute,
"unassign": 15 * time.Minute, "command": 30 * time.Minute, "doctor": 3 * time.Minute,
}
// callBound is a verb's bound.
func callBound(verb string) time.Duration {
if b, ok := callBounds[verb]; ok {
return b
}
return callDefault
}
// signalsTable is the table, in to-be 45's order.
var signalsTable = []signalRow{
{Row: "S1", Signal: "machine heartbeat", Emitter: "node-engine", Trigger: "its interval",
Bound: "3 × the interval it says (60 s when it says none), provisional; not raised while the machine " +
"said it is asleep or shutting down (ADR 0211); urgent after 30 min for a control node",
Kind: "silent", Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.machinesErr }, watch: watchHeartbeats,
newest: func(f *signalFacts) time.Time {
return newestOf(f.machines, func(m machineFacts) time.Time { return m.lastHeard })
}},
{Row: "S2", Signal: "report after a send", Emitter: "node-engine", Trigger: "each declaration sent",
Bound: "max(2 min, 3 × that machine's last apply duration), provisional; not while the machine is silent or asleep",
Kind: "sent-not-reported", Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.machinesErr }, watch: watchReports,
newest: func(f *signalFacts) time.Time {
return newestOf(f.machines, func(m machineFacts) time.Time { return m.reportedAt })
}},
{Row: "S3", Signal: "plan tier progress", Emitter: "controller's plan", Trigger: "each tier entered",
Bound: "max(30 min, 3 × the p90 of that repository's measured tiers), provisional; not while the " +
"build seat is paused under it",
Kind: "stalled", Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.plansErr }, watch: watchPlans,
newest: func(f *signalFacts) time.Time {
return newestOf(f.plans, func(p planFacts) time.Time { return p.entered })
}},
{Row: "S4", Signal: "the controller's event loop takes a message", Emitter: "controller",
Trigger: "while its consumers have pending messages", Bound: "2 min",
Kind: "controller-deaf", Severity: conditions.Urgent, Phase: 1,
needs: func(f *signalFacts) error { return f.loopErr }, watch: watchLoop,
newest: func(f *signalFacts) time.Time { return f.loop.took }},
{Row: "S5", Signal: "a merge announced becomes a plan, or nothing reads it", Emitter: "announcer → controller",
Trigger: "each merge", Bound: "10 min (the catch-up pass of issue 266)",
Kind: "merge-not-acted", Severity: conditions.Urgent, Phase: 1,
needs: func(f *signalFacts) error { return f.mergesErr }, watch: watchMerges,
newest: func(f *signalFacts) time.Time { return f.mergesPassed }},
{Row: "S6", Signal: "a build asked → its outcome", Emitter: "build seat", Trigger: "each ask",
Bound: "max(20 min, 3 × the p90 of measured builds), 1 h while nothing is measured, provisional — the " +
"build seat declares no timeout; and an ask the queue gave up on (dead) at once",
Kind: "ask-lost", Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.asksErr }, watch: watchAsks,
newest: func(f *signalFacts) time.Time { return newestOf(f.asks, func(a askFacts) time.Time { return a.since }) }},
{Row: "S7", Signal: "a call running → finished", Emitter: "controller", Trigger: "each call",
Bound: "the verb's bound: push, rotate and command 30 min, assign and unassign 15 min, doctor 3 min, " +
"any other 10 min",
Kind: "call-hung", Severity: conditions.Warning, Phase: 1,
needs: func(*signalFacts) error { return nil }, watch: watchCalls,
newest: func(f *signalFacts) time.Time {
return newestOf(f.calls, func(c link.Call) time.Time { return c.Started })
}},
{Row: "S8", Signal: "a provider's failing word repeated", Emitter: "provider",
Trigger: "every 15 min while failing (ADR 0224)", Bound: "30 min",
Kind: kindProviderSilent, Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.standingsErr }, watch: watchProviders,
newest: func(f *signalFacts) time.Time {
return newestOf(f.standings, func(c conditions.Condition) time.Time { return c.LastObserved })
}},
{Row: "S9", Signal: "bus advisories: maximum deliveries, consumer deleted; the controller's own slow " +
"consumer and refused subjects", Emitter: "bus server's advisory subjects; the controller's connection",
Trigger: "any", Bound: "any occurrence; clears after an hour without another, and a deleted consumer " +
"once it exists again or the mesh no longer expects it",
Kind: "slow-consumer, max-deliveries, refused, consumer-lost", Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.advisoriesErr }, watch: watchAdvisories,
newest: func(f *signalFacts) time.Time {
return newestOf(f.advisories, func(a link.Advisory) time.Time { return a.Last })
}},
{Row: "S10", Signal: "the self-check's heartbeat", Emitter: "controller's doctor", Trigger: "every run",
Bound: "2 × its interval; watched from a second machine by mesh-watcher, and here as well",
Kind: "self-check-silent", Severity: conditions.Urgent, Phase: 1,
needs: func(*signalFacts) error { return nil }, watch: watchSelfCheck,
newest: func(f *signalFacts) time.Time { return f.selfCheck.last }},
{Row: "S11", Signal: "node tools heartbeat", Emitter: "node tools", Trigger: "its interval",
Bound: "3 × the interval it says (60 s when it says none), provisional; only where node-tools is assigned",
Kind: "tools-silent", Severity: conditions.Warning, Phase: 1,
needs: func(f *signalFacts) error { return f.machinesErr }, watch: watchTools,
newest: func(f *signalFacts) time.Time {
return newestOf(f.machines, func(m machineFacts) time.Time { return m.toolsHeard })
}},
{Row: "S12", Signal: "the controller lease renewed", Emitter: "controller", Trigger: "every 5 s",
Bound: "15 s (the key's age); a holder that lost the lease, or stopped renewing and was taken over, and a " +
"lease bucket found raised again from nothing, are said for an hour after; a controller serving without " +
"the lease, for as long as it does",
Kind: "lease-lost", Severity: conditions.Urgent, Phase: 2,
needs: func(f *signalFacts) error { return f.leaseErr }, watch: watchLease,
newest: func(f *signalFacts) time.Time { return f.lease.renewed }},
{Row: "S13", Signal: "stale refusals", Emitter: "every receiver (rule 2)", Trigger: "each refusal",
Bound: "more than 5 from one writer in 5 min: a controller epoch, a controller that claimed none, or a " +
"machine's node-engine whose accounts the controller refused",
Kind: "stale-writer", Severity: conditions.Warning, Phase: 2,
needs: func(*signalFacts) error { return nil }, watch: watchStaleRefusals,
newest: func(f *signalFacts) time.Time {
return newestOf(f.staleRefusals, func(w link.WriterRefusals) time.Time { return w.Last })
}},
{Row: "S14", Signal: "facts snapshot exported", Emitter: "controller", Trigger: "when it moved, and daily",
Bound: "2 days: a snapshot older than that, or none kept since this controller began two days ago",
Kind: "facts-stale", Severity: conditions.Warning, Phase: 5,
needs: func(*signalFacts) error { return nil }, watch: watchFacts,
newest: func(f *signalFacts) time.Time { return f.facts.taken }},
{Row: "S15", Signal: "a hand act with a cause already recorded", Emitter: "hand-act log",
Trigger: "each act", Bound: "the second within 14 days; clears when fewer than two remain within 14 days",
Kind: "healer-wanted", Severity: conditions.Warning, Phase: 3,
needs: func(f *signalFacts) error { return f.handActsErr }, watch: watchHandActs,
newest: func(f *signalFacts) time.Time {
return newestOf(f.handActs, func(a link.HandAct) time.Time { return a.At })
}},
{Row: "S16", Signal: "a walk waiting for its delivery's word is let go", Emitter: "mesh-delivery, through deliver",
Trigger: "each merge whose walk waits (novox/hq ADR 0239)",
Bound: "30 min, then urgent after 4 h — raised by the controller whatever mesh-delivery says of itself, " +
"naming `plans go` as the way on",
Kind: kindWalkWaiting, Severity: conditions.Warning, Phase: 3,
needs: func(f *signalFacts) error { return f.plansErr }, watch: watchWaits,
newest: func(f *signalFacts) time.Time {
return newestOf(f.waits, func(w waitFacts) time.Time { return w.since })
}},
}
// watchFacts is S14: the snapshot a merge check is fed is older than its bound, or none was kept since
// this controller began that long ago (novox/hq to-be 45 §9). A check fed a stale snapshot judges a change
// against a mesh that no longer is, which is the fault the snapshot exists to end.
func watchFacts(f *signalFacts) []conditions.Observation {
since := f.facts.taken
if since.IsZero() {
since = f.facts.began
}
if since.IsZero() || f.now.Sub(since) <= factsStaleAfter {
return nil
}
said := "none kept since this controller began " + ago(f.now.Sub(since)) + " ago"
if !f.facts.taken.IsZero() {
said = "the newest kept was taken " + ago(f.now.Sub(f.facts.taken)) + " ago"
}
if f.facts.err != nil {
said += "; the last attempt: " + oneLine(f.facts.err.Error())
}
return []conditions.Observation{{Scope: conditions.ScopeCore, ID: "facts", Kind: "facts-stale", Token: "stale",
Severity: conditions.Warning,
Summary: fmt.Sprintf("the facts snapshot merge checks are fed is stale (bound %s): a change is judged against a "+
"mesh that no longer is", ago(factsStaleAfter)),
Said: said}}
}
// newestOf is the newest time among things.
func newestOf[T any](list []T, at func(T) time.Time) time.Time {
var newest time.Time
for _, x := range list {
if t := at(x); t.After(newest) {
newest = t
}
}
return newest
}
// ago is a duration as the summaries say it.
func ago(d time.Duration) string {
if d < time.Minute {
return d.Round(time.Second).String()
}
return d.Round(time.Minute).String()
}
// heartbeatBound is a machine's S1 or S11 bound from the interval it says.
func heartbeatBound(every time.Duration) time.Duration {
if every <= 0 {
every = heartbeatEvery
}
return heartbeatsMissed * every
}
func watchHeartbeats(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, m := range f.machines {
if m.lastHeard.IsZero() || m.asleep() {
// Never heard is a machine that has not joined, which status says; asleep is not lost.
continue
}
bound := heartbeatBound(m.every)
silent := f.now.Sub(m.lastHeard)
if silent <= bound {
continue
}
severity := conditions.Warning
if m.control && silent > controlNodeUrgentAfter {
severity = conditions.Urgent
}
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: m.name, Kind: "silent",
Machine: m.name, Severity: severity,
Summary: fmt.Sprintf("%s has not been heard from since %s (bound %s)", m.name,
m.lastHeard.UTC().Format("2006-01-02 15:04 MST"), bound),
Said: fmt.Sprintf("silent for %s", ago(silent))})
}
return out
}
func watchReports(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, m := range f.machines {
if m.sentAt.IsZero() || m.reportedCurrent || m.asleep() {
continue
}
if !m.lastHeard.IsZero() && f.now.Sub(m.lastHeard) > heartbeatBound(m.every) {
continue // silent: S1 says it, and a silent machine reports nothing
}
bound := max(reportAtLeast, 3*m.lastApply)
waited := f.now.Sub(m.sentAt)
if waited <= bound {
continue
}
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: m.name,
Kind: "sent-not-reported", Machine: m.name, Severity: conditions.Warning,
Summary: fmt.Sprintf("%s was sent a declaration at %s and has not reported applying it (bound %s)",
m.name, m.sentAt.UTC().Format("2006-01-02 15:04 MST"), ago(bound)),
Said: fmt.Sprintf("waiting %s for the report of the declaration sent", ago(waited))})
}
return out
}
func watchPlans(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, p := range f.plans {
if p.paused {
continue
}
in := f.now.Sub(p.entered)
if in <= p.bound {
continue
}
out = append(out, conditions.Observation{Scope: conditions.ScopePlan, ID: p.id, Kind: "stalled",
Headline: "Delivery of " + repoName(p.repository) + " is stuck halfway",
Explanation: fmt.Sprintf("Its walk across the machines has been at the same step for %s, longer than "+
"usual. Nothing is lost, and the mesh keeps it where it is.", humanDuration(f.now.Sub(p.entered))),
Resolved: "Delivery of " + repoName(p.repository) + " moves again",
Severity: conditions.Warning,
Summary: fmt.Sprintf("the plan for %s %s has been at tier %d of %d since %s (bound %s): %s",
p.repository, short(p.commit), p.tier+1, p.tiers, p.entered.UTC().Format("2006-01-02 15:04 MST"),
ago(p.bound), p.waiting),
Said: fmt.Sprintf("at tier %d for %s: %s", p.tier+1, ago(in), p.waiting)})
}
return out
}
// kindWalkWaiting is S16's kind: plan.<id>.waiting.
const kindWalkWaiting = "waiting"
// watchWaits is S16: a walk that has waited for its delivery's word past its bound, said by the controller
// itself — a delivery's owner that is up and never says go (a bug, a delivery stuck in its own table) would
// otherwise leave a merge waiting for ever with nothing open (novox/hq ADR 0239).
func watchWaits(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, w := range f.waits {
in := f.now.Sub(w.since)
if in <= waitBound {
continue
}
severity := conditions.Warning
if in > waitUrgentAfter {
severity = conditions.Urgent
}
out = append(out, conditions.Observation{Scope: conditions.ScopePlan, ID: w.id, Kind: kindWalkWaiting,
Severity: severity,
Summary: fmt.Sprintf("the walk of %s %s has waited %s for %s's word to start: `mesh-delivery.show` for the "+
"delivery that landed as %s says why; `plans go %s --why …` starts it by hand", w.repository,
short(w.commit), ago(in), w.awaits, short(w.commit), w.id),
Said: fmt.Sprintf("waiting since %s for %s", w.since.UTC().Format(time.RFC3339), w.awaits),
Headline: deliveryName(w.modules, w.repository) + " waiting to start",
Explanation: walkWaitingWords(w, in, severity),
Needs: waitingNeeds(severity),
Actions: waitingActions(w, severity),
Resolved: deliveryName(w.modules, w.repository) + " no longer waiting"})
}
return out
}
func watchLoop(f *signalFacts) []conditions.Observation {
if f.loop.pending == 0 {
return nil
}
since := f.loop.took
if since.IsZero() {
since = f.started
}
if f.now.Sub(since) <= loopDeafAfter {
return nil
}
return []conditions.Observation{{Scope: conditions.ScopeCore, ID: "controller", Kind: "controller-deaf",
Token: "deaf", Machine: f.host, Severity: conditions.Urgent,
Summary: fmt.Sprintf("the controller's event loop has taken nothing for %s while its consumers hold %d "+
"message(s): reports, builds and merges are not being acted on", ago(f.now.Sub(since)), f.loop.pending),
Said: fmt.Sprintf("%d pending (%s), last taken %s", f.loop.pending, f.loop.where, since.UTC().Format(time.RFC3339))}}
}
func watchMerges(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, m := range f.merges {
out = append(out, conditions.Observation{Scope: conditions.ScopeMerge, ID: m.Repo + "." + short(m.Commit),
Kind: "merge-not-acted", Token: "not-acted", Severity: conditions.Urgent,
Summary: fmt.Sprintf("%s/%s merged into %s (%s) was never handed to the controller by the bus; "+
"acted on late by the catch-up", m.Owner, m.Repo, m.Base, short(m.Commit)),
Said: fmt.Sprintf("announced %s, %s behind it", m.At.UTC().Format(time.RFC3339), readableList(m.Modules))})
}
return out
}
func watchAsks(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, a := range f.asks {
var said string
switch {
case a.state == link.AskDead:
said = fmt.Sprintf("the %s queue handed it out as often as it may and it was never settled", a.seat)
case a.state == link.AskInFlight && f.now.Sub(a.since) > a.bound:
said = fmt.Sprintf("in flight on %s for %s (bound %s)", orSomewhere(a.on), ago(f.now.Sub(a.since)), ago(a.bound))
default:
continue
}
out = append(out, conditions.Observation{Scope: conditions.ScopeBuild, ID: a.id, Kind: "ask-lost",
Token: "lost", Machine: a.on, Severity: conditions.Warning,
Summary: fmt.Sprintf("the build %s of %s has no outcome: %s", a.id, a.what, said), Said: said})
}
return out
}
func orSomewhere(node string) string {
if node == "" {
return "a machine that did not say which"
}
return node
}
func watchCalls(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, c := range f.calls {
bound := callBound(c.Verb)
running := f.now.Sub(c.Started)
if running <= bound {
continue
}
out = append(out, conditions.Observation{Scope: conditions.ScopeCall, ID: c.ID, Kind: "call-hung",
Token: "hung", Machine: f.host, Severity: conditions.Warning,
Summary: fmt.Sprintf("%s.%s (call %s) has been running since %s, past its bound of %s", c.Seat, c.Verb,
c.ID, c.Started.UTC().Format("2006-01-02 15:04 MST"), ago(bound)),
Said: fmt.Sprintf("running %s, asked by %s", ago(running), orSomebody(c.Caller))})
}
return out
}
func orSomebody(caller string) string {
if caller == "" {
return "a caller the bus did not name"
}
return caller
}
func watchProviders(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, c := range f.standings {
quiet := f.now.Sub(c.LastObserved)
if quiet <= providerSaysAgainWithin {
continue
}
module, node, consumer, ok := providerOf(c)
if !ok {
continue
}
out = append(out, conditions.Observation{Scope: conditions.ScopeProvider, ID: module + "." + node + "." + consumer,
Token: "silent", Kind: kindProviderSilent, Machine: node, Severity: conditions.Warning,
Summary: fmt.Sprintf("%s on %s said it keeps failing %s and has said nothing since %s: it stopped "+
"saying anything, so its last word is all the mesh has", module, node, consumer,
c.LastObserved.UTC().Format("2006-01-02 15:04 MST")),
Said: fmt.Sprintf("not said again for %s", ago(quiet))})
}
return out
}
func watchAdvisories(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, a := range f.advisories {
if f.now.Sub(a.Last) > advisoryQuiet {
continue
}
if a.Kind == link.AdvisoryConsumerLost && !f.lostConsumers[a.Stream+"."+a.Consumer] {
continue // it exists again, or the mesh no longer expects it: a removal, not a loss
}
severity := conditions.Warning
times := ""
if a.Count > 1 {
times = fmt.Sprintf(" (%d times since %s)", a.Count, a.First.UTC().Format("15:04 MST"))
}
machine := ""
if a.ID == "controller" {
machine = f.host
}
out = append(out, conditions.Observation{Scope: conditions.ScopeBus, ID: a.ID, Kind: a.Kind,
Machine: machine, Severity: severity, Summary: a.Said + times, Said: a.Said})
}
return out
}
func watchSelfCheck(f *signalFacts) []conditions.Observation {
every := f.selfCheck.every
if every <= 0 {
every = doctorEvery
}
since := f.selfCheck.last
if since.IsZero() {
since = f.started
}
if f.now.Sub(since) <= 2*every {
return nil
}
return []conditions.Observation{{Scope: conditions.ScopeCore, ID: "doctor", Kind: "self-check-silent",
Machine: f.host, Severity: conditions.Urgent,
Summary: fmt.Sprintf("the self-check has not finished a run since %s (it runs every %s): the mesh's "+
"invariants are not being checked", since.UTC().Format("2006-01-02 15:04 MST"), every),
Said: fmt.Sprintf("no run for %s", ago(f.now.Sub(since)))}}
}
func watchTools(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, m := range f.machines {
if !m.tools || m.asleep() {
continue
}
bound := heartbeatBound(m.toolsEvery)
since := m.toolsHeard
if since.IsZero() {
// Not heard since this controller started: silent since then at the most.
since = f.toolsHeardFrom
}
silent := f.now.Sub(since)
if silent <= bound {
continue
}
heard := "not since this controller started at " + since.UTC().Format("2006-01-02 15:04 MST")
if !m.toolsHeard.IsZero() {
heard = "since " + m.toolsHeard.UTC().Format("2006-01-02 15:04 MST")
}
out = append(out, conditions.Observation{Scope: conditions.ScopeMachine, ID: m.name, Kind: "tools-silent",
Machine: m.name, Severity: conditions.Warning,
Summary: fmt.Sprintf("the node tools on %s have not said they are there %s (bound %s): nothing can "+
"ask that machine anything", m.name, heard, bound),
Said: fmt.Sprintf("silent for %s", ago(silent))})
}
return out
}
func watchStaleRefusals(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
for _, w := range f.staleRefusals {
if w.Count <= staleRefusalsAllowed {
continue
}
scope, id, machine := conditions.ScopeCore, "controller.unnamed", ""
named := w.Writer
switch {
case w.Epoch > 0:
id = fmt.Sprintf("controller.epoch-%d", w.Epoch)
if e, ok := f.epochs[w.Epoch]; ok {
named = fmt.Sprintf("the controller of epoch %d (%s%s)", w.Epoch, e.Instance, endedWords(e))
}
case w.Writer == link.WriterNodeEngine(firstOf(w.Receivers)) || strings.HasPrefix(w.Writer, "the node-engine on "):
node := strings.TrimPrefix(w.Writer, "the node-engine on ")
scope, id, machine = conditions.ScopeMachine, node, node
}
if machine == "" && len(w.Receivers) > 0 {
machine = w.Receivers[0]
}
var also []string
for _, r := range w.Receivers {
if r != machine && r != "controller" {
also = append(also, r)
}
}
out = append(out, conditions.Observation{Scope: scope, ID: id, Kind: "stale-writer", Token: "stale-writer",
Machine: machine, Also: also, Severity: conditions.Warning,
Summary: fmt.Sprintf("%s was refused %d time(s) in %s as older than what its receivers hold (%s): a "+
"writer is sending what the mesh has moved past", named, w.Count, staleRefusalsWithin,
strings.Join(w.Receivers, ", ")),
Said: fmt.Sprintf("%d stale refusals in %s, the last at %s", w.Count, staleRefusalsWithin,
w.Last.UTC().Format(time.RFC3339))})
}
return out
}
// firstOf is a list's first, empty for none.
func firstOf(list []string) string {
if len(list) == 0 {
return ""
}
return list[0]
}
// endedWords is how an epoch ended, for a sentence naming it; nothing while it is held.
func endedWords(e inventory.Epoch) string {
if e.Ended == nil {
return ", still holding the lease"
}
return fmt.Sprintf(", %s at %s", e.How, e.Ended.UTC().Format("15:04:05 MST"))
}
// watchLease is S12: the lease held and renewed by this controller, and every holder that lost it.
func watchLease(f *signalFacts) []conditions.Observation {
var out []conditions.Observation
l := f.lease
if l.unleased != "" {
out = append(out, conditions.Observation{Scope: conditions.ScopeCore, ID: "controller.lease", Token: "unleased",
Kind: "lease-lost", Machine: f.host, Severity: conditions.Urgent,
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",
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",
Kind: "lease-lost", Machine: f.host, Severity: conditions.Urgent,
Summary: fmt.Sprintf("the controller has not renewed its lease (epoch %d) since %s, past the key's age "+
"of %s: another controller may take it", l.epoch, l.renewed.UTC().Format(time.RFC3339), leaseBound),
Said: fmt.Sprintf("not renewed for %s", ago(f.now.Sub(l.renewed)))})
}
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: whatever held the lease before " +
"does not know it lost it",
Said: l.resetSaid})
}
var lost []string
newest := inventory.Epoch{}
for _, e := range l.ended {
if e.Ended == nil || e.How == inventory.EpochReleased || f.now.Sub(*e.Ended) > advisoryQuiet {
continue
}
lost = append(lost, fmt.Sprintf("epoch %d (%s) %s at %s", e.Epoch, e.Instance, e.How,
e.Ended.UTC().Format("15:04:05 MST")))
if newest.Ended == nil || e.Ended.After(*newest.Ended) {
newest = e
}
}
if len(lost) > 0 {
how := "lost it: its renewal was refused or could not be made"
if newest.How == inventory.EpochExpired {
how = "stopped renewing it without giving it back, and was taken over"
}
out = append(out, conditions.Observation{Scope: conditions.ScopeCore, ID: "controller.lease", Token: "lost",
Kind: "lease-lost", Machine: newest.Host, Severity: conditions.Urgent,
Summary: fmt.Sprintf("the controller of epoch %d (%s) %s; the controller of epoch %d acts now", newest.Epoch,
newest.Instance, how, l.epoch),
Said: strings.Join(lost, "; ")})
}
return out
}
// handActsWithin is how far back a repeated cause counts (S15).
const handActsWithin = 14 * 24 * time.Hour
// watchHandActs is S15: a cause recorded by hand twice within a fortnight is a healer wanted, named by
// the cause. **A heal is never a hand act** (healers.go), so a cause a healer exists for and a person
// still repaired twice says the healer is not enough — its reach or its budget — and is said so. **An act
// that is a person's decision by design is no repair** (handActVerbs), so it never counts; one counted
// before this was so clears on the next tick, as any condition the row no longer sees.
func watchHandActs(f *signalFacts) []conditions.Observation {
recent := make([]link.HandAct, 0, len(f.handActs))
for _, a := range repairs(f.handActs) {
if f.now.Sub(a.At) <= handActsWithin {
recent = append(recent, a)
}
}
repeated := link.RepeatedCauses(recent, f.now)
causes := make([]string, 0, len(repeated))
for c := range repeated {
causes = append(causes, c)
}
sort.Strings(causes)
var out []conditions.Observation
for _, cause := range causes {
var acts []string
var newest link.HandAct
for _, a := range recent {
if a.Cause != cause {
continue
}
acts = append(acts, fmt.Sprintf("%s %s by %s: %s", a.At.UTC().Format("2006-01-02 15:04"),
strings.TrimSpace(a.Verb+" "+strings.Join(a.Args, " ")), a.By, a.Why))
if a.At.After(newest.At) {
newest = a
}
}
wanted := "a healer is wanted for it"
if h := healerNamedFor(cause); h != "" {
wanted = fmt.Sprintf("healer %s answers this cause and a person still repaired it: its reach or its "+
"budget is not enough", h)
}
out = append(out, conditions.Observation{Scope: conditions.ScopeMesh, ID: "hand-acts." + cause,
Token: "healer-wanted", Kind: "healer-wanted", Severity: conditions.Warning,
Summary: fmt.Sprintf("%q was repaired by hand %d times in %d days, the last by %s: %s", cause,
repeated[cause], int(handActsWithin.Hours()/24), newest.By, wanted),
Said: strings.Join(acts, "; "),
Headline: conditions.Capital(causeWords(cause)) + " keeps being fixed by hand",
Explanation: fmt.Sprintf("A person repaired %s by hand %d times in %d days, so an automatic repair is "+
"wanted for it. Nothing is broken now.", causeWords(cause), repeated[cause],
int(handActsWithin.Hours()/24)),
Resolved: "Resolved: no more hand repairs of " + causeWords(cause)})
}
return out
}
// watchedRows are the rows a watchdog runs for.
func watchedRows() []signalRow {
var out []signalRow
for _, r := range signalsTable {
if r.watch != nil {
out = append(out, r)
}
}
return out
}
// kindsOf is a row's condition kinds, one or several.
func kindsOf(r signalRow) []string {
var out []string
for _, k := range strings.Split(r.Kind, ",") {
out = append(out, strings.TrimSpace(k))
}
return out
}