Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -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<FlagEvaluationMetricCollector.FlagEvaluationMetric> {
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<String, AtomicLong> contextTruncationCounts =
new ConcurrentHashMap<>();
private final BlockingQueue<FlagEvaluationMetric> 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<String, AtomicLong> 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<FlagEvaluationMetric> drain() {
if (metricsQueue.isEmpty()) {
return Collections.emptyList();
}
List<FlagEvaluationMetric> 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);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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']
}
}
Original file line number Diff line number Diff line change
@@ -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<FlagEvaluationMetric> 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<String, Long> 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<FlagEvaluationMetric> staged = new ArrayList<>(collector.drain());
assertEquals(RAW_QUEUE_SIZE, staged.size());
assertEquals(RAW_QUEUE_SIZE, sum(staged, "flagevaluation.payload.splits"));

List<FlagEvaluationMetric> 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<Future<?>> 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<String, Long> 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<FlagEvaluationMetric> collect() {
collector.prepareMetrics();
return new ArrayList<>(collector.drain());
}

private static Map<String, Long> valuesByTag(Collection<FlagEvaluationMetric> metrics) {
Map<String, Long> 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<String, Long> totals, Collection<FlagEvaluationMetric> metrics) {
for (FlagEvaluationMetric metric : metrics) {
totals.merge(metric.metricName, metric.value.longValue(), Long::sum);
}
}

private static long sum(Collection<FlagEvaluationMetric> metrics, String name) {
return metrics.stream()
.filter(metric -> name.equals(metric.metricName))
.mapToLong(metric -> metric.value.longValue())
.sum();
}
}
Loading
Loading