Подія відповідності несла шкалу critical/high/medium, а поріг тригера
порівнювався шкалою алертів, де SeverityRank("critical") = 0. Тригер із
порогом «warning» пропускав середнє й відкидав найважче — тихо.
Подія тепер народжується в шкалі алертів; умова тригера приймає обидві.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
628 lines
30 KiB
Go
628 lines
30 KiB
Go
package store
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"net"
|
||
"regexp"
|
||
"strings"
|
||
|
||
"github.com/jackc/pgx/v5"
|
||
)
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Класифікація джерел
|
||
// ---------------------------------------------------------------------
|
||
|
||
// Джерела поділені не за темою, а за тим, звідки береться факт.
|
||
//
|
||
// Опитуване джерело має ряд вимірів: движок щотіку перепитує його й
|
||
// може відповісти і «так», і «ні». Подієве джерело ряду не має — є
|
||
// момент, коли щось сталося, і більше нічого. Уся різниця в поведінці
|
||
// алерту (як він гасне, як дедуплікується, чи можна його «не знайти»)
|
||
// випливає саме звідси.
|
||
var (
|
||
polledSources = map[string]bool{
|
||
"metric": true, "icmp": true, "interface": true,
|
||
}
|
||
eventSources = map[string]bool{
|
||
"syslog": true, "ncm": true, "compliance": true, "trap": true,
|
||
}
|
||
)
|
||
|
||
// IsEventSource — чи обробляється джерело в момент надходження події.
|
||
func IsEventSource(source string) bool { return eventSources[source] }
|
||
|
||
// UnsupportedSourceReason пояснює людині, чому джерело не працює.
|
||
//
|
||
// Порожній рядок означає «працює». Текст тут, а не в HTTP-шарі, бо ту
|
||
// саму відповідь має дати і збереження тригера шаблону: два різні
|
||
// пояснення тієї самої відмови розходяться на першій же правці.
|
||
func UnsupportedSourceReason(source string) string {
|
||
if polledSources[source] || eventSources[source] {
|
||
return ""
|
||
}
|
||
switch source {
|
||
case "link":
|
||
return "лінк на мапі не має власних вимірів — він живий рівно настільки, " +
|
||
"наскільки живі його кінці. Заведіть правило на пристрої або на інтерфейс"
|
||
case "agent":
|
||
return "«зонд не на звʼязку» — це стан, а не подія; він рахується опитуванням. " +
|
||
"Скористайтесь правилом «Пінг» з метрикою «Даних немає взагалі»"
|
||
default:
|
||
return "невідоме джерело правила"
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Перевірка умови при збереженні
|
||
// ---------------------------------------------------------------------
|
||
|
||
// ValidateRuleCondition відмовляє у правилі, яке ніколи не спрацює.
|
||
//
|
||
// Перевірка стоїть на записі, а не на обчисленні, з тієї ж причини, що
|
||
// й у правил відповідності: про друкарську помилку в регулярному виразі
|
||
// людина має дізнатися з форми, а не з тригера, який рік мовчав.
|
||
func ValidateRuleCondition(source string, raw []byte) error {
|
||
if reason := UnsupportedSourceReason(source); reason != "" {
|
||
return fmt.Errorf("%w: %s", ErrInvalid, reason)
|
||
}
|
||
if !eventSources[source] {
|
||
return nil
|
||
}
|
||
|
||
var c Condition
|
||
if len(raw) > 0 {
|
||
if err := json.Unmarshal(raw, &c); err != nil {
|
||
return fmt.Errorf("%w: умова не читається як JSON: %v", ErrInvalid, err)
|
||
}
|
||
}
|
||
|
||
switch source {
|
||
case "syslog":
|
||
if strings.TrimSpace(c.Regex) == "" {
|
||
return fmt.Errorf("%w: правило на журнал без зразка підпало б під кожен рядок "+
|
||
"— задайте, що саме шукати", ErrInvalid)
|
||
}
|
||
if _, err := regexp.Compile(c.Regex); err != nil {
|
||
return fmt.Errorf("%w: зразок не компілюється: %v", ErrInvalid, err)
|
||
}
|
||
if c.SeverityLTE != nil && (*c.SeverityLTE < 0 || *c.SeverityLTE > 7) {
|
||
return fmt.Errorf("%w: рівень syslog буває від 0 (emerg) до 7 (debug)", ErrInvalid)
|
||
}
|
||
case "ncm":
|
||
switch c.Event {
|
||
case "changed", "backup_failed":
|
||
case "":
|
||
return fmt.Errorf("%w: не вказано подію конфігу: changed або backup_failed", ErrInvalid)
|
||
default:
|
||
return fmt.Errorf("%w: невідома подія конфігу %q: буває changed або backup_failed",
|
||
ErrInvalid, c.Event)
|
||
}
|
||
case "compliance":
|
||
if c.Event != "" && c.Event != "violation" {
|
||
return fmt.Errorf("%w: для відповідності є лише подія violation", ErrInvalid)
|
||
}
|
||
// Приймаємо обидві шкали. Людина щойно дивилась на сторінку
|
||
// відповідності, де написано «critical», і саме це слово вона
|
||
// сюди й напише; вимагати подумки перекласти його в «disaster»
|
||
// означає роздавати відмови за власну незручність.
|
||
if c.MinSeverity != "" &&
|
||
!validSeverity[c.MinSeverity] && !validComplianceSeverity[c.MinSeverity] {
|
||
return fmt.Errorf(
|
||
"%w: невідома серйозність %q — буває critical, high, medium (як у правилі "+
|
||
"відповідності) або info, warning, average, high, disaster (як в алерті)",
|
||
ErrInvalid, c.MinSeverity)
|
||
}
|
||
case "trap":
|
||
return validateTrapCondition(c)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// validateTrapCondition перевіряє умову правила на трапи.
|
||
//
|
||
// Головна відповідальність цієї функції — не пропустити правило, яке
|
||
// виглядатиме працюючим. 0058 з'явився саме через таке: джерело `trap`
|
||
// зберігалось мовчки, показувалось увімкненим і не спрацьовувало
|
||
// ніколи. Повернути джерело й лишити хоч одну мовчазну гілку означало б
|
||
// повторити ту саму помилку в дрібнішому масштабі, а це гірше — дрібну
|
||
// довше не помічають.
|
||
//
|
||
// Тому кожна відмова тут не просто відмовляє, а каже, що робити далі.
|
||
func validateTrapCondition(c Condition) error {
|
||
// Поля чужого джерела в умові — не дрібниця. Людина, яка
|
||
// переключила джерело правила з «Syslog» на «Трапи» й лишила в
|
||
// формі зразок, має дізнатись, що зразок більше не діє. Мовчазне
|
||
// ігнорування дало б правило, яке ловить УСІ трапи замість тих, що
|
||
// підпадають під зразок.
|
||
if strings.TrimSpace(c.Regex) != "" {
|
||
return fmt.Errorf("%w: зразок (regex) до трапів не застосовується — трап це не рядок "+
|
||
"тексту, а набір типізованих полів. Задайте OID трапа, а за потреби "+
|
||
"звузьте його конкретним varbind-ом", ErrInvalid)
|
||
}
|
||
if c.SeverityLTE != nil {
|
||
return fmt.Errorf("%w: у трапа немає рівня severity — його має syslog. Серйозність "+
|
||
"алерту задається самим правилом", ErrInvalid)
|
||
}
|
||
if strings.TrimSpace(c.Tag) != "" {
|
||
return fmt.Errorf("%w: тега у трапа немає; те, «що саме сталося», задається OID трапа",
|
||
ErrInvalid)
|
||
}
|
||
|
||
oid := NormalizeOID(c.TrapOID)
|
||
src := strings.TrimSpace(c.SourceIP)
|
||
if oid == "" && src == "" {
|
||
// Правило без жодного звуження підпадає під КОЖЕН трап у
|
||
// кабінеті. Формально воно робоче, практично — це спосіб
|
||
// отримати алерт на кожен linkUp кожного порту й вимкнути
|
||
// сповіщення назавжди через тиждень.
|
||
return fmt.Errorf("%w: правило без OID трапа й без адреси джерела підпало б під "+
|
||
"кожен трап у мережі — вкажіть, що саме ловимо", ErrInvalid)
|
||
}
|
||
if oid != "" && !ValidOID(oid) {
|
||
return fmt.Errorf("%w: %q не схоже на OID. Очікуються числа через крапку "+
|
||
"(наприклад 1.3.6.1.6.3.1.1.5.3 — linkDown); назву трапа зі свого словника "+
|
||
"теж треба вказувати її OID-ом", ErrInvalid, c.TrapOID)
|
||
}
|
||
if src != "" && !validIPOrCIDR(src) {
|
||
return fmt.Errorf("%w: %q не схоже на адресу або підмережу (10.20.0.5 чи "+
|
||
"10.20.0.0/24)", ErrInvalid, c.SourceIP)
|
||
}
|
||
|
||
vbOID := NormalizeOID(c.VarbindOID)
|
||
if vbOID != "" && !ValidOID(vbOID) {
|
||
return fmt.Errorf("%w: %q не схоже на OID varbind-а", ErrInvalid, c.VarbindOID)
|
||
}
|
||
if strings.TrimSpace(c.VarbindValue) != "" && vbOID == "" {
|
||
// Порівнювати значення, не сказавши якого поля, ніде: у трапі
|
||
// їх десяток. Умова «будь-який varbind дорівнює 2» зривалась би
|
||
// на кожному другому трапі й виглядала б при цьому осмисленою.
|
||
return fmt.Errorf("%w: вказано значення varbind-а, але не вказано, якого саме. "+
|
||
"Додайте OID varbind-а — наприклад 1.3.6.1.2.1.2.2.1.1 (ifIndex)", ErrInvalid)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// validIPOrCIDR — поверхнева перевірка адреси джерела.
|
||
//
|
||
// Стоїть на записі, а не на спрацюванні, з тієї ж причини, що й
|
||
// перевірка регулярного виразу: про описку в адресі людина має
|
||
// дізнатись із форми, а не з правила, яке рік мовчало.
|
||
func validIPOrCIDR(s string) bool {
|
||
if _, _, err := net.ParseCIDR(s); err == nil {
|
||
return true
|
||
}
|
||
return net.ParseIP(s) != nil
|
||
}
|
||
|
||
var validSeverity = map[string]bool{
|
||
"info": true, "warning": true, "average": true, "high": true, "disaster": true,
|
||
}
|
||
|
||
// Шкала правил ВІДПОВІДНОСТІ. Інша, і це не недогляд: відповідність
|
||
// говорить мовою аудиту, алерти — мовою чергового. Спільне слово одне —
|
||
// «high».
|
||
var validComplianceSeverity = map[string]bool{
|
||
"critical": true, "high": true, "medium": true, "low": true,
|
||
}
|
||
|
||
// ComplianceSeverityToAlert переводить серйозність правила відповідності
|
||
// у шкалу алертів.
|
||
//
|
||
// Без цього переведення виходила тиха й ЗВОРОТНА помилка. Подія несла
|
||
// серйозність відповідності, а порівнювалась через SeverityRank, який
|
||
// знає лише шкалу алертів і на «critical» повертає нуль. Тобто тригер
|
||
// із порогом «warning» пропускав середні порушення й ВІДКИДАВ
|
||
// найважчі — рівно ті, через які цей механізм і вмикають.
|
||
//
|
||
// Напрямок очевидний: найгірше в одній шкалі стає найгіршим у другій.
|
||
func ComplianceSeverityToAlert(s string) string {
|
||
switch s {
|
||
case "critical":
|
||
return "disaster"
|
||
case "high":
|
||
return "high"
|
||
case "medium":
|
||
return "average"
|
||
case "low":
|
||
return "warning"
|
||
}
|
||
// Невідоме не піднімаємо й не занижуємо мовчки: «warning» означає
|
||
// «покажи, але не буди», і для незнайомого слова це найчесніше.
|
||
return "warning"
|
||
}
|
||
|
||
// NormalizeEventSeverity зводить будь-яку з двох шкал до шкали алертів.
|
||
func NormalizeEventSeverity(s string) string {
|
||
if validSeverity[s] {
|
||
return s
|
||
}
|
||
return ComplianceSeverityToAlert(s)
|
||
}
|
||
|
||
// SeverityRank — порядок серйозності, той самий, що в alr.severity.
|
||
func SeverityRank(s string) int {
|
||
switch s {
|
||
case "warning":
|
||
return 1
|
||
case "average":
|
||
return 2
|
||
case "high":
|
||
return 3
|
||
case "disaster":
|
||
return 4
|
||
default:
|
||
return 0
|
||
}
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Читання подієвих правил
|
||
// ---------------------------------------------------------------------
|
||
|
||
// EventRules читає увімкнені подієві правила одного тенанта.
|
||
//
|
||
// Окремо від ActiveRules навмисно: та вибірка наскрізна по всіх
|
||
// тенантах і робиться раз на тік одним процесом, а ця — гаряча. Її
|
||
// смикає приймач журналу, тобто найчастіший шлях у системі, і вона має
|
||
// віддавати рівно правила одного кабінету, щоб їх можна було закешувати
|
||
// поруч із його ж списком хостів.
|
||
func (s *Store) EventRules(ctx context.Context, tenantID string) ([]Rule, error) {
|
||
// Тенантна транзакція, хоча предикат r.tenant_id = $1 у запиті вже
|
||
// стоїть: alr.rules під tenant_isolation, і без app.tenant_id
|
||
// вибірка порожня. Наслідок був би тихий і найгірший з можливих —
|
||
// подієві алерти просто перестали б заводитись, а сторінка алертів
|
||
// виглядала б як «усе спокійно».
|
||
var out []Rule
|
||
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
|
||
rows, err := tx.Query(ctx, `
|
||
SELECT r.id::text, r.name, r.source::text, r.severity::text,
|
||
r.selector::text, r.condition::text,
|
||
r.auto_close_seconds, r.min_interval_seconds
|
||
FROM alr.rules r
|
||
JOIN core.tenants t ON t.id = r.tenant_id
|
||
WHERE r.tenant_id = $1
|
||
AND r.enabled
|
||
AND r.source IN ('syslog','ncm','compliance','trap')
|
||
AND t.status NOT IN ('suspended','cancelled')
|
||
ORDER BY r.name
|
||
`, tenantID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer rows.Close()
|
||
|
||
for rows.Next() {
|
||
r := Rule{TenantID: tenantID}
|
||
var sel, cond string
|
||
if err := rows.Scan(&r.ID, &r.Name, &r.Source, &r.Severity, &sel, &cond,
|
||
&r.AutoCloseSeconds, &r.MinIntervalSeconds); err != nil {
|
||
return err
|
||
}
|
||
if err := json.Unmarshal([]byte(sel), &r.Selector); err != nil {
|
||
return fmt.Errorf("правило %s: selector: %w", r.Name, err)
|
||
}
|
||
if err := json.Unmarshal([]byte(cond), &r.Condition); err != nil {
|
||
return fmt.Errorf("правило %s: condition: %w", r.Name, err)
|
||
}
|
||
out = append(out, r)
|
||
}
|
||
return rows.Err()
|
||
})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return out, nil
|
||
}
|
||
|
||
// SelectorDevices розгортає селектор правила у перелік хостів.
|
||
//
|
||
// Подієвий шлях не може перевіряти належність хоста до селектора
|
||
// запитом на кожну подію: рядків журналу за секунду більше, ніж хостів
|
||
// у кабінеті. Тому селектор розгортається один раз і живе в кеші поруч
|
||
// із правилами — ціною того, що щойно доданий хост підпадає під правило
|
||
// не миттєво, а з наступним оновленням кешу.
|
||
func (s *Store) SelectorDevices(ctx context.Context, tenantID string, sel Selector) ([]string, error) {
|
||
var out []string
|
||
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
|
||
ids, err := s.resolveSelectorDevices(ctx, tx, tenantID, sel)
|
||
out = ids
|
||
return err
|
||
})
|
||
return out, err
|
||
}
|
||
|
||
// ---------------------------------------------------------------------
|
||
// Підняття подієвого алерту
|
||
// ---------------------------------------------------------------------
|
||
|
||
// EventAlert — подія, яка вже визнана такою, що підпадає під правило.
|
||
type EventAlert struct {
|
||
RuleID string
|
||
DeviceID string
|
||
DeviceName string
|
||
Severity string
|
||
Title string
|
||
Message string
|
||
DedupKey string
|
||
Context json.RawMessage
|
||
// Скільки однакових подій злилося в це звернення. Одиниця — звичайна
|
||
// подія; більше — пачка, зведена ще до звернення до бази.
|
||
Count int
|
||
SuppressedBy string
|
||
}
|
||
|
||
// RaiseEventAlert створює подієвий алерт або доливає подію в наявний.
|
||
//
|
||
// Дедуплікація — той самий унікальний індекс, що й у метричних алертів,
|
||
// і це не економія коду. Ключ навмисно не містить нічого від самої
|
||
// події: один алерт на пару «правило + хост» незалежно від того, чи
|
||
// прийшов один рядок журналу, чи чотири тисячі. Інакше перший же
|
||
// мигаючий порт зробив би дошку алертів нечитабельною за хвилину, а
|
||
// саме дошка — те, заради чого все це існує.
|
||
//
|
||
// Ціна такого рішення чесна й видима: у алерті лишається останній текст
|
||
// і лічильник подій, а не весь їхній перелік. Перелік є в журналі, і
|
||
// шукати його треба там.
|
||
func (s *Store) RaiseEventAlert(ctx context.Context, tenantID string, e EventAlert) (FiredAlert, error) {
|
||
var out FiredAlert
|
||
ctxJSON := "{}"
|
||
if len(e.Context) > 0 {
|
||
ctxJSON = string(e.Context)
|
||
}
|
||
count := e.Count
|
||
if count < 1 {
|
||
count = 1
|
||
}
|
||
|
||
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
|
||
return tx.QueryRow(ctx, `
|
||
WITH prev AS (
|
||
SELECT id, state::text AS old_state
|
||
FROM alr.alerts
|
||
WHERE tenant_id = $1 AND dedup_key = $6
|
||
AND state IN ('firing','acknowledged','suppressed')
|
||
), ups AS (
|
||
INSERT INTO alr.alerts
|
||
(tenant_id, rule_id, device_id, severity, state, title, message,
|
||
dedup_key, context, suppressed_by, event_count, notify_pending,
|
||
started_at, last_seen_at)
|
||
VALUES ($1, $2, $3, $4::alr.severity,
|
||
CASE WHEN $8::text IS NULL THEN 'firing' ELSE 'suppressed' END::alr.alert_state,
|
||
$5, $9, $6, $7::jsonb, $8, $10,
|
||
-- Придушений алерт не ставиться в чергу на
|
||
-- розсилку: заглушення означає «не турбувати»,
|
||
-- а не «покажи пізніше».
|
||
$8::text IS NULL,
|
||
now(), now())
|
||
ON CONFLICT (tenant_id, dedup_key)
|
||
WHERE state IN ('firing','acknowledged','suppressed')
|
||
DO UPDATE SET
|
||
last_seen_at = now(),
|
||
event_count = alr.alerts.event_count + EXCLUDED.event_count,
|
||
message = EXCLUDED.message,
|
||
context = EXCLUDED.context,
|
||
suppressed_by = $8,
|
||
state = CASE
|
||
WHEN alr.alerts.state = 'acknowledged' THEN 'acknowledged'
|
||
WHEN $8::text IS NULL THEN 'firing'
|
||
ELSE 'suppressed' END::alr.alert_state
|
||
-- notify_pending навмисно не чіпаємо: продовження
|
||
-- вже відомої події не є новиною, і сотий рядок
|
||
-- журналу не має слати сотого повідомлення.
|
||
RETURNING id, state::text AS new_state, started_at, last_seen_at,
|
||
event_count, notify_count, (xmax = 0) AS inserted
|
||
)
|
||
SELECT ups.id::text, ups.new_state, ups.started_at, ups.last_seen_at,
|
||
ups.event_count, ups.notify_count, ups.inserted,
|
||
COALESCE(prev.old_state, '')
|
||
FROM ups LEFT JOIN prev ON prev.id = ups.id
|
||
`, tenantID, nullUUID(e.RuleID), nullUUID(e.DeviceID), e.Severity,
|
||
e.Title, e.DedupKey, ctxJSON, nullString(e.SuppressedBy),
|
||
nullString(e.Message), count,
|
||
).Scan(&out.ID, &out.State, &out.StartedAt, &out.LastSeenAt,
|
||
&out.EventCount, &out.NotifyCount, &out.IsNew, &out.PrevState)
|
||
})
|
||
if err != nil {
|
||
return out, fmt.Errorf("подієвий алерт %s: %w", e.DedupKey, err)
|
||
}
|
||
|
||
out.RuleID, out.DeviceID = e.RuleID, e.DeviceID
|
||
out.DeviceName, out.Severity, out.Title = e.DeviceName, e.Severity, e.Title
|
||
out.Message, out.DedupKey, out.SuppressedBy = e.Message, e.DedupKey, e.SuppressedBy
|
||
return out, nil
|
||
}
|
||
|
||
// ResolveEventAlerts закриває подієві алерти за їхніми ключами.
|
||
//
|
||
// Потрібно рівно там, де в події ВСЕ Ж таки є зворотний бік:
|
||
// відповідність перевіряється прогоном, і той самий прогін, у якому
|
||
// хост правило пройшов, — єдиний чесний сигнал «більше не порушено».
|
||
// Для журналу й конфігів такого сигналу не існує: рядок «конфіг
|
||
// змінився» ніщо не скасовує.
|
||
func (s *Store) ResolveEventAlerts(ctx context.Context, tenantID string, dedupKeys []string, reason string) (int, error) {
|
||
if len(dedupKeys) == 0 {
|
||
return 0, nil
|
||
}
|
||
var n int
|
||
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
|
||
return tx.QueryRow(ctx, `
|
||
WITH closed AS (
|
||
UPDATE alr.alerts
|
||
SET state = 'resolved', resolved_at = now(), notify_pending = false
|
||
WHERE tenant_id = $1
|
||
AND dedup_key = ANY($2::text[])
|
||
AND state IN ('firing','acknowledged','suppressed')
|
||
RETURNING id, device_id, severity, title
|
||
), published AS (
|
||
INSERT INTO core.event_outbox (tenant_id, topic, payload)
|
||
SELECT $1, 'alert.resolved',
|
||
jsonb_build_object('alert_id', id::text,
|
||
'device_id', COALESCE(device_id::text,''),
|
||
'severity', severity::text,
|
||
'title', title,
|
||
'state', 'resolved',
|
||
'reason', $3::text)
|
||
FROM closed
|
||
RETURNING 1
|
||
)
|
||
SELECT count(*)::int FROM published
|
||
`, tenantID, dedupKeys, reason).Scan(&n)
|
||
})
|
||
return n, err
|
||
}
|
||
|
||
// ExpireEventAlerts гасить подієві алерти, до яких давно не було подій.
|
||
//
|
||
// Це не прибирання й не косметика — це відповідь на питання «як
|
||
// закривається алерт, що не має умови».
|
||
//
|
||
// Метричний алерт закриває сама дійсність: умова перестала виконуватись
|
||
// — рядок зник із кандидатів. Подієвий такого шансу не має: рядок
|
||
// журналу стався, і «перестати ставатись» не може. Лишити його висіти
|
||
// назавжди означає за тиждень отримати дошку з двома сотнями старих
|
||
// подій, на яку ніхто не дивиться, — а тоді на ній не помітять і
|
||
// справжню аварію.
|
||
//
|
||
// Тому такий алерт має строк. Стан навмисно `expired`, а не `resolved`:
|
||
// ніхто не казав, що проблему полагодили, вона просто відстоялась. У
|
||
// журналі й на екрані це має виглядати по-різному, інакше «саме
|
||
// минулося» неможливо відрізнити від «розібрались».
|
||
func (s *Store) ExpireEventAlerts(ctx context.Context) (int64, error) {
|
||
tag, err := s.bg.Exec(ctx, `
|
||
WITH aged AS (
|
||
UPDATE alr.alerts a
|
||
SET state = 'expired', resolved_at = now(), notify_pending = false
|
||
FROM alr.rules r
|
||
WHERE r.id = a.rule_id
|
||
AND r.auto_close_seconds > 0
|
||
AND a.state IN ('firing','acknowledged','suppressed')
|
||
AND a.last_seen_at < now() - make_interval(secs => r.auto_close_seconds)
|
||
RETURNING a.id, a.tenant_id, a.device_id, a.severity, a.title
|
||
)
|
||
INSERT INTO core.event_outbox (tenant_id, topic, payload)
|
||
SELECT tenant_id, 'alert.resolved',
|
||
jsonb_build_object('alert_id', id::text,
|
||
'device_id', COALESCE(device_id::text,''),
|
||
'severity', severity::text,
|
||
'title', title,
|
||
'state', 'expired',
|
||
'reason', 'подій більше не було')
|
||
FROM aged
|
||
`)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
return tag.RowsAffected(), nil
|
||
}
|
||
|
||
// TakeNotifyPending забирає подієві алерти, які ще нікуди не пішли.
|
||
//
|
||
// Забирає, а не читає: позначка знімається тією ж командою, що й
|
||
// повертає рядки. Інакше два процеси, які випадково опинились у тіку
|
||
// одночасно, розіслали б одне й те саме двічі — а телефон о третій ночі
|
||
// не розрізняє «дублікат» і «друга аварія».
|
||
//
|
||
// Знята позначка означає «спробували», а не «доставили». Це свідомо:
|
||
// повторні спроби доставки — робота каналу, і робити їх звідси означало
|
||
// б слати вдруге те, що вже дійшло, щоразу як мовчить один із трьох
|
||
// каналів.
|
||
func (s *Store) TakeNotifyPending(ctx context.Context, limit int) ([]Alert, error) {
|
||
if limit <= 0 || limit > 500 {
|
||
limit = 200
|
||
}
|
||
rows, err := s.bg.Query(ctx, `
|
||
WITH taken AS (
|
||
UPDATE alr.alerts a
|
||
SET notify_pending = false
|
||
WHERE a.id IN (
|
||
SELECT id FROM alr.alerts
|
||
WHERE notify_pending
|
||
ORDER BY started_at
|
||
LIMIT $1
|
||
FOR UPDATE SKIP LOCKED
|
||
)
|
||
RETURNING a.*
|
||
)
|
||
SELECT t.tenant_id::text, t.id::text, COALESCE(t.rule_id::text,''),
|
||
COALESCE(t.device_id::text,''), COALESCE(d.name,''),
|
||
t.severity::text, t.state::text, t.title, COALESCE(t.message,''),
|
||
t.dedup_key, COALESCE(t.suppressed_by,''), t.started_at, t.last_seen_at,
|
||
t.event_count
|
||
FROM taken t
|
||
LEFT JOIN inv.devices d ON d.id = t.device_id
|
||
`, limit)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
var out []Alert
|
||
for rows.Next() {
|
||
var a Alert
|
||
if err := rows.Scan(&a.TenantID, &a.ID, &a.RuleID, &a.DeviceID, &a.DeviceName,
|
||
&a.Severity, &a.State, &a.Title, &a.Message, &a.DedupKey,
|
||
&a.SuppressedBy, &a.StartedAt, &a.LastSeenAt, &a.EventCount); err != nil {
|
||
return nil, err
|
||
}
|
||
out = append(out, a)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
// EventDedupKey — ключ подієвого алерту.
|
||
//
|
||
// Свідомо той самий вигляд, що й у Candidate.DedupKey: подія й вимір
|
||
// про той самий хост і те саме правило — це один алерт, і два різні
|
||
// формати ключа рано чи пізно дали б два.
|
||
func EventDedupKey(ruleID, deviceID string) string {
|
||
return ruleID + ":dev:" + deviceID
|
||
}
|
||
|
||
// TrapDedupKey — ключ алерту за трапом.
|
||
//
|
||
// Для трапа від відомого хоста це той самий ключ, що й для решти
|
||
// подієвих джерел: подія про хост — один алерт на пару «правило+хост».
|
||
//
|
||
// Для трапа від адреси, якої немає в інвентарі, хоста немає взагалі, і
|
||
// ключ будується від адреси. Без цього всі невпізнані відправники
|
||
// злилися б у ОДИН алерт на правило — тобто «щось у мережі шле трапи»,
|
||
// з чим неможливо нічого зробити. З адресою в ключі кожен незнайомець
|
||
// має власний рядок на дошці, і його видно як окреме питання.
|
||
func TrapDedupKey(ruleID, deviceID, sourceIP string) string {
|
||
if deviceID != "" {
|
||
return EventDedupKey(ruleID, deviceID)
|
||
}
|
||
return ruleID + ":ip:" + sourceIP
|
||
}
|
||
|
||
// DeviceNames — імена всіх живих хостів кабінету.
|
||
//
|
||
// Одним запитом на весь кабінет, а не по хосту на подію: ім'я потрібне
|
||
// лише в заголовку алерту, а заголовок пишеться раз при створенні.
|
||
// Платити за нього окремим запитом на кожен рядок журналу означало б
|
||
// зробити найдорожчою частиною шляху найдешевшу його потребу.
|
||
func (s *Store) DeviceNames(ctx context.Context, tenantID string) (map[string]string, error) {
|
||
out := map[string]string{}
|
||
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
|
||
rows, err := tx.Query(ctx, `
|
||
SELECT id::text, name FROM inv.devices
|
||
WHERE tenant_id = $1 AND deleted_at IS NULL
|
||
`, tenantID)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer rows.Close()
|
||
for rows.Next() {
|
||
var id, name string
|
||
if err := rows.Scan(&id, &name); err != nil {
|
||
return err
|
||
}
|
||
out[id] = name
|
||
}
|
||
return rows.Err()
|
||
})
|
||
return out, err
|
||
}
|