diff --git a/dd-java-agent/instrumentation-testing/src/main/groovy/datadog/trace/agent/test/base/HttpServerTest.groovy b/dd-java-agent/instrumentation-testing/src/main/groovy/datadog/trace/agent/test/base/HttpServerTest.groovy index 46abd76a6da..ed867b49643 100644 --- a/dd-java-agent/instrumentation-testing/src/main/groovy/datadog/trace/agent/test/base/HttpServerTest.groovy +++ b/dd-java-agent/instrumentation-testing/src/main/groovy/datadog/trace/agent/test/base/HttpServerTest.groovy @@ -1264,7 +1264,6 @@ abstract class HttpServerTest extends WithHttpServer { } } - @Flaky(value = "https://github.com/DataDog/dd-trace-java/issues/9396", suites = ["PekkoHttpServerInstrumentationAsyncHttp2Test"]) def "Instrumentation test exception"() { setup: def method = "GET" diff --git a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/build.gradle b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/build.gradle index a68a645badb..c181e3b2238 100644 --- a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/build.gradle +++ b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/build.gradle @@ -47,8 +47,10 @@ configurations { sourceSets { latestDepTest.groovy.srcDir sourceSets.baseTest.groovy latestDepTest.scala.srcDir sourceSets.baseTest.scala + latestDepTest.java.srcDir sourceSets.baseTest.java latestPekko10Test.groovy.srcDir sourceSets.baseTest.groovy latestPekko10Test.scala.srcDir sourceSets.baseTest.scala + latestPekko10Test.java.srcDir sourceSets.baseTest.java } dependencies { @@ -59,6 +61,7 @@ dependencies { // These are the common dependencies that are inherited by the other test sets testImplementation(libs.testing.okhttp3) + testImplementation libs.bundles.mockito testImplementation project(':dd-java-agent:instrumentation:datadog:tracing:trace-annotation') testImplementation project(':dd-java-agent:instrumentation:pekko:pekko-concurrent-1.0') testImplementation project(':dd-java-agent:instrumentation:scala:scala-concurrent-2.8') diff --git a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/java/AbstractPekkoHttpAsyncHandlerWrapperTest.java b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/java/AbstractPekkoHttpAsyncHandlerWrapperTest.java new file mode 100644 index 00000000000..5fdac160d90 --- /dev/null +++ b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/java/AbstractPekkoHttpAsyncHandlerWrapperTest.java @@ -0,0 +1,282 @@ +import static datadog.trace.agent.test.assertions.SpanMatcher.span; +import static datadog.trace.agent.test.assertions.TraceMatcher.trace; +import static datadog.trace.api.DDSpanTypes.HTTP_SERVER; +import static datadog.trace.bootstrap.instrumentation.api.AgentSpan.fromContext; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import datadog.context.Context; +import datadog.context.ContextScope; +import datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.agent.test.assertions.SpanMatcher; +import datadog.trace.instrumentation.pekkohttp.DatadogAsyncHandlerWrapper; +import datadog.trace.test.junit.utils.config.WithConfig; +import java.lang.reflect.Field; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; +import java.util.regex.Pattern; +import org.apache.pekko.http.scaladsl.model.HttpRequest; +import org.apache.pekko.http.scaladsl.model.HttpRequest$; +import org.apache.pekko.http.scaladsl.model.HttpResponse; +import org.apache.pekko.http.scaladsl.model.HttpResponse$; +import org.apache.pekko.http.scaladsl.model.Uri$; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; +import scala.Function1; +import scala.concurrent.ExecutionContext; +import scala.concurrent.ExecutionContext$; +import scala.concurrent.ExecutionContextExecutorService; +import scala.concurrent.Future; +import scala.concurrent.Promise; +import scala.concurrent.Promise$; +import scala.runtime.AbstractFunction1; +import scala.util.Failure; +import scala.util.Success; +import scala.util.Try; + +/** + * Verifies that completing the handler Future with the request context active does not propagate + * that context to the Future exposed to Pekko. Blocking the simulated Pekko callback makes trace + * retention deterministic. + */ +abstract class AbstractPekkoHttpAsyncHandlerWrapperTest extends AbstractInstrumentationTest { + + /** + * The operation name is a {@code UTF8BytesString}; {@link SpanMatcher#operationName(String)} + * compares object types, while the {@link Pattern} overload compares {@code CharSequence} + * content. + */ + private static final Pattern OPERATION_NAME = Pattern.compile("pekko-http\\.request"); + + protected abstract boolean expectedCompletionPriority(); + + @BeforeEach + void verifyPropagationMode() throws Exception { + // Load the injected helper only after the test agent is installed, then verify each concrete + // variant is exercising the intended Scala Promise propagation mode. + assertEquals( + expectedCompletionPriority(), + readStaticBoolean( + Class.forName("datadog.trace.instrumentation.scala.PromiseHelper"), + "completionPriority"), + "Unexpected Scala Promise propagation mode"); + } + + private static boolean readStaticBoolean(final Class type, final String name) + throws Exception { + final Field field = type.getDeclaredField(name); + field.setAccessible(true); + return (Boolean) field.get(null); + } + + /** + * Covers the failed-response path. On Scala 2.12 {@code Promise.resolveTry} allocates a fresh + * {@code Failure}, which incidentally drops any context associated with the completing {@code + * Try}, so this path exercises the root-context attachment only. Scala 2.13 passes the {@code + * Try} through, so there it exercises both defenses. + */ + @Test + void doesNotPropagateRequestContextWhenHandlerFails() throws Exception { + assertRequestTraceIsNotRetained( + new Failure<>(new Exception("controller exception")), + span().root().operationName(OPERATION_NAME).type(HTTP_SERVER).error()); + } + + /** + * Covers the successful-response path. Both Scala generations pass a {@code Success} through + * completion unchanged, so this is the case that pins the defensive {@code Try} copy in + * completion-priority mode. + */ + @Test + void doesNotPropagateRequestContextWhenHandlerSucceeds() throws Exception { + assertRequestTraceIsNotRetained( + new Success<>(emptyResponse()), + span().root().operationName(OPERATION_NAME).type(HTTP_SERVER).error(false)); + } + + private void assertRequestTraceIsNotRetained( + final Try handlerResult, final SpanMatcher expectedSpan) throws Exception { + try (AsyncHandlerWrapperReproducer reproducer = + new AsyncHandlerWrapperReproducer(handlerResult)) { + reproducer.start(); + + assertTrue(reproducer.awaitFrameworkCallback(), "Framework callback did not start"); + assertTrue( + writer.waitForTracesMax(1, 5), + "Request trace was held by the contextless framework callback"); + assertTraces(trace(expectedSpan)); + } + } + + @Test + void propagatesFatalCallbackFailure() { + assertCallbackFailure(new LinkageError("callback linkage failure"), true); + } + + @Test + void completesFutureWithNonFatalCallbackFailure() { + assertCallbackFailure(new IllegalStateException("callback failure"), false); + } + + @SuppressWarnings({"unchecked", "rawtypes"}) + private void assertCallbackFailure(final Throwable failure, final boolean fatal) { + Future handlerFuture = mock(Future.class); + Context[] requestContext = new Context[1]; + DatadogAsyncHandlerWrapper wrapper = + new DatadogAsyncHandlerWrapper( + new AbstractFunction1>() { + @Override + public Future apply(HttpRequest request) { + requestContext[0] = Context.current(); + return handlerFuture; + } + }, + mock(ExecutionContext.class)); + Future response = wrapper.apply(emptyRequest()); + try { + ArgumentCaptor callback = ArgumentCaptor.forClass(Function1.class); + verify(handlerFuture).onComplete(callback.capture(), any()); + // Inject an error while reading the result to exercise the callback's exception boundary. + Try result = mock(Try.class); + when(result.isSuccess()).thenReturn(true); + when(result.get()).thenThrow(failure); + if (fatal) { + assertSame( + failure, assertThrows(LinkageError.class, () -> callback.getValue().apply(result))); + assertFalse(response.isCompleted()); + } else { + callback.getValue().apply(result); + assertTrue(response.isCompleted()); + assertSame(failure, ((Failure) response.value().get()).exception()); + } + } finally { + // The injected failure happens before finishSpan can finish the request span. + fromContext(requestContext[0]).finish(); + } + assertTraces(trace(span().root().operationName(OPERATION_NAME).type(HTTP_SERVER).error(false))); + } + + private static HttpRequest emptyRequest() { + return HttpRequest$.MODULE$.apply( + HttpRequest$.MODULE$.apply$default$1(), + Uri$.MODULE$.apply("/exception"), + HttpRequest$.MODULE$.apply$default$3(), + HttpRequest$.MODULE$.apply$default$4(), + HttpRequest$.MODULE$.apply$default$5()); + } + + private static HttpResponse emptyResponse() { + return HttpResponse$.MODULE$.apply( + HttpResponse$.MODULE$.apply$default$1(), + HttpResponse$.MODULE$.apply$default$2(), + HttpResponse$.MODULE$.apply$default$3(), + HttpResponse$.MODULE$.apply$default$4()); + } + + private static final class AsyncHandlerWrapperReproducer implements AutoCloseable { + private final Try handlerResult; + private final Promise handlerPromise = Promise$.MODULE$.apply(); + private Context requestContext; + private final CountDownLatch callbackStarted = new CountDownLatch(1); + private final CountDownLatch releaseCallback = new CountDownLatch(1); + + private final ExecutionContextExecutorService handlerExecutor = + ExecutionContext$.MODULE$.fromExecutorService(Executors.newSingleThreadExecutor()); + private final ExecutionContextExecutorService frameworkExecutor = + ExecutionContext$.MODULE$.fromExecutorService(Executors.newSingleThreadExecutor()); + + AsyncHandlerWrapperReproducer(final Try handlerResult) { + this.handlerResult = handlerResult; + } + + void start() { + DatadogAsyncHandlerWrapper wrapper = + new DatadogAsyncHandlerWrapper( + new AbstractFunction1>() { + @Override + public Future apply(HttpRequest request) { + requestContext = Context.current(); + return handlerPromise.future(); + } + }, + handlerExecutor); + + Future response = wrapper.apply(emptyRequest()); + + // Model a callback registered by Pekko after the wrapper has closed the request scope. It + // should not inherit the request context when the handler Future completes. + response.onComplete( + new AbstractFunction1, Void>() { + @Override + public Void apply(Try result) { + callbackStarted.countDown(); + try { + releaseCallback.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + return null; + } + }, + frameworkExecutor); + + // Application futures normally complete from callbacks running with the propagated request + // context. Model that completion path so the test also covers context captured at promise + // dispatch time, not only at callback construction time. + try (ContextScope ignored = requestContext.attach()) { + handlerPromise.complete(handlerResult); + } + } + + boolean awaitFrameworkCallback() throws InterruptedException { + return callbackStarted.await(5, TimeUnit.SECONDS); + } + + @Override + public void close() throws InterruptedException { + releaseCallback.countDown(); + handlerExecutor.shutdown(); + frameworkExecutor.shutdown(); + assertTrue( + handlerExecutor.awaitTermination(5, TimeUnit.SECONDS), + "Handler executor did not terminate"); + assertTrue( + frameworkExecutor.awaitTermination(5, TimeUnit.SECONDS), + "Framework executor did not terminate"); + } + } +} + +/** Runs the async-handler context-retention reproducer with default Scala Promise propagation. */ +class PekkoHttpAsyncHandlerWrapperTest extends AbstractPekkoHttpAsyncHandlerWrapperTest { + + @Override + protected boolean expectedCompletionPriority() { + return false; + } +} + +/** + * Runs the async-handler context-retention reproducer with completion-priority propagation, which + * associates the completing context with the resolved {@code Try} instead of the thread. + * + *

{@link WithConfig} applies before the test agent is installed, so the Scala Promise + * instrumentation registers the advice that creates that association. + */ +@WithConfig(key = "trace.integration.scala_promise_completion_priority.enabled", value = "true") +class PekkoHttpAsyncHandlerWrapperForkedTest extends AbstractPekkoHttpAsyncHandlerWrapperTest { + + @Override + protected boolean expectedCompletionPriority() { + return true; + } +} diff --git a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogAsyncHandlerWrapper.java b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogAsyncHandlerWrapper.java index 0578eda8aba..0b5281bdf0e 100644 --- a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogAsyncHandlerWrapper.java +++ b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogAsyncHandlerWrapper.java @@ -1,18 +1,26 @@ package datadog.trace.instrumentation.pekkohttp; -import static datadog.trace.bootstrap.instrumentation.api.AgentSpan.fromContext; - +import datadog.context.Context; import datadog.context.ContextScope; -import datadog.trace.bootstrap.instrumentation.api.AgentSpan; +import datadog.trace.api.InstrumenterConfig; import org.apache.pekko.http.scaladsl.model.HttpRequest; import org.apache.pekko.http.scaladsl.model.HttpResponse; import scala.Function1; import scala.concurrent.ExecutionContext; import scala.concurrent.Future; +import scala.concurrent.Promise; +import scala.concurrent.Promise$; import scala.runtime.AbstractFunction1; +import scala.util.Failure; +import scala.util.Success; +import scala.util.Try; +import scala.util.control.NonFatal$; public class DatadogAsyncHandlerWrapper extends AbstractFunction1> { + private static final boolean SCALA_PROMISE_COMPLETION_PRIORITY_ENABLED = + InstrumenterConfig.get().isScalaPromiseCompletionPriorityEnabled(); + private final Function1> userHandler; private final ExecutionContext executionContext; @@ -26,33 +34,59 @@ public DatadogAsyncHandlerWrapper( @Override public Future apply(final HttpRequest request) { final ContextScope scope = DatadogWrapperHelper.createSpan(request); - AgentSpan span = fromContext(scope.context()); - Future futureResponse; + final Context context = scope.context(); + Future handlerFuture; try { - futureResponse = userHandler.apply(request); + handlerFuture = userHandler.apply(request); } catch (final Throwable t) { scope.close(); - DatadogWrapperHelper.finishSpan(scope.context(), t); + DatadogWrapperHelper.finishSpan(context, t); throw t; } - final Future wrapped = - futureResponse.transform( - new AbstractFunction1() { - @Override - public HttpResponse apply(final HttpResponse response) { - DatadogWrapperHelper.finishSpan(scope.context(), response); - return response; + scope.close(); + final Promise frameworkPromise = Promise$.MODULE$.apply(); + handlerFuture.onComplete( + new AbstractFunction1, Void>() { + @Override + public Void apply(final Try result) { + Try completion = result; + try { + if (result.isSuccess()) { + DatadogWrapperHelper.finishSpan(context, result.get()); + } else { + DatadogWrapperHelper.finishSpan( + context, ((Failure) result).exception()); } - }, - new AbstractFunction1() { - @Override - public Throwable apply(final Throwable t) { - DatadogWrapperHelper.finishSpan(scope.context(), t); - return t; + } catch (final Throwable t) { + if (!NonFatal$.MODULE$.apply(t)) { + throw t; } - }, - executionContext); - scope.close(); - return wrapped; + // Preserve non-fatal decoration failures in the returned Future. + // Pekko does not support response blocking, and this wrapper has no Materializer + // for discarding a successful response entity replaced by this failure. + completion = new Failure<>(t); + } + final Try frameworkResult; + if (SCALA_PROMISE_COMPLETION_PRIORITY_ENABLED) { + // Completion-priority mode can associate context directly with a Try. Copy the + // result so neither that association nor the active thread context reaches Pekko. + frameworkResult = + completion.isSuccess() + ? new Success<>(completion.get()) + : new Failure<>(((Failure) completion).exception()); + } else { + frameworkResult = completion; + } + // The application Future can complete while the request context is active. Complete + // the Future exposed to Pekko under the root context so its framework callbacks do not + // capture and retain the request trace. + try (ContextScope ignored = Context.root().attach()) { + frameworkPromise.complete(frameworkResult); + } + return null; + } + }, + executionContext); + return frameworkPromise.future(); } } diff --git a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogWrapperHelper.java b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogWrapperHelper.java index 3cee9b327b5..7766d5018e6 100644 --- a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogWrapperHelper.java +++ b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogWrapperHelper.java @@ -64,18 +64,22 @@ public void close() { public static void finishSpan(final Context context, final HttpResponse response) { final AgentSpan span = fromContext(context); - DECORATE.onResponse(span, response); - DECORATE.beforeFinish(context); - - span.finish(); + try { + DECORATE.onResponse(span, response); + DECORATE.beforeFinish(context); + } finally { + span.finish(); + } } public static void finishSpan(final Context context, final Throwable t) { final AgentSpan span = fromContext(context); - DECORATE.onError(span, t); - span.setHttpStatusCode(500); - DECORATE.beforeFinish(context); - - span.finish(); + try { + DECORATE.onError(span, t); + span.setHttpStatusCode(500); + DECORATE.beforeFinish(context); + } finally { + span.finish(); + } } } diff --git a/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-2.10/src/main/java/datadog/trace/instrumentation/scala210/concurrent/ScalaPromiseModule.java b/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-2.10/src/main/java/datadog/trace/instrumentation/scala210/concurrent/ScalaPromiseModule.java index 80a0f8cf947..a5916f54932 100644 --- a/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-2.10/src/main/java/datadog/trace/instrumentation/scala210/concurrent/ScalaPromiseModule.java +++ b/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-2.10/src/main/java/datadog/trace/instrumentation/scala210/concurrent/ScalaPromiseModule.java @@ -14,7 +14,6 @@ import datadog.trace.bootstrap.instrumentation.java.concurrent.State; import java.util.ArrayList; import java.util.Collection; -import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -65,8 +64,7 @@ public List typeInstrumentations() { } // Only enable this if integrations have been enabled and the extra "integration" // scala_promise_completion_priority has been enabled specifically - if (config.isIntegrationEnabled( - Collections.singletonList("scala_promise_completion_priority"), false)) { + if (config.isScalaPromiseCompletionPriorityEnabled()) { instrumenters.add(new PromiseObjectInstrumentation()); } return instrumenters; diff --git a/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-2.13/src/main/java/datadog/trace/instrumentation/scala213/concurrent/ScalaPromiseModule.java b/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-2.13/src/main/java/datadog/trace/instrumentation/scala213/concurrent/ScalaPromiseModule.java index 42d013039ab..0ec164f1379 100644 --- a/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-2.13/src/main/java/datadog/trace/instrumentation/scala213/concurrent/ScalaPromiseModule.java +++ b/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-2.13/src/main/java/datadog/trace/instrumentation/scala213/concurrent/ScalaPromiseModule.java @@ -14,7 +14,6 @@ import datadog.trace.bootstrap.instrumentation.java.concurrent.State; import java.util.ArrayList; import java.util.Collection; -import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -61,8 +60,7 @@ public List typeInstrumentations() { final InstrumenterConfig config = InstrumenterConfig.get(); // Only enable this if integrations have been enabled and the extra "integration" // scala_promise_completion_priority has been enabled specifically - if (config.isIntegrationEnabled( - Collections.singletonList("scala_promise_completion_priority"), false)) { + if (config.isScalaPromiseCompletionPriorityEnabled()) { ret.add(new DefaultPromiseInstrumentation()); ret.add(new PromiseObjectInstrumentation()); } diff --git a/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-common/src/main/java/datadog/trace/instrumentation/scala/PromiseHelper.java b/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-common/src/main/java/datadog/trace/instrumentation/scala/PromiseHelper.java index b3e6b27d7ac..6506bf44d37 100644 --- a/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-common/src/main/java/datadog/trace/instrumentation/scala/PromiseHelper.java +++ b/dd-java-agent/instrumentation/scala/scala-promise/scala-promise-common/src/main/java/datadog/trace/instrumentation/scala/PromiseHelper.java @@ -7,16 +7,13 @@ import datadog.trace.bootstrap.ContextStore; import datadog.trace.bootstrap.instrumentation.java.concurrent.AdviceUtils; import datadog.trace.bootstrap.instrumentation.java.concurrent.State; -import java.util.Collections; import scala.util.Failure; import scala.util.Success; import scala.util.Try; public class PromiseHelper { public static final boolean completionPriority = - InstrumenterConfig.get() - .isIntegrationEnabled( - Collections.singletonList("scala_promise_completion_priority"), false); + InstrumenterConfig.get().isScalaPromiseCompletionPriorityEnabled(); /** * Get the {@code Try} that should be associated with the {@code Context}. Will create a new copy diff --git a/internal-api/src/main/java/datadog/trace/api/InstrumenterConfig.java b/internal-api/src/main/java/datadog/trace/api/InstrumenterConfig.java index 84b407a6836..76b65d43036 100644 --- a/internal-api/src/main/java/datadog/trace/api/InstrumenterConfig.java +++ b/internal-api/src/main/java/datadog/trace/api/InstrumenterConfig.java @@ -100,6 +100,7 @@ import static datadog.trace.api.config.UsmConfig.USM_ENABLED; import static datadog.trace.util.CollectionUtils.tryMakeImmutableList; import static datadog.trace.util.CollectionUtils.tryMakeImmutableSet; +import static java.util.Collections.singletonList; import datadog.environment.JavaVirtualMachine; import datadog.trace.api.profiling.ProfilingEnablement; @@ -143,6 +144,15 @@ public class InstrumenterConfig { } } + /** + * Name of the opt-in Scala Promise integration that gives the context completing a {@code + * Promise} priority over the context that registered callbacks on it. Defined here so that the + * instrumentations enabling this mode and the ones that have to compensate for it cannot drift + * apart. + */ + private static final String SCALA_PROMISE_COMPLETION_PRIORITY = + "scala_promise_completion_priority"; + private final ConfigProvider configProvider; private final boolean triageEnabled; @@ -231,6 +241,7 @@ public class InstrumenterConfig { private final boolean appLogsCollectionEnabled; private final boolean legacyContextManagerEnabled; + private final boolean scalaPromiseCompletionPriorityEnabled; static { // Bind telemetry collector to config module before initializing ConfigProvider @@ -400,6 +411,9 @@ private InstrumenterConfig() { configProvider.getBoolean(APP_LOGS_COLLECTION_ENABLED, DEFAULT_APP_LOGS_COLLECTION_ENABLED); legacyContextManagerEnabled = configProvider.getBoolean(LEGACY_CONTEXT_MANAGER_ENABLED, true); + + scalaPromiseCompletionPriorityEnabled = + isIntegrationEnabled(singletonList(SCALA_PROMISE_COMPLETION_PRIORITY), false); } public boolean isCodeOriginEnabled() { @@ -453,6 +467,20 @@ public boolean isIntegrationEnabled( return anyEnabled; } + /** + * Whether the Scala Promise instrumentation gives the context completing a {@code Promise} + * priority over the context that registered callbacks on it. + * + *

This mode associates the completing context with the resolved {@code Try} itself, so it is + * also read by instrumentations that must keep such an association from reaching a framework + * callback. + * + * @return {@code true} if completion-priority propagation is enabled, else {@code false} + */ + public boolean isScalaPromiseCompletionPriorityEnabled() { + return scalaPromiseCompletionPriorityEnabled; + } + public boolean isIntegrationShortcutMatchingEnabled( final Iterable integrationNames, final boolean defaultEnabled) { return configProvider.isEnabled( @@ -882,6 +910,8 @@ public String toString() { + apiSecurityEndpointCollectionEnabled + ", legacyContextManagerEnabled=" + legacyContextManagerEnabled + + ", scalaPromiseCompletionPriorityEnabled=" + + scalaPromiseCompletionPriorityEnabled + '}'; } }