diff --git a/pkg/tsdb/prometheus/buffered/time_series_query.go b/pkg/tsdb/prometheus/buffered/time_series_query.go index 9f5357d73b5..716aea00240 100644 --- a/pkg/tsdb/prometheus/buffered/time_series_query.go +++ b/pkg/tsdb/prometheus/buffered/time_series_query.go @@ -3,6 +3,7 @@ package buffered import ( "context" "encoding/json" + "errors" "fmt" "math" "net/http" @@ -64,6 +65,11 @@ type Buffered struct { TimeInterval string } +type bufferedResponse struct { + Response interface{} + Warnings apiv1.Warnings +} + // New creates and object capable of executing and parsing a Prometheus queries. It's "buffered" because there is // another implementation capable of streaming parse the response. func New(roundTripper http.RoundTripper, tracer tracing.Tracer, settings backend.DataSourceInstanceSettings, plog log.Logger) (*Buffered, error) { @@ -131,7 +137,7 @@ func (b *Buffered) runQueries(ctx context.Context, queries []*PrometheusQuery) ( }) defer endSpan() - response := make(map[TimeSeriesQueryType]interface{}) + response := make(map[TimeSeriesQueryType]bufferedResponse) timeRange := apiv1.Range{ Step: query.Step, @@ -141,23 +147,42 @@ func (b *Buffered) runQueries(ctx context.Context, queries []*PrometheusQuery) ( } if query.RangeQuery { - rangeResponse, _, err := b.client.QueryRange(ctx, query.Expr, timeRange) + rangeResponse, warnings, err := b.client.QueryRange(ctx, query.Expr, timeRange) if err != nil { - b.log.Error("Range query failed", "query", query.Expr, "err", err) - result.Responses[query.RefId] = backend.DataResponse{Error: err} + var promErr *apiv1.Error + if errors.As(err, &promErr) { + b.log.Error("Range query failed", "query", query.Expr, "error", err, "detail", promErr.Detail) + result.Responses[query.RefId] = backend.DataResponse{Error: fmt.Errorf("%w: details: %s", err, promErr.Detail)} + } else { + b.log.Error("Range query failed", "query", query.Expr, "err", err) + result.Responses[query.RefId] = backend.DataResponse{Error: err} + } + continue } - response[RangeQueryType] = rangeResponse + response[RangeQueryType] = bufferedResponse{ + Response: rangeResponse, + Warnings: warnings, + } } if query.InstantQuery { - instantResponse, _, err := b.client.Query(ctx, query.Expr, query.End) + instantResponse, warnings, err := b.client.Query(ctx, query.Expr, query.End) if err != nil { - b.log.Error("Instant query failed", "query", query.Expr, "err", err) - result.Responses[query.RefId] = backend.DataResponse{Error: err} + var promErr *apiv1.Error + if errors.As(err, &promErr) { + b.log.Error("Instant query failed", "query", query.Expr, "error", err, "detail", promErr.Detail) + result.Responses[query.RefId] = backend.DataResponse{Error: fmt.Errorf("%w: details: %s", err, promErr.Detail)} + } else { + b.log.Error("Instant query failed", "query", query.Expr, "err", err) + result.Responses[query.RefId] = backend.DataResponse{Error: err} + } continue } - response[InstantQueryType] = instantResponse + response[InstantQueryType] = bufferedResponse{ + Response: instantResponse, + Warnings: warnings, + } } // This is a special case @@ -167,7 +192,10 @@ func (b *Buffered) runQueries(ctx context.Context, queries []*PrometheusQuery) ( if err != nil { b.log.Error("Exemplar query failed", "query", query.Expr, "err", err) } else { - response[ExemplarQueryType] = exemplarResponse + response[ExemplarQueryType] = bufferedResponse{ + Response: exemplarResponse, + Warnings: nil, + } } } @@ -226,7 +254,7 @@ func (b *Buffered) parseTimeSeriesQuery(req *backend.QueryDataRequest) ([]*Prome if err != nil { return nil, fmt.Errorf("error unmarshaling query model: %v", err) } - //Final interval value + // Final interval value interval, err := calculatePrometheusInterval(model, b.TimeInterval, query, b.intervalCalculator) if err != nil { return nil, fmt.Errorf("error calculating interval: %v", err) @@ -263,17 +291,17 @@ func (b *Buffered) parseTimeSeriesQuery(req *backend.QueryDataRequest) ([]*Prome return qs, nil } -func parseTimeSeriesResponse(value map[TimeSeriesQueryType]interface{}, query *PrometheusQuery) (data.Frames, error) { +func parseTimeSeriesResponse(value map[TimeSeriesQueryType]bufferedResponse, query *PrometheusQuery) (data.Frames, error) { var ( frames = data.Frames{} nextFrames = data.Frames{} ) - for _, value := range value { + for _, val := range value { // Zero out the slice to prevent data corruption. nextFrames = nextFrames[:0] - switch v := value.(type) { + switch v := val.Response.(type) { case model.Matrix: nextFrames = matrixToDataFrames(v, query, nextFrames) case model.Vector: @@ -286,16 +314,39 @@ func parseTimeSeriesResponse(value map[TimeSeriesQueryType]interface{}, query *P return nil, fmt.Errorf("unexpected result type: %s query: %s", v, query.Expr) } + if len(val.Warnings) > 0 { + for _, frame := range nextFrames { + if frame.Meta == nil { + frame.Meta = &data.FrameMeta{} + } + frame.Meta.Notices = readWarnings(val.Warnings) + } + } + frames = append(frames, nextFrames...) } return frames, nil } +func readWarnings(warnings apiv1.Warnings) []data.Notice { + notices := []data.Notice{} + + for _, w := range warnings { + notice := data.Notice{ + Severity: data.NoticeSeverityWarning, + Text: w, + } + notices = append(notices, notice) + } + + return notices +} + func calculatePrometheusInterval(model *QueryModel, timeInterval string, query backend.DataQuery, intervalCalculator intervalv2.Calculator) (time.Duration, error) { queryInterval := model.Interval - //If we are using variable for interval/step, we will replace it with calculated interval + // If we are using variable for interval/step, we will replace it with calculated interval if isVariableInterval(queryInterval) { queryInterval = "" } @@ -654,7 +705,7 @@ func isVariableInterval(interval string) bool { if interval == varInterval || interval == varIntervalMs || interval == varRateInterval { return true } - //Repetitive code, we should have functionality to unify these + // Repetitive code, we should have functionality to unify these if interval == varIntervalAlt || interval == varIntervalMsAlt || interval == varRateIntervalAlt { return true } diff --git a/pkg/tsdb/prometheus/buffered/time_series_query_test.go b/pkg/tsdb/prometheus/buffered/time_series_query_test.go index 5b27c2051f9..9f3b9279188 100644 --- a/pkg/tsdb/prometheus/buffered/time_series_query_test.go +++ b/pkg/tsdb/prometheus/buffered/time_series_query_test.go @@ -598,7 +598,7 @@ func TestPrometheus_timeSeriesQuery_parseTimeSeriesQuery(t *testing.T) { func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { t.Run("exemplars response should be sampled and parsed normally", func(t *testing.T) { - value := make(map[TimeSeriesQueryType]interface{}) + value := make(map[TimeSeriesQueryType]bufferedResponse) exemplars := []apiv1.ExemplarQueryResult{ { SeriesLabels: p.LabelSet{ @@ -631,7 +631,10 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { }, } - value[ExemplarQueryType] = exemplars + value[ExemplarQueryType] = bufferedResponse{ + Response: exemplars, + Warnings: nil, + } query := &PrometheusQuery{ LegendFormat: "legend {{app}}", } @@ -652,7 +655,7 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { }) t.Run("exemplars response with inconsistent labels should marshal json ok", func(t *testing.T) { - value := make(map[TimeSeriesQueryType]interface{}) + value := make(map[TimeSeriesQueryType]bufferedResponse) exemplars := []apiv1.ExemplarQueryResult{ { SeriesLabels: p.LabelSet{ @@ -685,7 +688,10 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { }, } - value[ExemplarQueryType] = exemplars + value[ExemplarQueryType] = bufferedResponse{ + Response: exemplars, + Warnings: nil, + } query := &PrometheusQuery{ LegendFormat: "legend {{app}}", } @@ -723,12 +729,15 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { {Value: 4, Timestamp: 4000}, {Value: 5, Timestamp: 5000}, } - value := make(map[TimeSeriesQueryType]interface{}) - value[RangeQueryType] = p.Matrix{ - &p.SampleStream{ - Metric: p.Metric{"app": "Application", "tag2": "tag2"}, - Values: values, + value := make(map[TimeSeriesQueryType]bufferedResponse) + value[RangeQueryType] = bufferedResponse{ + Response: p.Matrix{ + &p.SampleStream{ + Metric: p.Metric{"app": "Application", "tag2": "tag2"}, + Values: values, + }, }, + Warnings: nil, } query := &PrometheusQuery{ LegendFormat: "legend {{app}}", @@ -760,12 +769,15 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { {Value: 1, Timestamp: 1000}, {Value: 4, Timestamp: 4000}, } - value := make(map[TimeSeriesQueryType]interface{}) - value[RangeQueryType] = p.Matrix{ - &p.SampleStream{ - Metric: p.Metric{"app": "Application", "tag2": "tag2"}, - Values: values, + value := make(map[TimeSeriesQueryType]bufferedResponse) + value[RangeQueryType] = bufferedResponse{ + Response: p.Matrix{ + &p.SampleStream{ + Metric: p.Metric{"app": "Application", "tag2": "tag2"}, + Values: values, + }, }, + Warnings: nil, } query := &PrometheusQuery{ LegendFormat: "", @@ -791,12 +803,15 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { {Value: 1, Timestamp: 1000}, {Value: 4, Timestamp: 4000}, } - value := make(map[TimeSeriesQueryType]interface{}) - value[RangeQueryType] = p.Matrix{ - &p.SampleStream{ - Metric: p.Metric{"app": "Application", "tag2": "tag2"}, - Values: values, + value := make(map[TimeSeriesQueryType]bufferedResponse) + value[RangeQueryType] = bufferedResponse{ + Response: p.Matrix{ + &p.SampleStream{ + Metric: p.Metric{"app": "Application", "tag2": "tag2"}, + Values: values, + }, }, + Warnings: nil, } query := &PrometheusQuery{ LegendFormat: "", @@ -820,14 +835,17 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { }) t.Run("matrix response with NaN value should be changed to null", func(t *testing.T) { - value := make(map[TimeSeriesQueryType]interface{}) - value[RangeQueryType] = p.Matrix{ - &p.SampleStream{ - Metric: p.Metric{"app": "Application"}, - Values: []p.SamplePair{ - {Value: p.SampleValue(math.NaN()), Timestamp: 1000}, + value := make(map[TimeSeriesQueryType]bufferedResponse) + value[RangeQueryType] = bufferedResponse{ + Response: p.Matrix{ + &p.SampleStream{ + Metric: p.Metric{"app": "Application"}, + Values: []p.SamplePair{ + {Value: p.SampleValue(math.NaN()), Timestamp: 1000}, + }, }, }, + Warnings: nil, } query := &PrometheusQuery{ LegendFormat: "", @@ -844,13 +862,16 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { }) t.Run("vector response should be parsed normally", func(t *testing.T) { - value := make(map[TimeSeriesQueryType]interface{}) - value[RangeQueryType] = p.Vector{ - &p.Sample{ - Metric: p.Metric{"app": "Application", "tag2": "tag2"}, - Value: 1, - Timestamp: 123, + value := make(map[TimeSeriesQueryType]bufferedResponse) + value[RangeQueryType] = bufferedResponse{ + Response: p.Vector{ + &p.Sample{ + Metric: p.Metric{"app": "Application", "tag2": "tag2"}, + Value: 1, + Timestamp: 123, + }, }, + Warnings: nil, } query := &PrometheusQuery{ LegendFormat: "legend {{app}}", @@ -876,10 +897,13 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { }) t.Run("scalar response should be parsed normally", func(t *testing.T) { - value := make(map[TimeSeriesQueryType]interface{}) - value[RangeQueryType] = &p.Scalar{ - Value: 1, - Timestamp: 123, + value := make(map[TimeSeriesQueryType]bufferedResponse) + value[RangeQueryType] = bufferedResponse{ + Response: &p.Scalar{ + Value: 1, + Timestamp: 123, + }, + Warnings: nil, } query := &PrometheusQuery{} @@ -899,6 +923,33 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) { require.Equal(t, "UTC", testValue.(time.Time).Location().String()) require.Equal(t, int64(123), testValue.(time.Time).UnixMilli()) }) + + t.Run("warnings, if there is any, should be added to each frame", + func(t *testing.T) { + value := make(map[TimeSeriesQueryType]bufferedResponse) + value[RangeQueryType] = bufferedResponse{ + Response: &p.Scalar{ + Value: 1, + Timestamp: 123, + }, + Warnings: []string{"warning1", "warning2"}, + } + + query := &PrometheusQuery{} + res, err := parseTimeSeriesResponse(value, query) + require.NoError(t, err) + + require.Len(t, res, 1) + require.Equal(t, res[0].Name, "1") + require.Len(t, res[0].Fields, 2) + require.Len(t, res[0].Fields[0].Labels, 0) + require.Equal(t, res[0].Fields[0].Name, "Time") + require.Equal(t, res[0].Fields[1].Name, "Value") + require.Equal(t, res[0].Fields[1].Config.DisplayNameFromDS, "1") + + require.Equal(t, res[0].Meta.Notices[0].Text, "warning1") + require.Equal(t, res[0].Meta.Notices[1].Text, "warning2") + }) } func queryContext(json string, timeRange backend.TimeRange) *backend.QueryDataRequest {