Tempo: Migrates tags and tag values to datasource backend CallResource requests (#110511)

* Move tags and tag values request to datasource backend

* Remove outdated test

* Fix tests

* lint

* Refactor to use handlers for CallResource request

* lint

* Fix nit
This commit is contained in:
Andre Pereira
2025-09-09 17:31:32 +01:00
committed by GitHub
parent fc0db985c6
commit d26c6c112a
12 changed files with 268 additions and 311 deletions
@@ -235,7 +235,7 @@ func NewPlugin(pluginID string, cfg *setting.Cfg, httpClientProvider *httpclient
case Prometheus:
svc = prometheus.ProvideService(httpClientProvider)
case Tempo:
svc = tempo.ProvideService(httpClientProvider)
svc = tempo.ProvideService(httpClientProvider, tracer)
case PostgreSQL:
svc = postgres.ProvideService(cfg, features)
case MySQL:
+2 -2
View File
@@ -389,7 +389,7 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api
lokiService := loki.ProvideService(httpclientProvider, tracer)
opentsdbService := opentsdb.ProvideService(httpclientProvider)
prometheusService := prometheus.ProvideService(httpclientProvider)
tempoService := tempo.ProvideService(httpclientProvider)
tempoService := tempo.ProvideService(httpclientProvider, tracer)
testdatasourceService := testdatasource.ProvideService()
postgresService := postgres.ProvideService(cfg, featureToggles)
mysqlService := mysql.ProvideService()
@@ -976,7 +976,7 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac
lokiService := loki.ProvideService(httpclientProvider, tracer)
opentsdbService := opentsdb.ProvideService(httpclientProvider)
prometheusService := prometheus.ProvideService(httpclientProvider)
tempoService := tempo.ProvideService(httpclientProvider)
tempoService := tempo.ProvideService(httpclientProvider, tracer)
testdatasourceService := testdatasource.ProvideService()
postgresService := postgres.ProvideService(cfg, featureToggles)
mysqlService := mysql.ProvideService()
@@ -159,7 +159,7 @@ func TestIntegrationPluginManager(t *testing.T) {
lk := loki.ProvideService(hcp, tracer)
otsdb := opentsdb.ProvideService(hcp)
pr := prometheus.ProvideService(hcp)
tmpo := tempo.ProvideService(hcp)
tmpo := tempo.ProvideService(hcp, tracer)
td := testdatasource.ProvideService()
pg := postgres.ProvideService(cfg, features)
my := mysql.ProvideService()
+10 -3
View File
@@ -3,6 +3,8 @@ package main
import (
"context"
"go.opentelemetry.io/otel/trace/noop"
"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"
@@ -11,8 +13,9 @@ import (
)
var (
_ backend.QueryDataHandler = (*Datasource)(nil)
_ backend.StreamHandler = (*Datasource)(nil)
_ backend.QueryDataHandler = (*Datasource)(nil)
_ backend.StreamHandler = (*Datasource)(nil)
_ backend.CallResourceHandler = (*Datasource)(nil)
)
type Datasource struct {
@@ -21,7 +24,7 @@ type Datasource struct {
func NewDatasource(c context.Context, b backend.DataSourceInstanceSettings) (instancemgmt.Instance, error) {
return &Datasource{
Service: tempo.ProvideService(httpclient.NewProvider()),
Service: tempo.ProvideService(httpclient.NewProvider(), noop.NewTracerProvider().Tracer("tempo")),
}, nil
}
@@ -40,3 +43,7 @@ func (d *Datasource) PublishStream(ctx context.Context, req *backend.PublishStre
func (d *Datasource) RunStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender) error {
return d.Service.RunStream(ctx, req, sender)
}
func (d *Datasource) CallResource(ctx context.Context, req *backend.CallResourceRequest, sender backend.CallResourceResponseSender) error {
return d.Service.CallResource(ctx, req, sender)
}
+143 -4
View File
@@ -3,22 +3,38 @@ package tempo
import (
"context"
"fmt"
"io"
"net/http"
"net/url"
"path"
"runtime"
"strings"
"time"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes"
"go.opentelemetry.io/otel/trace"
"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/httpclient"
"github.com/grafana/grafana-plugin-sdk-go/backend/instancemgmt"
"github.com/grafana/grafana-plugin-sdk-go/backend/log"
"github.com/grafana/grafana-plugin-sdk-go/backend/resource/httpadapter"
"github.com/grafana/grafana/pkg/tsdb/tempo/kinds/dataquery"
"github.com/grafana/tempo/pkg/tempopb"
)
var (
_ backend.QueryDataHandler = (*Service)(nil)
_ backend.CallResourceHandler = (*Service)(nil)
)
type Service struct {
im instancemgmt.InstanceManager
logger log.Logger
im instancemgmt.InstanceManager
logger log.Logger
tracer trace.Tracer
resourceHandler backend.CallResourceHandler
}
type DatasourceInfo struct {
@@ -27,11 +43,20 @@ type DatasourceInfo struct {
URL string
}
func ProvideService(httpClientProvider *httpclient.Provider) *Service {
return &Service{
func ProvideService(httpClientProvider *httpclient.Provider, tracer trace.Tracer) *Service {
s := &Service{
im: datasource.NewInstanceManager(newInstanceSettings(httpClientProvider)),
logger: backend.NewLoggerWith("logger", "tsdb.tempo"),
tracer: tracer,
}
// Set up resource routes using httpadapter
mux := http.NewServeMux()
mux.HandleFunc("/tags", s.handleTags)
mux.HandleFunc("/tag-values", s.handleTagValues)
s.resourceHandler = httpadapter.New(mux)
return s
}
func newInstanceSettings(httpClientProvider *httpclient.Provider) datasource.InstanceFactoryFunc {
@@ -43,6 +68,8 @@ func newInstanceSettings(httpClientProvider *httpclient.Provider) datasource.Ins
return nil, err
}
opts.ForwardHTTPHeaders = true
client, err := httpClientProvider.New(opts)
if err != nil {
ctxLogger.Error("Failed to get HTTP client provider", "error", err, "function", logEntrypoint())
@@ -123,6 +150,118 @@ func (s *Service) getDSInfo(ctx context.Context, pluginCtx backend.PluginContext
return instance, nil
}
func (s *Service) CallResource(ctx context.Context, req *backend.CallResourceRequest, sender backend.CallResourceResponseSender) error {
return s.resourceHandler.CallResource(ctx, req, sender)
}
// handleTags handles requests to /tags resource
func (s *Service) handleTags(rw http.ResponseWriter, req *http.Request) {
s.proxyToTempo(rw, req, "api/v2/search/tags")
}
// handleTagValues handles requests to /tag-values resource
func (s *Service) handleTagValues(rw http.ResponseWriter, req *http.Request) {
// Extract the encoded tag from query parameters
encodedTag := req.URL.Query().Get("tag")
if encodedTag == "" {
http.Error(rw, "Missing required 'tag' parameter", http.StatusBadRequest)
return
}
tempoPath := fmt.Sprintf("api/v2/search/tag/%s/values", encodedTag)
s.proxyToTempo(rw, req, tempoPath)
}
// proxyToTempo is the shared function that builds the URL and proxies requests to Tempo
func (s *Service) proxyToTempo(rw http.ResponseWriter, req *http.Request, tempoPath string) {
ctx := req.Context()
pCtx := backend.PluginConfigFromContext(ctx)
// Get datasource info
dsInfo, err := s.getDSInfo(ctx, pCtx)
if err != nil {
s.logger.Error("Failed to get data source info", "error", err)
http.Error(rw, "Failed to get data source configuration", http.StatusInternalServerError)
return
}
ctx, span := s.tracer.Start(ctx, "datasource.tempo.proxyToTempo", trace.WithAttributes(
attribute.String("tempoPath", tempoPath),
))
defer span.End()
// Build the full URL to Tempo
parsedURL, err := url.Parse(dsInfo.URL)
if err != nil {
span.RecordError(err)
span.SetStatus(codes.Error, err.Error())
s.logger.Error("Failed to parse data source URL", "error", err, "url", dsInfo.URL)
http.Error(rw, "Invalid data source URL", http.StatusInternalServerError)
return
}
// Join the tempo path with the base URL
parsedURL.Path = path.Join(parsedURL.Path, tempoPath)
// Preserve query parameters from the original request
parsedURL.RawQuery = req.URL.RawQuery
s.logger.Debug("Making resource request to Tempo", "url", parsedURL.String())
start := time.Now()
// Create the request to Tempo
httpReq, err := http.NewRequestWithContext(ctx, req.Method, parsedURL.String(), req.Body)
if err != nil {
span.RecordError(err)
span.SetStatus(codes.Error, err.Error())
s.logger.Error("Failed to create HTTP request", "error", err)
http.Error(rw, "Failed to create request", http.StatusInternalServerError)
return
}
// Copy headers from the original request
for name, values := range req.Header {
for _, value := range values {
httpReq.Header.Add(name, value)
}
}
// Make the request to Tempo
resp, err := dsInfo.HTTPClient.Do(httpReq)
if err != nil {
span.RecordError(err)
span.SetStatus(codes.Error, err.Error())
s.logger.Error("Failed resource call to Tempo", "error", err, "url", parsedURL.String(), "duration", time.Since(start))
http.Error(rw, "Failed to connect to Tempo", http.StatusBadGateway)
return
}
defer func() {
if err := resp.Body.Close(); err != nil {
s.logger.Warn("Failed to close response body", "error", err)
}
}()
s.logger.Debug("Response received from Tempo", "statusCode", resp.StatusCode, "contentLength", resp.Header.Get("Content-Length"), "duration", time.Since(start))
// Copy response headers
for name, values := range resp.Header {
for _, value := range values {
rw.Header().Add(name, value)
}
}
// Set the status code
rw.WriteHeader(resp.StatusCode)
// Copy the response body
_, err = io.Copy(rw, resp.Body)
if err != nil {
span.RecordError(err)
span.SetStatus(codes.Error, err.Error())
s.logger.Error("Failed to copy response body", "error", err)
return
}
}
// Return the file, line, and (full-path) function name of the caller
func getRunContext() (string, int, string) {
pc := make([]uintptr, 10)