Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,7 @@
* [BUGFIX] Compactor: Fix spurious `bucket operation fail after retries` error logs emitted during partial block cleanup. #7749
* [BUGFIX] Alertmanager: Fix panic in `validateAlertmanagerConfig` when receiver config traversal encounters nil interface values. #7751
* [BUGFIX] Parquet Converter: Fix `auto_forget_delay` having no effect. The ring lifecycler was created without the auto-forget delegate, so unhealthy instances were never automatically removed from the ring. #7752
* [BUGFIX] Ingester: Fix tracing of write requests when `-distributor.use-stream-push` is enabled. Every push received over a `PushStream` connection was attached to the span of the connection itself, gluing all of them into a single trace growing for as long as the connection lived. #7753

## 1.21.1 2026-06-04

Expand Down
49 changes: 49 additions & 0 deletions pkg/cortexpb/carrier.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
package cortexpb

import (
"github.com/opentracing/opentracing-go"
)

// StreamWriteRequestCarrier is used to transfer trace
// information from/to a StreamWriteRequest.
type StreamWriteRequestCarrier StreamWriteRequest

func (c *StreamWriteRequestCarrier) Set(key, val string) {
c.TraceContext = append(c.TraceContext, &Header{
Key: key,
Values: []string{val},
})
}

func (c *StreamWriteRequestCarrier) ForeachKey(handler func(key, val string) error) error {
for _, h := range c.TraceContext {
for _, v := range h.Values {
if err := handler(h.Key, v); err != nil {
return err
}
}
}
return nil
}

// InjectSpanIntoStreamWriteRequest makes req carry the trace context of span.
func InjectSpanIntoStreamWriteRequest(tracer opentracing.Tracer, span opentracing.Span, req *StreamWriteRequest) error {
if tracer == nil || span == nil {
return nil
}

return tracer.Inject(span.Context(), opentracing.HTTPHeaders, (*StreamWriteRequestCarrier)(req))
}

// GetParentSpanForStreamWriteRequest returns the span context req carries.
func GetParentSpanForStreamWriteRequest(tracer opentracing.Tracer, req *StreamWriteRequest) (opentracing.SpanContext, error) {
if tracer == nil || len(req.TraceContext) == 0 {
return nil, nil
}

extracted, err := tracer.Extract(opentracing.HTTPHeaders, (*StreamWriteRequestCarrier)(req))
if err == opentracing.ErrSpanContextNotFound {
err = nil
}
return extracted, err
}
26 changes: 26 additions & 0 deletions pkg/cortexpb/carrier_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
package cortexpb

import (
"testing"

"github.com/opentracing/opentracing-go/mocktracer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

func TestStreamWriteRequestCarrier_RoundTrip(t *testing.T) {
tracer := mocktracer.New()
span := tracer.StartSpan("Distributor.Push")
defer span.Finish()

req := &StreamWriteRequest{TenantID: "user-1"}
require.NoError(t, InjectSpanIntoStreamWriteRequest(tracer, span, req))
require.NotEmpty(t, req.TraceContext, "the trace context must be carried by the message itself")

extracted, err := GetParentSpanForStreamWriteRequest(tracer, req)
require.NoError(t, err)
require.NotNil(t, extracted)

assert.Equal(t, span.Context().(mocktracer.MockSpanContext).TraceID, extracted.(mocktracer.MockSpanContext).TraceID)
assert.Equal(t, span.Context().(mocktracer.MockSpanContext).SpanID, extracted.(mocktracer.MockSpanContext).SpanID)
}
591 changes: 478 additions & 113 deletions pkg/cortexpb/cortex.pb.go

Large diffs are not rendered by default.

7 changes: 7 additions & 0 deletions pkg/cortexpb/cortex.proto
Original file line number Diff line number Diff line change
Expand Up @@ -117,9 +117,16 @@ message MetadataV2 {
uint32 unit_ref = 4;
}

message Header {
string Key = 1;
repeated string Values = 2;
}

message StreamWriteRequest {
string TenantID = 1;
WriteRequest Request = 2;
// TraceContext carries the tracing headers of the write request this message belongs to.
repeated Header TraceContext = 3;

MessageWithBufRef Ref = 1000 [(gogoproto.embed) = true, (gogoproto.customtype) = "MessageWithBufRef", (gogoproto.nullable) = false]; //set intentionally high to keep WriteRequest compatible with upstream Prometheus
}
Expand Down
5 changes: 5 additions & 0 deletions pkg/ingester/client/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"time"

"github.com/go-kit/log"
"github.com/opentracing/opentracing-go"
"github.com/pkg/errors"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promauto"
Expand Down Expand Up @@ -115,6 +116,10 @@ func (c *closableHealthAndIngesterClient) PushStreamConnection(ctx context.Conte
Request: in,
}

if err := cortexpb.InjectSpanIntoStreamWriteRequest(opentracing.GlobalTracer(), opentracing.SpanFromContext(ctx), streamReq); err != nil {
return nil, err
}

job := &streamWriteJob{
req: streamReq,
sendDone: make(chan struct{}),
Expand Down
18 changes: 17 additions & 1 deletion pkg/ingester/ingester.go
Original file line number Diff line number Diff line change
Expand Up @@ -1833,6 +1833,10 @@ func (i *Ingester) Push(ctx context.Context, req *cortexpb.WriteRequest) (*corte

func (i *Ingester) PushStream(srv client.Ingester_PushStreamServer) error {
ctx := srv.Context()
tracer := opentracing.GlobalTracer()
// Drop the span of the connection, which lives for as long as the connection does:
// requests are traced from the trace context carried by their own message.
baseCtx := opentracing.ContextWithSpan(ctx, nil)
for {
select {
case <-ctx.Done():
Expand Down Expand Up @@ -1860,7 +1864,16 @@ func (i *Ingester) PushStream(srv client.Ingester_PushStreamServer) error {
}
}

pushCtx := user.InjectOrgID(ctx, req.TenantID)
pushCtx := user.InjectOrgID(baseCtx, req.TenantID)

// Trace the push as part of its own write request, not of the connection it was
// sent on. Ignore errors: with no parent span we simply don't create one.
parentSpanContext, _ := cortexpb.GetParentSpanForStreamWriteRequest(tracer, req)
var requestSpan opentracing.Span
if parentSpanContext != nil {
requestSpan, pushCtx = opentracing.StartSpanFromContextWithTracer(pushCtx, tracer, "Ingester.PushStreamRequest", opentracing.ChildOf(parentSpanContext))
}

resp, err := i.Push(pushCtx, req.Request)
if resp == nil {
resp = &cortexpb.WriteResponse{}
Expand All @@ -1877,6 +1890,9 @@ func (i *Ingester) PushStream(srv client.Ingester_PushStreamServer) error {
}
err = srv.Send(resp)
req.Free()
if requestSpan != nil {
requestSpan.Finish()
}
if err != nil {
level.Error(logutil.WithContext(ctx, i.logger)).Log("msg", "error sending from PushStream", "err", err)
}
Expand Down
72 changes: 72 additions & 0 deletions pkg/ingester/pushstream_tracing_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
package ingester

import (
"context"
"os"
"testing"

"github.com/opentracing/opentracing-go"
"github.com/opentracing/opentracing-go/mocktracer"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/weaveworks/common/user"

"github.com/cortexproject/cortex/pkg/cortexpb"
)

// testTracer records the spans started by the tests. opentracing.SetGlobalTracer is an
// unsynchronized assignment and the ingesters keep reading the global tracer from their ring
// heartbeat, so it is registered once, before any test runs.
var testTracer = mocktracer.New()

func TestMain(m *testing.M) {
opentracing.SetGlobalTracer(testTracer)
os.Exit(m.Run())
}

// streamWriteReqWithTrace builds a request carrying its own trace context, as the client does.
func streamWriteReqWithTrace(t *testing.T, tenantID, metricName string) (*cortexpb.StreamWriteRequest, int) {
t.Helper()

span := testTracer.StartSpan("Distributor.Push")

req := &cortexpb.StreamWriteRequest{TenantID: tenantID, Request: makeWriteReq(metricName)}
require.NoError(t, cortexpb.InjectSpanIntoStreamWriteRequest(testTracer, span, req))
return req, span.Context().(mocktracer.MockSpanContext).TraceID
}

// TestPushStream_TracesEachRequestSeparately checks that each push is traced as part of the
// write request it belongs to, and not of the connection it happened to arrive on.
func TestPushStream_TracesEachRequestSeparately(t *testing.T) {
testTracer.Reset()
ing := newTestIngester(t)

// The interceptor puts the span of the connection in the stream context.
streamCtx := user.InjectOrgID(context.Background(), "ingester-127.0.0.1-9095-stream-push-worker-0")
streamSpan := testTracer.StartSpan("/cortex.Ingester/PushStream")
streamCtx = opentracing.ContextWithSpan(streamCtx, streamSpan)

tracedReq, tracedTraceID := streamWriteReqWithTrace(t, "user-1", "metric_one")
// A client that predates this propagation carries no trace context.
legacyReq := &cortexpb.StreamWriteRequest{TenantID: "user-2", Request: makeWriteReq("metric_two")}

srv := &pushStreamServer{ctx: streamCtx, requests: []*cortexpb.StreamWriteRequest{tracedReq, legacyReq}}
require.NoError(t, ing.PushStream(srv))

byTrace := map[int][]string{}
for _, span := range testTracer.FinishedSpans() {
byTrace[span.SpanContext.TraceID] = append(byTrace[span.SpanContext.TraceID], span.OperationName)
}

// A push joins the trace of its own write request.
assert.Equal(t, []string{"Ingester.Push", "Ingester.PushStreamRequest"}, byTrace[tracedTraceID])

var parentless []*mocktracer.MockSpan
for _, span := range testTracer.FinishedSpans() {
if span.OperationName == "Ingester.Push" && span.ParentID == 0 {
parentless = append(parentless, span)
}
}
require.Len(t, parentless, 1)
assert.NotContains(t, byTrace, streamSpan.Context().(mocktracer.MockSpanContext).TraceID)
}