Kuma: a fact per monitor, so she can name the service that is down #147

Merged
claude merged 4 commits from task/444-kuma-a-fact-per-monitor-so-she-can-name into master 2026-08-04 18:24:48 +02:00
5 changed files with 194 additions and 13 deletions
Showing only changes of commit 4f516657da - Show all commits
+14
View File
@@ -91,6 +91,20 @@ func (g *Gatherer) GatherState(ctx context.Context, now time.Time) (State, []sto
return State{}, nil, err
}
// prefix families — the keys a rule cannot name at wiring time (one fact
// per kuma monitor). Loaded into the same map; State.FactsUnder reads them.
for _, r := range g.rules {
for _, p := range r.WantPrefixes {
fam, err := g.store.LatestFactsByPrefix(ctx, p)
if err != nil {
return State{}, nil, err
}
for _, f := range fam {
facts[f.Key] = f
}
}
}
// last nudge per rule + cooldown-until derived from the active cooldown.
// "active" = the feedback tuner's persisted base if one exists, else the
// rule's static Base. LatestFactBySource is the trust-by-provenance read
+55 -13
View File
@@ -28,6 +28,13 @@ type Rule struct {
// that check itself, leave this empty. Otherwise set to the key(s) the rule
// needs and the gate will skip the rule when any are missing.
InertWhenNoData []string
// WantPrefixes — key prefixes whose whole family the gatherer must load.
// InertWhenNoData names keys that exist at wiring time; a rule over a key
// set that is only known at read time (one fact per kuma monitor) declares
// the prefix here instead. Prefixes never make a rule inert: an empty
// family is the predicate's own "no data" case.
WantPrefixes []string
}
// Cooldown — tunable bounded by the envelope so a weird week (auto-tuned) can't
@@ -98,23 +105,58 @@ func BreakRule() Rule {
}
}
// ServiceDownRule — sev4 ops hard: the `service_down` aggregate fact reads
// "down". Source must be poll:uptimekuma — kuma is the source of truth for
// service up/down (mavpoll writes this key). The predicate is provenance-scoped:
// a compromised poller writing under a different source can't forge the trigger.
// ServiceDownPrefix — mavpoll writes one fact per kuma monitor under this
// prefix, `service_down:<monitor name>`. The suffix is the name he hears.
const ServiceDownPrefix = "service_down:"
// ServiceDownSource — kuma is the source of truth for service up/down. The
// rule is provenance-scoped: a poller writing under a different source cannot
// forge the trigger.
const ServiceDownSource = "poll:uptimekuma"
// DownServices — the monitors currently reading "down", by name, in key order.
//
// Pure, and the rule and the phraser both call it, so the message can never
// name a service the predicate did not fire on.
func DownServices(s State) []string {
var out []string
for _, f := range s.FactsUnder(ServiceDownPrefix) {
if f.Source == ServiceDownSource && f.Value == `"down"` {
out = append(out, strings.TrimPrefix(f.Key, ServiceDownPrefix))
}
}
return out
}
// ServiceDownRule — sev4 ops hard: at least one kuma monitor reads "down".
//
// It used to read one aggregate `service_down` fact, which is why it was
// disabled in deploy: the nudge could say that something on homesrv was down
// but never which thing. Per-monitor facts fix that, and pausing a monitor in
// kuma now silences that monitor rather than nothing.
//
// Edge-triggered — see State.NudgedSince. Without it a service that stays down
// for a day qualifies on every tick and cooldown alone is the only brake.
func ServiceDownRule() Rule {
return Rule{
Name: "service_down",
Severity: Sev4,
Cooldown: Cooldown{Base: 15 * time.Minute, Min: 5 * time.Minute, Max: 1 * time.Hour},
InertWhenNoData: []string{"service_down"},
Name: "service_down",
Severity: Sev4,
Cooldown: Cooldown{Base: 15 * time.Minute, Min: 5 * time.Minute, Max: 1 * time.Hour},
WantPrefixes: []string{ServiceDownPrefix},
Predicate: func(s State) bool {
f, ok := s.Fact("service_down")
if !ok || f.Ts.IsZero() {
return false
var newest time.Time
for _, f := range s.FactsUnder(ServiceDownPrefix) {
if f.Source != ServiceDownSource || f.Value != `"down"` {
continue
}
if f.Ts.After(newest) {
newest = f.Ts
}
}
// value is json `"down"`; trivial check keyed off source provenance.
return f.Source == "poll:uptimekuma" && f.Value == `"down"`
if newest.IsZero() {
return false // nothing down, or no data at all → shut up
}
return !s.NudgedSince("service_down", newest)
},
}
}
+28
View File
@@ -23,6 +23,8 @@
package loop
import (
"sort"
"strings"
"time"
"github.com/kami/maven/internal/store"
@@ -95,6 +97,32 @@ func (s State) Fact(key string) (store.Fact, bool) {
return f, true
}
// FactsUnder returns every gathered fact whose key starts with prefix, ordered
// by key so a caller that names them speaks them in a stable order. Facts with
// a zero Ts are skipped, the same "no data" rule Fact applies.
func (s State) FactsUnder(prefix string) []store.Fact {
var out []store.Fact
for k, f := range s.Facts {
if strings.HasPrefix(k, prefix) && !f.Ts.IsZero() {
out = append(out, f)
}
}
sort.Slice(out, func(i, j int) bool { return out[i].Key < out[j].Key })
return out
}
// NudgedSince reports whether rule already sent a nudge at or after ts.
//
// It is what makes a rule edge-triggered. A polled fact is written only when
// the value changes, so its Ts is the moment the service went down — but the
// predicate reads the current value, so a service that stays down keeps
// qualifying forever and cooldown alone only slows the repetition. Asking
// whether he was already told about THIS transition stops it.
func (s State) NudgedSince(rule string, ts time.Time) bool {
n, ok := s.LastNudge[rule]
return ok && !n.Ts.Before(ts)
}
// Since returns the duration since the latest fact for key, or (0,false).
// "false" ⇒ no data ⇒ shuts up when uncertain.
func (s State) Since(key string) (time.Duration, bool) {
+35
View File
@@ -254,6 +254,41 @@ func (s *Store) LatestFactBySource(ctx context.Context, key, source string) (Fac
return scanFact(row)
}
// LatestFactsByPrefix — the latest non-voided fact for every key that starts
// with prefix, newest-per-key, ordered by key.
//
// The loop's gatherer loads the keys its rules declare, which works while the
// key set is static. Kuma's monitors are not: one fact per monitor means the
// keys are only known once the gauge is read, so the rule declares the prefix
// and this read resolves it per tick. `_` and `%` are escaped — a monitor name
// is user text and must not act as a LIKE wildcard.
func (s *Store) LatestFactsByPrefix(ctx context.Context, prefix string) ([]Fact, error) {
esc := strings.NewReplacer(`\`, `\\`, `%`, `\%`, `_`, `\_`).Replace(prefix)
rows, err := s.db.QueryContext(ctx, `
SELECT id, ts, kind, key, value, source, confidence, voids_id
FROM facts f
WHERE key LIKE ? ESCAPE '\'
AND id NOT IN (SELECT voids_id FROM facts WHERE voids_id IS NOT NULL)
AND id = (SELECT id FROM facts g
WHERE g.key = f.key
AND g.id NOT IN (SELECT voids_id FROM facts WHERE voids_id IS NOT NULL)
ORDER BY g.ts DESC, g.id DESC LIMIT 1)
ORDER BY key`, esc+"%")
if err != nil {
return nil, fmt.Errorf("facts by prefix %q: %w", prefix, err)
}
defer rows.Close()
var out []Fact
for rows.Next() {
f, err := scanFact(rows)
if err != nil {
return nil, err
}
out = append(out, f)
}
return out, rows.Err()
}
// Since returns how long ago the latest non-voided fact for key landed, or
// (0, ErrNoFact). Implements the `since(key)==null → don't fire` guard from
// the spec — silence on no-data is "shuts up when uncertain".
+62
View File
@@ -0,0 +1,62 @@
package store
import (
"context"
"testing"
"time"
)
// One fact per kuma monitor means the loop cannot name its keys at wiring time,
// so it asks for the family by prefix. The read must return the newest row per
// key and stop at the prefix boundary.
func TestLatestFactsByPrefix(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
now := time.Now().UTC().Truncate(time.Millisecond)
write := func(key, val string, at time.Time) int64 {
id, err := s.SetValue(ctx, KindEnv, key, "poll:uptimekuma", val, at)
if err != nil {
t.Fatalf("SetValue %s: %v", key, err)
}
return id
}
write("service_down:db", "up", now.Add(-2*time.Hour))
write("service_down:db", "down", now.Add(-time.Hour)) // newer wins
write("service_down:web", "up", now.Add(-time.Hour))
write("service_downtime", "irrelevant", now) // no colon, not in the family
write("water", "250ml", now)
got, err := s.LatestFactsByPrefix(ctx, "service_down:")
if err != nil {
t.Fatalf("LatestFactsByPrefix: %v", err)
}
if len(got) != 2 {
t.Fatalf("got %d facts, want 2: %+v", len(got), got)
}
if got[0].Key != "service_down:db" || got[0].Value != `"down"` {
t.Errorf("first = %s=%s, want the newest db row", got[0].Key, got[0].Value)
}
if got[1].Key != "service_down:web" {
t.Errorf("second = %s, want service_down:web", got[1].Key)
}
}
// A monitor name is user text. An underscore in it must match itself, not act
// as a LIKE wildcard and drag in every other monitor.
func TestLatestFactsByPrefixEscapesWildcards(t *testing.T) {
s := newTestStore(t)
ctx := context.Background()
now := time.Now().UTC().Truncate(time.Millisecond)
for _, k := range []string{"a_b:one", "axb:two"} {
if _, err := s.SetValue(ctx, KindEnv, k, "poll:uptimekuma", "down", now); err != nil {
t.Fatalf("SetValue %s: %v", k, err)
}
}
got, err := s.LatestFactsByPrefix(ctx, "a_b:")
if err != nil {
t.Fatalf("LatestFactsByPrefix: %v", err)
}
if len(got) != 1 || got[0].Key != "a_b:one" {
t.Fatalf("got %+v, want only a_b:one", got)
}
}