From 3fd40fcc952752b7233a6586859c110523728400 Mon Sep 17 00:00:00 2001 From: Brian Marks Date: Wed, 26 Aug 2026 20:23:31 -0400 Subject: [PATCH 1/7] Add OpenTelemetry metrics lifecycle controls --- .../trace/api/internal/InternalTracer.java | 9 + .../api/metrics/OpenTelemetryMetrics.java | 30 +++ .../api/metrics/OpenTelemetryMetricsTest.java | 49 ++++ .../java/datadog/trace/core/CoreTracer.java | 17 ++ .../core/otlp/metrics/OtlpMetricsService.java | 135 +++++++++-- .../otlp/metrics/OtlpMetricsServiceTest.java | 220 ++++++++++++++++++ 6 files changed, 440 insertions(+), 20 deletions(-) create mode 100644 dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java create mode 100644 dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java diff --git a/dd-trace-api/src/main/java/datadog/trace/api/internal/InternalTracer.java b/dd-trace-api/src/main/java/datadog/trace/api/internal/InternalTracer.java index 52b1adfc97e..2cfd4a264b1 100644 --- a/dd-trace-api/src/main/java/datadog/trace/api/internal/InternalTracer.java +++ b/dd-trace-api/src/main/java/datadog/trace/api/internal/InternalTracer.java @@ -2,6 +2,7 @@ import datadog.trace.api.experimental.DataStreamsCheckpointer; import datadog.trace.api.profiling.Profiling; +import java.util.concurrent.CompletableFuture; /** * Tracer internal features. Those features are not part of public API and can change or be removed @@ -20,6 +21,14 @@ public interface InternalTracer { void flushMetrics(); + default CompletableFuture forceFlushOtelMetrics() { + return CompletableFuture.completedFuture(false); + } + + default CompletableFuture shutdownOtelMetrics() { + return CompletableFuture.completedFuture(false); + } + void flushLogs(); Profiling getProfilingContext(); diff --git a/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java b/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java new file mode 100644 index 00000000000..12b52342191 --- /dev/null +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java @@ -0,0 +1,30 @@ +package datadog.trace.api.metrics; + +import datadog.trace.api.GlobalTracer; +import datadog.trace.api.Tracer; +import datadog.trace.api.internal.InternalTracer; +import java.util.concurrent.CompletableFuture; + +public final class OpenTelemetryMetrics { + private OpenTelemetryMetrics() {} + + public static CompletableFuture forceFlush() { + Tracer tracer = GlobalTracer.get(); + if (tracer instanceof InternalTracer) { + return ((InternalTracer) tracer).forceFlushOtelMetrics(); + } + return unavailable(); + } + + public static CompletableFuture shutdown() { + Tracer tracer = GlobalTracer.get(); + if (tracer instanceof InternalTracer) { + return ((InternalTracer) tracer).shutdownOtelMetrics(); + } + return unavailable(); + } + + private static CompletableFuture unavailable() { + return CompletableFuture.completedFuture(false); + } +} diff --git a/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java b/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java new file mode 100644 index 00000000000..507b2880d03 --- /dev/null +++ b/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java @@ -0,0 +1,49 @@ +package datadog.trace.api.metrics; + +import static java.lang.reflect.Modifier.isPublic; +import static java.lang.reflect.Modifier.isStatic; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.lang.reflect.Method; +import java.util.concurrent.CompletableFuture; +import org.junit.jupiter.api.Test; + +class OpenTelemetryMetricsTest { + + @Test + void lifecycleIsUnavailableWithoutAnInstalledTracer() { + CompletableFuture forceFlush = OpenTelemetryMetrics.forceFlush(); + CompletableFuture shutdown = OpenTelemetryMetrics.shutdown(); + + assertTrue(forceFlush.isDone()); + assertFalse(forceFlush.join()); + assertTrue(shutdown.isDone()); + assertFalse(shutdown.join()); + } + + @Test + void unavailableResultCannotBeChangedForLaterCalls() { + CompletableFuture first = OpenTelemetryMetrics.forceFlush(); + + first.obtrudeValue(true); + + assertFalse(OpenTelemetryMetrics.forceFlush().join()); + } + + @Test + void exposesPublicStaticLifecycleMethods() throws Exception { + assertLifecycleMethod("forceFlush"); + assertLifecycleMethod("shutdown"); + } + + private static void assertLifecycleMethod(String name) throws Exception { + Method method = OpenTelemetryMetrics.class.getMethod(name); + + assertTrue(isPublic(method.getModifiers())); + assertTrue(isStatic(method.getModifiers())); + assertEquals(CompletableFuture.class, method.getReturnType()); + assertEquals(0, method.getParameterCount()); + } +} diff --git a/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java b/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java index da89c0d073d..4a1fba0844c 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java @@ -129,6 +129,7 @@ import java.util.Properties; import java.util.ServiceConfigurationError; import java.util.ServiceLoader; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeoutException; @@ -1544,6 +1545,22 @@ public void flushMetrics() { } } + @Override + public CompletableFuture forceFlushOtelMetrics() { + if (initialConfig.isMetricsOtlpExporterEnabled()) { + return OtlpMetricsService.INSTANCE.forceFlush(); + } + return CompletableFuture.completedFuture(false); + } + + @Override + public CompletableFuture shutdownOtelMetrics() { + if (initialConfig.isMetricsOtlpExporterEnabled()) { + return OtlpMetricsService.INSTANCE.shutdown(); + } + return CompletableFuture.completedFuture(false); + } + @Override public void flushLogs() { if (initialConfig.isLogsOtlpExporterEnabled()) { diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java index 4a7580aa061..632ceea9e14 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java @@ -9,7 +9,12 @@ import datadog.trace.common.writer.RemoteApi; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; -import datadog.trace.util.AgentTaskScheduler; +import datadog.trace.util.AgentThreadFactory; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.Executors; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; import org.slf4j.Logger; @@ -18,19 +23,20 @@ /** Periodic service to collect OpenTelemetry metrics and export them over OTLP. */ public final class OtlpMetricsService { private static final Logger LOGGER = LoggerFactory.getLogger(OtlpMetricsService.class); - public static final OtlpMetricsService INSTANCE = new OtlpMetricsService(Config.get()); - private final AgentTaskScheduler scheduler; + private final ScheduledExecutorService executor; private final OtlpMetricsCollector collector; private final OtlpSender sender; - private final int intervalMillis; + private final Object lifecycleLock = new Object(); - private AgentTaskScheduler.Scheduled scheduledTask = null; + private ScheduledFuture scheduledTask; + private CompletableFuture shutdownFuture; OtlpMetricsService(Config config) { - this.scheduler = new AgentTaskScheduler(OTLP_METRICS_EXPORTER); + this.executor = + Executors.newSingleThreadScheduledExecutor(new AgentThreadFactory(OTLP_METRICS_EXPORTER)); this.sender = OtlpMetricsSenderFactory.create(config); if (this.sender == null) { LOGGER.debug("Unsupported OTLP metrics protocol: {}", config.getOtlpMetricsProtocol()); @@ -41,10 +47,20 @@ public final class OtlpMetricsService { ? new OtlpMetricsJsonCollector(SystemTimeSource.INSTANCE) : new OtlpMetricsProtoCollector(SystemTimeSource.INSTANCE); } - this.intervalMillis = config.getMetricsOtelInterval(); } + OtlpMetricsService( + ScheduledExecutorService executor, + OtlpMetricsCollector collector, + OtlpSender sender, + int intervalMillis) { + this.executor = executor; + this.collector = collector; + this.sender = sender; + this.intervalMillis = intervalMillis; + } + OtlpSender getSender() { return sender; } @@ -69,32 +85,111 @@ public void start() { / Math.log(1 - 0.25)), 5_000); - scheduledTask = - scheduler.scheduleAtFixedRate( - this::export, initialMillis, intervalMillis, TimeUnit.MILLISECONDS); + synchronized (lifecycleLock) { + if (shutdownFuture == null && scheduledTask == null) { + scheduledTask = + executor.scheduleAtFixedRate( + this::export, initialMillis, intervalMillis, TimeUnit.MILLISECONDS); + } + } + } + + public CompletableFuture forceFlush() { + synchronized (lifecycleLock) { + if (sender == null || shutdownFuture != null) { + return CompletableFuture.completedFuture(false); + } + CompletableFuture result = new CompletableFuture<>(); + try { + executor.execute(() -> result.complete(export())); + } catch (RejectedExecutionException e) { + LOGGER.debug("OTLP metrics executor rejected force flush", e); + result.complete(false); + } + return result; + } } public void flush() { - if (sender != null) { - scheduler.execute(this::export); + forceFlush(); + } + + public CompletableFuture shutdown() { + synchronized (lifecycleLock) { + if (shutdownFuture != null) { + return shutdownResult(); + } + + shutdownFuture = new CompletableFuture<>(); + if (scheduledTask != null) { + scheduledTask.cancel(false); + } + if (sender == null) { + executor.shutdown(); + shutdownFuture.complete(false); + return shutdownResult(); + } + + try { + executor.execute(this::finishShutdown); + } catch (RejectedExecutionException e) { + LOGGER.debug("OTLP metrics executor rejected shutdown", e); + closeSender(); + executor.shutdown(); + shutdownFuture.complete(false); + } + return shutdownResult(); } } - public void shutdown() { - if (scheduledTask != null) { - scheduledTask.cancel(); + private CompletableFuture shutdownResult() { + return shutdownFuture.thenApply(result -> result); + } + + private void finishShutdown() { + boolean result = export(); + if (!closeSender()) { + result = false; } - if (sender != null) { + try { + executor.shutdown(); + } catch (Throwable e) { + LOGGER.debug("Failed to shut down OTLP metrics executor", e); + result = false; + } + shutdownFuture.complete(result); + } + + private boolean closeSender() { + try { sender.shutdown(); + return true; + } catch (Throwable e) { + LOGGER.debug("Failed to shut down OTLP metrics sender", e); + return false; } } - private void export() { - OtlpPayload payload = collector.collectMetrics(); - if (payload != OtlpPayload.EMPTY) { + private boolean export() { + boolean attempted = false; + try { + OtlpPayload payload = collector.collectMetrics(); + if (payload == OtlpPayload.EMPTY) { + return true; + } + OtlpTelemetry.getInstance().onMetricsExportAttempt(); + attempted = true; RemoteApi.Response response = sender.send(payload); - OtlpTelemetry.getInstance().onMetricsExportComplete(response.success()); + boolean success = response != null && response.success(); + OtlpTelemetry.getInstance().onMetricsExportComplete(success); + return success; + } catch (Throwable e) { + if (attempted) { + OtlpTelemetry.getInstance().onMetricsExportComplete(false); + } + LOGGER.debug("Failed to export OTLP metrics", e); + return false; } } } diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java index 70e719a3dbb..98e53a8b574 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java @@ -2,15 +2,48 @@ import static datadog.trace.api.config.OtlpConfig.OTLP_METRICS_ENDPOINT; import static datadog.trace.api.config.OtlpConfig.OTLP_METRICS_PROTOCOL; +import static datadog.trace.common.writer.RemoteApi.Response.failed; +import static datadog.trace.common.writer.RemoteApi.Response.success; +import static java.util.concurrent.TimeUnit.SECONDS; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNotSame; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; import datadog.trace.api.Config; +import datadog.trace.api.telemetry.OtlpTelemetry; import datadog.trace.core.otlp.common.OtlpHttpSender; +import datadog.trace.core.otlp.common.OtlpPayload; +import datadog.trace.core.otlp.common.OtlpSender; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; import java.util.Properties; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; class OtlpMetricsServiceTest { + private static final OtlpPayload PAYLOAD = + new OtlpPayload(ByteBuffer.wrap(new byte[] {1}), OtlpPayload.PROTOBUF_CONTENT_TYPE); + private final List executors = new ArrayList<>(); + + @AfterEach + void stopExecutors() { + executors.forEach(ScheduledExecutorService::shutdownNow); + } @Test void httpJsonProtocolUsesJsonCollectorAndConfiguredEndpoint() { @@ -23,5 +56,192 @@ void httpJsonProtocolUsesJsonCollectorAndConfiguredEndpoint() { assertInstanceOf(OtlpMetricsJsonCollector.class, service.getCollector()); OtlpHttpSender sender = assertInstanceOf(OtlpHttpSender.class, service.getSender()); assertEquals("http://localhost:4318/v1/metrics", sender.url().toString()); + service.shutdown().join(); + } + + @Test + void forceFlushCompletesWithTransportResult() { + TestService test = service(PAYLOAD); + when(test.sender.send(PAYLOAD)).thenReturn(success(200), failed(500)); + + assertTrue(test.service.forceFlush().join()); + assertFalse(test.service.forceFlush().join()); + } + + @Test + void emptyFlushSucceedsWithoutTransport() { + TestService test = service(OtlpPayload.EMPTY); + + assertTrue(test.service.forceFlush().join()); + + verify(test.sender, never()).send(PAYLOAD); + } + + @Test + void collectionAndTransportExceptionsCompleteFalse() { + drainMetricsTelemetry(); + TestService collectionFailure = service(PAYLOAD); + when(collectionFailure.collector.collectMetrics()).thenThrow(new IllegalStateException("boom")); + + assertFalse(collectionFailure.service.forceFlush().join()); + + TestService transportFailure = service(PAYLOAD); + when(transportFailure.sender.send(PAYLOAD)).thenThrow(new IllegalStateException("boom")); + + assertFalse(transportFailure.service.forceFlush().join()); + + Map metrics = drainMetricsTelemetry(); + assertEquals(1L, metrics.get("otel.metrics_export_attempts").value); + assertEquals(1L, metrics.get("otel.metrics_export_failures").value); + } + + @Test + void forceFlushDoesNotCompleteBeforeTransport() throws Exception { + TestService test = service(PAYLOAD); + CountDownLatch entered = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + when(test.sender.send(PAYLOAD)) + .thenAnswer( + ignored -> { + entered.countDown(); + assertTrue(release.await(5, SECONDS)); + return success(200); + }); + + CompletableFuture result = test.service.forceFlush(); + + assertTrue(entered.await(5, SECONDS)); + assertFalse(result.isDone()); + release.countDown(); + assertTrue(result.get(5, SECONDS)); + } + + @Test + void concurrentFlushesAreSerialized() throws Exception { + TestService test = service(PAYLOAD); + AtomicInteger active = new AtomicInteger(); + AtomicInteger maximum = new AtomicInteger(); + CountDownLatch firstEntered = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + when(test.collector.collectMetrics()) + .thenAnswer( + ignored -> { + int count = active.incrementAndGet(); + maximum.accumulateAndGet(count, Math::max); + firstEntered.countDown(); + assertTrue(release.await(5, SECONDS)); + active.decrementAndGet(); + return PAYLOAD; + }); + when(test.sender.send(PAYLOAD)).thenReturn(success(200)); + + CompletableFuture first = test.service.forceFlush(); + assertTrue(firstEntered.await(5, SECONDS)); + CompletableFuture second = test.service.forceFlush(); + release.countDown(); + + assertTrue(first.get(5, SECONDS)); + assertTrue(second.get(5, SECONDS)); + assertEquals(1, maximum.get()); + } + + @Test + void shutdownFinalExportsClosesResourcesAndIsIdempotent() throws Exception { + TestService test = service(PAYLOAD); + when(test.sender.send(PAYLOAD)).thenReturn(success(200)); + + CompletableFuture first = test.service.shutdown(); + CompletableFuture second = test.service.shutdown(); + + assertTrue(first.join()); + assertTrue(second.join()); + assertNotSame(first, second); + verify(test.collector).collectMetrics(); + verify(test.sender).send(PAYLOAD); + verify(test.sender).shutdown(); + assertTrue(test.executor.isShutdown()); + assertTrue(test.executor.awaitTermination(5, SECONDS)); + assertFalse(test.service.forceFlush().join()); + } + + @Test + void shutdownClosesResourcesWhenFinalExportFails() throws Exception { + TestService test = service(PAYLOAD); + when(test.sender.send(PAYLOAD)).thenReturn(failed(500)); + + assertFalse(test.service.shutdown().join()); + + verify(test.sender).shutdown(); + assertTrue(test.executor.isShutdown()); + assertTrue(test.executor.awaitTermination(5, SECONDS)); + } + + @Test + void concurrentShutdownWaitsForInflightFlushAndExportsOnce() throws Exception { + TestService test = service(PAYLOAD); + CountDownLatch entered = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + when(test.sender.send(PAYLOAD)) + .thenAnswer( + ignored -> { + entered.countDown(); + assertTrue(release.await(5, SECONDS)); + return success(200); + }); + + CompletableFuture flush = test.service.forceFlush(); + assertTrue(entered.await(5, SECONDS)); + CompletableFuture shutdown = test.service.shutdown(); + assertFalse(shutdown.isDone()); + CompletableFuture repeated = test.service.shutdown(); + assertNotSame(shutdown, repeated); + shutdown.complete(false); + release.countDown(); + + assertTrue(flush.get(5, SECONDS)); + assertFalse(shutdown.get(5, SECONDS)); + assertTrue(repeated.get(5, SECONDS)); + assertTrue(test.service.shutdown().get(5, SECONDS)); + verify(test.sender, times(2)).send(PAYLOAD); + verify(test.sender).shutdown(); + } + + private TestService service(OtlpPayload payload) { + ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); + executors.add(executor); + OtlpMetricsCollector collector = mock(OtlpMetricsCollector.class); + OtlpSender sender = mock(OtlpSender.class); + when(collector.collectMetrics()).thenReturn(payload); + return new TestService( + new OtlpMetricsService(executor, collector, sender, 10_000), executor, collector, sender); + } + + private static Map drainMetricsTelemetry() { + Map byName = new HashMap<>(); + OtlpTelemetry.getInstance().prepareMetrics(); + for (OtlpTelemetry.OtlpMetric metric : OtlpTelemetry.getInstance().drain()) { + if (metric.metricName.startsWith("otel.metrics_")) { + byName.put(metric.metricName, metric); + } + } + return byName; + } + + private static final class TestService { + private final OtlpMetricsService service; + private final ScheduledExecutorService executor; + private final OtlpMetricsCollector collector; + private final OtlpSender sender; + + private TestService( + OtlpMetricsService service, + ScheduledExecutorService executor, + OtlpMetricsCollector collector, + OtlpSender sender) { + this.service = service; + this.executor = executor; + this.collector = collector; + this.sender = sender; + } } } From f585a26a93f8acd8951fb205cb7b289678dc5a3a Mon Sep 17 00:00:00 2001 From: Brian Marks Date: Thu, 27 Aug 2026 09:19:28 -0400 Subject: [PATCH 2/7] Test OpenTelemetry metrics lifecycle delegation --- .../api/metrics/OpenTelemetryMetricsTest.java | 44 +++++++++++++++++++ 1 file changed, 44 insertions(+) diff --git a/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java b/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java index 507b2880d03..0aba353ae88 100644 --- a/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java +++ b/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java @@ -4,13 +4,29 @@ import static java.lang.reflect.Modifier.isStatic; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Answers.CALLS_REAL_METHODS; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; +import static org.mockito.Mockito.withSettings; +import datadog.trace.api.GlobalTracer; +import datadog.trace.api.Tracer; +import datadog.trace.api.internal.InternalTracer; +import java.lang.reflect.Field; import java.lang.reflect.Method; import java.util.concurrent.CompletableFuture; +import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; class OpenTelemetryMetricsTest { + private final Tracer originalTracer = GlobalTracer.get(); + + @AfterEach + void restoreGlobalTracer() throws Exception { + setGlobalTracer(originalTracer); + } @Test void lifecycleIsUnavailableWithoutAnInstalledTracer() { @@ -32,6 +48,28 @@ void unavailableResultCannotBeChangedForLaterCalls() { assertFalse(OpenTelemetryMetrics.forceFlush().join()); } + @Test + void delegatesLifecycleToInstalledInternalTracer() throws Exception { + Tracer tracer = mock(Tracer.class, withSettings().extraInterfaces(InternalTracer.class)); + InternalTracer internalTracer = (InternalTracer) tracer; + CompletableFuture forceFlush = CompletableFuture.completedFuture(true); + CompletableFuture shutdown = CompletableFuture.completedFuture(false); + when(internalTracer.forceFlushOtelMetrics()).thenReturn(forceFlush); + when(internalTracer.shutdownOtelMetrics()).thenReturn(shutdown); + setGlobalTracer(tracer); + + assertSame(forceFlush, OpenTelemetryMetrics.forceFlush()); + assertSame(shutdown, OpenTelemetryMetrics.shutdown()); + } + + @Test + void internalTracerDefaultsReportLifecycleUnavailable() { + InternalTracer tracer = mock(InternalTracer.class, CALLS_REAL_METHODS); + + assertFalse(tracer.forceFlushOtelMetrics().join()); + assertFalse(tracer.shutdownOtelMetrics().join()); + } + @Test void exposesPublicStaticLifecycleMethods() throws Exception { assertLifecycleMethod("forceFlush"); @@ -46,4 +84,10 @@ private static void assertLifecycleMethod(String name) throws Exception { assertEquals(CompletableFuture.class, method.getReturnType()); assertEquals(0, method.getParameterCount()); } + + private static void setGlobalTracer(Tracer tracer) throws Exception { + Field provider = GlobalTracer.class.getDeclaredField("provider"); + provider.setAccessible(true); + provider.set(null, tracer); + } } From 31a8c31bade3aeab3e6a328a9ce9534764877b9b Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Tue, 1 Sep 2026 16:06:22 -0400 Subject: [PATCH 3/7] Clarify OpenTelemetry metrics lifecycle behavior --- .../api/metrics/OpenTelemetryMetrics.java | 14 ++++++++ .../otlp/metrics/OtlpMetricsServiceTest.java | 36 +++++++++++++++++++ 2 files changed, 50 insertions(+) diff --git a/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java b/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java index 12b52342191..6dead25274f 100644 --- a/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java @@ -5,9 +5,18 @@ import datadog.trace.api.internal.InternalTracer; import java.util.concurrent.CompletableFuture; +/** + * Controls Datadog's OTLP metrics export pipeline; shutdown does not disable the OpenTelemetry + * meter provider, and metrics recorded afterward are not exported. + */ public final class OpenTelemetryMetrics { private OpenTelemetryMetrics() {} + /** + * Exports pending metrics. The result is {@code true} after a successful or empty export and + * {@code false} if export is unavailable, fails, or shutdown has begun. The future has no + * deadline; a timed wait bounds only the caller and does not cancel export. + */ public static CompletableFuture forceFlush() { Tracer tracer = GlobalTracer.get(); if (tracer instanceof InternalTracer) { @@ -16,6 +25,11 @@ public static CompletableFuture forceFlush() { return unavailable(); } + /** + * Performs a final export and stops Datadog's metrics export pipeline. Repeated calls observe the + * first result; {@code false} means the pipeline was unavailable or export or cleanup failed. The + * future has no deadline; a timed wait bounds only the caller and does not cancel shutdown. + */ public static CompletableFuture shutdown() { Tracer tracer = GlobalTracer.get(); if (tracer instanceof InternalTracer) { diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java index 98e53a8b574..90bb3256d6d 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java @@ -10,6 +10,7 @@ import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; @@ -176,6 +177,41 @@ void shutdownClosesResourcesWhenFinalExportFails() throws Exception { assertTrue(test.executor.awaitTermination(5, SECONDS)); } + @Test + void shutdownReportsSenderCloseFailure() throws Exception { + TestService test = service(PAYLOAD); + when(test.sender.send(PAYLOAD)).thenReturn(success(200)); + doThrow(new IllegalStateException("boom")).when(test.sender).shutdown(); + + assertFalse(test.service.shutdown().join()); + + verify(test.sender).shutdown(); + assertTrue(test.executor.awaitTermination(5, SECONDS)); + } + + @Test + void rejectedLifecycleOperationsFailAndCloseSender() { + TestService test = service(PAYLOAD); + test.executor.shutdown(); + + assertFalse(test.service.forceFlush().join()); + assertFalse(test.service.shutdown().join()); + + verify(test.collector, never()).collectMetrics(); + verify(test.sender).shutdown(); + } + + @Test + void unavailablePipelineFailsLifecycleAndStopsExecutor() throws Exception { + ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); + executors.add(executor); + OtlpMetricsService service = new OtlpMetricsService(executor, null, null, 10_000); + + assertFalse(service.forceFlush().join()); + assertFalse(service.shutdown().join()); + assertTrue(executor.awaitTermination(5, SECONDS)); + } + @Test void concurrentShutdownWaitsForInflightFlushAndExportsOnce() throws Exception { TestService test = service(PAYLOAD); From cd8ab82b31aa54537815038a76e0261e454d09ca Mon Sep 17 00:00:00 2001 From: Brian Marks Date: Fri, 4 Sep 2026 12:05:45 -0400 Subject: [PATCH 4/7] Expose OpenTelemetry metrics shutdown through MeterProvider --- .../shim/metrics/OtelMeterProvider.java | 10 +- ...enTelemetryMetricsLifecycleForkedTest.java | 45 +++++ .../trace/api/internal/InternalTracer.java | 9 - .../api/metrics/CompletableResultCode.java | 121 ++++++++++++++ .../api/metrics/DatadogMeterProvider.java | 19 +++ .../api/metrics/OpenTelemetryMetrics.java | 44 ----- .../metrics/CompletableResultCodeTest.java | 118 +++++++++++++ .../api/metrics/DatadogMeterProviderTest.java | 25 +++ .../api/metrics/OpenTelemetryMetricsTest.java | 93 ----------- .../java/datadog/trace/core/CoreTracer.java | 25 ++- .../core/otlp/metrics/OtlpMetricsService.java | 109 +++++++----- .../otlp/metrics/OtlpMetricsServiceTest.java | 156 +++++++++++++----- .../java/datadog/opentracing/DDTracer.java | 11 -- .../datadog/opentracing/DDTracerTest.java | 17 -- .../instrumentation/api/AgentTracer.java | 10 ++ .../instrumentation/api/AgentTracerTest.java | 13 ++ 16 files changed, 556 insertions(+), 269 deletions(-) create mode 100644 dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java create mode 100644 dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java create mode 100644 dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java delete mode 100644 dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java create mode 100644 dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java create mode 100644 dd-trace-api/src/test/java/datadog/trace/api/metrics/DatadogMeterProviderTest.java delete mode 100644 dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java create mode 100644 internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentTracerTest.java diff --git a/dd-java-agent/agent-otel/otel-shim/src/main/java/datadog/opentelemetry/shim/metrics/OtelMeterProvider.java b/dd-java-agent/agent-otel/otel-shim/src/main/java/datadog/opentelemetry/shim/metrics/OtelMeterProvider.java index 245a5e19a1f..48e94589c7c 100644 --- a/dd-java-agent/agent-otel/otel-shim/src/main/java/datadog/opentelemetry/shim/metrics/OtelMeterProvider.java +++ b/dd-java-agent/agent-otel/otel-shim/src/main/java/datadog/opentelemetry/shim/metrics/OtelMeterProvider.java @@ -1,5 +1,8 @@ package datadog.opentelemetry.shim.metrics; +import datadog.trace.api.metrics.CompletableResultCode; +import datadog.trace.api.metrics.DatadogMeterProvider; +import datadog.trace.bootstrap.instrumentation.api.AgentTracer; import datadog.trace.bootstrap.otel.common.OtelInstrumentationScope; import datadog.trace.bootstrap.otel.metrics.data.OtelMetricStorage; import datadog.trace.util.Strings; @@ -15,7 +18,7 @@ import org.slf4j.LoggerFactory; @ParametersAreNonnullByDefault -public final class OtelMeterProvider implements MeterProvider { +public final class OtelMeterProvider implements MeterProvider, DatadogMeterProvider { private static final Logger LOGGER = LoggerFactory.getLogger(OtelMeterProvider.class); private static final String DEFAULT_METER_NAME = "unknown"; @@ -43,6 +46,11 @@ public MeterBuilder meterBuilder(String instrumentationScopeName) { return new OtelMeterBuilder(this, instrumentationScopeName); } + @Override + public CompletableResultCode shutdown() { + return AgentTracer.get().shutdownOtelMetrics(); + } + OtelMeter getMeterShim( String instrumentationScopeName, @Nullable String instrumentationScopeVersion, diff --git a/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java new file mode 100644 index 00000000000..eca54453c82 --- /dev/null +++ b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java @@ -0,0 +1,45 @@ +package opentelemetry147.metrics; + +import static datadog.trace.api.metrics.CompletableResultCode.ofSuccess; +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNotSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.trace.agent.test.AbstractInstrumentationTest; +import datadog.trace.api.GlobalTracer; +import datadog.trace.api.Tracer; +import datadog.trace.api.metrics.CompletableResultCode; +import datadog.trace.api.metrics.DatadogMeterProvider; +import datadog.trace.test.junit.utils.config.WithConfig; +import io.opentelemetry.api.GlobalOpenTelemetry; +import org.junit.jupiter.api.Test; + +@WithConfig(key = "metrics.otel.enabled", value = "true") +class OpenTelemetryMetricsLifecycleForkedTest extends AbstractInstrumentationTest { + + @Test + void globalMeterProviderExposesDatadogShutdown() { + DatadogMeterProvider meterProvider = + assertInstanceOf(DatadogMeterProvider.class, GlobalOpenTelemetry.get().getMeterProvider()); + Tracer originalTracer = GlobalTracer.get(); + Tracer replacementTracer = + (Tracer) + java.lang.reflect.Proxy.newProxyInstance( + Tracer.class.getClassLoader(), + new Class[] {Tracer.class}, + (proxy, method, arguments) -> + method.getReturnType() == boolean.class ? false : null); + + CompletableResultCode result; + try { + GlobalTracer.forceRegister(replacementTracer); + result = meterProvider.shutdown().join(10, SECONDS); + } finally { + GlobalTracer.forceRegister(originalTracer); + } + + assertNotSame(ofSuccess(), result); + assertTrue(result.isDone()); + } +} diff --git a/dd-trace-api/src/main/java/datadog/trace/api/internal/InternalTracer.java b/dd-trace-api/src/main/java/datadog/trace/api/internal/InternalTracer.java index 2cfd4a264b1..52b1adfc97e 100644 --- a/dd-trace-api/src/main/java/datadog/trace/api/internal/InternalTracer.java +++ b/dd-trace-api/src/main/java/datadog/trace/api/internal/InternalTracer.java @@ -2,7 +2,6 @@ import datadog.trace.api.experimental.DataStreamsCheckpointer; import datadog.trace.api.profiling.Profiling; -import java.util.concurrent.CompletableFuture; /** * Tracer internal features. Those features are not part of public API and can change or be removed @@ -21,14 +20,6 @@ public interface InternalTracer { void flushMetrics(); - default CompletableFuture forceFlushOtelMetrics() { - return CompletableFuture.completedFuture(false); - } - - default CompletableFuture shutdownOtelMetrics() { - return CompletableFuture.completedFuture(false); - } - void flushLogs(); Profiling getProfilingContext(); diff --git a/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java b/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java new file mode 100644 index 00000000000..05c87172b81 --- /dev/null +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java @@ -0,0 +1,121 @@ +package datadog.trace.api.metrics; + +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + +public final class CompletableResultCode { + private static final CompletableResultCode SUCCESS = new CompletableResultCode(true); + private static final CompletableResultCode FAILURE = new CompletableResultCode(false); + + private Boolean success; + private List callbacks; + + public CompletableResultCode() {} + + private CompletableResultCode(boolean success) { + this.success = success; + } + + public static CompletableResultCode ofSuccess() { + return SUCCESS; + } + + public static CompletableResultCode ofFailure() { + return FAILURE; + } + + public CompletableResultCode succeed() { + return complete(true); + } + + public CompletableResultCode fail() { + return complete(false); + } + + public synchronized boolean isSuccess() { + return Boolean.TRUE.equals(success); + } + + public synchronized boolean isDone() { + return success != null; + } + + /** + * Registers an action to run on the completing thread. If this result is already complete, the + * action runs immediately on the calling thread. + * + * @param callback action to run after completion + * @return this result + * @throws NullPointerException if {@code callback} is {@code null} + */ + public CompletableResultCode whenComplete(Runnable callback) { + Objects.requireNonNull(callback, "callback"); + synchronized (this) { + if (success == null) { + if (callbacks == null) { + callbacks = new ArrayList<>(); + } + callbacks.add(callback); + return this; + } + } + callback.run(); + return this; + } + + /** + * Waits up to the timeout for completion and returns this result. A timeout does not complete or + * cancel the operation; use {@link #isDone()} and {@link #isSuccess()} to inspect the outcome. + * + * @param timeout maximum time to wait + * @param unit unit of the timeout + * @return this result, which may still be incomplete after the timeout + */ + public CompletableResultCode join(long timeout, TimeUnit unit) { + if (isDone()) { + return this; + } + CountDownLatch completed = new CountDownLatch(1); + whenComplete(completed::countDown); + try { + completed.await(timeout, unit); + } catch (InterruptedException ignored) { + Thread.currentThread().interrupt(); + } + return this; + } + + private CompletableResultCode complete(boolean succeeded) { + List completionCallbacks; + synchronized (this) { + if (success != null) { + return this; + } + success = succeeded; + completionCallbacks = callbacks; + callbacks = null; + } + Throwable firstFailure = null; + if (completionCallbacks != null) { + for (Runnable callback : completionCallbacks) { + try { + callback.run(); + } catch (RuntimeException | Error failure) { + if (firstFailure == null) { + firstFailure = failure; + } + } + } + } + if (firstFailure instanceof RuntimeException) { + throw (RuntimeException) firstFailure; + } + if (firstFailure != null) { + throw (Error) firstFailure; + } + return this; + } +} diff --git a/dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java b/dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java new file mode 100644 index 00000000000..ad5d9f9a4ce --- /dev/null +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java @@ -0,0 +1,19 @@ +package datadog.trace.api.metrics; + +/** + * Datadog lifecycle controls implemented by the {@code MeterProvider} returned from {@code + * GlobalOpenTelemetry} when Datadog OpenTelemetry metrics support is enabled. + */ +public interface DatadogMeterProvider { + + /** + * Performs a final export and stops Datadog's OpenTelemetry metrics export pipeline. Metrics + * recorded after shutdown are not exported. Repeated calls observe the first result. + * + *

The operation has no deadline. A timed join bounds only the caller and does not cancel + * shutdown. + * + * @return the shutdown result; an unavailable or disabled pipeline succeeds as a no-op + */ + CompletableResultCode shutdown(); +} diff --git a/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java b/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java deleted file mode 100644 index 6dead25274f..00000000000 --- a/dd-trace-api/src/main/java/datadog/trace/api/metrics/OpenTelemetryMetrics.java +++ /dev/null @@ -1,44 +0,0 @@ -package datadog.trace.api.metrics; - -import datadog.trace.api.GlobalTracer; -import datadog.trace.api.Tracer; -import datadog.trace.api.internal.InternalTracer; -import java.util.concurrent.CompletableFuture; - -/** - * Controls Datadog's OTLP metrics export pipeline; shutdown does not disable the OpenTelemetry - * meter provider, and metrics recorded afterward are not exported. - */ -public final class OpenTelemetryMetrics { - private OpenTelemetryMetrics() {} - - /** - * Exports pending metrics. The result is {@code true} after a successful or empty export and - * {@code false} if export is unavailable, fails, or shutdown has begun. The future has no - * deadline; a timed wait bounds only the caller and does not cancel export. - */ - public static CompletableFuture forceFlush() { - Tracer tracer = GlobalTracer.get(); - if (tracer instanceof InternalTracer) { - return ((InternalTracer) tracer).forceFlushOtelMetrics(); - } - return unavailable(); - } - - /** - * Performs a final export and stops Datadog's metrics export pipeline. Repeated calls observe the - * first result; {@code false} means the pipeline was unavailable or export or cleanup failed. The - * future has no deadline; a timed wait bounds only the caller and does not cancel shutdown. - */ - public static CompletableFuture shutdown() { - Tracer tracer = GlobalTracer.get(); - if (tracer instanceof InternalTracer) { - return ((InternalTracer) tracer).shutdownOtelMetrics(); - } - return unavailable(); - } - - private static CompletableFuture unavailable() { - return CompletableFuture.completedFuture(false); - } -} diff --git a/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java b/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java new file mode 100644 index 00000000000..a671faba63b --- /dev/null +++ b/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java @@ -0,0 +1,118 @@ +package datadog.trace.api.metrics; + +import static datadog.trace.api.metrics.CompletableResultCode.ofFailure; +import static datadog.trace.api.metrics.CompletableResultCode.ofSuccess; +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.Test; + +class CompletableResultCodeTest { + + @Test + void completedResultsExposeTheirStatus() { + CompletableResultCode success = ofSuccess(); + CompletableResultCode failure = ofFailure(); + + assertTrue(success.isDone()); + assertTrue(success.isSuccess()); + assertSame(success, success.join(0, MILLISECONDS)); + assertTrue(failure.isDone()); + assertFalse(failure.isSuccess()); + assertSame(failure, failure.join(0, MILLISECONDS)); + } + + @Test + void firstCompletionWinsAndRunsCallbacksOnce() { + CompletableResultCode result = new CompletableResultCode(); + AtomicInteger callbacks = new AtomicInteger(); + + assertSame(result, result.whenComplete(callbacks::incrementAndGet)); + assertSame(result, result.whenComplete(callbacks::incrementAndGet)); + assertSame(result, result.succeed()); + assertSame(result, result.fail()); + + assertTrue(result.isDone()); + assertTrue(result.isSuccess()); + assertSame(result, result.whenComplete(callbacks::incrementAndGet)); + assertTrue(result.isSuccess()); + assertSame(result, result.join(0, MILLISECONDS)); + assertEquals(3, callbacks.get()); + } + + @Test + void canCompleteWithFailureWithoutCallbacks() { + CompletableResultCode result = new CompletableResultCode(); + + assertSame(result, result.fail()); + + assertTrue(result.isDone()); + assertFalse(result.isSuccess()); + } + + @Test + void callbackFailureDoesNotPreventRemainingCallbacks() { + CompletableResultCode result = new CompletableResultCode(); + AtomicInteger callbacks = new AtomicInteger(); + IllegalStateException failure = new IllegalStateException("boom"); + result.whenComplete( + () -> { + throw failure; + }); + result.whenComplete(callbacks::incrementAndGet); + + assertSame(failure, assertThrows(IllegalStateException.class, result::succeed)); + + assertTrue(result.isSuccess()); + assertEquals(1, callbacks.get()); + } + + @Test + void joinReturnsAfterCompletion() throws Exception { + CompletableResultCode result = new CompletableResultCode(); + CountDownLatch started = new CountDownLatch(1); + Thread completer = + new Thread( + () -> { + started.countDown(); + result.succeed(); + }); + completer.start(); + + assertTrue(started.await(5, SECONDS)); + assertSame(result, result.join(5, SECONDS)); + assertTrue(result.isSuccess()); + completer.join(); + } + + @Test + void joinReturnsIncompleteResultAfterTimeout() { + CompletableResultCode result = new CompletableResultCode(); + + assertSame(result, result.join(0, MILLISECONDS)); + + assertFalse(result.isDone()); + assertFalse(result.isSuccess()); + } + + @Test + void joinPreservesInterruptStatus() { + CompletableResultCode result = new CompletableResultCode(); + Thread.currentThread().interrupt(); + + try { + assertSame(result, result.join(1, SECONDS)); + assertTrue(Thread.currentThread().isInterrupted()); + assertFalse(result.isDone()); + } finally { + Thread.interrupted(); + } + } +} diff --git a/dd-trace-api/src/test/java/datadog/trace/api/metrics/DatadogMeterProviderTest.java b/dd-trace-api/src/test/java/datadog/trace/api/metrics/DatadogMeterProviderTest.java new file mode 100644 index 00000000000..a87572ff4b6 --- /dev/null +++ b/dd-trace-api/src/test/java/datadog/trace/api/metrics/DatadogMeterProviderTest.java @@ -0,0 +1,25 @@ +package datadog.trace.api.metrics; + +import static java.lang.reflect.Modifier.isAbstract; +import static java.lang.reflect.Modifier.isPublic; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.lang.reflect.Method; +import org.junit.jupiter.api.Test; + +class DatadogMeterProviderTest { + + @Test + void exposesOnlyShutdownLifecycleMethod() throws Exception { + Method shutdown = DatadogMeterProvider.class.getMethod("shutdown"); + + assertTrue(isPublic(shutdown.getModifiers())); + assertTrue(isAbstract(shutdown.getModifiers())); + assertEquals(CompletableResultCode.class, shutdown.getReturnType()); + assertEquals(0, shutdown.getParameterCount()); + assertThrows( + NoSuchMethodException.class, () -> DatadogMeterProvider.class.getMethod("forceFlush")); + } +} diff --git a/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java b/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java deleted file mode 100644 index 0aba353ae88..00000000000 --- a/dd-trace-api/src/test/java/datadog/trace/api/metrics/OpenTelemetryMetricsTest.java +++ /dev/null @@ -1,93 +0,0 @@ -package datadog.trace.api.metrics; - -import static java.lang.reflect.Modifier.isPublic; -import static java.lang.reflect.Modifier.isStatic; -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertSame; -import static org.junit.jupiter.api.Assertions.assertTrue; -import static org.mockito.Answers.CALLS_REAL_METHODS; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.when; -import static org.mockito.Mockito.withSettings; - -import datadog.trace.api.GlobalTracer; -import datadog.trace.api.Tracer; -import datadog.trace.api.internal.InternalTracer; -import java.lang.reflect.Field; -import java.lang.reflect.Method; -import java.util.concurrent.CompletableFuture; -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.Test; - -class OpenTelemetryMetricsTest { - private final Tracer originalTracer = GlobalTracer.get(); - - @AfterEach - void restoreGlobalTracer() throws Exception { - setGlobalTracer(originalTracer); - } - - @Test - void lifecycleIsUnavailableWithoutAnInstalledTracer() { - CompletableFuture forceFlush = OpenTelemetryMetrics.forceFlush(); - CompletableFuture shutdown = OpenTelemetryMetrics.shutdown(); - - assertTrue(forceFlush.isDone()); - assertFalse(forceFlush.join()); - assertTrue(shutdown.isDone()); - assertFalse(shutdown.join()); - } - - @Test - void unavailableResultCannotBeChangedForLaterCalls() { - CompletableFuture first = OpenTelemetryMetrics.forceFlush(); - - first.obtrudeValue(true); - - assertFalse(OpenTelemetryMetrics.forceFlush().join()); - } - - @Test - void delegatesLifecycleToInstalledInternalTracer() throws Exception { - Tracer tracer = mock(Tracer.class, withSettings().extraInterfaces(InternalTracer.class)); - InternalTracer internalTracer = (InternalTracer) tracer; - CompletableFuture forceFlush = CompletableFuture.completedFuture(true); - CompletableFuture shutdown = CompletableFuture.completedFuture(false); - when(internalTracer.forceFlushOtelMetrics()).thenReturn(forceFlush); - when(internalTracer.shutdownOtelMetrics()).thenReturn(shutdown); - setGlobalTracer(tracer); - - assertSame(forceFlush, OpenTelemetryMetrics.forceFlush()); - assertSame(shutdown, OpenTelemetryMetrics.shutdown()); - } - - @Test - void internalTracerDefaultsReportLifecycleUnavailable() { - InternalTracer tracer = mock(InternalTracer.class, CALLS_REAL_METHODS); - - assertFalse(tracer.forceFlushOtelMetrics().join()); - assertFalse(tracer.shutdownOtelMetrics().join()); - } - - @Test - void exposesPublicStaticLifecycleMethods() throws Exception { - assertLifecycleMethod("forceFlush"); - assertLifecycleMethod("shutdown"); - } - - private static void assertLifecycleMethod(String name) throws Exception { - Method method = OpenTelemetryMetrics.class.getMethod(name); - - assertTrue(isPublic(method.getModifiers())); - assertTrue(isStatic(method.getModifiers())); - assertEquals(CompletableFuture.class, method.getReturnType()); - assertEquals(0, method.getParameterCount()); - } - - private static void setGlobalTracer(Tracer tracer) throws Exception { - Field provider = GlobalTracer.class.getDeclaredField("provider"); - provider.setAccessible(true); - provider.set(null, tracer); - } -} diff --git a/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java b/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java index 237e8aa9066..d9378f2ba38 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java @@ -5,6 +5,7 @@ import static datadog.trace.api.DDTags.DSM_ENABLED; import static datadog.trace.api.DDTags.PROFILING_CONTEXT_ENGINE; import static datadog.trace.api.TracePropagationBehaviorExtract.IGNORE; +import static datadog.trace.api.metrics.CompletableResultCode.ofSuccess; import static datadog.trace.bootstrap.instrumentation.api.AgentPropagation.BAGGAGE_CONCERN; import static datadog.trace.bootstrap.instrumentation.api.AgentPropagation.DSM_CONCERN; import static datadog.trace.bootstrap.instrumentation.api.AgentPropagation.INFERRED_PROXY_CONCERN; @@ -58,6 +59,7 @@ import datadog.trace.api.interceptor.TraceInterceptor; import datadog.trace.api.internal.TraceSegment; import datadog.trace.api.internal.VisibleForTesting; +import datadog.trace.api.metrics.CompletableResultCode; import datadog.trace.api.metrics.SpanMetricRegistry; import datadog.trace.api.naming.SpanNaming; import datadog.trace.api.remoteconfig.ServiceNameCollector; @@ -129,7 +131,6 @@ import java.util.Properties; import java.util.ServiceConfigurationError; import java.util.ServiceLoader; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeoutException; @@ -1530,10 +1531,12 @@ public void close() { AgentMeter.statsDClient().close(); metricsAggregator.close(); if (initialConfig.isMetricsOtlpExporterEnabled()) { - try { - OtlpMetricsService.INSTANCE.shutdown().get(METRICS_FLUSH_TIMEOUT_MILLIS, MILLISECONDS); - } catch (InterruptedException | ExecutionException | TimeoutException e) { - log.debug("Failed to wait for OTLP metrics shutdown.", e); + CompletableResultCode result = + OtlpMetricsService.INSTANCE.shutdown().join(METRICS_FLUSH_TIMEOUT_MILLIS, MILLISECONDS); + if (!result.isDone()) { + log.debug("Timed out waiting for OTLP metrics shutdown."); + } else if (!result.isSuccess()) { + log.debug("OTLP metrics shutdown failed."); } } if (initialConfig.isLogsOtlpExporterEnabled()) { @@ -1581,19 +1584,11 @@ public void flushMetrics() { } @Override - public CompletableFuture forceFlushOtelMetrics() { - if (initialConfig.isMetricsOtlpExporterEnabled()) { - return OtlpMetricsService.INSTANCE.forceFlush(); - } - return CompletableFuture.completedFuture(false); - } - - @Override - public CompletableFuture shutdownOtelMetrics() { + public CompletableResultCode shutdownOtelMetrics() { if (initialConfig.isMetricsOtlpExporterEnabled()) { return OtlpMetricsService.INSTANCE.shutdown(); } - return CompletableFuture.completedFuture(false); + return ofSuccess(); } @Override diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java index 451d4e5e08d..09d0a57e2f1 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java @@ -6,13 +6,13 @@ import datadog.trace.api.Config; import datadog.trace.api.config.OtlpConfig; +import datadog.trace.api.metrics.CompletableResultCode; import datadog.trace.api.telemetry.OtlpTelemetry; import datadog.trace.api.time.SystemTimeSource; import datadog.trace.common.writer.RemoteApi; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; import datadog.trace.util.AgentThreadFactory; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; @@ -34,7 +34,7 @@ public final class OtlpMetricsService { private final Object lifecycleLock = new Object(); private ScheduledFuture scheduledTask; - private CompletableFuture shutdownFuture; + private CompletableResultCode shutdownResult; OtlpMetricsService(Config config) { this.executor = @@ -88,7 +88,7 @@ public void start() { 5_000); synchronized (lifecycleLock) { - if (shutdownFuture == null && scheduledTask == null) { + if (shutdownResult == null && scheduledTask == null) { scheduledTask = executor.scheduleAtFixedRate( this::export, initialMillis, intervalMillis, TimeUnit.MILLISECONDS); @@ -96,56 +96,60 @@ public void start() { } } - public CompletableFuture forceFlush() { + public void flush() { synchronized (lifecycleLock) { - if (sender == null || shutdownFuture != null) { - return CompletableFuture.completedFuture(false); + if (sender == null || shutdownResult != null) { + return; } - CompletableFuture result = new CompletableFuture<>(); try { - execute(() -> result.complete(export())); + execute(this::export); } catch (RejectedExecutionException e) { - LOGGER.debug("OTLP metrics executor rejected force flush", e); - result.complete(false); + LOGGER.debug("OTLP metrics executor rejected flush", e); } - return result; } } - public void flush() { - forceFlush(); - } - - public CompletableFuture shutdown() { + public CompletableResultCode shutdown() { synchronized (lifecycleLock) { - if (shutdownFuture != null) { - return shutdownResult(); + if (shutdownResult != null) { + return shutdownResultView(); } - shutdownFuture = new CompletableFuture<>(); - if (scheduledTask != null) { - scheduledTask.cancel(false); - } + shutdownResult = new CompletableResultCode(); + boolean cancellationSucceeded = cancelScheduledExport(); if (sender == null) { - executor.shutdown(); - shutdownFuture.complete(false); - return shutdownResult(); + boolean executorShutdown = shutdownExecutor(); + if (cancellationSucceeded && executorShutdown) { + shutdownResult.succeed(); + } else { + shutdownResult.fail(); + } + return shutdownResultView(); } try { - execute(this::finishShutdown); - } catch (RejectedExecutionException e) { - LOGGER.debug("OTLP metrics executor rejected shutdown", e); + execute(() -> finishShutdown(cancellationSucceeded)); + } catch (Throwable e) { + LOGGER.debug("Failed to submit OTLP metrics shutdown", e); closeSender(); - executor.shutdown(); - shutdownFuture.complete(false); + shutdownExecutor(); + shutdownResult.fail(); } - return shutdownResult(); + return shutdownResultView(); } } - private CompletableFuture shutdownResult() { - return shutdownFuture.thenApply(result -> result); + private CompletableResultCode shutdownResultView() { + CompletableResultCode result = new CompletableResultCode(); + shutdownResult.whenComplete( + () -> { + if (shutdownResult.isSuccess()) { + result.succeed(); + } else { + result.fail(); + } + }); + return result; } private void execute(Runnable task) { @@ -162,18 +166,35 @@ private void execute(Runnable task) { } } - private void finishShutdown() { + private boolean cancelScheduledExport() { + if (scheduledTask == null) { + return true; + } + try { + scheduledTask.cancel(false); + return true; + } catch (Throwable e) { + LOGGER.debug("Failed to cancel scheduled OTLP metrics export", e); + return false; + } + } + + private void finishShutdown(boolean cancellationSucceeded) { boolean result = export(); + if (!cancellationSucceeded) { + result = false; + } if (!closeSender()) { result = false; } - try { - executor.shutdown(); - } catch (Throwable e) { - LOGGER.debug("Failed to shut down OTLP metrics executor", e); + if (!shutdownExecutor()) { result = false; } - shutdownFuture.complete(result); + if (result) { + shutdownResult.succeed(); + } else { + shutdownResult.fail(); + } } private boolean closeSender() { @@ -186,6 +207,16 @@ private boolean closeSender() { } } + private boolean shutdownExecutor() { + try { + executor.shutdown(); + return true; + } catch (Throwable e) { + LOGGER.debug("Failed to shut down OTLP metrics executor", e); + return false; + } + } + private boolean export() { boolean attempted = false; try { diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java index e81ccf31709..5ee79519cdf 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java @@ -11,15 +11,20 @@ import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.inOrder; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; +import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import datadog.trace.api.Config; +import datadog.trace.api.metrics.CompletableResultCode; import datadog.trace.api.telemetry.OtlpTelemetry; import datadog.trace.bootstrap.instrumentation.api.AgentTracer; import datadog.trace.core.otlp.common.OtlpHttpSender; @@ -31,10 +36,11 @@ import java.util.List; import java.util.Map; import java.util.Properties; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Test; @@ -63,24 +69,26 @@ void httpJsonProtocolUsesJsonCollectorAndConfiguredEndpoint() { assertInstanceOf(OtlpMetricsJsonCollector.class, service.getCollector()); OtlpHttpSender sender = assertInstanceOf(OtlpHttpSender.class, service.getSender()); assertEquals("http://localhost:4318/v1/metrics", sender.url().toString()); - service.shutdown().join(); + assertTrue(service.shutdown().join(5, SECONDS).isSuccess()); } @Test - void forceFlushCompletesWithTransportResult() { + void flushExportsPendingMetrics() { TestService test = service(PAYLOAD); - when(test.sender.send(PAYLOAD)).thenReturn(success(200), failed(500)); + when(test.sender.send(PAYLOAD)).thenReturn(success(200)); + + test.service.flush(); - assertTrue(test.service.forceFlush().join()); - assertFalse(test.service.forceFlush().join()); + verify(test.sender, timeout(5_000)).send(PAYLOAD); } @Test - void emptyFlushSucceedsWithoutTransport() { + void emptyFlushSkipsTransport() { TestService test = service(OtlpPayload.EMPTY); - assertTrue(test.service.forceFlush().join()); + test.service.flush(); + verify(test.collector, timeout(5_000)).collectMetrics(); verify(test.sender, never()).send(PAYLOAD); } @@ -90,12 +98,12 @@ void collectionAndTransportExceptionsCompleteFalse() { TestService collectionFailure = service(PAYLOAD); when(collectionFailure.collector.collectMetrics()).thenThrow(new IllegalStateException("boom")); - assertFalse(collectionFailure.service.forceFlush().join()); + assertFalse(collectionFailure.service.shutdown().join(5, SECONDS).isSuccess()); TestService transportFailure = service(PAYLOAD); when(transportFailure.sender.send(PAYLOAD)).thenThrow(new IllegalStateException("boom")); - assertFalse(transportFailure.service.forceFlush().join()); + assertFalse(transportFailure.service.shutdown().join(5, SECONDS).isSuccess()); Map metrics = drainMetricsTelemetry(); assertEquals(1L, metrics.get("otel.metrics_export_attempts").value); @@ -103,7 +111,7 @@ void collectionAndTransportExceptionsCompleteFalse() { } @Test - void forceFlushDoesNotCompleteBeforeTransport() throws Exception { + void shutdownDoesNotCompleteBeforeTransport() throws Exception { TestService test = service(PAYLOAD); CountDownLatch entered = new CountDownLatch(1); CountDownLatch release = new CountDownLatch(1); @@ -115,12 +123,12 @@ void forceFlushDoesNotCompleteBeforeTransport() throws Exception { return success(200); }); - CompletableFuture result = test.service.forceFlush(); + CompletableResultCode result = test.service.shutdown(); assertTrue(entered.await(5, SECONDS)); assertFalse(result.isDone()); release.countDown(); - assertTrue(result.get(5, SECONDS)); + assertTrue(result.join(5, SECONDS).isSuccess()); } @Test @@ -142,14 +150,14 @@ void concurrentFlushesAreSerialized() throws Exception { }); when(test.sender.send(PAYLOAD)).thenReturn(success(200)); - CompletableFuture first = test.service.forceFlush(); + test.service.flush(); assertTrue(firstEntered.await(5, SECONDS)); - CompletableFuture second = test.service.forceFlush(); + test.service.flush(); release.countDown(); - assertTrue(first.get(5, SECONDS)); - assertTrue(second.get(5, SECONDS)); + assertTrue(test.service.shutdown().join(5, SECONDS).isSuccess()); assertEquals(1, maximum.get()); + verify(test.sender, times(3)).send(PAYLOAD); } @Test @@ -157,18 +165,19 @@ void shutdownFinalExportsClosesResourcesAndIsIdempotent() throws Exception { TestService test = service(PAYLOAD); when(test.sender.send(PAYLOAD)).thenReturn(success(200)); - CompletableFuture first = test.service.shutdown(); - CompletableFuture second = test.service.shutdown(); + CompletableResultCode first = test.service.shutdown(); + CompletableResultCode second = test.service.shutdown(); - assertTrue(first.join()); - assertTrue(second.join()); + assertTrue(first.join(5, SECONDS).isSuccess()); + assertTrue(second.join(5, SECONDS).isSuccess()); assertNotSame(first, second); verify(test.collector).collectMetrics(); verify(test.sender).send(PAYLOAD); verify(test.sender).shutdown(); assertTrue(test.executor.isShutdown()); assertTrue(test.executor.awaitTermination(5, SECONDS)); - assertFalse(test.service.forceFlush().join()); + test.service.flush(); + verify(test.collector).collectMetrics(); } @Test @@ -176,7 +185,7 @@ void shutdownClosesResourcesWhenFinalExportFails() throws Exception { TestService test = service(PAYLOAD); when(test.sender.send(PAYLOAD)).thenReturn(failed(500)); - assertFalse(test.service.shutdown().join()); + assertFalse(test.service.shutdown().join(5, SECONDS).isSuccess()); verify(test.sender).shutdown(); assertTrue(test.executor.isShutdown()); @@ -189,7 +198,7 @@ void shutdownReportsSenderCloseFailure() throws Exception { when(test.sender.send(PAYLOAD)).thenReturn(success(200)); doThrow(new IllegalStateException("boom")).when(test.sender).shutdown(); - assertFalse(test.service.shutdown().join()); + assertFalse(test.service.shutdown().join(5, SECONDS).isSuccess()); verify(test.sender).shutdown(); assertTrue(test.executor.awaitTermination(5, SECONDS)); @@ -200,13 +209,75 @@ void rejectedLifecycleOperationsFailAndCloseSender() { TestService test = service(PAYLOAD); test.executor.shutdown(); - assertFalse(test.service.forceFlush().join()); - assertFalse(test.service.shutdown().join()); + test.service.flush(); + assertFalse(test.service.shutdown().join(5, SECONDS).isSuccess()); verify(test.collector, never()).collectMetrics(); verify(test.sender).shutdown(); } + @Test + void executorFailuresCompleteShutdownResult() { + ScheduledExecutorService submissionFailure = mock(ScheduledExecutorService.class); + OtlpSender sender = mock(OtlpSender.class); + doThrow(new IllegalStateException("boom")).when(submissionFailure).execute(any(Runnable.class)); + OtlpMetricsService service = + new OtlpMetricsService(submissionFailure, mock(OtlpMetricsCollector.class), sender, 10_000); + + CompletableResultCode failedSubmission = service.shutdown(); + CompletableResultCode repeatedSubmission = service.shutdown(); + assertTrue(failedSubmission.isDone()); + assertFalse(failedSubmission.isSuccess()); + assertTrue(repeatedSubmission.isDone()); + assertFalse(repeatedSubmission.isSuccess()); + verify(sender).shutdown(); + verify(submissionFailure).shutdown(); + + ScheduledExecutorService cleanupFailure = mock(ScheduledExecutorService.class); + doThrow(new SecurityException("boom")).when(cleanupFailure).shutdown(); + OtlpMetricsService unavailable = new OtlpMetricsService(cleanupFailure, null, null, 10_000); + + CompletableResultCode failedCleanup = unavailable.shutdown(); + CompletableResultCode repeatedCleanup = unavailable.shutdown(); + assertTrue(failedCleanup.isDone()); + assertFalse(failedCleanup.isSuccess()); + assertTrue(repeatedCleanup.isDone()); + assertFalse(repeatedCleanup.isSuccess()); + } + + @Test + void scheduledExportCancellationFailureCompletesShutdownResult() { + ScheduledExecutorService executor = mock(ScheduledExecutorService.class); + ScheduledFuture scheduledExport = mock(ScheduledFuture.class); + OtlpMetricsCollector collector = mock(OtlpMetricsCollector.class); + OtlpSender sender = mock(OtlpSender.class); + doReturn(scheduledExport) + .when(executor) + .scheduleAtFixedRate(any(Runnable.class), anyLong(), anyLong(), any(TimeUnit.class)); + when(collector.collectMetrics()).thenReturn(OtlpPayload.EMPTY); + doAnswer( + invocation -> { + ((Runnable) invocation.getArgument(0)).run(); + return null; + }) + .when(executor) + .execute(any(Runnable.class)); + doThrow(new IllegalStateException("boom")).when(scheduledExport).cancel(false); + OtlpMetricsService service = new OtlpMetricsService(executor, collector, sender, 10_000); + service.start(); + + CompletableResultCode failedCancellation = service.shutdown(); + CompletableResultCode repeatedCancellation = service.shutdown(); + + assertTrue(failedCancellation.isDone()); + assertFalse(failedCancellation.isSuccess()); + assertTrue(repeatedCancellation.isDone()); + assertFalse(repeatedCancellation.isSuccess()); + verify(collector).collectMetrics(); + verify(sender).shutdown(); + verify(executor).shutdown(); + } + @Test void lifecycleSubmissionsDisableAsyncPropagation() { AgentTracer.TracerAPI tracer = mock(AgentTracer.TracerAPI.class); @@ -217,7 +288,7 @@ void lifecycleSubmissionsDisableAsyncPropagation() { new OtlpMetricsService( executor, mock(OtlpMetricsCollector.class), mock(OtlpSender.class), 10_000); - service.forceFlush(); + service.flush(); service.shutdown(); InOrder calls = inOrder(tracer, executor); @@ -230,18 +301,18 @@ void lifecycleSubmissionsDisableAsyncPropagation() { } @Test - void unavailablePipelineFailsLifecycleAndStopsExecutor() throws Exception { + void unavailablePipelineTreatsShutdownAsSuccessfulNoopAndStopsExecutor() throws Exception { ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); executors.add(executor); OtlpMetricsService service = new OtlpMetricsService(executor, null, null, 10_000); - assertFalse(service.forceFlush().join()); - assertFalse(service.shutdown().join()); + service.flush(); + assertTrue(service.shutdown().join(5, SECONDS).isSuccess()); assertTrue(executor.awaitTermination(5, SECONDS)); } @Test - void concurrentShutdownWaitsForInflightFlushAndExportsOnce() throws Exception { + void concurrentShutdownWaitsForInflightFlushAndCompletesAllViews() throws Exception { TestService test = service(PAYLOAD); CountDownLatch entered = new CountDownLatch(1); CountDownLatch release = new CountDownLatch(1); @@ -253,19 +324,24 @@ void concurrentShutdownWaitsForInflightFlushAndExportsOnce() throws Exception { return success(200); }); - CompletableFuture flush = test.service.forceFlush(); + test.service.flush(); assertTrue(entered.await(5, SECONDS)); - CompletableFuture shutdown = test.service.shutdown(); + CompletableResultCode shutdown = test.service.shutdown(); assertFalse(shutdown.isDone()); - CompletableFuture repeated = test.service.shutdown(); - assertNotSame(shutdown, repeated); - shutdown.complete(false); + CompletableResultCode throwing = test.service.shutdown(); + throwing.whenComplete( + () -> { + throw new IllegalStateException("boom"); + }); + CompletableResultCode unaffected = test.service.shutdown(); + assertNotSame(shutdown, throwing); + shutdown.fail(); release.countDown(); - assertTrue(flush.get(5, SECONDS)); - assertFalse(shutdown.get(5, SECONDS)); - assertTrue(repeated.get(5, SECONDS)); - assertTrue(test.service.shutdown().get(5, SECONDS)); + assertFalse(shutdown.join(5, SECONDS).isSuccess()); + assertTrue(throwing.join(5, SECONDS).isSuccess()); + assertTrue(unaffected.join(5, SECONDS).isSuccess()); + assertTrue(test.service.shutdown().join(5, SECONDS).isSuccess()); verify(test.sender, times(2)).send(PAYLOAD); verify(test.sender).shutdown(); } diff --git a/dd-trace-ot/src/main/java/datadog/opentracing/DDTracer.java b/dd-trace-ot/src/main/java/datadog/opentracing/DDTracer.java index e7c56a2c69a..9494c772aee 100644 --- a/dd-trace-ot/src/main/java/datadog/opentracing/DDTracer.java +++ b/dd-trace-ot/src/main/java/datadog/opentracing/DDTracer.java @@ -41,7 +41,6 @@ import java.util.Map; import java.util.Map.Entry; import java.util.Properties; -import java.util.concurrent.CompletableFuture; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -523,16 +522,6 @@ public void flushMetrics() { tracer.flushMetrics(); } - @Override - public CompletableFuture forceFlushOtelMetrics() { - return tracer.forceFlushOtelMetrics(); - } - - @Override - public CompletableFuture shutdownOtelMetrics() { - return tracer.shutdownOtelMetrics(); - } - @Override public void flushLogs() { tracer.flushLogs(); diff --git a/dd-trace-ot/src/test/java/datadog/opentracing/DDTracerTest.java b/dd-trace-ot/src/test/java/datadog/opentracing/DDTracerTest.java index dabe4f067fd..ff8704e305e 100644 --- a/dd-trace-ot/src/test/java/datadog/opentracing/DDTracerTest.java +++ b/dd-trace-ot/src/test/java/datadog/opentracing/DDTracerTest.java @@ -2,11 +2,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; -import static org.junit.jupiter.api.Assertions.assertSame; import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.when; -import datadog.trace.bootstrap.instrumentation.api.AgentTracer; import datadog.trace.common.sampling.Sampler; import datadog.trace.common.writer.DDAgentWriter; import datadog.trace.common.writer.ListWriter; @@ -15,7 +12,6 @@ import datadog.trace.test.util.DDJavaSpecification; import io.opentracing.Scope; import java.util.HashMap; -import java.util.concurrent.CompletableFuture; import org.junit.jupiter.api.Test; class DDTracerTest extends DDJavaSpecification { @@ -56,19 +52,6 @@ void testTracerBuilderWithDefaultWriter() throws Exception { tracer.close(); } - @Test - void delegatesOtelMetricsLifecycle() { - AgentTracer.TracerAPI delegate = mock(AgentTracer.TracerAPI.class); - CompletableFuture forceFlush = CompletableFuture.completedFuture(true); - CompletableFuture shutdown = CompletableFuture.completedFuture(false); - when(delegate.forceFlushOtelMetrics()).thenReturn(forceFlush); - when(delegate.shutdownOtelMetrics()).thenReturn(shutdown); - DDTracer tracer = new DDTracer(delegate); - - assertSame(forceFlush, tracer.forceFlushOtelMetrics()); - assertSame(shutdown, tracer.shutdownOtelMetrics()); - } - @Test void testAccessToTraceSegment() throws Exception { DDTracer tracer = DDTracer.builder().writer(DDAgentWriter.builder().build()).build(); diff --git a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentTracer.java b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentTracer.java index 406d6a015c2..38d3930001b 100644 --- a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentTracer.java +++ b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/AgentTracer.java @@ -1,5 +1,7 @@ package datadog.trace.bootstrap.instrumentation.api; +import static datadog.trace.api.metrics.CompletableResultCode.ofSuccess; + import datadog.context.Context; import datadog.context.ContextContinuation; import datadog.context.ContextListener; @@ -21,6 +23,7 @@ import datadog.trace.api.interceptor.TraceInterceptor; import datadog.trace.api.internal.InternalTracer; import datadog.trace.api.internal.TraceSegment; +import datadog.trace.api.metrics.CompletableResultCode; import datadog.trace.api.sampling.SamplingRule; import datadog.trace.api.scopemanager.ScopeListener; import datadog.trace.context.TraceScope; @@ -262,6 +265,8 @@ private AgentTracer() {} public interface TracerAPI extends datadog.trace.api.Tracer, InternalTracer, EndpointCheckpointer { + CompletableResultCode shutdownOtelMetrics(); + /** * Create and start a new span. * @@ -554,6 +559,11 @@ public void flush() {} @Override public void flushMetrics() {} + @Override + public CompletableResultCode shutdownOtelMetrics() { + return ofSuccess(); + } + @Override public void flushLogs() {} diff --git a/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentTracerTest.java b/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentTracerTest.java new file mode 100644 index 00000000000..71dfdafa492 --- /dev/null +++ b/internal-api/src/test/java/datadog/trace/bootstrap/instrumentation/api/AgentTracerTest.java @@ -0,0 +1,13 @@ +package datadog.trace.bootstrap.instrumentation.api; + +import static org.junit.jupiter.api.Assertions.assertTrue; + +import org.junit.jupiter.api.Test; + +class AgentTracerTest { + + @Test + void noopTracerTreatsMetricsShutdownAsSuccessful() { + assertTrue(new AgentTracer.NoopTracerAPI().shutdownOtelMetrics().isSuccess()); + } +} From 6b62b8c838a78fd7b52b88b6d251b471d83bd2f4 Mon Sep 17 00:00:00 2001 From: Brian Marks Date: Fri, 4 Sep 2026 13:13:55 -0400 Subject: [PATCH 5/7] Isolate OpenTelemetry shutdown result callbacks --- ...enTelemetryMetricsLifecycleForkedTest.java | 35 ++++-- .../api/metrics/CompletableResultCode.java | 119 ++++++++++++++---- .../api/metrics/DatadogMeterProvider.java | 7 +- .../metrics/CompletableResultCodeTest.java | 68 ++++++++++ .../core/otlp/metrics/OtlpMetricsService.java | 11 +- .../otlp/metrics/OtlpMetricsServiceTest.java | 38 ++++++ 6 files changed, 227 insertions(+), 51 deletions(-) diff --git a/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java index eca54453c82..d2f95ff355c 100644 --- a/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java +++ b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java @@ -1,18 +1,17 @@ package opentelemetry147.metrics; -import static datadog.trace.api.metrics.CompletableResultCode.ofSuccess; -import static java.util.concurrent.TimeUnit.SECONDS; import static org.junit.jupiter.api.Assertions.assertInstanceOf; -import static org.junit.jupiter.api.Assertions.assertNotSame; -import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.assertSame; import datadog.trace.agent.test.AbstractInstrumentationTest; import datadog.trace.api.GlobalTracer; import datadog.trace.api.Tracer; import datadog.trace.api.metrics.CompletableResultCode; import datadog.trace.api.metrics.DatadogMeterProvider; +import datadog.trace.bootstrap.instrumentation.api.AgentTracer; import datadog.trace.test.junit.utils.config.WithConfig; import io.opentelemetry.api.GlobalOpenTelemetry; +import java.lang.reflect.Proxy; import org.junit.jupiter.api.Test; @WithConfig(key = "metrics.otel.enabled", value = "true") @@ -22,24 +21,34 @@ class OpenTelemetryMetricsLifecycleForkedTest extends AbstractInstrumentationTes void globalMeterProviderExposesDatadogShutdown() { DatadogMeterProvider meterProvider = assertInstanceOf(DatadogMeterProvider.class, GlobalOpenTelemetry.get().getMeterProvider()); - Tracer originalTracer = GlobalTracer.get(); - Tracer replacementTracer = + Tracer originalGlobalTracer = GlobalTracer.get(); + AgentTracer.TracerAPI originalAgentTracer = AgentTracer.get(); + Tracer replacementGlobalTracer = (Tracer) - java.lang.reflect.Proxy.newProxyInstance( + Proxy.newProxyInstance( Tracer.class.getClassLoader(), new Class[] {Tracer.class}, (proxy, method, arguments) -> method.getReturnType() == boolean.class ? false : null); + Object expected = new CompletableResultCode(); + AgentTracer.TracerAPI replacementAgentTracer = + (AgentTracer.TracerAPI) + Proxy.newProxyInstance( + AgentTracer.TracerAPI.class.getClassLoader(), + new Class[] {AgentTracer.TracerAPI.class}, + (proxy, method, arguments) -> + method.getName().equals("shutdownOtelMetrics") ? expected : null); - CompletableResultCode result; + Object result; try { - GlobalTracer.forceRegister(replacementTracer); - result = meterProvider.shutdown().join(10, SECONDS); + GlobalTracer.forceRegister(replacementGlobalTracer); + AgentTracer.forceRegister(replacementAgentTracer); + result = meterProvider.shutdown(); } finally { - GlobalTracer.forceRegister(originalTracer); + AgentTracer.forceRegister(originalAgentTracer); + GlobalTracer.forceRegister(originalGlobalTracer); } - assertNotSame(ofSuccess(), result); - assertTrue(result.isDone()); + assertSame(expected, result); } } diff --git a/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java b/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java index 05c87172b81..478ff1be2e2 100644 --- a/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java @@ -1,22 +1,34 @@ package datadog.trace.api.metrics; +import static java.util.concurrent.TimeUnit.NANOSECONDS; + import java.util.ArrayList; import java.util.List; import java.util.Objects; -import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; public final class CompletableResultCode { private static final CompletableResultCode SUCCESS = new CompletableResultCode(true); private static final CompletableResultCode FAILURE = new CompletableResultCode(false); - private Boolean success; + private final SharedState sharedState; + private final boolean resultView; + + private Boolean resultViewSuccess; private List callbacks; - public CompletableResultCode() {} + public CompletableResultCode() { + this(new SharedState(), false); + } + + private CompletableResultCode(SharedState sharedState, boolean resultView) { + this.sharedState = sharedState; + this.resultView = resultView; + } private CompletableResultCode(boolean success) { - this.success = success; + this(); + sharedState.success = success; } public static CompletableResultCode ofSuccess() { @@ -27,6 +39,10 @@ public static CompletableResultCode ofFailure() { return FAILURE; } + public CompletableResultCode newResultView() { + return new CompletableResultCode(sharedState, true); + } + public CompletableResultCode succeed() { return complete(true); } @@ -35,12 +51,16 @@ public CompletableResultCode fail() { return complete(false); } - public synchronized boolean isSuccess() { - return Boolean.TRUE.equals(success); + public boolean isSuccess() { + synchronized (sharedState) { + return Boolean.TRUE.equals(outcome()); + } } - public synchronized boolean isDone() { - return success != null; + public boolean isDone() { + synchronized (sharedState) { + return outcome() != null; + } } /** @@ -53,10 +73,14 @@ public synchronized boolean isDone() { */ public CompletableResultCode whenComplete(Runnable callback) { Objects.requireNonNull(callback, "callback"); - synchronized (this) { - if (success == null) { + synchronized (sharedState) { + if (outcome() == null) { if (callbacks == null) { callbacks = new ArrayList<>(); + if (sharedState.callbackResults == null) { + sharedState.callbackResults = new ArrayList<>(); + } + sharedState.callbackResults.add(this); } callbacks.add(callback); return this; @@ -75,29 +99,45 @@ public CompletableResultCode whenComplete(Runnable callback) { * @return this result, which may still be incomplete after the timeout */ public CompletableResultCode join(long timeout, TimeUnit unit) { - if (isDone()) { - return this; - } - CountDownLatch completed = new CountDownLatch(1); - whenComplete(completed::countDown); - try { - completed.await(timeout, unit); - } catch (InterruptedException ignored) { - Thread.currentThread().interrupt(); + synchronized (sharedState) { + if (outcome() != null) { + return this; + } + + long remainingNanos = Objects.requireNonNull(unit, "unit").toNanos(timeout); + while (outcome() == null && remainingNanos > 0) { + long start = System.nanoTime(); + try { + NANOSECONDS.timedWait(sharedState, remainingNanos); + } catch (InterruptedException ignored) { + Thread.currentThread().interrupt(); + break; + } + remainingNanos -= Math.max(1, System.nanoTime() - start); + } } return this; } private CompletableResultCode complete(boolean succeeded) { List completionCallbacks; - synchronized (this) { - if (success != null) { + synchronized (sharedState) { + if (outcome() != null) { return this; } - success = succeeded; - completionCallbacks = callbacks; - callbacks = null; + + if (resultView) { + resultViewSuccess = succeeded; + completionCallbacks = callbacks; + callbacks = null; + removeCallbackResult(); + } else { + sharedState.success = succeeded; + completionCallbacks = collectCallbacks(); + } + sharedState.notifyAll(); } + Throwable firstFailure = null; if (completionCallbacks != null) { for (Runnable callback : completionCallbacks) { @@ -118,4 +158,35 @@ private CompletableResultCode complete(boolean succeeded) { } return this; } + + private Boolean outcome() { + return resultView && resultViewSuccess != null ? resultViewSuccess : sharedState.success; + } + + private List collectCallbacks() { + if (sharedState.callbackResults == null) { + return null; + } + List completionCallbacks = new ArrayList<>(); + for (CompletableResultCode result : sharedState.callbackResults) { + completionCallbacks.addAll(result.callbacks); + result.callbacks = null; + } + sharedState.callbackResults = null; + return completionCallbacks; + } + + private void removeCallbackResult() { + if (sharedState.callbackResults != null) { + sharedState.callbackResults.remove(this); + if (sharedState.callbackResults.isEmpty()) { + sharedState.callbackResults = null; + } + } + } + + private static final class SharedState { + private Boolean success; + private List callbackResults; + } } diff --git a/dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java b/dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java index ad5d9f9a4ce..339ab5545a2 100644 --- a/dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java @@ -7,11 +7,10 @@ public interface DatadogMeterProvider { /** - * Performs a final export and stops Datadog's OpenTelemetry metrics export pipeline. Metrics - * recorded after shutdown are not exported. Repeated calls observe the first result. + * Performs a final export and stops Datadog's OpenTelemetry metrics pipeline. Repeated calls + * observe the first result. * - *

The operation has no deadline. A timed join bounds only the caller and does not cancel - * shutdown. + *

A timed join bounds only the caller and does not cancel shutdown. * * @return the shutdown result; an unavailable or disabled pipeline succeeds as a no-op */ diff --git a/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java b/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java index a671faba63b..a44b6105d13 100644 --- a/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java +++ b/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java @@ -74,6 +74,74 @@ void callbackFailureDoesNotPreventRemainingCallbacks() { assertEquals(1, callbacks.get()); } + @Test + void resultViewFollowsSourceOutcome() { + CompletableResultCode source = new CompletableResultCode(); + CompletableResultCode success = source.newResultView(); + CompletableResultCode failure = ofFailure().newResultView(); + + source.succeed(); + + assertTrue(success.isSuccess()); + assertTrue(failure.isDone()); + assertFalse(failure.isSuccess()); + } + + @Test + void resultViewCanCompleteIndependently() { + CompletableResultCode source = new CompletableResultCode(); + CompletableResultCode result = source.newResultView(); + + result.fail(); + source.succeed(); + + assertTrue(source.isSuccess()); + assertFalse(result.isSuccess()); + } + + @Test + void resultViewChainsCompleteWithoutGrowingTheStack() { + CompletableResultCode source = new CompletableResultCode(); + CompletableResultCode result = source; + for (int i = 0; i < 100_000; i++) { + result = result.newResultView(); + } + CompletableResultCode sibling = source.newResultView(); + + source.succeed(); + + assertTrue(result.isSuccess()); + assertTrue(sibling.isSuccess()); + } + + @Test + void blockedCallbackDoesNotDelaySiblingResultView() throws Exception { + CompletableResultCode source = new CompletableResultCode(); + CompletableResultCode blocking = source.newResultView(); + CompletableResultCode unaffected = source.newResultView(); + CountDownLatch callbackEntered = new CountDownLatch(1); + CountDownLatch releaseCallback = new CountDownLatch(1); + blocking.whenComplete( + () -> { + callbackEntered.countDown(); + try { + assertTrue(releaseCallback.await(5, SECONDS)); + } catch (InterruptedException ignored) { + Thread.currentThread().interrupt(); + } + }); + Thread completer = new Thread(source::succeed); + + completer.start(); + assertTrue(callbackEntered.await(5, SECONDS)); + try { + assertTrue(unaffected.join(1, SECONDS).isSuccess()); + } finally { + releaseCallback.countDown(); + } + completer.join(); + } + @Test void joinReturnsAfterCompletion() throws Exception { CompletableResultCode result = new CompletableResultCode(); diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java index 09d0a57e2f1..b93610e19b2 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpMetricsService.java @@ -140,16 +140,7 @@ public CompletableResultCode shutdown() { } private CompletableResultCode shutdownResultView() { - CompletableResultCode result = new CompletableResultCode(); - shutdownResult.whenComplete( - () -> { - if (shutdownResult.isSuccess()) { - result.succeed(); - } else { - result.fail(); - } - }); - return result; + return shutdownResult.newResultView(); } private void execute(Runnable task) { diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java index 5ee79519cdf..8796acd3f37 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpMetricsServiceTest.java @@ -346,6 +346,44 @@ void concurrentShutdownWaitsForInflightFlushAndCompletesAllViews() throws Except verify(test.sender).shutdown(); } + @Test + void blockedCallbackOnOneShutdownViewDoesNotDelayAnotherView() throws Exception { + TestService test = service(PAYLOAD); + CountDownLatch exportEntered = new CountDownLatch(1); + CountDownLatch releaseExport = new CountDownLatch(1); + CountDownLatch callbackEntered = new CountDownLatch(1); + CountDownLatch releaseCallback = new CountDownLatch(1); + when(test.sender.send(PAYLOAD)) + .thenAnswer( + ignored -> { + exportEntered.countDown(); + assertTrue(releaseExport.await(5, SECONDS)); + return success(200); + }); + + CompletableResultCode blocking = test.service.shutdown(); + blocking.whenComplete( + () -> { + callbackEntered.countDown(); + try { + assertTrue(releaseCallback.await(5, SECONDS)); + } catch (InterruptedException ignored) { + Thread.currentThread().interrupt(); + } + }); + CompletableResultCode unaffected = test.service.shutdown(); + + assertTrue(exportEntered.await(5, SECONDS)); + releaseExport.countDown(); + assertTrue(callbackEntered.await(5, SECONDS)); + try { + assertTrue(unaffected.join(1, SECONDS).isSuccess()); + } finally { + releaseCallback.countDown(); + } + assertTrue(blocking.join(5, SECONDS).isSuccess()); + } + private TestService service(OtlpPayload payload) { ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); executors.add(executor); From 01e1241dd8085282040b523d91e89e17f6c18fb4 Mon Sep 17 00:00:00 2001 From: Brian Marks Date: Fri, 4 Sep 2026 13:35:32 -0400 Subject: [PATCH 6/7] Suppress false singleton warning --- .../java/datadog/trace/api/metrics/CompletableResultCode.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java b/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java index 478ff1be2e2..76d34683861 100644 --- a/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java @@ -2,11 +2,15 @@ import static java.util.concurrent.TimeUnit.NANOSECONDS; +import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import java.util.ArrayList; import java.util.List; import java.util.Objects; import java.util.concurrent.TimeUnit; +@SuppressFBWarnings( + value = "SING_SINGLETON_HAS_NONPRIVATE_CONSTRUCTOR", + justification = "Not a singleton") public final class CompletableResultCode { private static final CompletableResultCode SUCCESS = new CompletableResultCode(true); private static final CompletableResultCode FAILURE = new CompletableResultCode(false); From 2fd3467c9c8f90f17ce4b35e2f278f1a97d3017e Mon Sep 17 00:00:00 2001 From: Brian Marks Date: Fri, 4 Sep 2026 14:16:48 -0400 Subject: [PATCH 7/7] Improve OpenTelemetry shutdown test coverage --- ...enTelemetryMetricsLifecycleForkedTest.java | 12 ---- .../metrics/CompletableResultCodeTest.java | 66 +++++++++++++++---- 2 files changed, 54 insertions(+), 24 deletions(-) diff --git a/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java index d2f95ff355c..4357e8bffa1 100644 --- a/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java +++ b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java @@ -4,8 +4,6 @@ import static org.junit.jupiter.api.Assertions.assertSame; import datadog.trace.agent.test.AbstractInstrumentationTest; -import datadog.trace.api.GlobalTracer; -import datadog.trace.api.Tracer; import datadog.trace.api.metrics.CompletableResultCode; import datadog.trace.api.metrics.DatadogMeterProvider; import datadog.trace.bootstrap.instrumentation.api.AgentTracer; @@ -21,15 +19,7 @@ class OpenTelemetryMetricsLifecycleForkedTest extends AbstractInstrumentationTes void globalMeterProviderExposesDatadogShutdown() { DatadogMeterProvider meterProvider = assertInstanceOf(DatadogMeterProvider.class, GlobalOpenTelemetry.get().getMeterProvider()); - Tracer originalGlobalTracer = GlobalTracer.get(); AgentTracer.TracerAPI originalAgentTracer = AgentTracer.get(); - Tracer replacementGlobalTracer = - (Tracer) - Proxy.newProxyInstance( - Tracer.class.getClassLoader(), - new Class[] {Tracer.class}, - (proxy, method, arguments) -> - method.getReturnType() == boolean.class ? false : null); Object expected = new CompletableResultCode(); AgentTracer.TracerAPI replacementAgentTracer = (AgentTracer.TracerAPI) @@ -41,12 +31,10 @@ void globalMeterProviderExposesDatadogShutdown() { Object result; try { - GlobalTracer.forceRegister(replacementGlobalTracer); AgentTracer.forceRegister(replacementAgentTracer); result = meterProvider.shutdown(); } finally { AgentTracer.forceRegister(originalAgentTracer); - GlobalTracer.forceRegister(originalGlobalTracer); } assertSame(expected, result); diff --git a/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java b/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java index a44b6105d13..b9440df42c0 100644 --- a/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java +++ b/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java @@ -10,6 +10,7 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import datadog.trace.test.util.PollingConditions; import java.util.concurrent.CountDownLatch; import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.Test; @@ -74,6 +75,23 @@ void callbackFailureDoesNotPreventRemainingCallbacks() { assertEquals(1, callbacks.get()); } + @Test + void callbackErrorDoesNotPreventRemainingCallbacks() { + CompletableResultCode result = new CompletableResultCode(); + AtomicInteger callbacks = new AtomicInteger(); + AssertionError failure = new AssertionError("boom"); + result.whenComplete( + () -> { + throw failure; + }); + result.whenComplete(callbacks::incrementAndGet); + + assertSame(failure, assertThrows(AssertionError.class, result::succeed)); + + assertTrue(result.isSuccess()); + assertEquals(1, callbacks.get()); + } + @Test void resultViewFollowsSourceOutcome() { CompletableResultCode source = new CompletableResultCode(); @@ -99,6 +117,25 @@ void resultViewCanCompleteIndependently() { assertFalse(result.isSuccess()); } + @Test + void independentlyCompletedResultViewRunsItsCallbacksOnlyOnce() { + CompletableResultCode source = new CompletableResultCode(); + CompletableResultCode result = source.newResultView(); + CompletableResultCode sibling = source.newResultView(); + AtomicInteger resultCallbacks = new AtomicInteger(); + AtomicInteger siblingCallbacks = new AtomicInteger(); + result.whenComplete(resultCallbacks::incrementAndGet); + sibling.whenComplete(siblingCallbacks::incrementAndGet); + + result.fail(); + source.succeed(); + + assertEquals(1, resultCallbacks.get()); + assertEquals(1, siblingCallbacks.get()); + assertFalse(result.isSuccess()); + assertTrue(sibling.isSuccess()); + } + @Test void resultViewChainsCompleteWithoutGrowingTheStack() { CompletableResultCode source = new CompletableResultCode(); @@ -143,21 +180,26 @@ void blockedCallbackDoesNotDelaySiblingResultView() throws Exception { } @Test - void joinReturnsAfterCompletion() throws Exception { + void joinWaitsForCompletion() throws Exception { CompletableResultCode result = new CompletableResultCode(); - CountDownLatch started = new CountDownLatch(1); - Thread completer = - new Thread( - () -> { - started.countDown(); - result.succeed(); - }); - completer.start(); + Thread waiter = new Thread(() -> result.join(30, SECONDS)); + waiter.start(); + + try { + new PollingConditions(5) + .delay(0.01) + .eventually(() -> assertEquals(Thread.State.TIMED_WAITING, waiter.getState())); - assertTrue(started.await(5, SECONDS)); - assertSame(result, result.join(5, SECONDS)); + result.succeed(); + waiter.join(SECONDS.toMillis(5)); + } finally { + result.succeed(); + waiter.interrupt(); + waiter.join(SECONDS.toMillis(5)); + } + + assertFalse(waiter.isAlive()); assertTrue(result.isSuccess()); - completer.join(); } @Test