За другою рецензією:
* 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>
450 lines
18 KiB
Go
450 lines
18 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.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)
|
||
}
|
||
}
|
||
}
|
||
}
|