Система вміла малювати мапу, але мовчала, коли щось падало. Схема alr.* лежала готовою з Етапу 1 і повністю порожньою. Движок: - обчислює правила з icmp / interface / metric / no_data і не тримає стану між тіками: вікно for_seconds — це запит по часу до TSDB, тож перезапуск нічого не збиває - дві семантики вікна: без agg умова має триматися всі виміри (антифлап), з agg порівнюється агрегат - кореляція за топологією: причина аварії — той, у кого лишився живий сусід; хто оточений мертвими, той наслідок. Заради цього й будувалась topo.links: інакше падіння маршрутизатора дає сорок сповіщень - вікна обслуговування, ручне заглушення зі стелею в тиждень - advisory-блокування на тік: кілька API за балансувальником безпечні Доставка: - Telegram, webhook, SMTP; токени в core.secrets тим самим кільцем, що й паролі від обладнання - тенант без маршрутів отримує все в усі придатні канали — підключили Telegram, має працювати - тиха година не глушить disaster - вебхуки на внутрішні адреси заблоковано; для self-hosted знімається прапорцем процесу, бо там ця мережа своя UI: індикатор у шапці, панель зі списком, ack/mute/close, окремий фільтр придушених, спільна шина подій замість другого WebSocket. Знайдено живою роботою й виправлено: - вимкнене або видалене правило лишало алерти сиротами назавжди - перехід у suppressed не публікував події — UI дізнавався лише після перезавантаження сторінки - кнопки дій на телефоні були 26 px 20 нових тестів, go vet і tsc чисто. Живий прогін: поріг посередині розділив два справжні пристрої, 5 доставок на 5 подій без повторів. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
319 lines
10 KiB
Go
319 lines
10 KiB
Go
// Package alerting — движок правил і доставка сповіщень.
|
||
//
|
||
// Движок навмисно не тримає стану між тіками. Усе, що потрібно для
|
||
// рішення «піднімати чи ні», обчислюється з даних у TSDB: вікно
|
||
// for_seconds — це запит по часу, а не лічильник у пам'яті. Тому
|
||
// перезапуск процесу нічого не збиває, а два процеси, що випадково
|
||
// працюють одночасно, дадуть однаковий результат.
|
||
package alerting
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"log/slog"
|
||
"time"
|
||
|
||
"github.com/netpulse/netpulse/server/internal/crypto"
|
||
"github.com/netpulse/netpulse/server/internal/store"
|
||
)
|
||
|
||
// Порядок серйозності — той самий, що в alr.severity.
|
||
var severityRank = map[string]int{
|
||
"info": 0, "warning": 1, "average": 2, "high": 3, "disaster": 4,
|
||
}
|
||
|
||
type Engine struct {
|
||
st *store.Store
|
||
ring *crypto.Keyring
|
||
log *slog.Logger
|
||
notifier *Notifier
|
||
|
||
interval time.Duration
|
||
}
|
||
|
||
func New(st *store.Store, ring *crypto.Keyring, log *slog.Logger,
|
||
interval time.Duration, allowPrivateHooks bool) *Engine {
|
||
|
||
if interval <= 0 {
|
||
interval = 30 * time.Second
|
||
}
|
||
return &Engine{
|
||
st: st,
|
||
ring: ring,
|
||
log: log.With("component", "alerting"),
|
||
notifier: NewNotifier(st, log, allowPrivateHooks),
|
||
interval: interval,
|
||
}
|
||
}
|
||
|
||
// advisoryLockKey — довільна стала; важливо лише, щоб її не займав
|
||
// ніхто інший у цій же БД.
|
||
const advisoryLockKey = 0x6e70_616c // "npal"
|
||
|
||
// Run крутить цикл обчислення до скасування контексту.
|
||
//
|
||
// Перед кожним тіком береться advisory-блокування Postgres. Кілька
|
||
// серверів за балансувальником — норма, але правила має обчислювати
|
||
// рівно один: дедуплікація алертів захищена індексом, а от сповіщення
|
||
// пішли б у кількох копіях, і людина отримала б три однакові
|
||
// повідомлення про одну аварію.
|
||
func (e *Engine) Run(ctx context.Context) {
|
||
t := time.NewTicker(e.interval)
|
||
defer t.Stop()
|
||
|
||
e.log.Info("движок алертів запущено", "інтервал", e.interval)
|
||
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
e.log.Info("движок алертів зупинено")
|
||
return
|
||
case <-t.C:
|
||
start := time.Now()
|
||
n, err := e.tick(ctx)
|
||
if err != nil {
|
||
e.log.Error("тік алертів", "помилка", err)
|
||
continue
|
||
}
|
||
if n > 0 {
|
||
e.log.Debug("тік алертів", "правил", n, "тривалість", time.Since(start))
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
func (e *Engine) tick(ctx context.Context) (int, error) {
|
||
conn, err := e.st.Pool().Acquire(ctx)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
defer conn.Release()
|
||
|
||
var got bool
|
||
if err := conn.QueryRow(ctx, `SELECT pg_try_advisory_lock($1)`, int64(advisoryLockKey)).Scan(&got); err != nil {
|
||
return 0, err
|
||
}
|
||
if !got {
|
||
// Обчислює інший інстанс — це штатний стан, не помилка.
|
||
return 0, nil
|
||
}
|
||
defer func() {
|
||
_, _ = conn.Exec(context.WithoutCancel(ctx),
|
||
`SELECT pg_advisory_unlock($1)`, int64(advisoryLockKey))
|
||
}()
|
||
|
||
rules, err := e.st.ActiveRules(ctx)
|
||
if err != nil {
|
||
return 0, fmt.Errorf("читання правил: %w", err)
|
||
}
|
||
if len(rules) == 0 {
|
||
return 0, nil
|
||
}
|
||
|
||
// Правила групуються за тенантом, бо топологія, вікна обслуговування
|
||
// й канали читаються раз на тенант, а не раз на правило.
|
||
byTenant := map[string][]store.Rule{}
|
||
for _, r := range rules {
|
||
byTenant[r.TenantID] = append(byTenant[r.TenantID], r)
|
||
}
|
||
|
||
for tenantID, rs := range byTenant {
|
||
if err := e.tickTenant(ctx, tenantID, rs); err != nil {
|
||
// Зламаний тенант не має зупиняти решту: у SaaS це означало б,
|
||
// що одне криве правило одного клієнта гасить моніторинг усім.
|
||
e.log.Error("тік тенанта", "tenant", tenantID, "помилка", err)
|
||
}
|
||
}
|
||
return len(rules), nil
|
||
}
|
||
|
||
func (e *Engine) tickTenant(ctx context.Context, tenantID string, rules []store.Rule) error {
|
||
sup, err := e.st.LoadSuppression(ctx, tenantID)
|
||
if err != nil {
|
||
return fmt.Errorf("вікна обслуговування: %w", err)
|
||
}
|
||
|
||
// Топологія потрібна лише якщо хоч одне правило її враховує.
|
||
var roots map[string]string
|
||
for _, r := range rules {
|
||
if r.DependsOnTopology {
|
||
g, err := e.st.LoadTopology(ctx, tenantID)
|
||
if err != nil {
|
||
return fmt.Errorf("топологія: %w", err)
|
||
}
|
||
roots = g.RootCause()
|
||
break
|
||
}
|
||
}
|
||
|
||
var fired, resolved, changed []store.Alert
|
||
|
||
for _, r := range rules {
|
||
cands, err := e.st.EvaluateRule(ctx, r)
|
||
if err != nil {
|
||
e.log.Error("обчислення правила", "правило", r.Name, "помилка", err)
|
||
continue
|
||
}
|
||
|
||
keep := make([]string, 0, len(cands))
|
||
for _, c := range cands {
|
||
key := c.DedupKey(r.ID)
|
||
keep = append(keep, key)
|
||
|
||
reason := sup.For(c.DeviceID, r.ID)
|
||
if reason == "" && r.DependsOnTopology {
|
||
if _, collateral := roots[c.DeviceID]; collateral {
|
||
reason = "topology"
|
||
}
|
||
}
|
||
|
||
a := buildAlert(r, c, key, reason)
|
||
f, err := e.st.RaiseAlert(ctx, tenantID, a)
|
||
if err != nil {
|
||
e.log.Error("підняття алерту", "правило", r.Name, "помилка", err)
|
||
continue
|
||
}
|
||
switch {
|
||
case f.IsNew:
|
||
fired = append(fired, f.Alert)
|
||
case f.StateChanged():
|
||
// Наприклад, заглушення пристрою: проблема та сама, але
|
||
// на екрані алерт має піти в «придушені» без чекання на
|
||
// наступне перезавантаження сторінки.
|
||
changed = append(changed, f.Alert)
|
||
}
|
||
}
|
||
|
||
gone, err := e.st.ResolveMissing(ctx, tenantID, r.ID, keep)
|
||
if err != nil {
|
||
e.log.Error("закриття алертів", "правило", r.Name, "помилка", err)
|
||
continue
|
||
}
|
||
resolved = append(resolved, gone...)
|
||
}
|
||
|
||
e.publish(ctx, tenantID, fired, resolved, changed)
|
||
|
||
// Сповіщаються лише щойно підняті й не придушені: продовження вже
|
||
// відомої проблеми не є новиною, а придушене — за визначенням те,
|
||
// про що просили не турбувати.
|
||
var notify []store.Alert
|
||
for _, a := range fired {
|
||
if a.State == "firing" {
|
||
notify = append(notify, a)
|
||
}
|
||
}
|
||
if len(notify) > 0 {
|
||
e.notifier.Dispatch(ctx, tenantID, notify, e.ring)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func buildAlert(r store.Rule, c store.Candidate, dedupKey, suppressedBy string) store.Alert {
|
||
name := c.DeviceName
|
||
if name == "" {
|
||
name = c.DeviceID
|
||
}
|
||
title := fmt.Sprintf("%s: %s", name, r.Name)
|
||
if c.IfName != "" {
|
||
title = fmt.Sprintf("%s [%s]: %s", name, c.IfName, r.Name)
|
||
}
|
||
|
||
msg := describe(r, c)
|
||
value := c.Value
|
||
threshold := r.Condition.Value
|
||
|
||
meta, _ := json.Marshal(map[string]any{
|
||
"metric": r.Condition.Metric,
|
||
"op": r.Condition.Op,
|
||
"agg": r.Condition.Agg,
|
||
"samples": c.Samples,
|
||
"for": r.ForSeconds,
|
||
})
|
||
|
||
return store.Alert{
|
||
RuleID: r.ID,
|
||
RuleName: r.Name,
|
||
DeviceID: c.DeviceID,
|
||
DeviceName: c.DeviceName,
|
||
InterfaceID: c.InterfaceID,
|
||
Severity: r.Severity,
|
||
Title: title,
|
||
Message: msg,
|
||
DedupKey: dedupKey,
|
||
Value: &value,
|
||
Threshold: &threshold,
|
||
Context: meta,
|
||
SuppressedBy: suppressedBy,
|
||
}
|
||
}
|
||
|
||
// describe робить із умови людський текст.
|
||
//
|
||
// Повідомлення читає людина о третій ночі з телефона: «loss_pct > 20»
|
||
// там марне, а «втрати 100% при порозі 20% протягом 3 хв» — ні.
|
||
func describe(r store.Rule, c store.Candidate) string {
|
||
if r.Condition.Metric == "no_data" {
|
||
return fmt.Sprintf("даних немає понад %s (останні — %s)",
|
||
humanDur(r.ForSeconds), c.LastTS.Format("15:04:05"))
|
||
}
|
||
scope := "усі виміри"
|
||
if r.Condition.Agg != "" {
|
||
scope = r.Condition.Agg + " за"
|
||
}
|
||
return fmt.Sprintf("%s %s %s %g (поточне %.3g, %d вимірів за %s)",
|
||
r.Condition.Metric, scope, r.Condition.Op, r.Condition.Value,
|
||
c.Value, c.Samples, humanDur(r.ForSeconds))
|
||
}
|
||
|
||
func humanDur(sec int) string {
|
||
d := time.Duration(sec) * time.Second
|
||
switch {
|
||
case d >= time.Hour:
|
||
return fmt.Sprintf("%.0f год", d.Hours())
|
||
case d >= time.Minute:
|
||
return fmt.Sprintf("%.0f хв", d.Minutes())
|
||
default:
|
||
return fmt.Sprintf("%d с", sec)
|
||
}
|
||
}
|
||
|
||
// publish кладе події в outbox, щоб мапа й список алертів оновились без
|
||
// перезавантаження сторінки.
|
||
func (e *Engine) publish(ctx context.Context, tenantID string, fired, resolved, changed []store.Alert) {
|
||
emit := func(topic string, as []store.Alert) {
|
||
for _, a := range as {
|
||
if err := e.st.PublishEvent(ctx, tenantID, topic, map[string]any{
|
||
"alert_id": a.ID, "device_id": a.DeviceID, "severity": a.Severity,
|
||
"title": a.Title, "state": a.State, "suppressed_by": a.SuppressedBy,
|
||
}); err != nil {
|
||
e.log.Error("подія "+topic, "помилка", err)
|
||
}
|
||
}
|
||
}
|
||
emit("alert.fired", fired)
|
||
emit("alert.resolved", resolved)
|
||
emit("alert.updated", changed)
|
||
}
|
||
|
||
// RunHousekeeping переносить закриті алерти в історію.
|
||
func (e *Engine) RunHousekeeping(ctx context.Context, keepResolved time.Duration) {
|
||
t := time.NewTicker(15 * time.Minute)
|
||
defer t.Stop()
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-t.C:
|
||
n, err := e.st.ArchiveResolved(ctx, keepResolved)
|
||
if err != nil {
|
||
e.log.Error("архівація алертів", "помилка", err)
|
||
continue
|
||
}
|
||
if n > 0 {
|
||
e.log.Info("алерти заархівовано", "рядків", n)
|
||
}
|
||
}
|
||
}
|
||
}
|