From 3c15da8bd21f947317fc9a7acc1eaa17ff5ca3fe Mon Sep 17 00:00:00 2001 From: Stuart McCulloch Date: Mon, 20 Jul 2026 10:30:36 +0100 Subject: [PATCH] Fix W3C baggage propagation across reactor-netty connect hand-off The connect-span path carried only the active AgentSpan across the subscription -> I/O thread hand-off, dropping the rest of the Datadog Context (including baggage). Capture the full Context instead and build the continuation directly via Context.capture(). Co-Authored-By: Claude Sonnet 5 --- .../reactor/netty/CaptureConnectSpan.java | 16 ++- .../reactor/netty/TransferConnectSpan.java | 30 +++--- .../ReactorNettyBaggagePropagationTest.java | 100 ++++++++++++++++++ 3 files changed, 122 insertions(+), 24 deletions(-) create mode 100644 dd-java-agent/instrumentation/reactor-netty-1.0/src/test/java/ReactorNettyBaggagePropagationTest.java diff --git a/dd-java-agent/instrumentation/reactor-netty-1.0/src/main/java/datadog/trace/instrumentation/reactor/netty/CaptureConnectSpan.java b/dd-java-agent/instrumentation/reactor-netty-1.0/src/main/java/datadog/trace/instrumentation/reactor/netty/CaptureConnectSpan.java index 351a37affd3..41daf101e96 100644 --- a/dd-java-agent/instrumentation/reactor-netty-1.0/src/main/java/datadog/trace/instrumentation/reactor/netty/CaptureConnectSpan.java +++ b/dd-java-agent/instrumentation/reactor-netty-1.0/src/main/java/datadog/trace/instrumentation/reactor/netty/CaptureConnectSpan.java @@ -1,8 +1,6 @@ package datadog.trace.instrumentation.reactor.netty; -import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; - -import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.context.Context; import java.util.function.Function; import reactor.core.publisher.Mono; import reactor.netty.Connection; @@ -10,17 +8,17 @@ public class CaptureConnectSpan implements Function, Mono> { - static final String CONNECT_SPAN = "datadog.connect.span"; + static final String CONNECT_CONTEXT = "datadog.connect.context"; @Override public Mono apply(Mono mono) { return mono.contextWrite( - context -> { - final AgentSpan span = activeSpan(); - if (null != span) { - return context.put(CONNECT_SPAN, span); + reactorCtx -> { + final Context context = Context.current(); + if (context != Context.root()) { + return reactorCtx.put(CONNECT_CONTEXT, context); } else { - return context; + return reactorCtx; } }); } diff --git a/dd-java-agent/instrumentation/reactor-netty-1.0/src/main/java/datadog/trace/instrumentation/reactor/netty/TransferConnectSpan.java b/dd-java-agent/instrumentation/reactor-netty-1.0/src/main/java/datadog/trace/instrumentation/reactor/netty/TransferConnectSpan.java index 84a02afa36e..0e58a33b624 100644 --- a/dd-java-agent/instrumentation/reactor-netty-1.0/src/main/java/datadog/trace/instrumentation/reactor/netty/TransferConnectSpan.java +++ b/dd-java-agent/instrumentation/reactor-netty-1.0/src/main/java/datadog/trace/instrumentation/reactor/netty/TransferConnectSpan.java @@ -1,29 +1,29 @@ package datadog.trace.instrumentation.reactor.netty; -import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.captureSpan; import static datadog.trace.instrumentation.netty41.AttributeKeys.CONNECT_PARENT_CONTINUATION_ATTRIBUTE_KEY; -import static datadog.trace.instrumentation.reactor.netty.CaptureConnectSpan.CONNECT_SPAN; +import static datadog.trace.instrumentation.reactor.netty.CaptureConnectSpan.CONNECT_CONTEXT; +import datadog.context.Context; import datadog.context.ContextContinuation; -import datadog.trace.bootstrap.instrumentation.api.AgentSpan; import java.util.function.BiConsumer; import reactor.netty.Connection; import reactor.netty.http.client.HttpClientRequest; public class TransferConnectSpan implements BiConsumer { @Override - public void accept(HttpClientRequest httpClientRequest, Connection connection) { - final AgentSpan span = httpClientRequest.currentContextView().getOrDefault(CONNECT_SPAN, null); - final ContextContinuation continuation = null == span ? null : captureSpan(span); - if (null != continuation) { - ContextContinuation current = - connection - .channel() - .attr(CONNECT_PARENT_CONTINUATION_ATTRIBUTE_KEY) - .getAndSet(continuation); - if (null != current) { - current.release(); - } + public void accept(HttpClientRequest clientRequest, Connection connection) { + final Context context = clientRequest.currentContextView().getOrDefault(CONNECT_CONTEXT, null); + if (null == context) { + return; + } + ContextContinuation newContinuation = context.capture(); + ContextContinuation oldContinuation = + connection + .channel() + .attr(CONNECT_PARENT_CONTINUATION_ATTRIBUTE_KEY) + .getAndSet(newContinuation); + if (null != oldContinuation) { + oldContinuation.release(); } } } diff --git a/dd-java-agent/instrumentation/reactor-netty-1.0/src/test/java/ReactorNettyBaggagePropagationTest.java b/dd-java-agent/instrumentation/reactor-netty-1.0/src/test/java/ReactorNettyBaggagePropagationTest.java new file mode 100644 index 00000000000..c366d14d2d7 --- /dev/null +++ b/dd-java-agent/instrumentation/reactor-netty-1.0/src/test/java/ReactorNettyBaggagePropagationTest.java @@ -0,0 +1,100 @@ +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import com.sun.net.httpserver.HttpServer; +import datadog.context.Context; +import datadog.context.ContextScope; +import datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.bootstrap.instrumentation.api.AgentScope; +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.bootstrap.instrumentation.api.AgentTracer; +import datadog.trace.bootstrap.instrumentation.api.Baggage; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.time.Duration; +import java.util.Collections; +import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import reactor.netty.http.client.HttpClient; + +/** + * Regression test for the W3C baggage header not propagating on outgoing Reactor Netty requests. + * + *

The bug: the connect-span path carried only {@code activeSpan()} across the subscription -> + * I/O thread hand-off, dropping the rest of the Datadog {@code Context} (including baggage). The + * outgoing request then had no baggage to inject, so the {@code baggage} header was silently + * skipped. This test sets a baggage item in the active context, makes one outgoing request, and + * asserts the server received the {@code baggage} header. + * + *

It is a coupling test: producer (sets baggage) -> Reactor Netty carrier -> consumer + * (injects the header). The failure it guards against lives in the hand-off between integrations, + * not in any one of them, so per-integration tests would not catch it. + */ +class ReactorNettyBaggagePropagationTest extends AbstractInstrumentationTest { + + private static HttpServer mockServer; + private static String baseUrl; + private static final AtomicReference capturedBaggage = new AtomicReference<>(); + + @BeforeAll + static void startServer() throws IOException { + mockServer = HttpServer.create(new InetSocketAddress("localhost", 0), 0); + mockServer.createContext( + "/capture", + exchange -> { + capturedBaggage.set(exchange.getRequestHeaders().getFirst("baggage")); + byte[] body = "ok".getBytes(); + exchange.sendResponseHeaders(200, body.length); + exchange.getResponseBody().write(body); + exchange.close(); + }); + mockServer.setExecutor(Executors.newCachedThreadPool()); + mockServer.start(); + baseUrl = + "http://" + + mockServer.getAddress().getHostString() + + ":" + + mockServer.getAddress().getPort(); + } + + @AfterAll + static void stopServer() { + if (mockServer != null) { + mockServer.stop(0); + mockServer = null; + } + } + + @Test + void baggageHeaderPropagatedOnOutgoingRequest() { + capturedBaggage.set(null); + Baggage baggage = Baggage.create(Collections.singletonMap("user.id", "abc123")); + + AgentSpan span = AgentTracer.startSpan("test", "parent"); + try (AgentScope spanScope = AgentTracer.activateSpan(span)) { + // Active context now carries both the span and the baggage — the exact shape the connect-span + // path must carry across the subscription -> I/O thread hand-off. + try (ContextScope baggageScope = Context.current().with(baggage).attach()) { + HttpClient.create() + .get() + .uri(baseUrl + "/capture") + .response() + .block(Duration.ofSeconds(10)); + } + } finally { + span.finish(); + } + + String header = capturedBaggage.get(); + assertNotNull( + header, + "outgoing request must carry a W3C 'baggage' header when baggage is in the active context;" + + " null means the connect-span path dropped the context (carried only the span)"); + assertTrue( + header.contains("user.id=abc123"), + "baggage header should contain the propagated item, was: " + header); + } +}