Skip to content
Merged
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 @@ -29,6 +29,12 @@ private static final class RateLimiterHolder {
ProfilingConfig.PROFILING_UPLOAD_PERIOD_DEFAULT)));
}

/** Whether queue timing can safely start. */
public static boolean isReady() {
// Queue timing is unsupported in native images. JFR must initialise the TSC frequency first.
return !Platform.isNativeImage() && InstrumentationBasedProfiling.isJFRReady();
}

public static <T> void startQueuingTimer(
ContextStore<T, State> taskContextStore,
Class<?> schedulerClass,
Expand All @@ -41,14 +47,8 @@ public static <T> void startQueuingTimer(

public static void startQueuingTimer(
State state, Class<?> schedulerClass, Class<?> queueClass, int queueLength, Object task) {
if (Platform.isNativeImage()) {
// explicitly not supported for Graal native image
return;
}
// TODO consider queue length based sampling here to reduce overhead
// avoid calling this before JFR is initialised because it will lead to reading the wrong
// TSC frequency before JFR has set it up properly
if (task != null && state != null && InstrumentationBasedProfiling.isJFRReady()) {
if (task != null && state != null && isReady()) {
QueueTiming timing =
(QueueTiming) AgentTracer.get().getProfilingContext().start(Timer.TimerType.QUEUEING);
timing.setTask(task);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ public static void capture(
// queue time needs to be handled separately because there are RunnableFutures which are
// excluded as
// Runnables but it is not until now that they will be put on the executor's queue
if (!exclude(EXECUTOR, tpe)) {
if (QueueTimerHelper.isReady() && !exclude(EXECUTOR, tpe)) {
if (!exclude(RUNNABLE, task)) {
Queue<?> queue = tpe.getQueue();
QueueTimerHelper.startQueuingTimer(
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
import datadog.trace.agent.test.InstrumentationSpecification
import datadog.trace.bootstrap.instrumentation.api.AgentTracer
import datadog.trace.bootstrap.instrumentation.jfr.InstrumentationBasedProfiling

import java.util.concurrent.CountDownLatch
import java.util.concurrent.LinkedBlockingQueue
import java.util.concurrent.ThreadPoolExecutor
import java.util.concurrent.TimeUnit
import java.util.concurrent.atomic.AtomicInteger

import static datadog.trace.agent.test.utils.TraceUtils.basicSpan
import static datadog.trace.agent.test.utils.TraceUtils.runUnderTrace

class QueueTimingReadinessForkedTest extends InstrumentationSpecification {
@Override
protected void configurePreAgent() {
injectSysConfig("dd.profiling.enabled", "true")
injectSysConfig("dd.profiling.queueing.time.enabled", "true")
super.configurePreAgent()
}

def "queue metadata is read only after JFR is ready while tracing always propagates"() {
setup:
def queue = new CountingQueue()
def executor = new ThreadPoolExecutor(1, 1, 0, TimeUnit.SECONDS, queue)
assert !InstrumentationBasedProfiling.isJFRReady()

when:
runTasks(executor, "before")

then:
queue.sizeCalls.get() == 0
TEST_PROFILING_CONTEXT_INTEGRATION.closedTimings.isEmpty()
assertTraces(2) {
trace(1) {
basicSpan(it, "execute-before")
}
trace(1) {
basicSpan(it, "submit-before")
}
}

when:
TEST_WRITER.clear()
InstrumentationBasedProfiling.enableInstrumentationBasedProfiling()
runTasks(executor, "after")

then:
queue.sizeCalls.get() == 2
assertTraces(2) {
trace(1) {
basicSpan(it, "execute-after")
}
trace(1) {
basicSpan(it, "submit-after")
}
}
TEST_PROFILING_CONTEXT_INTEGRATION.isBalanced()
TEST_PROFILING_CONTEXT_INTEGRATION.closedTimings.size() == 2
TEST_PROFILING_CONTEXT_INTEGRATION.closedTimings.every {
it.queue == CountingQueue && it.scheduler == ThreadPoolExecutor && it.queueLength >= 0
}

cleanup:
executor.shutdownNow()
assert executor.awaitTermination(5, TimeUnit.SECONDS)
TEST_PROFILING_CONTEXT_INTEGRATION.closedTimings.clear()
}

private static void runTasks(ThreadPoolExecutor executor, String phase) {
runUnderTrace("execute-" + phase) {
def expectedSpanId = AgentTracer.activeSpan().spanId
def task = new ContextCheckingTask()
executor.execute(task)
assert task.done.await(5, TimeUnit.SECONDS)
assert task.spanId == expectedSpanId
}
runUnderTrace("submit-" + phase) {
def expectedSpanId = AgentTracer.activeSpan().spanId
def task = new ContextCheckingTask()
executor.submit(task).get(5, TimeUnit.SECONDS)
assert task.spanId == expectedSpanId
}
}

static class ContextCheckingTask implements Runnable {
final CountDownLatch done = new CountDownLatch(1)
volatile long spanId

@Override
void run() {
def span = AgentTracer.activeSpan()
spanId = span == null ? 0 : span.spanId
done.countDown()
}
}

static class CountingQueue extends LinkedBlockingQueue<Runnable> {
final AtomicInteger sizeCalls = new AtomicInteger()

@Override
int size() {
sizeCalls.incrementAndGet()
return super.size()
}

@Override
boolean isEmpty() {
// Executor housekeeping must not count as a queue-timing metadata read.
return super.size() == 0
}
}
}
Loading