From 88bc65d3839e0f539378e803ec5d501ec0ee13e2 Mon Sep 17 00:00:00 2001 From: Jocelyn Collado-Kuri Date: Wed, 7 Jan 2026 06:54:37 -0800 Subject: [PATCH] add headers to streaming --- pkg/tsdb/tempo/metrics_stream.go | 7 +++++-- pkg/tsdb/tempo/search_stream.go | 9 ++++++--- pkg/tsdb/tempo/stream_handler.go | 10 +++++++++- 3 files changed, 20 insertions(+), 6 deletions(-) diff --git a/pkg/tsdb/tempo/metrics_stream.go b/pkg/tsdb/tempo/metrics_stream.go index 5cefc4c5e33..7836395cb99 100644 --- a/pkg/tsdb/tempo/metrics_stream.go +++ b/pkg/tsdb/tempo/metrics_stream.go @@ -27,6 +27,7 @@ type PartialTempoQuery struct { func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *DatasourceInfo) error { ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.runMetricsStream") defer span.End() + backend.Logger.Warn("runMetricsStream called heeer") response := &backend.DataResponse{} @@ -58,7 +59,7 @@ func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRe } if qrr.GetQuery() == "" { - return backend.DownstreamErrorf("tempo search query cannot be empty") + return backend.DownstreamErrorf("tempo metrics stream search query cannot be empty") } qrr.Start = uint64(backendQuery.TimeRange.From.UnixNano()) @@ -68,7 +69,9 @@ func (s *Service) runMetricsStream(ctx context.Context, req *backend.RunStreamRe // changes or updates, so we have to get it from context. // 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()) - + for key, value := range req.Headers { + ctx = metadata.AppendToOutgoingContext(ctx, key, value) + } if isInstantQuery(tempoQuery.MetricsQueryType) { instantQuery := &tempopb.QueryInstantRequest{ Query: qrr.Query, diff --git a/pkg/tsdb/tempo/search_stream.go b/pkg/tsdb/tempo/search_stream.go index f9394ea382b..6740d2da40a 100644 --- a/pkg/tsdb/tempo/search_stream.go +++ b/pkg/tsdb/tempo/search_stream.go @@ -34,7 +34,7 @@ type StreamSender interface { func (s *Service) runSearchStream(ctx context.Context, req *backend.RunStreamRequest, sender *backend.StreamSender, datasource *DatasourceInfo) error { ctx, span := tracing.DefaultTracer().Start(ctx, "datasource.tempo.runSearchStream") defer span.End() - + backend.Logger.Warn("runSearchStream called heeer") response := &backend.DataResponse{} var backendQuery *backend.DataQuery @@ -56,7 +56,7 @@ func (s *Service) runSearchStream(ctx context.Context, req *backend.RunStreamReq } if sr.GetQuery() == "" { - return backend.DownstreamErrorf("tempo search query cannot be empty") + return backend.DownstreamErrorf("tempo run search query cannot be empty") } sr.Start = uint32(backendQuery.TimeRange.From.Unix()) @@ -66,7 +66,10 @@ func (s *Service) runSearchStream(ctx context.Context, req *backend.RunStreamReq // changes or updates, so we have to get it from context. // 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()) - + // append the rest of the headers + for key, value := range req.Headers { + ctx = metadata.AppendToOutgoingContext(ctx, key, value) + } stream, err := datasource.StreamingClient.Search(ctx, sr) if err != nil { span.RecordError(err) diff --git a/pkg/tsdb/tempo/stream_handler.go b/pkg/tsdb/tempo/stream_handler.go index d3d4a8d7e74..7f9690ba204 100644 --- a/pkg/tsdb/tempo/stream_handler.go +++ b/pkg/tsdb/tempo/stream_handler.go @@ -40,7 +40,15 @@ func (s *Service) PublishStream(_ context.Context, _ *backend.PublishStreamReque func (s *Service) RunStream(ctx context.Context, request *backend.RunStreamRequest, sender *backend.StreamSender) error { s.logger.Debug("New stream call", "path", request.Path) tempoDatasource, err := s.getDSInfo(ctx, request.PluginContext) - + plugin := backend.PluginConfigFromContext(ctx) + opts, err := plugin.DataSourceInstanceSettings.HTTPClientOptions(ctx) + headers := map[string]string{} + for name, values := range opts.Header { + for _, value := range values { + headers[name] = value + } + } + request.Headers = headers if strings.HasPrefix(request.Path, SearchPathPrefix) { if err != nil { return backend.DownstreamErrorf("failed to get datasource information: %w", err)