From 3fd40fcc952752b7233a6586859c110523728400 Mon Sep 17 00:00:00 2001 From: Brian Marks Date: Wed, 26 Aug 2026 20:23:31 -0400 Subject: [PATCH 1/2] 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/2] 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); + } }