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..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,12 +85,10 @@ 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 + 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 +97,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/main/java/datadog/trace/common/writer/WriterFactory.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java index 73875a0e408..1a5abc85dba 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 @@ -134,6 +134,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); @@ -147,6 +154,7 @@ public static Writer createWriter( .healthMetrics(healthMetrics) .monitoring(commObjects.monitoring) .singleSpanSampler(singleSpanSampler) + .alwaysFlush(isServerlessDefault) .flushIntervalMilliseconds(flushIntervalMilliseconds); if (config.isCiVisibilityEnabled()) { @@ -167,8 +175,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"); 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()); } }