Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -20,6 +21,14 @@ public interface InternalTracer {

void flushMetrics();

default CompletableFuture<Boolean> forceFlushOtelMetrics() {
return CompletableFuture.completedFuture(false);
}

default CompletableFuture<Boolean> shutdownOtelMetrics() {
return CompletableFuture.completedFuture(false);
}

void flushLogs();

Profiling getProfilingContext();
Expand Down
Original file line number Diff line number Diff line change
@@ -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<Boolean> forceFlush() {
Tracer tracer = GlobalTracer.get();
if (tracer instanceof InternalTracer) {
return ((InternalTracer) tracer).forceFlushOtelMetrics();
}
return unavailable();
}

public static CompletableFuture<Boolean> shutdown() {
Tracer tracer = GlobalTracer.get();
if (tracer instanceof InternalTracer) {
return ((InternalTracer) tracer).shutdownOtelMetrics();
}
return unavailable();
}

private static CompletableFuture<Boolean> unavailable() {
return CompletableFuture.completedFuture(false);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
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<Boolean> forceFlush = OpenTelemetryMetrics.forceFlush();
CompletableFuture<Boolean> shutdown = OpenTelemetryMetrics.shutdown();

assertTrue(forceFlush.isDone());
assertFalse(forceFlush.join());
assertTrue(shutdown.isDone());
assertFalse(shutdown.join());
}

@Test
void unavailableResultCannotBeChangedForLaterCalls() {
CompletableFuture<Boolean> 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<Boolean> forceFlush = CompletableFuture.completedFuture(true);
CompletableFuture<Boolean> 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);
}
}
17 changes: 17 additions & 0 deletions dd-trace-core/src/main/java/datadog/trace/core/CoreTracer.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -1544,6 +1545,22 @@ public void flushMetrics() {
}
}

@Override
public CompletableFuture<Boolean> forceFlushOtelMetrics() {
if (initialConfig.isMetricsOtlpExporterEnabled()) {
return OtlpMetricsService.INSTANCE.forceFlush();
}
return CompletableFuture.completedFuture(false);
}

@Override
public CompletableFuture<Boolean> shutdownOtelMetrics() {
if (initialConfig.isMetricsOtlpExporterEnabled()) {
return OtlpMetricsService.INSTANCE.shutdown();
}
return CompletableFuture.completedFuture(false);
}

@Override
public void flushLogs() {
if (initialConfig.isLogsOtlpExporterEnabled()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<Boolean> 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());
Expand All @@ -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;
}
Expand All @@ -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<Boolean> forceFlush() {
synchronized (lifecycleLock) {
if (sender == null || shutdownFuture != null) {
return CompletableFuture.completedFuture(false);
}
CompletableFuture<Boolean> 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<Boolean> 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<Boolean> 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;
}
}
}
Loading
Loading