From 7bebd624465e60654d3a0b9f5c312dc9b7ac0f19 Mon Sep 17 00:00:00 2001 From: Gareth Date: Thu, 14 Aug 2025 10:23:44 +0100 Subject: [PATCH] Tempo: small refactor to tempo backend (#109581) * update datasource info struct name * remove unnecessary abstraction * update error messages --- pkg/tsdb/tempo/metrics_stream.go | 2 +- pkg/tsdb/tempo/search_stream.go | 2 +- pkg/tsdb/tempo/standalone/datasource.go | 9 +-- pkg/tsdb/tempo/standalone/main.go | 1 - pkg/tsdb/tempo/tempo.go | 96 +++++++++++++------------ pkg/tsdb/tempo/trace.go | 4 +- pkg/tsdb/tempo/trace_test.go | 8 +-- pkg/tsdb/tempo/traceql_query.go | 4 +- pkg/tsdb/tempo/traceql_query_test.go | 6 +- 9 files changed, 70 insertions(+), 62 deletions(-) diff --git a/pkg/tsdb/tempo/metrics_stream.go b/pkg/tsdb/tempo/metrics_stream.go index 85fac2df83b..cd711ac80f9 100644 --- a/pkg/tsdb/tempo/metrics_stream.go +++ b/pkg/tsdb/tempo/metrics_stream.go @@ -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() diff --git a/pkg/tsdb/tempo/search_stream.go b/pkg/tsdb/tempo/search_stream.go index 527a1789f32..c78e138e4a7 100644 --- a/pkg/tsdb/tempo/search_stream.go +++ b/pkg/tsdb/tempo/search_stream.go @@ -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() diff --git a/pkg/tsdb/tempo/standalone/datasource.go b/pkg/tsdb/tempo/standalone/datasource.go index 087cb4dff71..8d371c07917 100644 --- a/pkg/tsdb/tempo/standalone/datasource.go +++ b/pkg/tsdb/tempo/standalone/datasource.go @@ -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()), diff --git a/pkg/tsdb/tempo/standalone/main.go b/pkg/tsdb/tempo/standalone/main.go index 6961ead4c20..748c9459e68 100644 --- a/pkg/tsdb/tempo/standalone/main.go +++ b/pkg/tsdb/tempo/standalone/main.go @@ -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) diff --git a/pkg/tsdb/tempo/tempo.go b/pkg/tsdb/tempo/tempo.go index 96a295b7fac..cbe4893f8e2 100644 --- a/pkg/tsdb/tempo/tempo.go +++ b/pkg/tsdb/tempo/tempo.go @@ -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) +} diff --git a/pkg/tsdb/tempo/trace.go b/pkg/tsdb/tempo/trace.go index 2536f3988ef..6cc705de502 100644 --- a/pkg/tsdb/tempo/trace.go +++ b/pkg/tsdb/tempo/trace.go @@ -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 diff --git a/pkg/tsdb/tempo/trace_test.go b/pkg/tsdb/tempo/trace_test.go index 78983040050..63932ee66bc 100644 --- a/pkg/tsdb/tempo/trace_test.go +++ b/pkg/tsdb/tempo/trace_test.go @@ -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()) diff --git a/pkg/tsdb/tempo/traceql_query.go b/pkg/tsdb/tempo/traceql_query.go index 1b5d31edad5..058a7944aa0 100644 --- a/pkg/tsdb/tempo/traceql_query.go +++ b/pkg/tsdb/tempo/traceql_query.go @@ -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" diff --git a/pkg/tsdb/tempo/traceql_query_test.go b/pkg/tsdb/tempo/traceql_query_test.go index cffe2c2273f..2dddec9696e 100644 --- a/pkg/tsdb/tempo/traceql_query_test.go +++ b/pkg/tsdb/tempo/traceql_query_test.go @@ -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\"}"