Prometheus: Remove cache, pass headers in request, simplify client creation for resource calls and custom client (#51436)

* Remove cache, pass headers in request, simplify client creation

* Add test for http options creation
This commit is contained in:
Andrej Ocenas
2022-07-04 11:18:45 +02:00
committed by GitHub
parent b7e22c37a8
commit 3df34fe064
17 changed files with 244 additions and 924 deletions
@@ -118,7 +118,10 @@ func loadStoredQuery(fileName string) (*backend.QueryDataRequest, error) {
}
func runQuery(response []byte, q *backend.QueryDataRequest, wide bool) (*backend.QueryDataResponse, error) {
tCtx := setup(wide)
tCtx, err := setup(wide)
if err != nil {
return nil, err
}
res := &http.Response{
StatusCode: 200,
Body: ioutil.NopCloser(bytes.NewReader(response)),
@@ -22,7 +22,9 @@ import (
// - go tool pprof -http=localhost:6061 memprofile.out
func BenchmarkJson(b *testing.B) {
body, q := createJsonTestData(1642000000, 1, 300, 400)
tCtx := setup(true)
tCtx, err := setup(true)
require.NoError(b, err)
b.ResetTimer()
for n := 0; n < b.N; n++ {
res := http.Response{
+39 -42
View File
@@ -2,21 +2,19 @@ package querydata
import (
"context"
"encoding/json"
"fmt"
"net/http"
"regexp"
"time"
"github.com/grafana/grafana-plugin-sdk-go/backend"
"github.com/grafana/grafana-plugin-sdk-go/data"
"github.com/grafana/grafana/pkg/infra/httpclient"
"github.com/grafana/grafana/pkg/infra/log"
"github.com/grafana/grafana/pkg/infra/tracing"
"github.com/grafana/grafana/pkg/services/featuremgmt"
"github.com/grafana/grafana/pkg/setting"
"github.com/grafana/grafana/pkg/tsdb/intervalv2"
"github.com/grafana/grafana/pkg/tsdb/prometheus/client"
"github.com/grafana/grafana/pkg/tsdb/prometheus/models"
"github.com/grafana/grafana/pkg/tsdb/prometheus/utils"
"github.com/grafana/grafana/pkg/util/maputil"
"go.opentelemetry.io/otel/attribute"
)
@@ -25,18 +23,18 @@ const legendFormatAuto = "__auto"
var legendFormatRegexp = regexp.MustCompile(`\{\{\s*(.+?)\s*\}\}`)
type clientGetter func(map[string]string) (*client.Client, error)
type ExemplarEvent struct {
Time time.Time
Value float64
Labels map[string]string
}
// QueryData handles querying but different from buffered package uses a custom client instead of default Go Prom
// client.
type QueryData struct {
intervalCalculator intervalv2.Calculator
tracer tracing.Tracer
getClient clientGetter
client *client.Client
log log.Logger
ID int64
URL string
@@ -45,34 +43,30 @@ type QueryData struct {
}
func New(
httpClientProvider httpclient.Provider,
cfg *setting.Cfg,
httpClient *http.Client,
features featuremgmt.FeatureToggles,
tracer tracing.Tracer,
settings backend.DataSourceInstanceSettings,
plog log.Logger,
) (*QueryData, error) {
var jsonData map[string]interface{}
if err := json.Unmarshal(settings.JSONData, &jsonData); err != nil {
return nil, fmt.Errorf("error reading settings: %w", err)
jsonData, err := utils.GetJsonData(settings)
if err != nil {
return nil, err
}
httpMethod, _ := maputil.GetStringOptional(jsonData, "httpMethod")
timeInterval, err := maputil.GetStringOptional(jsonData, "timeInterval")
if err != nil {
return nil, err
}
p := client.NewProvider(settings, jsonData, httpClientProvider, cfg, features, plog)
pc, err := client.NewProviderCache(p)
if err != nil {
return nil, err
}
promClient := client.NewClient(httpClient, httpMethod, settings.URL)
return &QueryData{
intervalCalculator: intervalv2.NewCalculator(),
tracer: tracer,
log: plog,
getClient: pc.GetClient,
client: promClient,
TimeInterval: timeInterval,
ID: settings.ID,
URL: settings.URL,
@@ -86,17 +80,12 @@ func (s *QueryData) Execute(ctx context.Context, req *backend.QueryDataRequest)
Responses: backend.Responses{},
}
client, err := s.getClient(req.Headers)
if err != nil {
return &result, err
}
for _, q := range req.Queries {
query, err := models.Parse(q, s.TimeInterval, s.intervalCalculator, fromAlert)
if err != nil {
return &result, err
}
r, err := s.fetch(ctx, client, query)
r, err := s.fetch(ctx, s.client, query, req.Headers)
if err != nil {
return &result, err
}
@@ -110,11 +99,11 @@ func (s *QueryData) Execute(ctx context.Context, req *backend.QueryDataRequest)
return &result, nil
}
func (s *QueryData) fetch(ctx context.Context, client *client.Client, q *models.Query) (*backend.DataResponse, error) {
func (s *QueryData) fetch(ctx context.Context, client *client.Client, q *models.Query, headers map[string]string) (*backend.DataResponse, error) {
s.log.Debug("Sending query", "start", q.Start, "end", q.End, "step", q.Step, "query", q.Expr)
traceCtx, span := s.trace(ctx, q)
defer span.End()
traceCtx, end := s.trace(ctx, q)
defer end()
response := &backend.DataResponse{
Frames: data.Frames{},
@@ -122,7 +111,7 @@ func (s *QueryData) fetch(ctx context.Context, client *client.Client, q *models.
}
if q.RangeQuery {
res, err := s.rangeQuery(traceCtx, client, q)
res, err := s.rangeQuery(traceCtx, client, q, headers)
if err != nil {
return nil, err
}
@@ -130,7 +119,7 @@ func (s *QueryData) fetch(ctx context.Context, client *client.Client, q *models.
}
if q.InstantQuery {
res, err := s.instantQuery(traceCtx, client, q)
res, err := s.instantQuery(traceCtx, client, q, headers)
if err != nil {
return nil, err
}
@@ -138,7 +127,7 @@ func (s *QueryData) fetch(ctx context.Context, client *client.Client, q *models.
}
if q.ExemplarQuery {
res, err := s.exemplarQuery(traceCtx, client, q)
res, err := s.exemplarQuery(traceCtx, client, q, headers)
if err != nil {
// If exemplar query returns error, we want to only log it and
// continue with other results processing
@@ -152,34 +141,42 @@ func (s *QueryData) fetch(ctx context.Context, client *client.Client, q *models.
return response, nil
}
func (s *QueryData) rangeQuery(ctx context.Context, c *client.Client, q *models.Query) (*backend.DataResponse, error) {
res, err := c.QueryRange(ctx, q)
func (s *QueryData) rangeQuery(ctx context.Context, c *client.Client, q *models.Query, headers map[string]string) (*backend.DataResponse, error) {
res, err := c.QueryRange(ctx, q, sdkHeaderToHttpHeader(headers))
if err != nil {
return nil, err
}
return s.parseResponse(ctx, q, res)
}
func (s *QueryData) instantQuery(ctx context.Context, c *client.Client, q *models.Query) (*backend.DataResponse, error) {
res, err := c.QueryInstant(ctx, q)
func (s *QueryData) instantQuery(ctx context.Context, c *client.Client, q *models.Query, headers map[string]string) (*backend.DataResponse, error) {
res, err := c.QueryInstant(ctx, q, sdkHeaderToHttpHeader(headers))
if err != nil {
return nil, err
}
return s.parseResponse(ctx, q, res)
}
func (s *QueryData) exemplarQuery(ctx context.Context, c *client.Client, q *models.Query) (*backend.DataResponse, error) {
res, err := c.QueryExemplars(ctx, q)
func (s *QueryData) exemplarQuery(ctx context.Context, c *client.Client, q *models.Query, headers map[string]string) (*backend.DataResponse, error) {
res, err := c.QueryExemplars(ctx, q, sdkHeaderToHttpHeader(headers))
if err != nil {
return nil, err
}
return s.parseResponse(ctx, q, res)
}
func (s *QueryData) trace(ctx context.Context, q *models.Query) (context.Context, tracing.Span) {
traceCtx, span := s.tracer.Start(ctx, "datasource.prometheus")
span.SetAttributes("expr", q.Expr, attribute.Key("expr").String(q.Expr))
span.SetAttributes("start_unixnano", q.Start, attribute.Key("start_unixnano").Int64(q.Start.UnixNano()))
span.SetAttributes("stop_unixnano", q.End, attribute.Key("stop_unixnano").Int64(q.End.UnixNano()))
return traceCtx, span
func (s *QueryData) trace(ctx context.Context, q *models.Query) (context.Context, func()) {
return utils.StartTrace(ctx, s.tracer, "datasource.prometheus", []utils.Attribute{
{Key: "expr", Value: q.Expr, Kv: attribute.Key("expr").String(q.Expr)},
{Key: "start_unixnano", Value: q.Start, Kv: attribute.Key("start_unixnano").Int64(q.Start.UnixNano())},
{Key: "stop_unixnano", Value: q.End, Kv: attribute.Key("stop_unixnano").Int64(q.End.UnixNano())},
})
}
func sdkHeaderToHttpHeader(headers map[string]string) http.Header {
httpHeader := make(http.Header)
for key, val := range headers {
httpHeader[key] = []string{val}
}
return httpHeader
}
+36 -18
View File
@@ -10,12 +10,13 @@ import (
"testing"
"time"
"github.com/grafana/grafana-azure-sdk-go/azsettings"
"github.com/grafana/grafana-plugin-sdk-go/backend"
sdkhttpclient "github.com/grafana/grafana-plugin-sdk-go/backend/httpclient"
"github.com/grafana/grafana-plugin-sdk-go/data"
"github.com/grafana/grafana/pkg/infra/httpclient"
"github.com/grafana/grafana/pkg/infra/tracing"
"github.com/grafana/grafana/pkg/setting"
"github.com/grafana/grafana/pkg/tsdb/prometheus/buffered"
"github.com/grafana/grafana/pkg/tsdb/prometheus/models"
"github.com/grafana/grafana/pkg/tsdb/prometheus/querydata"
apiv1 "github.com/prometheus/client_golang/api/prometheus/v1"
@@ -58,7 +59,8 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
},
}
tctx := setup(true)
tctx, err := setup(true)
require.NoError(t, err)
qm := models.QueryModel{
LegendFormat: "legend {{app}}",
@@ -119,7 +121,8 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
},
JSON: b,
}
tctx := setup(true)
tctx, err := setup(true)
require.NoError(t, err)
res, err := execute(tctx, query, result)
require.NoError(t, err)
@@ -165,7 +168,8 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
},
JSON: b,
}
tctx := setup(true)
tctx, err := setup(true)
require.NoError(t, err)
res, err := execute(tctx, query, result)
require.NoError(t, err)
@@ -207,7 +211,8 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
},
JSON: b,
}
tctx := setup(true)
tctx, err := setup(true)
require.NoError(t, err)
res, err := execute(tctx, query, result)
require.NoError(t, err)
@@ -248,7 +253,8 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
JSON: b,
}
tctx := setup(true)
tctx, err := setup(true)
require.NoError(t, err)
res, err := execute(tctx, query, result)
require.NoError(t, err)
@@ -277,7 +283,8 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
query := backend.DataQuery{
JSON: b,
}
tctx := setup(true)
tctx, err := setup(true)
require.NoError(t, err)
res, err := execute(tctx, query, qr)
require.NoError(t, err)
@@ -315,7 +322,8 @@ func TestPrometheus_parseTimeSeriesResponse(t *testing.T) {
query := backend.DataQuery{
JSON: b,
}
tctx := setup(true)
tctx, err := setup(true)
require.NoError(t, err)
res, err := execute(tctx, query, qr)
require.NoError(t, err)
@@ -389,7 +397,7 @@ type testContext struct {
queryData *querydata.QueryData
}
func setup(wideFrames bool) *testContext {
func setup(wideFrames bool) (*testContext, error) {
tracer := tracing.InitializeTracerForTest()
httpProvider := &fakeHttpClientProvider{
opts: sdkhttpclient.Options{
@@ -400,19 +408,29 @@ func setup(wideFrames bool) *testContext {
Body: ioutil.NopCloser(bytes.NewReader([]byte(`{}`))),
},
}
queryData, _ := querydata.New(
httpProvider,
setting.NewCfg(),
&fakeFeatureToggles{flags: map[string]bool{"prometheusStreamingJSONParser": true, "prometheusWideSeries": wideFrames}},
tracer,
backend.DataSourceInstanceSettings{URL: "http://localhost:9090", JSONData: json.RawMessage(`{"timeInterval": "15s"}`)},
&fakeLogger{},
)
settings := backend.DataSourceInstanceSettings{
URL: "http://localhost:9090",
JSONData: json.RawMessage(`{"timeInterval": "15s"}`),
}
features := &fakeFeatureToggles{flags: map[string]bool{"prometheusStreamingJSONParser": true, "prometheusWideSeries": wideFrames}}
opts, err := buffered.CreateTransportOptions(settings, &azsettings.AzureSettings{}, features, &fakeLogger{})
if err != nil {
return nil, err
}
httpClient, err := httpProvider.New(*opts)
if err != nil {
return nil, err
}
queryData, _ := querydata.New(httpClient, features, tracer, settings, &fakeLogger{})
return &testContext{
httpProvider: httpProvider,
queryData: queryData,
}
}, nil
}
type fakeFeatureToggles struct {