Skip to content
Open
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 @@ -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
Expand All @@ -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
Expand All @@ -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()
Expand All @@ -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)
Expand All @@ -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 {
Expand All @@ -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) {
Comment thread
amarziali marked this conversation as resolved.
if (System.nanoTime() >= deadline) {
throw new IllegalStateException('Jetty did not schedule the continuation timeout')
}
LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(1))
}
}
}

@Override
Expand Down
Loading