Netpulse_SasS/server/internal/alerting/engine.go
byrsapty 954be1d643
All checks were successful
CI / hygiene (push) Successful in 9s
CI / web (push) Successful in 1m15s
CI / server (push) Successful in 1m34s
CI / agent (push) Successful in 2m59s
Ескалації: закриваю те, що минулого разу закрив наполовину
За другою рецензією:

* scripts/dbtest.sh писав у шапці «не напрямляйте на робочу базу» й
  нічого для цього не робив — перевірено, пішов котити міграції на базу
  з бойовим іменем. Тепер вимагає probe/test в імені.
* sendText ковтав помилку, тож журнал ескалацій писав «надіслано» на
  сходинці, жодне повідомлення якої не дійшло. Три результати замість
  двох: no_channels, failed, sent.
* stopped_at IS NULL рятував лише від ack; гасіння правилом і
  ResolveMissing рядка драбини не чіпають, і сходинка дзвонила за
  погашеним алертом. Додано перевірку стану алерту в тому ж UPDATE.
* escalate() блокував весь тік движка — мертвий вебхук одного кабінету
  зупиняв обчислення правил усім. Винесено в RunEscalations.
* алерт, народжений під заглушенням, не сповіщався ніколи: ні при
  народженні, ні коли вікно скінчилось. Тепер перехід suppressed→firing
  сповіщається, а драбина рахує час від першого сповіщення.
* alr.rules.channel_ids приймав чужі канали, глушачи і сповіщення, і
  драбину. Перевірка як для сходинок; DeleteChannel чистить посилання.

І перше, що зловив прогін проти справжньої бази: nil-зріз каналів їде
явним NULL повз DEFAULT '{}' — правило без каналів давало 500.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 17:44:38 +03:00

450 lines
18 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// 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:3003: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)
}
}
}
}