195 lines
5.3 KiB
Go
195 lines
5.3 KiB
Go
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()
|
|
}
|