Chore: Add tracing to tempo, parca and pyroscope datasource backends (#76368)
* Added spans to trace.go
* Added spans to search_stream.go
* Added spans to parca datasource
* Added spans for pyroscope
* Fix tests
* Fix another test
* Lint
* Revert "Fix another test"
This reverts commit a1639049e3.
* Use grafana-sdk-go tracing
This commit is contained in:
@@ -8,9 +8,12 @@ import (
|
||||
"io"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend/tracing"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/grafana/grafana/pkg/tsdb/tempo/kinds/dataquery"
|
||||
"github.com/grafana/tempo/pkg/tempopb"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/codes"
|
||||
)
|
||||
|
||||
const SearchPathPrefix = "search/"
|
||||
@@ -27,12 +30,17 @@ type StreamSender interface {
|
||||
}
|
||||
|
||||
func (s *Service) runSearchStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *Datasource) error {
|
||||
ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.runSearchStream")
|
||||
defer span.End()
|
||||
|
||||
response := &backend.DataResponse{}
|
||||
|
||||
var backendQuery *backend.DataQuery
|
||||
err := json.Unmarshal(req.Data, &backendQuery)
|
||||
if err != nil {
|
||||
response.Error = fmt.Errorf("error unmarshaling backend query model: %v", err)
|
||||
span.RecordError(response.Error)
|
||||
span.SetStatus(codes.Error, response.Error.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -40,6 +48,8 @@ func (s *Service) runSearchStream(ctx context.Context, req *backend.RunStreamReq
|
||||
err = json.Unmarshal(req.Data, &sr)
|
||||
if err != nil {
|
||||
response.Error = fmt.Errorf("error unmarshaling Tempo query model: %v", err)
|
||||
span.RecordError(response.Error)
|
||||
span.SetStatus(codes.Error, response.Error.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -52,46 +62,60 @@ func (s *Service) runSearchStream(ctx context.Context, req *backend.RunStreamReq
|
||||
|
||||
stream, err := datasource.StreamingClient.Search(ctx, sr)
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
s.logger.Error("Error Search()", "err", err)
|
||||
return err
|
||||
}
|
||||
|
||||
return s.processStream(stream, sender)
|
||||
return s.processStream(ctx, stream, sender)
|
||||
}
|
||||
|
||||
func (s *Service) processStream(stream tempopb.StreamingQuerier_SearchClient, sender StreamSender) error {
|
||||
func (s *Service) processStream(ctx context.Context, stream tempopb.StreamingQuerier_SearchClient, sender StreamSender) error {
|
||||
ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.processStream")
|
||||
defer span.End()
|
||||
var traceList []*tempopb.TraceSearchMetadata
|
||||
var metrics *tempopb.SearchMetrics
|
||||
messageCount := 0
|
||||
for {
|
||||
msg, err := stream.Recv()
|
||||
messageCount++
|
||||
span.SetAttributes(attribute.Int("message_count", messageCount))
|
||||
if errors.Is(err, io.EOF) {
|
||||
if err := sendResponse(&ExtendedResponse{
|
||||
if err := s.sendResponse(ctx, &ExtendedResponse{
|
||||
State: dataquery.SearchStreamingStateDone,
|
||||
SearchResponse: &tempopb.SearchResponse{
|
||||
Metrics: metrics,
|
||||
Traces: traceList,
|
||||
},
|
||||
}, sender); err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return err
|
||||
}
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
s.logger.Error("Error receiving message", "err", err)
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
metrics = msg.Metrics
|
||||
traceList = append(traceList, msg.Traces...)
|
||||
traceList = removeDuplicates(traceList)
|
||||
span.SetAttributes(attribute.Int("traces_count", len(traceList)))
|
||||
|
||||
if err := sendResponse(&ExtendedResponse{
|
||||
if err := s.sendResponse(ctx, &ExtendedResponse{
|
||||
State: dataquery.SearchStreamingStateStreaming,
|
||||
SearchResponse: &tempopb.SearchResponse{
|
||||
Metrics: metrics,
|
||||
Traces: traceList,
|
||||
},
|
||||
}, sender); err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return err
|
||||
}
|
||||
}
|
||||
@@ -99,10 +123,14 @@ func (s *Service) processStream(stream tempopb.StreamingQuerier_SearchClient, se
|
||||
return nil
|
||||
}
|
||||
|
||||
func sendResponse(response *ExtendedResponse, sender StreamSender) error {
|
||||
func (s *Service) sendResponse(ctx context.Context, response *ExtendedResponse, sender StreamSender) error {
|
||||
_, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.sendResponse")
|
||||
defer span.End()
|
||||
frame := createResponseDataFrame()
|
||||
|
||||
if response != nil {
|
||||
span.SetAttributes(attribute.Int("trace_count", len(response.Traces)), attribute.String("state", string(response.State)))
|
||||
|
||||
tracesAsJson, err := json.Marshal(response.Traces)
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
@@ -21,7 +21,7 @@ func TestProcessStream_ValidInput_ReturnsNoError(t *testing.T) {
|
||||
service := &Service{}
|
||||
searchClient := &mockStreamer{}
|
||||
streamSender := &mockSender{}
|
||||
err := service.processStream(searchClient, streamSender)
|
||||
err := service.processStream(context.Background(), searchClient, streamSender)
|
||||
if err != nil {
|
||||
t.Errorf("Expected no error, but got %s", err)
|
||||
}
|
||||
@@ -33,7 +33,7 @@ func TestProcessStream_InvalidInput_ReturnsError(t *testing.T) {
|
||||
}
|
||||
searchClient := &mockStreamer{err: errors.New("invalid input")}
|
||||
streamSender := &mockSender{}
|
||||
err := service.processStream(searchClient, streamSender)
|
||||
err := service.processStream(context.Background(), searchClient, streamSender)
|
||||
if err != nil {
|
||||
if !strings.Contains(err.Error(), "invalid input") {
|
||||
t.Errorf("Expected error message to contain 'invalid input', but got %s", err)
|
||||
@@ -109,7 +109,7 @@ func TestProcessStream_ValidInput_ReturnsExpectedOutput(t *testing.T) {
|
||||
},
|
||||
}
|
||||
streamSender := &mockSender{}
|
||||
err := service.processStream(searchClient, streamSender)
|
||||
err := service.processStream(context.Background(), searchClient, streamSender)
|
||||
if err != nil {
|
||||
t.Errorf("Expected no error, but got %s", err)
|
||||
return
|
||||
|
||||
@@ -8,15 +8,24 @@ import (
|
||||
"net/http"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend/tracing"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/data"
|
||||
"github.com/grafana/grafana/pkg/tsdb/tempo/kinds/dataquery"
|
||||
"go.opentelemetry.io/collector/pdata/ptrace"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/codes"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
)
|
||||
|
||||
func (s *Service) getTrace(ctx context.Context, pCtx backend.PluginContext, query backend.DataQuery) (*backend.DataResponse, error) {
|
||||
result := &backend.DataResponse{}
|
||||
refID := query.RefID
|
||||
|
||||
ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.getTrace", trace.WithAttributes(
|
||||
attribute.String("queryType", query.QueryType),
|
||||
))
|
||||
defer span.End()
|
||||
|
||||
model := &dataquery.TempoQuery{}
|
||||
err := json.Unmarshal(query.JSON, model)
|
||||
if err != nil {
|
||||
@@ -34,11 +43,15 @@ func (s *Service) getTrace(ctx context.Context, pCtx backend.PluginContext, quer
|
||||
|
||||
request, err := s.createRequest(ctx, dsInfo, *model.Query, query.TimeRange.From.Unix(), query.TimeRange.To.Unix())
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return result, err
|
||||
}
|
||||
|
||||
resp, err := dsInfo.HTTPClient.Do(request)
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return result, fmt.Errorf("failed get to tempo: %w", err)
|
||||
}
|
||||
|
||||
@@ -55,6 +68,8 @@ func (s *Service) getTrace(ctx context.Context, pCtx backend.PluginContext, quer
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
result.Error = fmt.Errorf("failed to get trace with id: %v Status: %s Body: %s", model.Query, resp.Status, string(body))
|
||||
span.RecordError(result.Error)
|
||||
span.SetStatus(codes.Error, result.Error.Error())
|
||||
return result, nil
|
||||
}
|
||||
|
||||
@@ -62,11 +77,15 @@ func (s *Service) getTrace(ctx context.Context, pCtx backend.PluginContext, quer
|
||||
otTrace, err := pbUnmarshaler.UnmarshalTraces(body)
|
||||
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return &backend.DataResponse{}, fmt.Errorf("failed to convert tempo response to Otlp: %w", err)
|
||||
}
|
||||
|
||||
frame, err := TraceToFrame(otTrace)
|
||||
if err != nil {
|
||||
span.RecordError(err)
|
||||
span.SetStatus(codes.Error, err.Error())
|
||||
return &backend.DataResponse{}, fmt.Errorf("failed to transform trace %v to data frame: %w", model.Query, err)
|
||||
}
|
||||
frame.RefID = refID
|
||||
|
||||
Reference in New Issue
Block a user