diff --git a/dd-java-agent/agent-bootstrap/src/main/java/datadog/trace/bootstrap/ResourceMethodSpanTracker.java b/dd-java-agent/agent-bootstrap/src/main/java/datadog/trace/bootstrap/ResourceMethodSpanTracker.java new file mode 100644 index 00000000000..cbf94267328 --- /dev/null +++ b/dd-java-agent/agent-bootstrap/src/main/java/datadog/trace/bootstrap/ResourceMethodSpanTracker.java @@ -0,0 +1,53 @@ +package datadog.trace.bootstrap; + +import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import java.util.ArrayDeque; + +/** + * Tracks, per thread, which annotated resource-method invocation (identified by its span) is + * currently the innermost one still on the call stack -- i.e. whose exit advice has not yet run. + * + *

Used by the JAX-RS/Jakarta-RS instrumentation to tell a synchronous {@code + * AsyncResponse#resume()}/{@code cancel()} call -- one nested inside the still-running resource + * method that owns the response -- apart from a genuinely asynchronous one called later, from a + * different thread or after an intervening unrelated instrumented call. A generic "is any resource + * method open on this thread" counter is not enough for that: it can't tell one resource method's + * invocation apart from another's on a shared thread pool, and it can't see past a nested + * instrumented call (e.g. a {@code @Trace}-annotated helper) that becomes the current active span + * without popping this stack. Comparing the specific span object against the top of this stack + * answers the exact question that matters, without either failure mode. + * + *

Deliberately bootstrap-loaded (like {@link CallDepthThreadLocalMap}) rather than living on a + * per-instrumentation helper class: helper classes are injected once per target classloader, so a + * container that loads the resource-method advice and the AsyncResponse advice into different + * classloaders (e.g. a modular server where the JAX-RS runtime and the deployed application are in + * separate classloaders) would otherwise give each advice its own, disconnected copy of this state. + */ +public final class ResourceMethodSpanTracker { + + private static final ThreadLocal> STACK = new ThreadLocal<>(); + + private ResourceMethodSpanTracker() {} + + public static void enter(final AgentSpan span) { + ArrayDeque stack = STACK.get(); + if (stack == null) { + stack = new ArrayDeque<>(4); + STACK.set(stack); + } + stack.push(span); + } + + public static void exit() { + final ArrayDeque stack = STACK.get(); + if (stack != null) { + stack.pop(); + } + } + + /** True if {@code span} is the innermost still-open resource-method invocation on this thread. */ + public static boolean isInnermost(final AgentSpan span) { + final ArrayDeque stack = STACK.get(); + return stack != null && stack.peek() == span; + } +} diff --git a/dd-java-agent/instrumentation/cxf-2.1/src/test/groovy/CxfContextPropagationTest.groovy b/dd-java-agent/instrumentation/cxf-2.1/src/test/groovy/CxfContextPropagationTest.groovy index 9e8e660b781..c17215be81d 100644 --- a/dd-java-agent/instrumentation/cxf-2.1/src/test/groovy/CxfContextPropagationTest.groovy +++ b/dd-java-agent/instrumentation/cxf-2.1/src/test/groovy/CxfContextPropagationTest.groovy @@ -22,12 +22,20 @@ class CxfContextPropagationTest extends InstrumentationSpecification { @Override void setupSpec() { JAXRSServerFactoryBean sf = new JAXRSServerFactoryBean() - sf.setResourceClasses(TestResource) + sf.setResourceClasses(TestResource, AsyncResumeResource, TrueAsyncResumeResource, AsyncCancelResource, NestedResumeResource) List providers = [new TestExceptionMapper()] sf.setProviders(providers) sf.setResourceProvider(TestResource, new SingletonResourceProvider(new TestResource(), true)) + sf.setResourceProvider(AsyncResumeResource, + new SingletonResourceProvider(new AsyncResumeResource(), true)) + sf.setResourceProvider(TrueAsyncResumeResource, + new SingletonResourceProvider(new TrueAsyncResumeResource(), true)) + sf.setResourceProvider(AsyncCancelResource, + new SingletonResourceProvider(new AsyncCancelResource(), true)) + sf.setResourceProvider(NestedResumeResource, + new SingletonResourceProvider(new NestedResumeResource(), true)) sf.setAddress("http://localhost:0") server = sf.create() @@ -40,6 +48,15 @@ class CxfContextPropagationTest extends InstrumentationSpecification { server?.stop() } + @Override + protected boolean enabledFinishTimingChecks() { + // Regression guard for https://github.com/DataDog/dd-trace-java/issues/12597: + // fails the test with the exact "finished more than once" stack traces if the + // jax-rs.request span is ever finished twice (e.g. once from AsyncResponse#resume() + // and again from the resource method's own exit advice). + return true + } + def "should propagate context on async request resume"() { setup: def client = OkHttpUtils.client() @@ -97,4 +114,251 @@ class CxfContextPropagationTest extends InstrumentationSpecification { } } } + + def "resume() called synchronously from within the resource method finishes the span only once"() { + // Regression test for https://github.com/DataDog/dd-trace-java/issues/12597: when + // AsyncResponse#resume() is called synchronously (not truly suspended, the resource + // method keeps running), the resource method's own jax-rs.request scope is still on + // top of the scope stack. Before the fix, JakartaRsAsyncResponseInstrumentation / + // JaxRsAsyncResponseInstrumentation would eagerly finish the span there, so any work + // done afterwards (here: doWorkAfterResume()) would be attributed as a child of an + // already-finished span, and the resource method's own exit advice would finish the + // same span a second time (caught by enabledFinishTimingChecks()). + setup: + def client = OkHttpUtils.client() + when: + def response = client.newCall(new Request.Builder() + .url("http://localhost:$port/asyncresume") + .get().build()).execute() + then: + assert response.code() == 200 + assert response.body().string() == "OK" + + assertTraces(1) { + trace(3) { + sortSpansByStart() + span { + operationName "servlet.request" + resourceName "GET /asyncresume" + spanType DDSpanTypes.HTTP_SERVER + errored false + parent() + tags { + "$Tags.COMPONENT" "jax-rs" + "$Tags.SPAN_KIND" Tags.SPAN_KIND_SERVER + "$Tags.PEER_HOST_IPV4" "127.0.0.1" + "$Tags.PEER_PORT" Integer + "$Tags.HTTP_URL" "http://localhost:$port/asyncresume" + "$Tags.HTTP_HOSTNAME" "localhost" + "$Tags.HTTP_METHOD" "GET" + "$Tags.HTTP_STATUS" 200 + "$Tags.HTTP_ROUTE" String + "servlet.path" { it == null || it == "/asyncresume" } + "$Tags.HTTP_USER_AGENT" String + "$Tags.HTTP_CLIENT_IP" "127.0.0.1" + "$Tags.NETWORK_CLIENT_IP" "127.0.0.1" + withCustomIntegrationName("jetty-server") + defaultTags() + } + } + span { + operationName "jax-rs.request" + resourceName "AsyncResumeResource.resumeThenWork" + spanType DDSpanTypes.HTTP_SERVER + errored false + childOfPrevious() + tags { + "$Tags.COMPONENT" "jax-rs-controller" + defaultTags() + } + } + // Still parented under jax-rs.request: proves that span wasn't finished (and its + // scope wasn't popped) by resume() itself, before the resource method returned. + TraceUtils.basicSpan(it, "trace.annotation", "AsyncResumeResource.doWorkAfterResume", span(1), null, ["component": "trace"]) + } + } + } + + def "cancel() called synchronously from within the resource method finishes the span only once"() { + // Same regression as above, but for AsyncResponseCancelAdvice: cancel() is called + // synchronously and the resource method keeps running afterwards. + setup: + def client = OkHttpUtils.client() + when: + def response = client.newCall(new Request.Builder() + .url("http://localhost:$port/asynccancel") + .get().build()).execute() + then: + assert response.code() == 503 + + assertTraces(1) { + trace(3) { + sortSpansByStart() + span { + operationName "servlet.request" + resourceName "GET /asynccancel" + spanType DDSpanTypes.HTTP_SERVER + errored true + parent() + tags { + "$Tags.COMPONENT" "jax-rs" + "$Tags.SPAN_KIND" Tags.SPAN_KIND_SERVER + "$Tags.PEER_HOST_IPV4" "127.0.0.1" + "$Tags.PEER_PORT" Integer + "$Tags.HTTP_URL" "http://localhost:$port/asynccancel" + "$Tags.HTTP_HOSTNAME" "localhost" + "$Tags.HTTP_METHOD" "GET" + "$Tags.HTTP_STATUS" 503 + "$Tags.HTTP_ROUTE" String + "servlet.path" { it == null || it == "/asynccancel" } + "$Tags.HTTP_USER_AGENT" String + "$Tags.HTTP_CLIENT_IP" "127.0.0.1" + "$Tags.NETWORK_CLIENT_IP" "127.0.0.1" + withCustomIntegrationName("jetty-server") + defaultTags() + } + } + span { + operationName "jax-rs.request" + resourceName "AsyncCancelResource.cancelThenWork" + spanType DDSpanTypes.HTTP_SERVER + errored false + childOfPrevious() + tags { + "$Tags.COMPONENT" "jax-rs-controller" + "canceled" true + defaultTags() + } + } + // Still parented under jax-rs.request: proves the span wasn't finished (and its + // scope wasn't popped) by cancel() itself, before the resource method returned. + TraceUtils.basicSpan(it, "trace.annotation", "AsyncCancelResource.doWorkAfterCancel", span(1), null, ["component": "trace"]) + } + } + } + + def "resume() called from a genuinely different thread (textbook async pattern) is unaffected"() { + // Regression guard the other way: the fix must not change the standard cross-thread + // async pattern, where the resource method returns without resolving anything and a + // completely different thread calls resume() later. Here the span IS finished by the + // resume() advice (activeSpan() on that other thread is not this span), exactly as + // before the fix. + setup: + def client = OkHttpUtils.client() + when: + def response = client.newCall(new Request.Builder() + .url("http://localhost:$port/trueasyncresume") + .get().build()).execute() + then: + assert response.code() == 200 + assert response.body().string() == "OK" + + assertTraces(1) { + trace(3) { + sortSpansByStart() + span { + operationName "servlet.request" + resourceName "GET /trueasyncresume" + spanType DDSpanTypes.HTTP_SERVER + errored false + parent() + tags { + "$Tags.COMPONENT" "jax-rs" + "$Tags.SPAN_KIND" Tags.SPAN_KIND_SERVER + "$Tags.PEER_HOST_IPV4" "127.0.0.1" + "$Tags.PEER_PORT" Integer + "$Tags.HTTP_URL" "http://localhost:$port/trueasyncresume" + "$Tags.HTTP_HOSTNAME" "localhost" + "$Tags.HTTP_METHOD" "GET" + "$Tags.HTTP_STATUS" 200 + "$Tags.HTTP_ROUTE" String + "servlet.path" { it == null || it == "/trueasyncresume" } + "$Tags.HTTP_USER_AGENT" String + "$Tags.HTTP_CLIENT_IP" "127.0.0.1" + "$Tags.NETWORK_CLIENT_IP" "127.0.0.1" + withCustomIntegrationName("jetty-server") + defaultTags() + } + } + span { + operationName "jax-rs.request" + resourceName "TrueAsyncResumeResource.suspendThenResumeFromAnotherThread" + spanType DDSpanTypes.HTTP_SERVER + errored false + childOfPrevious() + tags { + "$Tags.COMPONENT" "jax-rs-controller" + defaultTags() + } + } + // Runs on the background thread, before resume() -- still correctly parented + // under the (still-open, cross-thread-propagated) jax-rs.request span. + TraceUtils.basicSpan(it, "trace.annotation", "TrueAsyncResumeResource.doWorkOnBackgroundThread", span(1), null, ["component": "trace"]) + } + } + } + + def "resume() called synchronously from a nested @Trace helper finishes the span only once"() { + // Regression test for a gap found reviewing the fix for GH-12597: resume() is called + // synchronously, but from a @Trace-annotated helper method rather than directly from + // the resource method's own body. At that moment, the *helper's* span is the active + // one on this thread, not the resource method's -- checking activeSpan() against the + // resource method's span directly (an earlier version of this fix) would wrongly treat + // this as a genuinely-async resume and finish the span right there, then finish it + // again when the resource method itself returns (caught by enabledFinishTimingChecks()). + setup: + def client = OkHttpUtils.client() + when: + def response = client.newCall(new Request.Builder() + .url("http://localhost:$port/nestedresume") + .get().build()).execute() + then: + assert response.code() == 200 + assert response.body().string() == "OK" + + assertTraces(1) { + trace(3) { + sortSpansByStart() + span { + operationName "servlet.request" + resourceName "GET /nestedresume" + spanType DDSpanTypes.HTTP_SERVER + errored false + parent() + tags { + "$Tags.COMPONENT" "jax-rs" + "$Tags.SPAN_KIND" Tags.SPAN_KIND_SERVER + "$Tags.PEER_HOST_IPV4" "127.0.0.1" + "$Tags.PEER_PORT" Integer + "$Tags.HTTP_URL" "http://localhost:$port/nestedresume" + "$Tags.HTTP_HOSTNAME" "localhost" + "$Tags.HTTP_METHOD" "GET" + "$Tags.HTTP_STATUS" 200 + "$Tags.HTTP_ROUTE" String + "servlet.path" { it == null || it == "/nestedresume" } + "$Tags.HTTP_USER_AGENT" String + "$Tags.HTTP_CLIENT_IP" "127.0.0.1" + "$Tags.NETWORK_CLIENT_IP" "127.0.0.1" + withCustomIntegrationName("jetty-server") + defaultTags() + } + } + span { + operationName "jax-rs.request" + resourceName "NestedResumeResource.resumeViaHelper" + spanType DDSpanTypes.HTTP_SERVER + errored false + childOfPrevious() + tags { + "$Tags.COMPONENT" "jax-rs-controller" + defaultTags() + } + } + // The helper that actually calls resume() -- still parented under jax-rs.request, + // proving the resource-method span wasn't finished/popped while the helper (and + // resume() inside it) was still running. + TraceUtils.basicSpan(it, "trace.annotation", "NestedResumeResource.resumeFromHelper", span(1), null, ["component": "trace"]) + } + } + } } diff --git a/dd-java-agent/instrumentation/cxf-2.1/src/test/java/AsyncCancelResource.java b/dd-java-agent/instrumentation/cxf-2.1/src/test/java/AsyncCancelResource.java new file mode 100644 index 00000000000..742b105d16e --- /dev/null +++ b/dd-java-agent/instrumentation/cxf-2.1/src/test/java/AsyncCancelResource.java @@ -0,0 +1,22 @@ +import datadog.trace.api.Trace; +import javax.ws.rs.GET; +import javax.ws.rs.Path; +import javax.ws.rs.container.AsyncResponse; +import javax.ws.rs.container.Suspended; + +/** + * Same GH-12597 pattern as {@link AsyncResumeResource}, but exercising {@code + * AsyncResponse#cancel()} instead of {@code resume()} -- the third advice touched by the fix + * ({@code AsyncResponseCancelAdvice}). + */ +@Path("/asynccancel") +public class AsyncCancelResource { + @GET + public void cancelThenWork(@Suspended final AsyncResponse response) { + response.cancel(); + doWorkAfterCancel(); + } + + @Trace + private void doWorkAfterCancel() {} +} diff --git a/dd-java-agent/instrumentation/cxf-2.1/src/test/java/AsyncResumeResource.java b/dd-java-agent/instrumentation/cxf-2.1/src/test/java/AsyncResumeResource.java new file mode 100644 index 00000000000..705a67ae38a --- /dev/null +++ b/dd-java-agent/instrumentation/cxf-2.1/src/test/java/AsyncResumeResource.java @@ -0,0 +1,17 @@ +import datadog.trace.api.Trace; +import javax.ws.rs.GET; +import javax.ws.rs.Path; +import javax.ws.rs.container.AsyncResponse; +import javax.ws.rs.container.Suspended; + +@Path("/asyncresume") +public class AsyncResumeResource { + @GET + public void resumeThenWork(@Suspended final AsyncResponse response) { + response.resume("OK"); + doWorkAfterResume(); + } + + @Trace + private void doWorkAfterResume() {} +} diff --git a/dd-java-agent/instrumentation/cxf-2.1/src/test/java/NestedResumeResource.java b/dd-java-agent/instrumentation/cxf-2.1/src/test/java/NestedResumeResource.java new file mode 100644 index 00000000000..e3cebedd916 --- /dev/null +++ b/dd-java-agent/instrumentation/cxf-2.1/src/test/java/NestedResumeResource.java @@ -0,0 +1,25 @@ +import datadog.trace.api.Trace; +import javax.ws.rs.GET; +import javax.ws.rs.Path; +import javax.ws.rs.container.AsyncResponse; +import javax.ws.rs.container.Suspended; + +/** + * Regression coverage for a gap found in code review of GH-12597's fix: {@code resume()} called + * from a {@code @Trace}-annotated helper, not directly from the resource method's own body. At the + * moment {@code resume()} runs, the helper's own span is the active one, not the resource method's + * -- an early version of the fix compared {@code activeSpan()} directly against the resource + * method's span and wrongly treated this as a genuinely-async call. + */ +@Path("/nestedresume") +public class NestedResumeResource { + @GET + public void resumeViaHelper(@Suspended final AsyncResponse response) { + resumeFromHelper(response); + } + + @Trace + private void resumeFromHelper(final AsyncResponse response) { + response.resume("OK"); + } +} diff --git a/dd-java-agent/instrumentation/cxf-2.1/src/test/java/TrueAsyncResumeResource.java b/dd-java-agent/instrumentation/cxf-2.1/src/test/java/TrueAsyncResumeResource.java new file mode 100644 index 00000000000..d26501de069 --- /dev/null +++ b/dd-java-agent/instrumentation/cxf-2.1/src/test/java/TrueAsyncResumeResource.java @@ -0,0 +1,43 @@ +import datadog.trace.api.Trace; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import javax.ws.rs.GET; +import javax.ws.rs.Path; +import javax.ws.rs.container.AsyncResponse; +import javax.ws.rs.container.Suspended; + +/** + * The "textbook" async pattern, for regression coverage alongside {@link AsyncResumeResource}: the + * resource method returns without resuming, and a different thread resumes it later. The GH-12597 + * fix must not change behavior on this path. + */ +@Path("/trueasyncresume") +public class TrueAsyncResumeResource { + + private static final ExecutorService EXECUTOR = Executors.newFixedThreadPool(2); + + @GET + public void suspendThenResumeFromAnotherThread(@Suspended final AsyncResponse response) { + EXECUTOR.submit( + () -> { + try { + // Wait for the actual condition (the container has genuinely suspended the + // response) instead of guessing a fixed delay, so this can't race under a slow + // or overloaded test runner. + final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); + while (!response.isSuspended() && System.nanoTime() < deadline) { + Thread.sleep(1); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + doWorkOnBackgroundThread(); + response.resume("OK"); + }); + } + + @Trace + private void doWorkOnBackgroundThread() {} +} diff --git a/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/build.gradle b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/build.gradle index 516456dd2da..ef80cd6480e 100644 --- a/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/build.gradle +++ b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/build.gradle @@ -31,6 +31,7 @@ dependencies { testImplementation project(':dd-java-agent:instrumentation:servlet:javax-servlet:javax-servlet-3.0') testImplementation group: 'jakarta.ws.rs', name: 'jakarta.ws.rs-api', version: '3.0.0' testImplementation group: 'jakarta.xml.bind', name: 'jakarta.xml.bind-api', version: '3.0.0' + testRuntimeOnly project(':dd-java-agent:instrumentation:datadog:tracing:trace-annotation') latestDepTestImplementation group: 'jakarta.ws.rs', name: 'jakarta.ws.rs-api', version: '3.0.+' latestDepTestImplementation group: 'jakarta.xml.bind', name: 'jakarta.xml.bind-api', version: '3.0.+' diff --git a/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/main/java/datadog/trace/instrumentation/jakarta3/JakartaRsAnnotationsInstrumentation.java b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/main/java/datadog/trace/instrumentation/jakarta3/JakartaRsAnnotationsInstrumentation.java index c0c2ec9031c..f01f56ded1f 100644 --- a/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/main/java/datadog/trace/instrumentation/jakarta3/JakartaRsAnnotationsInstrumentation.java +++ b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/main/java/datadog/trace/instrumentation/jakarta3/JakartaRsAnnotationsInstrumentation.java @@ -7,6 +7,8 @@ import static datadog.trace.agent.tooling.bytebuddy.matcher.HierarchyMatchers.isAnnotatedWith; import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.namedOneOf; +import static datadog.trace.bootstrap.ResourceMethodSpanTracker.enter; +import static datadog.trace.bootstrap.ResourceMethodSpanTracker.exit; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; @@ -128,6 +130,9 @@ public static ContextScope nameSpan( if (contextStore != null && asyncResponse != null) { contextStore.put(asyncResponse, span); + // Only tracked for methods that can hand off to AsyncResponse#resume()/cancel(); + // see ResourceMethodSpanTracker for why this can't just be a bare counter. + enter(span); } return scope; @@ -141,8 +146,16 @@ public static void stopSpan( if (scope == null) { return; } + if (asyncResponse != null) { + exit(); + } final AgentSpan span = spanFromScope(scope); if (throwable != null) { + if (asyncResponse != null) { + // Clear span from the asyncResponse so a later resume()/cancel() call (e.g. from + // container exception-mapping) doesn't find a stale mapping and double-finish it. + InstrumentationContext.get(AsyncResponse.class, AgentSpan.class).put(asyncResponse, null); + } DECORATE.onError(span, throwable); DECORATE.beforeFinish(span); scope.close(); diff --git a/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/main/java/datadog/trace/instrumentation/jakarta3/JakartaRsAsyncResponseInstrumentation.java b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/main/java/datadog/trace/instrumentation/jakarta3/JakartaRsAsyncResponseInstrumentation.java index 1ae11e76090..9f58e4e29c7 100644 --- a/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/main/java/datadog/trace/instrumentation/jakarta3/JakartaRsAsyncResponseInstrumentation.java +++ b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/main/java/datadog/trace/instrumentation/jakarta3/JakartaRsAsyncResponseInstrumentation.java @@ -2,6 +2,7 @@ import static datadog.trace.agent.tooling.bytebuddy.matcher.HierarchyMatchers.implementsInterface; import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; +import static datadog.trace.bootstrap.ResourceMethodSpanTracker.isInnermost; import static datadog.trace.instrumentation.jakarta3.JakartaRsAnnotationsDecorator.DECORATE; import static net.bytebuddy.matcher.ElementMatchers.isPublic; import static net.bytebuddy.matcher.ElementMatchers.takesArgument; @@ -74,8 +75,18 @@ public static void stopSpan( final AgentSpan span = contextStore.get(asyncResponse); if (span != null) { - contextStore.put(asyncResponse, null); DECORATE.onError(span, throwable); + if (isInnermost(span)) { + // resume()/cancel() was called synchronously, nested inside the still-running + // resource method that owns this span (ResourceMethodSpanTracker.isInnermost(span) + // says its invocation is still the innermost open one on this thread). Let that + // method's own exit advice close the scope and finish the span (it will see + // asyncResponse.isSuspended() == false) instead of finishing it here, which would + // both double-finish the span and finish it prematurely while the resource method + // may still be doing work under it. + return; + } + contextStore.put(asyncResponse, null); DECORATE.beforeFinish(span); span.finish(); } @@ -94,8 +105,12 @@ public static void stopSpan( final AgentSpan span = contextStore.get(asyncResponse); if (span != null) { - contextStore.put(asyncResponse, null); DECORATE.onError(span, throwable); + if (isInnermost(span)) { + // see comment in AsyncResponseAdvice#stopSpan + return; + } + contextStore.put(asyncResponse, null); DECORATE.beforeFinish(span); span.finish(); } @@ -113,12 +128,16 @@ public static void stopSpan( final AgentSpan span = contextStore.get(asyncResponse); if (span != null) { - contextStore.put(asyncResponse, null); if (throwable != null) { DECORATE.onError(span, throwable); } else { span.setTag("canceled", true); } + if (isInnermost(span)) { + // see comment in AsyncResponseAdvice#stopSpan + return; + } + contextStore.put(asyncResponse, null); DECORATE.beforeFinish(span); span.finish(); } diff --git a/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/test/groovy/JakartaRsAsyncResponseInstrumentationTest.groovy b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/test/groovy/JakartaRsAsyncResponseInstrumentationTest.groovy new file mode 100644 index 00000000000..b4fa8e036e4 --- /dev/null +++ b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/test/groovy/JakartaRsAsyncResponseInstrumentationTest.groovy @@ -0,0 +1,204 @@ +import datadog.trace.agent.test.InstrumentationSpecification +import datadog.trace.agent.test.utils.TraceUtils +import datadog.trace.api.Trace +import datadog.trace.bootstrap.instrumentation.api.Tags +import jakarta.ws.rs.GET +import jakarta.ws.rs.Path +import jakarta.ws.rs.container.AsyncResponse +import jakarta.ws.rs.container.Suspended + +import java.util.concurrent.ExecutorService +import java.util.concurrent.Executors +import java.util.concurrent.TimeUnit + +/** + * Regression coverage for GH-12597 (jax-rs.request span double-finish / premature finish on + * synchronous AsyncResponse#resume()/cancel()) directly against the jakarta.ws.rs advice, with no + * real JAX-RS server involved -- resource methods are invoked directly, exactly like the existing + * JakartaRsAnnotations3InstrumentationTest does. This exercises the same advice code as the + * cxf-2.1 module's CxfContextPropagationTest (which only covers the javax.ws.rs path). + */ +class JakartaRsAsyncResponseInstrumentationTest extends InstrumentationSpecification { + + @Override + protected boolean enabledFinishTimingChecks() { + // Fails the test with the exact "finished more than once" stack traces if a jax-rs.request + // span is ever finished twice -- see CxfContextPropagationTest for the same guard. + return true + } + + def "resume() called synchronously from within the resource method finishes the span only once"() { + setup: + def response = new FakeAsyncResponse() + + when: + new AsyncResumeResource().resumeThenWork(response) + + then: + assertTraces(1) { + trace(2) { + sortSpansByStart() + span { + operationName "jakarta-rs.request" + resourceName "GET /asyncresume" + spanType "web" + errored false + parent() + tags { + "$Tags.COMPONENT" "jakarta-rs-controller" + "$Tags.HTTP_ROUTE" "/asyncresume" + defaultTags() + } + } + // Still parented under jakarta-rs.request: proves the span wasn't finished (and its + // scope wasn't popped) by resume() itself, before the resource method returned. + TraceUtils.basicSpan(it, "trace.annotation", "AsyncResumeResource.doWorkAfterResume", span(0), null, ["component": "trace"]) + } + } + } + + def "cancel() called synchronously from within the resource method finishes the span only once"() { + setup: + def response = new FakeAsyncResponse() + + when: + new AsyncCancelResource().cancelThenWork(response) + + then: + assertTraces(1) { + trace(2) { + sortSpansByStart() + span { + operationName "jakarta-rs.request" + resourceName "GET /asynccancel" + spanType "web" + errored false + parent() + tags { + "$Tags.COMPONENT" "jakarta-rs-controller" + "$Tags.HTTP_ROUTE" "/asynccancel" + "canceled" true + defaultTags() + } + } + TraceUtils.basicSpan(it, "trace.annotation", "AsyncCancelResource.doWorkAfterCancel", span(0), null, ["component": "trace"]) + } + } + } + + def "resume() called synchronously from a nested @Trace helper finishes the span only once"() { + // The gap found in code review: resume() called from a @Trace-annotated helper, not + // directly from the resource method's own body. activeSpan() at that moment is the + // helper's span, not the resource method's. + setup: + def response = new FakeAsyncResponse() + + when: + new NestedResumeResource().resumeViaHelper(response) + + then: + assertTraces(1) { + trace(2) { + sortSpansByStart() + span { + operationName "jakarta-rs.request" + resourceName "GET /nestedresume" + spanType "web" + errored false + parent() + tags { + "$Tags.COMPONENT" "jakarta-rs-controller" + "$Tags.HTTP_ROUTE" "/nestedresume" + defaultTags() + } + } + TraceUtils.basicSpan(it, "trace.annotation", "NestedResumeResource.resumeFromHelper", span(0), null, ["component": "trace"]) + } + } + } + + def "resume() called from a genuinely different thread is unaffected"() { + setup: + def response = new FakeAsyncResponse() + + when: + new TrueAsyncResumeResource().suspendThenResumeFromAnotherThread(response) + + then: + assertTraces(1) { + trace(2) { + sortSpansByStart() + span { + operationName "jakarta-rs.request" + resourceName "GET /trueasyncresume" + spanType "web" + errored false + parent() + tags { + "$Tags.COMPONENT" "jakarta-rs-controller" + "$Tags.HTTP_ROUTE" "/trueasyncresume" + defaultTags() + } + } + TraceUtils.basicSpan(it, "trace.annotation", "TrueAsyncResumeResource.doWorkOnBackgroundThread", span(0), null, ["component": "trace"]) + } + } + } + + @Path("/asyncresume") + static class AsyncResumeResource { + @GET + void resumeThenWork(@Suspended final AsyncResponse response) { + response.resume("OK") + doWorkAfterResume() + } + + @Trace + private void doWorkAfterResume() {} + } + + @Path("/asynccancel") + static class AsyncCancelResource { + @GET + void cancelThenWork(@Suspended final AsyncResponse response) { + response.cancel() + doWorkAfterCancel() + } + + @Trace + private void doWorkAfterCancel() {} + } + + @Path("/nestedresume") + static class NestedResumeResource { + @GET + void resumeViaHelper(@Suspended final AsyncResponse response) { + resumeFromHelper(response) + } + + @Trace + private void resumeFromHelper(final AsyncResponse response) { + response.resume("OK") + } + } + + @Path("/trueasyncresume") + static class TrueAsyncResumeResource { + private static final ExecutorService EXECUTOR = Executors.newFixedThreadPool(2) + + @GET + void suspendThenResumeFromAnotherThread(@Suspended final AsyncResponse response) { + EXECUTOR.submit({ + def deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5) + while (!response.isSuspended() && System.nanoTime() < deadline) { + Thread.sleep(1) + } + doWorkOnBackgroundThread() + response.resume("OK") + }) + } + + @Trace + private void doWorkOnBackgroundThread() {} + } +} diff --git a/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/test/java/FakeAsyncResponse.java b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/test/java/FakeAsyncResponse.java new file mode 100644 index 00000000000..c9770551c50 --- /dev/null +++ b/dd-java-agent/instrumentation/rs/jakarta-rs-annotations-3.0/src/test/java/FakeAsyncResponse.java @@ -0,0 +1,102 @@ +import jakarta.ws.rs.container.AsyncResponse; +import jakarta.ws.rs.container.TimeoutHandler; +import java.util.Collection; +import java.util.Collections; +import java.util.Date; +import java.util.Map; +import java.util.concurrent.TimeUnit; + +/** + * Minimal {@link AsyncResponse} implementation for testing the + * JakartaRsAsyncResponseInstrumentation advice directly (no real JAX-RS container/server involved) + * -- only {@code resume}/{@code cancel}/{@code isSuspended} carry real semantics; everything else + * is a no-op. + */ +public class FakeAsyncResponse implements AsyncResponse { + + private boolean suspended = true; + private boolean cancelled = false; + + @Override + public boolean resume(final Object response) { + if (!suspended) { + return false; + } + suspended = false; + return true; + } + + @Override + public boolean resume(final Throwable response) { + if (!suspended) { + return false; + } + suspended = false; + return true; + } + + @Override + public boolean cancel() { + if (!suspended) { + return false; + } + suspended = false; + cancelled = true; + return true; + } + + @Override + public boolean cancel(final int retryAfter) { + return cancel(); + } + + @Override + public boolean cancel(final Date retryAfter) { + return cancel(); + } + + @Override + public boolean isSuspended() { + return suspended; + } + + @Override + public boolean isCancelled() { + return cancelled; + } + + @Override + public boolean isDone() { + return !suspended; + } + + @Override + public boolean setTimeout(final long time, final TimeUnit unit) { + return true; + } + + @Override + public void setTimeoutHandler(final TimeoutHandler handler) {} + + @Override + public Collection> register(final Class callback) { + return Collections.emptyList(); + } + + @Override + public Map, Collection>> register( + final Class callback, final Class... callbacks) { + return Collections.emptyMap(); + } + + @Override + public Collection> register(final Object callback) { + return Collections.emptyList(); + } + + @Override + public Map, Collection>> register( + final Object callback, final Object... callbacks) { + return Collections.emptyMap(); + } +} diff --git a/dd-java-agent/instrumentation/rs/jax-rs/jax-rs-annotations/jax-rs-annotations-2.0/src/main/java/datadog/trace/instrumentation/jaxrs2/JaxRsAnnotationsInstrumentation.java b/dd-java-agent/instrumentation/rs/jax-rs/jax-rs-annotations/jax-rs-annotations-2.0/src/main/java/datadog/trace/instrumentation/jaxrs2/JaxRsAnnotationsInstrumentation.java index 4f7605a508e..cb57a5ba42b 100644 --- a/dd-java-agent/instrumentation/rs/jax-rs/jax-rs-annotations/jax-rs-annotations-2.0/src/main/java/datadog/trace/instrumentation/jaxrs2/JaxRsAnnotationsInstrumentation.java +++ b/dd-java-agent/instrumentation/rs/jax-rs/jax-rs-annotations/jax-rs-annotations-2.0/src/main/java/datadog/trace/instrumentation/jaxrs2/JaxRsAnnotationsInstrumentation.java @@ -8,6 +8,8 @@ import static datadog.trace.agent.tooling.bytebuddy.matcher.HierarchyMatchers.isAnnotatedWith; import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.namedOneOf; +import static datadog.trace.bootstrap.ResourceMethodSpanTracker.enter; +import static datadog.trace.bootstrap.ResourceMethodSpanTracker.exit; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activateSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.startSpan; @@ -135,6 +137,9 @@ public static ContextScope nameSpan( if (contextStore != null && asyncResponse != null) { contextStore.put(asyncResponse, span); + // Only tracked for methods that can hand off to AsyncResponse#resume()/cancel(); + // see ResourceMethodSpanTracker for why this can't just be a bare counter. + enter(span); } return scope; @@ -148,8 +153,16 @@ public static void stopSpan( if (scope == null) { return; } + if (asyncResponse != null) { + exit(); + } final AgentSpan span = spanFromScope(scope); if (throwable != null) { + if (asyncResponse != null) { + // Clear span from the asyncResponse so a later resume()/cancel() call (e.g. from + // container exception-mapping) doesn't find a stale mapping and double-finish it. + InstrumentationContext.get(AsyncResponse.class, AgentSpan.class).put(asyncResponse, null); + } DECORATE.onError(span, throwable); DECORATE.beforeFinish(span); scope.close(); diff --git a/dd-java-agent/instrumentation/rs/jax-rs/jax-rs-annotations/jax-rs-annotations-2.0/src/main/java/datadog/trace/instrumentation/jaxrs2/JaxRsAsyncResponseInstrumentation.java b/dd-java-agent/instrumentation/rs/jax-rs/jax-rs-annotations/jax-rs-annotations-2.0/src/main/java/datadog/trace/instrumentation/jaxrs2/JaxRsAsyncResponseInstrumentation.java index 9972a720d8c..0a308de91cf 100644 --- a/dd-java-agent/instrumentation/rs/jax-rs/jax-rs-annotations/jax-rs-annotations-2.0/src/main/java/datadog/trace/instrumentation/jaxrs2/JaxRsAsyncResponseInstrumentation.java +++ b/dd-java-agent/instrumentation/rs/jax-rs/jax-rs-annotations/jax-rs-annotations-2.0/src/main/java/datadog/trace/instrumentation/jaxrs2/JaxRsAsyncResponseInstrumentation.java @@ -2,6 +2,7 @@ import static datadog.trace.agent.tooling.bytebuddy.matcher.HierarchyMatchers.implementsInterface; import static datadog.trace.agent.tooling.bytebuddy.matcher.NameMatchers.named; +import static datadog.trace.bootstrap.ResourceMethodSpanTracker.isInnermost; import static datadog.trace.instrumentation.jaxrs2.JaxRsAnnotationsDecorator.DECORATE; import static net.bytebuddy.matcher.ElementMatchers.isPublic; import static net.bytebuddy.matcher.ElementMatchers.takesArgument; @@ -74,8 +75,18 @@ public static void stopSpan( final AgentSpan span = contextStore.get(asyncResponse); if (span != null) { - contextStore.put(asyncResponse, null); DECORATE.onError(span, throwable); + if (isInnermost(span)) { + // resume()/cancel() was called synchronously, nested inside the still-running + // resource method that owns this span (ResourceMethodSpanTracker.isInnermost(span) + // says its invocation is still the innermost open one on this thread). Let that + // method's own exit advice close the scope and finish the span (it will see + // asyncResponse.isSuspended() == false) instead of finishing it here, which would + // both double-finish the span and finish it prematurely while the resource method + // may still be doing work under it. + return; + } + contextStore.put(asyncResponse, null); DECORATE.beforeFinish(span); span.finish(); } @@ -94,8 +105,12 @@ public static void stopSpan( final AgentSpan span = contextStore.get(asyncResponse); if (span != null) { - contextStore.put(asyncResponse, null); DECORATE.onError(span, throwable); + if (isInnermost(span)) { + // see comment in AsyncResponseAdvice#stopSpan + return; + } + contextStore.put(asyncResponse, null); DECORATE.beforeFinish(span); span.finish(); } @@ -113,12 +128,16 @@ public static void stopSpan( final AgentSpan span = contextStore.get(asyncResponse); if (span != null) { - contextStore.put(asyncResponse, null); if (throwable != null) { DECORATE.onError(span, throwable); } else { span.setTag("canceled", true); } + if (isInnermost(span)) { + // see comment in AsyncResponseAdvice#stopSpan + return; + } + contextStore.put(asyncResponse, null); DECORATE.beforeFinish(span); span.finish(); }