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..4357e8bffa1 --- /dev/null +++ b/dd-java-agent/instrumentation/opentelemetry/opentelemetry-1.47/src/test/java/opentelemetry147/metrics/OpenTelemetryMetricsLifecycleForkedTest.java @@ -0,0 +1,42 @@ +package opentelemetry147.metrics; + +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertSame; + +import datadog.trace.agent.test.AbstractInstrumentationTest; +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") +class OpenTelemetryMetricsLifecycleForkedTest extends AbstractInstrumentationTest { + + @Test + void globalMeterProviderExposesDatadogShutdown() { + DatadogMeterProvider meterProvider = + assertInstanceOf(DatadogMeterProvider.class, GlobalOpenTelemetry.get().getMeterProvider()); + AgentTracer.TracerAPI originalAgentTracer = AgentTracer.get(); + 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); + + Object result; + try { + AgentTracer.forceRegister(replacementAgentTracer); + result = meterProvider.shutdown(); + } finally { + AgentTracer.forceRegister(originalAgentTracer); + } + + 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 new file mode 100644 index 00000000000..76d34683861 --- /dev/null +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/CompletableResultCode.java @@ -0,0 +1,196 @@ +package datadog.trace.api.metrics; + +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); + + private final SharedState sharedState; + private final boolean resultView; + + private Boolean resultViewSuccess; + private List callbacks; + + public CompletableResultCode() { + this(new SharedState(), false); + } + + private CompletableResultCode(SharedState sharedState, boolean resultView) { + this.sharedState = sharedState; + this.resultView = resultView; + } + + private CompletableResultCode(boolean success) { + this(); + sharedState.success = success; + } + + public static CompletableResultCode ofSuccess() { + return SUCCESS; + } + + public static CompletableResultCode ofFailure() { + return FAILURE; + } + + public CompletableResultCode newResultView() { + return new CompletableResultCode(sharedState, true); + } + + public CompletableResultCode succeed() { + return complete(true); + } + + public CompletableResultCode fail() { + return complete(false); + } + + public boolean isSuccess() { + synchronized (sharedState) { + return Boolean.TRUE.equals(outcome()); + } + } + + public boolean isDone() { + synchronized (sharedState) { + return outcome() != 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 (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; + } + } + 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) { + 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 (sharedState) { + if (outcome() != null) { + return this; + } + + 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) { + 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; + } + + 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 new file mode 100644 index 00000000000..339ab5545a2 --- /dev/null +++ b/dd-trace-api/src/main/java/datadog/trace/api/metrics/DatadogMeterProvider.java @@ -0,0 +1,18 @@ +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 pipeline. Repeated calls + * observe the first result. + * + *

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/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..b9440df42c0 --- /dev/null +++ b/dd-trace-api/src/test/java/datadog/trace/api/metrics/CompletableResultCodeTest.java @@ -0,0 +1,228 @@ +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 datadog.trace.test.util.PollingConditions; +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 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(); + 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 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(); + 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 joinWaitsForCompletion() throws Exception { + CompletableResultCode result = new CompletableResultCode(); + 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())); + + result.succeed(); + waiter.join(SECONDS.toMillis(5)); + } finally { + result.succeed(); + waiter.interrupt(); + waiter.join(SECONDS.toMillis(5)); + } + + assertFalse(waiter.isAlive()); + assertTrue(result.isSuccess()); + } + + @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-core/src/main/java/datadog/trace/core/CoreTracer.java b/dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java index 6b2e41a9b1d..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; @@ -143,6 +145,7 @@ */ public class CoreTracer implements AgentTracer.TracerAPI, TracerFlare.Reporter { private static final Logger log = LoggerFactory.getLogger(CoreTracer.class); + private static final long METRICS_FLUSH_TIMEOUT_MILLIS = 2_500; public static CoreTracerBuilder builder() { return new CoreTracerBuilder(); @@ -1528,7 +1531,13 @@ public void close() { AgentMeter.statsDClient().close(); metricsAggregator.close(); if (initialConfig.isMetricsOtlpExporterEnabled()) { - OtlpMetricsService.INSTANCE.shutdown(); + 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()) { OtlpLogsService.INSTANCE.shutdown(); @@ -1564,7 +1573,7 @@ public void flush() { @Override public void flushMetrics() { try { - metricsAggregator.forceReport().get(2_500, MILLISECONDS); + metricsAggregator.forceReport().get(METRICS_FLUSH_TIMEOUT_MILLIS, MILLISECONDS); } catch (InterruptedException | ExecutionException | TimeoutException e) { log.debug("Failed to wait for metrics flush.", e); } @@ -1574,6 +1583,14 @@ public void flushMetrics() { } } + @Override + public CompletableResultCode shutdownOtelMetrics() { + if (initialConfig.isMetricsOtlpExporterEnabled()) { + return OtlpMetricsService.INSTANCE.shutdown(); + } + return ofSuccess(); + } + @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..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 @@ -1,15 +1,22 @@ package datadog.trace.core.otlp.metrics; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.isAsyncPropagationEnabled; +import static datadog.trace.bootstrap.instrumentation.api.AgentTracer.setAsyncPropagationEnabled; import static datadog.trace.util.AgentThreadFactory.AgentThread.OTLP_METRICS_EXPORTER; 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.AgentTaskScheduler; +import datadog.trace.util.AgentThreadFactory; +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 +25,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 CompletableResultCode shutdownResult; 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 +49,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 +87,147 @@ public void start() { / Math.log(1 - 0.25)), 5_000); - scheduledTask = - scheduler.scheduleAtFixedRate( - this::export, initialMillis, intervalMillis, TimeUnit.MILLISECONDS); + synchronized (lifecycleLock) { + if (shutdownResult == null && scheduledTask == null) { + scheduledTask = + executor.scheduleAtFixedRate( + this::export, initialMillis, intervalMillis, TimeUnit.MILLISECONDS); + } + } } public void flush() { - if (sender != null) { - scheduler.execute(this::export); + synchronized (lifecycleLock) { + if (sender == null || shutdownResult != null) { + return; + } + try { + execute(this::export); + } catch (RejectedExecutionException e) { + LOGGER.debug("OTLP metrics executor rejected flush", e); + } + } + } + + public CompletableResultCode shutdown() { + synchronized (lifecycleLock) { + if (shutdownResult != null) { + return shutdownResultView(); + } + + shutdownResult = new CompletableResultCode(); + boolean cancellationSucceeded = cancelScheduledExport(); + if (sender == null) { + boolean executorShutdown = shutdownExecutor(); + if (cancellationSucceeded && executorShutdown) { + shutdownResult.succeed(); + } else { + shutdownResult.fail(); + } + return shutdownResultView(); + } + + try { + execute(() -> finishShutdown(cancellationSucceeded)); + } catch (Throwable e) { + LOGGER.debug("Failed to submit OTLP metrics shutdown", e); + closeSender(); + shutdownExecutor(); + shutdownResult.fail(); + } + return shutdownResultView(); + } + } + + private CompletableResultCode shutdownResultView() { + return shutdownResult.newResultView(); + } + + private void execute(Runnable task) { + boolean restorePropagation = isAsyncPropagationEnabled(); + if (restorePropagation) { + setAsyncPropagationEnabled(false); + } + try { + executor.execute(task); + } finally { + if (restorePropagation) { + setAsyncPropagationEnabled(true); + } } } - public void shutdown() { - if (scheduledTask != null) { - scheduledTask.cancel(); + private boolean cancelScheduledExport() { + if (scheduledTask == null) { + return true; } - if (sender != null) { + 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; + } + if (!shutdownExecutor()) { + result = false; + } + if (result) { + shutdownResult.succeed(); + } else { + shutdownResult.fail(); + } + } + + 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 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 { + 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..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 @@ -2,15 +2,61 @@ 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.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; +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.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; +import org.mockito.InOrder; class OtlpMetricsServiceTest { + private static final OtlpPayload PAYLOAD = + new OtlpPayload(ByteBuffer.wrap(new byte[] {1}), OtlpPayload.PROTOBUF_CONTENT_TYPE); + private final List executors = new ArrayList<>(); + private final AgentTracer.TracerAPI originalTracer = AgentTracer.get(); + + @AfterEach + void stopExecutors() { + executors.forEach(ScheduledExecutorService::shutdownNow); + AgentTracer.forceRegister(originalTracer); + } @Test void httpJsonProtocolUsesJsonCollectorAndConfiguredEndpoint() { @@ -23,5 +69,357 @@ void httpJsonProtocolUsesJsonCollectorAndConfiguredEndpoint() { assertInstanceOf(OtlpMetricsJsonCollector.class, service.getCollector()); OtlpHttpSender sender = assertInstanceOf(OtlpHttpSender.class, service.getSender()); assertEquals("http://localhost:4318/v1/metrics", sender.url().toString()); + assertTrue(service.shutdown().join(5, SECONDS).isSuccess()); + } + + @Test + void flushExportsPendingMetrics() { + TestService test = service(PAYLOAD); + when(test.sender.send(PAYLOAD)).thenReturn(success(200)); + + test.service.flush(); + + verify(test.sender, timeout(5_000)).send(PAYLOAD); + } + + @Test + void emptyFlushSkipsTransport() { + TestService test = service(OtlpPayload.EMPTY); + + test.service.flush(); + + verify(test.collector, timeout(5_000)).collectMetrics(); + 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.shutdown().join(5, SECONDS).isSuccess()); + + TestService transportFailure = service(PAYLOAD); + when(transportFailure.sender.send(PAYLOAD)).thenThrow(new IllegalStateException("boom")); + + assertFalse(transportFailure.service.shutdown().join(5, SECONDS).isSuccess()); + + Map metrics = drainMetricsTelemetry(); + assertEquals(1L, metrics.get("otel.metrics_export_attempts").value); + assertEquals(1L, metrics.get("otel.metrics_export_failures").value); + } + + @Test + void shutdownDoesNotCompleteBeforeTransport() 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); + }); + + CompletableResultCode result = test.service.shutdown(); + + assertTrue(entered.await(5, SECONDS)); + assertFalse(result.isDone()); + release.countDown(); + assertTrue(result.join(5, SECONDS).isSuccess()); + } + + @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)); + + test.service.flush(); + assertTrue(firstEntered.await(5, SECONDS)); + test.service.flush(); + release.countDown(); + + assertTrue(test.service.shutdown().join(5, SECONDS).isSuccess()); + assertEquals(1, maximum.get()); + verify(test.sender, times(3)).send(PAYLOAD); + } + + @Test + void shutdownFinalExportsClosesResourcesAndIsIdempotent() throws Exception { + TestService test = service(PAYLOAD); + when(test.sender.send(PAYLOAD)).thenReturn(success(200)); + + CompletableResultCode first = test.service.shutdown(); + CompletableResultCode second = test.service.shutdown(); + + 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)); + test.service.flush(); + verify(test.collector).collectMetrics(); + } + + @Test + void shutdownClosesResourcesWhenFinalExportFails() throws Exception { + TestService test = service(PAYLOAD); + when(test.sender.send(PAYLOAD)).thenReturn(failed(500)); + + assertFalse(test.service.shutdown().join(5, SECONDS).isSuccess()); + + verify(test.sender).shutdown(); + assertTrue(test.executor.isShutdown()); + 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(5, SECONDS).isSuccess()); + + verify(test.sender).shutdown(); + assertTrue(test.executor.awaitTermination(5, SECONDS)); + } + + @Test + void rejectedLifecycleOperationsFailAndCloseSender() { + TestService test = service(PAYLOAD); + test.executor.shutdown(); + + 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); + when(tracer.isAsyncPropagationEnabled()).thenReturn(true); + AgentTracer.forceRegister(tracer); + ScheduledExecutorService executor = mock(ScheduledExecutorService.class); + OtlpMetricsService service = + new OtlpMetricsService( + executor, mock(OtlpMetricsCollector.class), mock(OtlpSender.class), 10_000); + + service.flush(); + service.shutdown(); + + InOrder calls = inOrder(tracer, executor); + for (int i = 0; i < 2; i++) { + calls.verify(tracer).isAsyncPropagationEnabled(); + calls.verify(tracer).setAsyncPropagationEnabled(false); + calls.verify(executor).execute(any(Runnable.class)); + calls.verify(tracer).setAsyncPropagationEnabled(true); + } + } + + @Test + void unavailablePipelineTreatsShutdownAsSuccessfulNoopAndStopsExecutor() throws Exception { + ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); + executors.add(executor); + OtlpMetricsService service = new OtlpMetricsService(executor, null, null, 10_000); + + service.flush(); + assertTrue(service.shutdown().join(5, SECONDS).isSuccess()); + assertTrue(executor.awaitTermination(5, SECONDS)); + } + + @Test + void concurrentShutdownWaitsForInflightFlushAndCompletesAllViews() 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); + }); + + test.service.flush(); + assertTrue(entered.await(5, SECONDS)); + CompletableResultCode shutdown = test.service.shutdown(); + assertFalse(shutdown.isDone()); + CompletableResultCode throwing = test.service.shutdown(); + throwing.whenComplete( + () -> { + throw new IllegalStateException("boom"); + }); + CompletableResultCode unaffected = test.service.shutdown(); + assertNotSame(shutdown, throwing); + shutdown.fail(); + release.countDown(); + + 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(); + } + + @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); + 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; + } } } 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()); + } +}