diff --git a/internal/loop/gather.go b/internal/loop/gather.go index c55919f..73e164f 100644 --- a/internal/loop/gather.go +++ b/internal/loop/gather.go @@ -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 diff --git a/internal/loop/rules.go b/internal/loop/rules.go index f144e2d..1fc5ff6 100644 --- a/internal/loop/rules.go +++ b/internal/loop/rules.go @@ -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:`. 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) }, } } diff --git a/internal/loop/state.go b/internal/loop/state.go index f17b074..1c42f6e 100644 --- a/internal/loop/state.go +++ b/internal/loop/state.go @@ -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) { diff --git a/internal/store/facts.go b/internal/store/facts.go index 6e0a2c6..0644fda 100644 --- a/internal/store/facts.go +++ b/internal/store/facts.go @@ -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". diff --git a/internal/store/facts_test.go b/internal/store/facts_test.go new file mode 100644 index 0000000..47e666e --- /dev/null +++ b/internal/store/facts_test.go @@ -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) + } +}