Skip to content
Open
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 @@ -4,7 +4,9 @@
import static datadog.trace.util.AgentThreadFactory.THREAD_JOIN_TIMOUT_MS;
import static datadog.trace.util.AgentThreadFactory.newAgentThread;
import static java.util.concurrent.TimeUnit.MILLISECONDS;
import static java.util.concurrent.TimeUnit.NANOSECONDS;

import datadog.common.container.ServerlessInfo;
import datadog.common.queue.MessagePassingBlockingQueue;
import datadog.common.queue.Queues;
import datadog.communication.ddagent.DroppingPolicy;
Expand Down Expand Up @@ -46,6 +48,11 @@ public class TraceProcessingWorker implements AutoCloseable {

private final SpanSamplingWorker spanSamplingWorker;

// In Lambda the environment can freeze as soon as a flush returns, so also flush the secondary
// queue there rather than leaving sampled-out traces for a periodic flush that may never run.
private final boolean flushSecondaryQueue =
ServerlessInfo.get().isRunningInServerlessEnvironment();

public TraceProcessingWorker(
final int capacity,
final HealthMetrics healthMetrics,
Expand Down Expand Up @@ -85,20 +92,32 @@ public void start() {
}

public boolean flush(long timeout, TimeUnit timeUnit) {
CountDownLatch latch = new CountDownLatch(1);
long deadline = System.nanoTime() + timeUnit.toNanos(timeout);
CountDownLatch latch = new CountDownLatch(flushSecondaryQueue ? 2 : 1);
FlushEvent flush = new FlushEvent(latch);
boolean offered;
do {
offered = primaryQueue.offer(flush);
} while (!offered && serializerThread.isAlive());
try {
return latch.await(timeout, timeUnit);
return offer(primaryQueue, flush, deadline)
&& (!flushSecondaryQueue || offer(secondaryQueue, flush, deadline))
&& latch.await(deadline - System.nanoTime(), NANOSECONDS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return false;
}
}

private boolean offer(MessagePassingBlockingQueue<Object> queue, FlushEvent flush, long deadline)
throws InterruptedException {
// offer before checking the serializer: it may not have started yet (e.g. on IBM JDK 8)
while (!queue.offer(flush)) {
long remaining = deadline - System.nanoTime();
if (remaining <= 0 || !serializerThread.isAlive()) {
return false;
}
NANOSECONDS.sleep(Math.min(remaining, MILLISECONDS.toNanos(1)));
}
return true;
}

@Override
public void close() {
spanSamplingWorker.close();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,13 @@ public static Writer createWriter(
}
}

// In Lambda the execution environment can freeze as soon as the handler returns, before a
// writer's periodic flush timer next fires -- flush synchronously so buffered events (e.g.
// LLM Observability spans) aren't lost when that happens.
boolean isServerlessDefault =
config.isAgentConfiguredUsingDefault()
&& ServerlessInfo.get().isRunningInServerlessEnvironment();

RemoteWriter remoteWriter;
if (DD_INTAKE_WRITER_TYPE.equals(configuredType)) {
final TrackType trackType = DDIntakeTrackTypeResolver.resolve(config);
Expand All @@ -148,6 +155,7 @@ public static Writer createWriter(
.healthMetrics(healthMetrics)
.monitoring(commObjects.monitoring)
.singleSpanSampler(singleSpanSampler)
.alwaysFlush(isServerlessDefault)
.flushIntervalMilliseconds(flushIntervalMilliseconds);

if (config.isCiVisibilityEnabled()) {
Expand All @@ -168,8 +176,7 @@ public static Writer createWriter(

} else { // configuredType == DDAgentWriter
boolean alwaysFlush = false;
if (config.isAgentConfiguredUsingDefault()
&& ServerlessInfo.get().isRunningInServerlessEnvironment()) {
if (isServerlessDefault) {
if (!ServerlessInfo.get().hasExtension()) {
log.info(
"Detected serverless environment. Serverless extension has not been detected, using PrintingWriter");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@

import datadog.trace.bootstrap.instrumentation.api.SpanPostProcessor;
import datadog.trace.common.sampling.SingleSpanSampler;
import datadog.trace.common.writer.ddagent.FlushEvent;
import datadog.trace.common.writer.ddagent.PrioritizationStrategy.PublishResult;
import datadog.trace.core.CoreSpan;
import datadog.trace.core.DDSpan;
Expand All @@ -25,6 +26,10 @@
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
Expand Down Expand Up @@ -612,4 +617,85 @@ public <T extends CoreSpan<T>> boolean setSamplingPriority(T span) {
assertEquals(sampledSingleSpans, sampledSpansCount.get());
}
}

// -------------------------------------------------------------------------
// Test 10: flush times out when the queue is full and the serializer is blocked
// -------------------------------------------------------------------------

@Test
void testFlushTimesOutWhenQueueIsFullAndSerializerIsBlocked() throws Exception {
CountDownLatch flushStarted = new CountDownLatch(1);
CountDownLatch releaseFlush = new CountDownLatch(1);
PayloadDispatcher dispatcher = mock(PayloadDispatcher.class);
doAnswer(
inv -> {
flushStarted.countDown();
releaseFlush.await();
return null;
})
.when(dispatcher)
.flush();
TraceProcessingWorker worker =
new TraceProcessingWorker(
2,
mock(HealthMetrics.class),
dispatcher,
() -> false,
FAST_LANE,
0, // disable periodic flushes
TimeUnit.SECONDS,
null);
ExecutorService caller = Executors.newSingleThreadExecutor();
try {
// block the serializer inside a flush, then fill the primary queue behind it
assertTrue(worker.getPrimaryQueue().offer(new FlushEvent(new CountDownLatch(1))));
worker.start();
assertTrue(flushStarted.await(5, TimeUnit.SECONDS));
while (worker.getPrimaryQueue().offer(Collections.singletonList(mock(DDSpan.class)))) {
// fill the queue
}

// the flush marker can't be enqueued, so the flush must give up once its timeout elapses
Future<Boolean> flushed = caller.submit(() -> worker.flush(100, TimeUnit.MILLISECONDS));
assertFalse(flushed.get(5, TimeUnit.SECONDS));
} finally {
releaseFlush.countDown();
worker.close();
caller.shutdownNow();
}
}

// -------------------------------------------------------------------------
// Test 11: flush succeeds when the serializer starts after the flush begins
// -------------------------------------------------------------------------

@Test
void testFlushSucceedsWhenSerializerStartsAfterFlushBegins() throws Exception {
AtomicInteger flushCount = new AtomicInteger();
TraceProcessingWorker worker =
new TraceProcessingWorker(
10,
mock(HealthMetrics.class),
flushCountingPayloadDispatcher(flushCount),
() -> false,
FAST_LANE,
100,
TimeUnit.SECONDS, // prevent heartbeats from helping the flush happen
null);
ExecutorService caller = Executors.newSingleThreadExecutor();
try {
// the agent can delay starting the writer (e.g. on IBM JDK 8) past a short-lived app's exit
Future<Boolean> flushed = caller.submit(() -> worker.flush(10, TimeUnit.SECONDS));
while (worker.getPrimaryQueue().isEmpty() && !flushed.isDone()) {
Thread.yield();
}
worker.start();

assertTrue(flushed.get(5, TimeUnit.SECONDS));
assertEquals(1, flushCount.get());
} finally {
worker.close();
caller.shutdownNow();
}
}
}
Loading