Правило «процесор вище 85% — це проблема» описує клас пристроїв, а не окремий хост. Тепер воно живе поруч із перевірками, які дають йому дані, а не окремою сторінкою, де його доводилось повторювати руками для кожного комутатора. Тригер розгортається в ОДНЕ правило alr.rules із селектором за шаблоном, а не в правило на кожен хост: призначили шаблон новому пристрою — він одразу під правилом, без перегенерації. Форма шаблону розкладена на вкладки Загальне/Перевірки/Графіки/Тригери з лічильниками; смуга вкладок винесена у спільний компонент. Умова тригера редагується полями, JSON лишився запасним виходом. Правила з шаблону в списку правил помічені й не редагуються там. Клон копіює тригери вимкненими, щоб копія не подвоїла сповіщення. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
502 lines
17 KiB
Go
502 lines
17 KiB
Go
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
|
||
}
|
||
|
||
// 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
|
||
}
|
||
|
||
// 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.pool.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
|
||
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); 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/trap/ncm/compliance обробляються не опитуванням,
|
||
// а подіями — цей движок їх свідомо не чіпає.
|
||
return nil, nil
|
||
}
|
||
}
|
||
|
||
// 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
|
||
}
|