diff --git a/dd-java-agent/agent-tooling/src/main/resources/datadog/trace/agent/tooling/bytebuddy/matcher/ignored_class_name.trie b/dd-java-agent/agent-tooling/src/main/resources/datadog/trace/agent/tooling/bytebuddy/matcher/ignored_class_name.trie index e86702c934c..e719f4684f1 100644 --- a/dd-java-agent/agent-tooling/src/main/resources/datadog/trace/agent/tooling/bytebuddy/matcher/ignored_class_name.trie +++ b/dd-java-agent/agent-tooling/src/main/resources/datadog/trace/agent/tooling/bytebuddy/matcher/ignored_class_name.trie @@ -133,6 +133,7 @@ 0 akka.http.impl.engine.client.PoolMasterActor 0 akka.http.impl.engine.client.pool.NewHostConnectionPool$* 0 akka.http.impl.engine.http2.Http2Ext +0 akka.http.impl.engine.server.HttpServerBluePrint$TimeoutAccessImpl 0 akka.http.impl.engine.server.HttpServerBluePrint$TimeoutAccessImpl$* 0 akka.http.impl.util.StreamUtils$* # saves ~0.1s skipping ~233 classes diff --git a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/groovy/AkkaHttpTimeoutCleanupTest.groovy b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/groovy/AkkaHttpTimeoutCleanupTest.groovy new file mode 100644 index 00000000000..99c1ce703cc --- /dev/null +++ b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/groovy/AkkaHttpTimeoutCleanupTest.groovy @@ -0,0 +1,77 @@ +import akka.actor.ActorSystem +import akka.actor.Cancellable +import akka.http.javadsl.model.HttpRequest +import akka.stream.ActorMaterializer +import datadog.trace.agent.test.InstrumentationSpecification +import scala.concurrent.Await +import scala.concurrent.Promise$ +import scala.concurrent.duration.Duration +import scala.runtime.BoxedUnit +import spock.util.concurrent.PollingConditions + +import java.util.concurrent.TimeUnit + +import static datadog.trace.agent.test.utils.TraceUtils.runUnderTrace +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.isAsyncPropagationEnabled +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.setAsyncPropagationEnabled + +class AkkaHttpTimeoutCleanupTest extends InstrumentationSpecification { + + def 'timeout cleanup does not retain a request while its entity is unfinished: propagation=#enabled'() { + setup: + def system = ActorSystem.create('timeout-cleanup') + def materializer = ActorMaterializer.create(system) + def requestEnd = Promise$.MODULE$.apply() + def type = Class.forName('akka.http.impl.engine.server.HttpServerBluePrint$TimeoutAccessImpl') + def constructor = type.declaredConstructors[0] + constructor.accessible = true + def arguments = [ + HttpRequest.create(), + Duration.create(1, TimeUnit.DAYS), + requestEnd.future(), + null, + materializer + ] + if (constructor.parameterCount == 6) { + arguments.add(system.log()) + } + def access = constructor.newInstance(arguments as Object[]) + def clear = type.getDeclaredMethod('clear') + clear.accessible = true + + when: + runUnderTrace('request') { + setAsyncPropagationEnabled(enabled) + clear.invoke(access) + assert isAsyncPropagationEnabled() == enabled + } + + then: + // The request trace must finish without waiting for the entity or its timeout cleanup. + TEST_WRITER.waitForTraces(1) + TEST_WRITER.get(0).size() == 1 + !requestEnd.isCompleted() + + when: + requestEnd.success(BoxedUnit.UNIT) + def setup = Await.result(access.get(), Duration.create(5, TimeUnit.SECONDS)) + def scheduledTask = setup.class.getDeclaredMethod('scheduledTask') + scheduledTask.accessible = true + def task = scheduledTask.invoke(setup) as Cancellable + + then: + new PollingConditions(timeout: 5).eventually { + assert task.isCancelled() + } + + cleanup: + requestEnd?.trySuccess(BoxedUnit.UNIT) + materializer?.shutdown() + if (system != null) { + Await.result(system.terminate(), Duration.create(10, TimeUnit.SECONDS)) + } + + where: + enabled << [true, false] + } +} 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 82578773d74..89dd1c93e43 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 @@ -118,7 +118,8 @@ public String[] knownMatchingTypes() { "io.netty.handler.timeout.IdleStateHandler", "io.grpc.netty.shaded.io.netty.handler.timeout.IdleStateHandler", "com.linecorp.armeria.client.HttpClientFactory", - "com.linecorp.armeria.client.HttpChannelPool" + "com.linecorp.armeria.client.HttpChannelPool", + "akka.http.impl.engine.server.HttpServerBluePrint$TimeoutAccessImpl" }; } @@ -141,6 +142,14 @@ public ElementMatcher hierarchyMatcher() { @Override public void methodAdvice(MethodTransformer transformer) { String advice = getClass().getName() + "$DisableAsyncAdvice"; + // Timeout cancellation is housekeeping; its request-entity future may never complete. + transformer.applyAdvice( + named("clear") + .and(takesNoArguments()) + .and( + isDeclaredBy( + named("akka.http.impl.engine.server.HttpServerBluePrint$TimeoutAccessImpl"))), + advice); transformer.applyAdvice(named("schedulePeriodically").and(isDeclaredBy(RX_WORKERS)), advice); transformer.applyAdvice( named("call").and(isDeclaredBy(named("rx.internal.operators.OperatorTimeoutBase"))),