From 9c27dfedc8e05711716450b5f312895e4fdfb675 Mon Sep 17 00:00:00 2001 From: Andrea Marziali Date: Fri, 11 Sep 2026 17:57:35 +0200 Subject: [PATCH 1/2] Prevent Pekko stream lifecycle context retention --- .../concurrent/AsyncPropagatingDisableInstrumentation.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/dd-java-agent/instrumentation/java/java-concurrent/java-concurrent-1.8/src/main/java/datadog/trace/instrumentation/java/concurrent/AsyncPropagatingDisableInstrumentation.java b/dd-java-agent/instrumentation/java/java-concurrent/java-concurrent-1.8/src/main/java/datadog/trace/instrumentation/java/concurrent/AsyncPropagatingDisableInstrumentation.java index 2eadfb0cfe9..0a27b4ce90d 100644 --- a/dd-java-agent/instrumentation/java/java-concurrent/java-concurrent-1.8/src/main/java/datadog/trace/instrumentation/java/concurrent/AsyncPropagatingDisableInstrumentation.java +++ b/dd-java-agent/instrumentation/java/java-concurrent/java-concurrent-1.8/src/main/java/datadog/trace/instrumentation/java/concurrent/AsyncPropagatingDisableInstrumentation.java @@ -64,6 +64,8 @@ public AsyncPropagatingDisableInstrumentation() { extendsClass(named("java.net.http.HttpClient")); private static final String LETTUCE_HANDSHAKE_HANDLER = "io.lettuce.core.protocol.RedisHandshakeHandler"; + private static final String PEKKO_HTTP_STREAM_STAGE = + "org.apache.pekko.http.impl.util.StreamUtils$$anon$4$$anon$5"; @Override public boolean onlyMatchKnownTypes() { @@ -101,6 +103,7 @@ public String[] knownMatchingTypes() { "io.reactivex.rxjava3.internal.schedulers.AbstractDirectTask", "jdk.internal.net.http.HttpClientImpl", LETTUCE_HANDSHAKE_HANDLER, + PEKKO_HTTP_STREAM_STAGE, "io.netty.util.concurrent.GlobalEventExecutor", "io.grpc.netty.shaded.io.netty.util.concurrent.GlobalEventExecutor", "com.linecorp.armeria.client.HttpClientFactory", @@ -213,6 +216,8 @@ public void methodAdvice(MethodTransformer transformer) { transformer.applyAdvice(namedOneOf("sendAsync").and(isDeclaredBy(JAVA_HTTP_CLIENT)), advice); transformer.applyAdvice( named("channelRegistered").and(isDeclaredBy(named(LETTUCE_HANDSHAKE_HANDLER))), advice); + transformer.applyAdvice( + named("preStart").and(isDeclaredBy(named(PEKKO_HTTP_STREAM_STAGE))), advice); // armeria runs its own codec/pipeline, so the active request span captured during connection // pool creation and channel connect will have no consumers. transformer.applyAdvice( From e82bb5b91d1e23192bc20a800e7ec76227642aa0 Mon Sep 17 00:00:00 2001 From: Andrea Marziali Date: Mon, 14 Sep 2026 11:12:13 +0200 Subject: [PATCH 2/2] change matcher --- .../AsyncPropagatingDisableInstrumentation.java | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/dd-java-agent/instrumentation/java/java-concurrent/java-concurrent-1.8/src/main/java/datadog/trace/instrumentation/java/concurrent/AsyncPropagatingDisableInstrumentation.java b/dd-java-agent/instrumentation/java/java-concurrent/java-concurrent-1.8/src/main/java/datadog/trace/instrumentation/java/concurrent/AsyncPropagatingDisableInstrumentation.java index 0a27b4ce90d..7e7d915e0dc 100644 --- a/dd-java-agent/instrumentation/java/java-concurrent/java-concurrent-1.8/src/main/java/datadog/trace/instrumentation/java/concurrent/AsyncPropagatingDisableInstrumentation.java +++ b/dd-java-agent/instrumentation/java/java-concurrent/java-concurrent-1.8/src/main/java/datadog/trace/instrumentation/java/concurrent/AsyncPropagatingDisableInstrumentation.java @@ -10,6 +10,7 @@ import static datadog.trace.instrumentation.java.concurrent.ConcurrentInstrumentationNames.EXECUTOR_INSTRUMENTATION_NAME; import static net.bytebuddy.matcher.ElementMatchers.isDeclaredBy; import static net.bytebuddy.matcher.ElementMatchers.isTypeInitializer; +import static net.bytebuddy.matcher.ElementMatchers.takesNoArguments; import com.google.auto.service.AutoService; import datadog.trace.agent.tooling.Instrumenter; @@ -64,8 +65,9 @@ public AsyncPropagatingDisableInstrumentation() { extendsClass(named("java.net.http.HttpClient")); private static final String LETTUCE_HANDSHAKE_HANDLER = "io.lettuce.core.protocol.RedisHandshakeHandler"; - private static final String PEKKO_HTTP_STREAM_STAGE = - "org.apache.pekko.http.impl.util.StreamUtils$$anon$4$$anon$5"; + private static final ElementMatcher PEKKO_HTTP_STREAM_STAGE = + nameStartsWith("org.apache.pekko.http.impl.util.StreamUtils$") + .and(extendsClass(named("org.apache.pekko.stream.stage.GraphStageLogic"))); @Override public boolean onlyMatchKnownTypes() { @@ -103,7 +105,6 @@ public String[] knownMatchingTypes() { "io.reactivex.rxjava3.internal.schedulers.AbstractDirectTask", "jdk.internal.net.http.HttpClientImpl", LETTUCE_HANDSHAKE_HANDLER, - PEKKO_HTTP_STREAM_STAGE, "io.netty.util.concurrent.GlobalEventExecutor", "io.grpc.netty.shaded.io.netty.util.concurrent.GlobalEventExecutor", "com.linecorp.armeria.client.HttpClientFactory", @@ -123,7 +124,8 @@ public ElementMatcher hierarchyMatcher() { .or(REACTOR_DISABLED_TYPE_INITIALIZERS) .or(RXJAVA2_DISABLED_TYPE_INITIALIZERS) .or(RXJAVA3_DISABLED_TYPE_INITIALIZERS) - .or(JAVA_HTTP_CLIENT); + .or(JAVA_HTTP_CLIENT) + .or(PEKKO_HTTP_STREAM_STAGE); } @Override @@ -217,7 +219,8 @@ public void methodAdvice(MethodTransformer transformer) { transformer.applyAdvice( named("channelRegistered").and(isDeclaredBy(named(LETTUCE_HANDSHAKE_HANDLER))), advice); transformer.applyAdvice( - named("preStart").and(isDeclaredBy(named(PEKKO_HTTP_STREAM_STAGE))), advice); + named("preStart").and(takesNoArguments()).and(isDeclaredBy(PEKKO_HTTP_STREAM_STAGE)), + advice); // armeria runs its own codec/pipeline, so the active request span captured during connection // pool creation and channel connect will have no consumers. transformer.applyAdvice(