From 9ad29ad96d248ec76a2782d6ec55c9d9f9749b1e Mon Sep 17 00:00:00 2001 From: Sarah Chen Date: Mon, 28 Sep 2026 09:39:49 -0400 Subject: [PATCH 1/3] Decouple flag eval metrics from CoreMetricCollector --- .../api/telemetry/CoreMetricCollector.java | 8 - .../telemetry/CoreMetricCollectorTest.groovy | 20 --- .../build.gradle.kts | 2 + .../gradle.lockfile | 3 +- .../flagevaluation/FlagEvaluationMetrics.java | 87 +++++++++++ .../FlagEvaluationMetricsTest.java | 110 +++++++++++++ .../featureflag/FlagEvaluationWriterImpl.java | 77 +++------ .../FlagEvaluationTestSupport.java | 24 +-- .../FlagEvaluationWriterImplTest.java | 114 ++++---------- telemetry/build.gradle.kts | 1 + .../datadog/telemetry/TelemetrySystem.java | 4 + .../FlagEvaluationMetricPeriodicAction.java | 55 +++++++ ...lagEvaluationMetricPeriodicActionTest.java | 146 ++++++++++++++++++ 13 files changed, 467 insertions(+), 184 deletions(-) create mode 100644 products/feature-flagging/feature-flagging-bootstrap/src/main/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetrics.java create mode 100644 products/feature-flagging/feature-flagging-bootstrap/src/test/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetricsTest.java create mode 100644 telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java create mode 100644 telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java diff --git a/internal-api/src/main/java/datadog/trace/api/telemetry/CoreMetricCollector.java b/internal-api/src/main/java/datadog/trace/api/telemetry/CoreMetricCollector.java index 9c8074dee70..f0bbf3c93fa 100644 --- a/internal-api/src/main/java/datadog/trace/api/telemetry/CoreMetricCollector.java +++ b/internal-api/src/main/java/datadog/trace/api/telemetry/CoreMetricCollector.java @@ -31,14 +31,6 @@ private CoreMetricCollector() { this.metricsQueue = new ArrayBlockingQueue<>(RAW_QUEUE_SIZE); } - public void count(String metricName, long value, String tag) { - if (value <= 0) { - return; - } - this.metricsQueue.offer( - new CoreMetric(METRIC_NAMESPACE, true, metricName, "count", value, tag)); - } - @Override public void prepareMetrics() { // Collect the bounded, high-value client-side trace-stats span-collapse counters first, tagged diff --git a/internal-api/src/test/groovy/datadog/trace/api/telemetry/CoreMetricCollectorTest.groovy b/internal-api/src/test/groovy/datadog/trace/api/telemetry/CoreMetricCollectorTest.groovy index 5d0b920bfb9..205da2bd0f5 100644 --- a/internal-api/src/test/groovy/datadog/trace/api/telemetry/CoreMetricCollectorTest.groovy +++ b/internal-api/src/test/groovy/datadog/trace/api/telemetry/CoreMetricCollectorTest.groovy @@ -53,24 +53,4 @@ class CoreMetricCollectorTest extends DDSpecification { collector.prepareMetrics() collector.drain().size() == limit } - - def "direct count core metric"() { - setup: - def collector = CoreMetricCollector.getInstance() - collector.drain() - - when: - collector.count('flagevaluation.rows.dropped', 3, 'reason:queue_overflow') - def metrics = collector.drain() - - then: - metrics.size() == 1 - - def metric = metrics[0] - metric.type == 'count' - metric.value == 3 - metric.namespace == 'tracers' - metric.metricName == 'flagevaluation.rows.dropped' - metric.tags == ['reason:queue_overflow'] - } } diff --git a/products/feature-flagging/feature-flagging-bootstrap/build.gradle.kts b/products/feature-flagging/feature-flagging-bootstrap/build.gradle.kts index 848b16a2841..9c0add8f901 100644 --- a/products/feature-flagging/feature-flagging-bootstrap/build.gradle.kts +++ b/products/feature-flagging/feature-flagging-bootstrap/build.gradle.kts @@ -29,6 +29,8 @@ extra["excludedClassesCoverage"] = listOf( ) dependencies { + implementation(project(":products:metrics:metrics-api")) + testImplementation(libs.bundles.junit5) testImplementation(libs.bundles.mockito) testImplementation(project(":utils:test-utils")) diff --git a/products/feature-flagging/feature-flagging-bootstrap/gradle.lockfile b/products/feature-flagging/feature-flagging-bootstrap/gradle.lockfile index 4e8cd8dbd72..66f40fe18c8 100644 --- a/products/feature-flagging/feature-flagging-bootstrap/gradle.lockfile +++ b/products/feature-flagging/feature-flagging-bootstrap/gradle.lockfile @@ -68,6 +68,7 @@ org.ow2.asm:asm:9.10.1=jacocoAnt,spotbugs org.slf4j:jcl-over-slf4j:1.7.30=testCompileClasspath,testRuntimeClasspath org.slf4j:jul-to-slf4j:1.7.30=testCompileClasspath,testRuntimeClasspath org.slf4j:log4j-over-slf4j:1.7.30=testCompileClasspath,testRuntimeClasspath +org.slf4j:slf4j-api:1.7.30=runtimeClasspath org.slf4j:slf4j-api:1.7.32=testCompileClasspath,testRuntimeClasspath org.slf4j:slf4j-api:2.0.17=spotbugsSlf4j org.slf4j:slf4j-api:2.0.18=spotbugs @@ -78,4 +79,4 @@ org.spockframework:spock-core:2.4-groovy-3.0=testCompileClasspath,testRuntimeCla org.tabletest:tabletest-junit:1.2.2=testCompileClasspath,testRuntimeClasspath org.tabletest:tabletest-parser:1.2.1=testCompileClasspath,testRuntimeClasspath org.xmlresolver:xmlresolver:5.3.3=spotbugs -empty=annotationProcessor,runtimeClasspath,spotbugsPlugins,testAnnotationProcessor +empty=annotationProcessor,spotbugsPlugins,testAnnotationProcessor diff --git a/products/feature-flagging/feature-flagging-bootstrap/src/main/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetrics.java b/products/feature-flagging/feature-flagging-bootstrap/src/main/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetrics.java new file mode 100644 index 00000000000..ad6631ac48e --- /dev/null +++ b/products/feature-flagging/feature-flagging-bootstrap/src/main/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetrics.java @@ -0,0 +1,87 @@ +package datadog.trace.api.featureflag.flagevaluation; + +import datadog.metrics.api.Accumulator; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicLong; + +/** Flag evaluation counters owned by the producer and drained periodically by telemetry. */ +public final class FlagEvaluationMetrics { + private static final FlagEvaluationMetrics INSTANCE = new FlagEvaluationMetrics(); + + private final Accumulator counts = Accumulator.of(Metric.class); + private final ConcurrentHashMap contextTruncations = + new ConcurrentHashMap<>(); + + private FlagEvaluationMetrics() {} + + public static FlagEvaluationMetrics getInstance() { + return INSTANCE; + } + + public enum Metric { + DROPPED_QUEUE_OVERFLOW("flagevaluation.rows.dropped", "queue_overflow"), + DROPPED_CLOSED("flagevaluation.rows.dropped", "closed"), + DROPPED_DEGRADED_CAP("flagevaluation.rows.dropped", "degraded_cap"), + DROPPED_PAYLOAD_LIMIT("flagevaluation.rows.dropped", "payload_limit"), + DEGRADED_CARDINALITY_CAP("flagevaluation.rows.degraded", "cardinality_cap"), + DEGRADED_PAYLOAD_LIMIT("flagevaluation.rows.degraded", "payload_limit"), + PAYLOAD_SPLITS("flagevaluation.payload.splits", null); + + private final String name; + private final String tag; + + Metric(String name, String reason) { + this.name = name; + this.tag = reason == null ? null : "reason:" + reason; + } + } + + public void count(Metric metric, long value) { + if (value > 0) { + counts.add(metric, value); + } + } + + /** Records the hook's canonical, comma-separated combination of truncation reasons. */ + public void countContextTruncated(String reason, long value) { + if (value > 0) { + contextTruncations.computeIfAbsent(reason, key -> new AtomicLong()).addAndGet(value); + } + } + + /** Returns nonzero deltas, resetting only the counts included in this snapshot. */ + public List drain() { + List drained = new ArrayList<>(); + Accumulator.Counts snapshot = counts.accumulateAndReset(); + for (Metric metric : snapshot.keys()) { + long value = snapshot.get(metric); + if (value > 0) { + drained.add(new Count(metric.name, metric.tag, value)); + } + } + for (Map.Entry entry : contextTruncations.entrySet()) { + long value = entry.getValue().getAndSet(0); + if (value > 0) { + drained.add( + new Count("flagevaluation.context.truncated", "reason:" + entry.getKey(), value)); + } + } + return drained; + } + + /** A counter delta with its metric name and optional reason tag, independent of transport. */ + public static final class Count { + public final String name; + public final String tag; + public final long value; + + private Count(String name, String tag, long value) { + this.name = name; + this.tag = tag; + this.value = value; + } + } +} diff --git a/products/feature-flagging/feature-flagging-bootstrap/src/test/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetricsTest.java b/products/feature-flagging/feature-flagging-bootstrap/src/test/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetricsTest.java new file mode 100644 index 00000000000..c83e3c0215a --- /dev/null +++ b/products/feature-flagging/feature-flagging-bootstrap/src/test/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetricsTest.java @@ -0,0 +1,110 @@ +package datadog.trace.api.featureflag.flagevaluation; + +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_CLOSED; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Count; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +class FlagEvaluationMetricsTest { + private final FlagEvaluationMetrics metrics = FlagEvaluationMetrics.getInstance(); + + @BeforeEach + @AfterEach + void resetMetrics() { + metrics.drain(); + } + + @Test + void accumulatesPositiveLongCountsAndDrainsOnlyOnce() { + metrics.count(DROPPED_CLOSED, Integer.MAX_VALUE); + metrics.count(DROPPED_CLOSED, 5); + metrics.count(DROPPED_CLOSED, 0); + metrics.count(DROPPED_CLOSED, -1); + + List counts = metrics.drain(); + assertEquals(1, counts.size()); + assertEquals("flagevaluation.rows.dropped", counts.get(0).name); + assertEquals("reason:closed", counts.get(0).tag); + assertEquals(Integer.MAX_VALUE + 5L, counts.get(0).value); + assertTrue(metrics.drain().isEmpty()); + + metrics.count(DROPPED_CLOSED, 2); + assertEquals(2, metrics.drain().get(0).value); + } + + @Test + void preservesCombinedTruncationReasonsAcrossDrains() { + metrics.countContextTruncated("max_key_length,max_value_length", 2); + metrics.countContextTruncated("max_key_length,max_value_length", 3); + metrics.countContextTruncated("max_depth", 1); + metrics.countContextTruncated("ignored", 0); + metrics.countContextTruncated("ignored", -1); + + Map counts = new HashMap<>(); + for (Count count : metrics.drain()) { + assertEquals("flagevaluation.context.truncated", count.name); + counts.put(count.tag, count.value); + } + assertEquals(2, counts.size()); + assertEquals(5L, counts.get("reason:max_key_length,max_value_length").longValue()); + assertEquals(1L, counts.get("reason:max_depth").longValue()); + assertTrue(metrics.drain().isEmpty()); + + metrics.countContextTruncated("max_depth", 4); + assertEquals(4, metrics.drain().get(0).value); + } + + @Test + void concurrentProducersAndDrainingPreserveAllCounts() throws Exception { + ExecutorService executor = Executors.newFixedThreadPool(4); + CountDownLatch start = new CountDownLatch(1); + List> producers = new ArrayList<>(); + try { + for (int producer = 0; producer < 4; producer++) { + producers.add( + executor.submit( + () -> { + start.await(); + for (int i = 0; i < 5000; i++) { + metrics.count(DROPPED_CLOSED, 1); + metrics.countContextTruncated("max_depth", 1); + } + return null; + })); + } + start.countDown(); + Map totals = new HashMap<>(); + for (int i = 0; i < 100; i++) { + addTo(totals, metrics.drain()); + } + for (Future producer : producers) { + producer.get(10, TimeUnit.SECONDS); + } + addTo(totals, metrics.drain()); + assertEquals(20000L, totals.get("flagevaluation.rows.dropped").longValue()); + assertEquals(20000L, totals.get("flagevaluation.context.truncated").longValue()); + assertTrue(metrics.drain().isEmpty()); + } finally { + executor.shutdownNow(); + } + } + + private static void addTo(Map totals, List counts) { + for (Count count : counts) { + totals.merge(count.name, count.value, Long::sum); + } + } +} diff --git a/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java b/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java index 2da8ffe0256..8236301d1d5 100644 --- a/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java +++ b/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java @@ -1,5 +1,12 @@ package com.datadog.featureflag; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DEGRADED_CARDINALITY_CAP; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DEGRADED_PAYLOAD_LIMIT; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_CLOSED; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_DEGRADED_CAP; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_PAYLOAD_LIMIT; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_QUEUE_OVERFLOW; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.PAYLOAD_SPLITS; import static datadog.trace.util.AgentThreadFactory.AgentThread.FEATURE_FLAG_EVALUATION_PROCESSOR; import static datadog.trace.util.AgentThreadFactory.newAgentThread; import static java.util.concurrent.TimeUnit.SECONDS; @@ -12,13 +19,12 @@ import datadog.trace.api.Config; import datadog.trace.api.featureflag.FeatureFlaggingGateway; import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; +import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationWriter; -import datadog.trace.api.telemetry.CoreMetricCollector; import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import java.util.ArrayList; import java.util.List; import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -59,18 +65,8 @@ public class FlagEvaluationWriterImpl implements FlagEvaluationWriter { static final int FLUSH_INTERVAL_SECONDS = 10; static final int FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES = EvpProxy.PAYLOAD_SIZE_LIMIT_BYTES; - static final String FLAG_EVALUATION_DROPPED_METRIC = "flagevaluation.rows.dropped"; - static final String FLAG_EVALUATION_DEGRADED_METRIC = "flagevaluation.rows.degraded"; - static final String FLAG_EVALUATION_SPLITS_METRIC = "flagevaluation.payload.splits"; - static final String FLAG_EVALUATION_CONTEXT_TRUNCATED_METRIC = "flagevaluation.context.truncated"; - static final String DROP_REASON_QUEUE_OVERFLOW = "queue_overflow"; - static final String DROP_REASON_CLOSED = "closed"; - static final String DROP_REASON_DEGRADED_CAP = "degraded_cap"; - static final String DROP_REASON_PAYLOAD_LIMIT = "payload_limit"; - static final String DEGRADED_REASON_CARDINALITY_CAP = "cardinality_cap"; - static final String DEGRADED_REASON_PAYLOAD_LIMIT = "payload_limit"; private static final String FLAG_EVALUATION_ROUTE = "flagevaluation"; - private static final CoreMetricCollector CORE_METRICS = CoreMetricCollector.getInstance(); + private static final FlagEvaluationMetrics METRICS = FlagEvaluationMetrics.getInstance(); private final MessagePassingBlockingQueue queue; private final FlagEvaluationSerializingHandler serializer; @@ -78,27 +74,12 @@ public class FlagEvaluationWriterImpl implements FlagEvaluationWriter { private final Object lifecycleLock = new Object(); private final AtomicBoolean closed = new AtomicBoolean(false); - private static void countMetric(final String metricName, final long value, final String reason) { - if (value <= 0) { - return; - } - CORE_METRICS.count(metricName, value, reason == null ? null : "reason:" + reason); - } - /** * Observable counter for events dropped because the bounded hand-off queue was full when the hook * tried to enqueue (backpressure). Incremented on the hook thread, surfaced on flush. */ private final AtomicLong droppedQueueOverflow = new AtomicLong(0); - /** - * Per-reason-tag counters for evaluations whose context was truncated by copyPrunedContext. Keyed - * by the sorted comma-separated reason string (e.g. "max_key_length,max_value_length"). - * Incremented on the hook thread, drained and emitted on flush. - */ - private final ConcurrentHashMap contextTruncatedCounts = - new ConcurrentHashMap<>(); - public FlagEvaluationWriterImpl(final SharedCommunicationObjects sco, final Config config) { this( DEFAULT_CAPACITY, @@ -123,7 +104,6 @@ public FlagEvaluationWriterImpl(final SharedCommunicationObjects sco, final Conf timeUnit, FeatureFlagEvpContext.from(config), droppedQueueOverflow, - contextTruncatedCounts, this::close, FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES); this.serializerThread = newAgentThread(FEATURE_FLAG_EVALUATION_PROCESSOR, serializer); @@ -203,7 +183,7 @@ private void sweepAndCountResidualEvents() { while (queue.poll() != null) { residual++; } - countMetric(FLAG_EVALUATION_DROPPED_METRIC, residual, DROP_REASON_CLOSED); + METRICS.count(DROPPED_CLOSED, residual); } @Override @@ -251,7 +231,7 @@ public void countPreQueueOverflow() { @Override public void countContextTruncated(final String reason) { - contextTruncatedCounts.computeIfAbsent(reason, k -> new AtomicLong(0)).incrementAndGet(); + METRICS.countContextTruncated(reason, 1); } private boolean isClosedOrEnqueueDisabled() { @@ -263,7 +243,7 @@ private void countClosedDrop() { // by the surrounding subsystem. FeatureFlaggingSystem.stop() flips the gate before this // writer's close() runs, so an in-flight enqueue could race that flip and see gate=false while // closed=false. Counting on either condition keeps shutdown-loss observable. - countMetric(FLAG_EVALUATION_DROPPED_METRIC, 1, DROP_REASON_CLOSED); + METRICS.count(DROPPED_CLOSED, 1); } /** Returns the count of events dropped due to queue-overflow backpressure (observable). */ @@ -296,7 +276,6 @@ static class FlagEvaluationSerializingHandler implements Runnable { evpPublisher; final Map context; private final AtomicLong droppedQueueOverflow; - private final ConcurrentHashMap contextTruncatedCounts; private final Runnable errorCallback; private final int payloadSizeLimitBytes; final FlagEvaluationAggregator aggregator = new FlagEvaluationAggregator(); @@ -312,7 +291,6 @@ static class FlagEvaluationSerializingHandler implements Runnable { final TimeUnit timeUnit, final Map context, final AtomicLong droppedQueueOverflow, - final ConcurrentHashMap contextTruncatedCounts, final Runnable errorCallback, final int payloadSizeLimitBytes) { this.queue = queue; @@ -321,7 +299,6 @@ static class FlagEvaluationSerializingHandler implements Runnable { backendApiSupplier, FlagEvaluationPayloads.FlagEvaluationsRequest.class); this.context = context; this.droppedQueueOverflow = droppedQueueOverflow; - this.contextTruncatedCounts = contextTruncatedCounts; this.payloadSizeLimitBytes = payloadSizeLimitBytes; this.lastTicks = System.nanoTime(); this.ticksRequiredToFlush = timeUnit.toNanos(flushInterval); @@ -410,7 +387,7 @@ void flush() { // Surface backpressure (queue-overflow) drops as an observable warning even when there is // nothing else to flush. final long qDrops = droppedQueueOverflow.getAndSet(0); - countMetric(FLAG_EVALUATION_DROPPED_METRIC, qDrops, DROP_REASON_QUEUE_OVERFLOW); + METRICS.count(DROPPED_QUEUE_OVERFLOW, qDrops); if (qDrops > 0) { LOGGER.warn( "flag evaluation queue full - dropped {} evaluation(s) under backpressure" @@ -418,7 +395,7 @@ void flush() { qDrops); } final long dgDrops = aggregator.droppedDegradedOverflow.getAndSet(0); - countMetric(FLAG_EVALUATION_DROPPED_METRIC, dgDrops, DROP_REASON_DEGRADED_CAP); + METRICS.count(DROPPED_DEGRADED_CAP, dgDrops); if (dgDrops > 0) { LOGGER.warn( "degraded aggregation tier full - dropped {} evaluation(s); raise degraded cap" @@ -426,38 +403,21 @@ void flush() { dgDrops); } - // Drain per-reason context-truncation counters and emit one metric per unique reason tag. - for (final Map.Entry entry : contextTruncatedCounts.entrySet()) { - final long count = entry.getValue().getAndSet(0); - if (count > 0) { - countMetric(FLAG_EVALUATION_CONTEXT_TRUNCATED_METRIC, count, entry.getKey()); - } - } - if (aggregator.isEmpty()) { return; } try { - countMetric( - FLAG_EVALUATION_DEGRADED_METRIC, - aggregator.degradedEvaluationCount(), - DEGRADED_REASON_CARDINALITY_CAP); + METRICS.count(DEGRADED_CARDINALITY_CAP, aggregator.degradedEvaluationCount()); final List events = buildEventList(); if (events.isEmpty()) { return; } final FlagEvaluationPayloads.EncodedPayloads payloads = FlagEvaluationPayloads.buildPayloads(events, context, payloadSizeLimitBytes); - countMetric( - FLAG_EVALUATION_DROPPED_METRIC, - payloads.droppedPayloadLimit, - DROP_REASON_PAYLOAD_LIMIT); - countMetric( - FLAG_EVALUATION_DEGRADED_METRIC, - payloads.degradedPayloadLimit, - DEGRADED_REASON_PAYLOAD_LIMIT); + METRICS.count(DROPPED_PAYLOAD_LIMIT, payloads.droppedPayloadLimit); + METRICS.count(DEGRADED_PAYLOAD_LIMIT, payloads.degradedPayloadLimit); if (payloads.bodies.size() > 1) { - countMetric(FLAG_EVALUATION_SPLITS_METRIC, payloads.bodies.size() - 1, null); + METRICS.count(PAYLOAD_SPLITS, payloads.bodies.size() - 1); } if (payloads.droppedPayloadLimit > 0) { LOGGER.warn( @@ -537,7 +497,6 @@ static class SerializingHandlerForTest extends FlagEvaluationSerializingHandler TimeUnit.NANOSECONDS, context, new AtomicLong(0), - new ConcurrentHashMap<>(), () -> {}, payloadSizeLimitBytes); } diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationTestSupport.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationTestSupport.java index 502f0e76c69..5ee7c21efd4 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationTestSupport.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationTestSupport.java @@ -16,9 +16,8 @@ import datadog.communication.BackendApi; import datadog.communication.BackendApiFactory; import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; +import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; import datadog.trace.api.intake.Intake; -import datadog.trace.api.telemetry.CoreMetricCollector; -import datadog.trace.api.telemetry.MetricCollector; import java.lang.reflect.Type; import java.util.ArrayList; import java.util.Arrays; @@ -26,6 +25,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.function.Supplier; import okhttp3.RequestBody; import okio.Buffer; @@ -47,8 +47,8 @@ static Supplier backendApiSupplier(final BackendApiFactory factory) return () -> factory.createBackendApi(Intake.EVENT_PLATFORM, false); } - static void clearCoreMetrics() { - CoreMetricCollector.getInstance().drain(); + static void clearFlagEvaluationMetrics() { + FlagEvaluationMetrics.getInstance().drain(); } static FlagEvalEvent event( @@ -160,22 +160,14 @@ static Map flushAndCaptureJson(final TestWriterSetup setup) thro } static long metricSum( - final Collection metrics, + final Collection metrics, final String metricName, final String tag) { long sum = 0; - for (final MetricCollector.Metric metric : metrics) { - if (!metricName.equals(metric.metricName)) { - continue; + for (final FlagEvaluationMetrics.Count metric : metrics) { + if (metricName.equals(metric.name) && Objects.equals(tag, metric.tag)) { + sum += metric.value; } - if (tag == null) { - if (!metric.tags.isEmpty()) { - continue; - } - } else if (!metric.tags.contains(tag)) { - continue; - } - sum += metric.value.longValue(); } return sum; } diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java index 75a767c5c3f..ab944cf7119 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java @@ -4,7 +4,7 @@ import static com.datadog.featureflag.FlagEvaluationTestSupport.backendApiSupplier; import static com.datadog.featureflag.FlagEvaluationTestSupport.buildTestWriter; import static com.datadog.featureflag.FlagEvaluationTestSupport.cfg; -import static com.datadog.featureflag.FlagEvaluationTestSupport.clearCoreMetrics; +import static com.datadog.featureflag.FlagEvaluationTestSupport.clearFlagEvaluationMetrics; import static com.datadog.featureflag.FlagEvaluationTestSupport.event; import static com.datadog.featureflag.FlagEvaluationTestSupport.eventForFlag; import static com.datadog.featureflag.FlagEvaluationTestSupport.flushAndCapture; @@ -40,10 +40,9 @@ import datadog.trace.api.Config; import datadog.trace.api.featureflag.FeatureFlaggingGateway; import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; +import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; import datadog.trace.api.featureflag.ufc.v1.ServerConfiguration; import datadog.trace.api.intake.Intake; -import datadog.trace.api.telemetry.CoreMetricCollector; -import datadog.trace.api.telemetry.MetricCollector; import datadog.trace.test.util.PollingConditions; import java.io.IOException; import java.lang.reflect.Field; @@ -68,14 +67,14 @@ class FlagEvaluationWriterImplTest { private static final double TIMEOUT_SECONDS = 5; @BeforeEach - void clearCoreMetricsBefore() { - clearCoreMetrics(); + void clearFlagEvaluationMetricsBefore() { + clearFlagEvaluationMetrics(); FeatureFlaggingGateway.setFlagEvaluationEnqueueEnabled(true); } @AfterEach - void clearCoreMetricsAfter() { - clearCoreMetrics(); + void clearFlagEvaluationMetricsAfter() { + clearFlagEvaluationMetrics(); FeatureFlaggingGateway.setFlagEvalWriter(null); FeatureFlaggingGateway.setFlagEvaluationEnqueueEnabled(true); // Reset the dispatched UFC state so observeFullEvaluationData can't leak into other tests. @@ -90,14 +89,9 @@ void degradedCapOverflowTelemetryIsEmittedOnFlush() { setup.handler.addDroppedDegradedOverflowForTest(3); setup.handler.flush(); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); - assertEquals( - 3, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_DROPPED_METRIC, - "reason:" + FlagEvaluationWriterImpl.DROP_REASON_DEGRADED_CAP)); + final Collection metrics = + FlagEvaluationMetrics.getInstance().drain(); + assertEquals(3, metricSum(metrics, "flagevaluation.rows.dropped", "reason:degraded_cap")); } @Test @@ -134,14 +128,10 @@ void queueOverflowIncrementsObservableDropCounter() { assertTrue(writer.droppedQueueOverflow() > 0); final long queueDrops = writer.droppedQueueOverflow(); writer.flushForTest(); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); + final Collection metrics = + FlagEvaluationMetrics.getInstance().drain(); assertEquals( - queueDrops, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_DROPPED_METRIC, - "reason:" + FlagEvaluationWriterImpl.DROP_REASON_QUEUE_OVERFLOW)); + queueDrops, metricSum(metrics, "flagevaluation.rows.dropped", "reason:queue_overflow")); } @Test @@ -156,14 +146,9 @@ void enqueueAfterCloseIsDroppedAndCounted() { writer.close(); writer.enqueue(simpleEvent("closed-flag", "on")); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); - assertEquals( - 1, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_DROPPED_METRIC, - "reason:" + FlagEvaluationWriterImpl.DROP_REASON_CLOSED)); + final Collection metrics = + FlagEvaluationMetrics.getInstance().drain(); + assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -181,14 +166,9 @@ void enqueueDisabledDropsAndCountsAsClosedDrop() { FeatureFlaggingGateway.setFlagEvaluationEnqueueEnabled(false); writer.enqueue(simpleEvent("disabled-flag", "on")); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); - assertEquals( - 1, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_DROPPED_METRIC, - "reason:" + FlagEvaluationWriterImpl.DROP_REASON_CLOSED)); + final Collection metrics = + FlagEvaluationMetrics.getInstance().drain(); + assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -208,14 +188,9 @@ void closeSweepsAndCountsEventsLeftInTheQueue() { writer.enqueue(simpleEvent("residual-flag-2", "on")); writer.close(); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); - assertEquals( - 2, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_DROPPED_METRIC, - "reason:" + FlagEvaluationWriterImpl.DROP_REASON_CLOSED)); + final Collection metrics = + FlagEvaluationMetrics.getInstance().drain(); + assertEquals(2, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -283,7 +258,6 @@ void flushIfNecessaryDoesNotReturnEarlyWhenOnlyQueueDropsArePending() { TimeUnit.NANOSECONDS, context(), queueDrops, - new java.util.concurrent.ConcurrentHashMap<>(), () -> {}, FlagEvaluationWriterImpl.FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES); @@ -314,7 +288,6 @@ void workerHandlesEmptyPolls() throws Exception { TimeUnit.NANOSECONDS, context(), new AtomicLong(0), - new java.util.concurrent.ConcurrentHashMap<>(), () -> {}, FlagEvaluationWriterImpl.FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES); @@ -359,14 +332,9 @@ void payloadLimitDropsAreCountedOnFlush() { setup.handler.drainAndAggregate(); setup.handler.flush(); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); - assertEquals( - 1, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_DROPPED_METRIC, - "reason:" + FlagEvaluationWriterImpl.DROP_REASON_PAYLOAD_LIMIT)); + final Collection metrics = + FlagEvaluationMetrics.getInstance().drain(); + assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:payload_limit")); } @Test @@ -728,7 +696,7 @@ void agentlessWritesFlagEvaluationsDirectlyWhenLocalProxyIsUnavailable() throws } @Test - void countContextTruncatedAccumulatesPerReason() { + void countContextTruncatedAccumulatesPerReasonWithoutFlushingWriter() { final BackendApi mockEvp = mock(BackendApi.class); final BackendApiFactory factory = mock(BackendApiFactory.class); when(factory.createBackendApi(any(), anyBoolean())).thenReturn(mockEvp); @@ -739,22 +707,13 @@ void countContextTruncatedAccumulatesPerReason() { writer.countContextTruncated("field_count"); writer.countContextTruncated("field_count"); writer.countContextTruncated("field_length"); - writer.flushForTest(); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); - assertEquals( - 2, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_CONTEXT_TRUNCATED_METRIC, - "reason:field_count")); - assertEquals( - 1, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_CONTEXT_TRUNCATED_METRIC, - "reason:field_length")); + final Collection metrics = + FlagEvaluationMetrics.getInstance().drain(); + assertEquals(2, metricSum(metrics, "flagevaluation.context.truncated", "reason:field_count")); + assertEquals(1, metricSum(metrics, "flagevaluation.context.truncated", "reason:field_length")); + writer.flushForTest(); + assertTrue(FlagEvaluationMetrics.getInstance().drain().isEmpty()); } @Test @@ -779,14 +738,9 @@ void hasCapacityForEnqueueReflectsQueueSaturationAndCountsPreQueueOverflow() { writer.countPreQueueOverflow(); writer.flushForTest(); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); - assertEquals( - 1, - metricSum( - metrics, - FlagEvaluationWriterImpl.FLAG_EVALUATION_DROPPED_METRIC, - "reason:" + FlagEvaluationWriterImpl.DROP_REASON_QUEUE_OVERFLOW)); + final Collection metrics = + FlagEvaluationMetrics.getInstance().drain(); + assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:queue_overflow")); writer.close(); } diff --git a/telemetry/build.gradle.kts b/telemetry/build.gradle.kts index f2efa690d6e..130a5bea872 100644 --- a/telemetry/build.gradle.kts +++ b/telemetry/build.gradle.kts @@ -34,6 +34,7 @@ dependencies { implementation(libs.slf4j) implementation(project(":internal-api")) + implementation(project(":products:feature-flagging:feature-flagging-bootstrap")) compileOnly(project(":dd-java-agent:agent-tooling")) testImplementation(project(":dd-java-agent:agent-tooling")) diff --git a/telemetry/src/main/java/datadog/telemetry/TelemetrySystem.java b/telemetry/src/main/java/datadog/telemetry/TelemetrySystem.java index 4b233aa3037..fe28c5ed3f7 100644 --- a/telemetry/src/main/java/datadog/telemetry/TelemetrySystem.java +++ b/telemetry/src/main/java/datadog/telemetry/TelemetrySystem.java @@ -13,6 +13,7 @@ import datadog.telemetry.metric.ConfigInversionMetricPeriodicAction; import datadog.telemetry.metric.CoreMetricsPeriodicAction; import datadog.telemetry.metric.DebuggerMetricPeriodicAction; +import datadog.telemetry.metric.FlagEvaluationMetricPeriodicAction; import datadog.telemetry.metric.IastMetricPeriodicAction; import datadog.telemetry.metric.LLMObsMetricPeriodicAction; import datadog.telemetry.metric.OtelEnvMetricPeriodicAction; @@ -61,6 +62,9 @@ static Thread createTelemetryRunnable( List actions = new ArrayList<>(); if (telemetryMetricsEnabled) { actions.add(new CoreMetricsPeriodicAction()); + if (Config.get().isFeatureFlaggingProviderEnabled()) { + actions.add(new FlagEvaluationMetricPeriodicAction()); + } actions.add(new OtelEnvMetricPeriodicAction()); if (InstrumenterConfig.get().getTraceExtensionsPath() != null) { actions.add(new OtelSpiMetricPeriodicAction()); diff --git a/telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java b/telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java new file mode 100644 index 00000000000..4429b7bde2f --- /dev/null +++ b/telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java @@ -0,0 +1,55 @@ +package datadog.telemetry.metric; + +import static java.util.Collections.emptyIterator; +import static java.util.Collections.emptyList; + +import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; +import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Count; +import datadog.trace.api.telemetry.MetricCollector; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Iterator; +import java.util.List; +import java.util.concurrent.ArrayBlockingQueue; + +/** Adapts product-owned flag evaluation counters to telemetry on its collection schedule. */ +public final class FlagEvaluationMetricPeriodicAction extends MetricPeriodicAction { + private final Collector collector = new Collector(); + + @Override + public MetricCollector collector() { + return collector; + } + + private static final class Collector implements MetricCollector { + private final FlagEvaluationMetrics metrics = FlagEvaluationMetrics.getInstance(); + private final ArrayBlockingQueue queue = new ArrayBlockingQueue<>(RAW_QUEUE_SIZE); + private Iterator pending = emptyIterator(); + + @Override + public void prepareMetrics() { + if (queue.remainingCapacity() == 0) { + return; + } + if (!pending.hasNext()) { + pending = metrics.drain().iterator(); + } + // Only the telemetry thread prepares metrics. Retain the rest of a snapshot when full; + // producer updates remain in the counters until the next snapshot can be collected. + while (queue.remainingCapacity() > 0 && pending.hasNext()) { + Count count = pending.next(); + queue.offer(new Metric("tracers", true, count.name, "count", count.value, count.tag)); + } + } + + @Override + public Collection drain() { + if (queue.isEmpty()) { + return emptyList(); + } + List drained = new ArrayList<>(queue.size()); + queue.drainTo(drained); + return drained; + } + } +} diff --git a/telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java b/telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java new file mode 100644 index 00000000000..0c63a5b57bf --- /dev/null +++ b/telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java @@ -0,0 +1,146 @@ +package datadog.telemetry.metric; + +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DEGRADED_CARDINALITY_CAP; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_CLOSED; +import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.PAYLOAD_SPLITS; +import static datadog.trace.api.telemetry.MetricCollector.RAW_QUEUE_SIZE; +import static java.util.Collections.emptyList; +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.ArgumentCaptor.forClass; +import static org.mockito.Mockito.mock; +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.featureflag.flagevaluation.FlagEvaluationMetrics; +import datadog.trace.api.telemetry.CoreMetricCollector; +import datadog.trace.api.telemetry.MetricCollector; +import java.util.ArrayList; +import java.util.Collection; +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; +import org.tabletest.junit.TableTest; + +class FlagEvaluationMetricPeriodicActionTest { + private final FlagEvaluationMetrics metrics = FlagEvaluationMetrics.getInstance(); + private final FlagEvaluationMetricPeriodicAction action = + new FlagEvaluationMetricPeriodicAction(); + private final TelemetryService service = mock(TelemetryService.class); + + @BeforeEach + @AfterEach + void resetMetrics() { + metrics.drain(); + } + + @Test + void defaultActionDrainsProductCountersIndependentlyOfCoreMetrics() { + metrics.count(DROPPED_CLOSED, 3); + CoreMetricCollector core = CoreMetricCollector.getInstance(); + core.prepareMetrics(); + assertTrue( + core.drain().stream().noneMatch(metric -> metric.metricName.startsWith("flagevaluation."))); + + action.collector().prepareMetrics(); + action.doIteration(service); + ArgumentCaptor captor = forClass(Metric.class); + verify(service).addMetric(captor.capture()); + assertEquals("flagevaluation.rows.dropped", captor.getValue().getMetric()); + assertEquals(3L, captor.getValue().getPoints().get(0).get(1).longValue()); + } + + @TableTest({ + "Scenario | Counter | Name | Tag ", + "queue overflow | DROPPED_QUEUE_OVERFLOW | flagevaluation.rows.dropped | reason:queue_overflow ", + "closed | DROPPED_CLOSED | flagevaluation.rows.dropped | reason:closed ", + "degraded cap | DROPPED_DEGRADED_CAP | flagevaluation.rows.dropped | reason:degraded_cap ", + "payload drop | DROPPED_PAYLOAD_LIMIT | flagevaluation.rows.dropped | reason:payload_limit ", + "cardinality cap | DEGRADED_CARDINALITY_CAP | flagevaluation.rows.degraded | reason:cardinality_cap", + "payload degraded | DEGRADED_PAYLOAD_LIMIT | flagevaluation.rows.degraded | reason:payload_limit ", + "split | PAYLOAD_SPLITS | flagevaluation.payload.splits | " + }) + void preservesWireFormatAcrossCollectionIntervals( + FlagEvaluationMetrics.Metric counter, String name, String tag) { + metrics.count(counter, 2); + metrics.count(counter, 3); + assertTrue(action.collector().drain().isEmpty()); + action.collector().prepareMetrics(); + metrics.count(counter, 7); + action.collector().prepareMetrics(); + action.doIteration(service); + + ArgumentCaptor captor = forClass(Metric.class); + verify(service).addMetric(captor.capture()); + Metric metric = captor.getValue(); + assertEquals(name, metric.getMetric()); + assertEquals("tracers", metric.getNamespace()); + assertEquals(true, metric.getCommon()); + assertEquals(Metric.TypeEnum.COUNT, metric.getType()); + assertEquals(tag == null ? emptyList() : singletonList(tag), metric.getTags()); + assertEquals(2, metric.getPoints().size()); + assertEquals(5L, metric.getPoints().get(0).get(1).longValue()); + assertEquals(7L, metric.getPoints().get(1).get(1).longValue()); + action.collector().prepareMetrics(); + action.doIteration(service); + verifyNoMoreInteractions(service); + } + + @Test + void preservesContextTruncationTags() { + metrics.countContextTruncated("max_key_length,max_value_length", 3); + action.collector().prepareMetrics(); + action.doIteration(service); + + ArgumentCaptor captor = forClass(Metric.class); + verify(service).addMetric(captor.capture()); + Metric metric = captor.getValue(); + assertEquals("flagevaluation.context.truncated", metric.getMetric()); + assertEquals(singletonList("reason:max_key_length,max_value_length"), metric.getTags()); + assertEquals(3L, metric.getPoints().get(0).get(1).longValue()); + } + + @Test + void fullQueueRetainsSnapshotRemainderAndSubsequentProducerUpdates() { + MetricCollector collector = action.collector(); + for (int i = 0; i < RAW_QUEUE_SIZE - 1; i++) { + metrics.count(PAYLOAD_SPLITS, 1); + collector.prepareMetrics(); + } + metrics.count(DROPPED_CLOSED, 5); + metrics.count(DEGRADED_CARDINALITY_CAP, 7); + metrics.countContextTruncated("max_depth", 3); + collector.prepareMetrics(); + metrics.count(DROPPED_CLOSED, 11); + collector.prepareMetrics(); + + List reported = new ArrayList<>(collector.drain()); + assertEquals(RAW_QUEUE_SIZE, reported.size()); + collector.prepareMetrics(); + Collection remainder = collector.drain(); + assertEquals(2, remainder.size()); + reported.addAll(remainder); + collector.prepareMetrics(); + reported.addAll(collector.drain()); + + assertEquals(RAW_QUEUE_SIZE - 1, sum(reported, "flagevaluation.payload.splits")); + assertEquals(16, sum(reported, "flagevaluation.rows.dropped")); + assertEquals(7, sum(reported, "flagevaluation.rows.degraded")); + assertEquals(3, sum(reported, "flagevaluation.context.truncated")); + collector.prepareMetrics(); + assertTrue(collector.drain().isEmpty()); + assertTrue(metrics.drain().isEmpty()); + } + + private static long sum(List metrics, String name) { + return metrics.stream() + .filter(metric -> name.equals(metric.metricName)) + .mapToLong(metric -> metric.value.longValue()) + .sum(); + } +} From e10920c229291425aa0f0f321c9efca48a66e9a6 Mon Sep 17 00:00:00 2001 From: Sarah Chen Date: Mon, 28 Sep 2026 16:12:07 -0400 Subject: [PATCH 2/3] Add startForTest coverage --- .../com/datadog/featureflag/FlagEvaluationWriterImplTest.java | 1 + 1 file changed, 1 insertion(+) diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java index ab944cf7119..1560e058bec 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java @@ -109,6 +109,7 @@ void startRegistersWriterAndCloseDeregistersIt() { writer.close(); writer.close(); writer.start(); + writer.startForTest(); assertNull(FeatureFlaggingGateway.getFlagEvalWriter()); } From 0dc5e8207f0d6c57db0c26ba72d7e8f49e419646 Mon Sep 17 00:00:00 2001 From: Sarah Chen Date: Thu, 1 Oct 2026 12:06:29 -0400 Subject: [PATCH 3/3] Add FlagEvaluationMetricCollector in internal-api --- .../FlagEvaluationMetricCollector.java | 120 ++++++++++++ .../FlagEvaluationMetricCollectorTest.java | 176 ++++++++++++++++++ .../build.gradle.kts | 2 - .../gradle.lockfile | 3 +- .../flagevaluation/FlagEvaluationMetrics.java | 87 --------- .../FlagEvaluationMetricsTest.java | 110 ----------- .../featureflag/FlagEvaluationWriterImpl.java | 18 +- .../FlagEvaluationTestSupport.java | 28 ++- .../FlagEvaluationWriterImplTest.java | 29 ++- telemetry/build.gradle.kts | 1 - .../FlagEvaluationMetricPeriodicAction.java | 51 +---- ...lagEvaluationMetricPeriodicActionTest.java | 121 ++---------- 12 files changed, 356 insertions(+), 390 deletions(-) create mode 100644 internal-api/src/main/java/datadog/trace/api/telemetry/FlagEvaluationMetricCollector.java create mode 100644 internal-api/src/test/java/datadog/trace/api/telemetry/FlagEvaluationMetricCollectorTest.java delete mode 100644 products/feature-flagging/feature-flagging-bootstrap/src/main/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetrics.java delete mode 100644 products/feature-flagging/feature-flagging-bootstrap/src/test/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetricsTest.java diff --git a/internal-api/src/main/java/datadog/trace/api/telemetry/FlagEvaluationMetricCollector.java b/internal-api/src/main/java/datadog/trace/api/telemetry/FlagEvaluationMetricCollector.java new file mode 100644 index 00000000000..e0b897a8a42 --- /dev/null +++ b/internal-api/src/main/java/datadog/trace/api/telemetry/FlagEvaluationMetricCollector.java @@ -0,0 +1,120 @@ +package datadog.trace.api.telemetry; + +import datadog.trace.api.internal.VisibleForTesting; +import java.util.ArrayList; +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.AtomicLong; +import java.util.concurrent.atomic.AtomicLongArray; + +/** Collects telemetry metrics produced by flag evaluation. */ +public final class FlagEvaluationMetricCollector + implements MetricCollector { + private static final String METRIC_NAMESPACE = "tracers"; + private static final String CONTEXT_TRUNCATED_METRIC = "flagevaluation.context.truncated"; + private static final Counter[] COUNTERS = Counter.values(); + private static final FlagEvaluationMetricCollector INSTANCE = new FlagEvaluationMetricCollector(); + + private final AtomicLongArray counts = new AtomicLongArray(COUNTERS.length); + private final ConcurrentHashMap contextTruncationCounts = + new ConcurrentHashMap<>(); + private final BlockingQueue metricsQueue = + new ArrayBlockingQueue<>(RAW_QUEUE_SIZE); + + private FlagEvaluationMetricCollector() {} + + public static FlagEvaluationMetricCollector get() { + return INSTANCE; + } + + public enum Counter { + DROPPED_QUEUE_OVERFLOW("flagevaluation.rows.dropped", "reason:queue_overflow"), + DROPPED_CLOSED("flagevaluation.rows.dropped", "reason:closed"), + DROPPED_DEGRADED_CAP("flagevaluation.rows.dropped", "reason:degraded_cap"), + DROPPED_PAYLOAD_LIMIT("flagevaluation.rows.dropped", "reason:payload_limit"), + DEGRADED_CARDINALITY_CAP("flagevaluation.rows.degraded", "reason:cardinality_cap"), + DEGRADED_PAYLOAD_LIMIT("flagevaluation.rows.degraded", "reason:payload_limit"), + PAYLOAD_SPLITS("flagevaluation.payload.splits", null); + + private final String name; + private final String tag; + + Counter(String name, String tag) { + this.name = name; + this.tag = tag; + } + } + + public void count(Counter counter, long value) { + if (value > 0) { + counts.addAndGet(counter.ordinal(), value); + } + } + + /** Records the hook's canonical, comma-separated combination of truncation reasons. */ + public void countContextTruncated(String reason, long value) { + if (value > 0) { + contextTruncationCounts.computeIfAbsent(reason, key -> new AtomicLong()).addAndGet(value); + } + } + + @Override + public void prepareMetrics() { + for (Counter counter : COUNTERS) { + if (metricsQueue.remainingCapacity() == 0) { + return; + } + long value = counts.getAndSet(counter.ordinal(), 0); + if (value > 0 + && !metricsQueue.offer(new FlagEvaluationMetric(counter.name, value, counter.tag))) { + counts.addAndGet(counter.ordinal(), value); + return; + } + } + + for (Map.Entry entry : contextTruncationCounts.entrySet()) { + if (metricsQueue.remainingCapacity() == 0) { + return; + } + long value = entry.getValue().getAndSet(0); + if (value > 0 + && !metricsQueue.offer( + new FlagEvaluationMetric( + CONTEXT_TRUNCATED_METRIC, value, "reason:" + entry.getKey()))) { + entry.getValue().addAndGet(value); + return; + } + } + } + + @Override + public Collection drain() { + if (metricsQueue.isEmpty()) { + return Collections.emptyList(); + } + List drained = new ArrayList<>(metricsQueue.size()); + metricsQueue.drainTo(drained); + return drained; + } + + /** Clears all pending counters and metrics. Visible for testing only. */ + @VisibleForTesting + public void resetForTesting() { + for (int i = 0; i < counts.length(); i++) { + counts.set(i, 0); + } + contextTruncationCounts.clear(); + metricsQueue.clear(); + } + + public static final class FlagEvaluationMetric extends MetricCollector.Metric { + private FlagEvaluationMetric(String name, long value, String tag) { + super(METRIC_NAMESPACE, true, name, "count", value, tag); + } + } +} diff --git a/internal-api/src/test/java/datadog/trace/api/telemetry/FlagEvaluationMetricCollectorTest.java b/internal-api/src/test/java/datadog/trace/api/telemetry/FlagEvaluationMetricCollectorTest.java new file mode 100644 index 00000000000..ef879c025bc --- /dev/null +++ b/internal-api/src/test/java/datadog/trace/api/telemetry/FlagEvaluationMetricCollectorTest.java @@ -0,0 +1,176 @@ +package datadog.trace.api.telemetry; + +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.DROPPED_CLOSED; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.PAYLOAD_SPLITS; +import static datadog.trace.api.telemetry.MetricCollector.RAW_QUEUE_SIZE; +import static java.util.Collections.emptyList; +import static java.util.Collections.singletonList; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter; +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector.FlagEvaluationMetric; +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.tabletest.junit.TableTest; + +class FlagEvaluationMetricCollectorTest { + private final FlagEvaluationMetricCollector collector = FlagEvaluationMetricCollector.get(); + + @BeforeEach + @AfterEach + void resetCollector() { + collector.resetForTesting(); + } + + @Test + void accumulatesPositiveLongCountsAndDrainsOnlyOnce() { + collector.count(DROPPED_CLOSED, Integer.MAX_VALUE); + collector.count(DROPPED_CLOSED, 5); + collector.count(DROPPED_CLOSED, 0); + collector.count(DROPPED_CLOSED, -1); + + List metrics = collect(); + assertEquals(1, metrics.size()); + assertEquals(Integer.MAX_VALUE + 5L, metrics.get(0).value.longValue()); + assertTrue(collect().isEmpty()); + + collector.count(DROPPED_CLOSED, 2); + assertEquals(2, collect().get(0).value.longValue()); + } + + @TableTest({ + "Scenario | Counter | Name | Tag ", + "queue overflow | DROPPED_QUEUE_OVERFLOW | flagevaluation.rows.dropped | reason:queue_overflow ", + "closed | DROPPED_CLOSED | flagevaluation.rows.dropped | reason:closed ", + "degraded cap | DROPPED_DEGRADED_CAP | flagevaluation.rows.dropped | reason:degraded_cap ", + "payload drop | DROPPED_PAYLOAD_LIMIT | flagevaluation.rows.dropped | reason:payload_limit ", + "cardinality cap | DEGRADED_CARDINALITY_CAP | flagevaluation.rows.degraded | reason:cardinality_cap", + "payload degraded | DEGRADED_PAYLOAD_LIMIT | flagevaluation.rows.degraded | reason:payload_limit ", + "split | PAYLOAD_SPLITS | flagevaluation.payload.splits | " + }) + void preservesWireFormat(Counter counter, String name, String tag) { + collector.count(counter, 3); + + FlagEvaluationMetric metric = collect().get(0); + assertEquals(name, metric.metricName); + assertEquals("tracers", metric.namespace); + assertTrue(metric.common); + assertEquals("count", metric.type); + assertEquals(tag == null ? emptyList() : singletonList(tag), metric.tags); + assertEquals(3L, metric.value.longValue()); + } + + @Test + void preservesCombinedTruncationReasonsAcrossIntervals() { + collector.countContextTruncated("max_key_length,max_value_length", 2); + collector.countContextTruncated("max_key_length,max_value_length", 3); + collector.countContextTruncated("max_depth", 1); + collector.countContextTruncated("ignored", 0); + collector.countContextTruncated("ignored", -1); + + Map counts = valuesByTag(collect()); + assertEquals(2, counts.size()); + assertEquals(5L, counts.get("reason:max_key_length,max_value_length").longValue()); + assertEquals(1L, counts.get("reason:max_depth").longValue()); + assertTrue(collect().isEmpty()); + + collector.countContextTruncated("max_depth", 4); + assertEquals(4L, collect().get(0).value.longValue()); + } + + @Test + void fullQueueLeavesCountersPendingForTheNextCollection() { + for (int i = 0; i < RAW_QUEUE_SIZE; i++) { + collector.count(PAYLOAD_SPLITS, 1); + collector.prepareMetrics(); + } + collector.count(DROPPED_CLOSED, 5); + collector.countContextTruncated("max_depth", 3); + collector.prepareMetrics(); + collector.count(DROPPED_CLOSED, 11); + + List staged = new ArrayList<>(collector.drain()); + assertEquals(RAW_QUEUE_SIZE, staged.size()); + assertEquals(RAW_QUEUE_SIZE, sum(staged, "flagevaluation.payload.splits")); + + List pending = collect(); + assertEquals(16, sum(pending, "flagevaluation.rows.dropped")); + assertEquals(3, sum(pending, "flagevaluation.context.truncated")); + assertTrue(collect().isEmpty()); + } + + @Test + void concurrentProducersAndCollectionPreserveAllCounts() throws Exception { + ExecutorService executor = Executors.newFixedThreadPool(4); + CountDownLatch start = new CountDownLatch(1); + List> producers = new ArrayList<>(); + try { + for (int producer = 0; producer < 4; producer++) { + producers.add( + executor.submit( + () -> { + start.await(); + for (int i = 0; i < 5000; i++) { + collector.count(DROPPED_CLOSED, 1); + collector.countContextTruncated("max_depth", 1); + } + return null; + })); + } + start.countDown(); + Map totals = new HashMap<>(); + for (int i = 0; i < 100; i++) { + addTo(totals, collect()); + } + for (Future producer : producers) { + producer.get(10, TimeUnit.SECONDS); + } + addTo(totals, collect()); + + assertEquals(20000L, totals.get("flagevaluation.rows.dropped").longValue()); + assertEquals(20000L, totals.get("flagevaluation.context.truncated").longValue()); + assertTrue(collect().isEmpty()); + } finally { + executor.shutdownNow(); + } + } + + private List collect() { + collector.prepareMetrics(); + return new ArrayList<>(collector.drain()); + } + + private static Map valuesByTag(Collection metrics) { + Map values = new HashMap<>(); + for (FlagEvaluationMetric metric : metrics) { + assertEquals("flagevaluation.context.truncated", metric.metricName); + values.put(metric.tags.get(0), metric.value.longValue()); + } + return values; + } + + private static void addTo(Map totals, Collection metrics) { + for (FlagEvaluationMetric metric : metrics) { + totals.merge(metric.metricName, metric.value.longValue(), Long::sum); + } + } + + private static long sum(Collection metrics, String name) { + return metrics.stream() + .filter(metric -> name.equals(metric.metricName)) + .mapToLong(metric -> metric.value.longValue()) + .sum(); + } +} diff --git a/products/feature-flagging/feature-flagging-bootstrap/build.gradle.kts b/products/feature-flagging/feature-flagging-bootstrap/build.gradle.kts index 9c0add8f901..848b16a2841 100644 --- a/products/feature-flagging/feature-flagging-bootstrap/build.gradle.kts +++ b/products/feature-flagging/feature-flagging-bootstrap/build.gradle.kts @@ -29,8 +29,6 @@ extra["excludedClassesCoverage"] = listOf( ) dependencies { - implementation(project(":products:metrics:metrics-api")) - testImplementation(libs.bundles.junit5) testImplementation(libs.bundles.mockito) testImplementation(project(":utils:test-utils")) diff --git a/products/feature-flagging/feature-flagging-bootstrap/gradle.lockfile b/products/feature-flagging/feature-flagging-bootstrap/gradle.lockfile index 66f40fe18c8..4e8cd8dbd72 100644 --- a/products/feature-flagging/feature-flagging-bootstrap/gradle.lockfile +++ b/products/feature-flagging/feature-flagging-bootstrap/gradle.lockfile @@ -68,7 +68,6 @@ org.ow2.asm:asm:9.10.1=jacocoAnt,spotbugs org.slf4j:jcl-over-slf4j:1.7.30=testCompileClasspath,testRuntimeClasspath org.slf4j:jul-to-slf4j:1.7.30=testCompileClasspath,testRuntimeClasspath org.slf4j:log4j-over-slf4j:1.7.30=testCompileClasspath,testRuntimeClasspath -org.slf4j:slf4j-api:1.7.30=runtimeClasspath org.slf4j:slf4j-api:1.7.32=testCompileClasspath,testRuntimeClasspath org.slf4j:slf4j-api:2.0.17=spotbugsSlf4j org.slf4j:slf4j-api:2.0.18=spotbugs @@ -79,4 +78,4 @@ org.spockframework:spock-core:2.4-groovy-3.0=testCompileClasspath,testRuntimeCla org.tabletest:tabletest-junit:1.2.2=testCompileClasspath,testRuntimeClasspath org.tabletest:tabletest-parser:1.2.1=testCompileClasspath,testRuntimeClasspath org.xmlresolver:xmlresolver:5.3.3=spotbugs -empty=annotationProcessor,spotbugsPlugins,testAnnotationProcessor +empty=annotationProcessor,runtimeClasspath,spotbugsPlugins,testAnnotationProcessor diff --git a/products/feature-flagging/feature-flagging-bootstrap/src/main/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetrics.java b/products/feature-flagging/feature-flagging-bootstrap/src/main/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetrics.java deleted file mode 100644 index ad6631ac48e..00000000000 --- a/products/feature-flagging/feature-flagging-bootstrap/src/main/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetrics.java +++ /dev/null @@ -1,87 +0,0 @@ -package datadog.trace.api.featureflag.flagevaluation; - -import datadog.metrics.api.Accumulator; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.atomic.AtomicLong; - -/** Flag evaluation counters owned by the producer and drained periodically by telemetry. */ -public final class FlagEvaluationMetrics { - private static final FlagEvaluationMetrics INSTANCE = new FlagEvaluationMetrics(); - - private final Accumulator counts = Accumulator.of(Metric.class); - private final ConcurrentHashMap contextTruncations = - new ConcurrentHashMap<>(); - - private FlagEvaluationMetrics() {} - - public static FlagEvaluationMetrics getInstance() { - return INSTANCE; - } - - public enum Metric { - DROPPED_QUEUE_OVERFLOW("flagevaluation.rows.dropped", "queue_overflow"), - DROPPED_CLOSED("flagevaluation.rows.dropped", "closed"), - DROPPED_DEGRADED_CAP("flagevaluation.rows.dropped", "degraded_cap"), - DROPPED_PAYLOAD_LIMIT("flagevaluation.rows.dropped", "payload_limit"), - DEGRADED_CARDINALITY_CAP("flagevaluation.rows.degraded", "cardinality_cap"), - DEGRADED_PAYLOAD_LIMIT("flagevaluation.rows.degraded", "payload_limit"), - PAYLOAD_SPLITS("flagevaluation.payload.splits", null); - - private final String name; - private final String tag; - - Metric(String name, String reason) { - this.name = name; - this.tag = reason == null ? null : "reason:" + reason; - } - } - - public void count(Metric metric, long value) { - if (value > 0) { - counts.add(metric, value); - } - } - - /** Records the hook's canonical, comma-separated combination of truncation reasons. */ - public void countContextTruncated(String reason, long value) { - if (value > 0) { - contextTruncations.computeIfAbsent(reason, key -> new AtomicLong()).addAndGet(value); - } - } - - /** Returns nonzero deltas, resetting only the counts included in this snapshot. */ - public List drain() { - List drained = new ArrayList<>(); - Accumulator.Counts snapshot = counts.accumulateAndReset(); - for (Metric metric : snapshot.keys()) { - long value = snapshot.get(metric); - if (value > 0) { - drained.add(new Count(metric.name, metric.tag, value)); - } - } - for (Map.Entry entry : contextTruncations.entrySet()) { - long value = entry.getValue().getAndSet(0); - if (value > 0) { - drained.add( - new Count("flagevaluation.context.truncated", "reason:" + entry.getKey(), value)); - } - } - return drained; - } - - /** A counter delta with its metric name and optional reason tag, independent of transport. */ - public static final class Count { - public final String name; - public final String tag; - public final long value; - - private Count(String name, String tag, long value) { - this.name = name; - this.tag = tag; - this.value = value; - } - } -} diff --git a/products/feature-flagging/feature-flagging-bootstrap/src/test/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetricsTest.java b/products/feature-flagging/feature-flagging-bootstrap/src/test/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetricsTest.java deleted file mode 100644 index c83e3c0215a..00000000000 --- a/products/feature-flagging/feature-flagging-bootstrap/src/test/java/datadog/trace/api/featureflag/flagevaluation/FlagEvaluationMetricsTest.java +++ /dev/null @@ -1,110 +0,0 @@ -package datadog.trace.api.featureflag.flagevaluation; - -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_CLOSED; -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertTrue; - -import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Count; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.Future; -import java.util.concurrent.TimeUnit; -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Test; - -class FlagEvaluationMetricsTest { - private final FlagEvaluationMetrics metrics = FlagEvaluationMetrics.getInstance(); - - @BeforeEach - @AfterEach - void resetMetrics() { - metrics.drain(); - } - - @Test - void accumulatesPositiveLongCountsAndDrainsOnlyOnce() { - metrics.count(DROPPED_CLOSED, Integer.MAX_VALUE); - metrics.count(DROPPED_CLOSED, 5); - metrics.count(DROPPED_CLOSED, 0); - metrics.count(DROPPED_CLOSED, -1); - - List counts = metrics.drain(); - assertEquals(1, counts.size()); - assertEquals("flagevaluation.rows.dropped", counts.get(0).name); - assertEquals("reason:closed", counts.get(0).tag); - assertEquals(Integer.MAX_VALUE + 5L, counts.get(0).value); - assertTrue(metrics.drain().isEmpty()); - - metrics.count(DROPPED_CLOSED, 2); - assertEquals(2, metrics.drain().get(0).value); - } - - @Test - void preservesCombinedTruncationReasonsAcrossDrains() { - metrics.countContextTruncated("max_key_length,max_value_length", 2); - metrics.countContextTruncated("max_key_length,max_value_length", 3); - metrics.countContextTruncated("max_depth", 1); - metrics.countContextTruncated("ignored", 0); - metrics.countContextTruncated("ignored", -1); - - Map counts = new HashMap<>(); - for (Count count : metrics.drain()) { - assertEquals("flagevaluation.context.truncated", count.name); - counts.put(count.tag, count.value); - } - assertEquals(2, counts.size()); - assertEquals(5L, counts.get("reason:max_key_length,max_value_length").longValue()); - assertEquals(1L, counts.get("reason:max_depth").longValue()); - assertTrue(metrics.drain().isEmpty()); - - metrics.countContextTruncated("max_depth", 4); - assertEquals(4, metrics.drain().get(0).value); - } - - @Test - void concurrentProducersAndDrainingPreserveAllCounts() throws Exception { - ExecutorService executor = Executors.newFixedThreadPool(4); - CountDownLatch start = new CountDownLatch(1); - List> producers = new ArrayList<>(); - try { - for (int producer = 0; producer < 4; producer++) { - producers.add( - executor.submit( - () -> { - start.await(); - for (int i = 0; i < 5000; i++) { - metrics.count(DROPPED_CLOSED, 1); - metrics.countContextTruncated("max_depth", 1); - } - return null; - })); - } - start.countDown(); - Map totals = new HashMap<>(); - for (int i = 0; i < 100; i++) { - addTo(totals, metrics.drain()); - } - for (Future producer : producers) { - producer.get(10, TimeUnit.SECONDS); - } - addTo(totals, metrics.drain()); - assertEquals(20000L, totals.get("flagevaluation.rows.dropped").longValue()); - assertEquals(20000L, totals.get("flagevaluation.context.truncated").longValue()); - assertTrue(metrics.drain().isEmpty()); - } finally { - executor.shutdownNow(); - } - } - - private static void addTo(Map totals, List counts) { - for (Count count : counts) { - totals.merge(count.name, count.value, Long::sum); - } - } -} diff --git a/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java b/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java index 8236301d1d5..b618038216b 100644 --- a/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java +++ b/products/feature-flagging/feature-flagging-lib/src/main/java/com/datadog/featureflag/FlagEvaluationWriterImpl.java @@ -1,12 +1,12 @@ package com.datadog.featureflag; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DEGRADED_CARDINALITY_CAP; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DEGRADED_PAYLOAD_LIMIT; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_CLOSED; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_DEGRADED_CAP; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_PAYLOAD_LIMIT; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_QUEUE_OVERFLOW; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.PAYLOAD_SPLITS; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.DEGRADED_CARDINALITY_CAP; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.DEGRADED_PAYLOAD_LIMIT; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.DROPPED_CLOSED; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.DROPPED_DEGRADED_CAP; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.DROPPED_PAYLOAD_LIMIT; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.DROPPED_QUEUE_OVERFLOW; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.PAYLOAD_SPLITS; import static datadog.trace.util.AgentThreadFactory.AgentThread.FEATURE_FLAG_EVALUATION_PROCESSOR; import static datadog.trace.util.AgentThreadFactory.newAgentThread; import static java.util.concurrent.TimeUnit.SECONDS; @@ -19,8 +19,8 @@ import datadog.trace.api.Config; import datadog.trace.api.featureflag.FeatureFlaggingGateway; import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; -import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationWriter; +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector; import edu.umd.cs.findbugs.annotations.SuppressFBWarnings; import java.util.ArrayList; import java.util.List; @@ -66,7 +66,7 @@ public class FlagEvaluationWriterImpl implements FlagEvaluationWriter { static final int FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES = EvpProxy.PAYLOAD_SIZE_LIMIT_BYTES; private static final String FLAG_EVALUATION_ROUTE = "flagevaluation"; - private static final FlagEvaluationMetrics METRICS = FlagEvaluationMetrics.getInstance(); + private static final FlagEvaluationMetricCollector METRICS = FlagEvaluationMetricCollector.get(); private final MessagePassingBlockingQueue queue; private final FlagEvaluationSerializingHandler serializer; diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationTestSupport.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationTestSupport.java index 5ee7c21efd4..fdfb30880eb 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationTestSupport.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationTestSupport.java @@ -16,8 +16,9 @@ import datadog.communication.BackendApi; import datadog.communication.BackendApiFactory; import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; -import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; import datadog.trace.api.intake.Intake; +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector; +import datadog.trace.api.telemetry.MetricCollector; import java.lang.reflect.Type; import java.util.ArrayList; import java.util.Arrays; @@ -25,7 +26,6 @@ import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.Objects; import java.util.function.Supplier; import okhttp3.RequestBody; import okio.Buffer; @@ -48,7 +48,13 @@ static Supplier backendApiSupplier(final BackendApiFactory factory) } static void clearFlagEvaluationMetrics() { - FlagEvaluationMetrics.getInstance().drain(); + FlagEvaluationMetricCollector.get().resetForTesting(); + } + + static Collection collectFlagEvaluationMetrics() { + final FlagEvaluationMetricCollector collector = FlagEvaluationMetricCollector.get(); + collector.prepareMetrics(); + return collector.drain(); } static FlagEvalEvent event( @@ -160,14 +166,22 @@ static Map flushAndCaptureJson(final TestWriterSetup setup) thro } static long metricSum( - final Collection metrics, + final Collection metrics, final String metricName, final String tag) { long sum = 0; - for (final FlagEvaluationMetrics.Count metric : metrics) { - if (metricName.equals(metric.name) && Objects.equals(tag, metric.tag)) { - sum += metric.value; + for (final MetricCollector.Metric metric : metrics) { + if (!metricName.equals(metric.metricName)) { + continue; + } + if (tag == null) { + if (!metric.tags.isEmpty()) { + continue; + } + } else if (!metric.tags.contains(tag)) { + continue; } + sum += metric.value.longValue(); } return sum; } diff --git a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java index 1560e058bec..da699e56d31 100644 --- a/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java +++ b/products/feature-flagging/feature-flagging-lib/src/test/java/com/datadog/featureflag/FlagEvaluationWriterImplTest.java @@ -5,6 +5,7 @@ import static com.datadog.featureflag.FlagEvaluationTestSupport.buildTestWriter; import static com.datadog.featureflag.FlagEvaluationTestSupport.cfg; import static com.datadog.featureflag.FlagEvaluationTestSupport.clearFlagEvaluationMetrics; +import static com.datadog.featureflag.FlagEvaluationTestSupport.collectFlagEvaluationMetrics; import static com.datadog.featureflag.FlagEvaluationTestSupport.event; import static com.datadog.featureflag.FlagEvaluationTestSupport.eventForFlag; import static com.datadog.featureflag.FlagEvaluationTestSupport.flushAndCapture; @@ -40,9 +41,9 @@ import datadog.trace.api.Config; import datadog.trace.api.featureflag.FeatureFlaggingGateway; import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; -import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; import datadog.trace.api.featureflag.ufc.v1.ServerConfiguration; import datadog.trace.api.intake.Intake; +import datadog.trace.api.telemetry.MetricCollector; import datadog.trace.test.util.PollingConditions; import java.io.IOException; import java.lang.reflect.Field; @@ -89,8 +90,7 @@ void degradedCapOverflowTelemetryIsEmittedOnFlush() { setup.handler.addDroppedDegradedOverflowForTest(3); setup.handler.flush(); - final Collection metrics = - FlagEvaluationMetrics.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); assertEquals(3, metricSum(metrics, "flagevaluation.rows.dropped", "reason:degraded_cap")); } @@ -129,8 +129,7 @@ void queueOverflowIncrementsObservableDropCounter() { assertTrue(writer.droppedQueueOverflow() > 0); final long queueDrops = writer.droppedQueueOverflow(); writer.flushForTest(); - final Collection metrics = - FlagEvaluationMetrics.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); assertEquals( queueDrops, metricSum(metrics, "flagevaluation.rows.dropped", "reason:queue_overflow")); } @@ -147,8 +146,7 @@ void enqueueAfterCloseIsDroppedAndCounted() { writer.close(); writer.enqueue(simpleEvent("closed-flag", "on")); - final Collection metrics = - FlagEvaluationMetrics.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -167,8 +165,7 @@ void enqueueDisabledDropsAndCountsAsClosedDrop() { FeatureFlaggingGateway.setFlagEvaluationEnqueueEnabled(false); writer.enqueue(simpleEvent("disabled-flag", "on")); - final Collection metrics = - FlagEvaluationMetrics.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -189,8 +186,7 @@ void closeSweepsAndCountsEventsLeftInTheQueue() { writer.enqueue(simpleEvent("residual-flag-2", "on")); writer.close(); - final Collection metrics = - FlagEvaluationMetrics.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); assertEquals(2, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -333,8 +329,7 @@ void payloadLimitDropsAreCountedOnFlush() { setup.handler.drainAndAggregate(); setup.handler.flush(); - final Collection metrics = - FlagEvaluationMetrics.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:payload_limit")); } @@ -709,12 +704,11 @@ void countContextTruncatedAccumulatesPerReasonWithoutFlushingWriter() { writer.countContextTruncated("field_count"); writer.countContextTruncated("field_length"); - final Collection metrics = - FlagEvaluationMetrics.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); assertEquals(2, metricSum(metrics, "flagevaluation.context.truncated", "reason:field_count")); assertEquals(1, metricSum(metrics, "flagevaluation.context.truncated", "reason:field_length")); writer.flushForTest(); - assertTrue(FlagEvaluationMetrics.getInstance().drain().isEmpty()); + assertTrue(collectFlagEvaluationMetrics().isEmpty()); } @Test @@ -739,8 +733,7 @@ void hasCapacityForEnqueueReflectsQueueSaturationAndCountsPreQueueOverflow() { writer.countPreQueueOverflow(); writer.flushForTest(); - final Collection metrics = - FlagEvaluationMetrics.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:queue_overflow")); writer.close(); diff --git a/telemetry/build.gradle.kts b/telemetry/build.gradle.kts index 130a5bea872..f2efa690d6e 100644 --- a/telemetry/build.gradle.kts +++ b/telemetry/build.gradle.kts @@ -34,7 +34,6 @@ dependencies { implementation(libs.slf4j) implementation(project(":internal-api")) - implementation(project(":products:feature-flagging:feature-flagging-bootstrap")) compileOnly(project(":dd-java-agent:agent-tooling")) testImplementation(project(":dd-java-agent:agent-tooling")) diff --git a/telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java b/telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java index 4429b7bde2f..f6e4595de49 100644 --- a/telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java +++ b/telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java @@ -1,55 +1,14 @@ package datadog.telemetry.metric; -import static java.util.Collections.emptyIterator; -import static java.util.Collections.emptyList; - -import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; -import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Count; +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector; import datadog.trace.api.telemetry.MetricCollector; -import java.util.ArrayList; -import java.util.Collection; -import java.util.Iterator; -import java.util.List; -import java.util.concurrent.ArrayBlockingQueue; +import javax.annotation.Nonnull; -/** Adapts product-owned flag evaluation counters to telemetry on its collection schedule. */ public final class FlagEvaluationMetricPeriodicAction extends MetricPeriodicAction { - private final Collector collector = new Collector(); @Override - public MetricCollector collector() { - return collector; - } - - private static final class Collector implements MetricCollector { - private final FlagEvaluationMetrics metrics = FlagEvaluationMetrics.getInstance(); - private final ArrayBlockingQueue queue = new ArrayBlockingQueue<>(RAW_QUEUE_SIZE); - private Iterator pending = emptyIterator(); - - @Override - public void prepareMetrics() { - if (queue.remainingCapacity() == 0) { - return; - } - if (!pending.hasNext()) { - pending = metrics.drain().iterator(); - } - // Only the telemetry thread prepares metrics. Retain the rest of a snapshot when full; - // producer updates remain in the counters until the next snapshot can be collected. - while (queue.remainingCapacity() > 0 && pending.hasNext()) { - Count count = pending.next(); - queue.offer(new Metric("tracers", true, count.name, "count", count.value, count.tag)); - } - } - - @Override - public Collection drain() { - if (queue.isEmpty()) { - return emptyList(); - } - List drained = new ArrayList<>(queue.size()); - queue.drainTo(drained); - return drained; - } + @Nonnull + public MetricCollector collector() { + return FlagEvaluationMetricCollector.get(); } } diff --git a/telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java b/telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java index 0c63a5b57bf..d82af7e35dc 100644 --- a/telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java +++ b/telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java @@ -1,13 +1,9 @@ package datadog.telemetry.metric; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DEGRADED_CARDINALITY_CAP; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.DROPPED_CLOSED; -import static datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics.Metric.PAYLOAD_SPLITS; -import static datadog.trace.api.telemetry.MetricCollector.RAW_QUEUE_SIZE; -import static java.util.Collections.emptyList; +import static datadog.trace.api.telemetry.FlagEvaluationMetricCollector.Counter.DROPPED_CLOSED; 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.junit.jupiter.api.Assertions.assertSame; import static org.mockito.ArgumentCaptor.forClass; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; @@ -15,132 +11,41 @@ import datadog.telemetry.TelemetryService; import datadog.telemetry.api.Metric; -import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationMetrics; -import datadog.trace.api.telemetry.CoreMetricCollector; -import datadog.trace.api.telemetry.MetricCollector; -import java.util.ArrayList; -import java.util.Collection; -import java.util.List; +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; -import org.tabletest.junit.TableTest; class FlagEvaluationMetricPeriodicActionTest { - private final FlagEvaluationMetrics metrics = FlagEvaluationMetrics.getInstance(); + private final FlagEvaluationMetricCollector collector = FlagEvaluationMetricCollector.get(); private final FlagEvaluationMetricPeriodicAction action = new FlagEvaluationMetricPeriodicAction(); private final TelemetryService service = mock(TelemetryService.class); @BeforeEach @AfterEach - void resetMetrics() { - metrics.drain(); + void resetCollector() { + collector.resetForTesting(); } @Test - void defaultActionDrainsProductCountersIndependentlyOfCoreMetrics() { - metrics.count(DROPPED_CLOSED, 3); - CoreMetricCollector core = CoreMetricCollector.getInstance(); - core.prepareMetrics(); - assertTrue( - core.drain().stream().noneMatch(metric -> metric.metricName.startsWith("flagevaluation."))); - - action.collector().prepareMetrics(); - action.doIteration(service); - ArgumentCaptor captor = forClass(Metric.class); - verify(service).addMetric(captor.capture()); - assertEquals("flagevaluation.rows.dropped", captor.getValue().getMetric()); - assertEquals(3L, captor.getValue().getPoints().get(0).get(1).longValue()); - } + void emitsFlagEvaluationMetricsFromTheDedicatedCollector() { + assertSame(collector, action.collector()); + collector.count(DROPPED_CLOSED, 3); + collector.prepareMetrics(); - @TableTest({ - "Scenario | Counter | Name | Tag ", - "queue overflow | DROPPED_QUEUE_OVERFLOW | flagevaluation.rows.dropped | reason:queue_overflow ", - "closed | DROPPED_CLOSED | flagevaluation.rows.dropped | reason:closed ", - "degraded cap | DROPPED_DEGRADED_CAP | flagevaluation.rows.dropped | reason:degraded_cap ", - "payload drop | DROPPED_PAYLOAD_LIMIT | flagevaluation.rows.dropped | reason:payload_limit ", - "cardinality cap | DEGRADED_CARDINALITY_CAP | flagevaluation.rows.degraded | reason:cardinality_cap", - "payload degraded | DEGRADED_PAYLOAD_LIMIT | flagevaluation.rows.degraded | reason:payload_limit ", - "split | PAYLOAD_SPLITS | flagevaluation.payload.splits | " - }) - void preservesWireFormatAcrossCollectionIntervals( - FlagEvaluationMetrics.Metric counter, String name, String tag) { - metrics.count(counter, 2); - metrics.count(counter, 3); - assertTrue(action.collector().drain().isEmpty()); - action.collector().prepareMetrics(); - metrics.count(counter, 7); - action.collector().prepareMetrics(); action.doIteration(service); ArgumentCaptor captor = forClass(Metric.class); verify(service).addMetric(captor.capture()); + verifyNoMoreInteractions(service); Metric metric = captor.getValue(); - assertEquals(name, metric.getMetric()); + assertEquals("flagevaluation.rows.dropped", metric.getMetric()); assertEquals("tracers", metric.getNamespace()); assertEquals(true, metric.getCommon()); assertEquals(Metric.TypeEnum.COUNT, metric.getType()); - assertEquals(tag == null ? emptyList() : singletonList(tag), metric.getTags()); - assertEquals(2, metric.getPoints().size()); - assertEquals(5L, metric.getPoints().get(0).get(1).longValue()); - assertEquals(7L, metric.getPoints().get(1).get(1).longValue()); - action.collector().prepareMetrics(); - action.doIteration(service); - verifyNoMoreInteractions(service); - } - - @Test - void preservesContextTruncationTags() { - metrics.countContextTruncated("max_key_length,max_value_length", 3); - action.collector().prepareMetrics(); - action.doIteration(service); - - ArgumentCaptor captor = forClass(Metric.class); - verify(service).addMetric(captor.capture()); - Metric metric = captor.getValue(); - assertEquals("flagevaluation.context.truncated", metric.getMetric()); - assertEquals(singletonList("reason:max_key_length,max_value_length"), metric.getTags()); + assertEquals(singletonList("reason:closed"), metric.getTags()); assertEquals(3L, metric.getPoints().get(0).get(1).longValue()); } - - @Test - void fullQueueRetainsSnapshotRemainderAndSubsequentProducerUpdates() { - MetricCollector collector = action.collector(); - for (int i = 0; i < RAW_QUEUE_SIZE - 1; i++) { - metrics.count(PAYLOAD_SPLITS, 1); - collector.prepareMetrics(); - } - metrics.count(DROPPED_CLOSED, 5); - metrics.count(DEGRADED_CARDINALITY_CAP, 7); - metrics.countContextTruncated("max_depth", 3); - collector.prepareMetrics(); - metrics.count(DROPPED_CLOSED, 11); - collector.prepareMetrics(); - - List reported = new ArrayList<>(collector.drain()); - assertEquals(RAW_QUEUE_SIZE, reported.size()); - collector.prepareMetrics(); - Collection remainder = collector.drain(); - assertEquals(2, remainder.size()); - reported.addAll(remainder); - collector.prepareMetrics(); - reported.addAll(collector.drain()); - - assertEquals(RAW_QUEUE_SIZE - 1, sum(reported, "flagevaluation.payload.splits")); - assertEquals(16, sum(reported, "flagevaluation.rows.dropped")); - assertEquals(7, sum(reported, "flagevaluation.rows.degraded")); - assertEquals(3, sum(reported, "flagevaluation.context.truncated")); - collector.prepareMetrics(); - assertTrue(collector.drain().isEmpty()); - assertTrue(metrics.drain().isEmpty()); - } - - private static long sum(List metrics, String name) { - return metrics.stream() - .filter(metric -> name.equals(metric.metricName)) - .mapToLong(metric -> metric.value.longValue()) - .sum(); - } }