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/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/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/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-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..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,5 +1,12 @@ package com.datadog.featureflag; +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; @@ -13,12 +20,11 @@ import datadog.trace.api.featureflag.FeatureFlaggingGateway; import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; import datadog.trace.api.featureflag.flagevaluation.FlagEvaluationWriter; -import datadog.trace.api.telemetry.CoreMetricCollector; +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector; 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 FlagEvaluationMetricCollector METRICS = FlagEvaluationMetricCollector.get(); 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..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 @@ -17,7 +17,7 @@ import datadog.communication.BackendApiFactory; import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; import datadog.trace.api.intake.Intake; -import datadog.trace.api.telemetry.CoreMetricCollector; +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector; import datadog.trace.api.telemetry.MetricCollector; import java.lang.reflect.Type; import java.util.ArrayList; @@ -47,8 +47,14 @@ static Supplier backendApiSupplier(final BackendApiFactory factory) return () -> factory.createBackendApi(Intake.EVENT_PLATFORM, false); } - static void clearCoreMetrics() { - CoreMetricCollector.getInstance().drain(); + static void clearFlagEvaluationMetrics() { + FlagEvaluationMetricCollector.get().resetForTesting(); + } + + static Collection collectFlagEvaluationMetrics() { + final FlagEvaluationMetricCollector collector = FlagEvaluationMetricCollector.get(); + collector.prepareMetrics(); + return collector.drain(); } static FlagEvalEvent event( 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..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 @@ -4,7 +4,8 @@ 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.collectFlagEvaluationMetrics; import static com.datadog.featureflag.FlagEvaluationTestSupport.event; import static com.datadog.featureflag.FlagEvaluationTestSupport.eventForFlag; import static com.datadog.featureflag.FlagEvaluationTestSupport.flushAndCapture; @@ -42,7 +43,6 @@ import datadog.trace.api.featureflag.flagevaluation.FlagEvalEvent; 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; @@ -68,14 +68,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 +90,8 @@ 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 = collectFlagEvaluationMetrics(); + assertEquals(3, metricSum(metrics, "flagevaluation.rows.dropped", "reason:degraded_cap")); } @Test @@ -115,6 +109,7 @@ void startRegistersWriterAndCloseDeregistersIt() { writer.close(); writer.close(); writer.start(); + writer.startForTest(); assertNull(FeatureFlaggingGateway.getFlagEvalWriter()); } @@ -134,14 +129,9 @@ void queueOverflowIncrementsObservableDropCounter() { assertTrue(writer.droppedQueueOverflow() > 0); final long queueDrops = writer.droppedQueueOverflow(); writer.flushForTest(); - final Collection metrics = - CoreMetricCollector.getInstance().drain(); + final Collection metrics = collectFlagEvaluationMetrics(); 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,8 @@ 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 = collectFlagEvaluationMetrics(); + assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -181,14 +165,8 @@ 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 = collectFlagEvaluationMetrics(); + assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -208,14 +186,8 @@ 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 = collectFlagEvaluationMetrics(); + assertEquals(2, metricSum(metrics, "flagevaluation.rows.dropped", "reason:closed")); assertNull(writer.pollQueuedEventForTest()); } @@ -283,7 +255,6 @@ void flushIfNecessaryDoesNotReturnEarlyWhenOnlyQueueDropsArePending() { TimeUnit.NANOSECONDS, context(), queueDrops, - new java.util.concurrent.ConcurrentHashMap<>(), () -> {}, FlagEvaluationWriterImpl.FLAG_EVALUATION_PAYLOAD_SIZE_LIMIT_BYTES); @@ -314,7 +285,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 +329,8 @@ 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 = collectFlagEvaluationMetrics(); + assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:payload_limit")); } @Test @@ -728,7 +692,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 +703,12 @@ 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 = collectFlagEvaluationMetrics(); + assertEquals(2, metricSum(metrics, "flagevaluation.context.truncated", "reason:field_count")); + assertEquals(1, metricSum(metrics, "flagevaluation.context.truncated", "reason:field_length")); + writer.flushForTest(); + assertTrue(collectFlagEvaluationMetrics().isEmpty()); } @Test @@ -779,14 +733,8 @@ 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 = collectFlagEvaluationMetrics(); + assertEquals(1, metricSum(metrics, "flagevaluation.rows.dropped", "reason:queue_overflow")); writer.close(); } 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..f6e4595de49 --- /dev/null +++ b/telemetry/src/main/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicAction.java @@ -0,0 +1,14 @@ +package datadog.telemetry.metric; + +import datadog.trace.api.telemetry.FlagEvaluationMetricCollector; +import datadog.trace.api.telemetry.MetricCollector; +import javax.annotation.Nonnull; + +public final class FlagEvaluationMetricPeriodicAction extends MetricPeriodicAction { + + @Override + @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 new file mode 100644 index 00000000000..d82af7e35dc --- /dev/null +++ b/telemetry/src/test/java/datadog/telemetry/metric/FlagEvaluationMetricPeriodicActionTest.java @@ -0,0 +1,51 @@ +package datadog.telemetry.metric; + +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.assertSame; +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.telemetry.FlagEvaluationMetricCollector; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +class FlagEvaluationMetricPeriodicActionTest { + private final FlagEvaluationMetricCollector collector = FlagEvaluationMetricCollector.get(); + private final FlagEvaluationMetricPeriodicAction action = + new FlagEvaluationMetricPeriodicAction(); + private final TelemetryService service = mock(TelemetryService.class); + + @BeforeEach + @AfterEach + void resetCollector() { + collector.resetForTesting(); + } + + @Test + void emitsFlagEvaluationMetricsFromTheDedicatedCollector() { + assertSame(collector, action.collector()); + collector.count(DROPPED_CLOSED, 3); + collector.prepareMetrics(); + + action.doIteration(service); + + ArgumentCaptor captor = forClass(Metric.class); + verify(service).addMetric(captor.capture()); + verifyNoMoreInteractions(service); + Metric metric = captor.getValue(); + assertEquals("flagevaluation.rows.dropped", metric.getMetric()); + assertEquals("tracers", metric.getNamespace()); + assertEquals(true, metric.getCommon()); + assertEquals(Metric.TypeEnum.COUNT, metric.getType()); + assertEquals(singletonList("reason:closed"), metric.getTags()); + assertEquals(3L, metric.getPoints().get(0).get(1).longValue()); + } +}