diff --git a/dd-java-agent/instrumentation/jetty/jetty-server/jetty-server-9.3/src/test/groovy/datadog/trace/instrumentation/jetty9/JettyContinuationHandlerTest.groovy b/dd-java-agent/instrumentation/jetty/jetty-server/jetty-server-9.3/src/test/groovy/datadog/trace/instrumentation/jetty9/JettyContinuationHandlerTest.groovy index d32f386a387..3b782cfbe45 100644 --- a/dd-java-agent/instrumentation/jetty/jetty-server/jetty-server-9.3/src/test/groovy/datadog/trace/instrumentation/jetty9/JettyContinuationHandlerTest.groovy +++ b/dd-java-agent/instrumentation/jetty/jetty-server/jetty-server-9.3/src/test/groovy/datadog/trace/instrumentation/jetty9/JettyContinuationHandlerTest.groovy @@ -6,6 +6,8 @@ import datadog.trace.bootstrap.instrumentation.api.Tags import datadog.trace.instrumentation.servlet3.AsyncDispatcherDecorator import org.eclipse.jetty.continuation.Continuation import org.eclipse.jetty.continuation.ContinuationSupport +import org.eclipse.jetty.server.AsyncContextEvent +import org.eclipse.jetty.server.AsyncContextState import org.eclipse.jetty.server.Request import org.eclipse.jetty.server.handler.AbstractHandler import org.eclipse.jetty.server.session.SessionHandler @@ -14,8 +16,12 @@ import javax.servlet.MultipartConfigElement import javax.servlet.ServletException import javax.servlet.http.HttpServletRequest import javax.servlet.http.HttpServletResponse +import java.lang.reflect.Field +import java.lang.reflect.Modifier import java.util.concurrent.ExecutorService import java.util.concurrent.Executors +import java.util.concurrent.TimeUnit +import java.util.concurrent.locks.LockSupport import static test.TestHandler.handleRequest import static datadog.trace.agent.test.base.HttpServerTest.ServerEndpoint.SUCCESS @@ -25,6 +31,9 @@ import static datadog.trace.instrumentation.servlet3.TestServlet3.SERVLET_TIMEOU abstract class JettyContinuationHandlerTest extends Jetty9Test { + private static final long CONTINUATION_TIMEOUT = TimeUnit.SECONDS.toMillis(30) + private static final long TIMEOUT_TASK_WAIT = TimeUnit.SECONDS.toNanos(10) + @Override AbstractHandler handler() { def ret = new SessionHandler() @@ -33,10 +42,18 @@ abstract class JettyContinuationHandlerTest extends Jetty9Test { } static class ContinuationTestHandler extends AbstractHandler { + private static final Field TIMEOUT_TASK = AsyncContextEvent.getDeclaredField('_timeoutTask') private static final MultipartConfigElement MULTIPART_CONFIG_ELEMENT = new MultipartConfigElement(System.getProperty('java.io.tmpdir')) static final ContinuationTestHandler INSTANCE = new ContinuationTestHandler() final ExecutorService executorService = Executors.newSingleThreadExecutor() + static { + TIMEOUT_TASK.accessible = true + if (!Modifier.isVolatile(TIMEOUT_TASK.modifiers)) { + throw new IllegalStateException('Jetty timeout task must be volatile') + } + } + @Override void handle(String target, Request baseRequest, HttpServletRequest request, HttpServletResponse response) throws IOException, ServletException { request.setAttribute('org.eclipse.jetty.multipartConfig', MULTIPART_CONFIG_ELEMENT) @@ -48,8 +65,12 @@ abstract class JettyContinuationHandlerTest extends Jetty9Test { // this happens in the /exception endpoint if (!request.getAttribute('javax.servlet.error.status_code')) { if (continuation.initial) { + continuation.setTimeout(CONTINUATION_TIMEOUT) continuation.suspend() executorService.execute { + // Jetty schedules the timeout after this handler returns. Its volatile task field is + // the publication signal that makes it safe for resume to cancel the scheduled task. + awaitTimeoutTask(request) continuation.resume() } } else { @@ -58,6 +79,18 @@ abstract class JettyContinuationHandlerTest extends Jetty9Test { } baseRequest.handled = true } + + private static void awaitTimeoutTask(HttpServletRequest request) { + AsyncContextState asyncContext = request.asyncContext as AsyncContextState + AsyncContextEvent event = asyncContext.httpChannelState.asyncContextEvent + long deadline = System.nanoTime() + TIMEOUT_TASK_WAIT + while (TIMEOUT_TASK.get(event) == null) { + if (System.nanoTime() >= deadline) { + throw new IllegalStateException('Jetty did not schedule the continuation timeout') + } + LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(1)) + } + } } @Override