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
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,9 @@ dependencies {

testImplementation group: 'io.projectreactor.netty', name: 'reactor-netty-http', version: '1.0.0'
testImplementation project(':dd-java-agent:instrumentation:netty:netty-4.1')
testImplementation project(':dd-java-agent:instrumentation:reactor-core-3.1')
testImplementation project(':dd-java-agent:instrumentation:reactive-streams-1.0')
testImplementation libs.mokito.core
testRuntimeOnly project(':dd-java-agent:instrumentation:reactor-core-3.1')
testRuntimeOnly project(':dd-java-agent:instrumentation:reactive-streams-1.0')

latestDepTestImplementation group: 'io.projectreactor.netty', name: 'reactor-netty-http', version: '+'
latestDepTestImplementation project(':dd-java-agent:instrumentation:netty:netty-4.1')
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -139,7 +139,7 @@ org.junit.platform:junit-platform-runner:1.14.1=latestDepTestRuntimeClasspath,te
org.junit.platform:junit-platform-suite-api:1.14.1=latestDepTestRuntimeClasspath,testRuntimeClasspath
org.junit.platform:junit-platform-suite-commons:1.14.1=latestDepTestRuntimeClasspath,testRuntimeClasspath
org.junit:junit-bom:5.14.1=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath
org.mockito:mockito-core:4.4.0=latestDepTestRuntimeClasspath,testRuntimeClasspath
org.mockito:mockito-core:4.4.0=latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath
org.objenesis:objenesis:3.3=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath
org.opentest4j:opentest4j:1.3.0=latestDepTestCompileClasspath,latestDepTestRuntimeClasspath,testCompileClasspath,testRuntimeClasspath
org.ow2.asm:asm-analysis:9.10.1=spotbugs
Expand Down
Original file line number Diff line number Diff line change
@@ -1,14 +1,24 @@
package datadog.trace.instrumentation.reactor.netty;

import datadog.context.Context;
import datadog.trace.bootstrap.ContextStore;
import datadog.trace.bootstrap.InstrumentationContext;
import net.bytebuddy.asm.Advice;
import reactor.netty.http.client.HttpClient;
import reactor.netty.http.client.HttpClientRequest;

public class AfterConstructorAdvice {
@Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class)
public static void onExit(
@Advice.Thrown Throwable throwable, @Advice.Return(readOnly = false) HttpClient client) {
if (null == throwable) {
client = client.mapConnect(new CaptureConnectSpan()).doOnRequest(new TransferConnectSpan());
ContextStore<HttpClientRequest, Context> requestContexts =
InstrumentationContext.get(HttpClientRequest.class, Context.class);
client =
client
.mapConnect(new CaptureConnectSpan())
.doOnRequest(new TransferConnectSpan(requestContexts))
.doOnDisconnected(new ClearRequestContext(requestContexts));
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
package datadog.trace.instrumentation.reactor.netty;

import static datadog.trace.instrumentation.netty41.AttributeKeys.CLIENT_PARENT_ATTRIBUTE_KEY;
import static datadog.trace.instrumentation.netty41.AttributeKeys.CONTEXT_ATTRIBUTE_KEY;

import datadog.context.Context;
import datadog.trace.bootstrap.ContextStore;
import datadog.trace.bootstrap.instrumentation.api.AgentSpan;
import io.netty.channel.Channel;
import io.netty.util.Attribute;
import java.util.function.Consumer;
import reactor.netty.Connection;
import reactor.netty.channel.ChannelOperations;
import reactor.netty.http.client.HttpClientRequest;

/** Clears request context when a connection is released or disconnected. */
public class ClearRequestContext implements Consumer<Connection> {
private final ContextStore<HttpClientRequest, Context> requestContexts;

public ClearRequestContext(ContextStore<HttpClientRequest, Context> requestContexts) {
this.requestContexts = requestContexts;
}

@Override
public void accept(Connection connection) {
if (!(connection instanceof HttpClientRequest)) {
return;
}
Context requestContext = requestContexts.remove((HttpClientRequest) connection);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Codex recommends to arrange cleanup after Netty finishes the span - I think the suggestion makes sense from the description below:

This section "removes the saved request context before checking whether the client span is still live. With an outgoing request-body failure, the disconnect callback runs while that span is live and skips clearing the attributes. Netty subsequently finishes the client span and restores the parent context. Later cleanup callbacks cannot remove it because the saved context is gone or the connection is no longer an HttpClientRequest. Subsequent channel events can still reactivate that parent."

Channel channel = connection.channel();
ChannelOperations<?, ?> current = ChannelOperations.get(channel);
if (current != null && current != connection) {
return; // The pooled channel already belongs to another request.
}
AgentSpan parent = AgentSpan.fromContext(requestContext);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

IIUC there may not always be an AgentSpan when context is captured, e.g. with only baggage? Should we handle saved contexts without spans too?

if (parent == null) {
return;
}
Attribute<Context> attribute = channel.attr(CONTEXT_ATTRIBUTE_KEY);
Context stored = attribute.get();
// Leave a live client span for Netty's error/close handler to finish.
if (stored != null
&& AgentSpan.fromContext(stored) == parent
&& attribute.compareAndSet(stored, null)) {
channel.attr(CLIENT_PARENT_ATTRIBUTE_KEY).compareAndSet(parent, null);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,17 +2,20 @@

import static datadog.trace.agent.tooling.bytebuddy.matcher.ClassLoaderMatchers.hasClassNamed;
import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.namedOneOf;
import static java.util.Collections.singletonMap;
import static net.bytebuddy.matcher.ElementMatchers.isMethod;
import static net.bytebuddy.matcher.ElementMatchers.isStatic;

import com.google.auto.service.AutoService;
import datadog.context.Context;
import datadog.trace.agent.tooling.Instrumenter;
import datadog.trace.agent.tooling.InstrumenterModule;
import java.util.Map;
import net.bytebuddy.matcher.ElementMatcher;

/**
* This instrumentation only supports transfer of the active span at connection time to the
* underlying Netty Channel.
* Transfers request context to the underlying Netty channel and clears it when Reactor releases the
* connection.
*
* <p>Based on the OpenTelemetry Reactor Netty instrumentation.
* https://github.com/open-telemetry/opentelemetry-java-instrumentation/blob/main/instrumentation/reactor-netty/reactor-netty-1.0/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/reactornetty/v1_0/HttpClientInstrumentation.java
Expand All @@ -31,12 +34,18 @@ public ElementMatcher.Junction<ClassLoader> classLoaderMatcher() {
return hasClassNamed("reactor.netty.transport.AddressUtils");
}

@Override
public Map<String, String> contextStore() {
return singletonMap("reactor.netty.http.client.HttpClientRequest", Context.class.getName());
}

@Override
public String[] helperClassNames() {
return new String[] {
"datadog.trace.instrumentation.netty41.AttributeKeys",
packageName + ".CaptureConnectSpan",
packageName + ".TransferConnectSpan",
packageName + ".ClearRequestContext",
};
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,17 +5,26 @@

import datadog.context.Context;
import datadog.context.ContextContinuation;
import datadog.trace.bootstrap.ContextStore;
import java.util.function.BiConsumer;
import reactor.netty.Connection;
import reactor.netty.http.client.HttpClientRequest;

public class TransferConnectSpan implements BiConsumer<HttpClientRequest, Connection> {
private final ContextStore<HttpClientRequest, Context> requestContexts;

public TransferConnectSpan(ContextStore<HttpClientRequest, Context> requestContexts) {
this.requestContexts = requestContexts;
}

@Override
public void accept(HttpClientRequest clientRequest, Connection connection) {
final Context context = clientRequest.currentContextView().getOrDefault(CONNECT_CONTEXT, null);
if (null == context) {
return;
}
// Reactor clears its owner context before the pool-release callback runs.
requestContexts.put(clientRequest, context);
ContextContinuation newContinuation = context.capture();
ContextContinuation oldContinuation =
connection
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,6 @@
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;

import com.sun.net.httpserver.HttpServer;
Expand All @@ -8,18 +10,25 @@
import datadog.trace.bootstrap.instrumentation.api.AgentSpan;
import datadog.trace.bootstrap.instrumentation.api.AgentTracer;
import datadog.trace.bootstrap.instrumentation.api.Baggage;
import io.netty.channel.Channel;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.Collections;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
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;
import reactor.netty.resources.ConnectionProvider;
import reactor.netty.resources.LoopResources;

/**
* Regression test for the W3C baggage header not propagating on outgoing Reactor Netty requests.
Expand Down Expand Up @@ -76,6 +85,81 @@ static void stopServer() {
}
}

@Test
void idlePooledChannelDoesNotReactivateRequestContext() throws Exception {
ConnectionProvider connections = ConnectionProvider.create("idle-context-test", 1);
LoopResources eventLoops = LoopResources.create("idle-context-test");
AtomicReference<Channel> channel = new AtomicReference<>();
AtomicReference<AgentSpan> eventSpan = new AtomicReference<>();
Object idleEvent = new Object();
HttpClient client =
HttpClient.create(connections)
.runOn(eventLoops)
.doOnConnected(
connection -> {
channel.set(connection.channel());
connection
.channel()
.pipeline()
.addLast(
new ChannelInboundHandlerAdapter() {
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object event)
throws Exception {
if (event == idleEvent) {
eventSpan.set(AgentTracer.activeSpan());
}
super.userEventTriggered(ctx, event);
}
});
});
Channel firstChannel = null;
try {
for (int request = 0; request < 2; request++) {
CountDownLatch released = new CountDownLatch(1);
AtomicReference<AgentSpan> responseSpan = new AtomicReference<>();
AgentSpan parent = AgentTracer.startSpan("test", "parent-" + request);
try (ContextScope scope = AgentTracer.activateSpan(parent)) {
client
.doOnResponse((response, connection) -> responseSpan.set(AgentTracer.activeSpan()))
.doOnDisconnected(connection -> released.countDown())
.get()
.uri(baseUrl + "/capture")
.responseContent()
.aggregate()
.asString()
.block(Duration.ofSeconds(10));
} finally {
parent.finish();
}
assertSame(parent, responseSpan.get());
assertTrue(released.await(10, TimeUnit.SECONDS));
Channel currentChannel = channel.get();
if (firstChannel == null) {
firstChannel = currentChannel;
} else {
assertSame(
firstChannel, currentChannel, "the second request must reuse the pooled channel");
}
currentChannel
.eventLoop()
.submit(
() -> {
assertNull(AgentTracer.activeSpan());
currentChannel.pipeline().fireUserEventTriggered(idleEvent);
})
.get(10, TimeUnit.SECONDS);
assertNull(eventSpan.get(), "an idle channel event must not activate the previous request");
}
} finally {
try {
connections.disposeLater().block(Duration.ofSeconds(10));
} finally {
eventLoops.disposeLater(Duration.ZERO, Duration.ofSeconds(5)).block(Duration.ofSeconds(10));
}
}
}

@Test
void baggageHeaderPropagatedOnOutgoingRequest() {
Baggage baggage = Baggage.create(Collections.singletonMap("user.id", "abc123"));
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
package datadog.trace.instrumentation.reactor.netty;

import static datadog.trace.instrumentation.netty41.AttributeKeys.CLIENT_PARENT_ATTRIBUTE_KEY;
import static datadog.trace.instrumentation.netty41.AttributeKeys.CONTEXT_ATTRIBUTE_KEY;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.CALLS_REAL_METHODS;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import static org.mockito.Mockito.withSettings;

import datadog.context.Context;
import datadog.trace.bootstrap.ContextStore;
import datadog.trace.bootstrap.instrumentation.api.AgentSpan;
import io.netty.channel.Channel;
import io.netty.util.DefaultAttributeMap;
import org.junit.jupiter.api.Test;
import reactor.netty.Connection;
import reactor.netty.channel.ChannelOperations;
import reactor.netty.http.client.HttpClientRequest;

class ClearRequestContextTest {
@Test
void clearsReleasedRequestContext() {
Fixture fixture = new Fixture();

fixture.cleanup.accept(fixture.connection);

assertNull(fixture.channel.attr(CONTEXT_ATTRIBUTE_KEY).get());
assertNull(fixture.channel.attr(CLIENT_PARENT_ATTRIBUTE_KEY).get());
}

@Test
void leavesUnfinishedClientSpanForNettyCleanup() {
Fixture fixture = new Fixture();
Context clientContext = Context.root().with(mock(AgentSpan.class, CALLS_REAL_METHODS));
fixture.channel.attr(CONTEXT_ATTRIBUTE_KEY).set(clientContext);

fixture.cleanup.accept(fixture.connection);

assertSame(clientContext, fixture.channel.attr(CONTEXT_ATTRIBUTE_KEY).get());
assertSame(fixture.parent, fixture.channel.attr(CLIENT_PARENT_ATTRIBUTE_KEY).get());
}

@Test
void preservesReusedChannelEvenWhenNextRequestHasTheSameParent() {
Fixture fixture = new Fixture();
ChannelOperations<?, ?> nextRequest = mock(ChannelOperations.class);
when(nextRequest.as(ChannelOperations.class)).thenReturn(nextRequest);
Connection.from(fixture.channel).rebind(nextRequest);

fixture.cleanup.accept(fixture.connection);

assertSame(fixture.context, fixture.channel.attr(CONTEXT_ATTRIBUTE_KEY).get());
assertSame(fixture.parent, fixture.channel.attr(CLIENT_PARENT_ATTRIBUTE_KEY).get());
}

private static class Fixture {
final AgentSpan parent = mock(AgentSpan.class, CALLS_REAL_METHODS);
final Context context = Context.root().with(parent);
final Channel channel = mock(Channel.class);
final HttpClientRequest request =
mock(HttpClientRequest.class, withSettings().extraInterfaces(Connection.class));
final Connection connection = (Connection) request;
final ClearRequestContext cleanup;

@SuppressWarnings("unchecked")
Fixture() {
DefaultAttributeMap attributes = new DefaultAttributeMap();
when(channel.hasAttr(any()))
.thenAnswer(invocation -> attributes.hasAttr(invocation.getArgument(0)));
when(channel.attr(any()))
.thenAnswer(invocation -> attributes.attr(invocation.getArgument(0)));
when(connection.channel()).thenReturn(channel);
channel.attr(CONTEXT_ATTRIBUTE_KEY).set(context);
channel.attr(CLIENT_PARENT_ATTRIBUTE_KEY).set(parent);
ContextStore<HttpClientRequest, Context> contexts = mock(ContextStore.class);
when(contexts.remove(request)).thenReturn(context);
cleanup = new ClearRequestContext(contexts);
}
}
}
Loading