From 1b800bc9a109c10c696a553a7e03474bcc02e7e9 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Thu, 20 Aug 2026 13:44:54 -0700 Subject: [PATCH 01/10] Emit LLM Obs spans over the intake track when using DDAgentWriter The DDAgentWriter branch of WriterFactory never wired up an LLM Obs DDIntakeWriter, so LLM Obs spans were silently dropped whenever the tracer picked DDAgentWriter (e.g. in Lambda/serverless mode with CI Visibility disabled). Broadcast to a dedicated LLM Obs DDIntakeWriter via MultiWriter, flushed synchronously in serverless environments just like the primary writer. --- .../trace/common/writer/WriterFactory.java | 21 ++++++++++++++++++- 1 file changed, 20 insertions(+), 1 deletion(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java index a32c096c2b0..ca23ab7b354 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java @@ -135,7 +135,7 @@ public static Writer createWriter( } } - RemoteWriter remoteWriter; + Writer remoteWriter; if (DD_INTAKE_WRITER_TYPE.equals(configuredType)) { final TrackType trackType = DDIntakeTrackTypeResolver.resolve(config); final RemoteApi remoteApi = @@ -220,6 +220,25 @@ public static Writer createWriter( } remoteWriter = builder.build(); + + // DDAgentWriter only speaks the regular trace protocol, so LLM Observability spans + // (tagged with DDSpanTypes.LLMOBS) need their own track sent via the EVP proxy -- without + // this, they're silently dropped by the agent/extension instead of reaching LLM Obs + // intake. Flush this track synchronously too when the primary writer is (i.e. in a + // serverless environment), so spans aren't lost when the execution environment freezes. + if (config.isLlmObsEnabled()) { + final RemoteApi llmObsApi = + createDDIntakeRemoteApi(config, commObjects, featuresDiscovery, TrackType.LLMOBS); + final DDIntakeWriter llmObsWriter = + DDIntakeWriter.builder() + .addTrack(TrackType.LLMOBS, llmObsApi) + .healthMetrics(healthMetrics) + .monitoring(commObjects.monitoring) + .alwaysFlush(alwaysFlush) + .flushIntervalMilliseconds(flushIntervalMilliseconds) + .build(); + remoteWriter = new MultiWriter(new Writer[] {remoteWriter, llmObsWriter}); + } } return remoteWriter; From 10fdd039c5570247706077d5a162396eef128265 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Mon, 24 Aug 2026 13:29:29 -0700 Subject: [PATCH 02/10] Flush the DDIntakeWriter track instead of adding a second LLM Obs writer The DDAgentWriter branch's dedicated LLM Obs DDIntakeWriter duplicated the LLM Obs track that the DD_INTAKE_WRITER_TYPE branch already builds whenever llmObsEnabled is true (Agent.java defaults to MultiWriter:DDIntakeWriter,DDAgentWriter for that case), so every span was sent twice. Verified live on Lambda: the pre-fix build sent LLMOBS/v2 payloads twice per invocation. The original bug this branch was fixing -- spans silently dropped in Lambda -- was actually caused by the DDIntakeWriter branch never setting alwaysFlush, so the periodic flush timer rarely beat the execution environment freezing after the handler returns. Verified live: without alwaysFlush, 0/8 invocations delivered a span; with it, 8/8 did, one send each. --- .../trace/common/writer/WriterFactory.java | 29 ++++++------------- 1 file changed, 9 insertions(+), 20 deletions(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java index ca23ab7b354..4c256fb4cb3 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java @@ -135,12 +135,19 @@ public static Writer createWriter( } } - Writer remoteWriter; + RemoteWriter remoteWriter; if (DD_INTAKE_WRITER_TYPE.equals(configuredType)) { final TrackType trackType = DDIntakeTrackTypeResolver.resolve(config); final RemoteApi remoteApi = createDDIntakeRemoteApi(config, commObjects, featuresDiscovery, trackType); + // In a serverless environment the execution environment can freeze as soon as the handler + // returns, before this writer's periodic flush timer next fires -- flush synchronously so + // buffered events (e.g. LLM Observability spans) aren't lost when that happens. + boolean alwaysFlush = + config.isAgentConfiguredUsingDefault() + && ServerlessInfo.get().isRunningInServerlessEnvironment(); + DDIntakeWriter.DDIntakeWriterBuilder builder = DDIntakeWriter.builder() .addTrack(trackType, remoteApi) @@ -148,6 +155,7 @@ public static Writer createWriter( .healthMetrics(healthMetrics) .monitoring(commObjects.monitoring) .singleSpanSampler(singleSpanSampler) + .alwaysFlush(alwaysFlush) .flushIntervalMilliseconds(flushIntervalMilliseconds); if (config.isCiVisibilityEnabled()) { @@ -220,25 +228,6 @@ public static Writer createWriter( } remoteWriter = builder.build(); - - // DDAgentWriter only speaks the regular trace protocol, so LLM Observability spans - // (tagged with DDSpanTypes.LLMOBS) need their own track sent via the EVP proxy -- without - // this, they're silently dropped by the agent/extension instead of reaching LLM Obs - // intake. Flush this track synchronously too when the primary writer is (i.e. in a - // serverless environment), so spans aren't lost when the execution environment freezes. - if (config.isLlmObsEnabled()) { - final RemoteApi llmObsApi = - createDDIntakeRemoteApi(config, commObjects, featuresDiscovery, TrackType.LLMOBS); - final DDIntakeWriter llmObsWriter = - DDIntakeWriter.builder() - .addTrack(TrackType.LLMOBS, llmObsApi) - .healthMetrics(healthMetrics) - .monitoring(commObjects.monitoring) - .alwaysFlush(alwaysFlush) - .flushIntervalMilliseconds(flushIntervalMilliseconds) - .build(); - remoteWriter = new MultiWriter(new Writer[] {remoteWriter, llmObsWriter}); - } } return remoteWriter; From 0fad3138ba018ee1c5ac56239249f996af828a26 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Tue, 25 Aug 2026 14:36:00 -0700 Subject: [PATCH 03/10] Clarify alwaysFlush comment refers to Lambda, not serverless generally --- .../java/datadog/trace/common/writer/WriterFactory.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java index 4c256fb4cb3..a7fa42b6d20 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java @@ -141,9 +141,9 @@ public static Writer createWriter( final RemoteApi remoteApi = createDDIntakeRemoteApi(config, commObjects, featuresDiscovery, trackType); - // In a serverless environment the execution environment can freeze as soon as the handler - // returns, before this writer's periodic flush timer next fires -- flush synchronously so - // buffered events (e.g. LLM Observability spans) aren't lost when that happens. + // In Lambda the execution environment can freeze as soon as the handler returns, before + // this writer's periodic flush timer next fires -- flush synchronously so buffered events + // (e.g. LLM Observability spans) aren't lost when that happens. boolean alwaysFlush = config.isAgentConfiguredUsingDefault() && ServerlessInfo.get().isRunningInServerlessEnvironment(); From 4b50f3a374a436b4152dc2578b7bf904dc457499 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Fri, 28 Aug 2026 11:06:46 -0700 Subject: [PATCH 04/10] Flush both queues in TraceProcessingWorker.flush() Sampled-out traces route to the secondary queue, which flush() never touched, so a Lambda invocation whose execution environment freezes right after the synchronous flush returns can lose them. Verified live: 1/15 delivered without this fix, 15/15 with it. Co-Authored-By: Claude Sonnet 5 --- .../common/writer/TraceProcessingWorker.java | 19 +++++++++++++------ .../writer/TraceProcessingWorkerTest.java | 5 +++-- 2 files changed, 16 insertions(+), 8 deletions(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java index 9b70b09bc9e..2a617fbe3dd 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java @@ -85,12 +85,12 @@ public void start() { } public boolean flush(long timeout, TimeUnit timeUnit) { - CountDownLatch latch = new CountDownLatch(1); - FlushEvent flush = new FlushEvent(latch); - boolean offered; - do { - offered = primaryQueue.offer(flush); - } while (!offered && serializerThread.isAlive()); + // flush both queues so sampled-out traces (routed to the secondary queue) aren't + // left behind, e.g. when a Lambda invocation's execution environment freezes + // right after this synchronous flush returns. + CountDownLatch latch = new CountDownLatch(2); + offer(primaryQueue, new FlushEvent(latch)); + offer(secondaryQueue, new FlushEvent(latch)); try { return latch.await(timeout, timeUnit); } catch (InterruptedException e) { @@ -99,6 +99,13 @@ public boolean flush(long timeout, TimeUnit timeUnit) { } } + private void offer(MessagePassingBlockingQueue queue, FlushEvent flush) { + boolean offered; + do { + offered = queue.offer(flush); + } while (!offered && serializerThread.isAlive()); + } + @Override public void close() { spanSamplingWorker.close(); diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java index 86dbfa61b22..c2ee8fd3c7d 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java @@ -141,9 +141,10 @@ void testAFlushShouldClearThePrimaryQueue() { worker.start(); boolean flushed = worker.flush(10, TimeUnit.SECONDS); - // the flush succeeds, triggers a dispatch, and the queue is empty + // the flush succeeds, triggers a dispatch for both the primary and secondary + // queue's flush events, and the primary queue is empty assertTrue(flushed); - assertEquals(1, flushCount.get()); + assertEquals(2, flushCount.get()); assertTrue(worker.getPrimaryQueue().isEmpty()); } } From 49c342ce76d62c873ed7671be6c1c32f732bdac1 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Fri, 28 Aug 2026 11:08:26 -0700 Subject: [PATCH 05/10] Trim flush() comment Co-Authored-By: Claude Sonnet 5 --- .../datadog/trace/common/writer/TraceProcessingWorker.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java index 2a617fbe3dd..1a6171998b9 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java @@ -85,9 +85,7 @@ public void start() { } public boolean flush(long timeout, TimeUnit timeUnit) { - // flush both queues so sampled-out traces (routed to the secondary queue) aren't - // left behind, e.g. when a Lambda invocation's execution environment freezes - // right after this synchronous flush returns. + // flush both queues so sampled-out traces (routed to the secondary queue) aren't left behind CountDownLatch latch = new CountDownLatch(2); offer(primaryQueue, new FlushEvent(latch)); offer(secondaryQueue, new FlushEvent(latch)); From 68ba2fc69f6e5dd39507b75446e0ce51075198b1 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Fri, 28 Aug 2026 14:37:31 -0700 Subject: [PATCH 06/10] Deduplicate serverless-default check in WriterFactory Both the DDIntakeWriter and DDAgentWriter branches computed the same config.isAgentConfiguredUsingDefault() && isRunningInServerlessEnvironment() condition separately; hoist it into a single isServerlessDefault local. Co-Authored-By: Claude Sonnet 5 --- .../trace/common/writer/WriterFactory.java | 19 +++++++++---------- 1 file changed, 9 insertions(+), 10 deletions(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java index a7fa42b6d20..0ad61b81e39 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java @@ -135,19 +135,19 @@ 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); final RemoteApi remoteApi = createDDIntakeRemoteApi(config, commObjects, featuresDiscovery, trackType); - // In Lambda the execution environment can freeze as soon as the handler returns, before - // this writer's periodic flush timer next fires -- flush synchronously so buffered events - // (e.g. LLM Observability spans) aren't lost when that happens. - boolean alwaysFlush = - config.isAgentConfiguredUsingDefault() - && ServerlessInfo.get().isRunningInServerlessEnvironment(); - DDIntakeWriter.DDIntakeWriterBuilder builder = DDIntakeWriter.builder() .addTrack(trackType, remoteApi) @@ -155,7 +155,7 @@ public static Writer createWriter( .healthMetrics(healthMetrics) .monitoring(commObjects.monitoring) .singleSpanSampler(singleSpanSampler) - .alwaysFlush(alwaysFlush) + .alwaysFlush(isServerlessDefault) .flushIntervalMilliseconds(flushIntervalMilliseconds); if (config.isCiVisibilityEnabled()) { @@ -176,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"); From ea181b568ab129fca9e2dc7f26b7071288220fcb Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Mon, 28 Sep 2026 14:27:03 -0700 Subject: [PATCH 07/10] Flush the secondary queue only in Lambda and bound flush() by its timeout Co-Authored-By: Claude Opus 5.5 --- .../common/writer/TraceProcessingWorker.java | 36 ++++++++---- .../writer/TraceProcessingWorkerTest.java | 57 ++++++++++++++++++- 2 files changed, 80 insertions(+), 13 deletions(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java index 1a6171998b9..60dc4bdcf89 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java @@ -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; @@ -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, @@ -85,23 +92,32 @@ public void start() { } public boolean flush(long timeout, TimeUnit timeUnit) { - // flush both queues so sampled-out traces (routed to the secondary queue) aren't left behind - CountDownLatch latch = new CountDownLatch(2); - offer(primaryQueue, new FlushEvent(latch)); - offer(secondaryQueue, new FlushEvent(latch)); + long deadline = System.nanoTime() + timeUnit.toNanos(timeout); + CountDownLatch latch = new CountDownLatch(flushSecondaryQueue ? 2 : 1); + FlushEvent flush = new FlushEvent(latch); 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 void offer(MessagePassingBlockingQueue queue, FlushEvent flush) { - boolean offered; - do { - offered = queue.offer(flush); - } while (!offered && serializerThread.isAlive()); + private boolean offer(MessagePassingBlockingQueue queue, FlushEvent flush, long deadline) + throws InterruptedException { + while (serializerThread.isAlive()) { + long remaining = deadline - System.nanoTime(); + if (remaining <= 0) { + return false; + } + if (queue.offer(flush)) { + return true; + } + NANOSECONDS.sleep(Math.min(remaining, MILLISECONDS.toNanos(1))); + } + return false; } @Override diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java index c2ee8fd3c7d..dbc057e2933 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java @@ -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; @@ -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; @@ -141,10 +146,9 @@ void testAFlushShouldClearThePrimaryQueue() { worker.start(); boolean flushed = worker.flush(10, TimeUnit.SECONDS); - // the flush succeeds, triggers a dispatch for both the primary and secondary - // queue's flush events, and the primary queue is empty + // the flush succeeds, triggers a dispatch, and the queue is empty assertTrue(flushed); - assertEquals(2, flushCount.get()); + assertEquals(1, flushCount.get()); assertTrue(worker.getPrimaryQueue().isEmpty()); } } @@ -613,4 +617,51 @@ public > 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 flushed = caller.submit(() -> worker.flush(100, TimeUnit.MILLISECONDS)); + assertFalse(flushed.get(5, TimeUnit.SECONDS)); + } finally { + releaseFlush.countDown(); + worker.close(); + caller.shutdownNow(); + } + } } From 3428a6987c68b7991eb980c72425fdf1524a56e1 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Mon, 28 Sep 2026 14:37:03 -0700 Subject: [PATCH 08/10] Always attempt the flush offer at least once before checking the deadline Co-Authored-By: Claude Opus 5.5 --- .../datadog/trace/common/writer/TraceProcessingWorker.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java index 60dc4bdcf89..664b00f23b6 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java @@ -108,13 +108,13 @@ public boolean flush(long timeout, TimeUnit timeUnit) { private boolean offer(MessagePassingBlockingQueue queue, FlushEvent flush, long deadline) throws InterruptedException { while (serializerThread.isAlive()) { + if (queue.offer(flush)) { + return true; + } long remaining = deadline - System.nanoTime(); if (remaining <= 0) { return false; } - if (queue.offer(flush)) { - return true; - } NANOSECONDS.sleep(Math.min(remaining, MILLISECONDS.toNanos(1))); } return false; From 9f3a63faa1e187b6c3f4b3c8693b5c0cf6d202a5 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Fri, 9 Oct 2026 12:58:52 -0700 Subject: [PATCH 09/10] Enqueue the flush before checking the serializer is alive The merge queue reported infra failures (git checkout timeouts) for some attempts, but every attempt that got far enough also failed DemoExecutorServiceTest on IBM JDK 8 (7 of 7 runs since Oct 2; the test passes on ibm8 elsewhere). The queue cancels the pipeline at the first failed job, so the bot comment does not always show this failure. Cause: on IBM JDK 8 the agent delays starting the writer by at least 100ms (okHttpDelayMillis, because of IBMSASL). The ExecutorService app exits before that, so the shutdown flush runs before the serializer thread starts. This branch checked serializerThread.isAlive() before the first offer, so flush() returned false at once, the writer closed, and the trace was never sent. Before this branch, flush() enqueued the FlushEvent and waited, and the late-starting serializer drained it. Fix: always offer once, and check the thread only while the queue is full. The new test fails without this change (expected true, was false) on any JVM. Co-Authored-By: Claude Opus 5.5 --- .../common/writer/TraceProcessingWorker.java | 10 +++--- .../writer/TraceProcessingWorkerTest.java | 34 +++++++++++++++++++ 2 files changed, 38 insertions(+), 6 deletions(-) diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java index 664b00f23b6..b2705c5fec4 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/TraceProcessingWorker.java @@ -107,17 +107,15 @@ public boolean flush(long timeout, TimeUnit timeUnit) { private boolean offer(MessagePassingBlockingQueue queue, FlushEvent flush, long deadline) throws InterruptedException { - while (serializerThread.isAlive()) { - if (queue.offer(flush)) { - return true; - } + // 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) { + if (remaining <= 0 || !serializerThread.isAlive()) { return false; } NANOSECONDS.sleep(Math.min(remaining, MILLISECONDS.toNanos(1))); } - return false; + return true; } @Override diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java index dbc057e2933..f25fdfe3473 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/TraceProcessingWorkerTest.java @@ -664,4 +664,38 @@ void testFlushTimesOutWhenQueueIsFullAndSerializerIsBlocked() throws Exception { 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 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(); + } + } } From b85d757ee444263733510c38b027109b4c99bc24 Mon Sep 17 00:00:00 2001 From: Rey Abolofia Date: Fri, 9 Oct 2026 13:11:08 -0700 Subject: [PATCH 10/10] Run CI on all JVM vendors [ci: NON_DEFAULT_JVMS] Co-Authored-By: Claude Opus 5.5