Skip to content
Closed
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
Expand Up @@ -41,6 +41,19 @@ public static ContextScope startTaskScope(State state) {
return null;
}

@Nullable
public static ContextScope startTpeTaskScope(State state) {
if (state != null) {
final State.TpeContinuation continuation = state.getAndResetTpeContinuation();
if (continuation != null) {
final ContextScope scope = continuation.resume();
continuation.stopTiming();
return scope;
}
}
return null;
}

public static void endTaskScope(final ContextScope scope) {
if (null != scope) {
scope.close();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,46 @@ public boolean captureAndSetContinuation(final Context context) {
return false;
}

/**
* Captures context for an unwrapped {@code ThreadPoolExecutor} submission.
*
* <p>The tagged continuation can only be consumed by {@code beforeExecute}; regular Runnable
* advice deliberately ignores it. This prevents a direct invocation, or another submission of the
* same Runnable, from stealing the queued submission's context.
*/
@Nullable
public TpeContinuation captureAndSetTpeContinuation(final Context context) {
while (true) {
ContextContinuation current = CONTINUATION.get(this);
if (current == null) {
if (!CONTINUATION.compareAndSet(this, null, CLAIMED)) {
continue;
}
try {
TpeContinuation continuation = new TpeContinuation(context.capture());
CONTINUATION.lazySet(this, continuation);
return continuation;
} catch (Throwable error) {
CONTINUATION.compareAndSet(this, CLAIMED, null);
throw error;
}
}
if (current == CLAIMED || current instanceof TpeContinuation) {
return null;
}
// Generic Executor advice can run before ThreadPoolExecutor advice for the same call.
// Transfer
// that continuation instead of treating the duplicate instrumentation as task reuse.
if (current.context() != context) {
return null;
}
TpeContinuation continuation = new TpeContinuation(current);
if (CONTINUATION.compareAndSet(this, current, continuation)) {
return continuation;
}
}
}

public boolean setOrCancelContinuation(final ContextContinuation continuation) {
if (CONTINUATION.compareAndSet(this, null, CLAIMED)) {
// lazy write is guaranteed to be seen by getAndSet
Expand All @@ -58,6 +98,28 @@ public void closeContinuation() {
}
}

@Nullable
public ContextContinuation getContinuation() {
ContextContinuation continuation = CONTINUATION.get(this);
return continuation == CLAIMED || continuation instanceof TpeContinuation ? null : continuation;
}

@Nullable
public ContextContinuation getCancellableContinuation() {
ContextContinuation continuation = CONTINUATION.get(this);
return continuation == CLAIMED ? null : continuation;
}

public void closeContinuation(ContextContinuation expected) {
if (expected != null && CONTINUATION.compareAndSet(this, expected, null)) {
if (expected instanceof TpeContinuation) {
((TpeContinuation) expected).cancel();
} else {
expected.release();
}
}
}

public Context getContext() {
ContextContinuation continuation = CONTINUATION.get(this);
if (null == continuation || CLAIMED == continuation) {
Expand All @@ -68,20 +130,54 @@ public Context getContext() {

@Nullable
public ContextContinuation getAndResetContinuation() {
ContextContinuation continuation = CONTINUATION.get(this);
if (null == continuation || CLAIMED == continuation) {
return null;
while (true) {
ContextContinuation continuation = CONTINUATION.get(this);
if (null == continuation
|| CLAIMED == continuation
|| continuation instanceof TpeContinuation) {
return null;
}
if (CONTINUATION.compareAndSet(this, continuation, null)) {
return continuation;
}
}
}

@Nullable
public TpeContinuation getAndResetTpeContinuation() {
while (true) {
ContextContinuation continuation = CONTINUATION.get(this);
if (!(continuation instanceof TpeContinuation)) {
return null;
}
if (CONTINUATION.compareAndSet(this, continuation, null)) {
return (TpeContinuation) continuation;
}
}
CONTINUATION.compareAndSet(this, continuation, null);
return continuation;
}

@Nullable
public TpeContinuation getTpeContinuation() {
ContextContinuation continuation = CONTINUATION.get(this);
return continuation instanceof TpeContinuation ? (TpeContinuation) continuation : null;
}

public void closeTpeContinuation(TpeContinuation expected) {
closeContinuation(expected);
}

public void setTiming(Timing timing) {
TIMING.lazySet(this, timing);
TpeContinuation continuation = getTpeContinuation();
if (continuation != null) {
continuation.setTiming(timing);
} else {
TIMING.lazySet(this, timing);
}
}

public boolean isTimed() {
return TIMING.get(this) != null;
TpeContinuation continuation = getTpeContinuation();
return TIMING.get(this) != null || (continuation != null && continuation.isTimed());
}

public void stopTiming() {
Expand All @@ -90,4 +186,55 @@ public void stopTiming() {
QueueTimerHelper.stopQueuingTimer(timing);
}
}

public static final class TpeContinuation implements ContextContinuation {
private final ContextContinuation delegate;
private volatile Timing timing;

private TpeContinuation(ContextContinuation delegate) {
this.delegate = delegate;
}

@Override
public ContextContinuation hold() {
delegate.hold();
return this;
}

@Override
public Context context() {
return delegate.context();
}

@Override
public datadog.context.ContextScope resume() {
return delegate.resume();
}

@Override
public void release() {
delegate.release();
}

private void setTiming(Timing timing) {
this.timing = timing;
}

private boolean isTimed() {
return timing != null;
}

public void stopTiming() {
Timing timing = this.timing;
this.timing = null;
if (timing != null) {
QueueTimerHelper.stopQueuingTimer(timing);
}
}

private void cancel() {
release();
stopTiming();
}
}
}
Original file line number Diff line number Diff line change
@@ -1,14 +1,23 @@
package datadog.trace.bootstrap.instrumentation.java.concurrent;

import static datadog.trace.bootstrap.instrumentation.java.concurrent.AdviceUtils.shouldCapture;
import static datadog.trace.bootstrap.instrumentation.java.concurrent.ExcludeFilter.ExcludeType.RUNNABLE;
import static datadog.trace.bootstrap.instrumentation.java.concurrent.ExcludeFilter.exclude;

import datadog.context.Context;
import datadog.context.ContextContinuation;
import datadog.context.ContextScope;
import datadog.trace.api.GenericClassValue;
import datadog.trace.api.InstrumenterConfig;
import datadog.trace.api.Platform;
import datadog.trace.bootstrap.ContextStore;
import java.util.Set;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingDeque;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.LinkedTransferQueue;
import java.util.concurrent.SynchronousQueue;
import java.util.concurrent.ThreadPoolExecutor;

/**
Expand Down Expand Up @@ -63,17 +72,43 @@ public static boolean shouldPropagate(ThreadPoolExecutor executor) {
&& PROPAGATE.get(executor.getClass());
}

public static void capture(ContextStore<Runnable, State> contextStore, Runnable task) {
public static Runnable captureOrWrap(
ContextStore<Runnable, State> contextStore,
Runnable task,
Context context,
ThreadPoolExecutor executor) {
if (task != null && !exclude(RUNNABLE, task) && shouldCapture(context)) {
State state = contextStore.getOrCreate(task, State.FACTORY);
if (state.captureAndSetTpeContinuation(context) == null) {
return canWrapCollision(executor.getQueue()) ? Wrapper.wrap(task, context) : null;
}
}
return task;
}

/**
* Retains the historical best-effort propagation for subclass overrides outside JDK admission.
*/
public static void captureLegacy(ContextStore<Runnable, State> contextStore, Runnable task) {
if (task != null && !exclude(RUNNABLE, task)) {
AdviceUtils.capture(contextStore, task);
}
}

private static boolean canWrapCollision(BlockingQueue<Runnable> queue) {
return queue instanceof ArrayBlockingQueue
|| queue instanceof LinkedBlockingQueue
|| queue instanceof LinkedBlockingDeque
|| queue instanceof LinkedTransferQueue
|| queue instanceof SynchronousQueue;
}

public static ContextScope startScope(ContextStore<Runnable, State> contextStore, Runnable task) {
if (task == null || exclude(RUNNABLE, task)) {
return null;
}
return AdviceUtils.startTaskScope(contextStore, task);
State state = contextStore.get(task);
return AdviceUtils.startTpeTaskScope(state);
}

public static void setThreadLocalScope(ContextScope scope, Runnable task) {
Expand Down Expand Up @@ -106,10 +141,53 @@ public static void endScope(ContextScope scope, Runnable task) {
AdviceUtils.endTaskScope(scope);
}

public static void cancelTask(ContextStore<Runnable, State> contextStore, Runnable task) {
public static final class RejectedTask {
private final Wrapper<?> wrapper;
private final ContextContinuation continuation;
private final ContextScope scope;

public RejectedTask(Wrapper<?> wrapper, ContextContinuation continuation) {
this.wrapper = wrapper;
this.continuation = continuation;
this.scope =
wrapper != null
? wrapper.activate()
: continuation == null ? null : continuation.resume();
}

public void close() {
try {
if (scope != null) {
scope.close();
}
} finally {
if (wrapper != null) {
wrapper.cancel();
} else if (continuation != null) {
continuation.release();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

P2 Do not release a resumed rejection twice

Executor saturation inflates continuation cancellation health metrics.

Assertion details
  • Input: Reject any propagated task through an instrumented rejection handler.
  • Expected: Rejection ends each continuation once and records one finished continuation.
  • Actual: Scope close finishes the resumed continuation. The following explicit release records a second cancellation.

Was this helpful? React 👍 or 👎
🤖 Datadog Autotest · What is Autotest? · @DataDog review to ask questions · Any feedback? Reach out in #autotest · Open Bits AI session

}
}
}
}

public static ContextContinuation prepareRejectedTask(
ContextStore<Runnable, State> contextStore, Runnable task) {
if (task == null || exclude(RUNNABLE, task)) {
return;
return null;
}
State state = contextStore.get(task);
if (state == null) {
return null;
}
State.TpeContinuation tpeContinuation = state.getAndResetTpeContinuation();
if (tpeContinuation != null) {
tpeContinuation.stopTiming();
return tpeContinuation;
}
ContextContinuation continuation = state.getAndResetContinuation();
if (continuation != null) {
state.stopTiming();
}
AdviceUtils.cancelTask(contextStore, task);
return continuation;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -13,23 +13,41 @@ public class Wrapper<T extends Runnable> implements Runnable, AutoCloseable {

@SuppressWarnings({"unchecked", "rawtypes"})
public static <T extends Runnable> Runnable wrap(T task) {
if (task instanceof Wrapper
|| task instanceof RunnableFuture
|| task == null
|| exclude(RUNNABLE, task)) {
if (!isWrappable(task)) {
return task;
}
ContextContinuation continuation = captureActiveSpan();
if (continuation.context() != Context.root()) {
if (task instanceof Comparable) {
return new ComparableRunnable(task, continuation);
}
return new Wrapper<>(task, continuation);
return newWrapper(task, continuation);
}
// don't wrap unless there is scope to propagate
return task;
}

@SuppressWarnings({"unchecked", "rawtypes"})
public static <T extends Runnable> Runnable wrap(T task, Context context) {
if (context == Context.root() || !isWrappable(task)) {
return task;
}
return newWrapper(task, context.capture());
}

private static boolean isWrappable(Runnable task) {
return task != null
&& !(task instanceof Wrapper)
&& !(task instanceof RunnableFuture)
&& !exclude(RUNNABLE, task);
}

@SuppressWarnings({"unchecked", "rawtypes"})
private static <T extends Runnable> Runnable newWrapper(
T task, ContextContinuation continuation) {
if (task instanceof Comparable) {
return new ComparableRunnable(task, continuation);
}
return new Wrapper<>(task, continuation);
}

public static Runnable unwrap(Runnable task) {
return task instanceof Wrapper ? ((Wrapper<?>) task).unwrap() : task;
}
Expand Down Expand Up @@ -59,7 +77,7 @@ public T unwrap() {
return delegate;
}

private ContextScope activate() {
public ContextScope activate() {
return null == continuation ? null : continuation.resume();
}

Expand Down
Loading