From b7dcfcedcb436f1a3bddb19278506de4198c32b7 Mon Sep 17 00:00:00 2001 From: Steve Simpson Date: Thu, 6 Mar 2025 14:09:17 +0100 Subject: [PATCH] Alerting: Extend recording rule definitions/interfaces with data source. (#101678) Extend the recording rule definition to include the target data source, allowing configuration of where the output of the recording rule is written to. Also extends the relevant interfaces in preparation for the next set of changes. --- pkg/services/ngalert/api/compat.go | 22 ++++++++++++++----- .../api/tooling/definitions/cortex-ruler.go | 4 ++++ .../definitions/provisioning_alert_rules.go | 5 +++-- pkg/services/ngalert/models/alert_rule.go | 3 +++ .../ngalert/schedule/recording_rule.go | 2 +- pkg/services/ngalert/schedule/schedule.go | 1 + pkg/services/ngalert/writer/fake.go | 13 +++++++++++ pkg/services/ngalert/writer/noop.go | 4 ++++ pkg/services/ngalert/writer/prom.go | 12 ++++++++++ 9 files changed, 57 insertions(+), 9 deletions(-) diff --git a/pkg/services/ngalert/api/compat.go b/pkg/services/ngalert/api/compat.go index ba78f5f613d..6353e1db458 100644 --- a/pkg/services/ngalert/api/compat.go +++ b/pkg/services/ngalert/api/compat.go @@ -494,13 +494,21 @@ func NotificationSettingsFromAlertRuleNotificationSettings(ns *definitions.Alert } } +func pointerOmitEmpty(s string) *string { + if s == "" { + return nil + } + return &s +} + func AlertRuleRecordExportFromRecord(r *models.Record) *definitions.AlertRuleRecordExport { if r == nil { return nil } return &definitions.AlertRuleRecordExport{ - Metric: r.Metric, - From: r.From, + Metric: r.Metric, + From: r.From, + TargetDatasourceUID: pointerOmitEmpty(r.TargetDatasourceUID), } } @@ -509,8 +517,9 @@ func ModelRecordFromApiRecord(r *definitions.Record) *models.Record { return nil } return &models.Record{ - Metric: r.Metric, - From: r.From, + Metric: r.Metric, + From: r.From, + TargetDatasourceUID: r.TargetDatasourceUID, } } @@ -519,8 +528,9 @@ func ApiRecordFromModelRecord(r *models.Record) *definitions.Record { return nil } return &definitions.Record{ - Metric: r.Metric, - From: r.From, + Metric: r.Metric, + From: r.From, + TargetDatasourceUID: r.TargetDatasourceUID, } } diff --git a/pkg/services/ngalert/api/tooling/definitions/cortex-ruler.go b/pkg/services/ngalert/api/tooling/definitions/cortex-ruler.go index 777603dbb82..fc8c1c9c4fe 100644 --- a/pkg/services/ngalert/api/tooling/definitions/cortex-ruler.go +++ b/pkg/services/ngalert/api/tooling/definitions/cortex-ruler.go @@ -533,6 +533,10 @@ type Record struct { // required: true // example: A From string `json:"from" yaml:"from"` + // Which data source should be used to write the output of the recording rule, specified by UID. + // required: false + // example: my-prom + TargetDatasourceUID string `json:"target_datasource_uid,omitempty" yaml:"target_datasource_uid,omitempty"` } // swagger:model diff --git a/pkg/services/ngalert/api/tooling/definitions/provisioning_alert_rules.go b/pkg/services/ngalert/api/tooling/definitions/provisioning_alert_rules.go index d38cb17e2e7..43c14790547 100644 --- a/pkg/services/ngalert/api/tooling/definitions/provisioning_alert_rules.go +++ b/pkg/services/ngalert/api/tooling/definitions/provisioning_alert_rules.go @@ -308,6 +308,7 @@ type AlertRuleNotificationSettingsExport struct { // Record is the provisioned export of models.Record. type AlertRuleRecordExport struct { - Metric string `json:"metric" yaml:"metric" hcl:"metric"` - From string `json:"from" yaml:"from" hcl:"from"` + Metric string `json:"metric" yaml:"metric" hcl:"metric"` + From string `json:"from" yaml:"from" hcl:"from"` + TargetDatasourceUID *string `json:"targetDatasourceUid,omitempty" yaml:"targetDatasourceUid,omitempty" hcl:"target_datasource_uid,optional"` } diff --git a/pkg/services/ngalert/models/alert_rule.go b/pkg/services/ngalert/models/alert_rule.go index 2536c900b1c..69219b8b09f 100644 --- a/pkg/services/ngalert/models/alert_rule.go +++ b/pkg/services/ngalert/models/alert_rule.go @@ -1010,6 +1010,8 @@ type Record struct { Metric string // From contains a query RefID, indicating which expression node is the output of the recording rule. From string + // TargetDatasourceUID is the data source to write the result of the recording rule. + TargetDatasourceUID string } func (r *Record) Fingerprint() data.Fingerprint { @@ -1024,6 +1026,7 @@ func (r *Record) Fingerprint() data.Fingerprint { writeString(r.Metric) writeString(r.From) + writeString(r.TargetDatasourceUID) return data.Fingerprint(h.Sum64()) } diff --git a/pkg/services/ngalert/schedule/recording_rule.go b/pkg/services/ngalert/schedule/recording_rule.go index 6c3f733a5ca..d5d42763d07 100644 --- a/pkg/services/ngalert/schedule/recording_rule.go +++ b/pkg/services/ngalert/schedule/recording_rule.go @@ -269,7 +269,7 @@ func (r *recordingRule) tryEvaluation(ctx context.Context, ev *Evaluation, logge } writeStart := r.clock.Now() - err = r.writer.Write(ctx, ev.rule.Record.Metric, ev.scheduledAt, frames, ev.rule.OrgID, ev.rule.Labels) + err = r.writer.WriteDatasource(ctx, ev.rule.Record.TargetDatasourceUID, ev.rule.Record.Metric, ev.scheduledAt, frames, ev.rule.OrgID, ev.rule.Labels) writeDur := r.clock.Now().Sub(writeStart) if err != nil { diff --git a/pkg/services/ngalert/schedule/schedule.go b/pkg/services/ngalert/schedule/schedule.go index 0416fa2fb7e..e3515ef14b4 100644 --- a/pkg/services/ngalert/schedule/schedule.go +++ b/pkg/services/ngalert/schedule/schedule.go @@ -50,6 +50,7 @@ type RulesStore interface { type RecordingWriter interface { Write(ctx context.Context, name string, t time.Time, frames data.Frames, orgID int64, extraLabels map[string]string) error + WriteDatasource(ctx context.Context, dsUID string, name string, t time.Time, frames data.Frames, orgID int64, extraLabels map[string]string) error } // AlertRuleStopReasonProvider is an interface for determining the reason why an alert rule was stopped. diff --git a/pkg/services/ngalert/writer/fake.go b/pkg/services/ngalert/writer/fake.go index bd688d65f7a..17a597bd43a 100644 --- a/pkg/services/ngalert/writer/fake.go +++ b/pkg/services/ngalert/writer/fake.go @@ -2,6 +2,7 @@ package writer import ( "context" + "errors" "time" "github.com/grafana/grafana-plugin-sdk-go/data" @@ -18,3 +19,15 @@ func (w FakeWriter) Write(ctx context.Context, name string, t time.Time, frames return w.WriteFunc(ctx, name, t, frames, orgID, extraLabels) } + +func (w FakeWriter) WriteDatasource(ctx context.Context, dsUID string, name string, t time.Time, frames data.Frames, orgID int64, extraLabels map[string]string) error { + if w.WriteFunc == nil { + return nil + } + + if dsUID != "" { + return errors.New("expected empty data source uid") + } + + return w.WriteFunc(ctx, name, t, frames, orgID, extraLabels) +} diff --git a/pkg/services/ngalert/writer/noop.go b/pkg/services/ngalert/writer/noop.go index 066f412fa1b..a2d32f13334 100644 --- a/pkg/services/ngalert/writer/noop.go +++ b/pkg/services/ngalert/writer/noop.go @@ -12,3 +12,7 @@ type NoopWriter struct{} func (w NoopWriter) Write(ctx context.Context, name string, t time.Time, frames data.Frames, orgID int64, extraLabels map[string]string) error { return nil } + +func (w NoopWriter) WriteDatasource(ctx context.Context, dsUID string, name string, t time.Time, frames data.Frames, orgID int64, extraLabels map[string]string) error { + return nil +} diff --git a/pkg/services/ngalert/writer/prom.go b/pkg/services/ngalert/writer/prom.go index ff735fecc26..fb759b0494d 100644 --- a/pkg/services/ngalert/writer/prom.go +++ b/pkg/services/ngalert/writer/prom.go @@ -192,6 +192,18 @@ func createAuthOpts(username, password string) *httpclient.BasicAuthOptions { } } +// Write writes the given frames to the Prometheus remote write endpoint. +func (w PrometheusWriter) WriteDatasource(ctx context.Context, dsUID string, name string, t time.Time, frames data.Frames, orgID int64, extraLabels map[string]string) error { + l := w.logger.FromContext(ctx) + + if dsUID != "" { + l.Error("Writing to specific data sources is not enabled", "org_id", orgID, "datasource_uid", dsUID) + return errors.New("writing to specific data sources is not enabled") + } + + return w.Write(ctx, name, t, frames, orgID, extraLabels) +} + // Write writes the given frames to the Prometheus remote write endpoint. func (w PrometheusWriter) Write(ctx context.Context, name string, t time.Time, frames data.Frames, orgID int64, extraLabels map[string]string) error { l := w.logger.FromContext(ctx)