From 7afe7deec79a610be951619f0b07d3ecd6fb873b Mon Sep 17 00:00:00 2001 From: Andrea Marziali Date: Wed, 30 Sep 2026 14:43:38 +0200 Subject: [PATCH 1/2] Fix Akka and Pekko HTTP flow scope ownership --- ...tadogServerRequestResponseFlowWrapper.java | 43 ++++++++----------- ...tadogServerRequestResponseFlowWrapper.java | 42 ++++++++---------- 2 files changed, 37 insertions(+), 48 deletions(-) diff --git a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogServerRequestResponseFlowWrapper.java b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogServerRequestResponseFlowWrapper.java index 545aeff7600..97f29ae5398 100644 --- a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogServerRequestResponseFlowWrapper.java +++ b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogServerRequestResponseFlowWrapper.java @@ -1,7 +1,6 @@ package datadog.trace.instrumentation.akkahttp; import static datadog.trace.bootstrap.instrumentation.api.AgentSpan.fromContext; -import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; import akka.http.scaladsl.model.HttpRequest; import akka.http.scaladsl.model.HttpResponse; @@ -14,6 +13,7 @@ import akka.stream.stage.AbstractOutHandler; import akka.stream.stage.GraphStage; import akka.stream.stage.GraphStageLogic; +import datadog.context.Context; import datadog.context.ContextScope; import datadog.trace.api.gateway.RequestContext; import datadog.trace.bootstrap.instrumentation.api.AgentSpan; @@ -57,7 +57,8 @@ public GraphStageLogic createLogic(final Attributes inheritedAttributes) throws // that this connection was created with. This means that we can safely // close the span at the front of the queue when we receive the response // from the user code, since it will match up to the request for that span. - final Queue scopes = new ArrayBlockingQueue<>(pipeliningLimit); + // Actor invocation cleanup owns the scopes; only contexts cross the response boundary. + final Queue contexts = new ArrayBlockingQueue<>(pipeliningLimit); boolean[] skipNextPull = new boolean[] {false}; // This is where the request comes in from the server and TCP layer @@ -85,7 +86,7 @@ public void onPush() throws Exception { } } - scopes.add(scope); + contexts.add(scope.context()); push(requestOutlet, request); // Legacy mode leaves the scope open so the surrounding actor can clean it up. // Context-manager mode swaps the context and the actor restores it on exit. @@ -129,9 +130,9 @@ public void onDownstreamFinish() throws Exception { @Override public void onPush() throws Exception { HttpResponse response = grab(responseInlet); - final ContextScope scope = scopes.poll(); - if (scope != null) { - AgentSpan span = fromContext(scope.context()); + final Context context = contexts.poll(); + if (context != null) { + AgentSpan span = fromContext(context); HttpResponse newResponse = BlockingResponseHelper.handleFinishForWaf(span, response); if (newResponse != response) { @@ -139,12 +140,7 @@ public void onPush() throws Exception { response.discardEntityBytes(materializer()); response = newResponse; } - DatadogWrapperHelper.finishSpan(scope.context(), response); - // Legacy mode may still own the scope when the response arrives. - AgentSpan activeSpan = activeSpan(); - if (activeSpan == span) { - scope.close(); - } + DatadogWrapperHelper.finishSpan(context, response); } push(responseOutlet, response); } @@ -153,28 +149,27 @@ public void onPush() throws Exception { public void onUpstreamFinish() throws Exception { // We will not receive any more responses from the user code, so clean up any // remaining spans - ContextScope scope = scopes.poll(); - while (scope != null) { - fromContext(scope.context()).finish(); - scope = scopes.poll(); + Context context = contexts.poll(); + while (context != null) { + fromContext(context).finish(); + context = contexts.poll(); } completeStage(); } @Override public void onUpstreamFailure(final Throwable ex) throws Exception { - ContextScope scope = scopes.poll(); - if (scope != null) { + Context context = contexts.poll(); + if (context != null) { // Mark the span as failed - AgentSpan span = fromContext(scope.context()); - DatadogWrapperHelper.finishSpan(scope.context(), ex); + DatadogWrapperHelper.finishSpan(context, ex); } // We will not receive any more responses from the user code, so clean up any // remaining spans - scope = scopes.poll(); - while (scope != null) { - fromContext(scope.context()).finish(); - scope = scopes.poll(); + context = contexts.poll(); + while (context != null) { + fromContext(context).finish(); + context = contexts.poll(); } fail(responseOutlet, ex); } diff --git a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogServerRequestResponseFlowWrapper.java b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogServerRequestResponseFlowWrapper.java index 0de81b07a08..862f4675948 100644 --- a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogServerRequestResponseFlowWrapper.java +++ b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogServerRequestResponseFlowWrapper.java @@ -1,10 +1,9 @@ package datadog.trace.instrumentation.pekkohttp; import static datadog.trace.bootstrap.instrumentation.api.AgentSpan.fromContext; -import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.activeSpan; +import datadog.context.Context; import datadog.context.ContextScope; -import datadog.trace.bootstrap.instrumentation.api.AgentSpan; import java.util.Queue; import java.util.concurrent.ArrayBlockingQueue; import org.apache.pekko.http.scaladsl.model.HttpRequest; @@ -55,7 +54,8 @@ public GraphStageLogic createLogic(final Attributes inheritedAttributes) throws // that this connection was created with. This means that we can safely // close the span at the front of the queue when we receive the response // from the user code, since it will match up to the request for that span. - final Queue scopes = new ArrayBlockingQueue<>(pipeliningLimit); + // Actor invocation cleanup owns the scopes; only contexts cross the response boundary. + final Queue contexts = new ArrayBlockingQueue<>(pipeliningLimit); // This is where the request comes in from the server and TCP layer setHandler( @@ -65,7 +65,7 @@ public GraphStageLogic createLogic(final Attributes inheritedAttributes) throws public void onPush() throws Exception { final HttpRequest request = grab(requestInlet); final ContextScope scope = DatadogWrapperHelper.createSpanForFlow(request); - scopes.add(scope); + contexts.add(scope.context()); push(requestOutlet, request); // Legacy mode leaves the scope open so the surrounding actor can clean it up. // Context-manager mode swaps the context and the actor restores it on exit. @@ -109,15 +109,9 @@ public void onDownstreamFinish() throws Exception { @Override public void onPush() throws Exception { final HttpResponse response = grab(responseInlet); - final ContextScope scope = scopes.poll(); - if (scope != null) { - DatadogWrapperHelper.finishSpan(scope.context(), response); - // Legacy mode may still own the scope when the response arrives. - AgentSpan activeSpan = activeSpan(); - AgentSpan span = fromContext(scope.context()); - if (activeSpan == span) { - scope.close(); - } + final Context context = contexts.poll(); + if (context != null) { + DatadogWrapperHelper.finishSpan(context, response); } push(responseOutlet, response); } @@ -126,27 +120,27 @@ public void onPush() throws Exception { public void onUpstreamFinish() throws Exception { // We will not receive any more responses from the user code, so clean up any // remaining spans - ContextScope scope = scopes.poll(); - while (scope != null) { - fromContext(scope.context()).finish(); - scope = scopes.poll(); + Context context = contexts.poll(); + while (context != null) { + fromContext(context).finish(); + context = contexts.poll(); } completeStage(); } @Override public void onUpstreamFailure(final Throwable ex) throws Exception { - ContextScope scope = scopes.poll(); - if (scope != null) { + Context context = contexts.poll(); + if (context != null) { // Mark the span as failed - DatadogWrapperHelper.finishSpan(scope.context(), ex); + DatadogWrapperHelper.finishSpan(context, ex); } // We will not receive any more responses from the user code, so clean up any // remaining spans - scope = scopes.poll(); - while (scope != null) { - fromContext(scope.context()).finish(); - scope = scopes.poll(); + context = contexts.poll(); + while (context != null) { + fromContext(context).finish(); + context = contexts.poll(); } fail(responseOutlet, ex); } From 4040e8fa19d5b90ed16882488cc4f20215fe8493 Mon Sep 17 00:00:00 2001 From: Andrea Marziali Date: Thu, 1 Oct 2026 10:07:09 +0200 Subject: [PATCH 2/2] Fix dangling scopes --- .../AkkaHttpServerInstrumentationTest.groovy | 8 +++++ .../scala/AkkaHttpTestWebServer.scala | 34 +++++++++++++++++++ ...tadogServerRequestResponseFlowWrapper.java | 2 ++ .../akkahttp/DatadogWrapperHelper.java | 16 +++++++++ .../PekkoHttpServerInstrumentationTest.groovy | 8 +++++ .../scala/PekkoHttpTestWebServer.scala | 34 +++++++++++++++++++ ...tadogServerRequestResponseFlowWrapper.java | 2 ++ .../pekkohttp/DatadogWrapperHelper.java | 16 +++++++++ 8 files changed, 120 insertions(+) diff --git a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/groovy/AkkaHttpServerInstrumentationTest.groovy b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/groovy/AkkaHttpServerInstrumentationTest.groovy index 2d807d5e215..f8500d99970 100644 --- a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/groovy/AkkaHttpServerInstrumentationTest.groovy +++ b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/groovy/AkkaHttpServerInstrumentationTest.groovy @@ -283,6 +283,14 @@ class AkkaHttpServerInstrumentationAsyncTest extends AkkaHttpServerInstrumentati } class AkkaHttpServerInstrumentationBindAndHandleTest extends AkkaHttpServerInstrumentationTest { + def "restore context before processing a response downstream"() { + when: + server.checkResponseContext() + + then: + TEST_WRITER.waitForTraces(1) + } + String akkaHttpVersion @Override diff --git a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/scala/AkkaHttpTestWebServer.scala b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/scala/AkkaHttpTestWebServer.scala index 8fce57838e2..05d732e7ef2 100644 --- a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/scala/AkkaHttpTestWebServer.scala +++ b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/baseTest/scala/AkkaHttpTestWebServer.scala @@ -51,6 +51,40 @@ class AkkaHttpTestWebServer(binder: Binder) extends HttpServer { private var port: Int = 0 private var portBinding: Future[ServerBinding] = _ + def checkResponseContext(): Unit = { + import akka.stream.scaladsl.{BidiFlow, Flow, Sink, Source} + import datadog.context.Context + import datadog.trace.bootstrap.instrumentation.api.AgentSpan + import datadog.trace.core.DDSpan + import datadog.trace.instrumentation.akkahttp.DatadogServerRequestResponseFlowWrapper + + var previousContext: Context = null + var requestSpan: AgentSpan = null + val handler = Flow[HttpRequest].map { _ => + requestSpan = activeSpan() + HttpResponse() + } + val flow = BidiFlow + .fromGraph(new DatadogServerRequestResponseFlowWrapper(ServerSettings(system))) + .reversed + .join(handler) + val result = Source + .single(HttpRequest(uri = "/response-context")) + .map { request => + previousContext = Context.current() + request + } + .via(flow) + .map { response => + assert(requestSpan.asInstanceOf[DDSpan].isFinished) + assert(Context.current() eq previousContext, "request context is still active downstream") + response + } + .runWith(Sink.ignore) + + Await.result(result, 10 seconds) + } + override def start(): Unit = { portBinding = Await.ready(binder.bind(0), 10 seconds) port = portBinding.value.get.get.localAddress.getPort diff --git a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogServerRequestResponseFlowWrapper.java b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogServerRequestResponseFlowWrapper.java index 97f29ae5398..585cba75bc9 100644 --- a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogServerRequestResponseFlowWrapper.java +++ b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogServerRequestResponseFlowWrapper.java @@ -141,6 +141,7 @@ public void onPush() throws Exception { response = newResponse; } DatadogWrapperHelper.finishSpan(context, response); + DatadogWrapperHelper.deactivateFlowContext(context); } push(responseOutlet, response); } @@ -163,6 +164,7 @@ public void onUpstreamFailure(final Throwable ex) throws Exception { if (context != null) { // Mark the span as failed DatadogWrapperHelper.finishSpan(context, ex); + DatadogWrapperHelper.deactivateFlowContext(context); } // We will not receive any more responses from the user code, so clean up any // remaining spans diff --git a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogWrapperHelper.java b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogWrapperHelper.java index cf0a5d68258..0a3cc867d89 100644 --- a/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogWrapperHelper.java +++ b/dd-java-agent/instrumentation/akka/akka-http/akka-http-10.0/src/main/java/datadog/trace/instrumentation/akkahttp/DatadogWrapperHelper.java @@ -1,6 +1,8 @@ package datadog.trace.instrumentation.akkahttp; import static datadog.trace.bootstrap.instrumentation.api.AgentSpan.fromContext; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.checkpointActiveForRollback; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.rollbackActiveToCheckpoint; import static datadog.trace.instrumentation.akkahttp.AkkaHttpServerDecorator.DECORATE; import akka.http.scaladsl.model.HttpRequest; @@ -78,4 +80,18 @@ public static void finishSpan(final Context context, final Throwable t) { span.finish(); } + + public static void deactivateFlowContext(final Context context) { + if (context != Context.current()) { + return; + } + if (LEGACY_CONTEXT_MANAGER_ENABLED) { + // Close request scopes left active for stream propagation, then restore the actor checkpoint. + rollbackActiveToCheckpoint(); + checkpointActiveForRollback(); + } else { + // There is one current context; detach it now and let actor exit restore its saved context. + Context.root().swap(); + } + } } diff --git a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/groovy/PekkoHttpServerInstrumentationTest.groovy b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/groovy/PekkoHttpServerInstrumentationTest.groovy index 1034cf83e26..21f10261533 100644 --- a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/groovy/PekkoHttpServerInstrumentationTest.groovy +++ b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/groovy/PekkoHttpServerInstrumentationTest.groovy @@ -116,6 +116,14 @@ class PekkoHttpServerInstrumentationAsyncTest extends PekkoHttpServerInstrumenta } class PekkoHttpServerInstrumentationBindAndHandleTest extends PekkoHttpServerInstrumentationTest { + def "restore context before processing a response downstream"() { + when: + server.checkResponseContext() + + then: + TEST_WRITER.waitForTraces(1) + } + @Override HttpServer server() { return new PekkoHttpTestWebServer(PekkoHttpTestWebServer.BindAndHandle()) diff --git a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/scala/PekkoHttpTestWebServer.scala b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/scala/PekkoHttpTestWebServer.scala index 7e070ac16a5..66f3915d763 100644 --- a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/scala/PekkoHttpTestWebServer.scala +++ b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/baseTest/scala/PekkoHttpTestWebServer.scala @@ -36,6 +36,40 @@ class PekkoHttpTestWebServer(binder: Binder) extends HttpServer { private var port: Int = 0 private var portBinding: Future[ServerBinding] = null + def checkResponseContext(): Unit = { + import datadog.context.Context + import datadog.trace.bootstrap.instrumentation.api.AgentSpan + import datadog.trace.core.DDSpan + import datadog.trace.instrumentation.pekkohttp.DatadogServerRequestResponseFlowWrapper + import org.apache.pekko.stream.scaladsl.{BidiFlow, Flow, Sink, Source} + + var previousContext: Context = null + var requestSpan: AgentSpan = null + val handler = Flow[HttpRequest].map { _ => + requestSpan = activeSpan() + HttpResponse() + } + val flow = BidiFlow + .fromGraph(new DatadogServerRequestResponseFlowWrapper(ServerSettings(system))) + .reversed + .join(handler) + val result = Source + .single(HttpRequest(uri = "/response-context")) + .map { request => + previousContext = Context.current() + request + } + .via(flow) + .map { response => + assert(requestSpan.asInstanceOf[DDSpan].isFinished) + assert(Context.current() eq previousContext, "request context is still active downstream") + response + } + .runWith(Sink.ignore) + + Await.result(result, 10 seconds) + } + override def start(): Unit = { portBinding = Await.ready(binder.bind(0), 10 seconds) port = portBinding.value.get.get.localAddress.getPort diff --git a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogServerRequestResponseFlowWrapper.java b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogServerRequestResponseFlowWrapper.java index 862f4675948..d870a42deed 100644 --- a/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogServerRequestResponseFlowWrapper.java +++ b/dd-java-agent/instrumentation/pekko/pekko-http-1.0/src/main/java/datadog/trace/instrumentation/pekkohttp/DatadogServerRequestResponseFlowWrapper.java @@ -112,6 +112,7 @@ public void onPush() throws Exception { final Context context = contexts.poll(); if (context != null) { DatadogWrapperHelper.finishSpan(context, response); + DatadogWrapperHelper.deactivateFlowContext(context); } push(responseOutlet, response); } @@ -134,6 +135,7 @@ public void onUpstreamFailure(final Throwable ex) throws Exception { if (context != null) { // Mark the span as failed DatadogWrapperHelper.finishSpan(context, ex); + DatadogWrapperHelper.deactivateFlowContext(context); } // We will not receive any more responses from the user code, so clean up any // remaining spans 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..65645d5784a 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 @@ -1,6 +1,8 @@ package datadog.trace.instrumentation.pekkohttp; import static datadog.trace.bootstrap.instrumentation.api.AgentSpan.fromContext; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.checkpointActiveForRollback; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.rollbackActiveToCheckpoint; import static datadog.trace.instrumentation.pekkohttp.PekkoHttpServerDecorator.DECORATE; import datadog.context.Context; @@ -78,4 +80,18 @@ public static void finishSpan(final Context context, final Throwable t) { span.finish(); } + + public static void deactivateFlowContext(final Context context) { + if (context != Context.current()) { + return; + } + if (LEGACY_CONTEXT_MANAGER_ENABLED) { + // Close request scopes left active for stream propagation, then restore the actor checkpoint. + rollbackActiveToCheckpoint(); + checkpointActiveForRollback(); + } else { + // There is one current context; detach it now and let actor exit restore its saved context. + Context.root().swap(); + } + } }