OPS-05: alerting-bei-schwellwert-ueberschreitung (internal/alerting: regel-store, evaluator gegen ops-03-metriken, cfg-02-zustellung, drosselung je regel+zeitreihe)
This commit is contained in:
@@ -0,0 +1,194 @@
|
|||||||
|
package alerting
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"sort"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
dto "github.com/prometheus/client_model/go"
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
|
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/internal/notify"
|
||||||
|
)
|
||||||
|
|
||||||
|
// DefaultDebounceInterval: wiederholte Alarmierung für denselben
|
||||||
|
// anhaltenden Zustand ist gedrosselt (Akzeptanzkriterium 3) — 15 Minuten
|
||||||
|
// ist ein üblicher Kompromiss zwischen "schnell genug informiert" und
|
||||||
|
// "kein Alarm-Spam bei dauerhaft überschrittenem Wert".
|
||||||
|
const DefaultDebounceInterval = 15 * time.Minute
|
||||||
|
|
||||||
|
// AlertChannel ist der CFG-02-Kanal, über den Schwellwert-Alarme zugestellt
|
||||||
|
// werden — ein eigener Kanalname, damit Zustellregeln/-vorlagen (CFG-03)
|
||||||
|
// unabhängig von anderen Benachrichtigungsarten konfiguriert werden können.
|
||||||
|
const AlertChannel = "alert"
|
||||||
|
|
||||||
|
// Evaluator prüft konfigurierte Regeln gegen aktuell gesammelte Metriken
|
||||||
|
// (aus internal/metrics.Aggregator.Gather) und löst bei Überschreitung eine
|
||||||
|
// Benachrichtigung über CFG-02 aus (Akzeptanzkriterium 2), gedrosselt je
|
||||||
|
// Regel+Zeitreihe (Akzeptanzkriterium 3).
|
||||||
|
type Evaluator struct {
|
||||||
|
rules *RuleStore
|
||||||
|
debounce *debounceStore
|
||||||
|
dispatcher *notify.Dispatcher
|
||||||
|
interval time.Duration
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewEvaluator(rules *RuleStore, dispatcher *notify.Dispatcher, debouncePool *pgxpool.Pool, interval time.Duration) *Evaluator {
|
||||||
|
if interval <= 0 {
|
||||||
|
interval = DefaultDebounceInterval
|
||||||
|
}
|
||||||
|
return &Evaluator{
|
||||||
|
rules: rules,
|
||||||
|
debounce: &debounceStore{pool: debouncePool},
|
||||||
|
dispatcher: dispatcher,
|
||||||
|
interval: interval,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// FiredAlert beschreibt einen tatsächlich ausgelösten (nicht gedrosselten)
|
||||||
|
// Alarm — fürs Testen/Logging, nicht Teil des öffentlichen Zustellwegs.
|
||||||
|
type FiredAlert struct {
|
||||||
|
RuleID string
|
||||||
|
MetricName string
|
||||||
|
Value float64
|
||||||
|
Threshold float64
|
||||||
|
Labels map[string]string
|
||||||
|
Skipped bool // true, wenn wegen Drosselung NICHT tatsaechlich zugestellt
|
||||||
|
}
|
||||||
|
|
||||||
|
// Evaluate prüft alle konfigurierten Regeln gegen families (Akzeptanzkriterium 1).
|
||||||
|
// Für jede Zeitreihe, die eine Regel verletzt, wird — sofern nicht gedrosselt
|
||||||
|
// — eine Benachrichtigung mit Metrik/Wert/Schwellwert/Labels (Tenant/Modul,
|
||||||
|
// falls als Label vorhanden) über CFG-02 eingereiht (Akzeptanzkriterium 2).
|
||||||
|
func (e *Evaluator) Evaluate(ctx context.Context, families []*dto.MetricFamily) ([]FiredAlert, error) {
|
||||||
|
rules, err := e.rules.ListRules(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("regeln laden: %w", err)
|
||||||
|
}
|
||||||
|
if len(rules) == 0 {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
byName := make(map[string]*dto.MetricFamily, len(families))
|
||||||
|
for _, f := range families {
|
||||||
|
if f.Name != nil {
|
||||||
|
byName[*f.Name] = f
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
var fired []FiredAlert
|
||||||
|
for _, rule := range rules {
|
||||||
|
family, ok := byName[rule.MetricName]
|
||||||
|
if !ok {
|
||||||
|
continue // Metrik (noch) nicht vorhanden -> keine Aussage moeglich, kein Fehler.
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, m := range family.Metric {
|
||||||
|
labels := labelMap(m)
|
||||||
|
if !matchesFilter(labels, rule.LabelFilters) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
value, ok := metricValue(m)
|
||||||
|
if !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if !violates(rule, value) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
ruleKey := ruleKeyFor(rule.ID, labels)
|
||||||
|
allowed, err := e.debounce.shouldFire(ctx, ruleKey, e.interval.Seconds())
|
||||||
|
if err != nil {
|
||||||
|
return fired, fmt.Errorf("drosselung pruefen: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
alert := FiredAlert{
|
||||||
|
RuleID: rule.ID, MetricName: rule.MetricName, Value: value,
|
||||||
|
Threshold: rule.Threshold, Labels: labels, Skipped: !allowed,
|
||||||
|
}
|
||||||
|
fired = append(fired, alert)
|
||||||
|
|
||||||
|
if !allowed {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
payload := map[string]any{
|
||||||
|
"metric": rule.MetricName,
|
||||||
|
"value": value,
|
||||||
|
"threshold": rule.Threshold,
|
||||||
|
"comparison": string(rule.Comparison),
|
||||||
|
"description": rule.Description,
|
||||||
|
"labels": labels,
|
||||||
|
}
|
||||||
|
if _, err := e.dispatcher.Enqueue(ctx, AlertChannel, rule.Recipient, payload); err != nil {
|
||||||
|
return fired, fmt.Errorf("alarm einreihen: %w", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return fired, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func violates(rule Rule, value float64) bool {
|
||||||
|
switch rule.Comparison {
|
||||||
|
case ComparisonGreaterThan:
|
||||||
|
return value > rule.Threshold
|
||||||
|
case ComparisonLessThan:
|
||||||
|
return value < rule.Threshold
|
||||||
|
default:
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func labelMap(m *dto.Metric) map[string]string {
|
||||||
|
out := make(map[string]string, len(m.Label))
|
||||||
|
for _, l := range m.Label {
|
||||||
|
if l.Name != nil && l.Value != nil {
|
||||||
|
out[*l.Name] = *l.Value
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
func matchesFilter(labels, filter map[string]string) bool {
|
||||||
|
for k, v := range filter {
|
||||||
|
if labels[k] != v {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
func metricValue(m *dto.Metric) (float64, bool) {
|
||||||
|
switch {
|
||||||
|
case m.Gauge != nil && m.Gauge.Value != nil:
|
||||||
|
return *m.Gauge.Value, true
|
||||||
|
case m.Counter != nil && m.Counter.Value != nil:
|
||||||
|
return *m.Counter.Value, true
|
||||||
|
case m.Untyped != nil && m.Untyped.Value != nil:
|
||||||
|
return *m.Untyped.Value, true
|
||||||
|
default:
|
||||||
|
return 0, false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ruleKeyFor macht die Drosselung unabhaengig je Regel UND je konkreter
|
||||||
|
// Zeitreihe (z. B. verschiedene Tenants/Module derselben Metrik loesen
|
||||||
|
// unabhaengig voneinander aus, siehe Migrationskommentar).
|
||||||
|
func ruleKeyFor(ruleID string, labels map[string]string) string {
|
||||||
|
keys := make([]string, 0, len(labels))
|
||||||
|
for k := range labels {
|
||||||
|
keys = append(keys, k)
|
||||||
|
}
|
||||||
|
sort.Strings(keys)
|
||||||
|
var b strings.Builder
|
||||||
|
b.WriteString(ruleID)
|
||||||
|
for _, k := range keys {
|
||||||
|
b.WriteString("|")
|
||||||
|
b.WriteString(k)
|
||||||
|
b.WriteString("=")
|
||||||
|
b.WriteString(labels[k])
|
||||||
|
}
|
||||||
|
return b.String()
|
||||||
|
}
|
||||||
@@ -0,0 +1,233 @@
|
|||||||
|
package alerting
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"os"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
|
dto "github.com/prometheus/client_model/go"
|
||||||
|
|
||||||
|
"gitea.perlbach24.de/scripte/nexarch/internal/notify"
|
||||||
|
)
|
||||||
|
|
||||||
|
func setupTest(t *testing.T) (*RuleStore, *notify.Dispatcher, *pgxpool.Pool, func()) {
|
||||||
|
t.Helper()
|
||||||
|
adminDSN := os.Getenv("TEST_ADMIN_DSN")
|
||||||
|
if adminDSN == "" {
|
||||||
|
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
|
||||||
|
}
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
pool, err := pgxpool.New(ctx, adminDSN)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("pool: %v", err)
|
||||||
|
}
|
||||||
|
if _, err := pool.Exec(ctx, `
|
||||||
|
CREATE TABLE IF NOT EXISTS alert_rules (
|
||||||
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), metric_name TEXT NOT NULL,
|
||||||
|
comparison TEXT NOT NULL CHECK (comparison IN ('gt','lt')), threshold DOUBLE PRECISION NOT NULL,
|
||||||
|
label_filters JSONB NOT NULL DEFAULT '{}'::jsonb, recipient TEXT NOT NULL,
|
||||||
|
description TEXT NOT NULL DEFAULT '', created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
CREATE TABLE IF NOT EXISTS alert_debounce_state (
|
||||||
|
rule_key TEXT PRIMARY KEY, last_fired_at TIMESTAMPTZ NOT NULL
|
||||||
|
);
|
||||||
|
CREATE TABLE IF NOT EXISTS notification_jobs (
|
||||||
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(), channel TEXT NOT NULL, recipient TEXT NOT NULL,
|
||||||
|
payload JSONB NOT NULL DEFAULT '{}'::jsonb, status TEXT NOT NULL DEFAULT 'pending', attempts INT NOT NULL DEFAULT 0,
|
||||||
|
max_attempts INT NOT NULL DEFAULT 5, next_attempt_at TIMESTAMPTZ NOT NULL DEFAULT now(), last_error TEXT,
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
`); err != nil {
|
||||||
|
t.Fatalf("schema: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
cleanup := func() {
|
||||||
|
_, _ = pool.Exec(ctx, `DELETE FROM alert_rules`)
|
||||||
|
_, _ = pool.Exec(ctx, `DELETE FROM alert_debounce_state`)
|
||||||
|
_, _ = pool.Exec(ctx, `DELETE FROM notification_jobs`)
|
||||||
|
pool.Close()
|
||||||
|
}
|
||||||
|
return NewRuleStore(pool), notify.NewDispatcher(pool), pool, cleanup
|
||||||
|
}
|
||||||
|
|
||||||
|
func gaugeFamily(name string, labels map[string]string, value float64) *dto.MetricFamily {
|
||||||
|
pairs := make([]*dto.LabelPair, 0, len(labels))
|
||||||
|
for k, v := range labels {
|
||||||
|
k, v := k, v
|
||||||
|
pairs = append(pairs, &dto.LabelPair{Name: &k, Value: &v})
|
||||||
|
}
|
||||||
|
n := name
|
||||||
|
return &dto.MetricFamily{
|
||||||
|
Name: &n,
|
||||||
|
Metric: []*dto.Metric{
|
||||||
|
{Label: pairs, Gauge: &dto.Gauge{Value: &value}},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Akzeptanzkriterium 1 + 2 / Pruefung 1 + 2: Überschreitung löst eine
|
||||||
|
// Benachrichtigung mit vollständigem Inhalt (Metrik/Tenant/Modul) aus.
|
||||||
|
func TestEvaluate_FiresAlertOnThresholdExceeded(t *testing.T) {
|
||||||
|
rules, dispatcher, pool, cleanup := setupTest(t)
|
||||||
|
defer cleanup()
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
rule, err := rules.CreateRule(ctx, Rule{
|
||||||
|
MetricName: "nexarch_core_error_rate", Comparison: ComparisonGreaterThan, Threshold: 0.05,
|
||||||
|
Recipient: "ops@acme.example", Description: "Fehlerrate zu hoch",
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("create rule: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
eval := NewEvaluator(rules, dispatcher, pool, time.Hour)
|
||||||
|
families := []*dto.MetricFamily{
|
||||||
|
gaugeFamily("nexarch_core_error_rate", map[string]string{"tenant": "acme", "module": "dms"}, 0.12),
|
||||||
|
}
|
||||||
|
|
||||||
|
fired, err := eval.Evaluate(ctx, families)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("evaluate: %v", err)
|
||||||
|
}
|
||||||
|
if len(fired) != 1 || fired[0].Skipped {
|
||||||
|
t.Fatalf("erwartet genau 1 tatsaechlich ausgeloesten alarm, habe %+v", fired)
|
||||||
|
}
|
||||||
|
if fired[0].RuleID != rule.ID {
|
||||||
|
t.Fatalf("rule id = %q, want %q", fired[0].RuleID, rule.ID)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Pruefung 2: Benachrichtigungsinhalt vollstaendig (Metrik/Tenant/Modul).
|
||||||
|
var payloadJSON []byte
|
||||||
|
if err := pool.QueryRow(ctx, `SELECT payload FROM notification_jobs LIMIT 1`).Scan(&payloadJSON); err != nil {
|
||||||
|
t.Fatalf("notification_jobs lesen: %v", err)
|
||||||
|
}
|
||||||
|
payload := string(payloadJSON)
|
||||||
|
for _, want := range []string{`"metric"`, `nexarch_core_error_rate`, `"tenant"`, `"acme"`, `"module"`, `"dms"`} {
|
||||||
|
if !strings.Contains(payload, want) {
|
||||||
|
t.Errorf("payload enthaelt nicht %q: %s", want, payload)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestEvaluate_DoesNotFireBelowThreshold(t *testing.T) {
|
||||||
|
rules, dispatcher, pool, cleanup := setupTest(t)
|
||||||
|
defer cleanup()
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
if _, err := rules.CreateRule(ctx, Rule{
|
||||||
|
MetricName: "nexarch_core_error_rate", Comparison: ComparisonGreaterThan, Threshold: 0.05,
|
||||||
|
Recipient: "ops@acme.example",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("create rule: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
eval := NewEvaluator(rules, dispatcher, pool, time.Hour)
|
||||||
|
families := []*dto.MetricFamily{
|
||||||
|
gaugeFamily("nexarch_core_error_rate", map[string]string{"tenant": "acme"}, 0.01),
|
||||||
|
}
|
||||||
|
|
||||||
|
fired, err := eval.Evaluate(ctx, families)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("evaluate: %v", err)
|
||||||
|
}
|
||||||
|
if len(fired) != 0 {
|
||||||
|
t.Fatalf("erwartet keinen alarm unterhalb des schwellwerts, habe %+v", fired)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Akzeptanzkriterium 3 / Pruefung 3: anhaltende Überschreitung erzeugt NICHT
|
||||||
|
// bei jeder Messung eine neue Benachrichtigung.
|
||||||
|
func TestEvaluate_DebouncesRepeatedFiring(t *testing.T) {
|
||||||
|
rules, dispatcher, pool, cleanup := setupTest(t)
|
||||||
|
defer cleanup()
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
if _, err := rules.CreateRule(ctx, Rule{
|
||||||
|
MetricName: "nexarch_core_error_rate", Comparison: ComparisonGreaterThan, Threshold: 0.05,
|
||||||
|
Recipient: "ops@acme.example",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("create rule: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Langes Debounce-Intervall: der zweite Evaluate-Lauf (simuliert die
|
||||||
|
// naechste Messung bei anhaltend ueberschrittenem Wert) darf keinen
|
||||||
|
// weiteren Job einreihen.
|
||||||
|
eval := NewEvaluator(rules, dispatcher, pool, time.Hour)
|
||||||
|
families := []*dto.MetricFamily{
|
||||||
|
gaugeFamily("nexarch_core_error_rate", map[string]string{"tenant": "acme"}, 0.5),
|
||||||
|
}
|
||||||
|
|
||||||
|
first, err := eval.Evaluate(ctx, families)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("erster evaluate-lauf: %v", err)
|
||||||
|
}
|
||||||
|
if len(first) != 1 || first[0].Skipped {
|
||||||
|
t.Fatalf("erster lauf haette feuern muessen, habe %+v", first)
|
||||||
|
}
|
||||||
|
|
||||||
|
second, err := eval.Evaluate(ctx, families)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("zweiter evaluate-lauf: %v", err)
|
||||||
|
}
|
||||||
|
if len(second) != 1 || !second[0].Skipped {
|
||||||
|
t.Fatalf("zweiter lauf haette gedrosselt werden muessen, habe %+v", second)
|
||||||
|
}
|
||||||
|
|
||||||
|
var count int
|
||||||
|
if err := pool.QueryRow(ctx, `SELECT count(*) FROM notification_jobs`).Scan(&count); err != nil {
|
||||||
|
t.Fatalf("notification_jobs zaehlen: %v", err)
|
||||||
|
}
|
||||||
|
if count != 1 {
|
||||||
|
t.Fatalf("erwartet genau 1 eingereihten job trotz zwei ueberschreitenden messungen, habe %d", count)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Verschiedene Zeitreihen derselben Regel (unterschiedlicher Tenant) werden
|
||||||
|
// unabhaengig voneinander gedrosselt.
|
||||||
|
func TestEvaluate_DebouncesIndependentlyPerLabelSet(t *testing.T) {
|
||||||
|
rules, dispatcher, pool, cleanup := setupTest(t)
|
||||||
|
defer cleanup()
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
if _, err := rules.CreateRule(ctx, Rule{
|
||||||
|
MetricName: "nexarch_core_error_rate", Comparison: ComparisonGreaterThan, Threshold: 0.05,
|
||||||
|
Recipient: "ops@acme.example",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("create rule: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
eval := NewEvaluator(rules, dispatcher, pool, time.Hour)
|
||||||
|
families := []*dto.MetricFamily{
|
||||||
|
{
|
||||||
|
Name: strPtr("nexarch_core_error_rate"),
|
||||||
|
Metric: []*dto.Metric{
|
||||||
|
metricWithLabel("tenant", "acme", 0.5),
|
||||||
|
metricWithLabel("tenant", "beta", 0.6),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
fired, err := eval.Evaluate(ctx, families)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("evaluate: %v", err)
|
||||||
|
}
|
||||||
|
if len(fired) != 2 {
|
||||||
|
t.Fatalf("erwartet 2 unabhaengige alarme (verschiedene tenants), habe %d", len(fired))
|
||||||
|
}
|
||||||
|
for _, a := range fired {
|
||||||
|
if a.Skipped {
|
||||||
|
t.Fatalf("beide tenants sollten beim ersten mal feuern, habe %+v", a)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func metricWithLabel(name, value string, gaugeValue float64) *dto.Metric {
|
||||||
|
n, v := name, value
|
||||||
|
return &dto.Metric{Label: []*dto.LabelPair{{Name: &n, Value: &v}}, Gauge: &dto.Gauge{Value: &gaugeValue}}
|
||||||
|
}
|
||||||
|
|
||||||
|
func strPtr(s string) *string { return &s }
|
||||||
@@ -0,0 +1,132 @@
|
|||||||
|
// Package alerting implementiert Core OPS-05: schwellwertbasierte
|
||||||
|
// Alarmierung auf den aus OPS-03 aggregierten Metriken, Zustellung über den
|
||||||
|
// Core-Benachrichtigungs-Dispatcher (CFG-02). Der Alertmanager-Gedanke von
|
||||||
|
// Prometheus/Grafana, aber auf das Nötigste reduziert (Schwellwert, Ziel,
|
||||||
|
// Drosselung) — keine eigene Ausdruckssprache.
|
||||||
|
package alerting
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Comparison legt fest, ob ein Schwellwert nach oben oder unten überwacht
|
||||||
|
// wird — bewusst nur zwei Operatoren, keine eigene Ausdruckssprache
|
||||||
|
// (Ticket-Produkt-DNA).
|
||||||
|
type Comparison string
|
||||||
|
|
||||||
|
const (
|
||||||
|
ComparisonGreaterThan Comparison = "gt"
|
||||||
|
ComparisonLessThan Comparison = "lt"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Rule ist eine Schwellwert-Regel auf einer beliebigen aggregierten Metrik
|
||||||
|
// (Akzeptanzkriterium 1). LabelFilters schränkt optional auf bestimmte
|
||||||
|
// Label-Werte ein (z. B. tenant/module), leer = alle Zeitreihen der Metrik.
|
||||||
|
type Rule struct {
|
||||||
|
ID string
|
||||||
|
MetricName string
|
||||||
|
Comparison Comparison
|
||||||
|
Threshold float64
|
||||||
|
LabelFilters map[string]string
|
||||||
|
Recipient string
|
||||||
|
Description string
|
||||||
|
}
|
||||||
|
|
||||||
|
// RuleStore verwaltet Alert-Regeln in der zentralen Registry-DB.
|
||||||
|
type RuleStore struct {
|
||||||
|
pool *pgxpool.Pool
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewRuleStore(pool *pgxpool.Pool) *RuleStore {
|
||||||
|
return &RuleStore{pool: pool}
|
||||||
|
}
|
||||||
|
|
||||||
|
// CreateRule legt eine neue Schwellwert-Regel an (Akzeptanzkriterium 1).
|
||||||
|
func (s *RuleStore) CreateRule(ctx context.Context, r Rule) (Rule, error) {
|
||||||
|
if r.MetricName == "" || r.Recipient == "" {
|
||||||
|
return Rule{}, errors.New("alerting: metricName und recipient duerfen nicht leer sein")
|
||||||
|
}
|
||||||
|
if r.Comparison != ComparisonGreaterThan && r.Comparison != ComparisonLessThan {
|
||||||
|
return Rule{}, fmt.Errorf("alerting: unbekannter comparison-operator %q", r.Comparison)
|
||||||
|
}
|
||||||
|
if r.LabelFilters == nil {
|
||||||
|
r.LabelFilters = map[string]string{}
|
||||||
|
}
|
||||||
|
filtersJSON, err := json.Marshal(r.LabelFilters)
|
||||||
|
if err != nil {
|
||||||
|
return Rule{}, fmt.Errorf("label-filter serialisieren: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = s.pool.QueryRow(ctx, `
|
||||||
|
INSERT INTO alert_rules (metric_name, comparison, threshold, label_filters, recipient, description)
|
||||||
|
VALUES ($1, $2, $3, $4, $5, $6)
|
||||||
|
RETURNING id
|
||||||
|
`, r.MetricName, string(r.Comparison), r.Threshold, filtersJSON, r.Recipient, r.Description).Scan(&r.ID)
|
||||||
|
if err != nil {
|
||||||
|
return Rule{}, fmt.Errorf("regel speichern: %w", err)
|
||||||
|
}
|
||||||
|
return r, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ListRules liefert alle konfigurierten Regeln — Grundlage für Evaluate.
|
||||||
|
func (s *RuleStore) ListRules(ctx context.Context) ([]Rule, error) {
|
||||||
|
rows, err := s.pool.Query(ctx, `
|
||||||
|
SELECT id, metric_name, comparison, threshold, label_filters, recipient, description
|
||||||
|
FROM alert_rules ORDER BY created_at
|
||||||
|
`)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("regeln auflisten: %w", err)
|
||||||
|
}
|
||||||
|
defer rows.Close()
|
||||||
|
|
||||||
|
var out []Rule
|
||||||
|
for rows.Next() {
|
||||||
|
var r Rule
|
||||||
|
var comparison string
|
||||||
|
var filtersJSON []byte
|
||||||
|
if err := rows.Scan(&r.ID, &r.MetricName, &comparison, &r.Threshold, &filtersJSON, &r.Recipient, &r.Description); err != nil {
|
||||||
|
return nil, fmt.Errorf("regel lesen: %w", err)
|
||||||
|
}
|
||||||
|
r.Comparison = Comparison(comparison)
|
||||||
|
if err := json.Unmarshal(filtersJSON, &r.LabelFilters); err != nil {
|
||||||
|
return nil, fmt.Errorf("label-filter lesen: %w", err)
|
||||||
|
}
|
||||||
|
out = append(out, r)
|
||||||
|
}
|
||||||
|
return out, rows.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
// DeleteRule entfernt eine Regel.
|
||||||
|
func (s *RuleStore) DeleteRule(ctx context.Context, id string) error {
|
||||||
|
_, err := s.pool.Exec(ctx, `DELETE FROM alert_rules WHERE id = $1`, id)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("regel loeschen: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// debounceStore kapselt die Drosselungs-Zustandstabelle (Akzeptanzkriterium 3).
|
||||||
|
type debounceStore struct {
|
||||||
|
pool *pgxpool.Pool
|
||||||
|
}
|
||||||
|
|
||||||
|
// shouldFire prueft, ob seit dem letzten Alarm fuer ruleKey mindestens
|
||||||
|
// interval vergangen ist — atomar ueber eine bedingte UPDATE/INSERT-
|
||||||
|
// Sequenz, damit zwei gleichzeitige Evaluate-Laeufe (z. B. bei mehreren
|
||||||
|
// Core-Instanzen) nicht beide gleichzeitig alarmieren.
|
||||||
|
func (d *debounceStore) shouldFire(ctx context.Context, ruleKey string, intervalSeconds float64) (bool, error) {
|
||||||
|
tag, err := d.pool.Exec(ctx, `
|
||||||
|
INSERT INTO alert_debounce_state (rule_key, last_fired_at) VALUES ($1, now())
|
||||||
|
ON CONFLICT (rule_key) DO UPDATE SET last_fired_at = now()
|
||||||
|
WHERE alert_debounce_state.last_fired_at <= now() - ($2 || ' seconds')::interval
|
||||||
|
`, ruleKey, intervalSeconds)
|
||||||
|
if err != nil {
|
||||||
|
return false, fmt.Errorf("drosselungszustand pruefen: %w", err)
|
||||||
|
}
|
||||||
|
return tag.RowsAffected() == 1, nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,2 @@
|
|||||||
|
DROP TABLE alert_debounce_state;
|
||||||
|
DROP TABLE alert_rules;
|
||||||
@@ -0,0 +1,24 @@
|
|||||||
|
-- OPS-05: Schwellwert-Regeln fuer Alerting auf den aus OPS-03 aggregierten
|
||||||
|
-- Metriken. Lebt wie config_values/notification_jobs (CFG-01/02) in der
|
||||||
|
-- zentralen Registry-DB — modulübergreifende Betriebskonfiguration, keine
|
||||||
|
-- Mandanten-Geschaeftsdaten.
|
||||||
|
CREATE TABLE alert_rules (
|
||||||
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
||||||
|
metric_name TEXT NOT NULL,
|
||||||
|
comparison TEXT NOT NULL CHECK (comparison IN ('gt', 'lt')),
|
||||||
|
threshold DOUBLE PRECISION NOT NULL,
|
||||||
|
label_filters JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||||
|
recipient TEXT NOT NULL,
|
||||||
|
description TEXT NOT NULL DEFAULT '',
|
||||||
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
|
||||||
|
-- Haelt fest, wann eine Regel zuletzt tatsaechlich einen Alarm ausgeloest
|
||||||
|
-- hat (Akzeptanzkriterium 3: Drosselung wiederholter Alarmierung fuer
|
||||||
|
-- denselben anhaltenden Zustand). rule_key kombiniert Regel-ID mit den
|
||||||
|
-- tatsaechlichen Label-Werten der ausloesenden Zeitreihe, damit dieselbe
|
||||||
|
-- Regel fuer unterschiedliche Tenants/Module unabhaengig gedrosselt wird.
|
||||||
|
CREATE TABLE alert_debounce_state (
|
||||||
|
rule_key TEXT PRIMARY KEY,
|
||||||
|
last_fired_at TIMESTAMPTZ NOT NULL
|
||||||
|
);
|
||||||
Reference in New Issue
Block a user