From 74e91ad97b1e4152ee88dc99bfaaf902c135ae53 Mon Sep 17 00:00:00 2001 From: Stuart McCulloch Date: Thu, 23 Jul 2026 15:30:36 +0100 Subject: [PATCH 1/4] Telemetry for OTLP traces/metrics/logs: * Datadog tracer health metrics * otel.metrics_export_{attempts,successes,failures} * otel.log_records OTLP telemetry metrics are tagged by protocol and encoding (per signal) Co-Authored-By: Claude Sonnet 5 --- .../common/writer/OtlpPayloadDispatcher.java | 20 +++- .../trace/common/writer/OtlpWriter.java | 3 +- .../core/otlp/common/OtlpGrpcSender.java | 20 +--- .../core/otlp/common/OtlpHttpSender.java | 20 +--- .../trace/core/otlp/common/OtlpSender.java | 4 +- .../core/otlp/common/OtlpSenderSupport.java | 39 +++++++ .../core/otlp/logs/OtlpLogsCollector.java | 3 + .../core/otlp/logs/OtlpLogsJsonCollector.java | 9 ++ .../otlp/logs/OtlpLogsProtoCollector.java | 8 ++ .../trace/core/otlp/logs/OtlpLogsService.java | 8 +- .../core/otlp/metrics/OtlpMetricsService.java | 11 +- .../otlp/metrics/OtlpStatsMetricWriter.java | 10 +- .../core/otlp/trace/OtlpTraceCollector.java | 3 + .../otlp/trace/OtlpTraceJsonCollector.java | 36 ++++-- .../otlp/trace/OtlpTraceProtoCollector.java | 37 +++++-- .../writer/OtlpPayloadDispatcherTest.java | 53 +++++++++ .../otlp/common/OtlpSenderSupportTest.java | 95 ++++++++++++++++ .../metrics/OtlpStatsMetricWriterTest.java | 5 +- .../trace/api/telemetry/OtlpTelemetry.java | 104 ++++++++++++++++++ .../api/telemetry/OtlpTelemetryTest.java | 76 +++++++++++++ telemetry/build.gradle.kts | 3 +- .../datadog/telemetry/TelemetrySystem.java | 2 + .../metric/OtlpTelemetryPeriodicAction.java | 14 +++ 23 files changed, 516 insertions(+), 67 deletions(-) create mode 100644 dd-trace-core/src/main/java/datadog/trace/core/otlp/common/OtlpSenderSupport.java create mode 100644 dd-trace-core/src/test/java/datadog/trace/core/otlp/common/OtlpSenderSupportTest.java create mode 100644 internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java create mode 100644 internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java create mode 100644 telemetry/src/main/java/datadog/telemetry/metric/OtlpTelemetryPeriodicAction.java 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..a4235e77ab4 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,6 +1,7 @@ package datadog.trace.common.writer; import datadog.trace.core.CoreSpan; +import datadog.trace.core.monitor.HealthMetrics; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; import datadog.trace.core.otlp.trace.OtlpTraceCollector; @@ -11,10 +12,17 @@ final class OtlpPayloadDispatcher implements PayloadDispatcher { private final OtlpTraceCollector collector; private final OtlpSender sender; + private final HealthMetrics healthMetrics; OtlpPayloadDispatcher(OtlpSender sender, OtlpTraceCollector collector) { + this(sender, collector, HealthMetrics.NO_OP); + } + + OtlpPayloadDispatcher( + OtlpSender sender, OtlpTraceCollector collector, HealthMetrics healthMetrics) { this.sender = sender; this.collector = collector; + this.healthMetrics = healthMetrics; } @Override @@ -25,14 +33,22 @@ public void addTrace(List> trace) { @Override public void flush() { OtlpPayload payload = collector.collectTraces(); + int traceCount = collector.getTraceCount(); if (payload != OtlpPayload.EMPTY) { - sender.send(payload); + int sizeInBytes = payload.getContentLength(); + healthMetrics.onSerialize(sizeInBytes); + RemoteApi.Response response = sender.send(payload); + if (response.success()) { + healthMetrics.onSend(traceCount, sizeInBytes, response); + } else { + healthMetrics.onFailedSend(traceCount, sizeInBytes, response); + } } } @Override public void onDroppedTrace(int spanCount) { - // TODO: surface drop counts via HealthMetrics + // RemoteWriter already updated healthMetrics, no further action required } @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..15a39b801df 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 @@ -147,7 +147,8 @@ public OtlpWriter build() { protocol == OtlpConfig.Protocol.HTTP_JSON ? new OtlpTraceJsonCollector() : new OtlpTraceProtoCollector(); - final OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); + final OtlpPayloadDispatcher dispatcher = + new OtlpPayloadDispatcher(sender, collector, healthMetrics); final TraceProcessingWorker worker = new TraceProcessingWorker( traceBufferSize, 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..9abdbf59af3 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,13 @@ 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); + if (response.success()) { + OtlpTelemetry.getInstance().onMetricsExportSuccess(); + } else { + OtlpTelemetry.getInstance().onMetricsExportFailure(); + } } } } 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..c30b49a6ab2 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,13 @@ 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); + if (response.success()) { + OtlpTelemetry.getInstance().onMetricsExportSuccess(); + } else { + OtlpTelemetry.getInstance().onMetricsExportFailure(); + } } } finally { pending.clear(); diff --git a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java index dda40ff74b2..87364a8475e 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java @@ -15,6 +15,9 @@ public abstract class OtlpTraceCollector { /** Collects all spans added since the last collection. */ public abstract OtlpPayload collectTraces(); + /** Number of traces collected since the last collection. */ + public abstract int getTraceCount(); + protected final boolean shouldExport(CoreSpan span) { return span.samplingPriority() > 0 // trace-level sampling priority || span.getTag(SPAN_SAMPLING_MECHANISM_TAG) != null; // span-level sampling priority 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..15c2576310e 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 @@ -41,6 +41,7 @@ public final class OtlpTraceJsonCollector extends OtlpTraceCollector { private boolean payloadStarted; private boolean anySpanWritten; private boolean firstSpanInScope; + private int traceCount; private OtelInstrumentationScope currentScope; private DDSpan currentSpan; @@ -54,8 +55,12 @@ public void addTrace(List> spans) { payloadStarted = true; } + boolean exported = false; for (CoreSpan span : spans) { - visitSpan(span); + exported |= visitSpan(span); + } + if (exported) { + traceCount++; } } @@ -78,6 +83,8 @@ public OtlpPayload collectTraces() { /** Prepare temporary elements to collect trace data. */ private void start() { + traceCount = 0; + writer = new JsonWriter(); metaWriter = new OtlpTraceJson.MetaWriter(writer); @@ -104,6 +111,11 @@ private void stop() { currentSpanLinks = Collections.emptyList(); } + @Override + public int getTraceCount() { + return traceCount; + } + private void visitScopedSpans(OtelInstrumentationScope scope) { if (currentScope != null) { completeScope(); @@ -116,18 +128,20 @@ private void visitScopedSpans(OtelInstrumentationScope scope) { writer.name("spans").beginArray(); } - 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(); + private boolean visitSpan(CoreSpan span) { + if (!shouldExport(span)) { + return false; + } + 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(); + return true; } // 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..56a469edb25 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 @@ -44,6 +44,7 @@ public final class OtlpTraceProtoCollector extends OtlpTraceCollector { private int payloadBytes; private int scopedBytes; private int spanBytes; + private int traceCount; private OtelInstrumentationScope currentScope; private DDSpan currentSpan; @@ -56,9 +57,13 @@ public void addTrace(List> spans) { payloadStarted = true; } + boolean exported = false; // OtlpProtoBuffer collects spans in reverse for (int i = spans.size() - 1; i >= 0; i--) { - visitSpan(spans.get(i)); + exported |= visitSpan(spans.get(i)); + } + if (exported) { + traceCount++; } } @@ -78,6 +83,7 @@ public OtlpPayload collectTraces() { /** Prepare temporary elements to collect trace data. */ private void start() { + traceCount = 0; // remove stale entries from caches OtlpCommonProto.recalibrateCaches(); @@ -86,6 +92,11 @@ private void start() { visitScopedSpans(DEFAULT_TRACE_SCOPE); } + @Override + public int getTraceCount() { + return traceCount; + } + /** Cleanup elements used to collect trace data. */ private void stop() { payloadStarted = false; @@ -108,19 +119,21 @@ private void visitScopedSpans(OtelInstrumentationScope scope) { currentScope = 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(); + private boolean visitSpan(CoreSpan span) { + if (!shouldExport(span)) { + return false; + } + 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); + return true; } 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..a1442ae0f87 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,6 +5,10 @@ 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.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.same; +import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; @@ -12,6 +16,7 @@ import datadog.trace.api.sampling.PrioritySampling; import datadog.trace.core.CoreSpan; +import datadog.trace.core.monitor.HealthMetrics; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; import datadog.trace.core.otlp.trace.OtlpTraceCollector; @@ -19,6 +24,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.List; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.ArgumentCaptor; @@ -28,9 +34,15 @@ @ExtendWith(MockitoExtension.class) class OtlpPayloadDispatcherTest { @Mock OtlpSender sender; + @Mock HealthMetrics healthMetrics; TestCollector collector = new TestCollector(); + @BeforeEach + void stubSuccessfulSend() { + lenient().when(sender.send(any())).thenReturn(RemoteApi.Response.success(200)); + } + @Test void sampledTraceForwardsAllSpans() { OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); @@ -95,6 +107,32 @@ void emptyTraceForwardsNothing() { verifyNoInteractions(sender); } + @Test + void flushRecordsSerializedSizeBeforeSending() { + RemoteApi.Response response = RemoteApi.Response.success(200); + when(sender.send(any())).thenReturn(response); + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector, healthMetrics); + + dispatcher.addTrace(Arrays.asList(sampledSpan(), sampledSpan())); + dispatcher.flush(); + + verify(healthMetrics).onSerialize(2 /*spans*/); + verify(healthMetrics).onSend(eq(1) /*trace*/, eq(2) /*spans*/, same(response)); + } + + @Test + void flushRecordsFailedSendHealthMetric() { + RemoteApi.Response response = RemoteApi.Response.failed(500); + when(sender.send(any())).thenReturn(response); + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector, healthMetrics); + + dispatcher.addTrace(Arrays.asList(sampledSpan(), sampledSpan())); + dispatcher.flush(); + + verify(healthMetrics).onSerialize(2 /*spans*/); + verify(healthMetrics).onFailedSend(eq(1) /*trace*/, eq(2) /*spans*/, same(response)); + } + @Test void getApisIsEmpty() { OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); @@ -142,14 +180,24 @@ private static CoreSpan unsetSpan() { /** Test collector that creates payloads whose size equals the number of exported spans. */ private static class TestCollector extends OtlpTraceCollector { final List> spansToExport = new ArrayList<>(); + int traceCount; @Override public void addTrace(List> spans) { + if (spansToExport.isEmpty()) { + // starting a new batch - reset the count left over from the last collection + traceCount = 0; + } + boolean exported = false; for (CoreSpan span : spans) { if (shouldExport(span)) { spansToExport.add(span); + exported = true; } } + if (exported) { + traceCount++; + } } @Override @@ -165,5 +213,10 @@ public OtlpPayload collectTraces() { spansToExport.clear(); } } + + @Override + public int getTraceCount() { + return traceCount; + } } } 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..0babcd1177c --- /dev/null +++ b/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java @@ -0,0 +1,104 @@ +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.List; +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[] metricsTags = tagsFor(Config.get().getOtlpMetricsProtocol()); + private final String[] logsTags = tagsFor(Config.get().getOtlpLogsProtocol()); + + private final ExportCounters metricsExport = new ExportCounters("metrics"); + private final LongAdder logRecords = new LongAdder(); + + private OtlpTelemetry() {} + + public void onMetricsExportAttempt() { + metricsExport.attempts.increment(); + } + + public void onMetricsExportSuccess() { + metricsExport.successes.increment(); + } + + public void onMetricsExportFailure() { + metricsExport.failures.increment(); + } + + 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() { + // metrics are accumulated directly as they happen; nothing to prepare + } + + @Override + public Collection drain() { + List drained = new ArrayList<>(); + metricsExport.drainInto(drained, metricsTags); + long logRecordCount = logRecords.sumThenReset(); + if (logRecordCount > 0) { + drained.add(new OtlpMetric("otel.log_records", logRecordCount, logsTags)); + } + 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 drainInto(List out, String[] tags) { + addIfNonZero(out, attemptsMetric, attempts, tags); + addIfNonZero(out, successesMetric, successes, tags); + addIfNonZero(out, failuresMetric, failures, tags); + } + + private static void addIfNonZero( + List out, String metricName, LongAdder counter, String[] tags) { + long value = counter.sumThenReset(); + if (value > 0) { + out.add(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..7c5472328a1 --- /dev/null +++ b/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java @@ -0,0 +1,76 @@ +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 org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +class OtlpTelemetryTest { + private final OtlpTelemetry collector = OtlpTelemetry.getInstance(); + + @BeforeEach + void drainStaleMetrics() { + collector.drain(); + } + + @Test + void onMetricsExportAttemptQueuesCountMetricWithTags() { + collector.onMetricsExportAttempt(); + + 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.onMetricsExportSuccess(); + collector.onMetricsExportFailure(); + + Collection metrics = collector.drain(); + + assertEquals(3, metrics.size()); + assertTrue(metrics.stream().anyMatch(m -> m.metricName.equals("otel.metrics_export_attempts"))); + assertTrue( + metrics.stream().anyMatch(m -> m.metricName.equals("otel.metrics_export_successes"))); + assertTrue(metrics.stream().anyMatch(m -> m.metricName.equals("otel.metrics_export_failures"))); + } + + @Test + void onLogRecordsSubmittedQueuesCountWithGivenValue() { + collector.onLogRecordsSubmitted(5); + + 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()); + } +} 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(); + } +} From 91c11f7537373fb4b1208d435fb17be2425427e8 Mon Sep 17 00:00:00 2001 From: Stuart McCulloch Date: Thu, 23 Jul 2026 23:03:15 +0100 Subject: [PATCH 2/4] Remove HealthMetrics for OTLP traces (since HealthMetrics relies on StatsD client) and use telemetry instead --- .../common/writer/OtlpPayloadDispatcher.java | 19 ++----- .../trace/common/writer/OtlpWriter.java | 16 ++---- .../trace/common/writer/WriterFactory.java | 1 - .../core/otlp/trace/OtlpTraceCollector.java | 3 -- .../otlp/trace/OtlpTraceJsonCollector.java | 19 ++----- .../otlp/trace/OtlpTraceProtoCollector.java | 19 ++----- .../writer/OtlpPayloadDispatcherTest.java | 51 +++++++++--------- .../trace/api/telemetry/OtlpTelemetry.java | 15 ++++++ .../api/telemetry/OtlpTelemetryTest.java | 54 ++++++++++++++++--- 9 files changed, 102 insertions(+), 95 deletions(-) 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 a4235e77ab4..398f278ba94 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,7 +1,7 @@ package datadog.trace.common.writer; +import datadog.trace.api.telemetry.OtlpTelemetry; import datadog.trace.core.CoreSpan; -import datadog.trace.core.monitor.HealthMetrics; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; import datadog.trace.core.otlp.trace.OtlpTraceCollector; @@ -12,17 +12,10 @@ final class OtlpPayloadDispatcher implements PayloadDispatcher { private final OtlpTraceCollector collector; private final OtlpSender sender; - private final HealthMetrics healthMetrics; OtlpPayloadDispatcher(OtlpSender sender, OtlpTraceCollector collector) { - this(sender, collector, HealthMetrics.NO_OP); - } - - OtlpPayloadDispatcher( - OtlpSender sender, OtlpTraceCollector collector, HealthMetrics healthMetrics) { this.sender = sender; this.collector = collector; - this.healthMetrics = healthMetrics; } @Override @@ -33,22 +26,20 @@ public void addTrace(List> trace) { @Override public void flush() { OtlpPayload payload = collector.collectTraces(); - int traceCount = collector.getTraceCount(); if (payload != OtlpPayload.EMPTY) { - int sizeInBytes = payload.getContentLength(); - healthMetrics.onSerialize(sizeInBytes); + OtlpTelemetry.getInstance().onTracesExportAttempt(); RemoteApi.Response response = sender.send(payload); if (response.success()) { - healthMetrics.onSend(traceCount, sizeInBytes, response); + OtlpTelemetry.getInstance().onTracesExportSuccess(); } else { - healthMetrics.onFailedSend(traceCount, sizeInBytes, response); + OtlpTelemetry.getInstance().onTracesExportFailure(); } } } @Override public void onDroppedTrace(int spanCount) { - // RemoteWriter already updated healthMetrics, no further action required + // 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 15a39b801df..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; @@ -147,12 +140,11 @@ public OtlpWriter build() { protocol == OtlpConfig.Protocol.HTTP_JSON ? new OtlpTraceJsonCollector() : new OtlpTraceProtoCollector(); - final OtlpPayloadDispatcher dispatcher = - new OtlpPayloadDispatcher(sender, collector, healthMetrics); + final OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); final TraceProcessingWorker worker = new TraceProcessingWorker( traceBufferSize, - healthMetrics, + HealthMetrics.NO_OP, dispatcher, DroppingPolicy.DISABLED, Prioritization.FAST_LANE, @@ -161,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/trace/OtlpTraceCollector.java b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java index 87364a8475e..dda40ff74b2 100644 --- a/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java +++ b/dd-trace-core/src/main/java/datadog/trace/core/otlp/trace/OtlpTraceCollector.java @@ -15,9 +15,6 @@ public abstract class OtlpTraceCollector { /** Collects all spans added since the last collection. */ public abstract OtlpPayload collectTraces(); - /** Number of traces collected since the last collection. */ - public abstract int getTraceCount(); - protected final boolean shouldExport(CoreSpan span) { return span.samplingPriority() > 0 // trace-level sampling priority || span.getTag(SPAN_SAMPLING_MECHANISM_TAG) != null; // span-level sampling priority 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 15c2576310e..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 @@ -41,7 +41,6 @@ public final class OtlpTraceJsonCollector extends OtlpTraceCollector { private boolean payloadStarted; private boolean anySpanWritten; private boolean firstSpanInScope; - private int traceCount; private OtelInstrumentationScope currentScope; private DDSpan currentSpan; @@ -55,12 +54,8 @@ public void addTrace(List> spans) { payloadStarted = true; } - boolean exported = false; for (CoreSpan span : spans) { - exported |= visitSpan(span); - } - if (exported) { - traceCount++; + visitSpan(span); } } @@ -83,8 +78,6 @@ public OtlpPayload collectTraces() { /** Prepare temporary elements to collect trace data. */ private void start() { - traceCount = 0; - writer = new JsonWriter(); metaWriter = new OtlpTraceJson.MetaWriter(writer); @@ -111,11 +104,6 @@ private void stop() { currentSpanLinks = Collections.emptyList(); } - @Override - public int getTraceCount() { - return traceCount; - } - private void visitScopedSpans(OtelInstrumentationScope scope) { if (currentScope != null) { completeScope(); @@ -128,9 +116,9 @@ private void visitScopedSpans(OtelInstrumentationScope scope) { writer.name("spans").beginArray(); } - private boolean visitSpan(CoreSpan span) { + private void visitSpan(CoreSpan span) { if (!shouldExport(span)) { - return false; + return; } if (currentSpan != null) { // ensure last span written at trace boundary includes sampling tags @@ -141,7 +129,6 @@ private boolean visitSpan(CoreSpan span) { } currentSpan = (DDSpan) span; currentSpanLinks = currentSpan.getLinks(); - return true; } // 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 56a469edb25..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 @@ -44,7 +44,6 @@ public final class OtlpTraceProtoCollector extends OtlpTraceCollector { private int payloadBytes; private int scopedBytes; private int spanBytes; - private int traceCount; private OtelInstrumentationScope currentScope; private DDSpan currentSpan; @@ -57,13 +56,9 @@ public void addTrace(List> spans) { payloadStarted = true; } - boolean exported = false; // OtlpProtoBuffer collects spans in reverse for (int i = spans.size() - 1; i >= 0; i--) { - exported |= visitSpan(spans.get(i)); - } - if (exported) { - traceCount++; + visitSpan(spans.get(i)); } } @@ -83,8 +78,6 @@ public OtlpPayload collectTraces() { /** Prepare temporary elements to collect trace data. */ private void start() { - traceCount = 0; - // remove stale entries from caches OtlpCommonProto.recalibrateCaches(); @@ -92,11 +85,6 @@ private void start() { visitScopedSpans(DEFAULT_TRACE_SCOPE); } - @Override - public int getTraceCount() { - return traceCount; - } - /** Cleanup elements used to collect trace data. */ private void stop() { payloadStarted = false; @@ -119,9 +107,9 @@ private void visitScopedSpans(OtelInstrumentationScope scope) { currentScope = scope; } - private boolean visitSpan(CoreSpan span) { + private void visitSpan(CoreSpan span) { if (!shouldExport(span)) { - return false; + return; } if (currentSpan != null) { // ensure last span written at trace boundary includes sampling tags @@ -133,7 +121,6 @@ private boolean visitSpan(CoreSpan span) { } currentSpan = (DDSpan) span; currentSpan.getLinks().forEach(this::visitSpanLink); - return true; } 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 a1442ae0f87..5441f868bc0 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 @@ -6,8 +6,6 @@ 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.ArgumentMatchers.eq; -import static org.mockito.ArgumentMatchers.same; import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; @@ -15,15 +13,17 @@ 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.monitor.HealthMetrics; import datadog.trace.core.otlp.common.OtlpPayload; import datadog.trace.core.otlp.common.OtlpSender; import datadog.trace.core.otlp.trace.OtlpTraceCollector; 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; @@ -34,13 +34,13 @@ @ExtendWith(MockitoExtension.class) class OtlpPayloadDispatcherTest { @Mock OtlpSender sender; - @Mock HealthMetrics healthMetrics; TestCollector collector = new TestCollector(); @BeforeEach void stubSuccessfulSend() { lenient().when(sender.send(any())).thenReturn(RemoteApi.Response.success(200)); + OtlpTelemetry.getInstance().drain(); } @Test @@ -108,29 +108,41 @@ void emptyTraceForwardsNothing() { } @Test - void flushRecordsSerializedSizeBeforeSending() { + void flushRecordsSuccessfulExportTelemetry() { RemoteApi.Response response = RemoteApi.Response.success(200); when(sender.send(any())).thenReturn(response); - OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector, healthMetrics); + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); dispatcher.addTrace(Arrays.asList(sampledSpan(), sampledSpan())); dispatcher.flush(); - verify(healthMetrics).onSerialize(2 /*spans*/); - verify(healthMetrics).onSend(eq(1) /*trace*/, eq(2) /*spans*/, same(response)); + 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 flushRecordsFailedSendHealthMetric() { + void flushRecordsFailedExportTelemetry() { RemoteApi.Response response = RemoteApi.Response.failed(500); when(sender.send(any())).thenReturn(response); - OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector, healthMetrics); + OtlpPayloadDispatcher dispatcher = new OtlpPayloadDispatcher(sender, collector); dispatcher.addTrace(Arrays.asList(sampledSpan(), sampledSpan())); dispatcher.flush(); - verify(healthMetrics).onSerialize(2 /*spans*/); - verify(healthMetrics).onFailedSend(eq(1) /*trace*/, eq(2) /*spans*/, same(response)); + 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<>(); + for (OtlpTelemetry.OtlpMetric metric : OtlpTelemetry.getInstance().drain()) { + byName.put(metric.metricName, metric); + } + return byName; } @Test @@ -180,24 +192,14 @@ private static CoreSpan unsetSpan() { /** Test collector that creates payloads whose size equals the number of exported spans. */ private static class TestCollector extends OtlpTraceCollector { final List> spansToExport = new ArrayList<>(); - int traceCount; @Override public void addTrace(List> spans) { - if (spansToExport.isEmpty()) { - // starting a new batch - reset the count left over from the last collection - traceCount = 0; - } - boolean exported = false; for (CoreSpan span : spans) { if (shouldExport(span)) { spansToExport.add(span); - exported = true; } } - if (exported) { - traceCount++; - } } @Override @@ -213,10 +215,5 @@ public OtlpPayload collectTraces() { spansToExport.clear(); } } - - @Override - public int getTraceCount() { - return traceCount; - } } } 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 index 0babcd1177c..2f9440b6fd2 100644 --- a/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java +++ b/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java @@ -17,14 +17,28 @@ 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 OtlpTelemetry() {} + public void onTracesExportAttempt() { + tracesExport.attempts.increment(); + } + + public void onTracesExportSuccess() { + tracesExport.successes.increment(); + } + + public void onTracesExportFailure() { + tracesExport.failures.increment(); + } + public void onMetricsExportAttempt() { metricsExport.attempts.increment(); } @@ -57,6 +71,7 @@ public void prepareMetrics() { @Override public Collection drain() { List drained = new ArrayList<>(); + tracesExport.drainInto(drained, tracesTags); metricsExport.drainInto(drained, metricsTags); long logRecordCount = logRecords.sumThenReset(); if (logRecordCount > 0) { 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 index 7c5472328a1..b24b2030d71 100644 --- a/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java +++ b/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java @@ -4,6 +4,8 @@ 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; @@ -15,6 +17,38 @@ void drainStaleMetrics() { collector.drain(); } + @Test + void onTracesExportAttemptQueuesCountMetricWithTags() { + collector.onTracesExportAttempt(); + + 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.onTracesExportSuccess(); + collector.onTracesExportFailure(); + + 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(); @@ -34,17 +68,17 @@ void onMetricsExportAttemptQueuesCountMetricWithTags() { @Test void metricsExportMetricsUseExpectedNames() { + collector.onMetricsExportAttempt(); collector.onMetricsExportAttempt(); collector.onMetricsExportSuccess(); collector.onMetricsExportFailure(); - Collection metrics = collector.drain(); + Map valuesByName = drainToMap(); - assertEquals(3, metrics.size()); - assertTrue(metrics.stream().anyMatch(m -> m.metricName.equals("otel.metrics_export_attempts"))); - assertTrue( - metrics.stream().anyMatch(m -> m.metricName.equals("otel.metrics_export_successes"))); - assertTrue(metrics.stream().anyMatch(m -> m.metricName.equals("otel.metrics_export_failures"))); + 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 @@ -73,4 +107,12 @@ void onLogRecordsSubmittedIgnoresNonPositiveCounts() { 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; + } } From f446803c0f77ea0ba94ff08bedcce2c7afe526c2 Mon Sep 17 00:00:00 2001 From: Stuart McCulloch Date: Thu, 23 Jul 2026 23:12:55 +0100 Subject: [PATCH 3/4] Simplify OTLP telemetry --- .../common/writer/OtlpPayloadDispatcher.java | 6 +----- .../core/otlp/metrics/OtlpMetricsService.java | 6 +----- .../otlp/metrics/OtlpStatsMetricWriter.java | 6 +----- .../trace/api/telemetry/OtlpTelemetry.java | 20 ++++++++----------- .../api/telemetry/OtlpTelemetryTest.java | 8 ++++---- 5 files changed, 15 insertions(+), 31 deletions(-) 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 398f278ba94..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 @@ -29,11 +29,7 @@ public void flush() { if (payload != OtlpPayload.EMPTY) { OtlpTelemetry.getInstance().onTracesExportAttempt(); RemoteApi.Response response = sender.send(payload); - if (response.success()) { - OtlpTelemetry.getInstance().onTracesExportSuccess(); - } else { - OtlpTelemetry.getInstance().onTracesExportFailure(); - } + OtlpTelemetry.getInstance().onTracesExportComplete(response.success()); } } 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 9abdbf59af3..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 @@ -94,11 +94,7 @@ private void export() { if (payload != OtlpPayload.EMPTY) { OtlpTelemetry.getInstance().onMetricsExportAttempt(); RemoteApi.Response response = sender.send(payload); - if (response.success()) { - OtlpTelemetry.getInstance().onMetricsExportSuccess(); - } else { - OtlpTelemetry.getInstance().onMetricsExportFailure(); - } + 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 c30b49a6ab2..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 @@ -171,11 +171,7 @@ public void finishBucket() { if (payload != OtlpPayload.EMPTY) { OtlpTelemetry.getInstance().onMetricsExportAttempt(); RemoteApi.Response response = sender.send(payload); - if (response.success()) { - OtlpTelemetry.getInstance().onMetricsExportSuccess(); - } else { - OtlpTelemetry.getInstance().onMetricsExportFailure(); - } + OtlpTelemetry.getInstance().onMetricsExportComplete(response.success()); } } finally { pending.clear(); 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 index 2f9440b6fd2..7d6da08caec 100644 --- a/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java +++ b/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java @@ -31,24 +31,16 @@ public void onTracesExportAttempt() { tracesExport.attempts.increment(); } - public void onTracesExportSuccess() { - tracesExport.successes.increment(); - } - - public void onTracesExportFailure() { - tracesExport.failures.increment(); + public void onTracesExportComplete(boolean success) { + tracesExport.complete(success); } public void onMetricsExportAttempt() { metricsExport.attempts.increment(); } - public void onMetricsExportSuccess() { - metricsExport.successes.increment(); - } - - public void onMetricsExportFailure() { - metricsExport.failures.increment(); + public void onMetricsExportComplete(boolean success) { + metricsExport.complete(success); } public void onLogRecordsSubmitted(long count) { @@ -96,6 +88,10 @@ private static final class ExportCounters { this.failuresMetric = "otel." + signal + "_export_failures"; } + void complete(boolean success) { + (success ? successes : failures).increment(); + } + void drainInto(List out, String[] tags) { addIfNonZero(out, attemptsMetric, attempts, tags); addIfNonZero(out, successesMetric, successes, 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 index b24b2030d71..790bba0d82f 100644 --- a/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java +++ b/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java @@ -38,8 +38,8 @@ void onTracesExportAttemptQueuesCountMetricWithTags() { void tracesExportMetricsUseExpectedNames() { collector.onTracesExportAttempt(); collector.onTracesExportAttempt(); - collector.onTracesExportSuccess(); - collector.onTracesExportFailure(); + collector.onTracesExportComplete(true); + collector.onTracesExportComplete(false); Map valuesByName = drainToMap(); @@ -70,8 +70,8 @@ void onMetricsExportAttemptQueuesCountMetricWithTags() { void metricsExportMetricsUseExpectedNames() { collector.onMetricsExportAttempt(); collector.onMetricsExportAttempt(); - collector.onMetricsExportSuccess(); - collector.onMetricsExportFailure(); + collector.onMetricsExportComplete(true); + collector.onMetricsExportComplete(false); Map valuesByName = drainToMap(); From 66cce1632cabe12fb15bec2c24b403e1314b9aaf Mon Sep 17 00:00:00 2001 From: Stuart McCulloch Date: Thu, 23 Jul 2026 23:45:07 +0100 Subject: [PATCH 4/4] Stage telemetry metrics separate to collection --- .../writer/OtlpPayloadDispatcherTest.java | 2 ++ .../trace/api/telemetry/OtlpTelemetry.java | 28 ++++++++++++------- .../api/telemetry/OtlpTelemetryTest.java | 6 ++++ 3 files changed, 26 insertions(+), 10 deletions(-) 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 5441f868bc0..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 @@ -40,6 +40,7 @@ class OtlpPayloadDispatcherTest { @BeforeEach void stubSuccessfulSend() { lenient().when(sender.send(any())).thenReturn(RemoteApi.Response.success(200)); + OtlpTelemetry.getInstance().prepareMetrics(); OtlpTelemetry.getInstance().drain(); } @@ -139,6 +140,7 @@ void flushRecordsFailedExportTelemetry() { private static Map drainTracesTelemetry() { Map byName = new HashMap<>(); + OtlpTelemetry.getInstance().prepareMetrics(); for (OtlpTelemetry.OtlpMetric metric : OtlpTelemetry.getInstance().drain()) { byName.put(metric.metricName, metric); } 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 index 7d6da08caec..a36bcff67c8 100644 --- a/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java +++ b/internal-api/src/main/java/datadog/trace/api/telemetry/OtlpTelemetry.java @@ -4,7 +4,10 @@ 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. */ @@ -25,6 +28,8 @@ public static OtlpTelemetry getInstance() { 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() { @@ -57,18 +62,21 @@ private static String[] tagsFor(OtlpConfig.Protocol protocol) { @Override public void prepareMetrics() { - // metrics are accumulated directly as they happen; nothing to prepare + 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() { - List drained = new ArrayList<>(); - tracesExport.drainInto(drained, tracesTags); - metricsExport.drainInto(drained, metricsTags); - long logRecordCount = logRecords.sumThenReset(); - if (logRecordCount > 0) { - drained.add(new OtlpMetric("otel.log_records", logRecordCount, logsTags)); + if (telemetryQueue.isEmpty()) { + return Collections.emptyList(); } + List drained = new ArrayList<>(telemetryQueue.size()); + telemetryQueue.drainTo(drained); return drained; } @@ -92,17 +100,17 @@ void complete(boolean success) { (success ? successes : failures).increment(); } - void drainInto(List out, String[] tags) { + 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( - List out, String metricName, LongAdder counter, String[] tags) { + BlockingQueue out, String metricName, LongAdder counter, String[] tags) { long value = counter.sumThenReset(); if (value > 0) { - out.add(new OtlpMetric(metricName, value, tags)); + out.offer(new OtlpMetric(metricName, 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 index 790bba0d82f..53877bd8dc6 100644 --- a/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java +++ b/internal-api/src/test/java/datadog/trace/api/telemetry/OtlpTelemetryTest.java @@ -14,12 +14,14 @@ class OtlpTelemetryTest { @BeforeEach void drainStaleMetrics() { + collector.prepareMetrics(); collector.drain(); } @Test void onTracesExportAttemptQueuesCountMetricWithTags() { collector.onTracesExportAttempt(); + collector.prepareMetrics(); Collection metrics = collector.drain(); @@ -40,6 +42,7 @@ void tracesExportMetricsUseExpectedNames() { collector.onTracesExportAttempt(); collector.onTracesExportComplete(true); collector.onTracesExportComplete(false); + collector.prepareMetrics(); Map valuesByName = drainToMap(); @@ -52,6 +55,7 @@ void tracesExportMetricsUseExpectedNames() { @Test void onMetricsExportAttemptQueuesCountMetricWithTags() { collector.onMetricsExportAttempt(); + collector.prepareMetrics(); Collection metrics = collector.drain(); @@ -72,6 +76,7 @@ void metricsExportMetricsUseExpectedNames() { collector.onMetricsExportAttempt(); collector.onMetricsExportComplete(true); collector.onMetricsExportComplete(false); + collector.prepareMetrics(); Map valuesByName = drainToMap(); @@ -84,6 +89,7 @@ void metricsExportMetricsUseExpectedNames() { @Test void onLogRecordsSubmittedQueuesCountWithGivenValue() { collector.onLogRecordsSubmitted(5); + collector.prepareMetrics(); Collection metrics = collector.drain();