Query Service: Combine SSE handling in single tenant and multi tenant paths (#108041)

* parse via sse

I need to figure out how to handle the pipeline.execute with our own
client. I think this is important for MT reasons, just like using our
own cache (via legacy) is important.

parsing is done though!

* WIP nonsense

* horrible code but i think it works

* Add support for sql expressions config settings

* Cleanup:
- remove spew from nodes.go
- uncomment out plugin context and use in single tenant flow
- make code more readable and add comments

* Cleanup:
- create separate file for mt ds client builder
- ensure error handling is the same for both expressions and regular queries
- other cleanup

* not working but good thoughts

* WIP, vector not working for non sse

* super hacky but i think vectors work now

* delete delete delete

* Comments for future ref

* break out query handling and start test

* add prom debugger

* clean up: remove comments and commented out bits

* fix query_test

* add prom debugger

* create table-driven tests with testsdata files

* Fix test

* Add test

* go mod??

* idk

* Remove comment

* go enterprise issue maybe

* Fix codeowners

* Delete

* Remove test data

* Clean up

* logger

* Remove go changes hopefully

* idk go man

* sad

* idk i ran go mod tidy and this is what it wants

* Fix readme, with much help from adam

* some linting and testing errors

* lint

* fix lint

* fix lint register.go

* another lint

* address lint in test

* fix dead code and linters for query_test

* Go mod?

* Struggling with go mod

* Fix test

* Fix another test

* Revert headers change

* Its difficult to test this in OSS as it depends on functionality defined in enterprise, let's bring these tests back in some form in enterprise

* Fix codeowners

---------

Co-authored-by: Adam Simpson <adam@adamsimpson.net>
This commit is contained in:
Sarah Zinger
2025-07-17 17:22:55 -04:00
committed by GitHub
co-authored by Adam Simpson
parent 8acdc7e37b
commit 3fad863fd1
34 changed files with 763 additions and 1662 deletions
+1 -1
View File
@@ -150,6 +150,7 @@
/pkg/services/hooks/ @grafana/grafana-backend-group
/pkg/services/kmsproviders/ @grafana/grafana-operator-experience-squad
/pkg/services/licensing/ @grafana/grafana-operator-experience-squad
/pkg/services/mtdsclient/ @grafana/grafana-datasources-core-services
/pkg/services/navtree/ @grafana/grafana-backend-group
/pkg/services/notifications/ @grafana/grafana-backend-group
/pkg/services/org/ @grafana/grafana-backend-group
@@ -178,7 +179,6 @@
/pkg/setting/ @grafana/grafana-backend-services-squad
/pkg/tests/ @grafana/grafana-backend-services-squad
/pkg/tests/apis/ @grafana/grafana-app-platform-squad
/pkg/tests/apis/query @grafana/grafana-datasources-core-services
/pkg/tests/apis/alerting @grafana/grafana-app-platform-squad @grafana/alerting-backend
/pkg/tests/api/correlations/ @grafana/datapro
/pkg/tsdb/grafanads/ @grafana/grafana-backend-group
-2
View File
@@ -14,8 +14,6 @@ import (
v1beta1 "github.com/grafana/grafana/apps/secret/pkg/apis/secret/v1beta1"
)
var ()
var appManifestData = app.ManifestData{
AppName: "secret",
Group: "secret.grafana.app",
+20 -7
View File
@@ -22,6 +22,7 @@ import (
"github.com/grafana/grafana/pkg/plugins/manager/registry"
"github.com/grafana/grafana/pkg/services/datasources"
fakeDatasources "github.com/grafana/grafana/pkg/services/datasources/fakes"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginconfig"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
pluginSettings "github.com/grafana/grafana/pkg/services/pluginsintegration/pluginsettings/service"
@@ -60,16 +61,27 @@ func TestAPIEndpoint_Metrics_QueryMetricsV2(t *testing.T) {
return &backend.QueryDataResponse{Responses: resp}, nil
},
},
plugincontext.ProvideService(cfg, localcache.ProvideService(), &pluginstore.FakePluginStore{
PluginList: []pluginstore.Plugin{
{
JSONData: plugins.JSONData{
ID: "grafana",
plugincontext.ProvideService(
cfg,
localcache.ProvideService(),
&pluginstore.FakePluginStore{
PluginList: []pluginstore.Plugin{
{
JSONData: plugins.JSONData{
ID: "grafana",
},
},
},
},
}, &fakeDatasources.FakeCacheService{}, &fakeDatasources.FakeDataSourceService{},
pluginSettings.ProvideService(dbtest.NewFakeDB(), secretstest.NewFakeSecretsService()), pluginconfig.NewFakePluginRequestConfigProvider()),
&fakeDatasources.FakeCacheService{},
&fakeDatasources.FakeDataSourceService{},
pluginSettings.ProvideService(
dbtest.NewFakeDB(),
secretstest.NewFakeSecretsService(),
),
pluginconfig.NewFakePluginRequestConfigProvider(),
),
mtdsclient.NewNullMTDatasourceClientBuilder(),
)
server := SetupAPITestServer(t, func(hs *HTTPServer) {
hs.queryDataService = qds
@@ -252,6 +264,7 @@ func TestDataSourceQueryError(t *testing.T) {
&fakeDatasources.FakeCacheService{}, ds,
pluginSettings.ProvideService(dbtest.NewFakeDB(),
secretstest.NewFakeSecretsService()), pluginconfig.NewFakePluginRequestConfigProvider()),
mtdsclient.NewNullMTDatasourceClientBuilder(),
)
hs.QuotaService = quotatest.New(false, nil)
})
+2
View File
@@ -18,6 +18,7 @@ import (
"github.com/grafana/grafana/pkg/services/datasources"
datafakes "github.com/grafana/grafana/pkg/services/datasources/fakes"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginconfig"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore"
@@ -70,6 +71,7 @@ func framesPassThroughService(t *testing.T, frames data.Frames) (data.Frames, er
Features: features,
Tracer: tracing.InitializeTracerForTest(),
},
mtDatasourceClientBuilder: mtdsclient.NewNullMTDatasourceClientBuilder(),
}
queries := []Query{{
RefID: "A",
+1 -1
View File
@@ -134,7 +134,7 @@ func (m *MLNode) Execute(ctx context.Context, now time.Time, _ mathexp.Vars, s *
return result, err
}
func (s *Service) buildMLNode(dp *simple.DirectedGraph, rn *rawNode, req *Request) (Node, error) {
func (s *Service) buildMLNode(_ *simple.DirectedGraph, rn *rawNode, req *Request) (Node, error) {
if rn.TimeRange == nil {
return nil, errors.New("time range must be specified")
}
+33 -9
View File
@@ -213,7 +213,7 @@ func (dn *DSNode) NeedsVars() []string {
return []string{}
}
func (s *Service) buildDSNode(dp *simple.DirectedGraph, rn *rawNode, req *Request) (*DSNode, error) {
func (s *Service) buildDSNode(_ *simple.DirectedGraph, rn *rawNode, req *Request) (*DSNode, error) {
if rn.TimeRange == nil {
return nil, fmt.Errorf("time range must be specified for refID %s", rn.RefID)
}
@@ -361,17 +361,12 @@ func (dn *DSNode) Execute(ctx context.Context, now time.Time, _ mathexp.Vars, s
ctx, span := s.tracer.Start(ctx, "SSE.ExecuteDatasourceQuery")
defer span.End()
pCtx, err := s.pCtxProvider.GetWithDataSource(ctx, dn.datasource.Type, dn.request.User, dn.datasource)
if err != nil {
return mathexp.Results{}, err
}
span.SetAttributes(
attribute.String("datasource.type", dn.datasource.Type),
attribute.String("datasource.uid", dn.datasource.UID),
)
req := &backend.QueryDataRequest{
PluginContext: pCtx,
Queries: []backend.DataQuery{
{
RefID: dn.refID,
@@ -399,9 +394,38 @@ func (dn *DSNode) Execute(ctx context.Context, now time.Time, _ mathexp.Vars, s
s.metrics.DSRequests.WithLabelValues(respStatus, fmt.Sprintf("%t", useDataplane), dn.datasource.Type).Inc()
}()
resp, err := s.dataService.QueryData(ctx, req)
if err != nil {
return mathexp.Results{}, MakeQueryError(dn.refID, dn.datasource.UID, err)
var resp *backend.QueryDataResponse
mtDSClient, ok := s.mtDatasourceClientBuilder.BuildClient(dn.datasource.Type, dn.datasource.UID)
if !ok { // use single tenant client
pCtx, err := s.pCtxProvider.GetWithDataSource(ctx, dn.datasource.Type, dn.request.User, dn.datasource)
if err != nil {
return mathexp.Results{}, err
}
req.PluginContext = pCtx
resp, err = s.dataService.QueryData(ctx, req)
if err != nil {
return mathexp.Results{}, MakeQueryError(dn.refID, dn.datasource.UID, err)
}
} else {
// transform request from backend.QueryDataRequest to k8s request
k8sReq := &data.QueryDataRequest{}
for _, q := range req.Queries {
var dataQuery data.DataQuery
err := json.Unmarshal(q.JSON, &dataQuery)
if err != nil {
return mathexp.Results{}, MakeQueryError(dn.refID, dn.datasource.UID, err)
}
k8sReq.Queries = append(k8sReq.Queries, dataQuery)
}
var err error
// make the query with a mt client
resp, err = mtDSClient.QueryData(ctx, *k8sReq)
// handle error
if err != nil {
return mathexp.Results{}, MakeQueryError(dn.refID, dn.datasource.UID, err)
}
}
dataFrames, err := getResponseFrame(logger, resp, dn.refID)
+6 -3
View File
@@ -16,6 +16,7 @@ import (
"github.com/grafana/grafana/pkg/plugins"
"github.com/grafana/grafana/pkg/services/datasources"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
"github.com/grafana/grafana/pkg/setting"
)
@@ -65,8 +66,9 @@ type Service struct {
pluginsClient backend.CallResourceHandler
tracer tracing.Tracer
metrics *metrics.ExprMetrics
tracer tracing.Tracer
metrics *metrics.ExprMetrics
mtDatasourceClientBuilder mtdsclient.MTDatasourceClientBuilder
}
type pluginContextProvider interface {
@@ -75,7 +77,7 @@ type pluginContextProvider interface {
}
func ProvideService(cfg *setting.Cfg, pluginClient plugins.Client, pCtxProvider *plugincontext.Provider,
features featuremgmt.FeatureToggles, registerer prometheus.Registerer, tracer tracing.Tracer) *Service {
features featuremgmt.FeatureToggles, registerer prometheus.Registerer, tracer tracing.Tracer, builder mtdsclient.MTDatasourceClientBuilder) *Service {
return &Service{
cfg: cfg,
dataService: pluginClient,
@@ -88,6 +90,7 @@ func ProvideService(cfg *setting.Cfg, pluginClient plugins.Client, pCtxProvider
Features: features,
Tracer: tracer,
},
mtDatasourceClientBuilder: builder,
}
}
+2
View File
@@ -20,6 +20,7 @@ import (
"github.com/grafana/grafana/pkg/services/datasources"
datafakes "github.com/grafana/grafana/pkg/services/datasources/fakes"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginconfig"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore"
@@ -255,5 +256,6 @@ func newMockQueryService(responses map[string]backend.DataResponse, queries []Qu
Features: features,
Tracer: tracing.InitializeTracerForTest(),
},
mtDatasourceClientBuilder: mtdsclient.NewNullMTDatasourceClientBuilder(),
}, &Request{Queries: queries, User: &user.SignedInUser{}}
}
+26 -19
View File
@@ -1,11 +1,6 @@
# Query service
This query service aims to replace the existing /api/ds/query.
The key differences are:
1. This service has a stronger type system (not simplejson)
2. Same workflow regardless if expressions exist
3. Datasource settings+access is managed in each datasource, not at the beginning
This query service aims to replace the existing /api/ds/query, while preserving the same parsing and expression handling as `/api/ds/query`
@@ -57,22 +52,34 @@ sequenceDiagram
autonumber
actor User as User or Process
participant api as /apis/query.grafana.app
participant ds as Datasource<br/>Handler/Plugin
participant db as Storage<br/>(SQL)
participant db as Storage<br/>(CloudConfig)
participant ds as Datasource<br/>Plugin
participant expr as Expression<br/>Engine
User->>api: POST Query
api->>api: Parse queries
api->>api: Calculate dependencies
loop Each datasource (concurrently)
api->>ds: QueryData
ds->>ds: Verify user access
ds->>db: Get settings <br> and secrets
loop Each query
api->>api: Parse query
api->>db: Get ds config<br>and secrets
db->>api:
end
loop Each expression
api->>expr: Execute
alt Expressions exist
api->>expr: Calculate expressions graph
loop Each node (eg, refID)
alt Is query
expr->>ds: QueryData
else Is expression
expr->>expr: Process
end
end
else No expressions
alt Single datasource
api->>ds: QueryData
else Multiple datasources
loop Each datasource (concurrently)
api->>ds: QueryData
end
api->>api: Wait for results
end
end
api->>api: Verify ResultExpectations
api->>User: return results
```
```
+1 -29
View File
@@ -18,7 +18,6 @@ import (
"github.com/grafana/grafana/pkg/plugins"
"github.com/grafana/grafana/pkg/registry/apis/query/clientapi"
"github.com/grafana/grafana/pkg/services/accesscontrol"
"github.com/grafana/grafana/pkg/services/contexthandler"
"github.com/grafana/grafana/pkg/services/datasources"
"github.com/grafana/grafana/pkg/services/pluginsintegration/adapters"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
@@ -110,34 +109,7 @@ func getGrafanaDataSourceSettings(ctx context.Context) (*backend.DataSourceInsta
return adapters.ModelToInstanceSettings(ds, decryptFunc)
}
func (d *pluginClient) QueryData(ctx context.Context, req data.QueryDataRequest) (*clientapi.Response, error) {
// middlewares may set response-http-headers through context, so we need to do extra steps
// first we create an isolated context to be used by query-data
isolatedCtx := contexthandler.CopyWithReqContext(ctx)
qdr, err := d.innerQueryData(isolatedCtx, req)
if err != nil {
return nil, err
}
rsp := &clientapi.Response{
QDR: qdr,
Headers: nil,
}
// we extract the response-headers from the isolated context, and return them explicitly
// in the clientapi.Response structure
reqCtx := contexthandler.FromContext(isolatedCtx)
if reqCtx != nil {
rsp.Headers = reqCtx.Resp.Header()
}
return rsp, nil
}
// ExecuteQueryData implements QueryHelper.
func (d *pluginClient) innerQueryData(ctx context.Context, req data.QueryDataRequest) (*backend.QueryDataResponse, error) {
func (d *pluginClient) QueryData(ctx context.Context, req data.QueryDataRequest) (*backend.QueryDataResponse, error) {
queries, dsRef, err := data.ToDataSourceQueries(req)
if err != nil {
return nil, err
+10 -5
View File
@@ -3,6 +3,7 @@ package clientapi
import (
"context"
"net/http"
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
data "github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
@@ -15,14 +16,18 @@ type Response struct {
}
type QueryDataClient interface {
QueryData(ctx context.Context, req data.QueryDataRequest) (*Response, error)
QueryData(ctx context.Context, req data.QueryDataRequest) (*backend.QueryDataResponse, error)
}
type InstanceConfigurationSettings struct {
StackID uint32
FeatureToggles featuremgmt.FeatureToggles
FullConfig map[string]map[string]string // configuration file settings
Options map[string]string // additional settings related to an instance as set by grafana
StackID uint32
FeatureToggles featuremgmt.FeatureToggles
FullConfig map[string]map[string]string // configuration file settings
Options map[string]string // additional settings related to an instance as set by grafana
SQLExpressionCellLimit int64
SQLExpressionOutputCellLimit int64
SQLExpressionTimeout time.Duration
ExpressionsEnabled bool
}
type DataSourceClientSupplier interface {
-35
View File
@@ -2,7 +2,6 @@ package query
import (
"errors"
"fmt"
"github.com/grafana/grafana/pkg/apimachinery/errutil"
)
@@ -44,40 +43,6 @@ func MakePublicQueryError(refID, err string) error {
return QueryError.Build(data)
}
var depErrStr = "did not execute expression [{{ .Public.refId }}] due to a failure of the dependent expression or query [{{.Public.depRefId}}]"
var dependencyError = errutil.BadRequest("sse.dependencyError").MustTemplate(
depErrStr,
errutil.WithPublic(depErrStr))
func makeDependencyError(refID, depRefID string) error {
data := errutil.TemplateData{
Public: map[string]interface{}{
"refId": refID,
"depRefId": depRefID,
},
Error: fmt.Errorf("did not execute expression %v due to a failure of the dependent expression or query %v", refID, depRefID),
}
return dependencyError.Build(data)
}
var cyclicErrStr = "cyclic reference in expression [{{ .Public.refId }}]"
var cyclicErr = errutil.BadRequest("sse.cyclic").MustTemplate(
cyclicErrStr,
errutil.WithPublic(cyclicErrStr))
func makeCyclicError(refID string) error {
data := errutil.TemplateData{
Public: map[string]interface{}{
"refId": refID,
},
Error: fmt.Errorf("cyclic reference in %s", refID),
}
return cyclicErr.Build(data)
}
type ErrorWithRefID struct {
err error
refId string
-259
View File
@@ -1,259 +0,0 @@
package query
import (
"context"
"encoding/json"
"fmt"
"github.com/grafana/grafana-plugin-sdk-go/data/utils/jsoniter"
data "github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
"gonum.org/v1/gonum/graph/simple"
"gonum.org/v1/gonum/graph/topo"
query "github.com/grafana/grafana/pkg/apis/query/v0alpha1"
"github.com/grafana/grafana/pkg/expr"
"github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/infra/tracing"
"github.com/grafana/grafana/pkg/services/datasources/service"
)
type datasourceRequest struct {
// The type
PluginId string `json:"pluginId"`
// The UID
UID string `json:"uid"`
// Optionally show the additional query properties
Request *data.QueryDataRequest `json:"request"`
// Headers that should be forwarded to the next request
Headers map[string]string `json:"headers,omitempty"`
}
type parsedRequestInfo struct {
// Datasource queries, one for each datasource
Requests []datasourceRequest `json:"requests,omitempty"`
// Expressions in required execution order
Expressions []expr.ExpressionQuery `json:"expressions,omitempty"`
// Expressions include explicit hacks for influx+prometheus
RefIDTypes map[string]string `json:"types,omitempty"`
// Hidden queries used as dependencies
HideBeforeReturn []string `json:"hide,omitempty"`
// SQL Inputs
SqlInputs map[string]struct{} `json:"sqlInputs,omitempty"`
}
type queryParser struct {
legacy service.LegacyDataSourceLookup
reader *expr.ExpressionQueryReader
tracer tracing.Tracer
logger log.Logger
}
func newQueryParser(reader *expr.ExpressionQueryReader, legacy service.LegacyDataSourceLookup, tracer tracing.Tracer, logger log.Logger) *queryParser {
return &queryParser{
reader: reader,
legacy: legacy,
tracer: tracer,
logger: logger,
}
}
// Split the main query into multiple
func (p *queryParser) parseRequest(ctx context.Context, input *query.QueryDataRequest) (parsedRequestInfo, error) {
ctx, span := p.tracer.Start(ctx, "QueryService.parseRequest")
defer span.End()
queryRefIDs := make(map[string]*data.DataQuery, len(input.Queries))
expressions := make(map[string]*expr.ExpressionQuery)
index := make(map[string]int) // index lookup
rsp := parsedRequestInfo{
RefIDTypes: make(map[string]string, len(input.Queries)),
SqlInputs: make(map[string]struct{}),
}
for _, q := range input.Queries {
_, found := queryRefIDs[q.RefID]
if found {
return rsp, MakePublicQueryError(q.RefID, "multiple queries with same refId")
}
_, found = expressions[q.RefID]
if found {
return rsp, MakePublicQueryError(q.RefID, "multiple queries with same refId")
}
ds, err := p.getValidDataSourceRef(ctx, q.Datasource, q.DatasourceID)
if err != nil {
p.logger.Error("Failed to get valid datasource ref", "error", err)
return rsp, err
}
// Process each query
// check if ds is expression
if expr.IsDataSource(ds.UID) {
// In order to process the query as a typed expression query, we
// are writing it back to JSON and parsing again. Alternatively we
// could construct it from the untyped map[string]any additional properties
// but this approach lets us focus on well typed behavior first
raw, err := json.Marshal(q)
if err != nil {
p.logger.Error("Failed to marshal query for expression", "error", err)
return rsp, err
}
iter, err := jsoniter.ParseBytes(jsoniter.ConfigDefault, raw)
if err != nil {
p.logger.Error("Failed to parse bytes for expression", "error", err)
return rsp, err
}
exp, err := p.reader.ReadQuery(q, iter)
if err != nil {
p.logger.Error("Failed to read query for expression", "error", err)
return rsp, NewErrorWithRefID(q.RefID, err)
}
exp.GraphID = int64(len(expressions) + 1)
expressions[q.RefID] = &exp
} else {
key := fmt.Sprintf("%s/%s", ds.Type, ds.UID)
idx, ok := index[key]
if !ok {
idx = len(index)
index[key] = idx
rsp.Requests = append(rsp.Requests, datasourceRequest{
PluginId: ds.Type,
UID: ds.UID,
Request: &data.QueryDataRequest{
TimeRange: getTimeRangeForQuery(&input.TimeRange, q.TimeRange),
Debug: input.Debug,
// no queries
},
})
}
req := rsp.Requests[idx].Request
req.Queries = append(req.Queries, q)
queryRefIDs[q.RefID] = &req.Queries[len(req.Queries)-1]
}
// Mark all the queries that should be hidden ()
if q.Hide {
rsp.HideBeforeReturn = append(rsp.HideBeforeReturn, q.RefID)
}
}
// Make sure all referenced variables exist and the expression order is stable
if len(expressions) > 0 {
queryNode := &expr.ExpressionQuery{
GraphID: -1,
}
// Build the graph for a request
dg := simple.NewDirectedGraph()
dg.AddNode(queryNode)
for _, exp := range expressions {
dg.AddNode(exp)
}
for _, exp := range expressions {
vars := exp.Command.NeedsVars()
for _, refId := range vars {
target := queryNode
q, ok := queryRefIDs[refId]
if !ok {
if target, ok = expressions[refId]; !ok {
return rsp, makeDependencyError(exp.RefID, refId)
}
}
// If the input is SQL, conversion is handled differently
if _, isSqlExp := exp.Command.(*expr.SQLCommand); isSqlExp {
if _, ifDepIsAlsoExpression := expressions[refId]; ifDepIsAlsoExpression {
// Only allow data source nodes as SQL expression inputs for now
return rsp, fmt.Errorf("only data source queries may be inputs to a sql expression, %v is the input for %v", refId, exp.RefID)
} else {
rsp.SqlInputs[refId] = struct{}{}
}
}
// Do not hide queries used in variables
if q != nil && q.Hide {
q.Hide = false
}
if target.ID() == exp.ID() {
return rsp, makeCyclicError(refId)
}
dg.SetEdge(dg.NewEdge(target, exp))
}
}
// Add the sorted expressions
sortedNodes, err := topo.SortStabilized(dg, nil)
if err != nil {
p.logger.Error("Error when sorting nodes", "error", err)
return rsp, makeCyclicError("")
}
for _, v := range sortedNodes {
if v.ID() > 0 {
rsp.Expressions = append(rsp.Expressions, *v.(*expr.ExpressionQuery))
}
}
}
return rsp, nil
}
func getTimeRangeForQuery(parentTimerange, queryTimerange *data.TimeRange) data.TimeRange {
if queryTimerange != nil && queryTimerange.From != "" && queryTimerange.To != "" {
return *queryTimerange
}
if parentTimerange != nil && parentTimerange.To != "" && parentTimerange.From != "" {
return *parentTimerange
}
return data.TimeRange{
From: "0",
To: "0",
}
}
func (p *queryParser) getValidDataSourceRef(ctx context.Context, ds *data.DataSourceRef, id int64) (*data.DataSourceRef, error) {
if ds == nil {
if id == 0 {
return nil, fmt.Errorf("missing datasource reference or id")
}
if p.legacy == nil {
return nil, fmt.Errorf("legacy datasource lookup unsupported (id:%d)", id)
}
return p.legacy.GetDataSourceFromDeprecatedFields(ctx, "", id)
}
// we need to special-case the "grafana" data source
if ds.UID == "grafana" {
return &data.DataSourceRef{
// it does not really matter what `type` we set here,
// we will always detect this case by `uid` later.
// here we go with what the data source's plugin.json says.
Type: "grafana",
UID: "grafana",
}, nil
}
if ds.Type == "" {
if ds.UID == "" {
return nil, fmt.Errorf("missing name/uid in data source reference")
}
if expr.IsDataSource(ds.UID) {
return ds, nil
}
if p.legacy == nil {
return nil, fmt.Errorf("legacy datasource lookup unsupported (name:%s)", ds.UID)
}
return p.legacy.GetDataSourceFromDeprecatedFields(ctx, ds.UID, 0)
}
return ds, nil
}
-398
View File
@@ -1,398 +0,0 @@
package query
import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"path"
"strings"
"testing"
data "github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
query "github.com/grafana/grafana/pkg/apis/query/v0alpha1"
"github.com/grafana/grafana/pkg/expr"
"github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/infra/tracing"
"github.com/grafana/grafana/pkg/services/featuremgmt"
)
type parserTestObject struct {
Description string `json:"description,omitempty"`
Request query.QueryDataRequest `json:"input"`
Expect parsedRequestInfo `json:"expect"`
Error string `json:"error,omitempty"`
}
func TestQuerySplitting(t *testing.T) {
ctx := context.Background()
parser := newQueryParser(expr.NewExpressionQueryReader(featuremgmt.WithFeatures()),
&legacyDataSourceRetriever{}, tracing.InitializeTracerForTest(), log.NewNopLogger())
t.Run("missing datasource flavors", func(t *testing.T) {
split, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
Queries: []data.DataQuery{{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "A",
},
}},
},
})
require.Error(t, err) // Missing datasource
require.Empty(t, split.Requests)
})
t.Run("applies zero time range if time range is missing", func(t *testing.T) {
split, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
TimeRange: data.TimeRange{}, // missing
Queries: []data.DataQuery{{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "A",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
},
}},
},
})
require.NoError(t, err)
require.Len(t, split.Requests, 1)
require.Equal(t, "0", split.Requests[0].Request.From)
require.Equal(t, "0", split.Requests[0].Request.To)
})
t.Run("forbid duplicate refId", func(t *testing.T) {
_, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
TimeRange: data.TimeRange{},
Queries: []data.DataQuery{
{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "A",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
},
},
{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "A",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
},
},
},
},
})
require.Error(t, err)
})
t.Run("forbid duplicate refId, when refId=''", func(t *testing.T) {
_, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
TimeRange: data.TimeRange{},
Queries: []data.DataQuery{
{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
},
},
{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
},
},
},
},
})
require.Error(t, err)
})
t.Run("allow empty refId", func(t *testing.T) {
_, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
TimeRange: data.TimeRange{},
Queries: []data.DataQuery{
{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
},
},
{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "B",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
},
},
},
},
})
require.NoError(t, err)
})
t.Run("applies query time range if present", func(t *testing.T) {
split, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
TimeRange: data.TimeRange{}, // missing
Queries: []data.DataQuery{{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "A",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
TimeRange: &data.TimeRange{
From: "now-1d",
To: "now",
},
},
}},
},
})
require.NoError(t, err)
require.Len(t, split.Requests, 1)
require.Equal(t, "now-1d", split.Requests[0].Request.From)
require.Equal(t, "now", split.Requests[0].Request.To)
})
t.Run("applies query time range if all time ranges are present", func(t *testing.T) {
split, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
TimeRange: data.TimeRange{
From: "now-1h",
To: "now",
},
Queries: []data.DataQuery{{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "A",
Datasource: &data.DataSourceRef{
Type: "x",
UID: "abc",
},
TimeRange: &data.TimeRange{
From: "now-1d",
To: "now",
},
},
}},
},
})
require.NoError(t, err)
require.Len(t, split.Requests, 1)
require.Equal(t, "now-1d", split.Requests[0].Request.From)
require.Equal(t, "now", split.Requests[0].Request.To)
})
t.Run("verify tests", func(t *testing.T) {
files, err := os.ReadDir("testdata")
require.NoError(t, err)
for _, file := range files {
if !strings.HasSuffix(file.Name(), ".json") {
continue
}
t.Run(file.Name(), func(t *testing.T) {
fpath := path.Join("testdata", file.Name())
// nolint:gosec
body, err := os.ReadFile(fpath)
require.NoError(t, err)
harness := &parserTestObject{}
err = json.Unmarshal(body, harness)
require.NoError(t, err)
changed := false
parsed, err := parser.parseRequest(ctx, &harness.Request)
if err != nil {
if !assert.Equal(t, harness.Error, err.Error(), "File %s", file) {
changed = true
}
} else {
x, _ := json.Marshal(parsed)
y, _ := json.Marshal(harness.Expect)
if !assert.JSONEq(t, string(y), string(x), "File %s", file) {
changed = true
}
}
if changed {
harness.Error = ""
harness.Expect = parsed
if err != nil {
harness.Error = err.Error()
}
jj, err := json.MarshalIndent(harness, "", " ")
require.NoError(t, err)
err = os.WriteFile(fpath, jj, 0600)
require.NoError(t, err)
}
})
}
})
}
func TestSqlInputs(t *testing.T) {
parser := newQueryParser(
expr.NewExpressionQueryReader(featuremgmt.WithFeatures(featuremgmt.FlagSqlExpressions)),
nil,
tracing.InitializeTracerForTest(),
log.NewNopLogger(),
)
parsedRequestInfo, err := parser.parseRequest(context.Background(), &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
Queries: []data.DataQuery{
data.NewDataQuery(map[string]any{
"refId": "A",
"datasource": &data.DataSourceRef{
Type: "prometheus",
UID: "local-prom",
},
}),
data.NewDataQuery(map[string]any{
"refId": "B",
"datasource": &data.DataSourceRef{
Type: "__expr__",
UID: "__expr__",
},
"type": "sql",
"expression": "Select time, value + 10 from A",
}),
},
},
})
require.NoError(t, err)
require.Equal(t, parsedRequestInfo.SqlInputs["B"], struct{}{})
}
func TestSqlCTE(t *testing.T) {
parser := newQueryParser(
expr.NewExpressionQueryReader(featuremgmt.WithFeatures(featuremgmt.FlagSqlExpressions)),
nil,
tracing.InitializeTracerForTest(),
log.NewNopLogger(),
)
parsedRequestInfo, err := parser.parseRequest(context.Background(), &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
Queries: []data.DataQuery{
data.NewDataQuery(map[string]any{
"refId": "A",
"datasource": &data.DataSourceRef{
Type: "prometheus",
UID: "local-prom",
},
}),
data.NewDataQuery(map[string]any{
"refId": "B",
"datasource": &data.DataSourceRef{
Type: "__expr__",
UID: "__expr__",
},
"type": "sql",
"expression": `WITH CTE AS (
SELECT
Month
FROM A
)
SELECT * FROM CTE`,
}),
},
},
})
require.NoError(t, err)
require.Equal(t, parsedRequestInfo.SqlInputs["B"], struct{}{})
}
func TestGrafanaDS(t *testing.T) {
ctx := context.Background()
parser := newQueryParser(expr.NewExpressionQueryReader(featuremgmt.WithFeatures()),
&noLegacyRetriever{}, tracing.InitializeTracerForTest(), log.NewNopLogger())
t.Run("grafana ds without type", func(t *testing.T) {
parsed, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
Queries: []data.DataQuery{{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "A",
Datasource: &data.DataSourceRef{
UID: "grafana",
},
},
}},
},
})
require.NoError(t, err)
require.Len(t, parsed.Requests, 1)
require.Equal(t, "grafana", parsed.Requests[0].PluginId)
require.Equal(t, "grafana", parsed.Requests[0].UID)
})
t.Run("grafana ds with different type", func(t *testing.T) {
parsed, err := parser.parseRequest(ctx, &query.QueryDataRequest{
QueryDataRequest: data.QueryDataRequest{
Queries: []data.DataQuery{{
CommonQueryProperties: data.CommonQueryProperties{
RefID: "A",
Datasource: &data.DataSourceRef{
UID: "grafana",
Type: "datasource",
},
},
}},
},
})
require.NoError(t, err)
require.Len(t, parsed.Requests, 1)
require.Equal(t, "grafana", parsed.Requests[0].PluginId)
require.Equal(t, "grafana", parsed.Requests[0].UID)
})
}
type noLegacyRetriever struct{}
var errNoLegacy = errors.New("legacy dds retriever reached, it should not")
func (s *noLegacyRetriever) GetDataSourceFromDeprecatedFields(ctx context.Context, name string, id int64) (*data.DataSourceRef, error) {
return nil, errNoLegacy
}
type legacyDataSourceRetriever struct{}
func (s *legacyDataSourceRetriever) GetDataSourceFromDeprecatedFields(ctx context.Context, name string, id int64) (*data.DataSourceRef, error) {
if id == 100 {
return &data.DataSourceRef{
Type: "plugin-aaaa",
UID: "AAA",
}, nil
}
if name != "" {
return &data.DataSourceRef{
Type: "plugin-bbb",
UID: name,
}, nil
}
return nil, fmt.Errorf("missing parameter")
}
+112 -381
View File
@@ -2,31 +2,32 @@ package query
import (
"context"
"errors"
"fmt"
"encoding/json"
"net/http"
"slices"
"strconv"
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
"github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
"github.com/grafana/grafana/pkg/registry/apis/query/clientapi"
"github.com/grafana/grafana/pkg/services/contexthandler"
"github.com/grafana/grafana/pkg/api/dtos"
"github.com/grafana/grafana/pkg/components/simplejson"
"github.com/grafana/grafana/pkg/expr"
"github.com/grafana/grafana/pkg/services/datasources"
"github.com/grafana/grafana/pkg/services/ngalert/models"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/setting"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes"
"golang.org/x/sync/errgroup"
errorsK8s "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apiserver/pkg/endpoints/request"
"k8s.io/apiserver/pkg/registry/rest"
"github.com/grafana/grafana/pkg/apimachinery/identity"
query "github.com/grafana/grafana/pkg/apis/query/v0alpha1"
"github.com/grafana/grafana/pkg/expr/mathexp"
"github.com/grafana/grafana/pkg/infra/log"
ds_service "github.com/grafana/grafana/pkg/services/datasources/service"
service "github.com/grafana/grafana/pkg/services/query"
"github.com/grafana/grafana/pkg/web"
)
@@ -35,6 +36,32 @@ type queryREST struct {
builder *QueryAPIBuilder
}
type MyCacheService struct {
legacy ds_service.LegacyDataSourceLookup
}
func (mcs *MyCacheService) GetDatasource(ctx context.Context, datasourceID int64, _ identity.Requester, _ bool) (*datasources.DataSource, error) {
ref, err := mcs.legacy.GetDataSourceFromDeprecatedFields(ctx, "", datasourceID)
if err != nil {
return nil, err
}
return &datasources.DataSource{
UID: ref.UID,
Type: ref.Type,
}, nil
}
func (mcs *MyCacheService) GetDatasourceByUID(ctx context.Context, datasourceUID string, _ identity.Requester, _ bool) (*datasources.DataSource, error) {
ref, err := mcs.legacy.GetDataSourceFromDeprecatedFields(ctx, datasourceUID, 0)
if err != nil {
return nil, err
}
return &datasources.DataSource{
UID: ref.UID,
Type: ref.Type,
}, nil
}
var (
_ rest.Storage = (*queryREST)(nil)
_ rest.SingularNameProvider = (*queryREST)(nil)
@@ -147,65 +174,12 @@ func (r *queryREST) Connect(connectCtx context.Context, name string, _ runtime.O
responder.Error(err)
return
}
// Parses the request and splits it into multiple sub queries (if necessary)
req, err := b.parser.parseRequest(ctx, raw)
qdr, err := handleQuery(ctx, *raw, *b, httpreq, *responder)
if err != nil {
var refError ErrorWithRefID
statusCode := http.StatusBadRequest
message := err
refID := "A"
if errors.Is(err, datasources.ErrDataSourceNotFound) {
statusCode = http.StatusNotFound
message = errors.New("datasource not found")
}
if errors.As(err, &refError) {
refID = refError.refId
}
qdr := &query.QueryDataResponse{
QueryDataResponse: backend.QueryDataResponse{
Responses: backend.Responses{
refID: {
Error: message,
Status: backend.Status(statusCode),
},
},
},
}
b.log.Error("Error parsing query", "refId", refID, "message", message)
responder.Object(statusCode, qdr)
return
}
logEmptyRefids(raw.Queries, b.log)
for i := range req.Requests {
req.Requests[i].Headers = ExtractKnownHeaders(httpreq.Header)
}
// Fetch information on the grafana instance (e.g. feature toggles)
instanceConfig, err := b.clientSupplier.GetInstanceConfigurationSettings(ctx)
if err != nil {
msg := "failed to get instance configuration settings"
b.log.Error(msg, "err", err)
responder.Error(errors.New(msg))
return
}
// Actually run the query (includes expressions)
rsp, err := b.execute(ctx, req, instanceConfig)
if err != nil {
// we extract the QDR, but the response may be nil
var qdr *backend.QueryDataResponse
if rsp != nil {
qdr = rsp.QDR
}
b.log.Error("execute error", "http code", query.GetResponseCode(qdr), "err", err)
logEmptyRefids(raw.Queries, b.log)
if qdr != nil { // if we have a response, we assume the err is set in the response
responder.Object(query.GetResponseCode(qdr), &query.QueryDataResponse{
QueryDataResponse: *qdr,
@@ -219,335 +193,78 @@ func (r *queryREST) Connect(connectCtx context.Context, name string, _ runtime.O
}
}
// response headers are communicated back using the context for some reason
reqCtx := contexthandler.FromContext(ctx)
if reqCtx != nil {
mergeHeaders(reqCtx.Resp.Header(), rsp.Headers, b.log)
}
responder.Object(query.GetResponseCode(rsp.QDR), &query.QueryDataResponse{
QueryDataResponse: *rsp.QDR, // wrap the backend response as a QueryDataResponse
responder.Object(query.GetResponseCode(qdr), &query.QueryDataResponse{
QueryDataResponse: *qdr, // wrap the backend response as a QueryDataResponse
})
}), nil
}
func logEmptyRefids(queries []v0alpha1.DataQuery, logger log.Logger) {
emptyCount := 0
for _, q := range queries {
if q.RefID == "" {
emptyCount += 1
}
}
if emptyCount > 0 {
logger.Info("empty refid found", "empty_count", emptyCount, "query_count", len(queries))
}
}
func mergeHeaders(main http.Header, extra http.Header, l log.Logger) {
for headerName, extraValues := range extra {
mainValues := main.Values(headerName)
for _, extraV := range extraValues {
if !slices.Contains(mainValues, extraV) {
main.Add(headerName, extraV)
} else {
l.Warn("skipped duplicate response header", "header", headerName, "value", extraV)
}
}
}
}
func (b *QueryAPIBuilder) execute(ctx context.Context, req parsedRequestInfo, instanceConfig clientapi.InstanceConfigurationSettings) (*clientapi.Response, error) {
var rsp *clientapi.Response
var err error
switch len(req.Requests) {
case 0:
b.log.Debug("executing empty query")
rsp = &clientapi.Response{QDR: &backend.QueryDataResponse{}}
case 1:
b.log.Debug("executing single query")
rsp, err = b.handleQuerySingleDatasource(ctx, req.Requests[0], instanceConfig)
func handleQuery(ctx context.Context, raw query.QueryDataRequest, b QueryAPIBuilder, httpreq *http.Request, responder responderWrapper) (*backend.QueryDataResponse, error) {
var jsonQueries = make([]*simplejson.Json, 0, len(raw.Queries))
for _, query := range raw.Queries {
jsonBytes, err := json.Marshal(query)
if err != nil {
b.log.Debug("handleQuerySingleDatasource failed", err)
b.log.Error("error marshalling", err)
}
if err == nil && isSingleAlertQuery(req) {
b.log.Debug("handling alert query with single query")
rsp.QDR, err = b.convertQueryFromAlerting(ctx, req.Requests[0], rsp.QDR)
if err != nil {
b.log.Debug("convertQueryFromAlerting failed", "err", err)
}
}
default:
b.log.Debug("executing concurrent queries")
rsp, err = b.executeConcurrentQueries(ctx, req.Requests, instanceConfig)
sjQuery, _ := simplejson.NewJson(jsonBytes)
if err != nil {
b.log.Debug("error in executeConcurrentQueries", "err", err)
b.log.Error("error unmarshalling", err)
}
jsonQueries = append(jsonQueries, sjQuery)
}
mReq := dtos.MetricRequest{
From: raw.From,
To: raw.To,
Queries: jsonQueries,
}
cache := &MyCacheService{
legacy: b.legacyDatasourceLookup,
}
headers := ExtractKnownHeaders(httpreq.Header)
instanceConfig, err := b.clientSupplier.GetInstanceConfigurationSettings(ctx)
if err != nil {
b.log.Debug("error in query phase, skipping expressions", "error", err)
return rsp, err //return early here to prevent expressions from being executed if we got an error during the query phase
b.log.Error("failed to get instance configuration settings", "err", err)
responder.Error(err)
return nil, err
}
if len(req.Expressions) > 0 {
b.log.Debug("executing expressions")
rsp.QDR, err = b.handleExpressions(ctx, req, rsp.QDR)
if err != nil {
b.log.Debug("handleExpressions failed", "err", err)
}
}
// Remove hidden results
for _, refId := range req.HideBeforeReturn {
r, ok := rsp.QDR.Responses[refId]
if ok && r.Error == nil {
delete(rsp.QDR.Responses, refId)
}
}
return rsp, err
}
// Process a single request
// See: https://github.com/grafana/grafana/blob/v10.2.3/pkg/services/query/query.go#L242
func (b *QueryAPIBuilder) handleQuerySingleDatasource(ctx context.Context, req datasourceRequest, instanceConfig clientapi.InstanceConfigurationSettings) (*clientapi.Response, error) {
ctx, span := b.tracer.Start(ctx, "Query.handleQuerySingleDatasource")
defer span.End()
span.SetAttributes(
attribute.String("datasource.type", req.PluginId),
attribute.String("datasource.uid", req.UID),
)
allHidden := true
for idx := range req.Request.Queries {
if !req.Request.Queries[idx].Hide {
allHidden = false
break
}
}
if allHidden {
return &clientapi.Response{}, nil
}
client, err := b.clientSupplier.GetDataSourceClient(
mtDsClientBuilder := mtdsclient.NewMtDatasourceClientBuilderWithClientSupplier(
b.clientSupplier,
ctx,
v0alpha1.DataSourceRef{
Type: req.PluginId,
UID: req.UID,
},
req.Headers,
headers,
instanceConfig,
b.log,
)
exprService := expr.ProvideService(
&setting.Cfg{
ExpressionsEnabled: instanceConfig.ExpressionsEnabled,
SQLExpressionCellLimit: instanceConfig.SQLExpressionCellLimit,
SQLExpressionOutputCellLimit: instanceConfig.SQLExpressionOutputCellLimit,
SQLExpressionTimeout: instanceConfig.SQLExpressionTimeout,
},
nil,
nil,
instanceConfig.FeatureToggles,
nil,
b.tracer,
mtDsClientBuilder,
)
qdr, err := service.QueryData(ctx, b.log, cache, exprService, mReq, mtDsClientBuilder, headers)
if err != nil {
b.log.Debug("error getting single datasource client", "error", err, "reqUid", req.UID)
qdr := buildErrorResponse(err, req)
return qdr, err
}
rsp, err := client.QueryData(ctx, *req.Request)
if err == nil && rsp != nil {
for _, q := range req.Request.Queries {
if q.ResultAssertions != nil {
result, ok := rsp.QDR.Responses[q.RefID]
if ok && result.Error == nil {
err = q.ResultAssertions.Validate(result.Frames)
if err != nil {
b.log.Error("Validate failed", "err", err)
result.Error = err
result.ErrorSource = backend.ErrorSourceDownstream
rsp.QDR.Responses[q.RefID] = result
}
}
}
}
}
if err != nil {
b.log.Debug("error in single datasource query", "error", err)
rsp = buildErrorResponse(err, req)
}
return rsp, err
}
// buildErrorResponses applies the provided error to each query response in the list. These queries should all belong to the same datasource.
func buildErrorResponse(err error, req datasourceRequest) *clientapi.Response {
rsp := backend.NewQueryDataResponse()
for _, query := range req.Request.Queries {
rsp.Responses[query.RefID] = backend.DataResponse{
Error: err,
}
}
return &clientapi.Response{QDR: rsp, Headers: nil}
}
// executeConcurrentQueries executes queries to multiple datasources concurrently and returns the aggregate result.
func (b *QueryAPIBuilder) executeConcurrentQueries(ctx context.Context, requests []datasourceRequest, instanceConfig clientapi.InstanceConfigurationSettings) (*clientapi.Response, error) {
ctx, span := b.tracer.Start(ctx, "Query.executeConcurrentQueries")
defer span.End()
g, ctx := errgroup.WithContext(ctx)
g.SetLimit(b.concurrentQueryLimit) // prevent too many concurrent requests
rchan := make(chan *clientapi.Response, len(requests))
// Create panic recovery function for loop below
recoveryFn := func(req datasourceRequest) {
if r := recover(); r != nil {
var err error
b.log.Error("query datasource panic", "error", r, "stack", log.Stack(1))
if theErr, ok := r.(error); ok {
err = theErr
} else if theErrString, ok := r.(string); ok {
err = errors.New(theErrString)
} else {
err = fmt.Errorf("unexpected error - %s", b.userFacingDefaultError)
}
// Due to the panic, there is no valid response for any query for this datasource. Append an error for each one.
rchan <- buildErrorResponse(err, req)
}
}
// Query each datasource concurrently
for idx := range requests {
req := requests[idx]
g.Go(func() error {
defer recoveryFn(req)
rsp, err := b.handleQuerySingleDatasource(ctx, req, instanceConfig)
if err == nil {
rchan <- rsp
} else {
rchan <- buildErrorResponse(err, req)
}
return nil
})
}
if err := g.Wait(); err != nil {
return nil, err
}
close(rchan)
// Merge the results from each response
rsp := &clientapi.Response{
QDR: backend.NewQueryDataResponse(),
Headers: http.Header{},
}
for result := range rchan {
for refId, dataResponse := range result.QDR.Responses {
rsp.QDR.Responses[refId] = dataResponse
}
mergeHeaders(rsp.Headers, result.Headers, b.log)
}
return rsp, nil
}
// Unlike the implementation in expr/node.go, all datasource queries have been processed first
func (b *QueryAPIBuilder) handleExpressions(ctx context.Context, req parsedRequestInfo, data *backend.QueryDataResponse) (qdr *backend.QueryDataResponse, err error) {
start := time.Now()
ctx, span := b.tracer.Start(ctx, "Query.handleExpressions")
traceId := span.SpanContext().TraceID()
expressionsLogger := b.log.New("traceId", traceId.String())
expressionsLogger.Debug("handling expressions")
defer func() {
var respStatus string
switch err {
case nil:
respStatus = "success"
default:
respStatus = "failure"
}
duration := float64(time.Since(start).Nanoseconds()) / float64(time.Millisecond)
b.metrics.ExpressionsQuerySummary.WithLabelValues(respStatus).Observe(duration)
span.End()
}()
qdr = data
if qdr == nil {
qdr = &backend.QueryDataResponse{}
}
if qdr.Responses == nil {
qdr.Responses = make(backend.Responses) // avoid NPE for lookup
}
now := start // <<< this should come from the original query parser
vars := make(mathexp.Vars)
for _, expression := range req.Expressions {
// Setup the variables
for _, refId := range expression.Command.NeedsVars() {
_, ok := vars[refId]
if !ok {
dr, ok := qdr.Responses[refId]
if ok {
_, isSqlInput := req.SqlInputs[refId]
_, res, err := b.converter.Convert(ctx, req.RefIDTypes[refId], dr.Frames, isSqlInput)
if err != nil {
expressionsLogger.Error("error converting frames for expressions", "error", err)
res.Error = err
}
vars[refId] = res
} else {
expressionsLogger.Error("missing variable in handle expressions", "refId", refId, "expressionRefId", expression.RefID)
// This should error in the parsing phase
err := fmt.Errorf("missing variable %s for %s", refId, expression.RefID)
qdr.Responses[refId] = backend.DataResponse{
Error: err,
}
return qdr, err
}
}
}
refId := expression.RefID
results, err := expression.Command.Execute(ctx, now, vars, b.tracer, b.metrics)
if err != nil {
expressionsLogger.Error("error executing expression", "error", err)
results.Error = err
}
qdr.Responses[refId] = backend.DataResponse{
Error: results.Error,
Frames: results.Values.AsDataFrames(refId),
}
}
return qdr, nil
}
func (b *QueryAPIBuilder) convertQueryFromAlerting(ctx context.Context, req datasourceRequest,
qdr *backend.QueryDataResponse) (*backend.QueryDataResponse, error) {
if len(req.Request.Queries) == 0 {
return nil, errors.New("no queries to convert")
}
if qdr == nil {
b.log.Debug("unexpected response of nil from datasource", "datasource.type", req.PluginId, "datasource.uid", req.UID)
return nil, errors.New("unexpected response of nil from datasource")
}
refID := req.Request.Queries[0].RefID
if _, exist := qdr.Responses[refID]; !exist {
return nil, fmt.Errorf("refID '%s' does not exist", refID)
}
frames := qdr.Responses[refID].Frames
_, results, err := b.converter.Convert(ctx, req.PluginId, frames, false)
if err != nil {
b.log.Error("issue converting query from alerting", "err", err)
results.Error = err
}
qdr = &backend.QueryDataResponse{
Responses: map[string]backend.DataResponse{
refID: {
Frames: results.Values.AsDataFrames(refID),
Error: results.Error,
},
},
}
return qdr, err
}
type responderWrapper struct {
wrapped rest.Responder
onObjectFn func(statusCode *int, obj runtime.Object)
@@ -578,15 +295,29 @@ func (r responderWrapper) Error(err error) {
r.wrapped.Error(err)
}
// Checks if the request only contains a single query and is from Alerting
func isSingleAlertQuery(req parsedRequestInfo) bool {
if len(req.Requests) != 1 {
return false
func logEmptyRefids(queries []v0alpha1.DataQuery, logger log.Logger) {
emptyCount := 0
for _, q := range queries {
if q.RefID == "" {
emptyCount += 1
}
}
headers := req.Requests[0].Headers
_, exist := headers[models.FromAlertHeaderName]
if exist && len(req.Requests[0].Request.Queries) == 1 {
return true
if emptyCount > 0 {
logger.Info("empty refid found", "empty_count", emptyCount, "query_count", len(queries))
}
}
func mergeHeaders(main http.Header, extra http.Header, l log.Logger) {
for headerName, extraValues := range extra {
mainValues := main.Values(headerName)
for _, extraV := range extraValues {
if !slices.Contains(mainValues, extraV) {
main.Add(headerName, extraV)
} else {
l.Warn("skipped duplicate response header", "header", headerName, "value", extraV)
}
}
}
return false
}
+240 -117
View File
@@ -3,163 +3,274 @@ package query
import (
"bytes"
"context"
"fmt"
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"testing"
"time"
"github.com/google/go-cmp/cmp"
"github.com/grafana/grafana-plugin-sdk-go/backend"
frameData "github.com/grafana/grafana-plugin-sdk-go/data"
data "github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
"github.com/grafana/grafana-plugin-sdk-go/data"
dataapi "github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
queryapi "github.com/grafana/grafana/pkg/apis/query/v0alpha1"
"github.com/grafana/grafana/pkg/expr"
"github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/infra/tracing"
"github.com/grafana/grafana/pkg/registry/apis/query/clientapi"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/ngalert/models"
"github.com/stretchr/testify/require"
"k8s.io/apimachinery/pkg/runtime"
)
func TestQueryRestConnectHandler(t *testing.T) {
b := &QueryAPIBuilder{
clientSupplier: mockClient{
lastCalledWithHeaders: &map[string]string{},
},
tracer: tracing.InitializeTracerForTest(),
parser: newQueryParser(expr.NewExpressionQueryReader(featuremgmt.WithFeatures()),
&legacyDataSourceRetriever{}, tracing.InitializeTracerForTest(), nil),
log: log.New("test"),
func loadTestdataFrames(t *testing.T, filename string) *backend.QueryDataResponse {
t.Helper()
// Validate filename doesn't contain path traversal
if strings.Contains(filename, "..") || strings.Contains(filename, "/") {
t.Fatalf("Invalid test filename: %s", filename)
}
qr := newQueryREST(b)
ctx := context.Background()
mr := mockResponder{}
handler, err := qr.Connect(ctx, "name", nil, mr)
require.NoError(t, err)
testdataPath := filepath.Join("testdata", filename)
data, err := os.ReadFile(testdataPath) // #nosec G304 -- testdata files in tests
require.NoError(t, err, "Failed to read testdata file: %s", filename)
rr := httptest.NewRecorder()
body := runtime.RawExtension{
Raw: []byte(`{
"queries": [
{
"datasource": {
"type": "prometheus",
"uid": "demo-prometheus"
},
"expr": "sum(go_gc_duration_seconds)",
"range": false,
"instant": true
}
],
"from": "now-1h",
"to": "now"}`),
}
req := httptest.NewRequest(http.MethodGet, "/some-path", bytes.NewReader(body.Raw))
req.Header.Set(models.FromAlertHeaderName, "true")
req.Header.Set(models.CacheSkipHeaderName, "true")
req.Header.Set("X-Rule-Name", "name-1")
req.Header.Set("X-Rule-Uid", "abc")
req.Header.Set("X-Rule-Folder", "folder-1")
req.Header.Set("X-Rule-Source", "grafana-ruler")
req.Header.Set("X-Rule-Type", "type-1")
req.Header.Set("X-Rule-Version", "version-1")
req.Header.Set("X-Grafana-Org-Id", "1")
req.Header.Set("Content-Type", "application/json")
req.Header.Set("some-unexpected-header", "some-value")
handler.ServeHTTP(rr, req)
var result *backend.QueryDataResponse
err = json.Unmarshal(data, &result)
require.NoError(t, err, "Failed to unmarshal testdata file: %s", filename)
require.Equal(t, map[string]string{
models.FromAlertHeaderName: "true",
models.CacheSkipHeaderName: "true",
"X-Rule-Name": "name-1",
"X-Rule-Uid": "abc",
"X-Rule-Folder": "folder-1",
"X-Rule-Source": "grafana-ruler",
"X-Rule-Type": "type-1",
"X-Rule-Version": "version-1",
"X-Grafana-Org-Id": "1",
}, *b.clientSupplier.(mockClient).lastCalledWithHeaders)
return result
}
func TestInstantQueryFromAlerting(t *testing.T) {
builder := &QueryAPIBuilder{
converter: &expr.ResultConverter{
Features: featuremgmt.WithFeatures(),
Tracer: tracing.InitializeTracerForTest(),
func TestQueryAPI(t *testing.T) {
testCases := []struct {
name string
queryJSON string
headers map[string]string
expectedStatus int
testdataFile string
stubbedFrame *data.Frame
}{
{
name: "single prometheus query",
queryJSON: `{
"queries": [
{
"datasource": {
"type": "prometheus",
"uid": "demo-prom"
},
"expr": "1 + 6",
"range": false,
"instant": true,
"refId": "A"
}
],
"from": "now-1h",
"to": "now"
}`,
expectedStatus: http.StatusOK,
testdataFile: "single_prometheus_query.json",
stubbedFrame: data.NewFrame("",
data.NewField("Time", nil, []time.Time{time.Unix(1704067200, 0)}),
data.NewField("Value", nil, []float64{7.0}),
),
},
{
name: "prometheus query with server side expression",
queryJSON: `{
"queries": [
{
"refId": "A",
"datasource": {
"type": "prometheus",
"uid": "demo-prom"
},
"expr": "7",
"range": false,
"instant": true,
"hide": true
},
{
"refId": "B",
"datasource": {
"uid": "__expr__",
"type": "__expr__"
},
"type": "math",
"expression": "$A * 3"
}
],
"from": "now-1h",
"to": "now"
}`,
testdataFile: "prometheus_with_sse.json",
expectedStatus: http.StatusOK,
stubbedFrame: data.NewFrame("",
data.NewField("Value", nil, []float64{7.0}),
),
},
{
name: "prometheus query with sql expression and hidden prom query",
queryJSON: `{
"queries": [
{
"datasource": {
"type": "prometheus",
"uid": "demo-prom"
},
"expr": "1 + 2",
"range": false,
"instant": true,
"refId": "A",
"hidden": true
},
{
"datasource": {
"uid": "__expr__",
"type": "__expr__"
},
"type": "sql",
"expression": "Select Value + 10 from A;",
"refId": "B"
}
],
"from": "now-1h",
"to": "now"
}`,
testdataFile: "prometheus_with_sql_expression.json",
expectedStatus: http.StatusOK,
stubbedFrame: data.NewFrame("",
data.NewField("Time", nil, []time.Time{time.Unix(1704067200, 0)}),
data.NewField("Value", nil, []float64{7.0}),
),
},
}
dq := data.DataQuery{}
dq.RefID = "A"
dr := datasourceRequest{
Headers: map[string]string{
models.FromAlertHeaderName: "true",
},
Request: &data.QueryDataRequest{
Queries: []data.DataQuery{
dq,
},
},
}
fakeFrame := frameData.NewFrame(
"A",
frameData.NewField("Time", nil, []time.Time{time.Now()}),
frameData.NewField("Value", nil, []int64{42}),
)
fakeFrame.Meta = &frameData.FrameMeta{TypeVersion: frameData.FrameTypeVersion{0, 1}, Type: "numeric-multi"}
inputQDR := &backend.QueryDataResponse{
Responses: map[string]backend.DataResponse{
"A": {
Frames: frameData.Frames{
fakeFrame,
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
builder := &QueryAPIBuilder{
converter: &expr.ResultConverter{
Features: featuremgmt.WithFeatures(featuremgmt.FlagSqlExpressions),
Tracer: tracing.InitializeTracerForTest(),
},
},
},
clientSupplier: mockClient{
stubbedFrame: tc.stubbedFrame,
},
tracer: tracing.InitializeTracerForTest(),
log: log.New("test"),
legacyDatasourceLookup: &mockLegacyDataSourceLookup{},
}
req := httptest.NewRequest(http.MethodPost, "/some-path", bytes.NewReader([]byte(tc.queryJSON)))
req.Header.Set("Content-Type", "application/json")
// Set optional headers
for key, value := range tc.headers {
req.Header.Set(key, value)
}
ctx := context.Background()
mr := &mockResponder{}
qr := newQueryREST(builder)
handler, err := qr.Connect(ctx, "name", nil, mr)
require.NoError(t, err)
rr := httptest.NewRecorder()
handler.ServeHTTP(rr, req)
require.NoError(t, mr.err, "Should not have error in responder")
require.Equal(t, tc.expectedStatus, mr.statusCode, "Should return expected status code")
require.NotNil(t, mr.response, "Should have a response object")
// Verify the response is the expected type
qdr, ok := mr.response.(*queryapi.QueryDataResponse)
require.True(t, ok, "Response should be QueryDataResponse type")
require.NotNil(t, qdr.Responses, "Should have responses")
// Load expected frames from testdata if provided
if tc.testdataFile != "" {
expectedResponse := loadTestdataFrames(t, tc.testdataFile)
// get refids from expected response
expectedRefIds := make([]string, 0, len(expectedResponse.Responses))
for refID := range expectedResponse.Responses {
expectedRefIds = append(expectedRefIds, refID)
}
// Verify all expected refIDs are present
for _, refID := range expectedRefIds {
require.Contains(t, qdr.Responses, refID, "Should contain response for refId %s", refID)
actualResponse := qdr.Responses[refID]
expectedFrameResponse := expectedResponse.Responses[refID]
// Verify frame structure matches testdata
require.Len(t, actualResponse.Frames, len(expectedFrameResponse.Frames), "Frame count should match testdata for refId %s", refID)
for i, actualFrame := range actualResponse.Frames {
expectedFrame := expectedFrameResponse.Frames[i]
if diff := cmp.Diff(expectedFrame, actualFrame, data.FrameTestCompareOptions()...); diff != "" {
require.FailNowf(t, "Result mismatch (-want +got):%s", diff)
}
}
}
} else {
t.Fatalf("No testdata file provided for test case %s", tc.name)
}
t.Logf("Test case '%s' completed successfully", tc.name)
})
}
request := parsedRequestInfo{
Requests: []datasourceRequest{
dr,
},
}
result, err := builder.convertQueryFromAlerting(context.Background(), dr, inputQDR)
require.NoError(t, err)
require.True(t, isSingleAlertQuery(request), "Expected a valid alert query with a single query to return true")
require.NotNil(t, result)
require.Equal(t, 1, len(result.Responses["A"].Frames[0].Fields), "Expected a single field not Time and Value")
require.Equal(t, "Value", result.Responses["A"].Frames[0].Fields[0].Name, "Expected the single field to be Value")
}
type mockResponder struct {
statusCode int
response runtime.Object
err error
}
// Object writes the provided object to the response. Invoking this method multiple times is undefined.
func (m mockResponder) Object(statusCode int, obj runtime.Object) {
func (m *mockResponder) Object(statusCode int, obj runtime.Object) {
m.statusCode = statusCode
m.response = obj
}
// Error writes the provided error to the response. This method may only be invoked once.
func (m mockResponder) Error(err error) {
func (m *mockResponder) Error(err error) {
m.err = err
}
type mockClient struct {
lastCalledWithHeaders *map[string]string
stubbedFrame *data.Frame
}
func (m mockClient) GetDataSourceClient(ctx context.Context, ref data.DataSourceRef, headers map[string]string, instanceConfig clientapi.InstanceConfigurationSettings) (clientapi.QueryDataClient, error) {
*m.lastCalledWithHeaders = headers
return nil, fmt.Errorf("mock error")
func (m mockClient) GetDataSourceClient(ctx context.Context, ref dataapi.DataSourceRef, headers map[string]string, instanceConfig clientapi.InstanceConfigurationSettings) (clientapi.QueryDataClient, error) {
mclient := mockClient{
stubbedFrame: m.stubbedFrame,
}
return mclient, nil
}
func (m mockClient) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
return nil, fmt.Errorf("mock error")
func (m mockClient) QueryData(ctx context.Context, req dataapi.QueryDataRequest) (*backend.QueryDataResponse, error) {
responses := make(backend.Responses)
for i := range req.Queries {
refID := req.Queries[i].RefID
frame := m.stubbedFrame
frame.RefID = refID
responses[refID] = backend.DataResponse{
Status: backend.StatusOK,
Frames: []*data.Frame{frame},
}
}
return &backend.QueryDataResponse{
Responses: responses,
}, nil
}
func (m mockClient) CallResource(ctx context.Context, req *backend.CallResourceRequest, sender backend.CallResourceResponseSender) error {
@@ -170,8 +281,20 @@ func (m mockClient) CheckHealth(ctx context.Context, req *backend.CheckHealthReq
return nil, nil
}
func (m mockClient) GetInstanceConfigurationSettings(_ context.Context) (clientapi.InstanceConfigurationSettings, error) {
return clientapi.InstanceConfigurationSettings{}, nil
func (m mockClient) GetInstanceConfigurationSettings(ctx context.Context) (clientapi.InstanceConfigurationSettings, error) {
return clientapi.InstanceConfigurationSettings{
ExpressionsEnabled: true,
FeatureToggles: featuremgmt.WithFeatures(featuremgmt.FlagSqlExpressions),
}, nil
}
type mockLegacyDataSourceLookup struct{}
func (m *mockLegacyDataSourceLookup) GetDataSourceFromDeprecatedFields(ctx context.Context, name string, id int64) (*dataapi.DataSourceRef, error) {
return &dataapi.DataSourceRef{
UID: "demo-prom",
Type: "prometheus",
}, nil
}
func TestMergeHeaders(t *testing.T) {
+19 -18
View File
@@ -36,32 +36,30 @@ import (
var _ builder.APIGroupBuilder = (*QueryAPIBuilder)(nil)
type QueryAPIBuilder struct {
log log.Logger
concurrentQueryLimit int
userFacingDefaultError string
features featuremgmt.FeatureToggles
log log.Logger
concurrentQueryLimit int
features featuremgmt.FeatureToggles
authorizer authorizer.Authorizer
tracer tracing.Tracer
metrics *metrics.ExprMetrics
parser *queryParser
clientSupplier clientapi.DataSourceClientSupplier
registry query.DataSourceApiServerRegistry
converter *expr.ResultConverter
queryTypes *query.QueryTypeDefinitionList
tracer tracing.Tracer
metrics *metrics.ExprMetrics
clientSupplier clientapi.DataSourceClientSupplier
registry query.DataSourceApiServerRegistry
converter *expr.ResultConverter
queryTypes *query.QueryTypeDefinitionList
legacyDatasourceLookup service.LegacyDataSourceLookup
}
func NewQueryAPIBuilder(features featuremgmt.FeatureToggles,
func NewQueryAPIBuilder(
features featuremgmt.FeatureToggles,
clientSupplier clientapi.DataSourceClientSupplier,
ar authorizer.Authorizer,
registry query.DataSourceApiServerRegistry,
legacy service.LegacyDataSourceLookup,
registerer prometheus.Registerer,
tracer tracing.Tracer,
legacyDatasourceLookup service.LegacyDataSourceLookup,
) (*QueryAPIBuilder, error) {
reader := expr.NewExpressionQueryReader(features)
// Include well typed query definitions
var queryTypes *query.QueryTypeDefinitionList
if features.IsEnabledGlobally(featuremgmt.FlagDatasourceQueryTypes) {
@@ -83,7 +81,6 @@ func NewQueryAPIBuilder(features featuremgmt.FeatureToggles,
clientSupplier: clientSupplier,
authorizer: ar,
registry: registry,
parser: newQueryParser(reader, legacy, tracer, log.New("query_parser")),
metrics: metrics.NewQueryServiceExpressionsMetrics(registerer),
tracer: tracer,
features: features,
@@ -92,6 +89,7 @@ func NewQueryAPIBuilder(features featuremgmt.FeatureToggles,
Features: features,
Tracer: tracer,
},
legacyDatasourceLookup: legacyDatasourceLookup,
}, nil
}
@@ -104,7 +102,8 @@ func RegisterAPIService(features featuremgmt.FeatureToggles,
pCtxProvider *plugincontext.Provider,
registerer prometheus.Registerer,
tracer tracing.Tracer,
legacy service.LegacyDataSourceLookup,
legacyDatasourceLookup service.LegacyDataSourceLookup,
exprService *expr.Service,
) (*QueryAPIBuilder, error) {
if !featuremgmt.AnyEnabled(features,
featuremgmt.FlagQueryService,
@@ -132,7 +131,9 @@ func RegisterAPIService(features featuremgmt.FeatureToggles,
},
ar,
client.NewDataSourceRegistryFromStore(pluginStore, dataSourcesService),
legacy, registerer, tracer,
registerer,
tracer,
legacyDatasourceLookup,
)
apiregistration.RegisterAPI(builder)
return builder, err
-29
View File
@@ -1,29 +0,0 @@
{
"description": "self dependencies",
"input": {
"from": "now-6",
"to": "now",
"queries": [
{
"refId": "A",
"datasource": {
"type": "",
"uid": "__expr__"
},
"expression": "$B",
"type": "math"
},
{
"refId": "B",
"datasource": {
"type": "",
"uid": "__expr__"
},
"type": "math",
"expression": "$A"
}
]
},
"expect": {},
"error": "[sse.cyclic] cyclic reference in expression []"
}
@@ -1,60 +0,0 @@
{
"input": {
"from": "now-6",
"to": "now",
"queries": [
{
"refId": "A",
"datasource": {
"type": "plugin-x",
"uid": "123"
}
},
{
"refId": "B",
"datasource": {
"type": "plugin-x",
"uid": "456"
}
}
]
},
"expect": {
"requests": [
{
"pluginId": "plugin-x",
"uid": "123",
"request": {
"from": "now-6",
"to": "now",
"queries": [
{
"refId": "A",
"datasource": {
"type": "plugin-x",
"uid": "123"
}
}
]
}
},
{
"pluginId": "plugin-x",
"uid": "456",
"request": {
"from": "now-6",
"to": "now",
"queries": [
{
"refId": "B",
"datasource": {
"type": "plugin-x",
"uid": "456"
}
}
]
}
}
]
}
}
@@ -0,0 +1,32 @@
{
"results": {
"B": {
"status": 200,
"frames": [
{
"schema": {
"name": "B",
"refId": "B",
"fields": [
{
"name": "Value + 10",
"type": "number",
"typeInfo": {
"frame": "float64",
"nullable": false
}
}
]
},
"data": {
"values": [
[
17
]
]
}
}
]
}
}
}
@@ -0,0 +1,39 @@
{
"results": {
"B": {
"status": 200,
"frames": [
{
"schema": {
"refId": "B",
"fields": [
{
"name": "B",
"type": "number",
"typeInfo": {
"frame": "float64",
"nullable": true
}
}
],
"meta": {
"type": "numeric-multi",
"typeVersion": [
0,
1
]
}
},
"data": {
"values": [
[
21
]
]
}
}
]
}
}
}
-20
View File
@@ -1,20 +0,0 @@
{
"description": "self dependencies",
"input": {
"from": "now-6",
"to": "now",
"queries": [
{
"refId": "A",
"datasource": {
"type": "",
"uid": "__expr__"
},
"type": "math",
"expression": "$A"
}
]
},
"expect": {},
"error": "[sse.cyclic] cyclic reference in expression [A]"
}
@@ -0,0 +1,41 @@
{
"results": {
"A": {
"status": 200,
"frames": [
{
"schema": {
"refId": "A",
"fields": [
{
"name": "Time",
"type": "time",
"typeInfo": {
"frame": "time.Time"
}
},
{
"name": "Value",
"type": "number",
"typeInfo": {
"frame": "float64"
},
"labels": {}
}
]
},
"data": {
"values": [
[
1704067200000
],
[
7
]
]
}
}
]
}
}
}
-79
View File
@@ -1,79 +0,0 @@
{
"description": "one hidden query with two expressions that start out-of-order",
"input": {
"from": "now-6",
"to": "now",
"queries": [
{
"refId": "C",
"datasource": {
"type": "",
"uid": "__expr__"
},
"type": "reduce",
"expression": "$B",
"reducer": "last"
},
{
"refId": "A",
"datasource": {
"type": "sql",
"uid": "123"
},
"hide": true
},
{
"refId": "B",
"datasource": {
"type": "",
"uid": "-100"
},
"type": "math",
"expression": "$A + 10"
}
]
},
"expect": {
"requests": [
{
"pluginId": "sql",
"uid": "123",
"request": {
"from": "now-6",
"to": "now",
"queries": [
{
"refId": "A",
"datasource": {
"type": "sql",
"uid": "123"
}
}
]
}
}
],
"expressions": [
{
"id": 2,
"refId": "B",
"type": "math",
"properties": {
"expression": "$A + 10"
}
},
{
"id": 1,
"refId": "C",
"type": "reduce",
"properties": {
"expression": "$B",
"reducer": "last"
}
}
],
"hide": [
"A"
]
}
}
+2
View File
@@ -101,6 +101,7 @@ import (
"github.com/grafana/grafana/pkg/services/login/authinfoimpl"
"github.com/grafana/grafana/pkg/services/loginattempt"
"github.com/grafana/grafana/pkg/services/loginattempt/loginattemptimpl"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/navtree/navtreeimpl"
"github.com/grafana/grafana/pkg/services/ngalert"
ngimage "github.com/grafana/grafana/pkg/services/ngalert/image"
@@ -323,6 +324,7 @@ var wireBasicSet = wire.NewSet(
serviceaccountsmanager.ProvideServiceAccountsService,
serviceaccountsproxy.ProvideServiceAccountsProxy,
wire.Bind(new(serviceaccounts.Service), new(*serviceaccountsproxy.ServiceAccountsProxy)),
mtdsclient.NewNullMTDatasourceClientBuilder,
expr.ProvideService,
featuremgmt.ProvideManagerService,
featuremgmt.ProvideToggles,
+10 -7
View File
File diff suppressed because one or more lines are too long
@@ -0,0 +1,66 @@
package mtdsclient
import (
"context"
"github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
"github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/registry/apis/query/clientapi"
)
type MTDatasourceClientBuilder interface {
BuildClient(pluginId string, uid string) (clientapi.QueryDataClient, bool)
}
type nullBuilder struct{}
func (m *nullBuilder) BuildClient(pluginId string, uid string) (clientapi.QueryDataClient, bool) {
return nil, false
}
// we use this noop for st flows
func NewNullMTDatasourceClientBuilder() MTDatasourceClientBuilder {
return &nullBuilder{}
}
type MtDatasourceClientBuilderWithClientSupplier struct {
clientSupplier clientapi.DataSourceClientSupplier
ctx context.Context
headers map[string]string
instanceConfig clientapi.InstanceConfigurationSettings
logger log.Logger
}
func (b *MtDatasourceClientBuilderWithClientSupplier) BuildClient(pluginId string, uid string) (clientapi.QueryDataClient, bool) {
dsClient, err := b.clientSupplier.GetDataSourceClient(
b.ctx,
v0alpha1.DataSourceRef{
Type: pluginId,
UID: uid,
},
b.headers,
b.instanceConfig,
)
if err != nil {
b.logger.Debug("failed to get mt ds client", "error", err)
return nil, false
}
return dsClient, true
}
// TODO: I think we might be able to refactor this to just use the client supplier directly
func NewMtDatasourceClientBuilderWithClientSupplier(
clientSupplier clientapi.DataSourceClientSupplier,
ctx context.Context,
headers map[string]string,
instanceConfig clientapi.InstanceConfigurationSettings,
logger log.Logger,
) MTDatasourceClientBuilder {
return &MtDatasourceClientBuilderWithClientSupplier{
clientSupplier: clientSupplier,
ctx: ctx,
headers: headers,
instanceConfig: instanceConfig,
logger: logger,
}
}
+23 -2
View File
@@ -20,6 +20,7 @@ import (
"github.com/grafana/grafana/pkg/services/datasources"
fakes "github.com/grafana/grafana/pkg/services/datasources/fakes"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/ngalert/models"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore"
"github.com/grafana/grafana/pkg/services/user"
@@ -591,7 +592,15 @@ func TestValidate(t *testing.T) {
pluginsStore: store,
})
expressions := expr.ProvideService(&setting.Cfg{ExpressionsEnabled: true}, nil, nil, featuremgmt.WithFeatures(), nil, tracing.InitializeTracerForTest())
expressions := expr.ProvideService(
&setting.Cfg{ExpressionsEnabled: true},
nil,
nil,
featuremgmt.WithFeatures(),
nil,
tracing.InitializeTracerForTest(),
mtdsclient.NewNullMTDatasourceClientBuilder(),
)
validator := NewConditionValidator(cacheService, expressions, store)
evalCtx := NewContext(context.Background(), u)
@@ -710,7 +719,19 @@ func TestCreate_HysteresisCommand(t *testing.T) {
cache: cacheService,
pluginsStore: store,
})
evaluator := NewEvaluatorFactory(setting.UnifiedAlertingSettings{}, cacheService, expr.ProvideService(&setting.Cfg{ExpressionsEnabled: true}, nil, nil, featuremgmt.WithFeatures(), nil, tracing.InitializeTracerForTest()))
evaluator := NewEvaluatorFactory(
setting.UnifiedAlertingSettings{},
cacheService,
expr.ProvideService(
&setting.Cfg{ExpressionsEnabled: true},
nil,
nil,
featuremgmt.WithFeatures(),
nil,
tracing.InitializeTracerForTest(),
mtdsclient.NewNullMTDatasourceClientBuilder(),
),
)
evalCtx := NewContextWithPreviousResults(context.Background(), u, testCase.reader)
eval, err := evaluator.Create(evalCtx, condition)
@@ -26,6 +26,7 @@ import (
"github.com/grafana/grafana/pkg/infra/tracing"
datasources "github.com/grafana/grafana/pkg/services/datasources/fakes"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/ngalert/eval"
"github.com/grafana/grafana/pkg/services/ngalert/metrics"
"github.com/grafana/grafana/pkg/services/ngalert/models"
@@ -66,7 +67,19 @@ func TestProcessTicks(t *testing.T) {
}
cacheServ := &datasources.FakeCacheService{}
evaluator := eval.NewEvaluatorFactory(setting.UnifiedAlertingSettings{}, cacheServ, expr.ProvideService(&setting.Cfg{ExpressionsEnabled: true}, nil, nil, featuremgmt.WithFeatures(), nil, tracing.InitializeTracerForTest()))
evaluator := eval.NewEvaluatorFactory(
setting.UnifiedAlertingSettings{},
cacheServ,
expr.ProvideService(
&setting.Cfg{ExpressionsEnabled: true},
nil,
nil,
featuremgmt.WithFeatures(),
nil,
tracing.InitializeTracerForTest(),
mtdsclient.NewNullMTDatasourceClientBuilder(),
),
)
rrSet := setting.RecordingRuleSettings{
Enabled: true,
}
@@ -1192,7 +1205,19 @@ func setupScheduler(t *testing.T, rs *fakeRulesStore, is *state.FakeInstanceStor
var evaluator = evalMock
if evalMock == nil {
evaluator = eval.NewEvaluatorFactory(setting.UnifiedAlertingSettings{}, &datasources.FakeCacheService{}, expr.ProvideService(&setting.Cfg{ExpressionsEnabled: true}, nil, nil, featuremgmt.WithFeatures(), nil, tracing.InitializeTracerForTest()))
evaluator = eval.NewEvaluatorFactory(
setting.UnifiedAlertingSettings{},
&datasources.FakeCacheService{},
expr.ProvideService(
&setting.Cfg{ExpressionsEnabled: true},
nil,
nil,
featuremgmt.WithFeatures(),
nil,
tracing.InitializeTracerForTest(),
mtdsclient.NewNullMTDatasourceClientBuilder(),
),
)
}
if registry == nil {
@@ -26,6 +26,7 @@ import (
datasourceService "github.com/grafana/grafana/pkg/services/datasources/service"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/licensing/licensingtest"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginconfig"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
pluginSettings "github.com/grafana/grafana/pkg/services/pluginsintegration/pluginsettings/service"
@@ -155,6 +156,7 @@ func buildQueryDataService(t *testing.T, cs datasources.CacheService, fpc *fakeP
&fakeDataSourceRequestValidator{},
fpc,
pCtxProvider,
mtdsclient.NewNullMTDatasourceClientBuilder(),
)
}
@@ -196,7 +196,7 @@ func TestAPIQueryPublicDashboard(t *testing.T) {
}
}`
setup := func(enabled bool) (*web.Mux, *publicdashboards.FakePublicDashboardService) {
setup := func(_ bool) (*web.Mux, *publicdashboards.FakePublicDashboardService) {
service := publicdashboards.NewFakePublicDashboardService(t)
testServer := setupTestServer(t, nil, service, anonymousUser)
+44 -9
View File
@@ -2,6 +2,7 @@ package query
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
@@ -13,6 +14,7 @@ import (
"github.com/grafana/grafana-plugin-sdk-go/backend/gtime"
"golang.org/x/sync/errgroup"
data "github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
"github.com/grafana/grafana/pkg/api/dtos"
"github.com/grafana/grafana/pkg/apimachinery/errutil"
"github.com/grafana/grafana/pkg/apimachinery/identity"
@@ -22,6 +24,7 @@ import (
"github.com/grafana/grafana/pkg/plugins"
"github.com/grafana/grafana/pkg/services/contexthandler"
"github.com/grafana/grafana/pkg/services/datasources"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
"github.com/grafana/grafana/pkg/services/validations"
"github.com/grafana/grafana/pkg/setting"
@@ -47,6 +50,7 @@ func ProvideService(
dataSourceRequestValidator validations.DataSourceRequestValidator,
pluginClient plugins.Client,
pCtxProvider *plugincontext.Provider,
mtDatasourceClientBuilder mtdsclient.MTDatasourceClientBuilder,
) *ServiceImpl {
g := &ServiceImpl{
cfg: cfg,
@@ -57,6 +61,7 @@ func ProvideService(
pCtxProvider: pCtxProvider,
log: log.New("query_data"),
concurrentQueryLimit: cfg.SectionWithEnvOverrides("query").Key("concurrent_query_limit").MustInt(runtime.NumCPU()),
mtDatasourceClientBuilder: mtDatasourceClientBuilder,
}
g.log.Info("Query Service initialization")
return g
@@ -80,6 +85,8 @@ type ServiceImpl struct {
pCtxProvider *plugincontext.Provider
log log.Logger
concurrentQueryLimit int
mtDatasourceClientBuilder mtdsclient.MTDatasourceClientBuilder
headers map[string]string
}
// Run ServiceImpl.
@@ -201,7 +208,19 @@ func buildErrorResponses(err error, queries []*simplejson.Json) splitResponse {
return splitResponse{er, http.Header{}}
}
// handleExpressions handles POST /api/ds/query when there is an expression.
func QueryData(ctx context.Context, log log.Logger, dscache datasources.CacheService, exprService *expr.Service, reqDTO dtos.MetricRequest, mtDatasourceClientBuilder mtdsclient.MTDatasourceClientBuilder, headers map[string]string) (*backend.QueryDataResponse, error) {
s := &ServiceImpl{
log: log,
dataSourceCache: dscache,
expressionService: exprService,
dataSourceRequestValidator: validations.ProvideValidator(),
mtDatasourceClientBuilder: mtDatasourceClientBuilder,
headers: headers,
}
return s.QueryData(ctx, nil, false, reqDTO)
}
// handleExpressions handles queries when there is an expression.
func (s *ServiceImpl) handleExpressions(ctx context.Context, user identity.Requester, parsedReq *parsedRequest) (*backend.QueryDataResponse, error) {
exprReq := expr.Request{
Queries: []expr.Query{},
@@ -257,21 +276,37 @@ func (s *ServiceImpl) handleQuerySingleDatasource(ctx context.Context, user iden
}
}
pCtx, err := s.pCtxProvider.GetWithDataSource(ctx, ds.Type, user, ds)
if err != nil {
return nil, err
}
req := &backend.QueryDataRequest{
PluginContext: pCtx,
Headers: map[string]string{},
Queries: []backend.DataQuery{},
Headers: map[string]string{},
Queries: []backend.DataQuery{},
}
for _, q := range queries {
req.Queries = append(req.Queries, q.query)
}
return s.pluginClient.QueryData(ctx, req)
mtDsClient, ok := s.mtDatasourceClientBuilder.BuildClient(ds.Type, ds.UID)
if !ok { // single tenant flow
pCtx, err := s.pCtxProvider.GetWithDataSource(ctx, ds.Type, user, ds)
if err != nil {
return nil, err
}
req.PluginContext = pCtx
return s.pluginClient.QueryData(ctx, req)
} else { // multi tenant flow
// transform request from backend.QueryDataRequest to k8s request
k8sReq := &data.QueryDataRequest{}
for _, q := range req.Queries {
var dataQuery data.DataQuery
err := json.Unmarshal(q.JSON, &dataQuery)
if err != nil {
return nil, err
}
k8sReq.Queries = append(k8sReq.Queries, dataQuery)
}
return mtDsClient.QueryData(ctx, *k8sReq)
}
}
// parseRequest parses a request into parsed queries grouped by datasource uid
+3 -2
View File
@@ -30,6 +30,7 @@ import (
"github.com/grafana/grafana/pkg/services/datasources"
fakeDatasources "github.com/grafana/grafana/pkg/services/datasources/fakes"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/services/mtdsclient"
"github.com/grafana/grafana/pkg/services/pluginsintegration/pluginconfig"
"github.com/grafana/grafana/pkg/services/pluginsintegration/plugincontext"
pluginSettings "github.com/grafana/grafana/pkg/services/pluginsintegration/pluginsettings/service"
@@ -483,8 +484,8 @@ func setup(t *testing.T) *testContext {
pluginSettings.ProvideService(sqlStore, secretsService), pluginconfig.NewFakePluginRequestConfigProvider(),
)
exprService := expr.ProvideService(&setting.Cfg{ExpressionsEnabled: true}, pc, pCtxProvider,
featuremgmt.WithFeatures(), nil, tracing.InitializeTracerForTest())
queryService := ProvideService(setting.NewCfg(), dc, exprService, rv, pc, pCtxProvider) // provider belonging to this package
featuremgmt.WithFeatures(), nil, tracing.InitializeTracerForTest(), mtdsclient.NewNullMTDatasourceClientBuilder())
queryService := ProvideService(setting.NewCfg(), dc, exprService, rv, pc, pCtxProvider, mtdsclient.NewNullMTDatasourceClientBuilder()) // provider belonging to this package
return &testContext{
pluginContext: pc,
secretStore: ss,
-167
View File
@@ -1,167 +0,0 @@
package dashboards
import (
"context"
"encoding/json"
"net/http"
"testing"
"github.com/grafana/grafana-plugin-sdk-go/backend"
data "github.com/grafana/grafana-plugin-sdk-go/experimental/apis/data/v0alpha1"
"github.com/stretchr/testify/require"
"k8s.io/apimachinery/pkg/runtime/schema"
"github.com/grafana/grafana/pkg/services/datasources"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/tests/apis"
"github.com/grafana/grafana/pkg/tests/testinfra"
"github.com/grafana/grafana/pkg/tests/testsuite"
)
func TestMain(m *testing.M) {
testsuite.Run(m)
}
func TestIntegrationSimpleQuery(t *testing.T) {
if testing.Short() {
t.Skip("skipping integration test")
}
helper := apis.NewK8sTestHelper(t, testinfra.GrafanaOpts{
AppModeProduction: false, // dev mode required for datasource connections
EnableFeatureToggles: []string{
featuremgmt.FlagGrafanaAPIServerWithExperimentalAPIs, // Required to start the example service
},
})
// Create a single datasource
ds := helper.CreateDS(&datasources.AddDataSourceCommand{
Name: "test",
Type: datasources.DS_TESTDATA,
UID: "test",
OrgID: int64(1),
})
require.Equal(t, "test", ds.UID)
t.Run("Call query with expression", func(t *testing.T) {
client := helper.Org1.Admin.RESTClient(t, &schema.GroupVersion{
Group: "query.grafana.app",
Version: "v0alpha1",
})
body, err := json.Marshal(&data.QueryDataRequest{
Queries: []data.DataQuery{
data.NewDataQuery(map[string]any{
"refId": "X",
"datasource": data.DataSourceRef{
Type: "grafana-testdata-datasource",
UID: ds.UID,
},
"scenarioId": "csv_content",
"csvContent": "a\n1",
}),
data.NewDataQuery(map[string]any{
"refId": "Y",
"datasource": data.DataSourceRef{
UID: "__expr__",
},
"type": "math",
"expression": "$X + 2",
}),
},
})
//t.Logf("%s", string(body))
require.NoError(t, err)
result := client.Post().
Namespace("default").
Suffix("query").
SetHeader("Content-type", "application/json").
Body(body).
Do(context.Background())
require.NoError(t, result.Error())
contentType := "?"
result.ContentType(&contentType)
require.Equal(t, "application/json", contentType)
body, err = result.Raw()
require.NoError(t, err)
t.Logf("OUT: %s", string(body))
rsp := &backend.QueryDataResponse{}
err = json.Unmarshal(body, rsp)
require.NoError(t, err)
require.Equal(t, 2, len(rsp.Responses))
frameX := rsp.Responses["X"].Frames[0]
frameY := rsp.Responses["Y"].Frames[0]
vX, _ := frameX.Fields[0].ConcreteAt(0)
vY, _ := frameY.Fields[0].ConcreteAt(0)
require.Equal(t, int64(1), vX)
require.Equal(t, float64(3), vY) // 1 + 2, but always float64
})
t.Run("Gets an error with invalid queries", func(t *testing.T) {
client := helper.Org1.Admin.RESTClient(t, &schema.GroupVersion{
Group: "query.grafana.app",
Version: "v0alpha1",
})
body, err := json.Marshal(&data.QueryDataRequest{
Queries: []data.DataQuery{
data.NewDataQuery(map[string]any{
"refId": "Y",
"datasource": data.DataSourceRef{
UID: "__expr__",
},
"type": "math",
"expression": "$X + 2", // invalid X does not exit
}),
},
})
require.NoError(t, err)
result := client.Post().
Namespace("default").
Suffix("query").
SetHeader("Content-type", "application/json").
Body(body).
Do(context.Background())
body, err = result.Raw()
//t.Logf("OUT: %s", string(body))
require.Error(t, err, "expecting a 400")
require.JSONEq(t, `{
"results": {
"A": {
"error": "[sse.dependencyError] did not execute expression [Y] due to a failure of the dependent expression or query [X]",
"status": 400,
"errorSource": ""
}
}
}`, string(body))
// require.JSONEq(t, `{
// "status": "Failure",
// "metadata": {},
// "message": "did not execute expression [Y] due to a failure of the dependent expression or query [X]",
// "reason": "BadRequest",
// "details": { "group": "query.grafana.app" },
// "code": 400,
// "messageId": "sse.dependencyError",
// "extra": { "depRefId": "X", "refId": "Y" }
// }`, string(body))
statusCode := -1
contentType := "?"
result.ContentType(&contentType)
result.StatusCode(&statusCode)
require.Equal(t, "application/json", contentType)
require.Equal(t, http.StatusBadRequest, statusCode)
})
}