diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java index b2817b46dbc..900de33f117 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpPayloadDispatcher.java @@ -1,5 +1,6 @@ package datadog.trace.common.writer; +import datadog.trace.api.telemetry.OtlpTelemetry; import datadog.trace.core.CoreSpan; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; @@ -26,13 +27,15 @@ public void addTrace(List> trace) { public void flush() { OtlpPayload payload = collector.collectTraces(); if (payload != OtlpPayload.EMPTY) { - sender.send(payload); + OtlpTelemetry.getInstance().onTracesExportAttempt(); + RemoteApi.Response response = sender.send(payload); + OtlpTelemetry.getInstance().onTracesExportComplete(response.success()); } } @Override public void onDroppedTrace(int spanCount) { - // TODO: surface drop counts via HealthMetrics + // no telemetry currently tracked for dropped traces } @Override diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpWriter.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpWriter.java index 8118ff7b2bb..1650a77b52a 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpWriter.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/OtlpWriter.java @@ -38,11 +38,10 @@ public static OtlpWriterBuilder builder() { TraceProcessingWorker worker, PayloadDispatcher dispatcher, OtlpSender sender, - HealthMetrics healthMetrics, int flushTimeout, TimeUnit flushTimeoutUnit, boolean alwaysFlush) { - super(worker, dispatcher, healthMetrics, flushTimeout, flushTimeoutUnit, alwaysFlush); + super(worker, dispatcher, HealthMetrics.NO_OP, flushTimeout, flushTimeoutUnit, alwaysFlush); this.sender = sender; } @@ -64,7 +63,6 @@ public static class OtlpWriterBuilder { private OtlpConfig.Protocol protocol = OtlpConfig.Protocol.HTTP_PROTOBUF; private OtlpConfig.Compression compression = OtlpConfig.Compression.NONE; private int traceBufferSize = BUFFER_SIZE; - private HealthMetrics healthMetrics = HealthMetrics.NO_OP; private int flushIntervalMilliseconds = 1000; private int flushTimeout = 1; private TimeUnit flushTimeoutUnit = TimeUnit.SECONDS; @@ -102,11 +100,6 @@ public OtlpWriterBuilder traceBufferSize(int traceBufferSize) { return this; } - public OtlpWriterBuilder healthMetrics(HealthMetrics healthMetrics) { - this.healthMetrics = healthMetrics; - return this; - } - public OtlpWriterBuilder flushIntervalMilliseconds(int flushIntervalMilliseconds) { this.flushIntervalMilliseconds = flushIntervalMilliseconds; return this; @@ -151,7 +144,7 @@ public OtlpWriter build() { final TraceProcessingWorker worker = new TraceProcessingWorker( traceBufferSize, - healthMetrics, + HealthMetrics.NO_OP, dispatcher, DroppingPolicy.DISABLED, Prioritization.FAST_LANE, @@ -160,7 +153,7 @@ public OtlpWriter build() { singleSpanSampler); return new OtlpWriter( - worker, dispatcher, sender, healthMetrics, flushTimeout, flushTimeoutUnit, alwaysFlush); + worker, dispatcher, sender, flushTimeout, flushTimeoutUnit, alwaysFlush); } } } diff --git a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java index 23b15c39544..73875a0e408 100644 --- a/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java +++ b/dd-trace-core/src/main/java/datadog/trace/common/writer/WriterFactory.java @@ -76,7 +76,6 @@ public static Writer createWriter( .protocol(config.getOtlpTracesProtocol()) .compression(config.getOtlpTracesCompression()) .timeoutMillis(config.getOtlpTracesTimeout()) - .healthMetrics(healthMetrics) .spanSamplingRules(singleSpanSampler) .flushIntervalMilliseconds(flushIntervalMilliseconds) .build(); diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpGrpcSender.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpGrpcSender.java index 5969569ab4f..73c73a7f626 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpGrpcSender.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpGrpcSender.java @@ -2,18 +2,16 @@ import static datadog.communication.http.OkHttpUtils.buildHttp2Client; import static datadog.communication.http.OkHttpUtils.isPlainHttp; -import static datadog.communication.http.OkHttpUtils.sendWithRetries; import datadog.communication.http.HttpRetryPolicy; import datadog.logging.RatelimitedLogger; import datadog.trace.api.config.OtlpConfig.Compression; -import java.io.IOException; +import datadog.trace.common.writer.RemoteApi; import java.util.Map; import java.util.concurrent.TimeUnit; import okhttp3.HttpUrl; import okhttp3.OkHttpClient; import okhttp3.Request; -import okhttp3.Response; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -55,20 +53,8 @@ public OtlpGrpcSender( } @Override - public void send(OtlpPayload payload) { - Request request = makeRequest(payload); - try (Response response = sendWithRetries(client, retryPolicy, request)) { - if (!response.isSuccessful()) { - RATELIMITED_LOGGER.warn( - "OTLP export to {} failed with status {}: {}", - request.url(), - response.code(), - response.message()); - } - } catch (IOException e) { - RATELIMITED_LOGGER.warn( - "OTLP export to {} failed with exception: {}", request.url(), e.toString()); - } + public RemoteApi.Response send(OtlpPayload payload) { + return OtlpSenderSupport.send(client, retryPolicy, makeRequest(payload), RATELIMITED_LOGGER); } @Override diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpHttpSender.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpHttpSender.java index 9f6bd3b05c3..8cac29e85d5 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpHttpSender.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpHttpSender.java @@ -2,18 +2,16 @@ import static datadog.communication.http.OkHttpUtils.buildHttpClient; import static datadog.communication.http.OkHttpUtils.isPlainHttp; -import static datadog.communication.http.OkHttpUtils.sendWithRetries; import datadog.communication.http.HttpRetryPolicy; import datadog.logging.RatelimitedLogger; import datadog.trace.api.config.OtlpConfig.Compression; -import java.io.IOException; +import datadog.trace.common.writer.RemoteApi; import java.util.Map; import java.util.concurrent.TimeUnit; import okhttp3.HttpUrl; import okhttp3.OkHttpClient; import okhttp3.Request; -import okhttp3.Response; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -59,20 +57,8 @@ public HttpUrl url() { } @Override - public void send(OtlpPayload payload) { - Request request = makeRequest(payload); - try (Response response = sendWithRetries(client, retryPolicy, request)) { - if (!response.isSuccessful()) { - RATELIMITED_LOGGER.warn( - "OTLP export to {} failed with status {}: {}", - request.url(), - response.code(), - response.message()); - } - } catch (IOException e) { - RATELIMITED_LOGGER.warn( - "OTLP export to {} failed with exception: {}", request.url(), e.toString()); - } + public RemoteApi.Response send(OtlpPayload payload) { + return OtlpSenderSupport.send(client, retryPolicy, makeRequest(payload), RATELIMITED_LOGGER); } @Override diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpSender.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpSender.java index 55cf57d053e..81172279b85 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpSender.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpSender.java @@ -1,8 +1,10 @@ package datadog.trace.core.otlp.common; +import datadog.trace.common.writer.RemoteApi; + /** Sends chunks of OTLP data. */ public interface OtlpSender { - void send(OtlpPayload payload); + RemoteApi.Response send(OtlpPayload payload); void shutdown(); } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpSenderSupport.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpSenderSupport.java new file mode 100644 index 00000000000..3e27cc9e020 --- /dev/null +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpSenderSupport.java @@ -0,0 +1,39 @@ +package datadog.trace.core.otlp.common; + +import static datadog.communication.http.OkHttpUtils.sendWithRetries; +import static datadog.trace.common.writer.RemoteApi.Response.failed; +import static datadog.trace.common.writer.RemoteApi.Response.success; + +import datadog.communication.http.HttpRetryPolicy; +import datadog.logging.RatelimitedLogger; +import datadog.trace.common.writer.RemoteApi; +import java.io.IOException; +import okhttp3.OkHttpClient; + +/** Shared request execution and response handling for {@link OtlpSender} implementations. */ +final class OtlpSenderSupport { + private OtlpSenderSupport() {} + + /** Executes the given request with retries, logging failures via the rate-limited logger. */ + static RemoteApi.Response send( + OkHttpClient client, + HttpRetryPolicy.Factory retryPolicy, + okhttp3.Request request, + RatelimitedLogger ratelimitedLogger) { + try (okhttp3.Response response = sendWithRetries(client, retryPolicy, request)) { + if (response.isSuccessful()) { + return success(response.code()); + } + ratelimitedLogger.warn( + "OTLP export to {} failed with status {}: {}", + request.url(), + response.code(), + response.message()); + return failed(response.code()); + } catch (IOException e) { + ratelimitedLogger.warn( + "OTLP export to {} failed with exception: {}", request.url(), e.toString()); + return failed(e); + } + } +} diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsCollector.java index c263efab665..0e699e22623 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsCollector.java @@ -7,4 +7,7 @@ public abstract class OtlpLogsCollector { /** Waits for logs to be batched within the given interval. */ public abstract OtlpPayload waitForLogs(int intervalMillis); + + /** Number of log records collected. */ + public abstract int getLogRecordCount(); } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsJsonCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsJsonCollector.java index 00136bbd80e..b81822b6705 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsJsonCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsJsonCollector.java @@ -36,6 +36,7 @@ public final class OtlpLogsJsonCollector extends OtlpLogsCollector private JsonWriter writer; private boolean anyLogRecordWritten; private boolean logRecordStarted; + private int logRecordCount; private final LazyJsonArray attributesArray = new LazyJsonArray(); @@ -65,6 +66,8 @@ OtlpPayload collectLogs(ObjIntConsumer processor, int intervalM /** Prepare temporary elements to collect logs data. */ private void start() { + logRecordCount = 0; + writer = new JsonWriter(); writer.beginObject(); writer.name("resourceLogs").beginArray(); @@ -86,6 +89,11 @@ private void stop() { currentScope = null; } + @Override + public int getLogRecordCount() { + return logRecordCount; + } + @Override public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) { if (currentScope != null) { @@ -119,6 +127,7 @@ public void visitLogRecord(OtlpLogRecord logRecord) { logRecordStarted = false; anyLogRecordWritten = true; + logRecordCount++; } // opens the log record object on first attribute or value written for it diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsProtoCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsProtoCollector.java index bdd8436f58f..94fb0c6d1fa 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsProtoCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsProtoCollector.java @@ -41,6 +41,7 @@ public final class OtlpLogsProtoCollector extends OtlpLogsCollector // total number of chunked bytes at different nesting levels private int payloadBytes; private int scopedBytes; + private int logRecordCount; private OtelInstrumentationScope currentScope; @@ -68,6 +69,7 @@ OtlpPayload collectLogs(ObjIntConsumer processor, int intervalM /** Prepare temporary elements to collect logs data. */ private void start() { + logRecordCount = 0; // remove stale entries from caches OtlpCommonProto.recalibrateCaches(); @@ -84,6 +86,11 @@ private void stop() { currentScope = null; } + @Override + public int getLogRecordCount() { + return logRecordCount; + } + @Override public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) { if (currentScope != null) { @@ -96,6 +103,7 @@ public OtlpScopedLogsVisitor visitScopedLogs(OtelInstrumentationScope scope) { @Override public void visitLogRecord(OtlpLogRecord logRecord) { scopedBytes += recordLogRecordMessage(buf, logRecord, protobuf); + logRecordCount++; } @Override diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsService.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsService.java index 195f77efdf7..e65a3757782 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsService.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/logs/OtlpLogsService.java @@ -4,6 +4,8 @@ import static datadog.trace.util.AgentThreadFactory.newAgentThread; import datadog.trace.api.Config; +import datadog.trace.api.telemetry.OtlpTelemetry; +import datadog.trace.common.writer.RemoteApi; import datadog.trace.core.otlp.common.OtlpGrpcSender; import datadog.trace.core.otlp.common.OtlpHttpSender; import datadog.trace.core.otlp.common.OtlpPayload; @@ -108,7 +110,11 @@ private void export() { try { OtlpPayload payload = collector.waitForLogs(intervalMillis); if (payload != OtlpPayload.EMPTY) { - sender.send(payload); + int logRecordCount = collector.getLogRecordCount(); + RemoteApi.Response response = sender.send(payload); + if (response.success()) { + OtlpTelemetry.getInstance().onLogRecordsSubmitted(logRecordCount); + } } } catch (RuntimeException e) { LOGGER.debug("Uncaught exception exporting logs", e); 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 98a65618d45..4a7580aa061 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 @@ -4,7 +4,9 @@ import datadog.trace.api.Config; import datadog.trace.api.config.OtlpConfig; +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; @@ -29,7 +31,6 @@ public final class OtlpMetricsService { OtlpMetricsService(Config config) { this.scheduler = new AgentTaskScheduler(OTLP_METRICS_EXPORTER); - this.sender = OtlpMetricsSenderFactory.create(config); if (this.sender == null) { LOGGER.debug("Unsupported OTLP metrics protocol: {}", config.getOtlpMetricsProtocol()); @@ -91,7 +92,9 @@ public void shutdown() { private void export() { OtlpPayload payload = collector.collectMetrics(); if (payload != OtlpPayload.EMPTY) { - sender.send(payload); + OtlpTelemetry.getInstance().onMetricsExportAttempt(); + RemoteApi.Response response = sender.send(payload); + OtlpTelemetry.getInstance().onMetricsExportComplete(response.success()); } } } diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpStatsMetricWriter.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpStatsMetricWriter.java index 5d926b527e5..a5a21302877 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpStatsMetricWriter.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/metrics/OtlpStatsMetricWriter.java @@ -7,6 +7,7 @@ import datadog.metrics.api.Histogram; import datadog.trace.api.Config; import datadog.trace.api.config.OtlpConfig; +import datadog.trace.api.telemetry.OtlpTelemetry; import datadog.trace.api.time.SystemTimeSource; import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString; import datadog.trace.bootstrap.otel.common.OtelInstrumentationScope; @@ -16,6 +17,7 @@ import datadog.trace.bootstrap.otlp.metrics.OtlpMetricsVisitor; import datadog.trace.common.metrics.AggregateEntry; import datadog.trace.common.metrics.MetricWriter; +import datadog.trace.common.writer.RemoteApi; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpResourceJson; import datadog.trace.core.otlp.common.OtlpResourceProto; @@ -167,7 +169,9 @@ public void finishBucket() { } OtlpPayload payload = collector.collectMetrics(this::emit, startNanos, endNanos); if (payload != OtlpPayload.EMPTY) { - sender.send(payload); + OtlpTelemetry.getInstance().onMetricsExportAttempt(); + RemoteApi.Response response = sender.send(payload); + OtlpTelemetry.getInstance().onMetricsExportComplete(response.success()); } } finally { pending.clear(); diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollector.java index c5da9bcddaf..7ed1f2f7563 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceJsonCollector.java @@ -117,17 +117,18 @@ private void visitScopedSpans(OtelInstrumentationScope scope) { } private void visitSpan(CoreSpan span) { - if (shouldExport(span)) { - if (currentSpan != null) { - // ensure last span written at trace boundary includes sampling tags - if (!span.getTraceId().equals(currentSpan.getTraceId())) { - metaWriter.includeSamplingTags(); - } - completeSpan(); + if (!shouldExport(span)) { + return; + } + if (currentSpan != null) { + // ensure last span written at trace boundary includes sampling tags + if (!span.getTraceId().equals(currentSpan.getTraceId())) { + metaWriter.includeSamplingTags(); } - currentSpan = (DDSpan) span; - currentSpanLinks = currentSpan.getLinks(); + completeSpan(); } + currentSpan = (DDSpan) span; + currentSpanLinks = currentSpan.getLinks(); } // called once we've processed all scopes and span messages diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProtoCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProtoCollector.java index fa1efa87fba..969fc69ff64 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProtoCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceProtoCollector.java @@ -78,7 +78,6 @@ public OtlpPayload collectTraces() { /** Prepare temporary elements to collect trace data. */ private void start() { - // remove stale entries from caches OtlpCommonProto.recalibrateCaches(); @@ -109,18 +108,19 @@ private void visitScopedSpans(OtelInstrumentationScope scope) { } private void visitSpan(CoreSpan span) { - if (shouldExport(span)) { - if (currentSpan != null) { - // ensure last span written at trace boundary includes sampling tags - // payload buffer is prepending, so last span written appears first! - if (!span.getTraceId().equals(currentSpan.getTraceId())) { - metaWriter.includeSamplingTags(); - } - completeSpan(); + if (!shouldExport(span)) { + return; + } + if (currentSpan != null) { + // ensure last span written at trace boundary includes sampling tags + // payload buffer is prepending, so last span written appears first! + if (!span.getTraceId().equals(currentSpan.getTraceId())) { + metaWriter.includeSamplingTags(); } - currentSpan = (DDSpan) span; - currentSpan.getLinks().forEach(this::visitSpanLink); + completeSpan(); } + currentSpan = (DDSpan) span; + currentSpan.getLinks().forEach(this::visitSpanLink); } private void visitSpanLink(AgentSpanLink spanLink) { diff --git a/dd-trace-core/src/test/java/datadog/trace/common/writer/OtlpPayloadDispatcherTest.java b/dd-trace-core/src/test/java/datadog/trace/common/writer/OtlpPayloadDispatcherTest.java index 6e55ebeadaa..b4ac7542946 100644 --- a/dd-trace-core/src/test/java/datadog/trace/common/writer/OtlpPayloadDispatcherTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/common/writer/OtlpPayloadDispatcherTest.java @@ -5,12 +5,15 @@ import static java.util.Collections.singletonList; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; import datadog.trace.api.sampling.PrioritySampling; +import datadog.trace.api.telemetry.OtlpTelemetry; import datadog.trace.core.CoreSpan; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; @@ -18,7 +21,10 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Arrays; +import java.util.HashMap; import java.util.List; +import java.util.Map; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.ArgumentCaptor; @@ -31,6 +37,13 @@ class OtlpPayloadDispatcherTest { TestCollector collector = new TestCollector(); + @BeforeEach + void stubSuccessfulSend() { + lenient().when(sender.send(any())).thenReturn(RemoteApi.Response.success(200)); + OtlpTelemetry.getInstance().prepareMetrics(); + OtlpTelemetry.getInstance().drain(); + } + @Test void sampledTraceForwardsAllSpans() { OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); @@ -95,6 +108,45 @@ void emptyTraceForwardsNothing() { verifyNoInteractions(sender); } + @Test + void flushRecordsSuccessfulExportTelemetry() { + RemoteApi.Response response = RemoteApi.Response.success(200); + when(sender.send(any())).thenReturn(response); + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); + + dispatcher.addTrace(Arrays.asList(sampledSpan(), sampledSpan())); + dispatcher.flush(); + + Map metrics = drainTracesTelemetry(); + assertEquals(2, metrics.size()); + assertEquals(1L, metrics.get("otel.traces_export_attempts").value); + assertEquals(1L, metrics.get("otel.traces_export_successes").value); + } + + @Test + void flushRecordsFailedExportTelemetry() { + RemoteApi.Response response = RemoteApi.Response.failed(500); + when(sender.send(any())).thenReturn(response); + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); + + dispatcher.addTrace(Arrays.asList(sampledSpan(), sampledSpan())); + dispatcher.flush(); + + Map metrics = drainTracesTelemetry(); + assertEquals(2, metrics.size()); + assertEquals(1L, metrics.get("otel.traces_export_attempts").value); + assertEquals(1L, metrics.get("otel.traces_export_failures").value); + } + + private static Map drainTracesTelemetry() { + Map byName = new HashMap<>(); + OtlpTelemetry.getInstance().prepareMetrics(); + for (OtlpTelemetry.OtlpMetric metric : OtlpTelemetry.getInstance().drain()) { + byName.put(metric.metricName, metric); + } + return byName; + } + @Test void getApisIsEmpty() { OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpSenderSupportTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpSenderSupportTest.java new file mode 100644 index 00000000000..2d033113213 --- /dev/null +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpSenderSupportTest.java @@ -0,0 +1,95 @@ +package datadog.trace.core.otlp.common; + +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 static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.mockStatic; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +import datadog.communication.http.HttpRetryPolicy; +import datadog.communication.http.OkHttpUtils; +import datadog.logging.RatelimitedLogger; +import datadog.trace.common.writer.RemoteApi; +import java.io.IOException; +import okhttp3.MediaType; +import okhttp3.OkHttpClient; +import okhttp3.Protocol; +import okhttp3.Request; +import okhttp3.Response; +import okhttp3.ResponseBody; +import org.junit.jupiter.api.Test; +import org.mockito.MockedStatic; + +class OtlpSenderSupportTest { + + private final OkHttpClient client = mock(OkHttpClient.class); + private final HttpRetryPolicy.Factory retryPolicy = HttpRetryPolicy.Factory.NEVER_RETRY; + private final Request request = + new Request.Builder().url("http://localhost:4318/v1/traces").build(); + private final RatelimitedLogger ratelimitedLogger = mock(RatelimitedLogger.class); + + @Test + void successfulResponseIsReturnedWithoutLogging() throws IOException { + Response response = responseWithCode(200); + try (MockedStatic okHttpUtils = mockStatic(OkHttpUtils.class)) { + okHttpUtils + .when(() -> OkHttpUtils.sendWithRetries(client, retryPolicy, request)) + .thenReturn(response); + + RemoteApi.Response result = + OtlpSenderSupport.send(client, retryPolicy, request, ratelimitedLogger); + + assertTrue(result.success()); + assertEquals(200, result.status().getAsInt()); + verify(ratelimitedLogger, never()).warn(any(String.class), any()); + } + } + + @Test + void unsuccessfulResponseIsReturnedAndLogged() throws IOException { + Response response = responseWithCode(500); + try (MockedStatic okHttpUtils = mockStatic(OkHttpUtils.class)) { + okHttpUtils + .when(() -> OkHttpUtils.sendWithRetries(client, retryPolicy, request)) + .thenReturn(response); + + RemoteApi.Response result = + OtlpSenderSupport.send(client, retryPolicy, request, ratelimitedLogger); + + assertFalse(result.success()); + assertEquals(500, result.status().getAsInt()); + verify(ratelimitedLogger).warn(any(String.class), any(), any(), any()); + } + } + + @Test + void ioExceptionIsReturnedAsFailureAndLogged() throws IOException { + IOException exception = new IOException("boom"); + try (MockedStatic okHttpUtils = mockStatic(OkHttpUtils.class)) { + okHttpUtils + .when(() -> OkHttpUtils.sendWithRetries(client, retryPolicy, request)) + .thenThrow(exception); + + RemoteApi.Response result = + OtlpSenderSupport.send(client, retryPolicy, request, ratelimitedLogger); + + assertFalse(result.success()); + assertTrue(result.exception().isPresent()); + assertEquals(exception, result.exception().get()); + verify(ratelimitedLogger).warn(any(String.class), any(), any()); + } + } + + private Response responseWithCode(int code) { + return new Response.Builder() + .request(request) + .protocol(Protocol.HTTP_1_1) + .code(code) + .message(code == 200 ? "OK" : "Server Error") + .body(ResponseBody.create(MediaType.get("text/plain"), "")) + .build(); + } +} diff --git a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpStatsMetricWriterTest.java b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpStatsMetricWriterTest.java index 670f9098a14..49e252459ac 100644 --- a/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpStatsMetricWriterTest.java +++ b/dd-trace-core/src/test/java/datadog/trace/core/otlp/metrics/OtlpStatsMetricWriterTest.java @@ -1,5 +1,6 @@ package datadog.trace.core.otlp.metrics; +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; @@ -14,6 +15,7 @@ import datadog.metrics.impl.DDSketchHistograms; import datadog.trace.common.metrics.AggregateEntry; import datadog.trace.common.metrics.AggregateEntryTestUtils; +import datadog.trace.common.writer.RemoteApi; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; import java.io.IOException; @@ -64,12 +66,13 @@ private static final class CapturingSender implements OtlpSender { byte[] lastPayload; @Override - public void send(OtlpPayload payload) { + public RemoteApi.Response send(OtlpPayload payload) { sendCount++; java.nio.ByteBuffer content = payload.getContent(); byte[] bytes = new byte[content.remaining()]; content.get(bytes); lastPayload = bytes; + return success(200); } @Override diff --git a/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java b/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java new file mode 100644 index 00000000000..a36bcff67c8 --- /dev/null +++ b/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java @@ -0,0 +1,123 @@ +package datadog.trace.api.telemetry; + +import datadog.trace.api.Config; +import datadog.trace.api.config.OtlpConfig; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.atomic.LongAdder; + +/** Collects telemetry metrics for the OTLP trace, metrics, and log exporters. */ +public class OtlpTelemetry implements MetricCollector { + private static final String NAMESPACE = "tracers"; + + private static final OtlpTelemetry INSTANCE = new OtlpTelemetry(); + + public static OtlpTelemetry getInstance() { + return INSTANCE; + } + + private final String[] tracesTags = tagsFor(Config.get().getOtlpTracesProtocol()); + private final String[] metricsTags = tagsFor(Config.get().getOtlpMetricsProtocol()); + private final String[] logsTags = tagsFor(Config.get().getOtlpLogsProtocol()); + + private final ExportCounters tracesExport = new ExportCounters("traces"); + private final ExportCounters metricsExport = new ExportCounters("metrics"); + private final LongAdder logRecords = new LongAdder(); + + private final BlockingQueue telemetryQueue = new ArrayBlockingQueue<>(RAW_QUEUE_SIZE); + + private OtlpTelemetry() {} + + public void onTracesExportAttempt() { + tracesExport.attempts.increment(); + } + + public void onTracesExportComplete(boolean success) { + tracesExport.complete(success); + } + + public void onMetricsExportAttempt() { + metricsExport.attempts.increment(); + } + + public void onMetricsExportComplete(boolean success) { + metricsExport.complete(success); + } + + public void onLogRecordsSubmitted(long count) { + if (count > 0) { + logRecords.add(count); + } + } + + private static String[] tagsFor(OtlpConfig.Protocol protocol) { + String protocolTag = protocol == OtlpConfig.Protocol.GRPC ? "grpc" : "http"; + String encodingTag = protocol == OtlpConfig.Protocol.HTTP_JSON ? "json" : "protobuf"; + return new String[] {"protocol:" + protocolTag, "encoding:" + encodingTag}; + } + + @Override + public void prepareMetrics() { + tracesExport.stageInto(telemetryQueue, tracesTags); + metricsExport.stageInto(telemetryQueue, metricsTags); + long logRecordCount = logRecords.sumThenReset(); + if (logRecordCount > 0) { + telemetryQueue.offer(new OtlpMetric("otel.log_records", logRecordCount, logsTags)); + } + } + + @Override + public Collection drain() { + if (telemetryQueue.isEmpty()) { + return Collections.emptyList(); + } + List drained = new ArrayList<>(telemetryQueue.size()); + telemetryQueue.drainTo(drained); + return drained; + } + + /** Counters for a single signal's export attempts/successes/failures. */ + private static final class ExportCounters { + final String attemptsMetric; + final String successesMetric; + final String failuresMetric; + + final LongAdder attempts = new LongAdder(); + final LongAdder successes = new LongAdder(); + final LongAdder failures = new LongAdder(); + + ExportCounters(String signal) { + this.attemptsMetric = "otel." + signal + "_export_attempts"; + this.successesMetric = "otel." + signal + "_export_successes"; + this.failuresMetric = "otel." + signal + "_export_failures"; + } + + void complete(boolean success) { + (success ? successes : failures).increment(); + } + + void stageInto(BlockingQueue out, String[] tags) { + addIfNonZero(out, attemptsMetric, attempts, tags); + addIfNonZero(out, successesMetric, successes, tags); + addIfNonZero(out, failuresMetric, failures, tags); + } + + private static void addIfNonZero( + BlockingQueue out, String metricName, LongAdder counter, String[] tags) { + long value = counter.sumThenReset(); + if (value > 0) { + out.offer(new OtlpMetric(metricName, value, tags)); + } + } + } + + public static class OtlpMetric extends MetricCollector.Metric { + public OtlpMetric(String metricName, long value, String... tags) { + super(NAMESPACE, true, metricName, "count", value, tags); + } + } +} diff --git a/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java b/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java new file mode 100644 index 00000000000..53877bd8dc6 --- /dev/null +++ b/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java @@ -0,0 +1,124 @@ +package datadog.trace.api.telemetry; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.Collection; +import java.util.HashMap; +import java.util.Map; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +class OtlpTelemetryTest { + private final OtlpTelemetry collector = OtlpTelemetry.getInstance(); + + @BeforeEach + void drainStaleMetrics() { + collector.prepareMetrics(); + collector.drain(); + } + + @Test + void onTracesExportAttemptQueuesCountMetricWithTags() { + collector.onTracesExportAttempt(); + collector.prepareMetrics(); + + Collection metrics = collector.drain(); + + assertEquals(1, metrics.size()); + OtlpTelemetry.OtlpMetric metric = metrics.iterator().next(); + assertEquals("tracers", metric.namespace); + assertEquals("otel.traces_export_attempts", metric.metricName); + assertEquals("count", metric.type); + assertEquals(1L, metric.value); + assertEquals(2, metric.tags.size()); + assertTrue(metric.tags.contains("protocol:http")); + assertTrue(metric.tags.contains("encoding:protobuf")); + } + + @Test + void tracesExportMetricsUseExpectedNames() { + collector.onTracesExportAttempt(); + collector.onTracesExportAttempt(); + collector.onTracesExportComplete(true); + collector.onTracesExportComplete(false); + collector.prepareMetrics(); + + Map valuesByName = drainToMap(); + + assertEquals(3, valuesByName.size()); + assertEquals(2L, valuesByName.get("otel.traces_export_attempts")); + assertEquals(1L, valuesByName.get("otel.traces_export_successes")); + assertEquals(1L, valuesByName.get("otel.traces_export_failures")); + } + + @Test + void onMetricsExportAttemptQueuesCountMetricWithTags() { + collector.onMetricsExportAttempt(); + collector.prepareMetrics(); + + Collection metrics = collector.drain(); + + assertEquals(1, metrics.size()); + OtlpTelemetry.OtlpMetric metric = metrics.iterator().next(); + assertEquals("tracers", metric.namespace); + assertEquals("otel.metrics_export_attempts", metric.metricName); + assertEquals("count", metric.type); + assertEquals(1L, metric.value); + assertEquals(2, metric.tags.size()); + assertTrue(metric.tags.contains("protocol:http")); + assertTrue(metric.tags.contains("encoding:protobuf")); + } + + @Test + void metricsExportMetricsUseExpectedNames() { + collector.onMetricsExportAttempt(); + collector.onMetricsExportAttempt(); + collector.onMetricsExportComplete(true); + collector.onMetricsExportComplete(false); + collector.prepareMetrics(); + + Map valuesByName = drainToMap(); + + assertEquals(3, valuesByName.size()); + assertEquals(2L, valuesByName.get("otel.metrics_export_attempts")); + assertEquals(1L, valuesByName.get("otel.metrics_export_successes")); + assertEquals(1L, valuesByName.get("otel.metrics_export_failures")); + } + + @Test + void onLogRecordsSubmittedQueuesCountWithGivenValue() { + collector.onLogRecordsSubmitted(5); + collector.prepareMetrics(); + + Collection metrics = collector.drain(); + + assertEquals(1, metrics.size()); + OtlpTelemetry.OtlpMetric metric = metrics.iterator().next(); + assertEquals("otel.log_records", metric.metricName); + assertEquals(5L, metric.value); + assertTrue(metric.tags.contains("protocol:http")); + assertTrue(metric.tags.contains("encoding:protobuf")); + } + + @Test + void onLogRecordsSubmittedIgnoresNonPositiveCounts() { + collector.onLogRecordsSubmitted(0); + collector.onLogRecordsSubmitted(-1); + + assertTrue(collector.drain().isEmpty()); + } + + @Test + void drainReturnsEmptyCollectionWhenNoMetricsQueued() { + assertTrue(collector.drain().isEmpty()); + } + + private Map drainToMap() { + Map valuesByName = new HashMap<>(); + for (OtlpTelemetry.OtlpMetric metric : collector.drain()) { + valuesByName.put(metric.metricName, metric.value); + } + return valuesByName; + } +} diff --git a/telemetry/build.gradle.kts b/telemetry/build.gradle.kts index c6e72af33ed..755fd138b8a 100644 --- a/telemetry/build.gradle.kts +++ b/telemetry/build.gradle.kts @@ -18,7 +18,8 @@ extra["excludedClassesCoverage"] = listOf( "datadog.telemetry.TelemetrySystem", "datadog.telemetry.api.*", "datadog.telemetry.metric.CiVisibilityMetricPeriodicAction", - "datadog.telemetry.metric.OtelSpiMetricPeriodicAction" + "datadog.telemetry.metric.OtelSpiMetricPeriodicAction", + "datadog.telemetry.metric.OtlpTelemetryPeriodicAction" ) extra["excludedClassesBranchCoverage"] = listOf( "datadog.telemetry.PolymorphicAdapterFactory.1", diff --git a/telemetry/src/main/java/datadog/telemetry/TelemetrySystem.java b/telemetry/src/main/java/datadog/telemetry/TelemetrySystem.java index de6b2feed19..31b833aa42a 100644 --- a/telemetry/src/main/java/datadog/telemetry/TelemetrySystem.java +++ b/telemetry/src/main/java/datadog/telemetry/TelemetrySystem.java @@ -16,6 +16,7 @@ import datadog.telemetry.metric.LLMObsMetricPeriodicAction; import datadog.telemetry.metric.OtelEnvMetricPeriodicAction; import datadog.telemetry.metric.OtelSpiMetricPeriodicAction; +import datadog.telemetry.metric.OtlpTelemetryPeriodicAction; import datadog.telemetry.metric.WafMetricPeriodicAction; import datadog.telemetry.products.ProductChangeAction; import datadog.telemetry.rum.RumPeriodicAction; @@ -65,6 +66,7 @@ static Thread createTelemetryRunnable( actions.add(new ConfigInversionMetricPeriodicAction()); actions.add(new IntegrationPeriodicAction()); actions.add(new WafMetricPeriodicAction()); + actions.add(new OtlpTelemetryPeriodicAction()); if (Verbosity.OFF != Config.get().getIastTelemetryVerbosity()) { actions.add(new IastMetricPeriodicAction()); } diff --git a/telemetry/src/main/java/datadog/telemetry/metric/OtlpTelemetryPeriodicAction.java b/telemetry/src/main/java/datadog/telemetry/metric/OtlpTelemetryPeriodicAction.java new file mode 100644 index 00000000000..0464a353541 --- /dev/null +++ b/telemetry/src/main/java/datadog/telemetry/metric/OtlpTelemetryPeriodicAction.java @@ -0,0 +1,14 @@ +package datadog.telemetry.metric; + +import datadog.trace.api.telemetry.MetricCollector; +import datadog.trace.api.telemetry.OtlpTelemetry; +import javax.annotation.Nonnull; + +public class OtlpTelemetryPeriodicAction extends MetricPeriodicAction { + + @Override + @Nonnull + public MetricCollector collector() { + return OtlpTelemetry.getInstance(); + } +}