add headers to streaming

This commit is contained in:
Jocelyn Collado-Kuri
2026-01-07 06:54:37 -08:00
parent 29b04bd2ed
commit 88bc65d383
3 changed files with 20 additions and 6 deletions
+5 -2
View File
@@ -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,
+6 -3
View File
@@ -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)
+9 -1
View File
@@ -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)