Added VL alerts

This commit is contained in:
Magnus Root 2026-04-15 17:27:47 +03:00
parent 3bf36c95dc
commit 67dd80951c
4 changed files with 483 additions and 30 deletions

380
alerts_vl.go Normal file
View file

@ -0,0 +1,380 @@
package main
import (
"fmt"
"log"
"sort"
"strings"
)
// buildRPSAlertRuleVL строит правило алерта на RPS для VictoriaLogs datasource.
func buildRPSAlertRuleVL(client ClientData, tmpl AlertRulesTemplate, config Config) map[string]interface{} {
datasourceUID := config.VLDatasourceUID
datasourceType := DefaultVLDatasourceType
// SID-запрос в формате LogsQL
sidMap := make(map[string]bool)
for _, d := range client.Domains {
sidMap[d.SID] = true
}
uniqueSIDs := make([]string, 0, len(sidMap))
for s := range sidMap {
uniqueSIDs = append(uniqueSIDs, s)
}
sort.Strings(uniqueSIDs)
sidQuery := buildTenantSIDQueryVL(client.Domains)
defaults := tmpl.Defaults
receiver := orDefault(config.AlertsReceiver, defaults.Receiver)
group := orDefault(config.AlertsGroup, defaults.Group)
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)
_ = uniqueSIDs // используется через sidQuery
uid := generateAlertUID(client.ClientTitle+"_vl", "rps")
dataRefID := "A"
rpsRefID := "RPS"
reduceRefID := "_Reduce"
thresholdRefID := "Превышение"
// LogsQL запрос: count за минуту
expr := fmt.Sprintf("%s log_type:nginx-access-PTAF | stats by (_time:1m) count() rps", sidQuery)
return map[string]interface{}{
"uid": uid,
"title": fmt.Sprintf("VL 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: VL statsRange запрос
map[string]interface{}{
"refId": dataRefID,
"queryType": "statsRange",
"relativeTimeRange": map[string]interface{}{
"from": timeRangeFrom,
"to": 0,
},
"datasourceUid": datasourceUID,
"model": map[string]interface{}{
"datasource": map[string]interface{}{
"type": datasourceType,
"uid": datasourceUID,
},
"editorMode": "code",
"expr": expr,
"legendFormat": fmt.Sprintf("RPS %s", client.ClientTitle),
"maxLines": 1000,
"queryType": "statsRange",
"refId": dataRefID,
"intervalMs": 60000,
"maxDataPoints": 43200,
},
},
// Шаг 2: math $A / 60 → RPS (интервал 1m = 60 секунд)
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",
},
},
},
}
}
// buildErrorRateAlertRuleVL строит правило алерта на процент ошибок для VictoriaLogs.
func buildErrorRateAlertRuleVL(client ClientData, errCode, limitPct int, tmpl AlertRulesTemplate, config Config) map[string]interface{} {
datasourceUID := config.VLDatasourceUID
datasourceType := DefaultVLDatasourceType
defaults := tmpl.Defaults
receiver := orDefault(config.AlertsReceiver, defaults.Receiver)
group := orDefault(config.AlertsGroup, defaults.Group)
sidQuery := buildTenantSIDQueryVL(client.Domains)
codeFrom, codeTo := 400, 499
codeName := "4xx"
if errCode == 5 {
codeFrom, codeTo = 500, 599
codeName = "5xx"
}
uid := generateAlertUID(client.ClientTitle+"_vl", codeName)
totalRefID := "Total"
errRefID := "Errors"
totalReduceID := "_TotalReduce"
errReduceID := "_ErrorsReduce"
percentRefID := "Значение"
thresholdRefID := "Превышение"
summary := fmt.Sprintf("🚨 %s — %s ошибок >%d%% трафика", client.ClientTitle, codeName, limitPct)
relFrom := int64(900)
forDuration := "1m"
if errCode == 4 {
forDuration = "3m"
}
totalExpr := fmt.Sprintf("%s log_type:nginx-access-PTAF | stats by (_time:1m) count() hits", sidQuery)
errExpr := fmt.Sprintf("%s log_type:nginx-access-PTAF response_status_code:>=%d response_status_code:<=%d | stats by (_time:1m) count() hits", sidQuery, codeFrom, codeTo)
makeVLModel := func(refID, expr, legend string) map[string]interface{} {
return map[string]interface{}{
"datasource": map[string]interface{}{
"type": datasourceType,
"uid": datasourceUID,
},
"editorMode": "code",
"expr": expr,
"legendFormat": legend,
"maxLines": 1000,
"queryType": "statsRange",
"refId": refID,
"intervalMs": 60000,
"maxDataPoints": 43200,
}
}
makeReduce := func(refID, expression, reducer string) map[string]interface{} {
return map[string]interface{}{
"datasource": map[string]interface{}{"name": "Expression", "type": "__expr__", "uid": "__expr__"},
"expression": expression,
"intervalMs": 1000,
"maxDataPoints": 43200,
"reducer": reducer,
"refId": refID,
"type": "reduce",
}
}
return map[string]interface{}{
"uid": uid,
"title": fmt.Sprintf("VL %s %s", codeName, client.ClientTitle),
"condition": thresholdRefID,
"noDataState": defaults.NoDataState,
"execErrState": defaults.ExecErrState,
"isPaused": false,
"folderUID": "",
"ruleGroup": group,
"for": forDuration,
"annotations": map[string]string{"summary": summary},
"notification_settings": map[string]string{"receiver": receiver},
"data": []interface{}{
map[string]interface{}{
"refId": totalRefID,
"queryType": "statsRange",
"relativeTimeRange": map[string]interface{}{"from": relFrom, "to": 0},
"datasourceUid": datasourceUID,
"model": makeVLModel(totalRefID, totalExpr, "Total"),
},
map[string]interface{}{
"refId": errRefID,
"queryType": "statsRange",
"relativeTimeRange": map[string]interface{}{"from": relFrom, "to": 0},
"datasourceUid": datasourceUID,
"model": makeVLModel(errRefID, errExpr, codeName),
},
map[string]interface{}{
"refId": totalReduceID,
"relativeTimeRange": map[string]interface{}{"from": 0, "to": 0},
"datasourceUid": "__expr__",
"model": makeReduce(totalReduceID, totalRefID, "sum"),
},
map[string]interface{}{
"refId": errReduceID,
"relativeTimeRange": map[string]interface{}{"from": 0, "to": 0},
"datasourceUid": "__expr__",
"model": makeReduce(errReduceID, errRefID, "sum"),
},
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",
},
},
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",
},
},
},
}
}
// generateAndSendAlertsVL генерирует и отправляет алерты для VictoriaLogs.
func generateAndSendAlertsVL(clients map[string]ClientData, tmpl AlertRulesTemplate, config Config, dryRun bool) error {
if config.VLDatasourceUID == "" {
log.Printf("Skipping VL alerts: VLDatasourceUID not configured")
return nil
}
folderUID, err := ensureAlertFolder(config)
if err != nil {
return fmt.Errorf("failed to ensure alert folder: %w", err)
}
var tasks []alertTask
for _, client := range clients {
if len(client.Domains) == 0 {
continue
}
// RPS алерт
if client.RPSLimit > 0 {
rule := buildRPSAlertRuleVL(client, tmpl, config)
rule["folderUID"] = folderUID
tasks = append(tasks, alertTask{
rule: rule,
label: fmt.Sprintf("VL RPS %s", client.ClientTitle),
})
}
// Error rate алерты (4xx и 5xx)
for _, errCode := range []int{4, 5} {
limitPct := client.Limit4xx
if errCode == 5 {
limitPct = client.Limit5xx
}
if limitPct <= 0 {
continue
}
rule := buildErrorRateAlertRuleVL(client, errCode, limitPct, tmpl, config)
rule["folderUID"] = folderUID
tasks = append(tasks, alertTask{
rule: rule,
label: fmt.Sprintf("VL %dxx %s", errCode, client.ClientTitle),
})
}
}
sent, failed := parallelUpsert(tasks, config, dryRun, 5)
log.Printf("VL Alerts: sent=%d failed=%d", sent, failed)
return nil
}

View file

@ -61,7 +61,7 @@ const (
// Alerts settings // Alerts settings
DefaultAlertsFolder = "WAF - PTAF" DefaultAlertsFolder = "WAF - PTAF"
DefaultAlertsReceiver = "Telegram PTAF Grafana" DefaultAlertsReceiver = "tg+mail+mattermost PTAF Grafana"
DefaultAlertsGroup = "PTAF Grafana" DefaultAlertsGroup = "PTAF Grafana"
DefaultAlertsDatasource = "af84zsvlp9blsa" // UID datasource OpenSearch DefaultAlertsDatasource = "af84zsvlp9blsa" // UID datasource OpenSearch
@ -104,6 +104,10 @@ type Config struct {
AlertsReceiver string // Receiver (канал уведомлений) для алертов AlertsReceiver string // Receiver (канал уведомлений) для алертов
AlertsGroup string // Имя группы алертов в Grafana AlertsGroup string // Имя группы алертов в Grafana
AlertsDatasourceUID string // UID datasource для алертов OpenSearch AlertsDatasourceUID string // UID datasource для алертов OpenSearch
AlertsOSEnabled bool // Генерировать алерты для OpenSearch (default: true)
AlertsVLEnabled bool // Генерировать алерты для VictoriaLogs (default: true)
DashboardOSEnabled bool // Генерировать дашборд для OpenSearch (default: true)
DashboardVLEnabled bool // Генерировать дашборд для VictoriaLogs (default: true)
// VictoriaLogs второй дашборд // VictoriaLogs второй дашборд
VLDatasourceUID string // UID VictoriaLogs datasource VLDatasourceUID string // UID VictoriaLogs datasource

64
main.go
View file

@ -131,36 +131,46 @@ func main() {
dashboard := generateSingleDashboard(validClients, template, config) dashboard := generateSingleDashboard(validClients, template, config)
if config.DryRun { if config.DryRun {
logDryRun(config.DashboardTitle, dashboard) if config.DashboardOSEnabled {
logDryRun(config.DashboardTitle, dashboard)
}
if alertTmpl != nil { if alertTmpl != nil {
log.Printf("DRY RUN: Generating alert rules preview...") log.Printf("DRY RUN: Generating alert rules preview...")
_ = generateAndSendAlerts(validClients, *alertTmpl, config, true) if config.AlertsOSEnabled {
_ = generateAndSendAlerts(validClients, *alertTmpl, config, true)
}
if config.AlertsVLEnabled {
_ = generateAndSendAlertsVL(validClients, *alertTmpl, config, true)
}
} }
log.Printf("=== Summary ===") log.Printf("=== Summary ===")
log.Printf("Clients: %d", len(validClients)) log.Printf("Clients: %d", len(validClients))
log.Printf("Skipped: %d", skippedCount) log.Printf("Skipped: %d", skippedCount)
log.Printf("DRY RUN: Dashboard generated but not sent to Grafana") log.Printf("DRY RUN: Dashboard generated but not sent to Grafana")
if config.VLDatasourceUID != "" { if config.VLDatasourceUID != "" && config.DashboardVLEnabled {
vlDashboard := generateVLDashboard(validClients, template, config) vlDashboard := generateVLDashboard(validClients, template, config)
logDryRun(config.VLDashboardTitle, vlDashboard) logDryRun(config.VLDashboardTitle, vlDashboard)
} }
os.Exit(0) os.Exit(0)
} }
// Send to Grafana // Send OS dashboard to Grafana
if err := sendToGrafana(dashboard, config); err != nil { if config.DashboardOSEnabled {
log.Printf("Error sending dashboard: %v", err) if err := sendToGrafana(dashboard, config); err != nil {
log.Printf("=== Summary ===") log.Printf("Error sending dashboard: %v", err)
log.Printf("Clients: %d", len(validClients)) log.Printf("=== Summary ===")
log.Printf("Skipped: %d", skippedCount) log.Printf("Clients: %d", len(validClients))
log.Printf("Status: FAILED") log.Printf("Skipped: %d", skippedCount)
os.Exit(1) log.Printf("Status: FAILED")
os.Exit(1)
}
log.Printf("Successfully created/updated dashboard: %s", config.DashboardTitle)
} else {
log.Printf("Dashboard OS: disabled (DASHBOARD_OS_ENABLED=false)")
} }
log.Printf("Successfully created/updated dashboard: %s", config.DashboardTitle)
// Generate VictoriaLogs dashboard if configured // Generate VictoriaLogs dashboard if configured
if config.VLDatasourceUID != "" { if config.VLDatasourceUID != "" && config.DashboardVLEnabled {
log.Printf("Generating VictoriaLogs dashboard: %s", config.VLDashboardTitle) log.Printf("Generating VictoriaLogs dashboard: %s", config.VLDashboardTitle)
vlDashboard := generateVLDashboard(validClients, template, config) vlDashboard := generateVLDashboard(validClients, template, config)
if err := sendVLDashboardToGrafana(vlDashboard, config); err != nil { if err := sendVLDashboardToGrafana(vlDashboard, config); err != nil {
@ -168,14 +178,32 @@ func main() {
} else { } else {
log.Printf("Successfully created/updated VL dashboard: %s", config.VLDashboardTitle) log.Printf("Successfully created/updated VL dashboard: %s", config.VLDashboardTitle)
} }
} else if config.VLDatasourceUID != "" {
log.Printf("Dashboard VL: disabled (DASHBOARD_VL_ENABLED=false)")
} }
// Generate and send alert rules // Generate and send alert rules
if alertTmpl != nil { if alertTmpl != nil {
if alertsChanged || config.ForceRun { if alertsChanged || config.ForceRun {
if err := generateAndSendAlerts(validClients, *alertTmpl, config, false); err != nil { osOK := true
log.Printf("Warning: alerts generation failed: %v", err) if config.AlertsOSEnabled {
if err := generateAndSendAlerts(validClients, *alertTmpl, config, false); err != nil {
log.Printf("Warning: alerts generation failed: %v", err)
osOK = false
}
} else { } else {
log.Printf("Alerts OS: disabled (ALERTS_OS_ENABLED=false)")
}
if config.AlertsVLEnabled {
if err := generateAndSendAlertsVL(validClients, *alertTmpl, config, false); err != nil {
log.Printf("Warning: VL alerts generation failed: %v", err)
osOK = false
}
} else {
log.Printf("Alerts VL: disabled (ALERTS_VL_ENABLED=false)")
}
// Сохраняем хэш если хотя бы одна из систем алертов не вернула ошибку
if osOK {
stateManager.UpdateAlertsHash(validClients) stateManager.UpdateAlertsHash(validClients)
stateManager.UpdateAlertTemplateHash(alertTemplateHash) stateManager.UpdateAlertTemplateHash(alertTemplateHash)
} }
@ -267,6 +295,10 @@ func parseFlags() Config {
flag.BoolVar(&config.DryRun, "dry-run", false, "Generate JSON but don't send to Grafana") flag.BoolVar(&config.DryRun, "dry-run", false, "Generate JSON but don't send to Grafana")
flag.StringVar(&config.StateFile, "state-file", DefaultStateFile, "Path to state file") flag.StringVar(&config.StateFile, "state-file", DefaultStateFile, "Path to state file")
flag.BoolVar(&config.ForceRun, "force", false, "Force regeneration even if no changes detected") flag.BoolVar(&config.ForceRun, "force", false, "Force regeneration even if no changes detected")
config.AlertsOSEnabled = getEnvOrDefault("ALERTS_OS_ENABLED", "true") != "false"
config.AlertsVLEnabled = getEnvOrDefault("ALERTS_VL_ENABLED", "true") != "false"
config.DashboardOSEnabled = getEnvOrDefault("DASHBOARD_OS_ENABLED", "true") != "false"
config.DashboardVLEnabled = getEnvOrDefault("DASHBOARD_VL_ENABLED", "true") != "false"
flag.Parse() flag.Parse()

View file

@ -194,7 +194,7 @@ func buildVLRPSTargets(cfg map[string]interface{}, sidQuery string, _ string, da
if alias == "" { if alias == "" {
alias = client.ClientTitle alias = client.ClientTitle
} }
expr := fmt.Sprintf("%s log_type:access | stats by (_time:$__interval) count() rps", vlSID) expr := fmt.Sprintf("%s log_type:nginx-access-PTAF | stats by (_time:$__interval) count() rps", vlSID)
return []interface{}{ return []interface{}{
makeVLHiddenStatsTarget("A", expr, datasourceUID), makeVLHiddenStatsTarget("A", expr, datasourceUID),
@ -208,7 +208,7 @@ func buildVLRPSVarTargets(cfg map[string]interface{}, _ string, _ string, dataso
if alias == "" { if alias == "" {
alias = "rps" alias = "rps"
} }
expr := "SID:${rps_sid} log_type:access | stats by (_time:$__interval) count() rps" expr := "SID:${rps_sid} log_type:nginx-access-PTAF | stats by (_time:$__interval) count() rps"
return []interface{}{ return []interface{}{
makeVLHiddenStatsTarget("A", expr, datasourceUID), makeVLHiddenStatsTarget("A", expr, datasourceUID),
@ -221,14 +221,14 @@ func buildVLRPSMultiTargets(cfg map[string]interface{}, _ string, _ string, data
// Суммарная линия // Суммарная линия
allSID := buildTenantSIDQueryVL(client.Domains) allSID := buildTenantSIDQueryVL(client.Domains)
totalExpr := fmt.Sprintf("%s log_type:access | stats by (_time:$__interval) count() rps", allSID) totalExpr := fmt.Sprintf("%s log_type:nginx-access-PTAF | stats by (_time:$__interval) count() rps", allSID)
targets = append(targets, makeVLHiddenStatsTarget("A", totalExpr, datasourceUID)) targets = append(targets, makeVLHiddenStatsTarget("A", totalExpr, datasourceUID))
targets = append(targets, makeMathTargetWithAlias("$A / $__interval_ms * 1000", "Total", "Total")) targets = append(targets, makeMathTargetWithAlias("$A / $__interval_ms * 1000", "Total", "Total"))
// Линия на каждый домен // Линия на каждый домен
for i, domain := range client.Domains { for i, domain := range client.Domains {
rawRefID := fmt.Sprintf("D%d", i) rawRefID := fmt.Sprintf("D%d", i)
domainExpr := fmt.Sprintf("SID:%s log_type:access | stats by (_time:$__interval) count() rps", domain.SID) domainExpr := fmt.Sprintf("SID:%s log_type:nginx-access-PTAF | stats by (_time:$__interval) count() rps", domain.SID)
legend := domain.SID + "_" + domain.DomainName legend := domain.SID + "_" + domain.DomainName
if domain.DomainName == "" { if domain.DomainName == "" {
legend = domain.SID legend = domain.SID
@ -258,16 +258,28 @@ func buildVLMultiTargets(cfg map[string]interface{}, sidQuery string, _ string,
} }
refID := getString(tMap, "ref_id") + "_" + uidSuffix refID := getString(tMap, "ref_id") + "_" + uidSuffix
alias := getString(tMap, "alias") alias := getString(tMap, "alias")
suffix := getString(tMap, "query_suffix")
// Конвертируем suffix из Lucene в LogsQL и добавляем к sidQuery
vlSuffix := luceneToLogsQL(strings.TrimPrefix(suffix, " "))
// Если есть vl_base — используем его как базовый запрос вместо sidQuery
vlBase := getString(tMap, "vl_base")
vlSuffixRaw := getString(tMap, "vl_suffix")
var expr string var expr string
if vlSuffix != "" { if vlBase != "" {
expr = fmt.Sprintf("%s %s | stats by (_time:$__interval) count() hits", sidQuery, vlSuffix) // Панель с отдельным base query на каждый target (например nginx vs angie)
if vlSuffixRaw != "" {
expr = fmt.Sprintf("%s %s | stats by (_time:$__interval) count() hits", vlBase, vlSuffixRaw)
} else {
expr = fmt.Sprintf("%s | stats by (_time:$__interval) count() hits", vlBase)
}
} else { } else {
expr = fmt.Sprintf("%s | stats by (_time:$__interval) count() hits", sidQuery) suffix := getString(tMap, "query_suffix")
vlSuffix := luceneToLogsQL(strings.TrimPrefix(suffix, " "))
// Добавляем log_type:nginx-access-PTAF чтобы не смешивать nginx и angie логи
baseWithType := sidQuery + " log_type:nginx-access-PTAF"
if vlSuffix != "" {
expr = fmt.Sprintf("%s %s | stats by (_time:$__interval) count() hits", baseWithType, vlSuffix)
} else {
expr = fmt.Sprintf("%s | stats by (_time:$__interval) count() hits", baseWithType)
}
} }
legend := alias legend := alias
@ -406,7 +418,8 @@ func buildVLBucketLogsTargets(cfg map[string]interface{}, _ string, _ string, da
// Иначе конвертируем base_query из Lucene // Иначе конвертируем base_query из Lucene
var baseQuery string var baseQuery string
if vlRaw := getString(cfg, "vl_base_query"); vlRaw != "" { if vlRaw := getString(cfg, "vl_base_query"); vlRaw != "" {
baseQuery = vlRaw // Убираем фильтры вида field:* — они ломают группировку в VL
baseQuery = removeWildcardFilters(vlRaw)
} else { } else {
raw := getString(cfg, "base_query") raw := getString(cfg, "base_query")
raw = strings.ReplaceAll(raw, "SID:(${SID:lucene})", "SID:in(${SID:csv})") raw = strings.ReplaceAll(raw, "SID:(${SID:lucene})", "SID:in(${SID:csv})")
@ -452,9 +465,18 @@ func buildVLBucketLogsTargets(cfg map[string]interface{}, _ string, _ string, da
// legendFormat с именем поля — плагин подставит значение поля группировки // legendFormat с именем поля — плагин подставит значение поля группировки
legendFormat := fmt.Sprintf("{{%s}}", termsField) legendFormat := fmt.Sprintf("{{%s}}", termsField)
return []interface{}{ targets := []interface{}{
makeVLStatsTarget("A", expr, datasourceUID, legendFormat), makeVLStatsTarget("A", expr, datasourceUID, legendFormat),
} }
// Если задан unknown_label — добавляем отдельный target для записей с пустым значением поля
unknownLabel := getString(cfg, "unknown_label")
if unknownLabel != "" {
unknownQuery := fmt.Sprintf("%s %s:=\"\" | stats count() hits", baseQuery, termsField)
targets = append(targets, makeVLStatsTarget("B", unknownQuery, datasourceUID, unknownLabel))
}
return targets
} }
// unescapeLucene убирает Lucene-экранирование символов вида \X → X // unescapeLucene убирает Lucene-экранирование символов вида \X → X
@ -473,3 +495,18 @@ func unescapeLucene(q string) string {
} }
return result.String() return result.String()
} }
// removeWildcardFilters убирает из LogsQL запроса фильтры вида field:* и field:"*"
// которые ломают группировку — VL возвращает один фрейм вместо разбивки по полю.
func removeWildcardFilters(q string) string {
words := strings.Fields(q)
result := []string{}
for _, w := range words {
// Пропускаем field:* и field:"*"
if strings.HasSuffix(w, ":*") || strings.HasSuffix(w, `:"*"`) {
continue
}
result = append(result, w)
}
return strings.Join(result, " ")
}