diff --git a/dd-java-agent/instrumentation/spray-1.3/build.gradle b/dd-java-agent/instrumentation/spray-1.3/build.gradle index 701afc6880a..53109fe28bf 100644 --- a/dd-java-agent/instrumentation/spray-1.3/build.gradle +++ b/dd-java-agent/instrumentation/spray-1.3/build.gradle @@ -23,4 +23,5 @@ dependencies { testImplementation group: 'io.spray', name: "spray-routing_$scalaVersion", version: '1.3.1' testImplementation project(':dd-java-agent:instrumentation-testing') + testImplementation libs.bundles.mockito } diff --git a/dd-java-agent/instrumentation/spray-1.3/src/main/scala/datadog/trace/instrumentation/spray/SprayHelper.scala b/dd-java-agent/instrumentation/spray-1.3/src/main/scala/datadog/trace/instrumentation/spray/SprayHelper.scala index 0c37f91fd66..f74035b4d1d 100644 --- a/dd-java-agent/instrumentation/spray-1.3/src/main/scala/datadog/trace/instrumentation/spray/SprayHelper.scala +++ b/dd-java-agent/instrumentation/spray-1.3/src/main/scala/datadog/trace/instrumentation/spray/SprayHelper.scala @@ -1,7 +1,6 @@ package datadog.trace.instrumentation.spray import datadog.context.Context; -import datadog.context.ContextScope; import datadog.trace.bootstrap.instrumentation.api.AgentSpan import datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan import datadog.trace.bootstrap.instrumentation.api.Java8BytecodeBridge.rootContext; @@ -9,14 +8,40 @@ import datadog.trace.instrumentation.spray.SprayHttpServerDecorator.DECORATE import spray.http.HttpResponse import spray.routing.{RequestContext, Route} +import java.util.concurrent.atomic.AtomicInteger import scala.util.control.NonFatal object SprayHelper { + private val ResponseComplete = 1 + private val ScopeClosed = 2 + private val Complete = ResponseComplete | ScopeClosed + + def responseComplete(span: AgentSpan, completion: AtomicInteger): Unit = + complete(span, completion, ResponseComplete) + + def scopeClosed(span: AgentSpan, completion: AtomicInteger): Unit = + complete(span, completion, ScopeClosed) + + private def complete(span: AgentSpan, completion: AtomicInteger, event: Int): Unit = { + var current = completion.get() + while ((current & event) == 0) { + if (completion.compareAndSet(current, current | event)) { + // The second event observes the first event's work and owns finishing the span. + if ((current | event) == Complete) { + span.finish() + } + return + } + current = completion.get() + } + } + def wrapRequestContext( ctx: RequestContext, span: AgentSpan, parentContext: Context, - scope: ContextScope + context: Context, + completion: AtomicInteger ): RequestContext = { ctx.withRouteResponseMapped(message => { DECORATE.onRequest(span, ctx, ctx.request, parentContext) @@ -25,9 +50,8 @@ object SprayHelper { case throwable: Throwable => DECORATE.onError(span, throwable) case x => } - DECORATE.beforeFinish(scope.context()) - scope.close() - span.finish() + DECORATE.beforeFinish(context) + responseComplete(span, completion) message }) } diff --git a/dd-java-agent/instrumentation/spray-1.3/src/main/scala/datadog/trace/instrumentation/spray/SprayHttpServerRunSealedRouteAdvice.java b/dd-java-agent/instrumentation/spray-1.3/src/main/scala/datadog/trace/instrumentation/spray/SprayHttpServerRunSealedRouteAdvice.java index f3ba5b60e1b..fde1b73608d 100644 --- a/dd-java-agent/instrumentation/spray-1.3/src/main/scala/datadog/trace/instrumentation/spray/SprayHttpServerRunSealedRouteAdvice.java +++ b/dd-java-agent/instrumentation/spray-1.3/src/main/scala/datadog/trace/instrumentation/spray/SprayHttpServerRunSealedRouteAdvice.java @@ -10,6 +10,7 @@ import datadog.context.Context; import datadog.context.ContextScope; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import java.util.concurrent.atomic.AtomicInteger; import net.bytebuddy.asm.Advice; import spray.http.HttpRequest; import spray.routing.RequestContext; @@ -17,7 +18,8 @@ public class SprayHttpServerRunSealedRouteAdvice { @Advice.OnMethodEnter(suppress = Throwable.class) public static ContextScope enter( - @Advice.Argument(value = 1, readOnly = false) RequestContext ctx) { + @Advice.Argument(value = 1, readOnly = false) RequestContext ctx, + @Advice.Local("completion") AtomicInteger completion) { final Context parentContext; final Context context; final AgentSpan span; @@ -37,16 +39,23 @@ public static ContextScope enter( ContextScope scope = context.attach(); DECORATE.afterStart(span); - ctx = SprayHelper.wrapRequestContext(ctx, span, parentContext, scope); + completion = new AtomicInteger(); + ctx = SprayHelper.wrapRequestContext(ctx, span, parentContext, context, completion); return scope; } @Advice.OnMethodExit(onThrowable = Throwable.class, suppress = Throwable.class) public static void exit( - @Advice.Enter final ContextScope scope, @Advice.Thrown final Throwable throwable) { - if (throwable != null) { - DECORATE.onError(scope, throwable); + @Advice.Enter final ContextScope scope, + @Advice.Thrown final Throwable throwable, + @Advice.Local("completion") AtomicInteger completion) { + try { + if (throwable != null) { + DECORATE.onError(scope, throwable); + } + } finally { + scope.close(); + SprayHelper.scopeClosed(spanFromContext(scope.context()), completion); } - scope.close(); } } diff --git a/dd-java-agent/instrumentation/spray-1.3/src/test/java/datadog/trace/instrumentation/spray/RequestCompletionTest.java b/dd-java-agent/instrumentation/spray-1.3/src/test/java/datadog/trace/instrumentation/spray/RequestCompletionTest.java new file mode 100644 index 00000000000..689de03aac5 --- /dev/null +++ b/dd-java-agent/instrumentation/spray-1.3/src/test/java/datadog/trace/instrumentation/spray/RequestCompletionTest.java @@ -0,0 +1,83 @@ +package datadog.trace.instrumentation.spray; + +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +class RequestCompletionTest { + @ParameterizedTest + @ValueSource(booleans = {true, false}) + void finishesOnceAfterBothEvents(boolean responseFirst) { + AgentSpan span = mock(AgentSpan.class); + AtomicInteger completion = new AtomicInteger(); + Runnable first = + responseFirst + ? () -> SprayHelper.responseComplete(span, completion) + : () -> SprayHelper.scopeClosed(span, completion); + Runnable second = + responseFirst + ? () -> SprayHelper.scopeClosed(span, completion) + : () -> SprayHelper.responseComplete(span, completion); + + first.run(); + first.run(); + verify(span, never()).finish(); + + second.run(); + first.run(); + second.run(); + verify(span).finish(); + } + + @Test + void finishesOnceWhenResponseAndScopeClosureRace() throws Exception { + ExecutorService executor = Executors.newFixedThreadPool(2); + try { + for (int i = 0; i < 100; i++) { + AgentSpan span = mock(AgentSpan.class); + AtomicInteger completion = new AtomicInteger(); + CountDownLatch ready = new CountDownLatch(2); + CountDownLatch start = new CountDownLatch(1); + Future response = + executor.submit( + () -> { + ready.countDown(); + assertTrue(start.await(5, TimeUnit.SECONDS)); + SprayHelper.responseComplete(span, completion); + return null; + }); + Future route = + executor.submit( + () -> { + ready.countDown(); + assertTrue(start.await(5, TimeUnit.SECONDS)); + SprayHelper.scopeClosed(span, completion); + return null; + }); + try { + assertTrue(ready.await(5, TimeUnit.SECONDS)); + } finally { + start.countDown(); + } + response.get(5, TimeUnit.SECONDS); + route.get(5, TimeUnit.SECONDS); + verify(span).finish(); + } + } finally { + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } +} diff --git a/dd-java-agent/instrumentation/spray-1.3/src/test/scala/SprayHttpTestWebServer.scala b/dd-java-agent/instrumentation/spray-1.3/src/test/scala/SprayHttpTestWebServer.scala index aaf13b1f25e..67e2c72a65f 100644 --- a/dd-java-agent/instrumentation/spray-1.3/src/test/scala/SprayHttpTestWebServer.scala +++ b/dd-java-agent/instrumentation/spray-1.3/src/test/scala/SprayHttpTestWebServer.scala @@ -6,6 +6,8 @@ import datadog.trace.agent.test.base.HttpServerTest.ServerEndpoint import datadog.trace.agent.test.base.HttpServerTest.ServerEndpoint._ import datadog.trace.agent.test.base.{HttpServer, HttpServerTest} import datadog.trace.agent.test.utils.PortUtils +import datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan +import datadog.trace.core.DDSpan import groovy.lang.Closure import spray.can.Http import spray.http.HttpHeaders.RawHeader @@ -50,10 +52,12 @@ class ServiceActor extends HttpServiceActor with ActorLogging { def receive = runRoute { path(SUCCESS.relativePath()) { get { ctx: RequestContext => + val requestSpan = activeSpan().asInstanceOf[DDSpan] HttpServerTest.controller( SUCCESS, new ControllerHttpResponseToClosureAdapter(ctx, SUCCESS) ) + assert(!requestSpan.isFinished, "request finished before the route scope closed") } } ~ path(FORWARDED.relativePath()) {