Prometheus: Various buffered and streaming parsing fixes (#55941)

This commit is contained in:
Todd Treece
2022-10-03 10:26:54 -04:00
committed by GitHub
parent 8984507291
commit 1c61c81dde
14 changed files with 1452 additions and 461 deletions
+50 -26
View File
@@ -4,6 +4,7 @@ import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"os"
@@ -34,6 +35,7 @@ func TestMatrixResponses(t *testing.T) {
{name: "parse a simple matrix response with value missing steps", filepath: "range_missing"},
{name: "parse a response with Infinity", filepath: "range_infinity"},
{name: "parse a response with NaN", filepath: "range_nan"},
{name: "parse a response with legendFormat __auto", filepath: "range_auto"},
}
for _, test := range tt {
@@ -94,39 +96,28 @@ func makeMockedApi(responseBytes []byte) (apiv1.API, error) {
// struct here, because it has `time.time` and `time.duration` fields that
// cannot be unmarshalled from JSON automatically.
type storedPrometheusQuery struct {
RefId string
RangeQuery bool
Start int64
End int64
Step int64
Expr string
RefId string
RangeQuery bool
Start int64
End int64
Step int64
Expr string
LegendFormat string
}
func loadStoredPrometheusQuery(fileName string) (PrometheusQuery, error) {
func loadStoredPrometheusQuery(fileName string) (storedPrometheusQuery, error) {
//nolint:gosec
bytes, err := os.ReadFile(fileName)
if err != nil {
return PrometheusQuery{}, err
return storedPrometheusQuery{}, err
}
var query storedPrometheusQuery
err = json.Unmarshal(bytes, &query)
if err != nil {
return PrometheusQuery{}, err
}
return PrometheusQuery{
RefId: query.RefId,
RangeQuery: query.RangeQuery,
Start: time.Unix(query.Start, 0),
End: time.Unix(query.End, 0),
Step: time.Second * time.Duration(query.Step),
Expr: query.Expr,
}, nil
var sq storedPrometheusQuery
err = json.Unmarshal(bytes, &sq)
return sq, err
}
func runQuery(response []byte, query PrometheusQuery) (*backend.QueryDataResponse, error) {
func runQuery(response []byte, sq storedPrometheusQuery) (*backend.QueryDataResponse, error) {
api, err := makeMockedApi(response)
if err != nil {
return nil, err
@@ -134,14 +125,47 @@ func runQuery(response []byte, query PrometheusQuery) (*backend.QueryDataRespons
tracer := tracing.InitializeTracerForTest()
s := Buffered{
qm := QueryModel{
RangeQuery: sq.RangeQuery,
Expr: sq.Expr,
Interval: fmt.Sprintf("%ds", sq.Step),
IntervalMS: sq.Step * 1000,
LegendFormat: sq.LegendFormat,
}
b := Buffered{
intervalCalculator: intervalv2.NewCalculator(),
tracer: tracer,
TimeInterval: "15s",
log: &fakeLogger{},
client: api,
}
return s.runQueries(context.Background(), []*PrometheusQuery{&query})
data, err := json.Marshal(&qm)
if err != nil {
return nil, err
}
req := &backend.QueryDataRequest{
Queries: []backend.DataQuery{
{
TimeRange: backend.TimeRange{
From: time.Unix(sq.Start, 0),
To: time.Unix(sq.End, 0),
},
RefID: sq.RefId,
Interval: time.Second * time.Duration(sq.Step),
JSON: json.RawMessage(data),
},
},
}
queries, err := b.parseTimeSeriesQuery(req)
if err != nil {
return nil, err
}
return b.runQueries(context.Background(), queries)
}
type fakeLogger struct {
@@ -379,11 +379,7 @@ func matrixToDataFrames(matrix model.Matrix, query *PrometheusQuery, frames data
for i, k := range v.Values {
timeField.Set(i, k.Timestamp.Time().UTC())
value := float64(k.Value)
if !math.IsNaN(value) {
valueField.Set(i, value)
}
valueField.Set(i, float64(k.Value))
}
name := formatLegend(v.Metric, query)
@@ -836,7 +836,7 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
require.NoError(t, err)
require.Equal(t, "Value", res[0].Fields[1].Name)
require.Equal(t, float64(0), res[0].Fields[1].At(0))
require.True(t, math.IsNaN(res[0].Fields[1].At(0).(float64)))
})
t.Run("vector response should be parsed normally", func(t *testing.T) {