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 { 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 { 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" } summary := fmt.Sprintf( "🔴 %s / %s (SID:%s) — %s %s ошибок трафика держится > %s.", client.ClientTitle, domain.DomainName, domain.SID, codeName, fieldShort, forDuration, ) 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, }, "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, не от origin)", 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, не от origin)", "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", }, }, }, } }