Skip to content
Draft
Show file tree
Hide file tree
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
@@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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<ContextScope> scopes = new ArrayBlockingQueue<>(pipeliningLimit);
// Actor invocation cleanup owns the scopes; only contexts cross the response boundary.
final Queue<Context> contexts = new ArrayBlockingQueue<>(pipeliningLimit);
boolean[] skipNextPull = new boolean[] {false};

// This is where the request comes in from the server and TCP layer
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -129,22 +130,17 @@ 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) {
span.getRequestContext().getTraceSegment().effectivelyBlocked();
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);
Comment on lines +143 to 145

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Restore the context before pushing synchronous responses

When a fused or synchronous handler returns a response within the same actor invocation, actor cleanup has not run yet, so this finishes the server span while its context is still current and then push(responseOutlet, response) executes downstream stages under that finished span. The actor advice restores the checkpoint only when ActorCell.invoke exits, which can cause response-stage work, logging, or nested instrumentation to use an incorrect parent; retain a safe way to detach the request context before the push when it is still current. The parallel Pekko change at lines 114–116 has the same issue.

AGENTS.md reference: AGENTS.md:L43-L46

Useful? React with 👍 / 👎.

}
Expand All @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -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<ContextScope> scopes = new ArrayBlockingQueue<>(pipeliningLimit);
// Actor invocation cleanup owns the scopes; only contexts cross the response boundary.
final Queue<Context> contexts = new ArrayBlockingQueue<>(pipeliningLimit);

// This is where the request comes in from the server and TCP layer
setHandler(
Expand All @@ -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.
Expand Down Expand Up @@ -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);
}
Expand All @@ -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);
}
Expand Down
Loading