diff --git a/pkg/tsdb/tempo/metrics_stream.go b/pkg/tsdb/tempo/metrics_stream.go index 1c4bba6f713..85fac2df83b 100644 --- a/pkg/tsdb/tempo/metrics_stream.go +++ b/pkg/tsdb/tempo/metrics_stream.go @@ -20,6 +20,10 @@ import ( const MetricsPathPrefix = "metrics/" +type PartialTempoQuery struct { + MetricsQueryType *dataquery.MetricsQueryType +} + func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *Datasource) error { ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.runMetricsStream") defer span.End() @@ -35,6 +39,15 @@ func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRe return err } + tempoQuery := &PartialTempoQuery{} + err = json.Unmarshal(req.Data, tempoQuery) + 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 + } + var qrr *tempopb.QueryRangeRequest err = json.Unmarshal(req.Data, &qrr) if err != nil { @@ -56,6 +69,24 @@ func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRe // Ideally this would be pushed higher, so it's set once for all rpc calls, but we have only one now. ctx = metadata.AppendToOutgoingContext(ctx, "User-Agent", backend.UserAgentFromContext(ctx).String()) + if isInstantQuery(tempoQuery.MetricsQueryType) { + instantQuery := &tempopb.QueryInstantRequest{ + Query: qrr.Query, + Start: qrr.Start, + End: qrr.End, + } + + stream, err := datasource.StreamingClient.MetricsQueryInstant(ctx, instantQuery) + if err != nil { + span.RecordError(err) + span.SetStatus(codes.Error, err.Error()) + s.logger.Error("Error Search()", "err", err) + return err + } + + return s.processInstantMetricsStream(ctx, stream, sender) + } + stream, err := datasource.StreamingClient.MetricsQueryRange(ctx, qrr) if err != nil { span.RecordError(err) @@ -101,3 +132,38 @@ func (s *Service) processMetricsStream(ctx context.Context, query string, stream return nil } + +func (s *Service) processInstantMetricsStream(ctx context.Context, stream tempopb.StreamingQuerier_MetricsQueryInstantClient, sender StreamSender) error { + ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.processStream") + defer span.End() + messageCount := 0 + for { + msg, err := stream.Recv() + messageCount++ + span.SetAttributes(attribute.Int("message_count", messageCount)) + if errors.Is(err, io.EOF) { + if err := s.sendResponse(ctx, nil, nil, dataquery.SearchStreamingStateDone, 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 + } + + transformed := traceql.TransformInstantMetricsResponse(*msg) + + if err := s.sendResponse(ctx, transformed, msg.Metrics, dataquery.SearchStreamingStateStreaming, sender); err != nil { + span.RecordError(err) + span.SetStatus(codes.Error, err.Error()) + return err + } + } + + return nil +} diff --git a/pkg/tsdb/tempo/traceql/metrics.go b/pkg/tsdb/tempo/traceql/metrics.go index ec057c80cca..f2f2fe778c9 100644 --- a/pkg/tsdb/tempo/traceql/metrics.go +++ b/pkg/tsdb/tempo/traceql/metrics.go @@ -9,7 +9,6 @@ import ( "time" "github.com/grafana/grafana-plugin-sdk-go/data" - "github.com/grafana/grafana/pkg/tsdb/tempo/kinds/dataquery" "github.com/grafana/tempo/pkg/tempopb" v1 "github.com/grafana/tempo/pkg/tempopb/common/v1" ) @@ -61,7 +60,7 @@ func TransformMetricsResponse(query string, resp tempopb.QueryRangeResponse) []* return append(frames, exemplarFrames...) } -func TransformInstantMetricsResponse(query *dataquery.TempoQuery, resp tempopb.QueryInstantResponse) []*data.Frame { +func TransformInstantMetricsResponse(resp tempopb.QueryInstantResponse) []*data.Frame { frames := make([]*data.Frame, len(resp.Series)) for i, series := range resp.Series { diff --git a/pkg/tsdb/tempo/traceql/metrics_test.go b/pkg/tsdb/tempo/traceql/metrics_test.go index 0cf956f88a2..e7c9e06e404 100644 --- a/pkg/tsdb/tempo/traceql/metrics_test.go +++ b/pkg/tsdb/tempo/traceql/metrics_test.go @@ -5,7 +5,6 @@ import ( "time" "github.com/grafana/grafana-plugin-sdk-go/data" - "github.com/grafana/grafana/pkg/tsdb/tempo/kinds/dataquery" "github.com/grafana/tempo/pkg/tempopb" v1 "github.com/grafana/tempo/pkg/tempopb/common/v1" "github.com/stretchr/testify/assert" @@ -113,7 +112,6 @@ func TestTransformMetricsResponse_MultipleSeries(t *testing.T) { } func TestTransformInstantMetricsResponse(t *testing.T) { - query := &dataquery.TempoQuery{} resp := tempopb.QueryInstantResponse{ Series: []*tempopb.InstantSeries{ { @@ -123,7 +121,7 @@ func TestTransformInstantMetricsResponse(t *testing.T) { }, } - frames := TransformInstantMetricsResponse(query, resp) + frames := TransformInstantMetricsResponse(resp) assert.Len(t, frames, 1) frame := frames[0] diff --git a/pkg/tsdb/tempo/traceql_query.go b/pkg/tsdb/tempo/traceql_query.go index 43eea31f04c..1b5d31edad5 100644 --- a/pkg/tsdb/tempo/traceql_query.go +++ b/pkg/tsdb/tempo/traceql_query.go @@ -97,7 +97,7 @@ func (s *Service) runTraceQlQueryMetrics(ctx context.Context, pCtx backend.Plugi return res, err } - frames := traceql.TransformInstantMetricsResponse(tempoQuery, queryResponse) + frames := traceql.TransformInstantMetricsResponse(queryResponse) result.Frames = frames } else { var queryResponse tempopb.QueryRangeResponse diff --git a/public/app/plugins/datasource/tempo/datasource.ts b/public/app/plugins/datasource/tempo/datasource.ts index 63cb1bba4ea..b7c3af81613 100644 --- a/public/app/plugins/datasource/tempo/datasource.ts +++ b/public/app/plugins/datasource/tempo/datasource.ts @@ -377,6 +377,7 @@ export class TempoDatasource extends DataSourceWithBackend 0) { reportInteraction('grafana_traces_service_graph_queried', { datasourceType: 'tempo', diff --git a/public/app/plugins/datasource/tempo/streaming.ts b/public/app/plugins/datasource/tempo/streaming.ts index fdc66ea7bf8..3ea011d3185 100644 --- a/public/app/plugins/datasource/tempo/streaming.ts +++ b/public/app/plugins/datasource/tempo/streaming.ts @@ -20,7 +20,7 @@ import { import { cloneQueryResponse, combineResponses } from '@grafana/o11y-ds-frontend'; import { getGrafanaLiveSrv } from '@grafana/runtime'; -import { SearchStreamingState } from './dataquery.gen'; +import { MetricsQueryType, SearchStreamingState } from './dataquery.gen'; import { DEFAULT_SPSS, TempoDatasource } from './datasource'; import { formatTraceQLResponse } from './resultTransformer'; import { SearchMetrics, TempoJsonData, TempoQuery } from './types'; @@ -177,7 +177,8 @@ export function doTempoMetricsStreaming( if (!curr) { return acc; } - if (!acc) { + // If the query is an instant query, we always want the latest result. + if (!acc || query.metricsQueryType === MetricsQueryType.Instant) { return cloneQueryResponse(curr); } return mergeFrames(acc, curr); diff --git a/public/app/plugins/datasource/tempo/traceql/TempoQueryBuilderOptions.tsx b/public/app/plugins/datasource/tempo/traceql/TempoQueryBuilderOptions.tsx index c591455d578..96c577d0c6e 100644 --- a/public/app/plugins/datasource/tempo/traceql/TempoQueryBuilderOptions.tsx +++ b/public/app/plugins/datasource/tempo/traceql/TempoQueryBuilderOptions.tsx @@ -37,6 +37,7 @@ export const TempoQueryBuilderOptions = React.memo( const styles = useStyles2(getStyles); const [isOpen, toggleOpen] = useToggle(false); const isAlerting = app === CoreApp.UnifiedAlerting; + const isMetricsStreamingEnabled = metricsStreaming && !isAlerting; if (!query.hasOwnProperty('limit')) { query.limit = DEFAULT_LIMIT; @@ -93,7 +94,7 @@ export const TempoQueryBuilderOptions = React.memo( `Step: ${query.step || 'auto'}`, `Type: ${query.metricsQueryType === MetricsQueryType.Range ? 'Range' : 'Instant'}`, '|', - `Streaming: ${metricsStreaming ? 'Enabled' : 'Disabled'}`, + `Streaming: ${isMetricsStreamingEnabled ? 'Enabled' : 'Disabled'}`, // `Exemplars: ${query.exemplars !== undefined ? query.exemplars : 'auto'}`, ]; @@ -179,7 +180,7 @@ export const TempoQueryBuilderOptions = React.memo( } tooltipInteractive> -
{metricsStreaming ? 'Enabled' : 'Disabled'}
+
{isMetricsStreamingEnabled ? 'Enabled' : 'Disabled'}
{/*