Chore: Update prometheus, loki, graphite and influx plugins to support contextual logs. (#57708)
This commit is contained in:
@@ -9,28 +9,30 @@ import (
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/influxdata/influxdb-client-go/v2/api"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
)
|
||||
|
||||
const maxPointsEnforceFactor float64 = 10
|
||||
|
||||
// executeQuery runs a flux query using the queryModel to interpolate the query and the runner to execute it.
|
||||
// maxSeries somehow limits the response.
|
||||
func executeQuery(ctx context.Context, query queryModel, runner queryRunner, maxSeries int) (dr backend.DataResponse) {
|
||||
func executeQuery(ctx context.Context, logger log.Logger, query queryModel, runner queryRunner, maxSeries int) (dr backend.DataResponse) {
|
||||
dr = backend.DataResponse{}
|
||||
|
||||
flux := interpolate(query)
|
||||
|
||||
glog.Debug("Executing Flux query", "flux", flux)
|
||||
logger.Debug("Executing Flux query", "flux", flux)
|
||||
|
||||
tables, err := runner.runQuery(ctx, flux)
|
||||
if err != nil {
|
||||
glog.Warn("Flux query failed", "err", err, "query", flux)
|
||||
logger.Warn("Flux query failed", "err", err, "query", flux)
|
||||
dr.Error = err
|
||||
} else {
|
||||
// we only enforce a larger number than maxDataPoints
|
||||
maxPointsEnforced := int(float64(query.MaxDataPoints) * maxPointsEnforceFactor)
|
||||
|
||||
dr = readDataFrames(tables, maxPointsEnforced, maxSeries)
|
||||
dr = readDataFrames(logger, tables, maxPointsEnforced, maxSeries)
|
||||
|
||||
if dr.Error != nil {
|
||||
// we check if a too-many-data-points error happened, and if it is so,
|
||||
@@ -62,8 +64,8 @@ func executeQuery(ctx context.Context, query queryModel, runner queryRunner, max
|
||||
return dr
|
||||
}
|
||||
|
||||
func readDataFrames(result *api.QueryTableResult, maxPoints int, maxSeries int) (dr backend.DataResponse) {
|
||||
glog.Debug("Reading data frames from query result", "maxPoints", maxPoints, "maxSeries", maxSeries)
|
||||
func readDataFrames(logger log.Logger, result *api.QueryTableResult, maxPoints int, maxSeries int) (dr backend.DataResponse) {
|
||||
logger.Debug("Reading data frames from query result", "maxPoints", maxPoints, "maxSeries", maxSeries)
|
||||
dr = backend.DataResponse{}
|
||||
|
||||
builder := &frameBuilder{
|
||||
|
||||
@@ -14,12 +14,13 @@ import (
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/experimental"
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"github.com/xorcare/pointer"
|
||||
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
|
||||
influxdb2 "github.com/influxdata/influxdb-client-go/v2"
|
||||
"github.com/influxdata/influxdb-client-go/v2/api"
|
||||
)
|
||||
@@ -62,7 +63,7 @@ func executeMockedQuery(t *testing.T, name string, query queryModel) *backend.Da
|
||||
testDataPath: name + ".csv",
|
||||
}
|
||||
|
||||
dr := executeQuery(context.Background(), query, runner, 50)
|
||||
dr := executeQuery(context.Background(), glog, query, runner, 50)
|
||||
return &dr
|
||||
}
|
||||
|
||||
@@ -226,7 +227,7 @@ func TestRealQuery(t *testing.T) {
|
||||
runner, err := runnerFromDataSource(dsInfo)
|
||||
require.NoError(t, err)
|
||||
|
||||
dr := executeQuery(context.Background(), queryModel{
|
||||
dr := executeQuery(context.Background(), glog, queryModel{
|
||||
MaxDataPoints: 100,
|
||||
RawQuery: "buckets()",
|
||||
}, runner, 50)
|
||||
|
||||
@@ -5,10 +5,11 @@ import (
|
||||
"fmt"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
influxdb2 "github.com/influxdata/influxdb-client-go/v2"
|
||||
"github.com/influxdata/influxdb-client-go/v2/api"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -18,8 +19,9 @@ var (
|
||||
// Query builds flux queries, executes them, and returns the results.
|
||||
func Query(ctx context.Context, dsInfo *models.DatasourceInfo, tsdbQuery backend.QueryDataRequest) (
|
||||
*backend.QueryDataResponse, error) {
|
||||
logger := glog.FromContext(ctx)
|
||||
tRes := backend.NewQueryDataResponse()
|
||||
glog.Debug("Received a query", "query", tsdbQuery)
|
||||
logger.Debug("Received a query", "query", tsdbQuery)
|
||||
r, err := runnerFromDataSource(dsInfo)
|
||||
if err != nil {
|
||||
return &backend.QueryDataResponse{}, err
|
||||
@@ -36,7 +38,7 @@ func Query(ctx context.Context, dsInfo *models.DatasourceInfo, tsdbQuery backend
|
||||
|
||||
// If the default changes also update labels/placeholder in config page.
|
||||
maxSeries := dsInfo.MaxSeries
|
||||
res := executeQuery(ctx, *qm, r, maxSeries)
|
||||
res := executeQuery(ctx, logger, *qm, r, maxSeries)
|
||||
|
||||
tRes.Responses[query.RefID] = res
|
||||
}
|
||||
|
||||
@@ -7,6 +7,8 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/flux"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
)
|
||||
@@ -17,28 +19,30 @@ const (
|
||||
|
||||
func (s *Service) CheckHealth(ctx context.Context, req *backend.CheckHealthRequest) (*backend.CheckHealthResult,
|
||||
error) {
|
||||
logger := logger.FromContext(ctx)
|
||||
dsInfo, err := s.getDSInfo(req.PluginContext)
|
||||
if err != nil {
|
||||
return getHealthCheckMessage(s, "error getting datasource info", err)
|
||||
return getHealthCheckMessage(logger, "error getting datasource info", err)
|
||||
}
|
||||
|
||||
if dsInfo == nil {
|
||||
return getHealthCheckMessage(s, "", errors.New("invalid datasource info received"))
|
||||
return getHealthCheckMessage(logger, "", errors.New("invalid datasource info received"))
|
||||
}
|
||||
|
||||
switch dsInfo.Version {
|
||||
case influxVersionFlux:
|
||||
return CheckFluxHealth(ctx, dsInfo, s, req)
|
||||
return CheckFluxHealth(ctx, dsInfo, req)
|
||||
case influxVersionInfluxQL:
|
||||
return CheckInfluxQLHealth(ctx, dsInfo, s)
|
||||
default:
|
||||
return getHealthCheckMessage(s, "", errors.New("unknown influx version"))
|
||||
return getHealthCheckMessage(logger, "", errors.New("unknown influx version"))
|
||||
}
|
||||
}
|
||||
|
||||
func CheckFluxHealth(ctx context.Context, dsInfo *models.DatasourceInfo, s *Service,
|
||||
func CheckFluxHealth(ctx context.Context, dsInfo *models.DatasourceInfo,
|
||||
req *backend.CheckHealthRequest) (*backend.CheckHealthResult,
|
||||
error) {
|
||||
logger := logger.FromContext(ctx)
|
||||
ds, err := flux.Query(ctx, dsInfo, backend.QueryDataRequest{
|
||||
PluginContext: req.PluginContext,
|
||||
Queries: []backend.DataQuery{
|
||||
@@ -56,40 +60,41 @@ func CheckFluxHealth(ctx context.Context, dsInfo *models.DatasourceInfo, s *Serv
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
return getHealthCheckMessage(s, "error performing flux query", err)
|
||||
return getHealthCheckMessage(logger, "error performing flux query", err)
|
||||
}
|
||||
if res, ok := ds.Responses[refID]; ok {
|
||||
if res.Error != nil {
|
||||
return getHealthCheckMessage(s, "error reading buckets", res.Error)
|
||||
return getHealthCheckMessage(logger, "error reading buckets", res.Error)
|
||||
}
|
||||
if len(res.Frames) > 0 && len(res.Frames[0].Fields) > 0 {
|
||||
return getHealthCheckMessage(s, fmt.Sprintf("%d buckets found", res.Frames[0].Fields[0].Len()), nil)
|
||||
return getHealthCheckMessage(logger, fmt.Sprintf("%d buckets found", res.Frames[0].Fields[0].Len()), nil)
|
||||
}
|
||||
}
|
||||
|
||||
return getHealthCheckMessage(s, "", errors.New("error getting flux query buckets"))
|
||||
return getHealthCheckMessage(logger, "", errors.New("error getting flux query buckets"))
|
||||
}
|
||||
|
||||
func CheckInfluxQLHealth(ctx context.Context, dsInfo *models.DatasourceInfo, s *Service) (*backend.CheckHealthResult, error) {
|
||||
logger := logger.FromContext(ctx)
|
||||
queryString := "SHOW measurements"
|
||||
hcRequest, err := s.createRequest(ctx, dsInfo, queryString)
|
||||
hcRequest, err := s.createRequest(ctx, logger, dsInfo, queryString)
|
||||
if err != nil {
|
||||
return getHealthCheckMessage(s, "error creating influxDB healthcheck request", err)
|
||||
return getHealthCheckMessage(logger, "error creating influxDB healthcheck request", err)
|
||||
}
|
||||
|
||||
res, err := dsInfo.HTTPClient.Do(hcRequest)
|
||||
if err != nil {
|
||||
return getHealthCheckMessage(s, "error performing influxQL query", err)
|
||||
return getHealthCheckMessage(logger, "error performing influxQL query", err)
|
||||
}
|
||||
|
||||
defer func() {
|
||||
if err := res.Body.Close(); err != nil {
|
||||
s.glog.Warn("failed to close response body", "err", err)
|
||||
logger.Warn("failed to close response body", "err", err)
|
||||
}
|
||||
}()
|
||||
|
||||
if res.StatusCode/100 != 2 {
|
||||
return getHealthCheckMessage(s, "", fmt.Errorf("error reading InfluxDB. Status Code: %d", res.StatusCode))
|
||||
return getHealthCheckMessage(logger, "", fmt.Errorf("error reading InfluxDB. Status Code: %d", res.StatusCode))
|
||||
}
|
||||
resp := s.responseParser.Parse(res.Body, []Query{{
|
||||
RefID: refID,
|
||||
@@ -98,22 +103,22 @@ func CheckInfluxQLHealth(ctx context.Context, dsInfo *models.DatasourceInfo, s *
|
||||
}})
|
||||
if res, ok := resp.Responses[refID]; ok {
|
||||
if res.Error != nil {
|
||||
return getHealthCheckMessage(s, "error reading influxDB", res.Error)
|
||||
return getHealthCheckMessage(logger, "error reading influxDB", res.Error)
|
||||
}
|
||||
|
||||
if len(res.Frames) == 0 {
|
||||
return getHealthCheckMessage(s, "0 measurements found", nil)
|
||||
return getHealthCheckMessage(logger, "0 measurements found", nil)
|
||||
}
|
||||
|
||||
if len(res.Frames) > 0 && len(res.Frames[0].Fields) > 0 {
|
||||
return getHealthCheckMessage(s, fmt.Sprintf("%d measurements found", res.Frames[0].Fields[0].Len()), nil)
|
||||
return getHealthCheckMessage(logger, fmt.Sprintf("%d measurements found", res.Frames[0].Fields[0].Len()), nil)
|
||||
}
|
||||
}
|
||||
|
||||
return getHealthCheckMessage(s, "", errors.New("error connecting influxDB influxQL"))
|
||||
return getHealthCheckMessage(logger, "", errors.New("error connecting influxDB influxQL"))
|
||||
}
|
||||
|
||||
func getHealthCheckMessage(s *Service, message string, err error) (*backend.CheckHealthResult, error) {
|
||||
func getHealthCheckMessage(logger log.Logger, message string, err error) (*backend.CheckHealthResult, error) {
|
||||
if err == nil {
|
||||
return &backend.CheckHealthResult{
|
||||
Status: backend.HealthStatusOk,
|
||||
@@ -121,7 +126,7 @@ func getHealthCheckMessage(s *Service, message string, err error) (*backend.Chec
|
||||
}, nil
|
||||
}
|
||||
|
||||
s.glog.Warn("error performing influxdb healthcheck", "err", err.Error())
|
||||
logger.Warn("error performing influxdb healthcheck", "err", err.Error())
|
||||
errorMessage := fmt.Sprintf("%s %s", err.Error(), message)
|
||||
|
||||
return &backend.CheckHealthResult{
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend/datasource"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend/instancemgmt"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/httpclient"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/setting"
|
||||
@@ -20,10 +21,11 @@ import (
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
)
|
||||
|
||||
var logger log.Logger = log.New("tsdb.influxdb")
|
||||
|
||||
type Service struct {
|
||||
queryParser *InfluxdbQueryParser
|
||||
responseParser *ResponseParser
|
||||
glog log.Logger
|
||||
|
||||
im instancemgmt.InstanceManager
|
||||
}
|
||||
@@ -34,7 +36,6 @@ func ProvideService(httpClient httpclient.Provider) *Service {
|
||||
return &Service{
|
||||
queryParser: &InfluxdbQueryParser{},
|
||||
responseParser: &ResponseParser{},
|
||||
glog: log.New("tsdb.influxdb"),
|
||||
im: datasource.NewInstanceManager(newInstanceSettings(httpClient)),
|
||||
}
|
||||
}
|
||||
@@ -85,7 +86,8 @@ func newInstanceSettings(httpClientProvider httpclient.Provider) datasource.Inst
|
||||
}
|
||||
|
||||
func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
s.glog.Debug("Received a query request", "numQueries", len(req.Queries))
|
||||
logger := logger.FromContext(ctx)
|
||||
logger.Debug("Received a query request", "numQueries", len(req.Queries))
|
||||
|
||||
dsInfo, err := s.getDSInfo(req.PluginContext)
|
||||
if err != nil {
|
||||
@@ -96,7 +98,7 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
|
||||
return flux.Query(ctx, dsInfo, *req)
|
||||
}
|
||||
|
||||
s.glog.Debug("Making a non-Flux type query")
|
||||
logger.Debug("Making a non-Flux type query")
|
||||
|
||||
var allRawQueries string
|
||||
var queries []Query
|
||||
@@ -119,10 +121,10 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
|
||||
}
|
||||
|
||||
if setting.Env == setting.Dev {
|
||||
s.glog.Debug("Influxdb query", "raw query", allRawQueries)
|
||||
logger.Debug("Influxdb query", "raw query", allRawQueries)
|
||||
}
|
||||
|
||||
request, err := s.createRequest(ctx, dsInfo, allRawQueries)
|
||||
request, err := s.createRequest(ctx, logger, dsInfo, allRawQueries)
|
||||
if err != nil {
|
||||
return &backend.QueryDataResponse{}, err
|
||||
}
|
||||
@@ -133,7 +135,7 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
|
||||
}
|
||||
defer func() {
|
||||
if err := res.Body.Close(); err != nil {
|
||||
s.glog.Warn("Failed to close response body", "err", err)
|
||||
logger.Warn("Failed to close response body", "err", err)
|
||||
}
|
||||
}()
|
||||
if res.StatusCode/100 != 2 {
|
||||
@@ -145,7 +147,7 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (s *Service) createRequest(ctx context.Context, dsInfo *models.DatasourceInfo, query string) (*http.Request, error) {
|
||||
func (s *Service) createRequest(ctx context.Context, logger log.Logger, dsInfo *models.DatasourceInfo, query string) (*http.Request, error) {
|
||||
u, err := url.Parse(dsInfo.URL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -187,7 +189,7 @@ func (s *Service) createRequest(ctx context.Context, dsInfo *models.DatasourceIn
|
||||
|
||||
req.URL.RawQuery = params.Encode()
|
||||
|
||||
s.glog.Debug("Influxdb request", "url", req.URL.String())
|
||||
logger.Debug("Influxdb request", "url", req.URL.String())
|
||||
return req, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -6,10 +6,10 @@ import (
|
||||
"net/url"
|
||||
"testing"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
)
|
||||
|
||||
func TestExecutor_createRequest(t *testing.T) {
|
||||
@@ -22,11 +22,10 @@ func TestExecutor_createRequest(t *testing.T) {
|
||||
s := &Service{
|
||||
queryParser: &InfluxdbQueryParser{},
|
||||
responseParser: &ResponseParser{},
|
||||
glog: log.New("test"),
|
||||
}
|
||||
|
||||
t.Run("createRequest with GET httpMode", func(t *testing.T) {
|
||||
req, err := s.createRequest(context.Background(), datasource, query)
|
||||
req, err := s.createRequest(context.Background(), logger, datasource, query)
|
||||
|
||||
require.NoError(t, err)
|
||||
|
||||
@@ -40,7 +39,7 @@ func TestExecutor_createRequest(t *testing.T) {
|
||||
|
||||
t.Run("createRequest with POST httpMode", func(t *testing.T) {
|
||||
datasource.HTTPMode = "POST"
|
||||
req, err := s.createRequest(context.Background(), datasource, query)
|
||||
req, err := s.createRequest(context.Background(), logger, datasource, query)
|
||||
require.NoError(t, err)
|
||||
|
||||
assert.Equal(t, "POST", req.Method)
|
||||
@@ -59,7 +58,7 @@ func TestExecutor_createRequest(t *testing.T) {
|
||||
|
||||
t.Run("createRequest with PUT httpMode", func(t *testing.T) {
|
||||
datasource.HTTPMode = "PUT"
|
||||
_, err := s.createRequest(context.Background(), datasource, query)
|
||||
_, err := s.createRequest(context.Background(), logger, datasource, query)
|
||||
require.EqualError(t, err, ErrInvalidHttpMode.Error())
|
||||
})
|
||||
}
|
||||
|
||||
@@ -11,8 +11,8 @@ import (
|
||||
"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/backend/instancemgmt"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/httpclient"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/tsdb/influxdb/models"
|
||||
)
|
||||
|
||||
@@ -114,7 +114,6 @@ func GetMockService(version string, rt RoundTripper) *Service {
|
||||
return &Service{
|
||||
queryParser: &InfluxdbQueryParser{},
|
||||
responseParser: &ResponseParser{},
|
||||
glog: log.New("tsdb.influxdb"),
|
||||
im: &fakeInstance{
|
||||
version: version,
|
||||
fakeRoundTripper: rt,
|
||||
|
||||
Reference in New Issue
Block a user