Netpulse_SasS/server/internal/store/alerts_channels.go
byrsapty cedf1d261d
All checks were successful
CI / hygiene (push) Successful in 9s
CI / web (push) Successful in 1m18s
CI / server (push) Successful in 1m50s
CI / agent (push) Successful in 1m2s
Ескалації: журнал більше не бреше, ack не воскрешає драбину
Три вади, знайдені рецензією, яких щасливий шлях показати не міг:

* outcome='sent' писався до доставки; помилка читання каналів клала в
  кеш порожню мапу й з'їдала всі сходинки кабінету за тік — усі зі
  слідом «надіслано». Канали тепер читаються до просування стану,
  журнал пишеться після доставки, з правдою.
* UPDATE не мав stopped_at IS NULL — підтвердження алерту посеред
  партії не рятувало людину від дзвінка.
* час брався раз на партію.

Плюс суміжне: UpdateRule не гасив алертів вимкненого правила, сервер
домислював enabled на оновленні, channel_ids сходинок не звірялись із
каналами кабінету (зокрема чужого).

І scripts/dbtest.sh — тести проти бази перестали мовчки пропускатись.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 17:20:18 +03:00

621 lines
22 KiB
Go
Raw 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"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
"github.com/netpulse/netpulse/server/internal/crypto"
)
// Channel — куди слати. Секрет уже розшифрований: движок доставки не
// має доступу до БД секретів і не має його потребувати.
type Channel struct {
ID string `json:"id"`
TenantID string `json:"-"`
Kind string `json:"kind"`
Name string `json:"name"`
Config json.RawMessage `json:"config"`
Template string `json:"template,omitempty"`
MinSeverity string `json:"min_severity"`
Enabled bool `json:"enabled"`
Secret string `json:"-"`
HasSecret bool `json:"has_secret"`
// Назви драбин ескалації, сходинки яких посилаються на цей канал.
//
// Заповнюється лише для переліку в UI (ChannelEscalationRefs), а не
// в LoadChannels: движку доставки це не потрібно, а рахувати на
// кожному тіку — платити без причини.
Escalations []string `json:"escalations,omitempty"`
}
// Route — правило маршрутизації алерту в канали.
type Route struct {
ID string
Name string
Priority int
Matcher RouteMatcher
ChannelIDs []string
Schedule *RouteSchedule
}
type RouteMatcher struct {
SeverityGte string `json:"severity_gte"`
RuleIDs []string `json:"rule_ids"`
SiteIDs []string `json:"site_ids"`
DeviceIDs []string `json:"device_ids"`
}
// RouteSchedule — тихі години.
type RouteSchedule struct {
TZ string `json:"tz"`
Quiet []struct {
Days []int `json:"days"` // 0=неділя, як у time.Weekday
From string `json:"from"` // "22:00"
To string `json:"to"`
} `json:"quiet"`
}
// IsQuiet каже, чи момент t потрапляє в тиху годину.
//
// Інтервал через північ («22:0008:00») — не окремий випадок, а норма
// для чергувань, тому обробляється явно: без цього нічні сповіщення
// тихими годинами не глушились би взагалі.
func (s *RouteSchedule) IsQuiet(t time.Time) bool {
if s == nil || len(s.Quiet) == 0 {
return false
}
loc := time.UTC
if s.TZ != "" {
if l, err := time.LoadLocation(s.TZ); err == nil {
loc = l
}
}
lt := t.In(loc)
mins := lt.Hour()*60 + lt.Minute()
day := int(lt.Weekday())
for _, q := range s.Quiet {
if len(q.Days) > 0 && !containsInt(q.Days, day) {
continue
}
from, ok1 := parseHM(q.From)
to, ok2 := parseHM(q.To)
if !ok1 || !ok2 {
continue
}
if from <= to {
if mins >= from && mins < to {
return true
}
} else if mins >= from || mins < to {
return true
}
}
return false
}
func parseHM(s string) (int, bool) {
var h, m int
if _, err := fmt.Sscanf(s, "%d:%d", &h, &m); err != nil {
return 0, false
}
if h < 0 || h > 23 || m < 0 || m > 59 {
return 0, false
}
return h*60 + m, true
}
func containsInt(xs []int, v int) bool {
for _, x := range xs {
if x == v {
return true
}
}
return false
}
// LoadChannels читає канали тенанта й розшифровує їхні секрети.
func (s *Store) LoadChannels(ctx context.Context, tenantID string, ring *crypto.Keyring) ([]Channel, error) {
var out []Channel
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
rows, err := tx.Query(ctx, `
SELECT c.id::text, c.kind::text, c.name, c.config::text,
COALESCE(c.template,''), c.min_severity::text, c.enabled,
s.key_id, s.nonce, s.ciphertext, s.auth_tag, COALESCE(s.aad,'')
FROM alr.channels c
LEFT JOIN core.secrets s ON s.id = c.secret_id
WHERE c.tenant_id = $1
ORDER BY c.name
`, tenantID)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var c Channel
var cfg string
var keyID, aad *string
var nonce, ct, tag []byte
if err := rows.Scan(&c.ID, &c.Kind, &c.Name, &cfg, &c.Template,
&c.MinSeverity, &c.Enabled, &keyID, &nonce, &ct, &tag, &aad); err != nil {
return err
}
c.TenantID = tenantID
c.Config = json.RawMessage(cfg)
// Наявність секрету — факт про канал, а не наслідок того,
// чи його зараз розшифровують. Перелік для UI викликається
// без кільця навмисно, і без цього рядка канал із токеном
// показувався б у ньому як «секрету немає».
c.HasSecret = keyID != nil
if keyID != nil && ring != nil {
plain, err := ring.Decrypt(&crypto.Secret{
KeyID: *keyID, Nonce: nonce, Ciphertext: ct, AuthTag: tag,
}, derefStr(aad))
if err != nil {
// Канал із нечитабельним секретом не має валити
// доставку решти: одна зіпсована інтеграція гірша
// за мовчання лише для себе самої.
return fmt.Errorf("канал %s: розшифровка секрету: %w", c.Name, err)
}
c.Secret = string(plain)
}
out = append(out, c)
}
return rows.Err()
})
return out, err
}
func derefStr(p *string) string {
if p == nil {
return ""
}
return *p
}
// LoadRoutes читає маршрути тенанта в порядку пріоритету.
func (s *Store) LoadRoutes(ctx context.Context, tenantID string) ([]Route, error) {
var out []Route
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
rows, err := tx.Query(ctx, `
SELECT id::text, name, priority, matcher::text,
channel_ids::text[], COALESCE(schedule::text,'')
FROM alr.routes
WHERE tenant_id = $1 AND enabled
ORDER BY priority, name
`, tenantID)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var r Route
var matcher, sched string
if err := rows.Scan(&r.ID, &r.Name, &r.Priority, &matcher,
&r.ChannelIDs, &sched); err != nil {
return err
}
if err := json.Unmarshal([]byte(matcher), &r.Matcher); err != nil {
return fmt.Errorf("маршрут %s: matcher: %w", r.Name, err)
}
if sched != "" {
var sc RouteSchedule
if err := json.Unmarshal([]byte(sched), &sc); err == nil {
r.Schedule = &sc
}
}
out = append(out, r)
}
return rows.Err()
})
return out, err
}
// Matches — чи цей маршрут бере цей алерт.
func (r Route) Matches(a Alert) bool {
if r.Matcher.SeverityGte != "" {
if severityRank[a.Severity] < severityRank[r.Matcher.SeverityGte] {
return false
}
}
if len(r.Matcher.RuleIDs) > 0 && !containsStr(r.Matcher.RuleIDs, a.RuleID) {
return false
}
if len(r.Matcher.DeviceIDs) > 0 && !containsStr(r.Matcher.DeviceIDs, a.DeviceID) {
return false
}
return true
}
func containsStr(xs []string, v string) bool {
for _, x := range xs {
if x == v {
return true
}
}
return false
}
// RecordNotification пише спробу доставки в журнал і рахує сповіщення
// на алерті.
//
// Журнал ведеться до відправлення, а не після: інакше падіння процесу
// між HTTP-запитом і записом лишало б слід «не надсилали» на вже
// доставленому повідомленні, і ретрай слав би дубль.
func (s *Store) RecordNotification(ctx context.Context, tenantID, alertID, channelID,
status, errMsg, externalID string, payload any) error {
data, _ := json.Marshal(payload)
return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
if _, err := tx.Exec(ctx, `
INSERT INTO alr.notifications
(tenant_id, alert_id, channel_id, status, error, external_id, payload)
VALUES ($1, $2, $3, $4::alr.delivery_status, $5, $6, $7::jsonb)
`, tenantID, nullUUID(alertID), nullUUID(channelID), status,
nullString(errMsg), nullString(externalID), string(data)); err != nil {
return err
}
if status != "sent" || alertID == "" {
return nil
}
_, err := tx.Exec(ctx, `
UPDATE alr.alerts SET notify_count = notify_count + 1
WHERE tenant_id = $1 AND id = $2
`, tenantID, alertID)
return err
})
}
// ---------------------------------------------------------------------
// CRUD каналів для UI
// ---------------------------------------------------------------------
type ChannelInput struct {
Kind string
Name string
Config string
Template string
MinSeverity string
Secret string
Enabled bool
}
// CreateChannel зберігає канал; секрет шифрується тим самим кільцем, що
// й паролі від обладнання — окремого сховища для нього немає навмисно.
func (s *Store) CreateChannel(ctx context.Context, tenantID string, in ChannelInput, ring *crypto.Keyring) (string, error) {
var id string
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
var secretID any
if in.Secret != "" {
if ring == nil {
return errors.New("сервер запущено без ключа шифрування — зберегти секрет каналу ніяк")
}
aad := tenantID + "|alr.channel"
sec, err := ring.Encrypt([]byte(in.Secret), aad)
if err != nil {
return err
}
var sid string
if err := tx.QueryRow(ctx, `
INSERT INTO core.secrets (tenant_id, kind, key_id, nonce, ciphertext, auth_tag, aad)
VALUES ($1, 'api_token', $2, $3, $4, $5, $6)
RETURNING id::text
`, tenantID, sec.KeyID, sec.Nonce, sec.Ciphertext, sec.AuthTag, aad).Scan(&sid); err != nil {
return err
}
secretID = sid
}
return tx.QueryRow(ctx, `
INSERT INTO alr.channels
(tenant_id, kind, name, config, secret_id, template, min_severity, enabled)
VALUES ($1, $2::alr.channel_kind, $3, $4::jsonb, $5, $6, $7::alr.severity, $8)
RETURNING id::text
`, tenantID, in.Kind, in.Name, in.Config, secretID,
nullString(in.Template), in.MinSeverity, in.Enabled).Scan(&id)
})
return id, err
}
// escalationLadder — драбина в тому вигляді, в якому її читає чистка
// посилань: ідентифікатор, назва й розібрані сходинки.
type escalationLadder struct {
id string
name string
steps []EscalationStep
}
// readLadders читає драбини кабінету з уже розібраними сходинками.
//
// Читання окремо від запису навмисно: pgx не дає слати новий запит,
// поки не дочитано попередній, а чистка посилань — це саме «прочитати
// всі, переписати деякі».
func readLadders(ctx context.Context, tx pgx.Tx, tenantID string) ([]escalationLadder, error) {
rows, err := tx.Query(ctx, `
SELECT id::text, name, steps::text
FROM alr.escalation_policies WHERE tenant_id = $1 ORDER BY name
`, tenantID)
if err != nil {
return nil, err
}
defer rows.Close()
var out []escalationLadder
for rows.Next() {
var l escalationLadder
var steps string
if err := rows.Scan(&l.id, &l.name, &steps); err != nil {
return nil, err
}
if err := json.Unmarshal([]byte(steps), &l.steps); err != nil {
return nil, fmt.Errorf("політика %s: сходинки: %w", l.name, err)
}
out = append(out, l)
}
return out, rows.Err()
}
// removeChannelFromSteps прибирає канал зі сходинок і каже, чи щось
// змінилось.
//
// Сходинка, яка лишилась без каналів, ЛИШАЄТЬСЯ порожньою, а не
// зникає. Викинута сходинка мовчки зсунула б усе чергування нижче,
// якого людина не міняла; порожню видно і в переліку драбин, і у формі
// (вона не збережеться, поки канал не оберуть), а движок пише про неї
// в журнал окремим рядком.
func removeChannelFromSteps(steps []EscalationStep, channelID string) ([]EscalationStep, bool) {
changed := false
out := make([]EscalationStep, len(steps))
for i, s := range steps {
out[i] = s
kept := make([]string, 0, len(s.ChannelIDs))
for _, id := range s.ChannelIDs {
if id == channelID {
changed = true
continue
}
kept = append(kept, id)
}
out[i].ChannelIDs = kept
}
return out, changed
}
// ChannelEscalationRefs каже, які драбини посилаються на які канали:
// ідентифікатор каналу → назви драбин.
//
// Потрібне рівно для одного: щоб «Видалити канал» показало те саме, що
// вже показує «Видалити драбину» — скільки чужих налаштувань зараз
// перестане працювати. Мовчазне видалення каналу, на який спирається
// нічне чергування, коштує однієї пропущеної аварії, і дізнаються про
// це не в момент видалення.
func (s *Store) ChannelEscalationRefs(ctx context.Context, tenantID string) (map[string][]string, error) {
refs := map[string][]string{}
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
ladders, err := readLadders(ctx, tx, tenantID)
if err != nil {
return err
}
for _, l := range ladders {
// Драбину називаємо один раз, скільки б сходинок у неї не
// вело в цей канал: людині перед видаленням цікаво, ЩО
// зламається, а не скільки разів воно згадане.
for _, id := range ladderChannelIDs(l.steps) {
refs[id] = append(refs[id], l.name)
}
}
return nil
})
return refs, err
}
// ladderChannelIDs — канали драбини без повторів, у порядку появи.
func ladderChannelIDs(steps []EscalationStep) []string {
var out []string
seen := map[string]bool{}
for _, s := range steps {
for _, id := range s.ChannelIDs {
if id == "" || seen[id] {
continue
}
seen[id] = true
out = append(out, id)
}
}
return out
}
func (s *Store) DeleteChannel(ctx context.Context, tenantID, channelID string) error {
return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
// Секрет видаляється разом із каналом: залишений «на всякий
// випадок» токен у core.secrets — це чинний доступ, про який
// уже ніхто не пам'ятає.
//
// Лічильник беремо з видалення каналу, а не секрету: канал без
// секрету (webhook без підпису) дав би 0 рядків у другому
// операторі, і відповідь була б «немає такого» на успішне
// видалення.
var deleted int
err := tx.QueryRow(ctx, `
WITH gone AS (
DELETE FROM alr.channels
WHERE tenant_id = $1 AND id = $2
RETURNING secret_id
), dropped AS (
DELETE FROM core.secrets
WHERE id IN (SELECT secret_id FROM gone WHERE secret_id IS NOT NULL)
)
SELECT count(*)::int FROM gone
`, tenantID, channelID).Scan(&deleted)
if err != nil {
return err
}
if deleted == 0 {
return ErrAlertNotFound
}
// Посилання зі сходинок драбин прибираємо руками, бо прибрати їх
// нікому: 0066 чистить драбину з правила через ON DELETE SET
// NULL, але на масив усередині JSONB зовнішнього ключа немає.
// Залишений UUID видаленого каналу — це сходинка, яка виглядає
// налаштованою й не йде нікуди, тобто рівно та мовчазна
// обіцянка, заради якої ескалацію й заводили.
//
// Переписуємо в Go, а не запитом по JSONB: рішення «що саме
// лишається в сходинці» перевіряється тоді без Postgres, а
// запити зводяться до читання й запису.
ladders, err := readLadders(ctx, tx, tenantID)
if err != nil {
return err
}
for _, l := range ladders {
steps, changed := removeChannelFromSteps(l.steps, channelID)
if !changed {
continue
}
raw, err := json.Marshal(steps)
if err != nil {
return err
}
if _, err := tx.Exec(ctx, `
UPDATE alr.escalation_policies
SET steps = $3::jsonb, updated_at = now()
WHERE tenant_id = $1 AND id = $2
`, tenantID, l.id, string(raw)); err != nil {
return err
}
}
return nil
})
}
// RuleAction — куди й коли шле саме це правило.
//
// Zabbix розводить «тригер» і «дію над тригером» на дві сутності;
// у переважній більшості випадків це одна думка, розірвана надвоє. Тут
// дія живе в самому правилі, а alr.routes лишаються для спільної
// політики на всі правила разом.
type RuleAction struct {
ChannelIDs []string
Schedule *RouteSchedule
NotifyOnResolve bool
// Драбина ескалації правила. Порожньо — без ескалації.
EscalationPolicyID string
// Джерело правила. Потрібне рівно для одного рішення: подієвий
// алерт (0058) проходить драбину без повторів — див. PlanEscalation.
Source string
}
// LoadRuleActions читає маршрутизацію всіх увімкнених правил тенанта.
//
// Одним запитом на партію алертів, а не по правилу на алерт: під час
// масової аварії партія — це сотні алертів на десяток правил.
func (s *Store) LoadRuleActions(ctx context.Context, tenantID string) (map[string]RuleAction, error) {
out := map[string]RuleAction{}
err := s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
rows, err := tx.Query(ctx, `
SELECT id::text, channel_ids::text[],
COALESCE(notify_schedule::text,''), notify_on_resolve,
COALESCE(escalation_policy_id::text,''), source::text
FROM alr.rules WHERE tenant_id = $1
`, tenantID)
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var id, sched string
var a RuleAction
if err := rows.Scan(&id, &a.ChannelIDs, &sched, &a.NotifyOnResolve,
&a.EscalationPolicyID, &a.Source); err != nil {
return err
}
if sched != "" {
var sc RouteSchedule
if err := json.Unmarshal([]byte(sched), &sc); err == nil {
a.Schedule = &sc
}
}
out[id] = a
}
return rows.Err()
})
return out, err
}
// UpdateChannel змінює канал.
//
// Порожній Secret означає «лишити токен як є»: розшифрувати збережений
// заради показу означає віддати його туди, звідки він уже не
// повернеться, тож форма надіслати незмінений не може навіть теоретично.
//
// Вид каналу не змінюється: telegram із конфігом вебхука — це інший
// об'єкт, і чесніше завести новий, ніж мовчки лишити несумісні поля.
func (s *Store) UpdateChannel(ctx context.Context, tenantID, channelID string, in ChannelInput, ring *crypto.Keyring) error {
return s.InTenantTx(ctx, tenantID, func(tx pgx.Tx) error {
var oldSecret *string
if err := tx.QueryRow(ctx, `
SELECT secret_id::text FROM alr.channels WHERE id = $1 AND tenant_id = $2
`, channelID, tenantID).Scan(&oldSecret); err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return ErrNotFound
}
return err
}
secretID := any(nil)
if oldSecret != nil {
secretID = *oldSecret
}
if in.Secret != "" {
if ring == nil {
return errors.New("сервер запущено без ключа шифрування — зберегти секрет каналу ніяк")
}
aad := tenantID + "|alr.channel"
sec, err := ring.Encrypt([]byte(in.Secret), aad)
if err != nil {
return err
}
var sid string
if err := tx.QueryRow(ctx, `
INSERT INTO core.secrets (tenant_id, kind, key_id, nonce, ciphertext, auth_tag, aad)
VALUES ($1, 'api_token', $2, $3, $4, $5, $6)
RETURNING id::text
`, tenantID, sec.KeyID, sec.Nonce, sec.Ciphertext, sec.AuthTag, aad).Scan(&sid); err != nil {
return err
}
// Старий прибираємо лише після того, як новий ліг: зворотний
// порядок на помилці шифрування лишив би канал без токена.
if oldSecret != nil {
if _, err := tx.Exec(ctx, `DELETE FROM core.secrets WHERE id = $1`, *oldSecret); err != nil {
return err
}
}
secretID = sid
}
ct, err := tx.Exec(ctx, `
UPDATE alr.channels SET
name = $3, config = $4::jsonb, secret_id = $5,
template = $6, min_severity = $7::alr.severity, enabled = $8
WHERE id = $1 AND tenant_id = $2
`, channelID, tenantID, in.Name, in.Config, secretID,
nullString(in.Template), in.MinSeverity, in.Enabled)
if err != nil {
return err
}
if ct.RowsAffected() == 0 {
return ErrNotFound
}
return nil
})
}