Chore: Update sqleng, elasticsearch, tempo and opentsdb plugins to support contextual logs. (#57777)

* make sql engine use pick log context for logs
* update tempo to get log context
* update opentsdb to use log context
* update es client to use log context
This commit is contained in:
Yuriy Tseretyan
2022-11-02 10:03:50 -04:00
committed by GitHub
parent 17ebeab02c
commit d9c40ca41e
13 changed files with 84 additions and 106 deletions
+17 -15
View File
@@ -23,15 +23,15 @@ import (
"github.com/grafana/grafana/pkg/setting"
)
var logger = log.New("tsdb.opentsdb")
type Service struct {
logger log.Logger
im instancemgmt.InstanceManager
im instancemgmt.InstanceManager
}
func ProvideService(httpClientProvider httpclient.Provider) *Service {
return &Service{
logger: log.New("tsdb.opentsdb"),
im: datasource.NewInstanceManager(newInstanceSettings(httpClientProvider)),
im: datasource.NewInstanceManager(newInstanceSettings(httpClientProvider)),
}
}
@@ -66,6 +66,8 @@ func newInstanceSettings(httpClientProvider httpclient.Provider) datasource.Inst
func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
var tsdbQuery OpenTsdbQuery
logger := logger.FromContext(ctx)
q := req.Queries[0]
tsdbQuery.Start = q.TimeRange.From.UnixNano() / int64(time.Millisecond)
@@ -78,7 +80,7 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
// TODO: Don't use global variable
if setting.Env == setting.Dev {
s.logger.Debug("OpenTsdb request", "params", tsdbQuery)
logger.Debug("OpenTsdb request", "params", tsdbQuery)
}
dsInfo, err := s.getDSInfo(req.PluginContext)
@@ -86,7 +88,7 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
return nil, err
}
request, err := s.createRequest(ctx, dsInfo, tsdbQuery)
request, err := s.createRequest(ctx, logger, dsInfo, tsdbQuery)
if err != nil {
return &backend.QueryDataResponse{}, err
}
@@ -96,7 +98,7 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
return &backend.QueryDataResponse{}, err
}
result, err := s.parseResponse(res)
result, err := s.parseResponse(logger, res)
if err != nil {
return &backend.QueryDataResponse{}, err
}
@@ -104,7 +106,7 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
return result, nil
}
func (s *Service) createRequest(ctx context.Context, dsInfo *datasourceInfo, data OpenTsdbQuery) (*http.Request, error) {
func (s *Service) createRequest(ctx context.Context, logger log.Logger, dsInfo *datasourceInfo, data OpenTsdbQuery) (*http.Request, error) {
u, err := url.Parse(dsInfo.URL)
if err != nil {
return nil, err
@@ -113,13 +115,13 @@ func (s *Service) createRequest(ctx context.Context, dsInfo *datasourceInfo, dat
postData, err := json.Marshal(data)
if err != nil {
s.logger.Info("Failed marshaling data", "error", err)
logger.Info("Failed marshaling data", "error", err)
return nil, fmt.Errorf("failed to create request: %w", err)
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, u.String(), strings.NewReader(string(postData)))
if err != nil {
s.logger.Info("Failed to create request", "error", err)
logger.Info("Failed to create request", "error", err)
return nil, fmt.Errorf("failed to create request: %w", err)
}
@@ -127,7 +129,7 @@ func (s *Service) createRequest(ctx context.Context, dsInfo *datasourceInfo, dat
return req, nil
}
func (s *Service) parseResponse(res *http.Response) (*backend.QueryDataResponse, error) {
func (s *Service) parseResponse(logger log.Logger, res *http.Response) (*backend.QueryDataResponse, error) {
resp := backend.NewQueryDataResponse()
body, err := io.ReadAll(res.Body)
@@ -136,19 +138,19 @@ func (s *Service) parseResponse(res *http.Response) (*backend.QueryDataResponse,
}
defer func() {
if err := res.Body.Close(); err != nil {
s.logger.Warn("Failed to close response body", "err", err)
logger.Warn("Failed to close response body", "err", err)
}
}()
if res.StatusCode/100 != 2 {
s.logger.Info("Request failed", "status", res.Status, "body", string(body))
logger.Info("Request failed", "status", res.Status, "body", string(body))
return nil, fmt.Errorf("request failed, status: %s", res.Status)
}
var responseData []OpenTsdbResponse
err = json.Unmarshal(body, &responseData)
if err != nil {
s.logger.Info("Failed to unmarshal opentsdb response", "error", err, "status", res.Status, "body", string(body))
logger.Info("Failed to unmarshal opentsdb response", "error", err, "status", res.Status, "body", string(body))
return nil, err
}
@@ -162,7 +164,7 @@ func (s *Service) parseResponse(res *http.Response) (*backend.QueryDataResponse,
for timeString, value := range val.DataPoints {
timestamp, err := strconv.ParseInt(timeString, 10, 64)
if err != nil {
s.logger.Info("Failed to unmarshal opentsdb timestamp", "timestamp", timeString)
logger.Info("Failed to unmarshal opentsdb timestamp", "timestamp", timeString)
return nil, err
}
timeVector = append(timeVector, time.Unix(timestamp, 0).UTC())
+4 -8
View File
@@ -13,17 +13,13 @@ import (
"github.com/grafana/grafana-plugin-sdk-go/data"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/grafana/grafana/pkg/infra/log"
)
func TestOpenTsdbExecutor(t *testing.T) {
service := &Service{
logger: log.New("test"),
}
service := &Service{}
t.Run("create request", func(t *testing.T) {
req, err := service.createRequest(context.Background(), &datasourceInfo{}, OpenTsdbQuery{})
req, err := service.createRequest(context.Background(), logger, &datasourceInfo{}, OpenTsdbQuery{})
require.NoError(t, err)
assert.Equal(t, "POST", req.Method)
@@ -37,7 +33,7 @@ func TestOpenTsdbExecutor(t *testing.T) {
t.Run("Parse response should handle invalid JSON", func(t *testing.T) {
response := `{ invalid }`
result, err := service.parseResponse(&http.Response{Body: io.NopCloser(strings.NewReader(response))})
result, err := service.parseResponse(logger, &http.Response{Body: io.NopCloser(strings.NewReader(response))})
require.Nil(t, result)
require.Error(t, err)
})
@@ -67,7 +63,7 @@ func TestOpenTsdbExecutor(t *testing.T) {
resp := http.Response{Body: io.NopCloser(strings.NewReader(response))}
resp.StatusCode = 200
result, err := service.parseResponse(&resp)
result, err := service.parseResponse(logger, &resp)
require.NoError(t, err)
frame := result.Responses["A"]