Skip to content

Commit 378072a

Browse files
stats/opentelemetry: add e2e test for baggage propagation (#9160)
### Summary Fixes #8520. OpenTelemetry baggage already propagates correctly through the gRPC pipeline when the W3C Baggage propagator is configured — the tracing carriers (`IncomingCarrier`/`OutgoingCarrier`) are generic `TextMapCarrier`s used by the configured propagator's `Inject`/`Extract`, so baggage rides in the `baggage` metadata header exactly like trace context. No production code change is needed; the gap was the lack of a test proving it. This PR adds `TestRPC_BaggagePropagation`, which: - Sets baggage (`key1=value1`, `key2=value2`) on the client context. - Configures `propagation.Baggage{}` in the trace options. - Asserts the **server handler** observes the same baggage, for both a unary and a streaming RPC. ### Verification - Passes with `-race`. - Confirmed meaningful: removing `propagation.Baggage{}` makes the test fail (server sees empty baggage), so it exercises propagation rather than passing vacuously. - `go vet` and `gofmt` clean. RELEASE NOTES: none
1 parent 0e45140 commit 378072a

1 file changed

Lines changed: 93 additions & 0 deletions

File tree

stats/opentelemetry/e2e_test.go

Lines changed: 93 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,7 @@ import (
5252
"google.golang.org/protobuf/types/known/wrapperspb"
5353

5454
"go.opentelemetry.io/otel/attribute"
55+
"go.opentelemetry.io/otel/baggage"
5556
"go.opentelemetry.io/otel/propagation"
5657
"go.opentelemetry.io/otel/sdk/metric"
5758
"go.opentelemetry.io/otel/sdk/metric/metricdata"
@@ -2549,6 +2550,98 @@ func (s) TestRelayContextCollisionTracing(t *testing.T) {
25492550
}
25502551
}
25512552

2553+
// baggageToMap converts the members of an OpenTelemetry baggage into a
2554+
// key/value map for easy comparison in tests.
2555+
func baggageToMap(b baggage.Baggage) map[string]string {
2556+
m := make(map[string]string, b.Len())
2557+
for _, member := range b.Members() {
2558+
m[member.Key()] = member.Value()
2559+
}
2560+
return m
2561+
}
2562+
2563+
// TestRPC_BaggagePropagation verifies that OpenTelemetry baggage set on the
2564+
// client-side context is propagated through the gRPC pipeline and is available
2565+
// in the server handler's context, when the W3C Baggage propagator is
2566+
// configured in the trace options. It exercises both a unary and a streaming
2567+
// RPC.
2568+
func (s) TestRPC_BaggagePropagation(t *testing.T) {
2569+
member1, err := baggage.NewMember("key1", "value1")
2570+
if err != nil {
2571+
t.Fatalf("baggage.NewMember(key1, value1) failed: %v", err)
2572+
}
2573+
member2, err := baggage.NewMember("key2", "value2")
2574+
if err != nil {
2575+
t.Fatalf("baggage.NewMember(key2, value2) failed: %v", err)
2576+
}
2577+
bag, err := baggage.New(member1, member2)
2578+
if err != nil {
2579+
t.Fatalf("baggage.New() failed: %v", err)
2580+
}
2581+
wantBaggage := map[string]string{"key1": "value1", "key2": "value2"}
2582+
2583+
// serverBaggage receives the baggage observed by the server handler for
2584+
// each RPC. It is verified after each RPC, so a buffer of one suffices.
2585+
serverBaggage := make(chan map[string]string, 1)
2586+
ss := &stubserver.StubServer{
2587+
UnaryCallF: func(ctx context.Context, _ *testpb.SimpleRequest) (*testpb.SimpleResponse, error) {
2588+
serverBaggage <- baggageToMap(baggage.FromContext(ctx))
2589+
return &testpb.SimpleResponse{}, nil
2590+
},
2591+
FullDuplexCallF: func(stream testgrpc.TestService_FullDuplexCallServer) error {
2592+
serverBaggage <- baggageToMap(baggage.FromContext(stream.Context()))
2593+
for {
2594+
if _, err := stream.Recv(); err != nil {
2595+
if err == io.EOF {
2596+
return nil
2597+
}
2598+
return err
2599+
}
2600+
}
2601+
},
2602+
}
2603+
2604+
to, _ := defaultTraceOptions(t)
2605+
// Configure the W3C Baggage propagator so baggage is carried in metadata.
2606+
to.TextMapPropagator = propagation.NewCompositeTextMapPropagator(propagation.TraceContext{}, propagation.Baggage{})
2607+
otelOptions := opentelemetry.Options{TraceOptions: *to}
2608+
if err := ss.Start([]grpc.ServerOption{opentelemetry.ServerOption(otelOptions)}, opentelemetry.DialOption(otelOptions)); err != nil {
2609+
t.Fatalf("Error starting endpoint server: %v", err)
2610+
}
2611+
defer ss.Stop()
2612+
2613+
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
2614+
defer cancel()
2615+
ctx = baggage.ContextWithBaggage(ctx, bag)
2616+
2617+
verifyServerBaggage := func(rpc string) {
2618+
t.Helper()
2619+
select {
2620+
case got := <-serverBaggage:
2621+
if diff := cmp.Diff(wantBaggage, got); diff != "" {
2622+
t.Errorf("Baggage received by server for %s RPC mismatch (-want +got):\n%s", rpc, diff)
2623+
}
2624+
case <-ctx.Done():
2625+
t.Fatalf("Timed out waiting for server to receive baggage for %s RPC", rpc)
2626+
}
2627+
}
2628+
2629+
if _, err := ss.Client.UnaryCall(ctx, &testpb.SimpleRequest{}); err != nil {
2630+
t.Fatalf("Unexpected error from UnaryCall: %v", err)
2631+
}
2632+
verifyServerBaggage("unary")
2633+
2634+
stream, err := ss.Client.FullDuplexCall(ctx)
2635+
if err != nil {
2636+
t.Fatalf("ss.Client.FullDuplexCall failed: %v", err)
2637+
}
2638+
stream.CloseSend()
2639+
if _, err = stream.Recv(); err != io.EOF {
2640+
t.Fatalf("stream.Recv received an unexpected error: %v, expected an EOF error", err)
2641+
}
2642+
verifyServerBaggage("streaming")
2643+
}
2644+
25522645
// checkMetricWithMethod verifies that a metric with the specified name contains
25532646
// a data point matching the target grpc.method. It does not poll.
25542647
func checkMetricWithMethod(ctx context.Context, reader *metric.ManualReader, metricName, method string) error {

0 commit comments

Comments
 (0)