// 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.WorkerPool().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)) }() // Дві дії нижче стосуються подієвих алертів, які піднімає не цей // цикл, а приймачі подій (events.go). Вони мають статись навіть у // кабінеті без жодного метричного правила, тому стоять до вибірки // й до перевірки на порожньо. e.expireEvents(ctx) e.deliverPending(ctx) 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 { // Подієві джерела цей цикл не чіпає — і не «просто пропускає», // а мусить пропустити. Обчислення дало б порожній список // кандидатів, а ResolveMissing слідом закрив би щойно піднятий // подієвий алерт: із погляду опитування він «зник», хоча // зникнути він не може за побудовою. if store.IsEventSource(r.Source) { continue } 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, unsuppressed []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) if leftSuppression(f.PrevState, f.State) { // Вікно обслуговування скінчилось, а проблема — ні. // // Досі такий алерт лише перемальовувався на екрані: // сповіщення при народженні не пішло (бо заглушено), // а тут не йшло, бо «не новий». Тобто аварія, яка // почалась о 03:00 усередині вікна 02:30–03:30, не // будила нікого й ніколи — ні першим сповіщенням, ні // драбиною. Це та сама тиха відмова, тільки з // поважним на вигляд приводом. unsuppressed = append(unsuppressed, 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) } } // Ті, з кого щойно зняли заглушення, сповіщаються нарівні з новими: // для людини це перша звістка про проблему, хоч би скільки вона вже // тривала за зачиненими дверима. notify = append(notify, unsuppressed...) 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) } // expireEvents гасить подієві алерти, до яких давно не було подій. // // Робиться щотіку, а не в прибиральнику раз на чверть години: строк // життя правило задає в секундах, і «доба» з точністю до п'ятнадцяти // хвилин виглядала б на екрані як несправність. func (e *Engine) expireEvents(ctx context.Context) { n, err := e.st.ExpireEventAlerts(ctx) if err != nil { e.log.Error("гасіння подієвих алертів", "помилка", err) return } if n > 0 { e.log.Info("подієві алерти прострочено", "рядків", n) } } // deliverPending розсилає алерти, підняті подієвим шляхом. // // Подія приходить у процес, який не має ані ключів шифрування каналів, // ані маршрутів, ані тихих годин — тож розсилка лишається тут, де все // це вже прочитано, і під тим самим advisory-блокуванням: інакше два // інстанси API розбудили б людину двічі. // // Плата — затримка до одного тіку. Для «конфіг змінився» чи «порушено // стандарт» це прийнятно: жодне з них не є аварією, на яку біжать за // секунди. Для метричних алертів затримки як була, так і немає. func (e *Engine) deliverPending(ctx context.Context) { pending, err := e.st.TakeNotifyPending(ctx, 200) if err != nil { e.log.Error("черга розсилки подієвих алертів", "помилка", err) return } if len(pending) == 0 { return } byTenant := map[string][]store.Alert{} for _, a := range pending { // Придушене не турбує нікого — рівно як у метричному шляху. if a.State != "firing" { continue } byTenant[a.TenantID] = append(byTenant[a.TenantID], a) } for tenantID, as := range byTenant { e.notifier.Dispatch(ctx, tenantID, as, e.ring) } } // escalationLogKeep — той самий строк, що в alr.notifications (0007). // Розходження тут означало б, що на питання «чому мене розбудили» одна // половина відповіді ще є, а друга вже стерта. const escalationLogKeep = 90 * 24 * time.Hour // leftSuppression — чи алерт щойно вийшов із заглушення в бойовий стан. // // Винесено окремою функцією не заради краси: це рішення про те, кого // розбудити, а перевірити його всередині тіку можна лише піднявши базу, // вікно обслуговування й годинник. Пари станів тут коштують чийогось // сну в обидва боки — і «не сповістили, бо вважали продовженням», і // «сповістили вдруге про те саме». func leftSuppression(prev, cur string) bool { return prev == "suppressed" && cur == "firing" } // RunEscalations — окремий такт для драбин ескалації. // // Окремий, а не всередині tick(), і це не косметика. Доставка сходинки // синхронна, кожен канал має свій таймаут (десятки секунд на мертвому // вебхуці), а сходинок у партії до сотні. Поки escalate() жив у тіку // движка, один кабінет із непрацюючим каналом зупиняв обчислення правил // УСІМ: нові аварії не піднімались, перші сповіщення не йшли. Тобто // механізм, який існує, щоб аварію точно помітили, робив аварії // непомітними. // // Своє блокування тут не потрібне: черга розбирається через // FOR UPDATE SKIP LOCKED плюс оренда рядка, тож два інстанси не візьмуть // ту саму сходинку — на відміну від обчислення правил, яке саме тому й // сидить під advisory-блокуванням. func (e *Engine) RunEscalations(ctx context.Context) { t := time.NewTicker(e.interval) defer t.Stop() for { select { case <-ctx.Done(): return case <-t.C: e.escalate(ctx) } } } // 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) } // Журнал сходинок ескалації не гіпертаблиця, тож політики // ретеншену TimescaleDB в нього немає — прибираємо тут, тим // самим строком, що й у журналу доставки (0007). if k, err := e.st.PurgeEscalationLog(ctx, escalationLogKeep); err != nil { e.log.Error("прибирання журналу ескалацій", "помилка", err) } else if k > 0 { e.log.Info("журнал ескалацій прибрано", "рядків", k) } } } }