API: return query results as JSON rather than base64 encoded Arrow (#32303)
This commit is contained in:
+16
-16
@@ -5,6 +5,7 @@ import (
|
||||
"errors"
|
||||
"net/http"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana/pkg/expr"
|
||||
"github.com/grafana/grafana/pkg/models"
|
||||
"github.com/grafana/grafana/pkg/plugins"
|
||||
@@ -78,16 +79,25 @@ func (hs *HTTPServer) QueryMetricsV2(c *models.ReqContext, reqDTO dtos.MetricReq
|
||||
return response.Error(http.StatusInternalServerError, "Metric request error", err)
|
||||
}
|
||||
|
||||
// This is insanity... but ¯\_(ツ)_/¯, the current query path looks like:
|
||||
// encodeJson( decodeBase64( encodeBase64( decodeArrow( encodeArrow(frame)) ) )
|
||||
// this will soon change to a more direct route
|
||||
qdr, err := resp.ToBackendDataResponse()
|
||||
if err != nil {
|
||||
return response.Error(http.StatusInternalServerError, "error converting results", err)
|
||||
}
|
||||
return toMacronResponse(qdr)
|
||||
}
|
||||
|
||||
func toMacronResponse(qdr *backend.QueryDataResponse) response.Response {
|
||||
statusCode := http.StatusOK
|
||||
for _, res := range resp.Results {
|
||||
for _, res := range qdr.Responses {
|
||||
if res.Error != nil {
|
||||
res.ErrorString = res.Error.Error()
|
||||
resp.Message = res.ErrorString
|
||||
statusCode = http.StatusBadRequest
|
||||
}
|
||||
}
|
||||
|
||||
return response.JSONStreaming(statusCode, resp)
|
||||
return response.JSONStreaming(statusCode, qdr)
|
||||
}
|
||||
|
||||
// handleExpressions handles POST /api/ds/query when there is an expression.
|
||||
@@ -131,21 +141,11 @@ func (hs *HTTPServer) handleExpressions(c *models.ReqContext, reqDTO dtos.Metric
|
||||
Cfg: hs.Cfg,
|
||||
DataService: hs.DataService,
|
||||
}
|
||||
resp, err := exprService.WrapTransformData(c.Req.Context(), request)
|
||||
qdr, err := exprService.WrapTransformData(c.Req.Context(), request)
|
||||
if err != nil {
|
||||
return response.Error(500, "expression request error", err)
|
||||
}
|
||||
|
||||
statusCode := 200
|
||||
for _, res := range resp.Results {
|
||||
if res.Error != nil {
|
||||
res.ErrorString = res.Error.Error()
|
||||
resp.Message = res.ErrorString
|
||||
statusCode = 400
|
||||
}
|
||||
}
|
||||
|
||||
return response.JSONStreaming(statusCode, resp)
|
||||
return toMacronResponse(qdr)
|
||||
}
|
||||
|
||||
func (hs *HTTPServer) handleGetDataSourceError(err error, datasourceID int64) *response.NormalResponse {
|
||||
|
||||
+4
-58
@@ -35,7 +35,7 @@ func init() {
|
||||
}
|
||||
|
||||
// WrapTransformData creates and executes transform requests
|
||||
func (s *Service) WrapTransformData(ctx context.Context, query plugins.DataQuery) (plugins.DataResponse, error) {
|
||||
func (s *Service) WrapTransformData(ctx context.Context, query plugins.DataQuery) (*backend.QueryDataResponse, error) {
|
||||
sdkReq := &backend.QueryDataRequest{
|
||||
PluginContext: backend.PluginContext{
|
||||
OrgID: query.User.OrgId,
|
||||
@@ -46,7 +46,7 @@ func (s *Service) WrapTransformData(ctx context.Context, query plugins.DataQuery
|
||||
for _, q := range query.Queries {
|
||||
modelJSON, err := q.Model.MarshalJSON()
|
||||
if err != nil {
|
||||
return plugins.DataResponse{}, err
|
||||
return nil, err
|
||||
}
|
||||
sdkReq.Queries = append(sdkReq.Queries, backend.DataQuery{
|
||||
JSON: modelJSON,
|
||||
@@ -60,30 +60,7 @@ func (s *Service) WrapTransformData(ctx context.Context, query plugins.DataQuery
|
||||
},
|
||||
})
|
||||
}
|
||||
pbRes, err := s.TransformData(ctx, sdkReq)
|
||||
if err != nil {
|
||||
return plugins.DataResponse{}, err
|
||||
}
|
||||
|
||||
tR := plugins.DataResponse{
|
||||
Results: make(map[string]plugins.DataQueryResult, len(pbRes.Responses)),
|
||||
}
|
||||
for refID, res := range pbRes.Responses {
|
||||
tRes := plugins.DataQueryResult{
|
||||
RefID: refID,
|
||||
Dataframes: plugins.NewDecodedDataFrames(res.Frames),
|
||||
}
|
||||
// if len(res.JsonMeta) != 0 {
|
||||
// tRes.Meta = simplejson.NewFromAny(res.JsonMeta)
|
||||
// }
|
||||
if res.Error != nil {
|
||||
tRes.Error = res.Error
|
||||
tRes.ErrorString = res.Error.Error()
|
||||
}
|
||||
tR.Results[refID] = tRes
|
||||
}
|
||||
|
||||
return tR, nil
|
||||
return s.TransformData(ctx, sdkReq)
|
||||
}
|
||||
|
||||
// TransformData takes Queries which are either expressions nodes
|
||||
@@ -214,37 +191,6 @@ func (s *Service) queryData(ctx context.Context, req *backend.QueryDataRequest)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// Convert tsdb results (map) to plugin-model/datasource (slice) results.
|
||||
// Only error, Series, and encoded Dataframes responses are mapped.
|
||||
responses := make(map[string]backend.DataResponse, len(tsdbRes.Results))
|
||||
for refID, res := range tsdbRes.Results {
|
||||
pRes := backend.DataResponse{}
|
||||
if res.Error != nil {
|
||||
pRes.Error = res.Error
|
||||
}
|
||||
|
||||
if res.Dataframes != nil {
|
||||
decoded, err := res.Dataframes.Decoded()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
pRes.Frames = decoded
|
||||
responses[refID] = pRes
|
||||
continue
|
||||
}
|
||||
|
||||
for _, series := range res.Series {
|
||||
frame, err := plugins.SeriesToFrame(series)
|
||||
frame.RefID = refID
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
pRes.Frames = append(pRes.Frames, frame)
|
||||
}
|
||||
|
||||
responses[refID] = pRes
|
||||
}
|
||||
return &backend.QueryDataResponse{
|
||||
Responses: responses,
|
||||
}, nil
|
||||
return tsdbRes.ToBackendDataResponse()
|
||||
}
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/grafana/grafana/pkg/components/null"
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
@@ -189,6 +190,44 @@ type DataResponse struct {
|
||||
Message string `json:"message,omitempty"`
|
||||
}
|
||||
|
||||
// ToBackendDataResponse converts the legacy format to the standard SDK format
|
||||
func (r DataResponse) ToBackendDataResponse() (*backend.QueryDataResponse, error) {
|
||||
qdr := &backend.QueryDataResponse{
|
||||
Responses: make(map[string]backend.DataResponse, len(r.Results)),
|
||||
}
|
||||
|
||||
// Convert tsdb results (map) to plugin-model/datasource (slice) results.
|
||||
// Only error, Series, and encoded Dataframes responses are mapped.
|
||||
for refID, res := range r.Results {
|
||||
pRes := backend.DataResponse{}
|
||||
if res.Error != nil {
|
||||
pRes.Error = res.Error
|
||||
}
|
||||
|
||||
if res.Dataframes != nil {
|
||||
decoded, err := res.Dataframes.Decoded()
|
||||
if err != nil {
|
||||
return qdr, err
|
||||
}
|
||||
pRes.Frames = decoded
|
||||
qdr.Responses[refID] = pRes
|
||||
continue
|
||||
}
|
||||
|
||||
for _, series := range res.Series {
|
||||
frame, err := SeriesToFrame(series)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
frame.RefID = refID
|
||||
pRes.Frames = append(pRes.Frames, frame)
|
||||
}
|
||||
|
||||
qdr.Responses[refID] = pRes
|
||||
}
|
||||
return qdr, nil
|
||||
}
|
||||
|
||||
type DataPlugin interface {
|
||||
DataQuery(ctx context.Context, ds *models.DataSource, query DataQuery) (DataResponse, error)
|
||||
}
|
||||
|
||||
@@ -14,9 +14,9 @@ import (
|
||||
"github.com/aws/aws-sdk-go/aws/session"
|
||||
"github.com/aws/aws-sdk-go/service/cloudwatch/cloudwatchiface"
|
||||
"github.com/aws/aws-sdk-go/service/cloudwatchlogs/cloudwatchlogsiface"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/grafana/grafana/pkg/models"
|
||||
"github.com/grafana/grafana/pkg/plugins"
|
||||
"github.com/grafana/grafana/pkg/services/sqlstore"
|
||||
"github.com/grafana/grafana/pkg/tests/testinfra"
|
||||
"github.com/grafana/grafana/pkg/tsdb/cloudwatch"
|
||||
@@ -69,7 +69,7 @@ func TestQueryCloudWatchMetrics(t *testing.T) {
|
||||
}
|
||||
result := makeCWRequest(t, req, addr)
|
||||
|
||||
dataFrames := plugins.NewDecodedDataFrames(data.Frames{
|
||||
dataFrames := data.Frames{
|
||||
&data.Frame{
|
||||
RefID: "A",
|
||||
Fields: []*data.Field{
|
||||
@@ -82,21 +82,13 @@ func TestQueryCloudWatchMetrics(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
// Have to call this so that dataFrames.encoded is non-nil, for the comparison
|
||||
// In the future we should use gocmp instead and ignore this field
|
||||
_, err := dataFrames.Encoded()
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t, plugins.DataResponse{
|
||||
Results: map[string]plugins.DataQueryResult{
|
||||
"A": {
|
||||
RefID: "A",
|
||||
Dataframes: dataFrames,
|
||||
},
|
||||
},
|
||||
}, result)
|
||||
expect := backend.NewQueryDataResponse()
|
||||
expect.Responses["A"] = backend.DataResponse{
|
||||
Frames: dataFrames,
|
||||
}
|
||||
assert.Equal(t, *expect, result)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -130,7 +122,7 @@ func TestQueryCloudWatchLogs(t *testing.T) {
|
||||
}
|
||||
tr := makeCWRequest(t, req, addr)
|
||||
|
||||
dataFrames := plugins.NewDecodedDataFrames(data.Frames{
|
||||
dataFrames := data.Frames{
|
||||
&data.Frame{
|
||||
Name: "logGroups",
|
||||
RefID: "A",
|
||||
@@ -141,23 +133,17 @@ func TestQueryCloudWatchLogs(t *testing.T) {
|
||||
PreferredVisualization: "logs",
|
||||
},
|
||||
},
|
||||
})
|
||||
// Have to call this so that dataFrames.encoded is non-nil, for the comparison
|
||||
// In the future we should use gocmp instead and ignore this field
|
||||
_, err := dataFrames.Encoded()
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, plugins.DataResponse{
|
||||
Results: map[string]plugins.DataQueryResult{
|
||||
"A": {
|
||||
RefID: "A",
|
||||
Dataframes: dataFrames,
|
||||
},
|
||||
},
|
||||
}, tr)
|
||||
}
|
||||
|
||||
expect := backend.NewQueryDataResponse()
|
||||
expect.Responses["A"] = backend.DataResponse{
|
||||
Frames: dataFrames,
|
||||
}
|
||||
assert.Equal(t, *expect, tr)
|
||||
})
|
||||
}
|
||||
|
||||
func makeCWRequest(t *testing.T, req dtos.MetricRequest, addr string) plugins.DataResponse {
|
||||
func makeCWRequest(t *testing.T, req dtos.MetricRequest, addr string) backend.QueryDataResponse {
|
||||
t.Helper()
|
||||
|
||||
buf := bytes.Buffer{}
|
||||
@@ -180,7 +166,7 @@ func makeCWRequest(t *testing.T, req dtos.MetricRequest, addr string) plugins.Da
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 200, resp.StatusCode)
|
||||
|
||||
var tr plugins.DataResponse
|
||||
var tr backend.QueryDataResponse
|
||||
err = json.Unmarshal(buf.Bytes(), &tr)
|
||||
require.NoError(t, err)
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
jsoniter "github.com/json-iterator/go"
|
||||
|
||||
"github.com/grafana/grafana/pkg/cmd/grafana-cli/logger"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
@@ -87,7 +88,7 @@ func (p *testStreamHandler) runTestStream(ctx context.Context, path string, conf
|
||||
frame.Fields[2].Set(0, walker-((rand.Float64()*spread)+0.01)) // Min
|
||||
frame.Fields[3].Set(0, walker+((rand.Float64()*spread)+0.01)) // Max
|
||||
|
||||
bytes, err := data.FrameToJSON(frame, true, true)
|
||||
bytes, err := jsoniter.Marshal(frame) // schema + points
|
||||
if err != nil {
|
||||
logger.Warn("unable to marshal line", "error", err)
|
||||
continue
|
||||
|
||||
Reference in New Issue
Block a user