From 873f0af6e35432d07d67a9d656a640daee4800ae Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Andr=C3=A9=20Pereira?= Date: Wed, 11 Oct 2023 14:39:34 +0100 Subject: [PATCH] Added spans to search_stream.go --- pkg/tsdb/tempo/search_stream.go | 36 ++++++++++++++++++++++++++++----- 1 file changed, 31 insertions(+), 5 deletions(-) diff --git a/pkg/tsdb/tempo/search_stream.go b/pkg/tsdb/tempo/search_stream.go index bec6b987048..299885477ee 100644 --- a/pkg/tsdb/tempo/search_stream.go +++ b/pkg/tsdb/tempo/search_stream.go @@ -11,6 +11,9 @@ import ( "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" + "go.opentelemetry.io/otel/trace" ) 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 := s.tracer.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 := s.tracer.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,7 +123,9 @@ 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 { + ctx, span := s.tracer.Start(ctx, "datasource.tempo.sendResponse", trace.WithAttributes(attribute.Int("trace_count", len(response.Traces)), attribute.String("state", string(response.State)))) + defer span.End() frame := createResponseDataFrame() if response != nil {