Added spans to search_stream.go

This commit is contained in:
André Pereira
2023-10-11 14:39:34 +01:00
parent bf5259e0b2
commit 873f0af6e3
+31 -5
View File
@@ -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 {