diff --git a/docs/sources/setup-grafana/configure-grafana/feature-toggles/index.md b/docs/sources/setup-grafana/configure-grafana/feature-toggles/index.md index 0c9d33a8e8f..4dc078e5742 100644 --- a/docs/sources/setup-grafana/configure-grafana/feature-toggles/index.md +++ b/docs/sources/setup-grafana/configure-grafana/feature-toggles/index.md @@ -111,6 +111,7 @@ Alpha features might be changed or removed without prior notice. | `pyroscopeFlameGraph` | Changes flame graph to pyroscope one | | `pluginsAPIManifestKey` | Use grafana.com API to retrieve the public manifest key | | `opensearchDetectVersion` | Enable version detection in OpenSearch | +| `alertingLokiRangeToInstant` | Rewrites eligible loki range queries to instant queries | ## Development feature toggles diff --git a/packages/grafana-data/src/types/featureToggles.gen.ts b/packages/grafana-data/src/types/featureToggles.gen.ts index a6c0acb89d3..9980736ca3a 100644 --- a/packages/grafana-data/src/types/featureToggles.gen.ts +++ b/packages/grafana-data/src/types/featureToggles.gen.ts @@ -97,4 +97,5 @@ export interface FeatureToggles { advancedDataSourcePicker?: boolean; opensearchDetectVersion?: boolean; enableDatagridEditing?: boolean; + alertingLokiRangeToInstant?: boolean; } diff --git a/pkg/services/featuremgmt/registry.go b/pkg/services/featuremgmt/registry.go index 82648556b2f..867265d0eb2 100644 --- a/pkg/services/featuremgmt/registry.go +++ b/pkg/services/featuremgmt/registry.go @@ -533,5 +533,12 @@ var ( State: FeatureStateBeta, Owner: grafanaBiSquad, }, + { + Name: "alertingLokiRangeToInstant", + Description: "Rewrites eligible loki range queries to instant queries", + State: FeatureStateAlpha, + FrontendOnly: false, + Owner: grafanaAlertingSquad, + }, } ) diff --git a/pkg/services/featuremgmt/toggles_gen.csv b/pkg/services/featuremgmt/toggles_gen.csv index 4eefdac12a2..249e9edc4d8 100644 --- a/pkg/services/featuremgmt/toggles_gen.csv +++ b/pkg/services/featuremgmt/toggles_gen.csv @@ -78,3 +78,4 @@ pluginsAPIManifestKey,alpha,@grafana/plugins-platform-backend,false,false,false, advancedDataSourcePicker,stable,@grafana/dashboards-squad,false,false,false,true opensearchDetectVersion,alpha,@grafana/aws-plugins,false,false,false,true enableDatagridEditing,beta,@grafana/grafana-bi-squad,false,false,false,true +alertingLokiRangeToInstant,alpha,@grafana/alerting-squad,false,false,false,false diff --git a/pkg/services/featuremgmt/toggles_gen.go b/pkg/services/featuremgmt/toggles_gen.go index 0f22332ae3e..7e0e6a4d10c 100644 --- a/pkg/services/featuremgmt/toggles_gen.go +++ b/pkg/services/featuremgmt/toggles_gen.go @@ -322,4 +322,8 @@ const ( // FlagEnableDatagridEditing // Enables the edit functionality in the datagrid panel FlagEnableDatagridEditing = "enableDatagridEditing" + + // FlagAlertingLokiRangeToInstant + // Rewrites eligible loki range queries to instant queries + FlagAlertingLokiRangeToInstant = "alertingLokiRangeToInstant" ) diff --git a/pkg/services/ngalert/store/alert_rule.go b/pkg/services/ngalert/store/alert_rule.go index 04b8c5ac404..2785c015d85 100644 --- a/pkg/services/ngalert/store/alert_rule.go +++ b/pkg/services/ngalert/store/alert_rule.go @@ -8,6 +8,7 @@ import ( "github.com/grafana/grafana/pkg/infra/db" "github.com/grafana/grafana/pkg/services/dashboards" + "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/services/folder" "github.com/grafana/grafana/pkg/services/guardian" ngmodels "github.com/grafana/grafana/pkg/services/ngalert/models" @@ -476,7 +477,7 @@ func (st DBstore) GetAlertRulesForScheduling(ctx context.Context, query *ngmodel st.Logger.Error("unable to close rows session", "error", err) } }() - + lokiRangeToInstantEnabled := st.FeatureToggles.IsEnabled(featuremgmt.FlagAlertingLokiRangeToInstant) // Deserialize each rule separately in case any of them contain invalid JSON. for rows.Next() { rule := new(ngmodels.AlertRule) @@ -485,6 +486,17 @@ func (st DBstore) GetAlertRulesForScheduling(ctx context.Context, query *ngmodel st.Logger.Error("Invalid rule found in DB store, ignoring it", "func", "GetAlertRulesForScheduling", "error", err) continue } + // This was added to mitigate the high load that could be created by loki range queries. + // In previous versions of Grafana, Loki datasources would default to range queries + // instead of instant queries, sometimes creating unnecessary load. This is only + // done for Grafana Cloud. + if lokiRangeToInstantEnabled && canBeInstant(rule) { + if err := migrateToInstant(rule); err != nil { + st.Logger.Error("Could not migrate rule from range to instant query", "rule", rule.UID, "err", err) + } else { + st.Logger.Info("Migrated rule from range to instant query", "rule", rule.UID) + } + } rules = append(rules, rule) } diff --git a/pkg/services/ngalert/store/alert_rule_test.go b/pkg/services/ngalert/store/alert_rule_test.go index 7bc13cafd57..5bb1865d8de 100644 --- a/pkg/services/ngalert/store/alert_rule_test.go +++ b/pkg/services/ngalert/store/alert_rule_test.go @@ -9,6 +9,7 @@ import ( "github.com/grafana/grafana/pkg/bus" "github.com/grafana/grafana/pkg/infra/tracing" + "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/services/folder" "github.com/grafana/grafana/pkg/services/folder/folderimpl" "github.com/grafana/grafana/pkg/services/org" @@ -92,7 +93,8 @@ func TestIntegration_GetAlertRulesForScheduling(t *testing.T) { Cfg: setting.UnifiedAlertingSettings{ BaseInterval: time.Duration(rand.Int63n(100)) * time.Second, }, - FolderService: setupFolderService(t, sqlStore, cfg), + FolderService: setupFolderService(t, sqlStore, cfg), + FeatureToggles: featuremgmt.WithFeatures(), } rule1 := createRule(t, store) diff --git a/pkg/services/ngalert/store/loki_range_to_instant.go b/pkg/services/ngalert/store/loki_range_to_instant.go new file mode 100644 index 00000000000..4524ed066bc --- /dev/null +++ b/pkg/services/ngalert/store/loki_range_to_instant.go @@ -0,0 +1,63 @@ +package store + +import ( + "encoding/json" + + "github.com/grafana/grafana/pkg/services/ngalert/models" +) + +const ( + grafanaCloudLogs = "grafanacloud-logs" + grafanaCloudUsageInsights = "grafanacloud-usage-insights" + grafanaCloudStateHistory = "grafanacloud-loki-alert-state-history" +) + +func canBeInstant(r *models.AlertRule) bool { + if len(r.Data) < 2 { + return false + } + // First query part should be range query. + if r.Data[0].QueryType != "range" { + return false + } + // First query part should go to cloud logs or insights. + if r.Data[0].DatasourceUID != grafanaCloudLogs && + r.Data[0].DatasourceUID != grafanaCloudUsageInsights && + r.Data[0].DatasourceUID != grafanaCloudStateHistory { + return false + } + // Second query part should be and expression, '-100' is the legacy way to define it. + if r.Data[1].DatasourceUID != "__expr__" && r.Data[1].DatasourceUID != "-100" { + return false + } + exprRaw := make(map[string]interface{}) + if err := json.Unmarshal(r.Data[1].Model, &exprRaw); err != nil { + return false + } + // Second query part should be "last()" + if val, ok := exprRaw["reducer"].(string); !ok || val != "last" { + return false + } + // Second query part should use first query part as expression. + if ref, ok := exprRaw["expression"].(string); !ok || ref != r.Data[0].RefID { + return false + } + return true +} + +// migrateToInstant will move a range-query to an instant query. This should only +// be used for loki. +func migrateToInstant(r *models.AlertRule) error { + modelRaw := make(map[string]interface{}) + if err := json.Unmarshal(r.Data[0].Model, &modelRaw); err != nil { + return err + } + modelRaw["queryType"] = "instant" + model, err := json.Marshal(modelRaw) + if err != nil { + return err + } + r.Data[0].Model = model + r.Data[0].QueryType = "instant" + return nil +} diff --git a/pkg/services/ngalert/store/loki_range_to_instant_test.go b/pkg/services/ngalert/store/loki_range_to_instant_test.go new file mode 100644 index 00000000000..46aa34133de --- /dev/null +++ b/pkg/services/ngalert/store/loki_range_to_instant_test.go @@ -0,0 +1,182 @@ +package store + +import ( + "encoding/json" + "testing" + + "github.com/grafana/grafana/pkg/services/ngalert/models" + "github.com/stretchr/testify/require" +) + +func TestCanBeInstant(t *testing.T) { + tcs := []struct { + name string + expected bool + rule *models.AlertRule + }{ + { + name: "valid rule that can be migrated from range to instant", + expected: true, + rule: createMigrateableLokiRule(t), + }, + { + name: "invalid rule where the data array is too short to be migrateable", + expected: false, + rule: createMigrateableLokiRule(t, func(r *models.AlertRule) { + r.Data = []models.AlertQuery{r.Data[0]} + }), + }, + { + name: "invalid rule that is not a range query", + expected: false, + rule: createMigrateableLokiRule(t, func(r *models.AlertRule) { + r.Data[0].QueryType = "something-else" + }), + }, + { + name: "invalid rule that does not use a cloud datasource", + expected: false, + rule: createMigrateableLokiRule(t, func(r *models.AlertRule) { + r.Data[0].DatasourceUID = "something-else" + }), + }, + { + name: "invalid rule that has no aggregation as second item", + expected: false, + rule: createMigrateableLokiRule(t, func(r *models.AlertRule) { + r.Data[1].DatasourceUID = "something-else" + }), + }, + { + name: "invalid rule that has not last() as aggregation", + expected: false, + rule: createMigrateableLokiRule(t, func(r *models.AlertRule) { + raw := make(map[string]interface{}) + err := json.Unmarshal(r.Data[1].Model, &raw) + require.NoError(t, err) + raw["reducer"] = "avg" + r.Data[1].Model, err = json.Marshal(raw) + require.NoError(t, err) + }), + }, + { + name: "invalid rule that has not last() pointing to range query", + expected: false, + rule: createMigrateableLokiRule(t, func(r *models.AlertRule) { + raw := make(map[string]interface{}) + err := json.Unmarshal(r.Data[1].Model, &raw) + require.NoError(t, err) + raw["expression"] = "C" + r.Data[1].Model, err = json.Marshal(raw) + require.NoError(t, err) + }), + }, + } + for _, tc := range tcs { + t.Run(tc.name, func(t *testing.T) { + require.Equal(t, tc.expected, canBeInstant(tc.rule)) + }) + } +} + +func TestMigrateLokiQueryToInstant(t *testing.T) { + original := createMigrateableLokiRule(t) + mirgrated := createMigrateableLokiRule(t, func(r *models.AlertRule) { + r.Data[0].QueryType = "instant" + r.Data[0].Model = []byte(`{ + "datasource": { + "type": "loki", + "uid": "grafanacloud-logs" + }, + "editorMode": "code", + "expr": "1", + "hide": false, + "intervalMs": 1000, + "maxDataPoints": 43200, + "queryType": "instant", + "refId": "A" + }`) + }) + + require.True(t, canBeInstant(original)) + require.NoError(t, migrateToInstant(original)) + + require.Equal(t, mirgrated.Data[0].QueryType, original.Data[0].QueryType) + + originalModel := make(map[string]interface{}) + require.NoError(t, json.Unmarshal(original.Data[0].Model, &originalModel)) + migratedModel := make(map[string]interface{}) + require.NoError(t, json.Unmarshal(mirgrated.Data[0].Model, &migratedModel)) + + require.Equal(t, migratedModel, originalModel) + + require.False(t, canBeInstant(original)) +} + +func createMigrateableLokiRule(t *testing.T, muts ...func(*models.AlertRule)) *models.AlertRule { + t.Helper() + r := &models.AlertRule{ + Data: []models.AlertQuery{ + { + RefID: "A", + QueryType: "range", + DatasourceUID: grafanaCloudLogs, + Model: []byte(`{ + "datasource": { + "type": "loki", + "uid": "grafanacloud-logs" + }, + "editorMode": "code", + "expr": "1", + "hide": false, + "intervalMs": 1000, + "maxDataPoints": 43200, + "queryType": "range", + "refId": "A" + }`), + }, + { + RefID: "B", + DatasourceUID: "__expr__", + Model: []byte(`{ + "conditions": [ + { + "evaluator": { + "params": [], + "type": "gt" + }, + "operator": { + "type": "and" + }, + "query": { + "params": [ + "B" + ] + }, + "reducer": { + "params": [], + "type": "last" + }, + "type": "query" + } + ], + "datasource": { + "type": "__expr__", + "uid": "__expr__" + }, + "expression": "A", + "hide": false, + "intervalMs": 1000, + "maxDataPoints": 43200, + "reducer": "last", + "refId": "B", + "type": "reduce" + }`), + }, + }, + } + for _, m := range muts { + m(r) + } + return r +}