InfluxDB: backend migration (run query in explore) (#43352)

* InfluxDB backend migration

* Multiple queries and more

* Added types

* Updated preferredVisualisationType

* Updated model parser test to include limit,slimit,orderByTime

* Added test for building query with limit, slimit

* Added test for building query with limit, slimit, orderByTime and puts them in the correct order

* Add test: Influxdb response parser should parse two responses with different refIDs

* Moved methods to responds parser

* Add test to ensure ExecutedQueryString is populated

* Move functions out of response parser class

* Test for getSelectedParams

* Merge cases

* Change to const

* Test get table columns correctly

* Removed unecessary fields

* Test get table rows correctly

* Removed getSeries function

* Added test for preferredVisualisationType

* Added test for executedQueryString

* Modified response parser

* Removed test

* Improvements

* Tests

* Review changes

* Feature flag rename and code gen
This commit is contained in:
Joey Tawadrous
2022-02-09 18:26:16 +00:00
committed by GitHub
parent 7ef43fb959
commit 10232c7857
14 changed files with 447 additions and 79 deletions
+6
View File
@@ -94,6 +94,12 @@ var (
Description: "Use azure authentication for prometheus datasource",
State: FeatureStateBeta,
},
{
Name: "influxdbBackendMigration",
Description: "Query InfluxDB InfluxQL without the proxy",
State: FeatureStateAlpha,
FrontendOnly: true,
},
{
Name: "newNavigation",
Description: "Try the next gen navigation model",
+4
View File
@@ -71,6 +71,10 @@ const (
// Use azure authentication for prometheus datasource
FlagPrometheusAzureAuth = "prometheus_azure_auth"
// FlagInfluxdbBackendMigration
// Query InfluxDB InfluxQL without the proxy
FlagInfluxdbBackendMigration = "influxdbBackendMigration"
// FlagNewNavigation
// Try the next gen navigation model
FlagNewNavigation = "newNavigation"
+19 -26
View File
@@ -94,24 +94,31 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
s.glog.Debug("Making a non-Flux type query")
// NOTE: the following path is currently only called from alerting queries
// In dashboards, the request runs through proxy and are managed in the frontend
var allRawQueries string
var queries []Query
query, err := s.getQuery(dsInfo, req)
if err != nil {
return &backend.QueryDataResponse{}, err
}
for _, reqQuery := range req.Queries {
query, err := s.queryParser.Parse(reqQuery)
if err != nil {
return &backend.QueryDataResponse{}, err
}
rawQuery, err := query.Build(req)
if err != nil {
return &backend.QueryDataResponse{}, err
rawQuery, err := query.Build(req)
if err != nil {
return &backend.QueryDataResponse{}, err
}
allRawQueries = allRawQueries + rawQuery + ";"
query.RefID = reqQuery.RefID
query.RawQuery = rawQuery
queries = append(queries, *query)
}
if setting.Env == setting.Dev {
s.glog.Debug("Influxdb query", "raw query", rawQuery)
s.glog.Debug("Influxdb query", "raw query", allRawQueries)
}
request, err := s.createRequest(ctx, dsInfo, rawQuery)
request, err := s.createRequest(ctx, dsInfo, allRawQueries)
if err != nil {
return &backend.QueryDataResponse{}, err
}
@@ -129,25 +136,11 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
return &backend.QueryDataResponse{}, fmt.Errorf("InfluxDB returned error status: %s", res.Status)
}
resp := s.responseParser.Parse(res.Body, query)
resp := s.responseParser.Parse(res.Body, queries)
return resp, nil
}
func (s *Service) getQuery(dsInfo *models.DatasourceInfo, query *backend.QueryDataRequest) (*Query, error) {
queryCount := len(query.Queries)
// The model supports multiple queries, but right now this is only used from
// alerting so we only needed to support batch executing 1 query at a time.
if queryCount != 1 {
return nil, fmt.Errorf("query request should contain exactly 1 query, it contains: %d", queryCount)
}
q := query.Queries[0]
return s.queryParser.Parse(q)
}
func (s *Service) createRequest(ctx context.Context, dsInfo *models.DatasourceInfo, query string) (*http.Request, error) {
u, err := url.Parse(dsInfo.URL)
if err != nil {
+6
View File
@@ -22,6 +22,9 @@ func (qp *InfluxdbQueryParser) Parse(query backend.DataQuery) (*Query, error) {
useRawQuery := model.Get("rawQuery").MustBool(false)
alias := model.Get("alias").MustString("")
tz := model.Get("tz").MustString("")
limit := model.Get("limit").MustString("")
slimit := model.Get("slimit").MustString("")
orderByTime := model.Get("orderByTime").MustString("")
measurement := model.Get("measurement").MustString("")
@@ -60,6 +63,9 @@ func (qp *InfluxdbQueryParser) Parse(query backend.DataQuery) (*Query, error) {
Alias: alias,
UseRawQuery: useRawQuery,
Tz: tz,
Limit: limit,
Slimit: slimit,
OrderByTime: orderByTime,
}, nil
}
+6
View File
@@ -36,6 +36,9 @@ func TestInfluxdbQueryParser_Parse(t *testing.T) {
],
"measurement": "logins.count",
"tz": "Europe/Paris",
"limit": "1",
"slimit": "1",
"orderByTime": "ASC",
"policy": "default",
"refId": "B",
"resultFormat": "time_series",
@@ -113,6 +116,9 @@ func TestInfluxdbQueryParser_Parse(t *testing.T) {
require.Len(t, res.Selects, 3)
require.Len(t, res.Tags, 2)
require.Equal(t, "Europe/Paris", res.Tz)
require.Equal(t, "1", res.Limit)
require.Equal(t, "1", res.Slimit)
require.Equal(t, "ASC", res.OrderByTime)
require.Equal(t, time.Second*20, res.Interval)
require.Equal(t, "series alias", res.Alias)
})
+4
View File
@@ -13,6 +13,10 @@ type Query struct {
Alias string
Interval time.Duration
Tz string
Limit string
Slimit string
OrderByTime string
RefID string
}
type Tag struct {
+27
View File
@@ -26,6 +26,9 @@ func (query *Query) Build(queryContext *backend.QueryDataRequest) (string, error
res += query.renderWhereClause()
res += query.renderTimeFilter(queryContext)
res += query.renderGroupBy(queryContext)
res += query.renderOrderByTime()
res += query.renderLimit()
res += query.renderSlimit()
res += query.renderTz()
}
@@ -151,6 +154,14 @@ func (query *Query) renderGroupBy(queryContext *backend.QueryDataRequest) string
return groupBy
}
func (query *Query) renderOrderByTime() string {
orderByTime := query.OrderByTime
if orderByTime == "" {
return ""
}
return fmt.Sprintf(" ORDER BY time %s", orderByTime)
}
func (query *Query) renderTz() string {
tz := query.Tz
if tz == "" {
@@ -159,6 +170,22 @@ func (query *Query) renderTz() string {
return fmt.Sprintf(" tz('%s')", tz)
}
func (query *Query) renderLimit() string {
limit := query.Limit
if limit == "" {
return ""
}
return fmt.Sprintf(" limit %s", limit)
}
func (query *Query) renderSlimit() string {
slimit := query.Slimit
if slimit == "" {
return ""
}
return fmt.Sprintf(" slimit %s", slimit)
}
func epochMStoInfluxTime(tr *backend.TimeRange) (string, string) {
from := tr.From.UnixNano() / int64(time.Millisecond)
to := tr.To.UnixNano() / int64(time.Millisecond)
+17
View File
@@ -66,6 +66,23 @@ func TestInfluxdbQueryBuilder(t *testing.T) {
require.Equal(t, rawQuery, `SELECT mean("value") FROM "cpu" WHERE time > 1596240000000ms and time < 1596240300000ms GROUP BY time(5s) tz('Europe/Paris')`)
})
t.Run("can build query with tz, limit, slimit, orderByTime and puts them in the correct order", func(t *testing.T) {
query := &Query{
Selects: []*Select{{*qp1, *qp2}},
Measurement: "cpu",
GroupBy: []*QueryPart{groupBy1},
Tz: "Europe/Paris",
Limit: "1",
Slimit: "1",
OrderByTime: "ASC",
Interval: time.Second * 5,
}
rawQuery, err := query.Build(queryContext)
require.NoError(t, err)
require.Equal(t, rawQuery, `SELECT mean("value") FROM "cpu" WHERE time > 1596240000000ms and time < 1596240300000ms GROUP BY time(5s) ORDER BY time ASC limit 1 slimit 1 tz('Europe/Paris')`)
})
t.Run("can build query with group bys", func(t *testing.T) {
query := &Query{
Selects: []*Select{{*qp1, *qp2}},
+19 -15
View File
@@ -19,32 +19,27 @@ var (
legendFormat = regexp.MustCompile(`\[\[([\@\/\w-]+)(\.[\@\/\w-]+)*\]\]*|\$(\s*([\@\w-]+?))*`)
)
func (rp *ResponseParser) Parse(buf io.ReadCloser, query *Query) *backend.QueryDataResponse {
func (rp *ResponseParser) Parse(buf io.ReadCloser, queries []Query) *backend.QueryDataResponse {
resp := backend.NewQueryDataResponse()
queryRes := backend.DataResponse{}
response, jsonErr := parseJSON(buf)
if jsonErr != nil {
queryRes.Error = jsonErr
resp.Responses["A"] = queryRes
resp.Responses["A"] = backend.DataResponse{Error: jsonErr}
return resp
}
if response.Error != "" {
queryRes.Error = fmt.Errorf(response.Error)
resp.Responses["A"] = queryRes
resp.Responses["A"] = backend.DataResponse{Error: fmt.Errorf(response.Error)}
return resp
}
frames := data.Frames{}
for _, result := range response.Results {
frames = append(frames, transformRows(result.Series, query)...)
for i, result := range response.Results {
if result.Error != "" {
queryRes.Error = fmt.Errorf(result.Error)
resp.Responses[queries[i].RefID] = backend.DataResponse{Error: fmt.Errorf(result.Error)}
} else {
resp.Responses[queries[i].RefID] = backend.DataResponse{Frames: transformRows(result.Series, queries[i])}
}
}
queryRes.Frames = frames
resp.Responses["A"] = queryRes
return resp
}
@@ -58,7 +53,7 @@ func parseJSON(buf io.ReadCloser) (Response, error) {
return response, err
}
func transformRows(rows []Row, query *Query) data.Frames {
func transformRows(rows []Row, query Query) data.Frames {
frames := data.Frames{}
for _, row := range rows {
for columnIndex, column := range row.Columns {
@@ -86,14 +81,23 @@ func transformRows(rows []Row, query *Query) data.Frames {
// set a nice name on the value-field
valueField.SetConfig(&data.FieldConfig{DisplayNameFromDS: name})
frames = append(frames, data.NewFrame(name, timeField, valueField))
frames = append(frames, newDataFrame(name, query.RawQuery, timeField, valueField))
}
}
return frames
}
func formatFrameName(row Row, column string, query *Query) string {
func newDataFrame(name string, queryString string, timeField *data.Field, valueField *data.Field) *data.Frame {
frame := data.NewFrame(name, timeField, valueField)
frame.Meta = &data.FrameMeta{
ExecutedQueryString: queryString,
}
return frame
}
func formatFrameName(row Row, column string, query Query) string {
if query.Alias == "" {
return buildFrameNameFromQuery(row, column)
}
+99 -26
View File
@@ -10,6 +10,7 @@ import (
"github.com/google/go-cmp/cmp"
"github.com/grafana/grafana-plugin-sdk-go/data"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/xorcare/pointer"
)
@@ -18,6 +19,14 @@ func prepare(text string) io.ReadCloser {
return ioutil.NopCloser(strings.NewReader(text))
}
func addQueryToQueries(query Query) []Query {
var queries []Query
query.RefID = "A"
query.RawQuery = "Test raw query"
queries = append(queries, query)
return queries
}
func TestInfluxdbResponseParser(t *testing.T) {
t.Run("Influxdb response parser should handle invalid JSON", func(t *testing.T) {
parser := &ResponseParser{}
@@ -26,7 +35,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
query := &Query{}
result := parser.Parse(prepare(response), query)
result := parser.Parse(prepare(response), addQueryToQueries(*query))
require.Nil(t, result.Responses["A"].Frames)
require.Error(t, result.Responses["A"].Error)
@@ -72,8 +81,9 @@ func TestInfluxdbResponseParser(t *testing.T) {
}),
newField,
)
testFrame.Meta = &data.FrameMeta{ExecutedQueryString: "Test raw query"}
result := parser.Parse(prepare(response), query)
result := parser.Parse(prepare(response), addQueryToQueries(*query))
frame := result.Responses["A"]
if diff := cmp.Diff(testFrame, frame.Frames[0], data.FrameTestCompareOptions()...); diff != "" {
@@ -81,6 +91,61 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
})
t.Run("Influxdb response parser should parse two responses with different refIDs", func(t *testing.T) {
parser := &ResponseParser{}
response := `
{
"results": [
{
"series": [{}]
},
{
"series": [{}]
}
]
}
`
query := &Query{}
var queries = addQueryToQueries(*query)
queryB := &Query{}
queryB.RefID = "B"
queries = append(queries, *queryB)
result := parser.Parse(prepare(response), queries)
assert.Len(t, result.Responses, 2)
assert.Contains(t, result.Responses, "A")
assert.Contains(t, result.Responses, "B")
assert.NotContains(t, result.Responses, "C")
})
t.Run("Influxdb response parser populates the RawQuery in the response meta ExecutedQueryString", func(t *testing.T) {
parser := &ResponseParser{}
response := `
{
"results": [
{
"series": [
{
"name": "cpu",
"columns": ["time","mean"]
}
]
}
]
}
`
query := &Query{}
query.RawQuery = "Test raw query"
result := parser.Parse(prepare(response), addQueryToQueries(*query))
frame := result.Responses["A"]
assert.Equal(t, frame.Frames[0].Meta.ExecutedQueryString, "Test raw query")
})
t.Run("Influxdb response parser with invalid value-format", func(t *testing.T) {
parser := &ResponseParser{}
@@ -119,8 +184,9 @@ func TestInfluxdbResponseParser(t *testing.T) {
}),
newField,
)
testFrame.Meta = &data.FrameMeta{ExecutedQueryString: "Test raw query"}
result := parser.Parse(prepare(response), query)
result := parser.Parse(prepare(response), addQueryToQueries(*query))
frame := result.Responses["A"]
if diff := cmp.Diff(testFrame, frame.Frames[0], data.FrameTestCompareOptions()...); diff != "" {
@@ -166,8 +232,9 @@ func TestInfluxdbResponseParser(t *testing.T) {
}),
newField,
)
testFrame.Meta = &data.FrameMeta{ExecutedQueryString: "Test raw query"}
result := parser.Parse(prepare(response), query)
result := parser.Parse(prepare(response), addQueryToQueries(*query))
frame := result.Responses["A"]
if diff := cmp.Diff(testFrame, frame.Frames[0], data.FrameTestCompareOptions()...); diff != "" {
@@ -217,7 +284,8 @@ func TestInfluxdbResponseParser(t *testing.T) {
}),
newField,
)
result := parser.Parse(prepare(response), query)
testFrame.Meta = &data.FrameMeta{ExecutedQueryString: "Test raw query"}
result := parser.Parse(prepare(response), addQueryToQueries(*query))
t.Run("should parse aliases", func(t *testing.T) {
frame := result.Responses["A"]
if diff := cmp.Diff(testFrame, frame.Frames[0], data.FrameTestCompareOptions()...); diff != "" {
@@ -225,7 +293,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias $m $measurement", Measurement: "10m"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name := "alias 10m 10m"
@@ -236,7 +304,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias $col", Measurement: "10m"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias mean"
testFrame.Name = name
@@ -256,7 +324,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias $tag_datacenter"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias America"
testFrame.Name = name
@@ -270,7 +338,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias $tag_datacenter/$tag_datacenter"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias America/America"
testFrame.Name = name
@@ -284,7 +352,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias [[col]]", Measurement: "10m"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias mean"
testFrame.Name = name
@@ -294,7 +362,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias $0 $1 $2 $3 $4"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias cpu upc $2 $3 $4"
testFrame.Name = name
@@ -304,7 +372,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias $1"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias upc"
testFrame.Name = name
@@ -314,7 +382,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias $5"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias $5"
testFrame.Name = name
@@ -324,7 +392,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "series alias"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "series alias"
testFrame.Name = name
@@ -334,7 +402,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias [[m]] [[measurement]]", Measurement: "10m"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias 10m 10m"
testFrame.Name = name
@@ -344,7 +412,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias [[tag_datacenter]]"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias America"
testFrame.Name = name
@@ -354,7 +422,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias [[tag_dc.region.name]]"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias Northeast"
testFrame.Name = name
@@ -364,7 +432,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias [[tag_cluster-name]]"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias Cluster"
testFrame.Name = name
@@ -374,7 +442,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias [[tag_/cluster/name/]]"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias Cluster/"
testFrame.Name = name
@@ -384,7 +452,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias [[tag_@cluster@name@]]"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias Cluster@"
testFrame.Name = name
@@ -395,7 +463,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
})
t.Run("shouldn't parse aliases", func(t *testing.T) {
query = &Query{Alias: "alias words with no brackets"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame := result.Responses["A"]
name := "alias words with no brackets"
testFrame.Name = name
@@ -405,7 +473,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias Test 1.5"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias Test 1.5"
testFrame.Name = name
@@ -415,7 +483,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
}
query = &Query{Alias: "alias Test -1"}
result = parser.Parse(prepare(response), query)
result = parser.Parse(prepare(response), addQueryToQueries(*query))
frame = result.Responses["A"]
name = "alias Test -1"
testFrame.Name = name
@@ -454,6 +522,10 @@ func TestInfluxdbResponseParser(t *testing.T) {
`
query := &Query{}
var queries = addQueryToQueries(*query)
queryB := &Query{}
queryB.RefID = "B"
queries = append(queries, *queryB)
labels, err := data.LabelsFromString("datacenter=America")
require.Nil(t, err)
newField := data.NewField("value", labels, []*float64{
@@ -469,14 +541,15 @@ func TestInfluxdbResponseParser(t *testing.T) {
}),
newField,
)
result := parser.Parse(prepare(response), query)
testFrame.Meta = &data.FrameMeta{ExecutedQueryString: "Test raw query"}
result := parser.Parse(prepare(response), queries)
frame := result.Responses["A"]
if diff := cmp.Diff(testFrame, frame.Frames[0], data.FrameTestCompareOptions()...); diff != "" {
t.Errorf("Result mismatch (-want +got):\n%s", diff)
}
require.EqualError(t, result.Responses["A"].Error, "query-timeout limit exceeded")
require.EqualError(t, result.Responses["B"].Error, "query-timeout limit exceeded")
})
t.Run("Influxdb response parser with top-level error", func(t *testing.T) {
@@ -490,7 +563,7 @@ func TestInfluxdbResponseParser(t *testing.T) {
query := &Query{}
result := parser.Parse(prepare(response), query)
result := parser.Parse(prepare(response), addQueryToQueries(*query))
require.Nil(t, result.Responses["A"].Frames)