1059 lines
37 KiB
Go
Executable file
1059 lines
37 KiB
Go
Executable file
package main
|
||
|
||
import (
|
||
"bytes"
|
||
"crypto/md5"
|
||
"encoding/json"
|
||
"fmt"
|
||
"io"
|
||
"log"
|
||
"net/http"
|
||
"sort"
|
||
"strings"
|
||
"sync"
|
||
)
|
||
|
||
// AlertRulesTemplate — структура шаблона алертов из alert_rules_template.json
|
||
type AlertRulesTemplate struct {
|
||
Version string `json:"version"`
|
||
Description string `json:"description"`
|
||
Defaults AlertDefaults `json:"defaults"`
|
||
RPSAlert map[string]interface{} `json:"rps_alert"`
|
||
StaticAlerts []map[string]interface{} `json:"static_alerts"` // Алерты без привязки к БД
|
||
ContactPoint ContactPointConfig `json:"contact_point"` // Настройки шаблона уведомлений
|
||
}
|
||
|
||
// ContactPointConfig описывает кастомный шаблон сообщения для Contact Point
|
||
type ContactPointConfig struct {
|
||
ReceiverName string `json:"receiver_name"` // Имя contact point в Grafana
|
||
MessageTemplate string `json:"message_template"` // Go-шаблон сообщения Telegram
|
||
}
|
||
|
||
type AlertDefaults struct {
|
||
NoDataState string `json:"no_data_state"`
|
||
ExecErrState string `json:"exec_err_state"`
|
||
Receiver string `json:"receiver"`
|
||
Folder string `json:"folder"`
|
||
Group string `json:"group"`
|
||
GroupInterval string `json:"group_interval"`
|
||
}
|
||
|
||
// generateAlertUID генерирует детерминированный UID алерта из имени клиента и типа.
|
||
// Одинаковый UID при повторном запуске = Grafana обновит существующий алерт, не создаст новый.
|
||
func generateAlertUID(clientTitle, alertType string) string {
|
||
hash := md5.Sum([]byte(clientTitle + ":" + alertType))
|
||
return fmt.Sprintf("%x", hash[:8])
|
||
}
|
||
|
||
// buildRPSAlertRule строит JSON правила алерта для одного клиента по шаблону
|
||
func buildRPSAlertRule(client ClientData, tmpl AlertRulesTemplate, config Config) map[string]interface{} {
|
||
datasourceUID := config.AlertsDatasourceUID
|
||
|
||
// SID-запрос: один SID или несколько через OR
|
||
sids := make([]string, 0, len(client.Domains))
|
||
for _, d := range client.Domains {
|
||
sids = append(sids, d.SID)
|
||
}
|
||
// Убираем дубликаты и сортируем для детерминированности
|
||
sidMap := make(map[string]bool)
|
||
for _, s := range sids {
|
||
sidMap[s] = true
|
||
}
|
||
uniqueSIDs := make([]string, 0, len(sidMap))
|
||
for s := range sidMap {
|
||
uniqueSIDs = append(uniqueSIDs, s)
|
||
}
|
||
sort.Strings(uniqueSIDs)
|
||
|
||
sidParts := make([]string, len(uniqueSIDs))
|
||
for i, s := range uniqueSIDs {
|
||
sidParts[i] = "SID:" + s
|
||
}
|
||
sidQuery := strings.Join(sidParts, " OR ")
|
||
|
||
// Получаем настройки из шаблона
|
||
defaults := tmpl.Defaults
|
||
receiver := orDefault(config.AlertsReceiver, defaults.Receiver)
|
||
group := client.ClientTitle
|
||
|
||
timeRangeFrom := int64(1800)
|
||
if v, ok := tmpl.RPSAlert["relative_time_range_from"].(float64); ok {
|
||
timeRangeFrom = int64(v)
|
||
}
|
||
|
||
summaryTemplate := "Превышен порог в {rps_limit} RPS в {client_title}"
|
||
if v, ok := tmpl.RPSAlert["annotation_summary"].(string); ok {
|
||
summaryTemplate = v
|
||
}
|
||
summary := strings.NewReplacer(
|
||
"{rps_limit}", fmt.Sprintf("%d", client.RPSLimit),
|
||
"{client_title}", client.ClientTitle,
|
||
).Replace(summaryTemplate)
|
||
|
||
uid := generateAlertUID(client.ClientTitle, "rps")
|
||
dataRefID := "A"
|
||
rpsRefID := "RPS " + client.ClientTitle
|
||
reduceRefID := "Значение"
|
||
thresholdRefID := "Превышение"
|
||
|
||
rule := map[string]interface{}{
|
||
"uid": uid,
|
||
"title": fmt.Sprintf("RPS %s", client.ClientTitle),
|
||
"condition": thresholdRefID,
|
||
"noDataState": defaults.NoDataState,
|
||
"execErrState": defaults.ExecErrState,
|
||
"isPaused": false,
|
||
"folderUID": "", // будет заполнен при отправке
|
||
"ruleGroup": group,
|
||
"annotations": map[string]string{
|
||
"summary": summary,
|
||
},
|
||
"notification_settings": map[string]string{
|
||
"receiver": receiver,
|
||
},
|
||
"data": []interface{}{
|
||
// Шаг 1: count-запрос по SID
|
||
map[string]interface{}{
|
||
"refId": dataRefID,
|
||
"queryType": "lucene",
|
||
"relativeTimeRange": map[string]interface{}{
|
||
"from": timeRangeFrom,
|
||
"to": 0,
|
||
},
|
||
"datasourceUid": datasourceUID,
|
||
"model": map[string]interface{}{
|
||
"alias": fmt.Sprintf("RPS %s", client.ClientTitle),
|
||
"bucketAggs": []interface{}{
|
||
map[string]interface{}{
|
||
"field": "@timestamp",
|
||
"id": "2",
|
||
"settings": map[string]interface{}{
|
||
"interval": "1m",
|
||
"min_doc_count": "0",
|
||
"trimEdges": "0",
|
||
},
|
||
"type": "date_histogram",
|
||
},
|
||
},
|
||
"datasource": map[string]interface{}{
|
||
"type": "grafana-opensearch-datasource",
|
||
"uid": datasourceUID,
|
||
},
|
||
"format": "table",
|
||
"instant": false,
|
||
"intervalMs": 1000,
|
||
"luceneQueryType": "Metric",
|
||
"maxDataPoints": 43200,
|
||
"metrics": []interface{}{
|
||
map[string]interface{}{"hide": false, "id": "1", "type": "count"},
|
||
},
|
||
"query": sidQuery,
|
||
"queryType": "lucene",
|
||
"range": true,
|
||
"refId": dataRefID,
|
||
"timeField": "@timestamp",
|
||
},
|
||
},
|
||
// Шаг 2: math $A / 60 → RPS
|
||
map[string]interface{}{
|
||
"refId": rpsRefID,
|
||
"relativeTimeRange": map[string]interface{}{
|
||
"from": timeRangeFrom,
|
||
"to": 0,
|
||
},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"datasource": map[string]interface{}{
|
||
"name": "Expression",
|
||
"type": "__expr__",
|
||
"uid": "__expr__",
|
||
},
|
||
"expression": fmt.Sprintf("$%s / 60", dataRefID),
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"refId": rpsRefID,
|
||
"type": "math",
|
||
"window": "",
|
||
},
|
||
},
|
||
// Шаг 3: reduce median
|
||
map[string]interface{}{
|
||
"refId": reduceRefID,
|
||
"relativeTimeRange": map[string]interface{}{
|
||
"from": 0,
|
||
"to": 0,
|
||
},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"conditions": []interface{}{
|
||
map[string]interface{}{
|
||
"evaluator": map[string]interface{}{"params": []int{0, 0}, "type": "gt"},
|
||
"operator": map[string]interface{}{"type": "and"},
|
||
"query": map[string]interface{}{"params": []string{}},
|
||
"reducer": map[string]interface{}{"params": []string{}, "type": "avg"},
|
||
"type": "query",
|
||
},
|
||
},
|
||
"datasource": map[string]interface{}{
|
||
"name": "Expression",
|
||
"type": "__expr__",
|
||
"uid": "__expr__",
|
||
},
|
||
"expression": rpsRefID,
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"reducer": "median",
|
||
"refId": reduceRefID,
|
||
"settings": map[string]interface{}{"mode": ""},
|
||
"type": "reduce",
|
||
},
|
||
},
|
||
// Шаг 4: threshold > rps_limit
|
||
map[string]interface{}{
|
||
"refId": thresholdRefID,
|
||
"relativeTimeRange": map[string]interface{}{
|
||
"from": 0,
|
||
"to": 0,
|
||
},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"conditions": []interface{}{
|
||
map[string]interface{}{
|
||
"evaluator": map[string]interface{}{
|
||
"params": []interface{}{client.RPSLimit, 0},
|
||
"type": "gt",
|
||
},
|
||
"operator": map[string]interface{}{"type": "and"},
|
||
"query": map[string]interface{}{"params": []string{}},
|
||
"reducer": map[string]interface{}{"params": []string{}, "type": "avg"},
|
||
"type": "query",
|
||
},
|
||
},
|
||
"datasource": map[string]interface{}{
|
||
"name": "Expression",
|
||
"type": "__expr__",
|
||
"uid": "__expr__",
|
||
},
|
||
"expression": reduceRefID,
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"refId": thresholdRefID,
|
||
"type": "threshold",
|
||
},
|
||
},
|
||
},
|
||
}
|
||
|
||
return rule
|
||
}
|
||
|
||
// alertTask описывает одно правило алерта для отправки
|
||
type alertTask struct {
|
||
rule map[string]interface{}
|
||
label string // для логирования
|
||
}
|
||
|
||
// parallelUpsert отправляет список алертов параллельно с ограничением concurrency.
|
||
// dryRun=true только логирует без отправки.
|
||
// Возвращает количество успешных и неуспешных отправок.
|
||
func parallelUpsert(tasks []alertTask, config Config, dryRun bool, concurrency int) (sent, failed int) {
|
||
sem := make(chan struct{}, concurrency)
|
||
var mu sync.Mutex
|
||
var wg sync.WaitGroup
|
||
|
||
for _, task := range tasks {
|
||
wg.Add(1)
|
||
sem <- struct{}{} // захватить слот
|
||
go func(t alertTask) {
|
||
defer wg.Done()
|
||
defer func() { <-sem }() // освободить слот
|
||
|
||
if dryRun {
|
||
data, _ := json.MarshalIndent(t.rule, "", " ")
|
||
log.Printf("DRY RUN: Alert rule '%s':\n%s", t.label, string(data))
|
||
mu.Lock()
|
||
sent++
|
||
mu.Unlock()
|
||
return
|
||
}
|
||
|
||
if err := upsertAlertRule(t.rule, config); err != nil {
|
||
log.Printf("Warning: failed to upsert alert '%s': %v", t.label, err)
|
||
mu.Lock()
|
||
failed++
|
||
mu.Unlock()
|
||
} else {
|
||
log.Printf("Alert upserted: %s", t.label)
|
||
mu.Lock()
|
||
sent++
|
||
mu.Unlock()
|
||
}
|
||
}(task)
|
||
}
|
||
|
||
wg.Wait()
|
||
return
|
||
}
|
||
|
||
// generateAndSendAlerts генерирует и отправляет RPS-алерты для клиентов с rps_limit > 0
|
||
func generateAndSendAlerts(clients map[string]ClientData, tmpl AlertRulesTemplate, config Config, dryRun bool) error {
|
||
// Собираем клиентов с лимитом, сортируем для детерминированности
|
||
type clientWithLimit struct {
|
||
title string
|
||
client ClientData
|
||
}
|
||
var targets []clientWithLimit
|
||
for _, client := range clients {
|
||
if client.RPSLimit > 0 {
|
||
targets = append(targets, clientWithLimit{client.ClientTitle, client})
|
||
}
|
||
}
|
||
|
||
if len(targets) == 0 {
|
||
log.Printf("Alerts: no clients with rps_limit set, skipping")
|
||
return nil
|
||
}
|
||
|
||
sort.Slice(targets, func(i, j int) bool {
|
||
return targets[i].title < targets[j].title
|
||
})
|
||
|
||
log.Printf("Alerts: generating RPS rules for %d clients", len(targets))
|
||
|
||
// Получить UID папки для алертов
|
||
folderUID := ""
|
||
if !dryRun {
|
||
uid, err := ensureAlertFolder(config)
|
||
if err != nil {
|
||
return fmt.Errorf("failed to ensure alerts folder: %w", err)
|
||
}
|
||
folderUID = uid
|
||
}
|
||
|
||
// Собираем все динамические задачи: RPS + 4xx + 5xx
|
||
var tasks []alertTask
|
||
for _, t := range targets {
|
||
if t.client.RPSLimit > 0 {
|
||
rule := buildRPSAlertRule(t.client, tmpl, config)
|
||
rule["folderUID"] = folderUID
|
||
tasks = append(tasks, alertTask{
|
||
rule: rule,
|
||
label: fmt.Sprintf("RPS %s (limit: %d)", t.client.ClientTitle, t.client.RPSLimit),
|
||
})
|
||
}
|
||
// WAF block алерт — по одному на каждый SID (пропускаем если ptaf_fallback_code = "pass" или пустой)
|
||
if t.client.PtafFallbackCode != "" && strings.ToLower(t.client.PtafFallbackCode) != "pass" {
|
||
for _, domain := range t.client.Domains {
|
||
rule := buildWAFBlockAlertRule(t.client, domain, tmpl, config)
|
||
rule["folderUID"] = folderUID
|
||
tasks = append(tasks, alertTask{
|
||
rule: rule,
|
||
label: fmt.Sprintf("%s WAF %s / %s SID:%s", t.client.PtafFallbackCode, t.client.ClientTitle, domain.DomainName, domain.SID),
|
||
})
|
||
}
|
||
}
|
||
|
||
// Error rate алерты — пороги берутся из каждого домена (apps_settings)
|
||
// Генерируем для response_status_code и upstream_status_code
|
||
for _, domain := range t.client.Domains {
|
||
for _, statusField := range []string{"response_status_code", "upstream_status"} {
|
||
if domain.Limit4xx > 0 && domain.Limit4xx < 100 {
|
||
for _, rule := range buildErrorRateAlertRules(t.client, domain, 4, domain.Limit4xx, statusField, tmpl, config) {
|
||
rule["folderUID"] = folderUID
|
||
tasks = append(tasks, alertTask{
|
||
rule: rule,
|
||
label: fmt.Sprintf("4xx %s %s / %s SID:%s", statusField, t.client.ClientTitle, domain.DomainName, domain.SID),
|
||
})
|
||
}
|
||
}
|
||
if domain.Limit5xx > 0 && domain.Limit5xx < 100 {
|
||
for _, rule := range buildErrorRateAlertRules(t.client, domain, 5, domain.Limit5xx, statusField, tmpl, config) {
|
||
rule["folderUID"] = folderUID
|
||
tasks = append(tasks, alertTask{
|
||
rule: rule,
|
||
label: fmt.Sprintf("5xx %s %s / %s SID:%s", statusField, t.client.ClientTitle, domain.DomainName, domain.SID),
|
||
})
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
log.Printf("Alerts: sending %d dynamic rules (concurrency: 5)", len(tasks))
|
||
sent, failed := parallelUpsert(tasks, config, dryRun, 20)
|
||
log.Printf("Alerts: sent=%d, failed=%d", sent, failed)
|
||
if failed > 0 {
|
||
return fmt.Errorf("%d alert rules failed to upsert", failed)
|
||
}
|
||
|
||
// Статические алерты — отправляются всегда вместе с динамическими
|
||
if err := sendStaticAlerts(tmpl, config, dryRun); err != nil {
|
||
log.Printf("Warning: static alerts error: %v", err)
|
||
}
|
||
|
||
// Обновляем шаблон сообщения Contact Point
|
||
if err := upsertContactPoint(tmpl, config, dryRun); err != nil {
|
||
log.Printf("Warning: failed to update contact point message template: %v", err)
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// buildErrorRateAlertRules строит алерт на процент 4xx или 5xx ответов для конкретного домена.
|
||
// statusField: "response_status_code" или "upstream_status"
|
||
func buildErrorRateAlertRules(client ClientData, domain DomainInfo, errCode, limitPct int, statusField string, tmpl AlertRulesTemplate, config Config) []map[string]interface{} {
|
||
datasourceUID := config.AlertsDatasourceUID
|
||
defaults := tmpl.Defaults
|
||
receiver := orDefault(config.AlertsReceiver, defaults.Receiver)
|
||
group := client.ClientTitle
|
||
|
||
codeRange := "[400 TO 499]"
|
||
codeName := "4xx"
|
||
if errCode == 5 {
|
||
codeRange = "[500 TO 599]"
|
||
codeName = "5xx"
|
||
}
|
||
|
||
forDuration := "1m"
|
||
if errCode == 4 {
|
||
forDuration = "3m"
|
||
}
|
||
|
||
relFrom := int64(900)
|
||
|
||
var rules []map[string]interface{}
|
||
|
||
{
|
||
sidQuery := "SID:" + domain.SID
|
||
fieldShort := "response"
|
||
if statusField == "upstream_status" {
|
||
fieldShort = "upstream"
|
||
}
|
||
uid := generateAlertUID(client.ClientTitle+"_"+domain.SID+"_"+fieldShort, codeName)
|
||
title := fmt.Sprintf("%s %s %s / %s SID:%s", codeName, fieldShort, client.ClientTitle, domain.DomainName, domain.SID)
|
||
fieldDesc := "ответ от PTAF"
|
||
if statusField == "upstream_status" {
|
||
fieldDesc = "ответ от origin"
|
||
}
|
||
circle := "🟡"
|
||
summary := fmt.Sprintf(
|
||
"%s %s / %s (SID:%s) — %s ошибок > %d%% держится > %s. [%s]",
|
||
circle, client.ClientTitle, domain.DomainName, domain.SID, codeName, limitPct, forDuration, fieldShort,
|
||
)
|
||
thresholdAnnotation := fmt.Sprintf(">%d%% для %s (%s)", limitPct, statusField, fieldDesc)
|
||
fieldInfoAnnotation := fmt.Sprintf("для %s (%s)", statusField, fieldDesc)
|
||
thresholdValueAnnotation := fmt.Sprintf("%d", limitPct)
|
||
|
||
totalRefID := "Total"
|
||
errRefID := "Errors"
|
||
totalReduceID := "_TotalReduce"
|
||
errReduceID := "_ErrorsReduce"
|
||
percentRefID := "Значение"
|
||
thresholdRefID := "Превышение"
|
||
|
||
rule := map[string]interface{}{
|
||
"uid": uid,
|
||
"title": title,
|
||
"condition": thresholdRefID,
|
||
"noDataState": defaults.NoDataState,
|
||
"execErrState": defaults.ExecErrState,
|
||
"isPaused": false,
|
||
"folderUID": "",
|
||
"ruleGroup": group,
|
||
"for": forDuration,
|
||
"annotations": map[string]string{
|
||
"summary": summary,
|
||
"client_domain": fmt.Sprintf("%s / %s (SID:%s)", client.ClientTitle, domain.DomainName, domain.SID),
|
||
"code_field": codeName,
|
||
"for_duration": forDuration,
|
||
"threshold": thresholdAnnotation,
|
||
"threshold_value": thresholdValueAnnotation,
|
||
"field_info": fieldInfoAnnotation,
|
||
"severity": "warning",
|
||
},
|
||
"notification_settings": map[string]string{"receiver": receiver},
|
||
"data": []interface{}{
|
||
// Total запрос
|
||
map[string]interface{}{
|
||
"refId": totalRefID,
|
||
"queryType": "lucene",
|
||
"relativeTimeRange": map[string]interface{}{"from": relFrom, "to": 0},
|
||
"datasourceUid": datasourceUID,
|
||
"model": map[string]interface{}{
|
||
"alias": "Total",
|
||
"bucketAggs": []interface{}{map[string]interface{}{"field": "@timestamp", "id": "2", "settings": map[string]interface{}{"interval": "1m"}, "type": "date_histogram"}},
|
||
"datasource": map[string]interface{}{"type": "grafana-opensearch-datasource", "uid": datasourceUID},
|
||
"format": "table",
|
||
"instant": false,
|
||
"intervalMs": 1000,
|
||
"luceneQueryType": "Metric",
|
||
"maxDataPoints": 43200,
|
||
"metrics": []interface{}{map[string]interface{}{"id": "1", "type": "count"}},
|
||
"query": fmt.Sprintf("response_status_code:[0 TO 599] AND %s", sidQuery),
|
||
"queryType": "lucene",
|
||
"range": true,
|
||
"refId": totalRefID,
|
||
"timeField": "@timestamp",
|
||
},
|
||
},
|
||
// Error запрос
|
||
map[string]interface{}{
|
||
"refId": errRefID,
|
||
"queryType": "lucene",
|
||
"relativeTimeRange": map[string]interface{}{"from": relFrom, "to": 0},
|
||
"datasourceUid": datasourceUID,
|
||
"model": map[string]interface{}{
|
||
"alias": codeName,
|
||
"bucketAggs": []interface{}{map[string]interface{}{"field": "@timestamp", "id": "2", "settings": map[string]interface{}{"interval": "1m"}, "type": "date_histogram"}},
|
||
"datasource": map[string]interface{}{"type": "grafana-opensearch-datasource", "uid": datasourceUID},
|
||
"format": "table",
|
||
"instant": false,
|
||
"intervalMs": 1000,
|
||
"luceneQueryType": "Metric",
|
||
"maxDataPoints": 43200,
|
||
"metrics": []interface{}{map[string]interface{}{"id": "1", "type": "count"}},
|
||
"query": fmt.Sprintf("%s:%s AND %s", statusField, codeRange, sidQuery),
|
||
"queryType": "lucene",
|
||
"range": true,
|
||
"refId": errRefID,
|
||
"timeField": "@timestamp",
|
||
},
|
||
},
|
||
// Reduce total
|
||
map[string]interface{}{
|
||
"refId": totalReduceID,
|
||
"relativeTimeRange": map[string]interface{}{"from": 0, "to": 0},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"datasource": map[string]interface{}{"name": "Expression", "type": "__expr__", "uid": "__expr__"},
|
||
"expression": totalRefID,
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"reducer": "last",
|
||
"refId": totalReduceID,
|
||
"type": "reduce",
|
||
},
|
||
},
|
||
// Reduce errors
|
||
map[string]interface{}{
|
||
"refId": errReduceID,
|
||
"relativeTimeRange": map[string]interface{}{"from": 0, "to": 0},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"datasource": map[string]interface{}{"name": "Expression", "type": "__expr__", "uid": "__expr__"},
|
||
"expression": errRefID,
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"reducer": "last",
|
||
"refId": errReduceID,
|
||
"type": "reduce",
|
||
},
|
||
},
|
||
// Math: процент ошибок
|
||
map[string]interface{}{
|
||
"refId": percentRefID,
|
||
"relativeTimeRange": map[string]interface{}{"from": 0, "to": 0},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"datasource": map[string]interface{}{"name": "Expression", "type": "__expr__", "uid": "__expr__"},
|
||
"expression": fmt.Sprintf("($%s / $%s) * 100", errReduceID, totalReduceID),
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"refId": percentRefID,
|
||
"type": "math",
|
||
},
|
||
},
|
||
// Threshold
|
||
map[string]interface{}{
|
||
"refId": thresholdRefID,
|
||
"relativeTimeRange": map[string]interface{}{"from": 0, "to": 0},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"conditions": []interface{}{
|
||
map[string]interface{}{
|
||
"evaluator": map[string]interface{}{"params": []interface{}{limitPct, 0}, "type": "gt"},
|
||
"operator": map[string]interface{}{"type": "and"},
|
||
"query": map[string]interface{}{"params": []string{}},
|
||
"reducer": map[string]interface{}{"params": []string{}, "type": "avg"},
|
||
"type": "query",
|
||
},
|
||
},
|
||
"datasource": map[string]interface{}{"name": "Expression", "type": "__expr__", "uid": "__expr__"},
|
||
"expression": percentRefID,
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"refId": thresholdRefID,
|
||
"type": "threshold",
|
||
},
|
||
},
|
||
},
|
||
}
|
||
rules = append(rules, rule)
|
||
}
|
||
|
||
return rules
|
||
}
|
||
|
||
// sendStaticAlerts отправляет статические алерты из шаблона (не зависят от БД)
|
||
func sendStaticAlerts(tmpl AlertRulesTemplate, config Config, dryRun bool) error {
|
||
if len(tmpl.StaticAlerts) == 0 {
|
||
return nil
|
||
}
|
||
|
||
defaults := tmpl.Defaults
|
||
folderUID := ""
|
||
|
||
if !dryRun {
|
||
uid, err := ensureAlertFolder(config)
|
||
if err != nil {
|
||
return fmt.Errorf("failed to ensure alerts folder for static alerts: %w", err)
|
||
}
|
||
folderUID = uid
|
||
}
|
||
|
||
group := "[SYSTEM]"
|
||
receiver := orDefault(config.AlertsReceiver, defaults.Receiver)
|
||
|
||
log.Printf("Alerts: sending %d static alert rules (concurrency: 5)", len(tmpl.StaticAlerts))
|
||
|
||
var tasks []alertTask
|
||
for _, alertDef := range tmpl.StaticAlerts {
|
||
rule := deepCopyMap(alertDef)
|
||
|
||
// Применяем defaults если не заданы в шаблоне
|
||
if _, ok := rule["noDataState"]; !ok {
|
||
rule["noDataState"] = defaults.NoDataState
|
||
}
|
||
if _, ok := rule["execErrState"]; !ok {
|
||
rule["execErrState"] = defaults.ExecErrState
|
||
}
|
||
rule["folderUID"] = folderUID
|
||
// Используем группу из шаблона если задана, иначе дефолтную
|
||
if rg, ok := rule["ruleGroup"].(string); !ok || rg == "" {
|
||
rule["ruleGroup"] = group
|
||
}
|
||
if _, ok := rule["notification_settings"]; !ok {
|
||
rule["notification_settings"] = map[string]string{"receiver": receiver}
|
||
}
|
||
|
||
title, _ := rule["title"].(string)
|
||
tasks = append(tasks, alertTask{rule: rule, label: title})
|
||
}
|
||
|
||
sent, failed := parallelUpsert(tasks, config, dryRun, 20)
|
||
log.Printf("Static alerts: sent=%d, failed=%d", sent, failed)
|
||
return nil
|
||
}
|
||
|
||
// deepCopyMap делает глубокую копию map через JSON
|
||
func deepCopyMap(src map[string]interface{}) map[string]interface{} {
|
||
data, _ := json.Marshal(src)
|
||
var dst map[string]interface{}
|
||
json.Unmarshal(data, &dst)
|
||
return dst
|
||
}
|
||
|
||
// upsertAlertRule создаёт или обновляет правило алерта по UID
|
||
func upsertAlertRule(rule map[string]interface{}, config Config) error {
|
||
uid, _ := rule["uid"].(string)
|
||
baseURL := strings.TrimRight(config.GrafanaURL, "/")
|
||
|
||
// Проверяем существует ли алерт с таким UID
|
||
existing, err := getAlertRule(uid, config)
|
||
if err == nil && existing != nil {
|
||
// Обновляем — PUT /api/v1/provisioning/alert-rules/{uid}
|
||
// Сохраняем version из существующего для оптимистичной блокировки
|
||
if v, ok := existing["version"]; ok {
|
||
rule["version"] = v
|
||
}
|
||
return putAlertRule(uid, rule, baseURL, config.GrafanaAPIKey)
|
||
}
|
||
|
||
// Создаём — POST /api/v1/provisioning/alert-rules
|
||
return postAlertRule(rule, baseURL, config.GrafanaAPIKey)
|
||
}
|
||
|
||
func getAlertRule(uid string, config Config) (map[string]interface{}, error) {
|
||
url := fmt.Sprintf("%s/api/v1/provisioning/alert-rules/%s",
|
||
strings.TrimRight(config.GrafanaURL, "/"), uid)
|
||
|
||
req, _ := http.NewRequest("GET", url, nil)
|
||
req.Header.Set("Authorization", "Bearer "+config.GrafanaAPIKey)
|
||
|
||
resp, err := http.DefaultClient.Do(req)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
if resp.StatusCode == 404 {
|
||
return nil, fmt.Errorf("not found")
|
||
}
|
||
|
||
body, _ := io.ReadAll(resp.Body)
|
||
var result map[string]interface{}
|
||
if err := json.Unmarshal(body, &result); err != nil {
|
||
return nil, err
|
||
}
|
||
return result, nil
|
||
}
|
||
|
||
func postAlertRule(rule map[string]interface{}, baseURL, apiKey string) error {
|
||
return sendAlertRequest("POST",
|
||
baseURL+"/api/v1/provisioning/alert-rules",
|
||
rule, apiKey)
|
||
}
|
||
|
||
func putAlertRule(uid string, rule map[string]interface{}, baseURL, apiKey string) error {
|
||
return sendAlertRequest("PUT",
|
||
fmt.Sprintf("%s/api/v1/provisioning/alert-rules/%s", baseURL, uid),
|
||
rule, apiKey)
|
||
}
|
||
|
||
func sendAlertRequest(method, url string, payload map[string]interface{}, apiKey string) error {
|
||
data, err := json.Marshal(payload)
|
||
if err != nil {
|
||
return fmt.Errorf("marshal error: %w", err)
|
||
}
|
||
|
||
req, err := http.NewRequest(method, url, bytes.NewBuffer(data))
|
||
if err != nil {
|
||
return err
|
||
}
|
||
req.Header.Set("Authorization", "Bearer "+apiKey)
|
||
req.Header.Set("Content-Type", "application/json")
|
||
req.Header.Set("X-Disable-Provenance", "true") // разрешить редактирование в UI
|
||
|
||
resp, err := http.DefaultClient.Do(req)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
body, _ := io.ReadAll(resp.Body)
|
||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||
return fmt.Errorf("status %d: %s", resp.StatusCode, string(body))
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// ensureAlertFolder возвращает UID папки для алертов, создаёт если нет
|
||
func ensureAlertFolder(config Config) (string, error) {
|
||
folder := orDefault(config.AlertsFolder, DefaultAlertsFolder)
|
||
if folder == "" {
|
||
return "", nil
|
||
}
|
||
baseURL := strings.TrimRight(config.GrafanaURL, "/")
|
||
|
||
// Ищем существующую папку
|
||
req, _ := http.NewRequest("GET", baseURL+"/api/folders", nil)
|
||
req.Header.Set("Authorization", "Bearer "+config.GrafanaAPIKey)
|
||
|
||
resp, err := http.DefaultClient.Do(req)
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
body, _ := io.ReadAll(resp.Body)
|
||
var folders []map[string]interface{}
|
||
if err := json.Unmarshal(body, &folders); err != nil {
|
||
return "", err
|
||
}
|
||
|
||
for _, f := range folders {
|
||
if title, _ := f["title"].(string); title == folder {
|
||
if uid, _ := f["uid"].(string); uid != "" {
|
||
return uid, nil
|
||
}
|
||
}
|
||
}
|
||
|
||
// Создаём папку
|
||
payload := map[string]interface{}{"title": folder}
|
||
data, _ := json.Marshal(payload)
|
||
req, _ = http.NewRequest("POST", baseURL+"/api/folders", bytes.NewBuffer(data))
|
||
req.Header.Set("Authorization", "Bearer "+config.GrafanaAPIKey)
|
||
req.Header.Set("Content-Type", "application/json")
|
||
|
||
resp, err = http.DefaultClient.Do(req)
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
body, _ = io.ReadAll(resp.Body)
|
||
var created map[string]interface{}
|
||
if err := json.Unmarshal(body, &created); err != nil {
|
||
return "", err
|
||
}
|
||
if uid, _ := created["uid"].(string); uid != "" {
|
||
log.Printf("Created alerts folder: %s (uid: %s)", folder, uid)
|
||
return uid, nil
|
||
}
|
||
return "", fmt.Errorf("failed to get folder UID")
|
||
}
|
||
|
||
// upsertContactPoint обновляет шаблон сообщения у существующего Contact Point.
|
||
// Не создаёт contact point — только обновляет message у уже настроенного в Grafana.
|
||
func upsertContactPoint(tmpl AlertRulesTemplate, config Config, dryRun bool) error {
|
||
if tmpl.ContactPoint.ReceiverName == "" || tmpl.ContactPoint.MessageTemplate == "" {
|
||
return nil
|
||
}
|
||
|
||
receiverName := orDefault(config.AlertsReceiver, tmpl.ContactPoint.ReceiverName)
|
||
baseURL := strings.TrimRight(config.GrafanaURL, "/")
|
||
|
||
// Получить список contact points и найти нужный по имени
|
||
req, _ := http.NewRequest("GET", baseURL+"/api/v1/provisioning/contact-points", nil)
|
||
req.Header.Set("Authorization", "Bearer "+config.GrafanaAPIKey)
|
||
|
||
resp, err := http.DefaultClient.Do(req)
|
||
if err != nil {
|
||
return fmt.Errorf("failed to list contact points: %w", err)
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
body, _ := io.ReadAll(resp.Body)
|
||
var contactPoints []map[string]interface{}
|
||
if err := json.Unmarshal(body, &contactPoints); err != nil {
|
||
return fmt.Errorf("failed to parse contact points: %w", err)
|
||
}
|
||
|
||
// Ищем все telegram интеграции с нужным именем и обновляем message
|
||
var updated int
|
||
for _, cp := range contactPoints {
|
||
name, _ := cp["name"].(string)
|
||
if name != receiverName {
|
||
continue
|
||
}
|
||
// Обновляем все типы интеграций (telegram, email, slack и т.д.)
|
||
uid, _ := cp["uid"].(string)
|
||
cpType, _ := cp["type"].(string)
|
||
log.Printf("Contact point integration found: name=%s type=%s uid=%s", name, cpType, uid)
|
||
settings, ok := cp["settings"].(map[string]interface{})
|
||
if !ok {
|
||
settings = make(map[string]interface{})
|
||
cp["settings"] = settings
|
||
}
|
||
|
||
// Удаляем секретные поля которые Grafana вернула как "configured"
|
||
// чтобы не затереть реальные значения при PUT
|
||
for _, secretField := range []string{"token", "bot_token", "api_token", "password", "secret", "url", "recipient"} {
|
||
if val, exists := settings[secretField]; exists {
|
||
if strVal, ok := val.(string); ok && (strVal == "configured" || strVal == "") {
|
||
delete(settings, secretField)
|
||
}
|
||
}
|
||
}
|
||
|
||
// Поле для шаблона зависит от типа интеграции
|
||
// Slack/Mattermost не поддерживают кастомный шаблон через API — пропускаем
|
||
switch cpType {
|
||
case "slack":
|
||
settings["text"] = tmpl.ContactPoint.MessageTemplate
|
||
settings["title"] = "" // очищаем автоматический заголовок
|
||
default: // telegram, email и др.
|
||
settings["message"] = tmpl.ContactPoint.MessageTemplate
|
||
}
|
||
|
||
if dryRun {
|
||
data, _ := json.MarshalIndent(cp, "", " ")
|
||
log.Printf("DRY RUN: Contact point update for '%s' (uid: %s):\n%s", receiverName, uid, string(data))
|
||
updated++
|
||
continue
|
||
}
|
||
|
||
payload, _ := json.Marshal(cp)
|
||
updReq, _ := http.NewRequest("PUT", baseURL+"/api/v1/provisioning/contact-points/"+uid, bytes.NewBuffer(payload))
|
||
updReq.Header.Set("Authorization", "Bearer "+config.GrafanaAPIKey)
|
||
updReq.Header.Set("Content-Type", "application/json")
|
||
updReq.Header.Set("X-Disable-Provenance", "true")
|
||
|
||
updResp, err := http.DefaultClient.Do(updReq)
|
||
if err != nil {
|
||
log.Printf("Warning: failed to update contact point '%s' (uid: %s): %v", receiverName, uid, err)
|
||
continue
|
||
}
|
||
updResp.Body.Close()
|
||
updated++
|
||
log.Printf("Contact point '%s' (uid: %s) message template updated", receiverName, uid)
|
||
}
|
||
|
||
if updated == 0 {
|
||
log.Printf("No telegram integrations found for contact point '%s'", receiverName)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// deleteObsoleteAlerts удаляет алерты из папки которые не входят в актуальный набор UID.
|
||
func deleteObsoleteAlerts(activeUIDs map[string]bool, folderUID string, config Config) {
|
||
baseURL := strings.TrimRight(config.GrafanaURL, "/")
|
||
|
||
req, err := http.NewRequest("GET", baseURL+"/api/v1/provisioning/alert-rules", nil)
|
||
if err != nil {
|
||
log.Printf("Warning: failed to list alert rules: %v", err)
|
||
return
|
||
}
|
||
req.Header.Set("Authorization", "Bearer "+config.GrafanaAPIKey)
|
||
|
||
resp, err := http.DefaultClient.Do(req)
|
||
if err != nil {
|
||
log.Printf("Warning: failed to list alert rules: %v", err)
|
||
return
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
body, _ := io.ReadAll(resp.Body)
|
||
var rules []map[string]interface{}
|
||
if err := json.Unmarshal(body, &rules); err != nil {
|
||
log.Printf("Warning: failed to parse alert rules list: %v", err)
|
||
return
|
||
}
|
||
|
||
deleted := 0
|
||
for _, rule := range rules {
|
||
uid, _ := rule["uid"].(string)
|
||
ruleFolderUID, _ := rule["folderUID"].(string)
|
||
if ruleFolderUID != folderUID {
|
||
continue
|
||
}
|
||
if activeUIDs[uid] {
|
||
continue
|
||
}
|
||
delReq, _ := http.NewRequest("DELETE", fmt.Sprintf("%s/api/v1/provisioning/alert-rules/%s", baseURL, uid), nil)
|
||
delReq.Header.Set("Authorization", "Bearer "+config.GrafanaAPIKey)
|
||
delResp, err := http.DefaultClient.Do(delReq)
|
||
if err != nil {
|
||
log.Printf("Warning: failed to delete alert %s: %v", uid, err)
|
||
continue
|
||
}
|
||
delResp.Body.Close()
|
||
title, _ := rule["title"].(string)
|
||
log.Printf("Deleted obsolete alert: %s (%s)", title, uid)
|
||
deleted++
|
||
}
|
||
if deleted > 0 {
|
||
log.Printf("Alerts: deleted %d obsolete rules", deleted)
|
||
}
|
||
}
|
||
|
||
// orDefault возвращает val если не пустой, иначе def
|
||
func orDefault(val, def string) string {
|
||
if val != "" {
|
||
return val
|
||
}
|
||
return def
|
||
}
|
||
|
||
// buildWAFBlockAlertRule строит алерт на количество 418 ответов за минуту (OpenSearch).
|
||
func buildWAFBlockAlertRule(client ClientData, domain DomainInfo, tmpl AlertRulesTemplate, config Config) map[string]interface{} {
|
||
datasourceUID := config.AlertsDatasourceUID
|
||
defaults := tmpl.Defaults
|
||
receiver := orDefault(config.AlertsReceiver, defaults.Receiver)
|
||
group := client.ClientTitle
|
||
|
||
uid := generateAlertUID(client.ClientTitle+"_"+domain.SID, "waf_block")
|
||
title := fmt.Sprintf("%s WAF %s / %s SID:%s", client.PtafFallbackCode, client.ClientTitle, domain.DomainName, domain.SID)
|
||
|
||
clientDomain := fmt.Sprintf("%s / %s (SID:%s)", client.ClientTitle, domain.DomainName, domain.SID)
|
||
summary := fmt.Sprintf("🟡 %s — %s ошибка WAF - кол-во штук за 1 минуту.", clientDomain, client.PtafFallbackCode)
|
||
thresholdAnnotation := fmt.Sprintf(">%d штук(и) для response_status_code (ответ от PTAF)", domain.PtafFallbackCodeAlertCount)
|
||
thresholdCountAnnotation := fmt.Sprintf("%d шт.", domain.PtafFallbackCodeAlertCount)
|
||
|
||
sidQuery := "SID:" + domain.SID
|
||
dataRefID := "A"
|
||
reduceRefID := "_Reduce"
|
||
thresholdRefID := "Превышение"
|
||
relFrom := int64(120)
|
||
|
||
return map[string]interface{}{
|
||
"uid": uid,
|
||
"title": title,
|
||
"condition": thresholdRefID,
|
||
"noDataState": defaults.NoDataState,
|
||
"execErrState": defaults.ExecErrState,
|
||
"isPaused": false,
|
||
"folderUID": "",
|
||
"ruleGroup": group,
|
||
"for": "1m",
|
||
"annotations": map[string]string{
|
||
"summary": summary,
|
||
"client_domain": clientDomain,
|
||
"code_field": fmt.Sprintf("%s WAF", client.PtafFallbackCode),
|
||
"for_duration": "1m",
|
||
"threshold": thresholdAnnotation,
|
||
"field_info": "для response_status_code (ответ от PTAF)",
|
||
"runbook": "проверить логи docker и эскалировать на аналитика.",
|
||
"severity": "warning",
|
||
"value_unit": "шт.",
|
||
"threshold_count": thresholdCountAnnotation,
|
||
},
|
||
"notification_settings": map[string]string{"receiver": receiver},
|
||
"data": []interface{}{
|
||
map[string]interface{}{
|
||
"refId": dataRefID,
|
||
"queryType": "lucene",
|
||
"relativeTimeRange": map[string]interface{}{"from": relFrom, "to": 0},
|
||
"datasourceUid": datasourceUID,
|
||
"model": map[string]interface{}{
|
||
"datasource": map[string]interface{}{"type": "grafana-opensearch-datasource", "uid": datasourceUID},
|
||
"query": fmt.Sprintf("response_status_code:%s AND %s", client.PtafFallbackCode, sidQuery),
|
||
"luceneQueryType": "Metric",
|
||
"timeField": "@timestamp",
|
||
"metrics": []interface{}{
|
||
map[string]interface{}{"id": "1", "type": "count"},
|
||
},
|
||
"bucketAggs": []interface{}{
|
||
map[string]interface{}{
|
||
"id": "2",
|
||
"type": "date_histogram",
|
||
"field": "@timestamp",
|
||
"settings": map[string]interface{}{
|
||
"interval": "1m",
|
||
"min_doc_count": "1",
|
||
},
|
||
},
|
||
},
|
||
"refId": dataRefID,
|
||
},
|
||
},
|
||
map[string]interface{}{
|
||
"refId": reduceRefID,
|
||
"relativeTimeRange": map[string]interface{}{"from": 0, "to": 0},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"datasource": map[string]interface{}{"name": "Expression", "type": "__expr__", "uid": "__expr__"},
|
||
"expression": dataRefID,
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"reducer": "last",
|
||
"refId": reduceRefID,
|
||
"type": "reduce",
|
||
},
|
||
},
|
||
map[string]interface{}{
|
||
"refId": thresholdRefID,
|
||
"relativeTimeRange": map[string]interface{}{"from": 0, "to": 0},
|
||
"datasourceUid": "__expr__",
|
||
"model": map[string]interface{}{
|
||
"conditions": []interface{}{
|
||
map[string]interface{}{
|
||
"evaluator": map[string]interface{}{"params": []interface{}{domain.PtafFallbackCodeAlertCount, 0}, "type": "gt"},
|
||
"operator": map[string]interface{}{"type": "and"},
|
||
"query": map[string]interface{}{"params": []string{}},
|
||
"reducer": map[string]interface{}{"params": []string{}, "type": "avg"},
|
||
"type": "query",
|
||
},
|
||
},
|
||
"datasource": map[string]interface{}{"name": "Expression", "type": "__expr__", "uid": "__expr__"},
|
||
"expression": reduceRefID,
|
||
"intervalMs": 1000,
|
||
"maxDataPoints": 43200,
|
||
"refId": thresholdRefID,
|
||
"type": "threshold",
|
||
},
|
||
},
|
||
},
|
||
}
|
||
}
|