diff --git a/internal-api/src/main/java/datadog/trace/api/telemetry/LLMObsMetricCollector.java b/internal-api/src/main/java/datadog/trace/api/telemetry/LLMObsMetricCollector.java index f43d92cb741..4da5a62b0b9 100644 --- a/internal-api/src/main/java/datadog/trace/api/telemetry/LLMObsMetricCollector.java +++ b/internal-api/src/main/java/datadog/trace/api/telemetry/LLMObsMetricCollector.java @@ -7,11 +7,24 @@ import java.util.Collection; import java.util.Collections; import java.util.List; +import java.util.Map; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.LongAdder; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +/** + * Collects telemetry metrics for LLM Observability spans. + * + *

Counts are aggregated per tag combination in-process and emitted as one point per metrics + * interval. Emitting one point of value 1 per span instead would under-report badly: points are + * timestamped at second granularity, so every point a series produces within the same second + * collapses to a single value at the metrics intake, capping the reported rate at roughly 1/s per + * series regardless of the real span rate. This matches how {@link CoreMetricCollector} and the + * other tracers (dd-trace-py, dd-trace-js) report counts. + */ public final class LLMObsMetricCollector implements MetricCollector { private static final String METRIC_NAMESPACE = "mlobs"; @@ -35,19 +48,38 @@ public static LLMObsMetricCollector get() { private static final String HAS_SESSION_ID_TRUE = "has_session_id:1"; private static final String HAS_SESSION_ID_FALSE = "has_session_id:0"; + /** + * Upper bound on the number of distinct tag combinations tracked concurrently. Tag values are + * drawn from bounded sets (a handful of integrations and span kinds, plus four booleans), so this + * is only a guard against an unexpected high-cardinality source. It is also kept low enough that + * several {@link #prepareMetrics()} intervals can be staged without overflowing {@link + * MetricCollector#RAW_QUEUE_SIZE} before the next {@link #drain()}. + */ + static final int MAX_TAG_COMBINATIONS = 128; + private final BlockingQueue metricsQueue; private final DDCache integrationTagCache; private final DDCache spanKindTagCache; + /** + * Counter per tag combination, aggregated in-process and flushed once per metrics interval by + * {@link #prepareMetrics()}. Counting here rather than enqueuing one entry per span is what keeps + * the reported count accurate at high span rates. + */ + private final ConcurrentHashMap, LongAdder> spanFinishedCounters; private LLMObsMetricCollector() { this.metricsQueue = new ArrayBlockingQueue<>(RAW_QUEUE_SIZE); this.integrationTagCache = DDCaches.newFixedSizeCache(8); this.spanKindTagCache = DDCaches.newFixedSizeCache(8); + this.spanFinishedCounters = new ConcurrentHashMap<>(); } /** * Record a span finished metric for LLMObs telemetry. * + *

This only increments an in-process counter. The counter is converted into a single telemetry + * metric per tag combination by {@link #prepareMetrics()}, once per metrics interval. + * * @param integration the integration name (e.g., "openai") * @param spanKind the span kind (e.g., "llm", "embedding") * @param isRootSpan whether this is a root span @@ -74,16 +106,49 @@ public void recordSpanFinished( isAutoInstrumented ? AUTOINSTRUMENTED_TRUE : AUTOINSTRUMENTED_FALSE, hasError ? ERROR_TRUE : ERROR_FALSE, hasSessionId ? HAS_SESSION_ID_TRUE : HAS_SESSION_ID_FALSE); - LLMObsMetric metric = - new LLMObsMetric(METRIC_NAMESPACE, true, SPAN_FINISHED_METRIC, COUNT_METRIC_TYPE, 1L, tags); - if (!metricsQueue.offer(metric)) { - log.debug("Unable to add telemetry metric {} for {}", SPAN_FINISHED_METRIC, integration); + + LongAdder counter = spanFinishedCounters.get(tags); + if (counter == null) { + // Soft bound: concurrent recorders may overshoot slightly, which is fine for a guard. + if (spanFinishedCounters.size() >= MAX_TAG_COMBINATIONS) { + log.debug( + "Dropping telemetry metric {} for {}: tag combination limit ({}) reached", + SPAN_FINISHED_METRIC, + integration, + MAX_TAG_COMBINATIONS); + return; + } + counter = spanFinishedCounters.computeIfAbsent(tags, key -> new LongAdder()); } + counter.increment(); } @Override public void prepareMetrics() { - // metrics are added directly via recordSpanFinished; no additional preparation needed + // Entries are never removed: a recorder thread may already hold a reference to a LongAdder, so + // removing it here would silently drop a concurrent increment. Tag values come from bounded + // sets, so retaining idle combinations costs at most MAX_TAG_COMBINATIONS entries. + for (Map.Entry, LongAdder> entry : spanFinishedCounters.entrySet()) { + long value = entry.getValue().sumThenReset(); + if (value == 0) { + continue; + } + LLMObsMetric metric = + new LLMObsMetric( + METRIC_NAMESPACE, + true, + SPAN_FINISHED_METRIC, + COUNT_METRIC_TYPE, + value, + entry.getKey()); + if (!metricsQueue.offer(metric)) { + // Queue is full; give the count back to the counter so it is reported in a later interval + // instead of being lost, and stop staging for now. + entry.getValue().add(value); + log.debug("Unable to add telemetry metric {}: queue is full", SPAN_FINISHED_METRIC); + break; + } + } } @Override @@ -96,6 +161,12 @@ public Collection drain() { return drained; } + /** Clears all staged counters and metrics. Visible for testing only. */ + public void resetForTesting() { + spanFinishedCounters.clear(); + metricsQueue.clear(); + } + public static class LLMObsMetric extends MetricCollector.Metric { public LLMObsMetric( String namespace, diff --git a/internal-api/src/test/groovy/datadog/trace/api/telemetry/LLMObsMetricCollectorTest.groovy b/internal-api/src/test/groovy/datadog/trace/api/telemetry/LLMObsMetricCollectorTest.groovy deleted file mode 100644 index e3927c03900..00000000000 --- a/internal-api/src/test/groovy/datadog/trace/api/telemetry/LLMObsMetricCollectorTest.groovy +++ /dev/null @@ -1,75 +0,0 @@ -package datadog.trace.api.telemetry - -import datadog.trace.test.util.DDSpecification - -class LLMObsMetricCollectorTest extends DDSpecification { - LLMObsMetricCollector collector = LLMObsMetricCollector.get() - - void setup() { - // clear any previous metrics - collector.drain() - } - - def "no metrics - drain empty list"() { - when: - collector.prepareMetrics() - - then: - collector.drain().isEmpty() - } - - def "record and drain span finished metrics"() { - when: - collector.recordSpanFinished("openai", "llm", true, true, false, false) - collector.recordSpanFinished("openai", "llm", false, true, false, true) - collector.recordSpanFinished("anthropic", "embedding", true, false, true, false) - collector.prepareMetrics() - def metrics = collector.drain() - - then: - metrics.size() == 3 - - def metric1 = metrics[0] - metric1.type == 'count' - metric1.value == 1 - metric1.namespace == 'mlobs' - metric1.metricName == 'span.finished' - metric1.tags.sort() == [ - 'integration:openai', - 'span_kind:llm', - 'is_root_span:1', - 'autoinstrumented:1', - 'error:0', - 'has_session_id:0' - ].sort() - - def metric2 = metrics[1] - metric2.type == 'count' - metric2.value == 1 - metric2.namespace == 'mlobs' - metric2.metricName == 'span.finished' - metric2.tags.toSet() == [ - 'integration:openai', - 'span_kind:llm', - 'is_root_span:0', - 'autoinstrumented:1', - 'error:0', - 'has_session_id:1' - ].toSet() - - def metric3 = metrics[2] - metric3.type == 'count' - metric3.value == 1 - metric3.namespace == 'mlobs' - metric3.metricName == 'span.finished' - metric3.tags.toSet() == [ - 'integration:anthropic', - 'span_kind:embedding', - 'is_root_span:1', - 'autoinstrumented:0', - 'error:1', - 'has_session_id:0' - ].toSet() - } -} - diff --git a/internal-api/src/test/java/datadog/trace/api/telemetry/LLMObsMetricCollectorTest.java b/internal-api/src/test/java/datadog/trace/api/telemetry/LLMObsMetricCollectorTest.java new file mode 100644 index 00000000000..efce543cee3 --- /dev/null +++ b/internal-api/src/test/java/datadog/trace/api/telemetry/LLMObsMetricCollectorTest.java @@ -0,0 +1,167 @@ +package datadog.trace.api.telemetry; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.HashSet; +import java.util.List; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +class LLMObsMetricCollectorTest { + private final LLMObsMetricCollector collector = LLMObsMetricCollector.get(); + + @BeforeEach + void clearStaleMetrics() { + collector.resetForTesting(); + } + + @AfterEach + void clearMetrics() { + collector.resetForTesting(); + } + + @Test + void noMetricsDrainsEmptyList() { + collector.prepareMetrics(); + + assertTrue(collector.drain().isEmpty()); + } + + @Test + void recordsOneMetricPerDistinctTagCombination() { + collector.recordSpanFinished("openai", "llm", true, true, false, false); + collector.recordSpanFinished("openai", "llm", false, true, false, true); + collector.recordSpanFinished("anthropic", "embedding", true, false, true, false); + collector.prepareMetrics(); + + List metrics = sorted(collector.drain()); + + assertEquals(3, metrics.size()); + for (LLMObsMetricCollector.LLMObsMetric metric : metrics) { + assertEquals("mlobs", metric.namespace); + assertEquals("span.finished", metric.metricName); + assertEquals("count", metric.type); + assertEquals(1L, metric.value); + } + assertEquals( + new HashSet<>( + Arrays.asList( + new HashSet<>( + Arrays.asList( + "integration:openai", + "span_kind:llm", + "is_root_span:1", + "autoinstrumented:1", + "error:0", + "has_session_id:0")), + new HashSet<>( + Arrays.asList( + "integration:openai", + "span_kind:llm", + "is_root_span:0", + "autoinstrumented:1", + "error:0", + "has_session_id:1")), + new HashSet<>( + Arrays.asList( + "integration:anthropic", + "span_kind:embedding", + "is_root_span:1", + "autoinstrumented:0", + "error:1", + "has_session_id:0")))), + tagSets(metrics)); + } + + @Test + void aggregatesIdenticalTagCombinationsIntoASingleCount() { + collector.recordSpanFinished("openai", "llm", true, true, false, false); + collector.recordSpanFinished("openai", "llm", true, true, false, false); + collector.recordSpanFinished("openai", "llm", true, true, false, false); + collector.prepareMetrics(); + + Collection metrics = collector.drain(); + + assertEquals(1, metrics.size()); + LLMObsMetricCollector.LLMObsMetric metric = metrics.iterator().next(); + assertEquals("count", metric.type); + assertEquals(3L, metric.value); + } + + /** + * Regression test for the reported ~100x under-count: emitting one raw metric per span both + * overflowed the bounded raw queue and collapsed to ~1/s per series at the metrics intake, since + * points are timestamped at second granularity. The count must survive well past {@link + * MetricCollector#RAW_QUEUE_SIZE} spans in a single interval. + */ + @Test + void reportsExactCountWellBeyondRawQueueSize() { + int spans = MetricCollector.RAW_QUEUE_SIZE * 10; + for (int i = 0; i < spans; i++) { + collector.recordSpanFinished("openai", "llm", false, true, false, false); + } + collector.prepareMetrics(); + + Collection metrics = collector.drain(); + + assertEquals(1, metrics.size()); + assertEquals((long) spans, metrics.iterator().next().value); + } + + @Test + void countersResetBetweenIntervals() { + collector.recordSpanFinished("openai", "llm", true, true, false, false); + collector.prepareMetrics(); + assertEquals(1, collector.drain().size()); + + collector.prepareMetrics(); + + assertTrue(collector.drain().isEmpty(), "an idle interval must not re-report a stale count"); + } + + @Test + void reportsCountsAccumulatedAcrossIntervalsWithoutADrain() { + collector.recordSpanFinished("openai", "llm", true, true, false, false); + collector.prepareMetrics(); + collector.recordSpanFinished("openai", "llm", true, true, false, false); + collector.recordSpanFinished("openai", "llm", true, true, false, false); + collector.prepareMetrics(); + + List metrics = sorted(collector.drain()); + + assertEquals(2, metrics.size()); + assertEquals(1L, metrics.get(0).value); + assertEquals(2L, metrics.get(1).value); + } + + @Test + void boundsTheNumberOfTrackedTagCombinations() { + for (int i = 0; i < LLMObsMetricCollector.MAX_TAG_COMBINATIONS * 2; i++) { + collector.recordSpanFinished("integration-" + i, "llm", true, true, false, false); + } + collector.prepareMetrics(); + + assertEquals(LLMObsMetricCollector.MAX_TAG_COMBINATIONS, collector.drain().size()); + } + + private static List sorted( + Collection metrics) { + List sorted = new ArrayList<>(metrics); + sorted.sort((a, b) -> Long.compare(a.value.longValue(), b.value.longValue())); + return sorted; + } + + private static HashSet> tagSets( + Collection metrics) { + HashSet> tagSets = new HashSet<>(); + for (LLMObsMetricCollector.LLMObsMetric metric : metrics) { + tagSets.add(new HashSet<>(metric.tags)); + } + return tagSets; + } +} diff --git a/telemetry/src/test/groovy/datadog/telemetry/metric/LLMObsMetricPeriodicActionTest.groovy b/telemetry/src/test/groovy/datadog/telemetry/metric/LLMObsMetricPeriodicActionTest.groovy deleted file mode 100644 index 887cf766280..00000000000 --- a/telemetry/src/test/groovy/datadog/telemetry/metric/LLMObsMetricPeriodicActionTest.groovy +++ /dev/null @@ -1,90 +0,0 @@ -package datadog.telemetry.metric - -import datadog.telemetry.TelemetryService -import datadog.telemetry.api.Metric -import datadog.trace.api.telemetry.LLMObsMetricCollector -import datadog.trace.test.util.DDSpecification - -class LLMObsMetricPeriodicActionTest extends DDSpecification { - LLMObsMetricPeriodicAction periodicAction = new LLMObsMetricPeriodicAction() - TelemetryService telemetryService = Mock() - LLMObsMetricCollector collector = LLMObsMetricCollector.get() - - void setup() { - // clear any previous metrics - collector.drain() - } - - void 'test multiple span finished metrics with different tags'() { - when: - collector.recordSpanFinished('openai', 'llm', true, true, false, true) - collector.recordSpanFinished('openai', 'llm', false, true, false, false) - collector.recordSpanFinished('anthropic', 'embedding', true, false, true, false) - periodicAction.doIteration(telemetryService) - - then: - 1 * telemetryService.addMetric({ Metric metric -> - metric.namespace == 'mlobs' && - metric.metric == 'span.finished' && - metric.tags.toSet() == [ - 'integration:openai', - 'span_kind:llm', - 'is_root_span:1', - 'autoinstrumented:1', - 'error:0', - 'has_session_id:1' - ].toSet() - }) - 1 * telemetryService.addMetric({ Metric metric -> - metric.namespace == 'mlobs' && - metric.metric == 'span.finished' && - metric.tags.toSet() == [ - 'integration:openai', - 'span_kind:llm', - 'is_root_span:0', - 'autoinstrumented:1', - 'error:0', - 'has_session_id:0' - ].toSet() - }) - 1 * telemetryService.addMetric({ Metric metric -> - metric.namespace == 'mlobs' && - metric.metric == 'span.finished' && - metric.tags.toSet() == [ - 'integration:anthropic', - 'span_kind:embedding', - 'is_root_span:1', - 'autoinstrumented:0', - 'error:1', - 'has_session_id:0' - ].toSet() - }) - 0 * _ - } - - void 'test aggregation of identical metrics'() { - when: - collector.recordSpanFinished('openai', 'llm', true, true, false, false) - collector.recordSpanFinished('openai', 'llm', true, true, false, false) - collector.recordSpanFinished('openai', 'llm', true, true, false, false) - periodicAction.doIteration(telemetryService) - - then: - 1 * telemetryService.addMetric({ Metric metric -> - metric.namespace == 'mlobs' && - metric.metric == 'span.finished' && - metric.points.size() == 3 && - metric.points.every { it[1] == 1 } && - metric.tags.toSet() == [ - 'integration:openai', - 'span_kind:llm', - 'is_root_span:1', - 'autoinstrumented:1', - 'error:0', - 'has_session_id:0' - ].toSet() - }) - 0 * _ - } -} - diff --git a/telemetry/src/test/java/datadog/telemetry/metric/LLMObsMetricPeriodicActionTest.java b/telemetry/src/test/java/datadog/telemetry/metric/LLMObsMetricPeriodicActionTest.java new file mode 100644 index 00000000000..271bd4a2cd8 --- /dev/null +++ b/telemetry/src/test/java/datadog/telemetry/metric/LLMObsMetricPeriodicActionTest.java @@ -0,0 +1,120 @@ +package datadog.telemetry.metric; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.ArgumentCaptor.forClass; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; + +import datadog.telemetry.TelemetryService; +import datadog.telemetry.api.Metric; +import datadog.trace.api.telemetry.LLMObsMetricCollector; +import java.util.Arrays; +import java.util.HashSet; +import java.util.List; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +class LLMObsMetricPeriodicActionTest { + private final LLMObsMetricPeriodicAction periodicAction = new LLMObsMetricPeriodicAction(); + private final LLMObsMetricCollector collector = LLMObsMetricCollector.get(); + private TelemetryService telemetryService; + + @BeforeEach + void setUp() { + collector.resetForTesting(); + telemetryService = mock(TelemetryService.class); + } + + @AfterEach + void tearDown() { + collector.resetForTesting(); + } + + @Test + void emitsOneMetricPerDistinctTagCombination() { + collector.recordSpanFinished("openai", "llm", true, true, false, true); + collector.recordSpanFinished("openai", "llm", false, true, false, false); + collector.recordSpanFinished("anthropic", "embedding", true, false, true, false); + + collector.prepareMetrics(); + periodicAction.doIteration(telemetryService); + + ArgumentCaptor captor = forClass(Metric.class); + verify(telemetryService, times(3)).addMetric(captor.capture()); + verifyNoMoreInteractions(telemetryService); + + HashSet> tagSets = new HashSet<>(); + for (Metric metric : captor.getAllValues()) { + assertEquals("mlobs", metric.getNamespace()); + assertEquals("span.finished", metric.getMetric()); + assertEquals(1, metric.getPoints().size()); + assertEquals(1L, metric.getPoints().get(0).get(1).longValue()); + tagSets.add(new HashSet<>(metric.getTags())); + } + assertEquals( + new HashSet<>( + Arrays.asList( + new HashSet<>( + Arrays.asList( + "integration:openai", + "span_kind:llm", + "is_root_span:1", + "autoinstrumented:1", + "error:0", + "has_session_id:1")), + new HashSet<>( + Arrays.asList( + "integration:openai", + "span_kind:llm", + "is_root_span:0", + "autoinstrumented:1", + "error:0", + "has_session_id:0")), + new HashSet<>( + Arrays.asList( + "integration:anthropic", + "span_kind:embedding", + "is_root_span:1", + "autoinstrumented:0", + "error:1", + "has_session_id:0")))), + tagSets); + } + + /** + * Identical spans must produce a single point carrying the summed count. Emitting one point of + * value 1 per span instead loses all but one of the points that share a second-granularity + * timestamp once they reach the metrics intake. + */ + @Test + void emitsASinglePointCarryingTheAggregatedCount() { + for (int i = 0; i < 5000; i++) { + collector.recordSpanFinished("openai", "llm", true, true, false, false); + } + + collector.prepareMetrics(); + periodicAction.doIteration(telemetryService); + + ArgumentCaptor captor = forClass(Metric.class); + verify(telemetryService).addMetric(captor.capture()); + verifyNoMoreInteractions(telemetryService); + + Metric metric = captor.getValue(); + assertEquals(Metric.TypeEnum.COUNT, metric.getType()); + List> points = metric.getPoints(); + assertEquals(1, points.size()); + assertEquals(5000L, points.get(0).get(1).longValue()); + } + + @Test + void emitsNothingWhenNoSpansFinished() { + collector.prepareMetrics(); + periodicAction.doIteration(telemetryService); + + verifyNoMoreInteractions(telemetryService); + } +}