Tempo: small refactor to tempo backend (#109581)
* update datasource info struct name * remove unnecessary abstraction * update error messages
This commit is contained in:
@@ -24,7 +24,7 @@ type PartialTempoQuery struct {
|
||||
MetricsQueryType *dataquery.MetricsQueryType
|
||||
}
|
||||
|
||||
func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *Datasource) error {
|
||||
func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *DatasourceInfo) error {
|
||||
ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.runMetricsStream")
|
||||
defer span.End()
|
||||
|
||||
|
||||
@@ -31,7 +31,7 @@ type StreamSender interface {
|
||||
SendBytes(data []byte) error
|
||||
}
|
||||
|
||||
func (s *Service) runSearchStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *Datasource) error {
|
||||
func (s *Service) runSearchStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *DatasourceInfo) error {
|
||||
ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.runSearchStream")
|
||||
defer span.End()
|
||||
|
||||
|
||||
@@ -6,18 +6,19 @@ import (
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend/httpclient"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend/instancemgmt"
|
||||
|
||||
tempo "github.com/grafana/grafana/pkg/tsdb/tempo"
|
||||
)
|
||||
|
||||
type Datasource struct {
|
||||
Service *tempo.Service
|
||||
}
|
||||
|
||||
var (
|
||||
_ backend.QueryDataHandler = (*Datasource)(nil)
|
||||
_ backend.StreamHandler = (*Datasource)(nil)
|
||||
)
|
||||
|
||||
type Datasource struct {
|
||||
Service *tempo.Service
|
||||
}
|
||||
|
||||
func NewDatasource(c context.Context, b backend.DataSourceInstanceSettings) (instancemgmt.Instance, error) {
|
||||
return &Datasource{
|
||||
Service: tempo.ProvideService(httpclient.NewProvider()),
|
||||
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
)
|
||||
|
||||
func main() {
|
||||
// Created as described at https://grafana.com/developers/plugin-tools/introduction/backend-plugins
|
||||
if err := datasource.Manage("tempo", NewDatasource, datasource.ManageOpts{}); err != nil {
|
||||
log.DefaultLogger.Error(err.Error())
|
||||
os.Exit(1)
|
||||
|
||||
+52
-44
@@ -21,36 +21,19 @@ type Service struct {
|
||||
logger log.Logger
|
||||
}
|
||||
|
||||
// Return the file, line, and (full-path) function name of the caller
|
||||
func getRunContext() (string, int, string) {
|
||||
pc := make([]uintptr, 10)
|
||||
runtime.Callers(2, pc)
|
||||
f := runtime.FuncForPC(pc[0])
|
||||
file, line := f.FileLine(pc[0])
|
||||
return file, line, f.Name()
|
||||
}
|
||||
|
||||
// Return a formatted string representing the execution context for the logger
|
||||
func logEntrypoint() string {
|
||||
file, line, pathToFunction := getRunContext()
|
||||
parts := strings.Split(pathToFunction, "/")
|
||||
functionName := parts[len(parts)-1]
|
||||
return fmt.Sprintf("%s:%d[%s]", file, line, functionName)
|
||||
type DatasourceInfo struct {
|
||||
HTTPClient *http.Client
|
||||
StreamingClient tempopb.StreamingQuerierClient
|
||||
URL string
|
||||
}
|
||||
|
||||
func ProvideService(httpClientProvider *httpclient.Provider) *Service {
|
||||
return &Service{
|
||||
logger: backend.NewLoggerWith("logger", "tsdb.tempo"),
|
||||
im: datasource.NewInstanceManager(newInstanceSettings(httpClientProvider)),
|
||||
logger: backend.NewLoggerWith("logger", "tsdb.tempo"),
|
||||
}
|
||||
}
|
||||
|
||||
type Datasource struct {
|
||||
HTTPClient *http.Client
|
||||
StreamingClient tempopb.StreamingQuerierClient
|
||||
URL string
|
||||
}
|
||||
|
||||
func newInstanceSettings(httpClientProvider *httpclient.Provider) datasource.InstanceFactoryFunc {
|
||||
return func(ctx context.Context, settings backend.DataSourceInstanceSettings) (instancemgmt.Instance, error) {
|
||||
ctxLogger := backend.NewLoggerWith("logger", "tsdb.tempo").FromContext(ctx)
|
||||
@@ -72,7 +55,7 @@ func newInstanceSettings(httpClientProvider *httpclient.Provider) datasource.Ins
|
||||
return nil, err
|
||||
}
|
||||
|
||||
model := &Datasource{
|
||||
model := &DatasourceInfo{
|
||||
HTTPClient: client,
|
||||
StreamingClient: streamingClient,
|
||||
URL: settings.URL,
|
||||
@@ -91,16 +74,34 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
|
||||
// loop over queries and execute them individually.
|
||||
for i, q := range req.Queries {
|
||||
ctxLogger.Debug("Processing query", "counter", i, "function", logEntrypoint())
|
||||
if res, err := s.query(ctx, req.PluginContext, q); err != nil {
|
||||
ctxLogger.Error("Error processing query", "error", err)
|
||||
return response, err
|
||||
} else {
|
||||
if res != nil {
|
||||
ctxLogger.Debug("Query processed", "counter", i, "function", logEntrypoint())
|
||||
response.Responses[q.RefID] = *res
|
||||
} else {
|
||||
ctxLogger.Debug("Query resulted in empty response", "counter", i, "function", logEntrypoint())
|
||||
|
||||
var res *backend.DataResponse
|
||||
var err error
|
||||
|
||||
switch q.QueryType {
|
||||
case string(dataquery.TempoQueryTypeTraceId):
|
||||
res, err = s.getTrace(ctx, req.PluginContext, q)
|
||||
if err != nil {
|
||||
ctxLogger.Error("Error processing TraceId query", "error", err)
|
||||
return response, err
|
||||
}
|
||||
|
||||
case string(dataquery.TempoQueryTypeTraceql):
|
||||
res, err = s.runTraceQlQuery(ctx, req.PluginContext, q)
|
||||
if err != nil {
|
||||
ctxLogger.Error("Error processing TraceQL query", "error", err)
|
||||
return response, err
|
||||
}
|
||||
|
||||
default:
|
||||
return nil, fmt.Errorf("unsupported query type: '%s' for query with refID '%s'", q.QueryType, q.RefID)
|
||||
}
|
||||
|
||||
if res != nil {
|
||||
ctxLogger.Debug("Query processed", "counter", i, "function", logEntrypoint())
|
||||
response.Responses[q.RefID] = *res
|
||||
} else {
|
||||
ctxLogger.Debug("Query resulted in empty response", "counter", i, "function", logEntrypoint())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -108,26 +109,33 @@ func (s *Service) QueryData(ctx context.Context, req *backend.QueryDataRequest)
|
||||
return response, nil
|
||||
}
|
||||
|
||||
func (s *Service) query(ctx context.Context, pCtx backend.PluginContext, query backend.DataQuery) (*backend.DataResponse, error) {
|
||||
switch query.QueryType {
|
||||
case string(dataquery.TempoQueryTypeTraceId):
|
||||
return s.getTrace(ctx, pCtx, query)
|
||||
case string(dataquery.TempoQueryTypeTraceql):
|
||||
return s.runTraceQlQuery(ctx, pCtx, query)
|
||||
}
|
||||
return nil, fmt.Errorf("unsupported query type: '%s' for query with refID '%s'", query.QueryType, query.RefID)
|
||||
}
|
||||
|
||||
func (s *Service) getDSInfo(ctx context.Context, pluginCtx backend.PluginContext) (*Datasource, error) {
|
||||
func (s *Service) getDSInfo(ctx context.Context, pluginCtx backend.PluginContext) (*DatasourceInfo, error) {
|
||||
i, err := s.im.Get(ctx, pluginCtx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
instance, ok := i.(*Datasource)
|
||||
instance, ok := i.(*DatasourceInfo)
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("failed to cast datsource info")
|
||||
}
|
||||
|
||||
return instance, nil
|
||||
}
|
||||
|
||||
// Return the file, line, and (full-path) function name of the caller
|
||||
func getRunContext() (string, int, string) {
|
||||
pc := make([]uintptr, 10)
|
||||
runtime.Callers(2, pc)
|
||||
f := runtime.FuncForPC(pc[0])
|
||||
file, line := f.FileLine(pc[0])
|
||||
return file, line, f.Name()
|
||||
}
|
||||
|
||||
// Return a formatted string representing the execution context for the logger
|
||||
func logEntrypoint() string {
|
||||
file, line, pathToFunction := getRunContext()
|
||||
parts := strings.Split(pathToFunction, "/")
|
||||
functionName := parts[len(parts)-1]
|
||||
return fmt.Sprintf("%s:%d[%s]", file, line, functionName)
|
||||
}
|
||||
|
||||
@@ -135,7 +135,7 @@ func (s *Service) getTrace(ctx context.Context, pCtx backend.PluginContext, quer
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (s *Service) performTraceRequest(ctx context.Context, dsInfo *Datasource, apiVersion TraceRequestApiVersion, model *dataquery.TempoQuery, query backend.DataQuery, span trace.Span) (*http.Response, []byte, error) {
|
||||
func (s *Service) performTraceRequest(ctx context.Context, dsInfo *DatasourceInfo, apiVersion TraceRequestApiVersion, model *dataquery.TempoQuery, query backend.DataQuery, span trace.Span) (*http.Response, []byte, error) {
|
||||
ctxLogger := s.logger.FromContext(ctx)
|
||||
request, err := s.createRequest(ctx, dsInfo, apiVersion, *model.Query, query.TimeRange.From.Unix(), query.TimeRange.To.Unix())
|
||||
|
||||
@@ -177,7 +177,7 @@ const (
|
||||
TraceRequestApiVersionV2
|
||||
)
|
||||
|
||||
func (s *Service) createRequest(ctx context.Context, dsInfo *Datasource, apiVersion TraceRequestApiVersion, traceID string, start int64, end int64) (*http.Request, error) {
|
||||
func (s *Service) createRequest(ctx context.Context, dsInfo *DatasourceInfo, apiVersion TraceRequestApiVersion, traceID string, start int64, end int64) (*http.Request, error) {
|
||||
ctxLogger := s.logger.FromContext(ctx)
|
||||
var baseUrl string
|
||||
var tempoQuery string
|
||||
|
||||
@@ -12,7 +12,7 @@ import (
|
||||
func TestTempo(t *testing.T) {
|
||||
t.Run("createRequest v1 without time range - success", func(t *testing.T) {
|
||||
service := &Service{logger: backend.NewLoggerWith("logger", "tempo-test")}
|
||||
req, err := service.createRequest(context.Background(), &Datasource{}, TraceRequestApiVersionV1, "traceID", 0, 0)
|
||||
req, err := service.createRequest(context.Background(), &DatasourceInfo{}, TraceRequestApiVersionV1, "traceID", 0, 0)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, len(req.Header))
|
||||
assert.Equal(t, "/api/traces/traceID", req.URL.String())
|
||||
@@ -20,7 +20,7 @@ func TestTempo(t *testing.T) {
|
||||
|
||||
t.Run("createRequest v1 with time range - success", func(t *testing.T) {
|
||||
service := &Service{logger: backend.NewLoggerWith("logger", "tempo-test")}
|
||||
req, err := service.createRequest(context.Background(), &Datasource{}, TraceRequestApiVersionV1, "traceID", 1, 2)
|
||||
req, err := service.createRequest(context.Background(), &DatasourceInfo{}, TraceRequestApiVersionV1, "traceID", 1, 2)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, len(req.Header))
|
||||
assert.Equal(t, "/api/traces/traceID?start=1&end=2", req.URL.String())
|
||||
@@ -28,7 +28,7 @@ func TestTempo(t *testing.T) {
|
||||
|
||||
t.Run("createRequest v2 without time range - success", func(t *testing.T) {
|
||||
service := &Service{logger: backend.NewLoggerWith("logger", "tempo-test")}
|
||||
req, err := service.createRequest(context.Background(), &Datasource{}, TraceRequestApiVersionV2, "traceID", 0, 0)
|
||||
req, err := service.createRequest(context.Background(), &DatasourceInfo{}, TraceRequestApiVersionV2, "traceID", 0, 0)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, len(req.Header))
|
||||
assert.Equal(t, "/api/v2/traces/traceID", req.URL.String())
|
||||
@@ -36,7 +36,7 @@ func TestTempo(t *testing.T) {
|
||||
|
||||
t.Run("createRequest v2 with time range - success", func(t *testing.T) {
|
||||
service := &Service{logger: backend.NewLoggerWith("logger", "tempo-test")}
|
||||
req, err := service.createRequest(context.Background(), &Datasource{}, TraceRequestApiVersionV2, "traceID", 1, 2)
|
||||
req, err := service.createRequest(context.Background(), &DatasourceInfo{}, TraceRequestApiVersionV2, "traceID", 1, 2)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, 1, len(req.Header))
|
||||
assert.Equal(t, "/api/v2/traces/traceID?start=1&end=2", req.URL.String())
|
||||
|
||||
@@ -130,7 +130,7 @@ func handleConversionError(ctxLogger log.Logger, span trace.Span, err error) (*b
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (s *Service) performMetricsQuery(ctx context.Context, dsInfo *Datasource, model *dataquery.TempoQuery, query backend.DataQuery, span trace.Span) (*http.Response, []byte, error) {
|
||||
func (s *Service) performMetricsQuery(ctx context.Context, dsInfo *DatasourceInfo, model *dataquery.TempoQuery, query backend.DataQuery, span trace.Span) (*http.Response, []byte, error) {
|
||||
ctxLogger := s.logger.FromContext(ctx)
|
||||
request, err := s.createMetricsQuery(ctx, dsInfo, model, query.TimeRange.From.Unix(), query.TimeRange.To.Unix())
|
||||
if err != nil {
|
||||
@@ -156,7 +156,7 @@ func (s *Service) performMetricsQuery(ctx context.Context, dsInfo *Datasource, m
|
||||
return resp, body, nil
|
||||
}
|
||||
|
||||
func (s *Service) createMetricsQuery(ctx context.Context, dsInfo *Datasource, query *dataquery.TempoQuery, start int64, end int64) (*http.Request, error) {
|
||||
func (s *Service) createMetricsQuery(ctx context.Context, dsInfo *DatasourceInfo, query *dataquery.TempoQuery, start int64, end int64) (*http.Request, error) {
|
||||
ctxLogger := s.logger.FromContext(ctx)
|
||||
|
||||
queryType := "query_range"
|
||||
|
||||
@@ -14,7 +14,7 @@ func TestCreateMetricsQuery_Success(t *testing.T) {
|
||||
service := &Service{
|
||||
logger: logger,
|
||||
}
|
||||
dsInfo := &Datasource{
|
||||
dsInfo := &DatasourceInfo{
|
||||
URL: "http://tempo:3100",
|
||||
}
|
||||
queryVal := "{attribute=\"value\"}"
|
||||
@@ -40,7 +40,7 @@ func TestCreateMetricsQuery_OnlyQuery(t *testing.T) {
|
||||
service := &Service{
|
||||
logger: logger,
|
||||
}
|
||||
dsInfo := &Datasource{
|
||||
dsInfo := &DatasourceInfo{
|
||||
URL: "http://tempo:3100",
|
||||
}
|
||||
queryVal := "{attribute=\"value\"}"
|
||||
@@ -60,7 +60,7 @@ func TestCreateMetricsQuery_URLParseError(t *testing.T) {
|
||||
service := &Service{
|
||||
logger: logger,
|
||||
}
|
||||
dsInfo := &Datasource{
|
||||
dsInfo := &DatasourceInfo{
|
||||
URL: "http://[::1]:namedport",
|
||||
}
|
||||
queryVal := "{attribute=\"value\"}"
|
||||
|
||||
Reference in New Issue
Block a user