Подія відповідності несла шкалу critical/high/medium, а поріг тригера
порівнювався шкалою алертів, де SeverityRank("critical") = 0. Тригер із
порогом «warning» пропускав середнє й відкидав найважче — тихо.
Подія тепер народжується в шкалі алертів; умова тригера приймає обидві.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
826 lines
32 KiB
Go
826 lines
32 KiB
Go
package alerting
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"log/slog"
|
||
"net"
|
||
"regexp"
|
||
"strings"
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/netpulse/netpulse/server/internal/store"
|
||
)
|
||
|
||
// Подієві алерти.
|
||
//
|
||
// Движок у engine.go працює тактами: раз на пів хвилини перепитує ряди
|
||
// вимірів і питає в них «чи виконується умова зараз». Для трьох джерел
|
||
// (metric, icmp, interface) це єдиний можливий спосіб — ряд є, питання
|
||
// осмислене, відповідь може змінитися будь-коли.
|
||
//
|
||
// Для журналу, конфігів і відповідності такого ряду немає. Питати
|
||
// «чи змінився конфіг зараз» безглуздо: він змінився о 10:42 і більше
|
||
// нічого про це не скаже. Тому ці джерела не опитуються взагалі —
|
||
// правило перевіряється рівно в ту мить, коли подія надійшла, у тому
|
||
// процесі, який її прийняв.
|
||
//
|
||
// Наслідки цієї різниці треба було вирішити явно, і вони вирішені так:
|
||
//
|
||
// дедуплікація — один алерт на пару «правило + хост», скільки б подій
|
||
// під нього не підпало. Ключ не містить нічого від самої події;
|
||
// замість переліку в алерті лічильник. Інакше потік syslog за
|
||
// хвилину зробив би дошку алертів нечитабельною — тобто зламав би
|
||
// саме те, заради чого вона є.
|
||
//
|
||
// частота — між двома зверненнями до одного алерту витримується
|
||
// min_interval_seconds правила. Пропущені за цей час події не
|
||
// губляться: вони накопичуються тут і доливаються в лічильник
|
||
// наступним зверненням. Ми економимо запити до бази, а не правду.
|
||
//
|
||
// гасіння — подієвий алерт не зникає сам, бо зникати нема чому.
|
||
// Його закриває або людина, або строк (ExpireEventAlerts, стан
|
||
// `expired`). Виняток один — відповідність: там прогін, у якому
|
||
// хост правило пройшов, і є чесний сигнал «більше не порушено».
|
||
//
|
||
// доставка — подія приходить у netpulse-server, а канали, маршрути й
|
||
// тихі години живуть у netpulse-api. Тому тут алерт лише
|
||
// піднімається з позначкою notify_pending, а розсилає його движок
|
||
// наступним тіком — під тим самим advisory-блокуванням, тобто в
|
||
// одному екземплярі.
|
||
|
||
// SyslogEvent — рядок журналу у вигляді, потрібному правилам.
|
||
//
|
||
// Власний тип, а не protobuf: пакет алертів не має знати про транспорт
|
||
// зондів, інакше кожна зміна .proto тягла б за собою правку движка.
|
||
type SyslogEvent struct {
|
||
DeviceID string
|
||
Message string
|
||
Tag string
|
||
Severity int
|
||
}
|
||
|
||
// ConfigEvent — те, що сталося з конфігом хоста.
|
||
//
|
||
// Kind: "changed" — приїхала версія, відмінна від попередньої;
|
||
// "backup_failed" — збір не вдався.
|
||
type ConfigEvent struct {
|
||
DeviceID string
|
||
ConfigType string
|
||
Kind string
|
||
Detail string
|
||
}
|
||
|
||
// ComplianceEvent — результат перевірки одного правила на одному хості.
|
||
type ComplianceEvent struct {
|
||
RuleID string
|
||
RuleName string
|
||
Severity string
|
||
DeviceID string
|
||
Passed bool
|
||
Line string
|
||
LineNumber int
|
||
}
|
||
|
||
// EventSink приймає події й піднімає за ними алерти.
|
||
//
|
||
// Безпечний для конкурентного використання: приймач журналу викликає
|
||
// його з кожного стріму зонда.
|
||
type EventSink struct {
|
||
st *store.Store
|
||
log *slog.Logger
|
||
|
||
// Як довго живе кеш правил і хостів тенанта.
|
||
//
|
||
// Кеш тут не оптимізація, а умова існування: без нього кожен рядок
|
||
// журналу коштував би читання правил, розгортання селектора й
|
||
// вибірки вікон обслуговування. Ціна — щойно створене правило
|
||
// починає діяти не миттєво, і це чесний розмін: подія, яка сталася
|
||
// за півхвилини до появи правила, і так під нього не підпадає.
|
||
ttl time.Duration
|
||
|
||
mu sync.Mutex
|
||
cache map[string]*tenantView
|
||
rate map[string]*rateEntry
|
||
}
|
||
|
||
func NewEventSink(st *store.Store, log *slog.Logger) *EventSink {
|
||
if log == nil {
|
||
log = slog.Default()
|
||
}
|
||
return &EventSink{
|
||
st: st,
|
||
log: log.With("component", "alerting.events"),
|
||
ttl: 30 * time.Second,
|
||
cache: map[string]*tenantView{},
|
||
rate: map[string]*rateEntry{},
|
||
}
|
||
}
|
||
|
||
// tenantView — усе, що потрібно знати про кабінет, щоб вирішити долю
|
||
// події, не звертаючись до бази.
|
||
type tenantView struct {
|
||
at time.Time
|
||
rules []compiledRule
|
||
devices map[string]string // device_id → ім'я
|
||
sup store.Suppression
|
||
// Власний словник трапів кабінету. Потрібен лише для тексту
|
||
// алерту: «linkDown на sw-core-01» замість
|
||
// «1.3.6.1.6.3.1.1.5.3 на sw-core-01». Читається лише коли в
|
||
// кабінеті є хоч одне правило на трапи — зайвий запит раз на пів
|
||
// хвилини на кожного клієнта, який трапами не користується, нічого
|
||
// не вартий рівно доти, доки клієнтів мало.
|
||
trapNames map[string]store.TrapMeaning
|
||
}
|
||
|
||
type compiledRule struct {
|
||
rule store.Rule
|
||
re *regexp.Regexp
|
||
// Хости під селектором. nil означає «усі»: порожній селектор — це
|
||
// найчастіший випадок, і перетворювати його на перелік означало б
|
||
// щоразу відставати від щойно доданого хоста.
|
||
scope map[string]bool
|
||
}
|
||
|
||
func (c compiledRule) covers(deviceID string) bool {
|
||
return c.scope == nil || c.scope[deviceID]
|
||
}
|
||
|
||
// rateEntry — стан обмежувача частоти для одного алерту.
|
||
type rateEntry struct {
|
||
last time.Time
|
||
// Події, що надійшли, поки діяв проміжок. Не викидаються: людині
|
||
// важлива не кожна з них окремо, а те, що їх було багато.
|
||
carried int
|
||
touched time.Time
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Журнал
|
||
// ---------------------------------------------------------------------
|
||
|
||
// OnSyslog звіряє пачку рядків журналу з правилами джерела `syslog`.
|
||
//
|
||
// Зведення робиться до звернення до бази: пачка від зонда — це сотні
|
||
// рядків, і половина з них зазвичай про один і той самий порт, що
|
||
// мигає. Одна подія на пару «правило + хост» замість сотні запитів —
|
||
// різниця між приймачем, який справляється, і тим, який гальмує самі
|
||
// зонди.
|
||
func (s *EventSink) OnSyslog(ctx context.Context, tenantID string, events []SyslogEvent) {
|
||
if len(events) == 0 {
|
||
return
|
||
}
|
||
view := s.view(ctx, tenantID)
|
||
if view == nil {
|
||
return
|
||
}
|
||
|
||
// ключ пари «правило+хост» → скільки збігів і останній текст
|
||
type hit struct {
|
||
rule compiledRule
|
||
device string
|
||
count int
|
||
last string
|
||
tag string
|
||
sevSeen int
|
||
}
|
||
hits := map[string]*hit{}
|
||
|
||
for _, ev := range events {
|
||
if ev.DeviceID == "" || ev.Message == "" {
|
||
// Подія з невідомої адреси не належить нікому. Піднімати
|
||
// алерт «десь у мережі щось сталося» — гірше, ніж мовчати:
|
||
// з ним нічого не можна зробити.
|
||
continue
|
||
}
|
||
for _, c := range view.rules {
|
||
if c.rule.Source != "syslog" || c.re == nil || !c.covers(ev.DeviceID) {
|
||
continue
|
||
}
|
||
if lte := c.rule.Condition.SeverityLTE; lte != nil && ev.Severity > *lte {
|
||
continue
|
||
}
|
||
if t := c.rule.Condition.Tag; t != "" && !strings.EqualFold(t, ev.Tag) {
|
||
continue
|
||
}
|
||
if !c.re.MatchString(ev.Message) {
|
||
continue
|
||
}
|
||
key := store.EventDedupKey(c.rule.ID, ev.DeviceID)
|
||
h, ok := hits[key]
|
||
if !ok {
|
||
h = &hit{rule: c, device: ev.DeviceID, sevSeen: ev.Severity}
|
||
hits[key] = h
|
||
}
|
||
h.count++
|
||
h.last = ev.Message
|
||
h.tag = ev.Tag
|
||
}
|
||
}
|
||
|
||
for _, h := range hits {
|
||
meta := map[string]any{
|
||
"kind": "syslog",
|
||
"pattern": h.rule.rule.Condition.Regex,
|
||
"sample": trimLine(h.last),
|
||
"tag": h.tag,
|
||
}
|
||
s.raise(ctx, tenantID, view, h.rule, h.device,
|
||
trimLine(h.last), h.count, meta)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Конфіги
|
||
// ---------------------------------------------------------------------
|
||
|
||
// OnConfig піднімає алерти правил джерела `ncm`.
|
||
//
|
||
// Саме той сценарій, заради якого все це писалося: людина заводить
|
||
// тригер «конфіг змінився», і він має спрацювати тоді, коли конфіг
|
||
// змінився, — а не ніколи.
|
||
func (s *EventSink) OnConfig(ctx context.Context, tenantID string, ev ConfigEvent) {
|
||
if ev.DeviceID == "" || ev.Kind == "" {
|
||
return
|
||
}
|
||
view := s.view(ctx, tenantID)
|
||
if view == nil {
|
||
return
|
||
}
|
||
|
||
for _, c := range view.rules {
|
||
if c.rule.Source != "ncm" || c.rule.Condition.Event != ev.Kind || !c.covers(ev.DeviceID) {
|
||
continue
|
||
}
|
||
msg := "конфіг змінився (" + orDefault(ev.ConfigType, "running") + ")"
|
||
if ev.Kind == "backup_failed" {
|
||
msg = "збір конфігу не вдався: " + trimLine(ev.Detail)
|
||
}
|
||
meta := map[string]any{
|
||
"kind": "ncm",
|
||
"event": ev.Kind,
|
||
"config_type": ev.ConfigType,
|
||
"detail": trimLine(ev.Detail),
|
||
}
|
||
s.raise(ctx, tenantID, view, c, ev.DeviceID, msg, 1, meta)
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Відповідність
|
||
// ---------------------------------------------------------------------
|
||
|
||
// OnCompliance переносить результат прогону у стан алертів.
|
||
//
|
||
// Єдине з подієвих джерел, у якого є зворотний бік. Прогін перевіряє
|
||
// всі хости під правилом і каже про кожен «пройшов» або «ні» — отже,
|
||
// «пройшов» і є той самий сигнал зняття, якого немає в журналі. Тому
|
||
// тут алерт закривається сам, і це не виняток із правила, а наслідок
|
||
// того, що дані інші.
|
||
func (s *EventSink) OnCompliance(ctx context.Context, tenantID string, events []ComplianceEvent) {
|
||
if len(events) == 0 {
|
||
return
|
||
}
|
||
view := s.view(ctx, tenantID)
|
||
if view == nil {
|
||
return
|
||
}
|
||
|
||
var healed []string
|
||
for _, ev := range events {
|
||
if ev.DeviceID == "" {
|
||
continue
|
||
}
|
||
for _, c := range view.rules {
|
||
if c.rule.Source != "compliance" || !c.covers(ev.DeviceID) {
|
||
continue
|
||
}
|
||
if !matchesComplianceRule(c.rule.Condition, ev) {
|
||
continue
|
||
}
|
||
key := store.EventDedupKey(c.rule.ID, ev.DeviceID)
|
||
if ev.Passed {
|
||
healed = append(healed, key)
|
||
continue
|
||
}
|
||
msg := fmt.Sprintf("порушено вимогу «%s»", ev.RuleName)
|
||
if ev.Line != "" {
|
||
msg = fmt.Sprintf("%s: рядок %d — %s", msg, ev.LineNumber, trimLine(ev.Line))
|
||
}
|
||
meta := map[string]any{
|
||
"kind": "compliance",
|
||
"compliance_rule": ev.RuleName,
|
||
"compliance_id": ev.RuleID,
|
||
"line": trimLine(ev.Line),
|
||
"line_number": ev.LineNumber,
|
||
"finding_severity": ev.Severity,
|
||
}
|
||
s.raise(ctx, tenantID, view, c, ev.DeviceID, msg, 1, meta)
|
||
}
|
||
}
|
||
|
||
if len(healed) > 0 {
|
||
if _, err := s.st.ResolveEventAlerts(ctx, tenantID, healed,
|
||
"хост пройшов перевірку відповідності"); err != nil {
|
||
s.log.Error("закриття алертів відповідності", "tenant", tenantID, "помилка", err)
|
||
}
|
||
}
|
||
}
|
||
|
||
// matchesComplianceRule звужує тригер до частини знахідок.
|
||
//
|
||
// Порожня умова означає «будь-яке порушення»: тригер «скажи мені, коли
|
||
// щось поїхало» — найчастіший і найкорисніший, і вимагати для нього
|
||
// переліку правил означало б, що новий стандарт, доданий завтра, під
|
||
// нього не підпаде.
|
||
func matchesComplianceRule(cond store.Condition, ev ComplianceEvent) bool {
|
||
if len(cond.RuleIDs) > 0 {
|
||
var found bool
|
||
for _, id := range cond.RuleIDs {
|
||
if id == ev.RuleID {
|
||
found = true
|
||
break
|
||
}
|
||
}
|
||
if !found {
|
||
return false
|
||
}
|
||
}
|
||
// Обидві сторони зводяться до шкали алертів: подія може прийти зі
|
||
// шкалою відповідності (critical/medium), і поріг людина теж могла
|
||
// написати нею. Без зведення «critical» отримував ранг нуль і
|
||
// відкидався порогом «warning» — тобто фільтр працював навпаки.
|
||
if cond.MinSeverity != "" &&
|
||
store.SeverityRank(store.NormalizeEventSeverity(ev.Severity)) <
|
||
store.SeverityRank(store.NormalizeEventSeverity(cond.MinSeverity)) {
|
||
return false
|
||
}
|
||
return true
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// SNMP-трапи
|
||
// ---------------------------------------------------------------------
|
||
|
||
// TrapVarbind — одне поле трапа у вигляді, потрібному правилам.
|
||
type TrapVarbind struct {
|
||
OID string
|
||
Value string
|
||
}
|
||
|
||
// TrapEvent — трап, зведений до того, про що можна запитати в умові.
|
||
//
|
||
// Власний тип, а не protobuf: пакет алертів не має знати про транспорт
|
||
// зондів. DeviceID порожній, якщо адресу відправника не вдалося
|
||
// зіставити з хостом, — і це не помилка, див. OnTrap.
|
||
type TrapEvent struct {
|
||
DeviceID string
|
||
SourceIP string
|
||
TrapOID string
|
||
Varbinds []TrapVarbind
|
||
}
|
||
|
||
// OnTrap звіряє пачку трапів із правилами джерела `trap`.
|
||
//
|
||
// Зведення робиться до звернення до бази — так само, як для журналу:
|
||
// порт, що фліпає, дає linkDown/linkUp пачками, і сто UPSERT-ів замість
|
||
// одного тут нічого не додають.
|
||
//
|
||
// Окреме рішення, яке варто знати. Трап БЕЗ хоста піднімає алерт лише
|
||
// тоді, коли правило явно назвало адресу джерела (умова source_ip).
|
||
// Причина в тому, що алерт без хоста нікуди не маршрутизується, не
|
||
// глушиться вікном обслуговування й майже нічого не каже черговому:
|
||
// «трап від 10.20.0.77» — це питання, а не аварія. Робити з кожного
|
||
// такого питання алерт означало б залити дошку тим, з чим о третій ночі
|
||
// не можна зробити нічого.
|
||
//
|
||
// Але й губити їх не можна: незнайома адреса, що шле трапи, — часто
|
||
// перший слід нового заліза в мережі. Тому вони не зникають, а
|
||
// потрапляють у власний перелік (inv.trap_unknown_sources), який видно
|
||
// на сторінці трапів окремим блоком. Алерт — для того, що вже знаєш;
|
||
// перелік — для того, чого ще не знаєш.
|
||
func (s *EventSink) OnTrap(ctx context.Context, tenantID string, events []TrapEvent) {
|
||
if len(events) == 0 {
|
||
return
|
||
}
|
||
view := s.view(ctx, tenantID)
|
||
if view == nil {
|
||
return
|
||
}
|
||
|
||
hits := map[string]*trapHit{}
|
||
for _, ev := range events {
|
||
for _, c := range view.rules {
|
||
if c.rule.Source != "trap" || !trapMatches(c, ev) {
|
||
continue
|
||
}
|
||
key := store.TrapDedupKey(c.rule.ID, ev.DeviceID, ev.SourceIP)
|
||
h, ok := hits[key]
|
||
if !ok {
|
||
h = &trapHit{rule: c, device: ev.DeviceID, ip: ev.SourceIP, oid: ev.TrapOID}
|
||
hits[key] = h
|
||
}
|
||
h.count++
|
||
h.last = ev
|
||
}
|
||
}
|
||
|
||
for key, h := range hits {
|
||
meaning := store.ResolveTrapOID(view.trapNames, h.oid)
|
||
what := meaning.Name
|
||
if what == "" {
|
||
// Назви немає — так і кажемо. Вигадана за схожістю префікса
|
||
// назва в заголовку алерту була б найгіршим із можливих
|
||
// варіантів: саме заголовок читають, коли вирішують, чи
|
||
// вставати.
|
||
what = "невідомий трап " + orDefault(h.oid, "без OID")
|
||
}
|
||
msg := what
|
||
if h.device == "" {
|
||
msg += " від " + h.ip + " (адреси немає серед хостів)"
|
||
}
|
||
if detail := trapDetail(view, h.last); detail != "" {
|
||
msg += " · " + detail
|
||
}
|
||
|
||
meta := map[string]any{
|
||
"kind": "trap",
|
||
"trap_oid": h.oid,
|
||
"trap_name": meaning.Name,
|
||
"source_ip": h.ip,
|
||
"varbinds": trapVarbindMeta(view, h.last),
|
||
}
|
||
// Підпис для алерту без хоста — сама адреса: це єдине, що про
|
||
// такого відправника взагалі відомо.
|
||
label := ""
|
||
if h.device == "" {
|
||
label = h.ip
|
||
}
|
||
s.raiseKeyed(ctx, tenantID, view, h.rule, key, h.device, label, msg, h.count, meta)
|
||
}
|
||
}
|
||
|
||
// trapHit — зведення однакових трапів до одного звернення до бази.
|
||
type trapHit struct {
|
||
rule compiledRule
|
||
device string
|
||
ip string
|
||
oid string
|
||
count int
|
||
last TrapEvent
|
||
}
|
||
|
||
// trapMatches — чи підпадає трап під умову правила.
|
||
func trapMatches(c compiledRule, ev TrapEvent) bool {
|
||
cond := c.rule.Condition
|
||
|
||
if oid := store.NormalizeOID(cond.TrapOID); oid != "" {
|
||
if store.NormalizeOID(ev.TrapOID) != oid {
|
||
return false
|
||
}
|
||
}
|
||
|
||
if src := strings.TrimSpace(cond.SourceIP); src != "" {
|
||
if !ipMatches(src, ev.SourceIP) {
|
||
return false
|
||
}
|
||
} else if ev.DeviceID == "" {
|
||
// Трап без хоста й без явно названої адреси — не алерт.
|
||
// Пояснення в коментарі до OnTrap.
|
||
return false
|
||
}
|
||
|
||
// Селектор перевіряємо лише там, де хост є: він оперує хостами, і
|
||
// застосувати його до адреси, якої немає в інвентарі, неможливо.
|
||
// Тому правило з адресою джерела працює й для незнайомців — інакше
|
||
// саме той випадок, заради якого адресу й вписали, не спрацював би
|
||
// ніколи.
|
||
if ev.DeviceID != "" && !c.covers(ev.DeviceID) {
|
||
return false
|
||
}
|
||
|
||
if vbOID := store.NormalizeOID(cond.VarbindOID); vbOID != "" {
|
||
want := strings.TrimSpace(cond.VarbindValue)
|
||
var found bool
|
||
for _, vb := range ev.Varbinds {
|
||
if !varbindIs(vb.OID, vbOID) {
|
||
continue
|
||
}
|
||
// Порожнє очікуване значення означає «щоб такий varbind
|
||
// узагалі був»: умова «трап, у якому є ifIndex» осмислена й
|
||
// відсіює половину службового шуму.
|
||
if want == "" || vb.Value == want {
|
||
found = true
|
||
break
|
||
}
|
||
}
|
||
if !found {
|
||
return false
|
||
}
|
||
}
|
||
return true
|
||
}
|
||
|
||
// varbindIs порівнює OID varbind-а з OID стовпця.
|
||
//
|
||
// Збіг рахується й за префіксом: у трапі приходить ifOperStatus.7 —
|
||
// конкретний порт, — а в умові людина пише ifOperStatus без індексу, бо
|
||
// індексу наперед не знає. Вимагати повного збігу означало б, що умова
|
||
// «ifOperStatus = down» працює рівно для сьомого порту.
|
||
func varbindIs(got, want string) bool {
|
||
return got == want || strings.HasPrefix(got, want+".")
|
||
}
|
||
|
||
// ipMatches — чи належить адреса відправника тому, що написано в умові.
|
||
func ipMatches(pattern, ip string) bool {
|
||
addr := net.ParseIP(ip)
|
||
if addr == nil {
|
||
return false
|
||
}
|
||
if _, netw, err := net.ParseCIDR(pattern); err == nil {
|
||
return netw.Contains(addr)
|
||
}
|
||
return net.ParseIP(pattern).Equal(addr)
|
||
}
|
||
|
||
// trapDetail добирає з varbind-ів те, що варто показати в тексті.
|
||
//
|
||
// Не всі підряд: у повідомленні алерту (а звідти — у Telegram) десяток
|
||
// OID-ів займе весь екран і не пояснить нічого. Беремо ті, у яких є
|
||
// людська назва, — тобто ті, які словник упізнав. Решта лежить у
|
||
// контексті алерту й на сторінці трапів.
|
||
func trapDetail(view *tenantView, ev TrapEvent) string {
|
||
var parts []string
|
||
for _, vb := range ev.Varbinds {
|
||
name := store.ResolveVarbindOID(view.trapNames, vb.OID)
|
||
if name == "" || name == "sysUpTime" || name == "snmpTrapOID" {
|
||
continue
|
||
}
|
||
val := vb.Value
|
||
if lbl := store.DescribeVarbindValue(vb.OID, vb.Value); lbl != "" {
|
||
val = lbl
|
||
}
|
||
parts = append(parts, name+"="+val)
|
||
if len(parts) == 4 {
|
||
break
|
||
}
|
||
}
|
||
return strings.Join(parts, ", ")
|
||
}
|
||
|
||
// trapVarbindMeta кладе varbind-и в контекст алерту.
|
||
func trapVarbindMeta(view *tenantView, ev TrapEvent) []map[string]string {
|
||
out := make([]map[string]string, 0, len(ev.Varbinds))
|
||
for _, vb := range ev.Varbinds {
|
||
m := map[string]string{"oid": vb.OID, "value": vb.Value}
|
||
if name := store.ResolveVarbindOID(view.trapNames, vb.OID); name != "" {
|
||
m["name"] = name
|
||
}
|
||
out = append(out, m)
|
||
}
|
||
return out
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Спільне
|
||
// ---------------------------------------------------------------------
|
||
|
||
// raise доводить один збіг до алерту.
|
||
func (s *EventSink) raise(ctx context.Context, tenantID string, view *tenantView,
|
||
c compiledRule, deviceID, message string, count int, meta map[string]any) {
|
||
|
||
s.raiseKeyed(ctx, tenantID, view, c,
|
||
store.EventDedupKey(c.rule.ID, deviceID), deviceID, "", message, count, meta)
|
||
}
|
||
|
||
// raiseKeyed — те саме, але з явним ключем дедуплікації й підписом.
|
||
//
|
||
// Знадобилось рівно одному джерелу — трапам. Усі інші події приходять
|
||
// від хоста, і хост дає і ключ, і назву в заголовку. Трап приходить від
|
||
// АДРЕСИ, і адреса не завжди є хостом: саме такі трапи найцікавіші
|
||
// (у мережі з'явилось щось, чого інвентар не знає), і зводити їх усі до
|
||
// одного безіменного алерту означало б показати «щось десь сталося».
|
||
func (s *EventSink) raiseKeyed(ctx context.Context, tenantID string, view *tenantView,
|
||
c compiledRule, key, deviceID, label, message string, count int, meta map[string]any) {
|
||
|
||
allowed, carried := s.throttle(key, c.rule.MinIntervalSeconds, count)
|
||
if !allowed {
|
||
return
|
||
}
|
||
|
||
name := view.devices[deviceID]
|
||
if name == "" {
|
||
name = label
|
||
}
|
||
if name == "" {
|
||
name = deviceID
|
||
}
|
||
meta["events"] = carried
|
||
ctxJSON, err := json.Marshal(meta)
|
||
if err != nil {
|
||
ctxJSON = []byte("{}")
|
||
}
|
||
|
||
fired, err := s.st.RaiseEventAlert(ctx, tenantID, store.EventAlert{
|
||
RuleID: c.rule.ID,
|
||
DeviceID: deviceID,
|
||
DeviceName: name,
|
||
Severity: c.rule.Severity,
|
||
Title: name + ": " + c.rule.Name,
|
||
Message: message,
|
||
DedupKey: key,
|
||
Context: ctxJSON,
|
||
Count: carried,
|
||
SuppressedBy: view.sup.For(deviceID, c.rule.ID),
|
||
})
|
||
if err != nil {
|
||
s.log.Error("подієвий алерт", "правило", c.rule.Name, "помилка", err)
|
||
return
|
||
}
|
||
if !fired.IsNew {
|
||
// Продовження вже відомої події не показуємо окремо: лічильник
|
||
// у самому алерті вже виріс, а список алертів перечитується за
|
||
// подією `alert.fired`, якої тут навмисно немає.
|
||
return
|
||
}
|
||
|
||
if err := s.st.PublishEvent(ctx, tenantID, "alert.fired", map[string]any{
|
||
"alert_id": fired.ID, "device_id": deviceID, "severity": fired.Severity,
|
||
"title": fired.Title, "state": fired.State, "suppressed_by": fired.SuppressedBy,
|
||
}); err != nil {
|
||
s.log.Error("подія alert.fired", "помилка", err)
|
||
}
|
||
}
|
||
|
||
// throttle вирішує, чи йти в базу зараз.
|
||
//
|
||
// Обмежувач у пам'яті, а не в SQL, бо захищати треба саме звернення до
|
||
// бази: у потоці журналу дорогим є не сам UPSERT, а те, що їх сотня на
|
||
// секунду з кожного зонда. Кілька процесів матимуть кожен свій
|
||
// обмежувач — і це нормально: остаточну дедуплікацію все одно робить
|
||
// унікальний індекс, а тут йдеться лише про кількість спроб.
|
||
//
|
||
// Повертає, скільки подій слід записати: власні плюс усі, що набігли,
|
||
// поки проміжок не минув.
|
||
func (s *EventSink) throttle(key string, minInterval, count int) (bool, int) {
|
||
if minInterval <= 0 {
|
||
return true, count
|
||
}
|
||
now := time.Now()
|
||
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
|
||
e, ok := s.rate[key]
|
||
if !ok {
|
||
e = &rateEntry{}
|
||
s.rate[key] = e
|
||
}
|
||
e.touched = now
|
||
if ok && now.Sub(e.last) < time.Duration(minInterval)*time.Second {
|
||
e.carried += count
|
||
return false, 0
|
||
}
|
||
e.last = now
|
||
total := e.carried + count
|
||
e.carried = 0
|
||
return true, total
|
||
}
|
||
|
||
// view віддає стан кабінету з кешу, оновлюючи його за потреби.
|
||
func (s *EventSink) view(ctx context.Context, tenantID string) *tenantView {
|
||
now := time.Now()
|
||
|
||
s.mu.Lock()
|
||
v, ok := s.cache[tenantID]
|
||
s.mu.Unlock()
|
||
if ok && now.Sub(v.at) < s.ttl {
|
||
return v
|
||
}
|
||
|
||
rules, err := s.st.EventRules(ctx, tenantID)
|
||
if err != nil {
|
||
s.log.Error("читання подієвих правил", "tenant", tenantID, "помилка", err)
|
||
// Стухлий кеш кращий за жодного: правила рідко міняються, а
|
||
// перебій у базі не має вимикати алерти на весь час перебою.
|
||
return v
|
||
}
|
||
if len(rules) == 0 {
|
||
fresh := &tenantView{at: now, devices: map[string]string{}}
|
||
s.remember(tenantID, fresh)
|
||
return fresh
|
||
}
|
||
|
||
devices, err := s.st.DeviceNames(ctx, tenantID)
|
||
if err != nil {
|
||
s.log.Error("читання хостів", "tenant", tenantID, "помилка", err)
|
||
return v
|
||
}
|
||
sup, err := s.st.LoadSuppression(ctx, tenantID)
|
||
if err != nil {
|
||
// Не привід не піднімати алерти: гірше показати те, про що
|
||
// просили не турбувати, ніж проґавити справжню подію.
|
||
s.log.Warn("вікна обслуговування", "tenant", tenantID, "помилка", err)
|
||
}
|
||
|
||
fresh := &tenantView{at: now, devices: devices, sup: sup}
|
||
if hasSource(rules, "trap") {
|
||
// Помилка тут не має вимикати правила: без словника трап
|
||
// підпаде під умову так само (умова написана OID-ом), просто в
|
||
// заголовку алерту стоятиме число замість назви. Зворотний
|
||
// розмін — тиша замість неідеального тексту — був би гіршим.
|
||
names, err := s.st.TrapNames(ctx, tenantID)
|
||
if err != nil {
|
||
s.log.Warn("словник трапів", "tenant", tenantID, "помилка", err)
|
||
}
|
||
fresh.trapNames = names
|
||
}
|
||
for _, r := range rules {
|
||
c := compiledRule{rule: r}
|
||
if r.Source == "syslog" {
|
||
re, err := regexp.Compile(r.Condition.Regex)
|
||
if err != nil {
|
||
// Зразок перевіряється при збереженні, тож сюди можна
|
||
// дістатись лише правкою в обхід API. Мовчати не можна:
|
||
// правило виглядає ввімкненим.
|
||
s.log.Error("зразок правила не компілюється",
|
||
"правило", r.Name, "помилка", err)
|
||
continue
|
||
}
|
||
c.re = re
|
||
}
|
||
if !emptySelector(r.Selector) {
|
||
ids, err := s.st.SelectorDevices(ctx, tenantID, r.Selector)
|
||
if err != nil {
|
||
s.log.Error("розгортання селектора", "правило", r.Name, "помилка", err)
|
||
continue
|
||
}
|
||
c.scope = make(map[string]bool, len(ids))
|
||
for _, id := range ids {
|
||
c.scope[id] = true
|
||
}
|
||
}
|
||
fresh.rules = append(fresh.rules, c)
|
||
}
|
||
|
||
s.remember(tenantID, fresh)
|
||
return fresh
|
||
}
|
||
|
||
func (s *EventSink) remember(tenantID string, v *tenantView) {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
s.cache[tenantID] = v
|
||
|
||
// Обмежувач частоти тримає по рядку на кожен алерт, який колись
|
||
// піднімався. Без прибирання це повільний витік у процесі, що
|
||
// живе місяцями; година тиші означає, що алерт давно закритий.
|
||
cutoff := time.Now().Add(-time.Hour)
|
||
for k, e := range s.rate {
|
||
if e.touched.Before(cutoff) {
|
||
delete(s.rate, k)
|
||
}
|
||
}
|
||
}
|
||
|
||
func hasSource(rules []store.Rule, source string) bool {
|
||
for _, r := range rules {
|
||
if r.Source == source {
|
||
return true
|
||
}
|
||
}
|
||
return false
|
||
}
|
||
|
||
func emptySelector(s store.Selector) bool {
|
||
return len(s.DeviceIDs) == 0 && len(s.GroupIDs) == 0 && len(s.SiteIDs) == 0 &&
|
||
len(s.Kinds) == 0 && len(s.TemplateIDs) == 0 && len(s.Vendors) == 0 &&
|
||
len(s.Tags) == 0
|
||
}
|
||
|
||
// trimLine готує текст події до показу людині.
|
||
//
|
||
// Рядок журналу буває довжиною в кілограм: у заголовку алерту й у
|
||
// повідомленні в Telegram від цього немає користі, а є втрата решти
|
||
// тексту.
|
||
func trimLine(s string) string {
|
||
s = strings.TrimSpace(strings.ReplaceAll(s, "\n", " "))
|
||
const max = 300
|
||
if len(s) <= max {
|
||
return s
|
||
}
|
||
r := []rune(s)
|
||
if len(r) <= max {
|
||
return s
|
||
}
|
||
return string(r[:max]) + "…"
|
||
}
|
||
|
||
func orDefault(s, def string) string {
|
||
if s == "" {
|
||
return def
|
||
}
|
||
return s
|
||
}
|