Netpulse_SasS/server/internal/store/alerts.go
byrsapty ed8fc831bf Дві сесії роботи: 0058–0068, розгортання однією командою, тести
Один коміт, а не десяток тематичних, свідомо: теми переплетені в
спільних файлах (store.go, docker-compose.yml, deploy/README.md), і
розділити їх можна було б лише індексуванням шматків. Коміти, які не
збираються, гірші за один великий — тим паче що це рівно той стан, який
перевірявся разом.

ЩО ПРАЦЮЄ НА СТЕНДІ Й ПЕРЕВІРЕНО ТАМ

  0058  подієві алерти: syslog, ncm, compliance спрацьовують у мить
        події; правило з нереалізованим джерелом більше не зберігається
        мовчки
  0059  snmp.walk і прототипи шаблонів — таблиці з динамічним індексом
        описуються шаблоном, а не Go
  0060  відкат конфігу: план як різниця, маскування паролів із підписом
        плану, обов'язковий контрольний збір, verifying при обриві
  0061  кнопки Telegram: довге опитування, авторизація не з callback_data
  0062  аудит і архів хостів; тест на AST, що падає на ключі без назви
  0063  RLS: три ролі, окремий пул для фонових тактів
  0064  строки зберігання даних і сторінка сховища
  0065  приймач SNMP-трапів; перевірено справжніми пакетами по дроту,
        переклад v1→v2 за RFC 3584 дає правильний OID
  0066  ескалації сповіщень
  0067  алерт про вичерпання диска
  0068  поля заливки конфігу переїхали в каталог профілів

Плюс: 137 тестів вебу з нуля (їх не було взагалі), одинадцять справжніх
вад, знайдених ними й виправлених, і виправлення двох інтеграційних
тестів grpcapi, які мовчки пропускались півтора року.

ЩО ЩЕ НЕ ЗАПУСКАЛОСЬ

  netpulse            установник: одна команда замість 18 змінних і
                      593 рядків інструкції
  RLS з першого запуску  нова інсталяція під політиками одразу;
                      RLS-EXISTING-INSTALL.md лишається тільки для
                      старих інсталяцій
  .forgejo + CI       раннер не зареєстрований

Ці три перевірені компіляцією й міркуванням, але не виконанням.

ГОЛОВНИЙ ВИСНОВОК ДВОХ СЕСІЙ

Зелена перевірка доводить рівно те, що вона перевіряє. Тест ізоляції RLS
був правильний і зелений — і пропустив зламаний вхід, бо перевіряв «чи
не видно чужого», коли зламалось «чи видно своє». Інтеграційні тести
grpcapi були зелені, бо не виконувались. Схема, довідник і протокол
описували те, чого в коді не існувало, і виглядало це як готове.

Тому в кожному завданні цих сесій стояла вимога назвати НЕПОКРИТЕ, а
чотири задачі закінчились не можливістю, а відмовою: правило з
нереалізованим джерелом не зберігається, профіль без команд заливки
каже про це замість мовчазної кнопки, міграція RLS валить сама себе на
таблиці без політики, тест словника аудиту падає на ключі без назви.

Подробиці — HISTORY.md, розділи за 26 і 27 серпня.
2026-08-27 17:32:49 +03:00

552 lines
20 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 store
import (
"context"
"encoding/json"
"fmt"
"sort"
"strings"
"time"
"github.com/jackc/pgx/v5"
)
// Rule — правило з alr.rules у вигляді, придатному для обчислення.
type Rule struct {
ID string
TenantID string
Name string
Description string
Source string
Severity string
Selector Selector
Condition Condition
ForSeconds int
DependsOnTopology bool
// Строк життя подієвого алерту й мінімальний проміжок між
// зверненнями до нього. Для опитуваних джерел обидва не мають
// сенсу й лишаються нулями.
AutoCloseSeconds int
MinIntervalSeconds int
}
// Selector — до чого застосовується правило. Порожній означає «до всього»:
// це найчастіший випадок («ping усіх»), і вимагати для нього явного
// переліку означало б ламати правило щоразу, коли додається пристрій.
type Selector struct {
DeviceIDs []string `json:"device_ids"`
GroupIDs []string `json:"group_ids"`
SiteIDs []string `json:"site_ids"`
Kinds []string `json:"kinds"`
// Хости, яким призначено ці шаблони. Так тригер шаблону стає одним
// правилом замість правила на кожен хост: призначили шаблон новому
// комутатору — він одразу під правилом, без перегенерації.
TemplateIDs []string `json:"template_ids"`
Vendors []string `json:"vendors"`
Tags map[string]string `json:"tags"`
}
// Condition — умова спрацювання.
//
// {"metric":"loss_pct","op":">","value":20} — усі семпли вікна
// {"metric":"rtt_avg_ms","op":">","value":150,"agg":"avg"} — середнє за вікно
// {"metric":"no_data"} — даних немає взагалі
//
// Без `agg` умова має триматися ВСІ семпли вікна `for_seconds`. Це і є
// антифлап: одна втрачена відповідь не піднімає алерт, а вікно вимірюється
// в даних, а не в стані движка, тому перезапуск нічого не збиває.
type Condition struct {
Metric string `json:"metric"`
Op string `json:"op"`
Value float64 `json:"value"`
Agg string `json:"agg"`
MetricKey string `json:"metric_key"` // для source=metric: ts.series.metric_key
// --- подієві джерела ---
//
// У них немає ні порогу, ні вікна: подія або сталася, або ні.
// Тому й поля інші — вони описують не «скільки», а «яка саме».
// syslog: {"regex":"%LINK-3-UPDOWN.*down","severity_lte":4,"tag":"LINK"}
Regex string `json:"regex,omitempty"`
Tag string `json:"tag,omitempty"`
SeverityLTE *int `json:"severity_lte,omitempty"`
// ncm: {"event":"changed"} | {"event":"backup_failed"}
// compliance: {"event":"violation"}
Event string `json:"event,omitempty"`
// compliance: звузити до конкретних правил відповідності й до
// порога їхньої серйозності. Порожньо — усі.
RuleIDs []string `json:"rule_ids,omitempty"`
MinSeverity string `json:"min_severity,omitempty"`
// trap: {"trap_oid":"1.3.6.1.6.3.1.1.5.3",
// "varbind_oid":"1.3.6.1.2.1.2.2.1.1","varbind_value":"7",
// "source_ip":"10.20.0.0/24"}
//
// Три питання, і рівно ті, які до трапа ставлять: ЩО сталося (OID),
// ЗВІДКИ прийшло (адреса — потрібна окремо від селектора, бо
// селектор оперує хостами, а трап приходить і з адрес, яких в
// інвентарі немає) і З ЯКИМ значенням (varbind).
//
// Чого тут немає й не буде мовчки: зразка (regex) по тексту трапа.
// Трап — це не рядок, а набір типізованих полів, і «пошук по
// трапу» довелося б визначати як пошук по конкатенації чогось із
// чимось. Форма про це каже вголос (ValidateRuleCondition), а не
// приймає regex і не використовує його.
TrapOID string `json:"trap_oid,omitempty"`
VarbindOID string `json:"varbind_oid,omitempty"`
VarbindValue string `json:"varbind_value,omitempty"`
SourceIP string `json:"source_ip,omitempty"`
}
// Candidate — об'єкт, який щойно задовольнив умову правила.
type Candidate struct {
DeviceID string
DeviceName string
InterfaceID string
IfName string
Value float64
Samples int
LastTS time.Time
}
// DedupKey — ключ, за яким алерт вважається тим самим.
//
// Він же primary key дедуплікації в БД, тому має бути стабільним між
// тіками: інакше кожне обчислення відкривало б новий алерт і слало
// сповіщення заново.
func (c Candidate) DedupKey(ruleID string) string {
if c.InterfaceID != "" {
return ruleID + ":if:" + c.InterfaceID
}
return ruleID + ":dev:" + c.DeviceID
}
// ---------------------------------------------------------------------
// Правила
// ---------------------------------------------------------------------
// ActiveRules читає всі увімкнені правила всіх тенантів.
//
// Движок один на інсталяцію, тому вибірка навмисно наскрізна: тримати
// по горутині на тенант означало б платити з'єднанням до БД за кожного
// клієнта, у якого може не бути жодного правила.
func (s *Store) ActiveRules(ctx context.Context) ([]Rule, error) {
rows, err := s.bg.Query(ctx, `
SELECT r.id::text, r.tenant_id::text, r.name, COALESCE(r.description,''),
r.source::text, r.severity::text,
r.selector::text, r.condition::text,
r.for_seconds, r.depends_on_topology,
r.auto_close_seconds, r.min_interval_seconds
FROM alr.rules r
JOIN core.tenants t ON t.id = r.tenant_id
WHERE r.enabled AND t.status NOT IN ('suspended','cancelled')
ORDER BY r.tenant_id, r.name
`)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Rule
for rows.Next() {
var r Rule
var sel, cond string
if err := rows.Scan(&r.ID, &r.TenantID, &r.Name, &r.Description,
&r.Source, &r.Severity, &sel, &cond,
&r.ForSeconds, &r.DependsOnTopology,
&r.AutoCloseSeconds, &r.MinIntervalSeconds); err != nil {
return nil, err
}
if err := json.Unmarshal([]byte(sel), &r.Selector); err != nil {
return nil, fmt.Errorf("правило %s: selector: %w", r.Name, err)
}
if err := json.Unmarshal([]byte(cond), &r.Condition); err != nil {
return nil, fmt.Errorf("правило %s: condition: %w", r.Name, err)
}
out = append(out, r)
}
return out, rows.Err()
}
// ---------------------------------------------------------------------
// Обчислення умов
// ---------------------------------------------------------------------
// Колонки, дозволені в умовах. Список закритий навмисно: значення з
// condition потрапляє в текст запиту, і будь-яке послаблення тут
// перетворюється на SQL-ін'єкцію через JSON у таблиці правил.
var icmpMetrics = map[string]string{
"rtt_avg_ms": "rtt_avg_ms",
"rtt_min_ms": "rtt_min_ms",
"rtt_max_ms": "rtt_max_ms",
"jitter_ms": "jitter_ms",
"loss_pct": "loss_pct",
"reachable": "(CASE WHEN reachable THEN 1 ELSE 0 END)",
}
var ifMetrics = map[string]string{
"in_bps": "in_bps",
"out_bps": "out_bps",
"in_pps": "in_pps",
"out_pps": "out_pps",
"util_in_pct": "util_in_pct",
"util_out_pct": "util_out_pct",
"in_errors": "in_errors",
"out_errors": "out_errors",
"in_discards": "in_discards",
"out_discards": "out_discards",
"oper_up": "(CASE WHEN oper_up THEN 1 ELSE 0 END)",
}
var sqlOps = map[string]string{
">": ">", ">=": ">=", "<": "<", "<=": "<=",
"==": "=", "=": "=", "!=": "<>", "<>": "<>",
}
var sqlAggs = map[string]string{
"avg": "avg(%s)", "min": "min(%s)", "max": "max(%s)",
"sum": "sum(%s)", "count": "count(%s)",
"last": "(array_agg(%s ORDER BY ts DESC))[1]",
}
// EvaluateRule повертає об'єкти, для яких умова правила виконується.
func (s *Store) EvaluateRule(ctx context.Context, r Rule) ([]Candidate, error) {
if r.Condition.Metric == "no_data" {
return s.evalNoData(ctx, r)
}
switch r.Source {
case "icmp":
return s.evalSamples(ctx, r, icmpMetrics, `ts.icmp_samples s`, "", false)
case "interface":
return s.evalSamples(ctx, r, ifMetrics, `ts.if_counters s`, "", true)
case "metric":
return s.evalSeries(ctx, r)
default:
// Джерела syslog/ncm/compliance обробляються не опитуванням, а
// в момент надходження події (alerts_events.go). Сюди вони
// доходити не мають: движок відсіює їх раніше, бо інакше
// ResolveMissing із порожнім списком кандидатів закривав би
// щойно піднятий подієвий алерт на наступному ж тіку.
return nil, fmt.Errorf("правило %s: джерело %s не обчислюється опитуванням",
r.Name, r.Source)
}
}
// evalSamples — спільний шлях для icmp і interface: обидва лежать у
// гіпертаблицях із колонкою device_id і однаковою семантикою вікна.
func (s *Store) evalSamples(ctx context.Context, r Rule, allowed map[string]string,
from, extraJoin string, perInterface bool) ([]Candidate, error) {
col, ok := allowed[r.Condition.Metric]
if !ok {
return nil, fmt.Errorf("правило %s: невідома метрика %q для джерела %s",
r.Name, r.Condition.Metric, r.Source)
}
op, ok := sqlOps[r.Condition.Op]
if !ok {
return nil, fmt.Errorf("правило %s: невідомий оператор %q", r.Name, r.Condition.Op)
}
a := &args{}
tenant := a.add(r.TenantID)
window := a.add(fmt.Sprintf("%d seconds", maxInt(r.ForSeconds, 1)))
threshold := a.add(r.Condition.Value)
where, err := s.selectorSQL(r.Selector, a, "s.device_id")
if err != nil {
return nil, err
}
// Дві семантики вікна. З agg — агрегат порівнюється один раз;
// без agg — умова має триматися кожен семпл вікна.
var having, valueExpr string
if r.Condition.Agg != "" {
tmpl, ok := sqlAggs[r.Condition.Agg]
if !ok {
return nil, fmt.Errorf("правило %s: невідома агрегація %q", r.Name, r.Condition.Agg)
}
valueExpr = fmt.Sprintf(tmpl, col)
having = fmt.Sprintf("%s %s %s", valueExpr, op, threshold)
} else {
valueExpr = fmt.Sprintf("avg(%s)", col)
having = fmt.Sprintf("bool_and(%s %s %s)", col, op, threshold)
}
group, sel, join := "s.device_id", "s.device_id::text, ''::text", ""
if perInterface {
group = "s.device_id, s.interface_id"
sel = "s.device_id::text, s.interface_id::text"
join = "JOIN inv.interfaces i ON i.id = s.interface_id AND i.monitored"
}
q := fmt.Sprintf(`
SELECT %s, %s AS val, count(*)::int AS n, max(s.ts) AS last_ts
FROM %s
JOIN inv.devices d ON d.id = s.device_id AND d.deleted_at IS NULL AND d.enabled
%s %s
WHERE s.tenant_id = %s
AND s.ts >= now() - %s::interval
%s
GROUP BY %s
HAVING %s
`, sel, valueExpr, from, join, extraJoin, tenant, window, where, group, having)
return s.scanCandidates(ctx, r.TenantID, q, a.vals)
}
// evalSeries — довільні метрики з ts.samples через ts.series.
func (s *Store) evalSeries(ctx context.Context, r Rule) ([]Candidate, error) {
if r.Condition.MetricKey == "" {
return nil, fmt.Errorf("правило %s: для джерела metric потрібен metric_key", r.Name)
}
op, ok := sqlOps[r.Condition.Op]
if !ok {
return nil, fmt.Errorf("правило %s: невідомий оператор %q", r.Name, r.Condition.Op)
}
a := &args{}
tenant := a.add(r.TenantID)
window := a.add(fmt.Sprintf("%d seconds", maxInt(r.ForSeconds, 1)))
threshold := a.add(r.Condition.Value)
key := a.add(r.Condition.MetricKey)
where, err := s.selectorSQL(r.Selector, a, "se.device_id")
if err != nil {
return nil, err
}
var having, valueExpr string
if r.Condition.Agg != "" {
tmpl, ok := sqlAggs[r.Condition.Agg]
if !ok {
return nil, fmt.Errorf("правило %s: невідома агрегація %q", r.Name, r.Condition.Agg)
}
valueExpr = fmt.Sprintf(tmpl, "s.value")
having = fmt.Sprintf("%s %s %s", valueExpr, op, threshold)
} else {
valueExpr = "avg(s.value)"
having = fmt.Sprintf("bool_and(s.value %s %s)", op, threshold)
}
q := fmt.Sprintf(`
SELECT se.device_id::text, COALESCE(se.interface_id::text, ''),
%s AS val, count(*)::int AS n, max(s.ts) AS last_ts
FROM ts.samples s
JOIN ts.series se ON se.id = s.series_id
JOIN inv.devices d ON d.id = se.device_id AND d.deleted_at IS NULL AND d.enabled
WHERE se.tenant_id = %s
AND se.metric_key = %s
AND s.ts >= now() - %s::interval
%s
GROUP BY se.device_id, se.interface_id
HAVING %s
`, valueExpr, tenant, key, window, where, having)
return s.scanCandidates(ctx, r.TenantID, q, a.vals)
}
// evalNoData шукає протилежне решті движка: пристрої, від яких даних
// НЕ надходить.
//
// Це не додаткова зручність. Моніторинг, який мовчить, коли замовк
// зонд, показує зелену мапу мертвої мережі — стан, гірший за будь-яку
// хибну тривогу.
func (s *Store) evalNoData(ctx context.Context, r Rule) ([]Candidate, error) {
a := &args{}
tenant := a.add(r.TenantID)
window := a.add(fmt.Sprintf("%d seconds", maxInt(r.ForSeconds, 1)))
where, err := s.selectorSQL(r.Selector, a, "d.id")
if err != nil {
return nil, err
}
// Пристрої, які колись відповідали (last_seen_at не NULL), але
// замовкли. Ті, що не відповідали ніколи, — це помилка налаштування,
// а не збій, і вони не мають будити людину вночі.
q := fmt.Sprintf(`
SELECT d.id::text, ''::text, 0::double precision, 0::int,
COALESCE(d.last_seen_at, now())
FROM inv.devices d
WHERE d.tenant_id = %s
AND d.deleted_at IS NULL AND d.enabled
AND d.last_seen_at IS NOT NULL
AND d.last_seen_at < now() - %s::interval
%s
`, tenant, window, where)
return s.scanCandidates(ctx, r.TenantID, q, a.vals)
}
func (s *Store) scanCandidates(ctx context.Context, tenantID, q string, vals []any) ([]Candidate, error) {
var out []Candidate
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
rows, err := tx.Query(ctx, q, vals...)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var c Candidate
if err := rows.Scan(&c.DeviceID, &c.InterfaceID, &c.Value, &c.Samples, &c.LastTS); err != nil {
return err
}
out = append(out, c)
}
return rows.Err()
})
if err != nil {
return nil, err
}
return s.fillNames(ctx, tenantID, out)
}
// fillNames добирає імена одним запитом.
//
// Ім'я потрібне в заголовку алерту, а заголовок пишеться один раз при
// створенні: якщо взяти його з JOIN у гарячому запиті обчислення, той
// платитиме за це на кожному тіку для всіх пристроїв, а не лише для
// тих кількох, що справді спрацювали.
func (s *Store) fillNames(ctx context.Context, tenantID string, cs []Candidate) ([]Candidate, error) {
if len(cs) == 0 {
return cs, nil
}
devIDs := make([]string, 0, len(cs))
ifIDs := make([]string, 0, len(cs))
for _, c := range cs {
devIDs = append(devIDs, c.DeviceID)
if c.InterfaceID != "" {
ifIDs = append(ifIDs, c.InterfaceID)
}
}
devNames := map[string]string{}
ifNames := 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 id = ANY($2::uuid[])`,
tenantID, devIDs)
if err != nil {
return err
}
for rows.Next() {
var id, name string
if err := rows.Scan(&id, &name); err != nil {
rows.Close()
return err
}
devNames[id] = name
}
rows.Close()
if err := rows.Err(); err != nil {
return err
}
if len(ifIDs) == 0 {
return nil
}
rows, err = tx.Query(ctx,
`SELECT id::text, name FROM inv.interfaces WHERE tenant_id = $1 AND id = ANY($2::uuid[])`,
tenantID, ifIDs)
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
}
ifNames[id] = name
}
return rows.Err()
})
if err != nil {
return nil, err
}
for i := range cs {
cs[i].DeviceName = devNames[cs[i].DeviceID]
cs[i].IfName = ifNames[cs[i].InterfaceID]
}
return cs, nil
}
// ---------------------------------------------------------------------
// Побудова предикатів
// ---------------------------------------------------------------------
// args нумерує плейсхолдери. Значення з правил ніколи не вставляються
// в текст запиту — лише через параметри; у текст іде тільки те, що
// пройшло whitelist вище.
type args struct{ vals []any }
func (a *args) add(v any) string {
a.vals = append(a.vals, v)
return fmt.Sprintf("$%d", len(a.vals))
}
func (s *Store) selectorSQL(sel Selector, a *args, deviceCol string) (string, error) {
var parts []string
if len(sel.DeviceIDs) > 0 {
parts = append(parts, fmt.Sprintf("%s = ANY(%s::uuid[])", deviceCol, a.add(sel.DeviceIDs)))
}
if len(sel.SiteIDs) > 0 {
parts = append(parts, fmt.Sprintf("d.site_id = ANY(%s::uuid[])", a.add(sel.SiteIDs)))
}
if len(sel.Kinds) > 0 {
parts = append(parts, fmt.Sprintf("d.kind::text = ANY(%s::text[])", a.add(sel.Kinds)))
}
if len(sel.Vendors) > 0 {
parts = append(parts, fmt.Sprintf("d.vendor = ANY(%s::text[])", a.add(sel.Vendors)))
}
if len(sel.TemplateIDs) > 0 {
parts = append(parts, fmt.Sprintf(`EXISTS (
SELECT 1 FROM tpl.device_templates dt
WHERE dt.device_id = d.id AND dt.template_id = ANY(%s::uuid[]))`,
a.add(sel.TemplateIDs)))
}
if len(sel.GroupIDs) > 0 {
parts = append(parts, fmt.Sprintf(`EXISTS (
SELECT 1 FROM inv.device_group_members m
WHERE m.device_id = d.id AND m.group_id = ANY(%s::uuid[]))`,
a.add(sel.GroupIDs)))
}
// Ключі сортуються, щоб текст запиту не залежав від порядку обходу
// мапи: інакше одне й те саме правило породжувало б різні запити й
// щоразу промахувалось повз кеш підготовлених виразів.
for _, k := range sortedKeys(sel.Tags) {
parts = append(parts, fmt.Sprintf(`EXISTS (
SELECT 1 FROM inv.device_tags dt JOIN inv.tags t ON t.id = dt.tag_id
WHERE dt.device_id = d.id AND t.key = %s AND t.value = %s)`,
a.add(k), a.add(sel.Tags[k])))
}
if len(parts) == 0 {
return "", nil
}
return "AND " + strings.Join(parts, "\n AND "), nil
}
func sortedKeys(m map[string]string) []string {
out := make([]string, 0, len(m))
for k := range m {
out = append(out, k)
}
sort.Strings(out)
return out
}
func maxInt(a, b int) int {
if a > b {
return a
}
return b
}