diff --git a/products/metrics/metrics-api/build.gradle.kts b/products/metrics/metrics-api/build.gradle.kts
index 679b3aba56c..72f26a9d42a 100644
--- a/products/metrics/metrics-api/build.gradle.kts
+++ b/products/metrics/metrics-api/build.gradle.kts
@@ -1,15 +1,23 @@
plugins {
`java-library`
id("dd-trace-java.module.internal-api")
+ id("dd-trace-java.jmh-conventions")
}
description = "Metrics API"
dependencies {
implementation(libs.slf4j)
+ api(project(":components:environment"))
testImplementation(libs.bundles.junit5)
testImplementation(libs.bundles.mockito)
+ testImplementation(libs.jol.core)
+}
+
+jmh {
+ jmhVersion = libs.versions.jmh.get()
+ duplicateClassesStrategy = DuplicatesStrategy.EXCLUDE
}
extra["excludedClassesCoverage"] = listOf(
diff --git a/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorBenchmark.java b/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorBenchmark.java
new file mode 100644
index 00000000000..f5ab5147c1c
--- /dev/null
+++ b/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorBenchmark.java
@@ -0,0 +1,515 @@
+package datadog.metrics.api;
+
+import static java.util.concurrent.TimeUnit.MICROSECONDS;
+
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicLongArray;
+import java.util.concurrent.atomic.LongAdder;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Group;
+import org.openjdk.jmh.annotations.GroupThreads;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.Threads;
+import org.openjdk.jmh.annotations.Warmup;
+import org.openjdk.jmh.infra.Blackhole;
+
+/**
+ * {@link Accumulator} vs the alternatives it actually displaces: a single {@code LongAdder} (the
+ * collision-free baseline it can never beat, only approach), an independent {@code LongAdder} per
+ * counter guarded by a per-counter lock (the "just fix it with LongAdder" natural migration target
+ * -- {@code longAdderGroup*}), and the {@code ConcurrentHashMap.computeIfAbsent(key, k -> new
+ * AtomicLong())} anti-pattern ({@code chmAtomicLongIncrement*}) that {@link Accumulator} exists to
+ * avoid, and an unstriped {@code AtomicLongArray} ({@code atomicLongArray*}) -- one shared array,
+ * no per-thread distribution, isolating the cost of striping itself from the cost of correctness
+ * (see {@link #arrayAccumulateAndReset}). The CHM variant allocates its counter under the bucket's
+ * bin lock the first time its one constant key is seen, but since the map is a
+ * {@code @State(Scope.Benchmark)} field shared across the whole run, that allocation happens
+ * exactly once; every sampled op after it hits the warmed, already-present fast path. So this
+ * measures steady-state {@code computeIfAbsent} lookup overhead on an already-populated map, not
+ * the one-time allocation-under-lock cost -- still a useful number (a fixed, small key set that's
+ * allocated once and hit for the life of the process, as {@code WafMetricCollector}-style CHM
+ * counters are, spends nearly all its time in this same warmed path), just not the pathology the
+ * name of this benchmark might suggest.
+ *
+ *
{@code longAdderGroup*}: is a "just fix it with LongAdder" helper actually cheaper?
+ * {@code groupInc}/{@code groupAccumulateAnd} are the natural correct fix using {@code LongAdder}
+ * as the payload: one {@code LongAdder} per counter, with a per-counter lock guarding both
+ * the increment and the drain (locking only the drain does nothing -- {@code sumThenReset()}'s
+ * internal race is against the {@code LongAdder}'s own CAS-based {@code add()}, not against any
+ * lock a caller takes). The single-counter {@code *_*} benchmarks below collapse that per-counter
+ * lock to one lock shared by every thread -- the degenerate worst case for {@code longAdderGroup},
+ * with no thread-based distribution at all. The {@code *8_*} benchmarks fix that: each JMH worker
+ * thread is pinned to one of 8 counters for its lifetime (see {@link #threadCounterIndex}), so
+ * {@code longAdderGroup8}'s threads split into up to 8 groups each contending their own lock -- the
+ * topology where distributed locking should actually pay off, forcing {@link Accumulator}'s
+ * thread-striped design to earn its write-side win rather than facing a single-counter worst case.
+ * Fork(5), 15 samples per benchmark, Apple M1 Max, 10 CPUs - macOS/aarch64 - JDK 25 (Zulu):
+ * AccumulatorBenchmark.accumulatorAccumulateAndReset_highContention avgt 15 2.760 ± 0.052 us/op
+ * AccumulatorBenchmark.accumulatorAccumulateAndReset_lowContention avgt 15 0.049 ± 0.001 us/op
+ * AccumulatorBenchmark.accumulatorAccumulateAndReset8_highContention avgt 15 6.939 ± 0.303 us/op
+ * AccumulatorBenchmark.accumulatorAccumulateAndReset8_lowContention avgt 15 0.364 ± 0.003 us/op
+ * AccumulatorBenchmark.accumulatorIncrement_highContention avgt 15 0.009 ± 0.001 us/op
+ * AccumulatorBenchmark.accumulatorIncrement_lowContention avgt 15 0.007 ± 0.001 us/op
+ * AccumulatorBenchmark.accumulatorIncrement8_highContention avgt 15 0.016 ± 0.006 us/op
+ * AccumulatorBenchmark.accumulatorIncrement8_lowContention avgt 15 0.007 ± 0.001 us/op
+ * AccumulatorBenchmark.longAdderGroupAccumulateAnd_highContention avgt 15 1.703 ± 1.785 us/op
+ * AccumulatorBenchmark.longAdderGroupAccumulateAnd_lowContention avgt 15 0.024 ± 0.001 us/op
+ * AccumulatorBenchmark.longAdderGroupAccumulateAnd8_highContention avgt 15 5.989 ± 0.241 us/op
+ * AccumulatorBenchmark.longAdderGroupAccumulateAnd8_lowContention avgt 15 0.074 ± 0.008 us/op
+ * AccumulatorBenchmark.longAdderGroupIncrement_highContention avgt 15 2.775 ± 0.531 us/op
+ * AccumulatorBenchmark.longAdderGroupIncrement_lowContention avgt 15 0.012 ± 0.001 us/op
+ * AccumulatorBenchmark.longAdderGroupIncrement8_highContention avgt 15 0.513 ± 0.094 us/op
+ * AccumulatorBenchmark.longAdderGroupIncrement8_lowContention avgt 15 0.012 ± 0.001 us/op
+ * On the write side, {@link Accumulator} beats {@code longAdderGroup} at high contention by
+ * ~310x in the degenerate single-shared-lock case and still by ~32x once counters are fairly spread
+ * across 8 locks -- a large, reproducible win either way, on the call that runs on every event. On
+ * the drain side, the two designs remain close and the comparison stays noisy under contention for
+ * both: at width 1 {@link Accumulator}'s drain (2.760 us/op) reads slower than {@code
+ * longAdderGroup}'s (1.703 ± 1.785 us/op), but that error bar spans {@link Accumulator}'s own
+ * result, so the two aren't distinguishable at this sample size; at width 8 it's ~1.16x slower
+ * (6.939 vs 5.989 us/op), consistent with the earlier ~1.14x reading. Reproduced on a second,
+ * independent Fork(5) run after the stripe-count cap ({@code MAX_STRIPES = 64}) landed: a large,
+ * robust win on the call that fires on every event, and no confirmed cost on the call that fires
+ * once per reporting cycle.
+ *
+ *
{@code longAdderDelta*}: the fair single-counter baseline. {@code longAdderGroup}'s
+ * ~310x/~32x win above is real, but it's the cost of buying {@code sumThenReset()}'s no-lost-update
+ * guarantee via a lock -- not the cost of a single {@code LongAdder} on its own.
+ * Pre-migration {@code TracerHealthMetrics} never took that lock: it kept a single {@code
+ * previousCounts}/{@code countIndex}-tracked differ computing {@code sum() - previous} by hand,
+ * which is already lock-free and never loses an update, because a missed delta on one {@code sum()}
+ * just shows up whole on the next one (see {@link #deltaSumAndReset}). {@code longAdderDelta_*}/
+ * {@code longAdderDeltaMixed_*} below reproduce exactly that pattern -- one {@code LongAdder}, one
+ * dedicated differ, no lock anywhere -- as the fairest single-counter comparison to {@link
+ * Accumulator}: both designs give the same no-lost-update guarantee, just by different means
+ * (per-thread striping vs. a single differ's own unsynchronized bookkeeping), so neither pays for a
+ * lock the other doesn't need. Fork(5), 15 samples per benchmark, same machine:
+ * AccumulatorBenchmark.accumulatorAccumulateAndReset_lowContention avgt 15 0.050 ± 0.001 us/op
+ * AccumulatorBenchmark.accumulatorMixed avgt 15 0.158 ± 0.022 us/op
+ * AccumulatorBenchmark.accumulatorMixed:accumulatorMixed_drain avgt 15 0.692 ± 0.077 us/op
+ * AccumulatorBenchmark.accumulatorMixed:accumulatorMixed_write avgt 15 0.024 ± 0.009 us/op
+ * AccumulatorBenchmark.longAdderDelta_lowContention avgt 15 0.018 ± 0.001 us/op
+ * AccumulatorBenchmark.longAdderDeltaMixed avgt 15 0.118 ± 0.006 us/op
+ * AccumulatorBenchmark.longAdderDeltaMixed:longAdderDeltaMixed_drain avgt 15 0.122 ± 0.012 us/op
+ * AccumulatorBenchmark.longAdderDeltaMixed:longAdderDeltaMixed_write avgt 15 0.117 ± 0.006 us/op
+ * Here {@link Accumulator} does not win outright. In the single-threaded
+ * inc-then-diff-per-call shape, the safe {@code LongAdder} delta is ~2.8x cheaper (0.018 vs 0.050
+ * us/op) -- no striping to fan out or fold back in when there's only one thread. In the realistic
+ * many-writers/one-drainer topology ({@code accumulatorMixed} vs {@code longAdderDeltaMixed}),
+ * {@link Accumulator} wins the write side by ~4.9x (0.024 vs 0.117 us/op, the call on the hot path)
+ * but loses the drain side by ~5.7x (0.692 vs 0.122 us/op) and the combined total by ~1.34x (0.158
+ * vs 0.118 us/op) -- the striping that makes writes cheap has to be folded back together somewhere,
+ * and that fold costs more than one differ's plain subtraction.
+ *
+ *
That's not a mark against {@link Accumulator}: it's the expected shape of a primitive whose
+ * value isn't raw single-counter throughput. {@link Accumulator}'s real win shows up one level up,
+ * in a from-scratch before/after of an actual migrated caller -- {@code
+ * TracerHealthMetricsBenchmark} (see {@code
+ * dd-trace-core/src/jmh/java/datadog/trace/core/monitor/}), which replaced ~49 individual {@code
+ * LongAdder} fields and their hand-rolled {@code previousCounts}/{@code countIndex} delta tracking
+ * with one {@code Accumulator}. There, every hot single-counter call is at
+ * parity with or faster than the legacy code (0.6-0.7x of legacy on {@code onSend}, the most
+ * frequent real call site), and the batch drain -- the actual shape {@link Accumulator} is for,
+ * many counters read and reset together once per reporting cycle -- is where the real-code evidence
+ * is decisive, not just competitive. The fair single-counter numbers above are worth publishing
+ * precisely because they're not a clean win: they show {@link Accumulator} doesn't need to dominate
+ * every synthetic one-counter microbenchmark to be the right call once counters are plural and the
+ * drain is the thing that matters, which {@code TracerHealthMetricsBenchmark} demonstrates directly
+ * on real code rather than a synthetic stand-in.
+ */
+@State(Scope.Benchmark)
+@Warmup(iterations = 1, time = 10)
+@Measurement(iterations = 3, time = 10)
+@BenchmarkMode(Mode.AverageTime)
+@OutputTimeUnit(MICROSECONDS)
+@Fork(5)
+public class AccumulatorBenchmark {
+
+ enum Counter {
+ HITS
+ }
+
+ /**
+ * An 8-constant counterpart to {@link Counter}, used only by the {@code *8_*} benchmarks below.
+ * Unlike {@link Counter}, where every thread hits the single {@code HITS} constant (the worst
+ * case for {@code longAdderGroup}'s per-counter locking -- one lock shared by every thread,
+ * regardless of core count), these benchmarks spread writes across all 8 constants: each JMH
+ * worker thread is pinned to one fixed counter for its lifetime (see {@link
+ * #threadCounterIndex}), so under high contention, threads split into up to 8 groups each
+ * contending on their own lock instead of all threads sharing one. This is the topology where
+ * {@code longAdderGroup}'s distributed locking should actually pay off, and where {@link
+ * Accumulator}'s thread-striped design has to earn its win on the write side rather than facing a
+ * single-counter worst case. {@code accumulateAndReset}/{@code groupAccumulateAnd} also now walk
+ * 8 slots per drain instead of 1, sizing the drain cost closer to {@code TracerHealthMetric}'s
+ * 54-constant production shape.
+ */
+ enum Counter8 {
+ COUNTER_0,
+ COUNTER_1,
+ COUNTER_2,
+ COUNTER_3,
+ COUNTER_4,
+ COUNTER_5,
+ COUNTER_6,
+ COUNTER_7
+ }
+
+ private static final Counter8[] COUNTER8_VALUES = Counter8.values();
+
+ private final LongAdder adder = new LongAdder();
+ private final LongAdder deltaAdder = new LongAdder();
+ private long deltaPrevious;
+ private final Accumulator accumulator = Accumulator.of(Counter.class);
+ private final Accumulator accumulator8 = Accumulator.of(Counter8.class);
+ private final ConcurrentHashMap chm = new ConcurrentHashMap<>();
+ private final AtomicLongArray atomicLongArray = new AtomicLongArray(1);
+ private final AtomicLongArray atomicLongArray8 = new AtomicLongArray(8);
+ private final LongAdder[] longAdderGroup = {new LongAdder()};
+ private final LongAdder[] longAdderGroup8 = {
+ new LongAdder(),
+ new LongAdder(),
+ new LongAdder(),
+ new LongAdder(),
+ new LongAdder(),
+ new LongAdder(),
+ new LongAdder(),
+ new LongAdder()
+ };
+
+ /**
+ * Assigns each JMH worker thread a fixed {@code Counter8} index (round-robin over 8) the first
+ * time it calls into any {@code *8_*} benchmark, and keeps returning that same index for the
+ * thread's lifetime -- so under {@code Threads.MAX}, writes spread across all 8 counters instead
+ * of every thread hammering one.
+ */
+ private final AtomicInteger threadIndexAssigner = new AtomicInteger();
+
+ private final ThreadLocal threadCounterIndex =
+ ThreadLocal.withInitial(() -> threadIndexAssigner.getAndIncrement() % COUNTER8_VALUES.length);
+
+ /**
+ * The natural "just use LongAdder" fix for the reset hazard: one {@code LongAdder} per counter,
+ * with a per-counter lock guarding both the increment and the drain -- external locking around
+ * only the drain does nothing, since {@code sumThenReset()}'s internal race is against the {@code
+ * LongAdder}'s own CAS-based {@code add()}, not against any lock a caller takes. This is the fair
+ * comparison point: it closes the same reset hazard {@link Accumulator} does, but stripes by
+ * counter (one lock per enum constant) instead of by thread (one shared table
+ * across all counters) -- so N threads hammering the *same* counter contend on one lock
+ * regardless of core count, with no thread-bucket distribution at all.
+ */
+ private static void groupInc(LongAdder[] group, int ordinal) {
+ LongAdder counter = group[ordinal];
+ synchronized (counter) {
+ counter.add(1L);
+ }
+ }
+
+ private static long[] groupAccumulateAnd(LongAdder[] group) {
+ long[] acc = new long[group.length];
+ for (int i = 0; i < group.length; i++) {
+ LongAdder counter = group[i];
+ synchronized (counter) {
+ acc[i] = counter.sumThenReset();
+ }
+ }
+ return acc;
+ }
+
+ /**
+ * The unstriped baseline: a single shared {@code AtomicLongArray}, one slot per counter, with no
+ * per-thread distribution at all -- isolates the cost of {@link Accumulator}'s thread-striping
+ * itself from the cost of correctness (unlike {@code longAdderGroup}, this has no lock: {@code
+ * getAndAdd} and {@code getAndSet} are each already atomic per-slot, so no coordination is needed
+ * to give the same "no increment lost across a drain" guarantee).
+ */
+ private static long[] arrayAccumulateAndReset(AtomicLongArray array) {
+ long[] acc = new long[array.length()];
+ for (int i = 0; i < array.length(); i++) {
+ acc[i] = array.getAndSet(i, 0L);
+ }
+ return acc;
+ }
+
+ /**
+ * The safe, lock-free alternative to {@code sumThenReset()} that pre-migration {@code
+ * TracerHealthMetrics} actually used (via {@code previousCounts}/{@code countIndex}): never reset
+ * the {@code LongAdder} at all, and have a single differ thread track the last observed {@code
+ * sum()} to compute its own delta. {@code sum()} alone never loses an update permanently -- a
+ * miss just shows up in the next {@code sum()} -- so this closes {@code sumThenReset()}'s reset
+ * race without any lock, at the cost of one extra subtraction per drain. The catch is the "single
+ * differ" part: {@code deltaPrevious} is unsynchronized plain state, correct only because exactly
+ * one thread ever calls this method between increments. Unlike {@code sumThenReset()}, which
+ * degrades gracefully (just an occasional dropped delta) if called from multiple threads at once,
+ * concurrent callers here would race on {@code deltaPrevious} itself and corrupt it -- so there
+ * is deliberately no {@code longAdderDelta_highContention} mirroring {@code
+ * longAdderSumThenReset_highContention}'s "every thread both writes and drains" shape; see {@code
+ * longAdderDeltaMixed_write}/{@code _drain} below for the one topology (many writers, one
+ * dedicated drainer) this baseline is actually valid under.
+ */
+ private long deltaSumAndReset() {
+ long current = deltaAdder.sum();
+ long delta = current - deltaPrevious;
+ deltaPrevious = current;
+ return delta;
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void longAdderIncrement_lowContention() {
+ adder.increment();
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void longAdderIncrement_highContention() {
+ adder.increment();
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void accumulatorIncrement_lowContention() {
+ accumulator.inc(Counter.HITS);
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void accumulatorIncrement_highContention() {
+ accumulator.inc(Counter.HITS);
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void chmAtomicLongIncrement_lowContention() {
+ chm.computeIfAbsent("hits", k -> new AtomicLong()).incrementAndGet();
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void chmAtomicLongIncrement_highContention() {
+ chm.computeIfAbsent("hits", k -> new AtomicLong()).incrementAndGet();
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void longAdderSumThenReset_lowContention(Blackhole blackhole) {
+ adder.increment();
+ blackhole.consume(adder.sumThenReset());
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void longAdderSumThenReset_highContention(Blackhole blackhole) {
+ adder.increment();
+ blackhole.consume(adder.sumThenReset());
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void longAdderDelta_lowContention(Blackhole blackhole) {
+ deltaAdder.increment();
+ blackhole.consume(deltaSumAndReset());
+ }
+
+ /**
+ * The realistic, valid topology for {@link #deltaSumAndReset} -- many writers, one dedicated
+ * differ -- mirroring {@code accumulatorMixed_write}/{@code _drain} below so the two can be
+ * compared directly: this is the fairest single-counter match for {@link Accumulator}, since both
+ * are lock-free on the write side and both give the same no-lost-update guarantee, just by
+ * different means (per-slot atomic {@code getAndSet} vs. a single differ's own bookkeeping).
+ */
+ @Benchmark
+ @Group("longAdderDeltaMixed")
+ @GroupThreads(4)
+ public void longAdderDeltaMixed_write() {
+ deltaAdder.increment();
+ }
+
+ @Benchmark
+ @Group("longAdderDeltaMixed")
+ @GroupThreads(1)
+ public void longAdderDeltaMixed_drain(Blackhole blackhole) {
+ blackhole.consume(deltaSumAndReset());
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void accumulatorAccumulateAndReset_lowContention(Blackhole blackhole) {
+ accumulator.inc(Counter.HITS);
+ blackhole.consume(accumulator.accumulateAndReset());
+ }
+
+ /**
+ * A deliberately pessimistic topology: every thread both writes and drains on every op, so {@code
+ * Threads.MAX} threads are all draining concurrently. Real callers don't do this -- see {@code
+ * accumulatorMixed-write}/{@code accumulatorMixed-drain} below for the "many writers, one rare
+ * drainer" shape this class actually targets. Kept as the worst-case upper bound: no production
+ * topology should be more contended on {@link Accumulator#accumulateAndReset} than this.
+ */
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void accumulatorAccumulateAndReset_highContention(Blackhole blackhole) {
+ accumulator.inc(Counter.HITS);
+ blackhole.consume(accumulator.accumulateAndReset());
+ }
+
+ /**
+ * The realistic counterpart to {@code accumulatorAccumulateAndReset_highContention}: many writer
+ * threads incrementing, and a single dedicated thread polling {@link
+ * Accumulator#accumulateAndReset} -- not every thread doing both on every op. {@code
+ * accumulatorMixed-write} measures increment cost while a drain is actively running; {@code
+ * accumulatorMixed-drain} measures the drain's own cost under that same live write pressure. The
+ * 4:1 writer:drainer ratio is illustrative of "many writers, rare drain," not tuned to a specific
+ * core count.
+ */
+ @Benchmark
+ @Group("accumulatorMixed")
+ @GroupThreads(4)
+ public void accumulatorMixed_write() {
+ accumulator.inc(Counter.HITS);
+ }
+
+ @Benchmark
+ @Group("accumulatorMixed")
+ @GroupThreads(1)
+ public void accumulatorMixed_drain(Blackhole blackhole) {
+ blackhole.consume(accumulator.accumulateAndReset());
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void longAdderGroupIncrement_lowContention() {
+ groupInc(longAdderGroup, Counter.HITS.ordinal());
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void longAdderGroupIncrement_highContention() {
+ groupInc(longAdderGroup, Counter.HITS.ordinal());
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void longAdderGroupAccumulateAnd_lowContention(Blackhole blackhole) {
+ groupInc(longAdderGroup, Counter.HITS.ordinal());
+ blackhole.consume(groupAccumulateAnd(longAdderGroup));
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void longAdderGroupAccumulateAnd_highContention(Blackhole blackhole) {
+ groupInc(longAdderGroup, Counter.HITS.ordinal());
+ blackhole.consume(groupAccumulateAnd(longAdderGroup));
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void atomicLongArrayIncrement_lowContention() {
+ atomicLongArray.getAndAdd(Counter.HITS.ordinal(), 1L);
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void atomicLongArrayIncrement_highContention() {
+ atomicLongArray.getAndAdd(Counter.HITS.ordinal(), 1L);
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void atomicLongArrayAccumulateAndReset_lowContention(Blackhole blackhole) {
+ atomicLongArray.getAndAdd(Counter.HITS.ordinal(), 1L);
+ blackhole.consume(arrayAccumulateAndReset(atomicLongArray));
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void atomicLongArrayAccumulateAndReset_highContention(Blackhole blackhole) {
+ atomicLongArray.getAndAdd(Counter.HITS.ordinal(), 1L);
+ blackhole.consume(arrayAccumulateAndReset(atomicLongArray));
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void atomicLongArrayIncrement8_lowContention() {
+ atomicLongArray8.getAndAdd(threadCounterIndex.get(), 1L);
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void atomicLongArrayIncrement8_highContention() {
+ atomicLongArray8.getAndAdd(threadCounterIndex.get(), 1L);
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void atomicLongArrayAccumulateAndReset8_lowContention(Blackhole blackhole) {
+ atomicLongArray8.getAndAdd(threadCounterIndex.get(), 1L);
+ blackhole.consume(arrayAccumulateAndReset(atomicLongArray8));
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void atomicLongArrayAccumulateAndReset8_highContention(Blackhole blackhole) {
+ atomicLongArray8.getAndAdd(threadCounterIndex.get(), 1L);
+ blackhole.consume(arrayAccumulateAndReset(atomicLongArray8));
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void accumulatorIncrement8_lowContention() {
+ accumulator8.inc(COUNTER8_VALUES[threadCounterIndex.get()]);
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void accumulatorIncrement8_highContention() {
+ accumulator8.inc(COUNTER8_VALUES[threadCounterIndex.get()]);
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void accumulatorAccumulateAndReset8_lowContention(Blackhole blackhole) {
+ accumulator8.inc(COUNTER8_VALUES[threadCounterIndex.get()]);
+ blackhole.consume(accumulator8.accumulateAndReset());
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void accumulatorAccumulateAndReset8_highContention(Blackhole blackhole) {
+ accumulator8.inc(COUNTER8_VALUES[threadCounterIndex.get()]);
+ blackhole.consume(accumulator8.accumulateAndReset());
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void longAdderGroupIncrement8_lowContention() {
+ groupInc(longAdderGroup8, threadCounterIndex.get());
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void longAdderGroupIncrement8_highContention() {
+ groupInc(longAdderGroup8, threadCounterIndex.get());
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void longAdderGroupAccumulateAnd8_lowContention(Blackhole blackhole) {
+ groupInc(longAdderGroup8, threadCounterIndex.get());
+ blackhole.consume(groupAccumulateAnd(longAdderGroup8));
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void longAdderGroupAccumulateAnd8_highContention(Blackhole blackhole) {
+ groupInc(longAdderGroup8, threadCounterIndex.get());
+ blackhole.consume(groupAccumulateAnd(longAdderGroup8));
+ }
+}
diff --git a/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorVsCounterBenchmark.java b/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorVsCounterBenchmark.java
new file mode 100644
index 00000000000..46306a4fffc
--- /dev/null
+++ b/products/metrics/metrics-api/src/jmh/java/datadog/metrics/api/AccumulatorVsCounterBenchmark.java
@@ -0,0 +1,192 @@
+package datadog.metrics.api;
+
+import static java.util.concurrent.TimeUnit.MICROSECONDS;
+
+import datadog.metrics.api.statsd.StatsDClient;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.Threads;
+import org.openjdk.jmh.annotations.Warmup;
+
+/**
+ * A decision-support benchmark, not a design-proof one: {@link AccumulatorBenchmark} exists to
+ * justify {@link Accumulator}'s internal striping design against synthetic alternatives ({@code
+ * LongAdder}, a CHM of {@code AtomicLong}, an unstriped {@code AtomicLongArray}) that nobody in
+ * this codebase would actually reach for instead -- that question is settled and doesn't need
+ * re-litigating on every read. This class answers a different, durable one: given both {@link
+ * Counter} and {@link Accumulator} live in {@code metrics-api}, which do you actually use?
+ *
+ * {@link Counter} is the one most callers reach for first, and for good reason -- it's the
+ * advertised, general-purpose metrics API. But its real implementation ({@code StatsDCounter} in
+ * {@code metrics-lib}) calls {@link StatsDClient#count} synchronously on every {@link
+ * Counter#increment}, with no batching or striping of its own, and the real {@code StatsDClient}
+ * underneath serializes that call through one shared connection (a lock or an offer to a single
+ * queue) -- the same "one shared serialization point, no thread distribution" shape {@link
+ * AccumulatorBenchmark}'s {@code longAdderGroup} single-lock variants model. {@link
+ * #counterIncrement} below mirrors that shape with a {@link StatsDClient} that takes a real lock on
+ * every call (see {@link LockingStatsDClient}) rather than a true no-op -- a true no-op measures
+ * only virtual-dispatch overhead and understates {@code StatsDCounter}'s actual per-call cost to
+ * the point of not answering this benchmark's own question ({@code StatsDCounter} itself has
+ * package-private construction, so this reimplements its shape rather than depending on {@code
+ * metrics-lib}). {@link Accumulator} exists because that per-call cost is too high to pay on every
+ * request/span/event -- it stripes by thread and defers reporting to a periodic drain instead.
+ *
+ *
Rule of thumb: a counter incremented on a hot path (every request, span, or event)
+ * should use {@link Accumulator}, drained on a reporting cadence. A counter incremented rarely
+ * (startup, config changes, an error path already off the hot path) can use {@link Counter}
+ * directly, no ceremony required. This is the concrete case behind the (not yet built)
+ * {@code @ForegroundSafe}/{@code @BackgroundOnly} annotations: {@link Counter#increment} is
+ * {@code @BackgroundOnly}, {@link Accumulator#inc} is {@code @ForegroundSafe}.
+ *
+ *
Fork(5), 15 samples per benchmark, Apple M1 Max, 10 CPUs - macOS/aarch64 - JDK 25 (Zulu):
+ *
+ * AccumulatorVsCounterBenchmark.accumulatorIncrement_highContention avgt 15 0.009 ± 0.001 us/op
+ * AccumulatorVsCounterBenchmark.accumulatorIncrement_lowContention avgt 15 0.007 ± 0.001 us/op
+ * AccumulatorVsCounterBenchmark.counterIncrement_highContention avgt 15 1.559 ± 0.941 us/op
+ * AccumulatorVsCounterBenchmark.counterIncrement_lowContention avgt 15 0.009 ± 0.001 us/op
+ * At low contention the two are indistinguishable (0.007 vs 0.009 us/op) -- an uncontended
+ * lock costs almost nothing, so with only one thread ever calling in, {@link Counter}'s per-call
+ * cost and {@link Accumulator}'s are both dominated by the same handful of instructions. At high
+ * contention {@link Accumulator} wins by ~173x (0.009 vs 1.559 us/op, itself high-variance from run
+ * to run) -- {@link Counter}'s one shared lock serializes every calling thread, while {@link
+ * Accumulator}'s per-thread striping doesn't. This is the expected result for a counter hit from
+ * many concurrent threads, and it's why the Rule of thumb above exists.
+ *
+ *
What this does and doesn't measure. {@link LockingStatsDClient} is a conservative
+ * synthetic lower bound on {@link Counter}'s real cost, not a measurement of exact production
+ * overhead -- a real {@code StatsDClient} also encodes the metric line and offers it to a queue (or
+ * blocks on socket I/O) under that same lock/queue, work this stand-in skips entirely. That means
+ * the true gap between {@link Counter} and {@link Accumulator} at high contention in production is
+ * at least as large as the ~173x shown here, quite possibly larger; this benchmark establishes a
+ * floor, not a ceiling, on the win.
+ */
+@State(Scope.Benchmark)
+@Warmup(iterations = 1, time = 10)
+@Measurement(iterations = 3, time = 10)
+@BenchmarkMode(Mode.AverageTime)
+@OutputTimeUnit(MICROSECONDS)
+@Fork(5)
+public class AccumulatorVsCounterBenchmark {
+
+ enum Metric {
+ HITS
+ }
+
+ private final Accumulator accumulator = Accumulator.of(Metric.class);
+ private final Counter counter = new SynchronousStatsDCounter("hits", new LockingStatsDClient());
+
+ /**
+ * Mirrors {@code StatsDCounter}'s real shape: every {@link #increment} forwards straight to the
+ * client, with no batching of its own. Reimplemented here rather than depending on {@code
+ * metrics-lib} because {@code StatsDCounter}'s constructor is package-private.
+ */
+ private static final class SynchronousStatsDCounter implements Counter {
+ private static final String[] NO_TAGS = new String[0];
+ private final String name;
+ private final StatsDClient statsd;
+
+ SynchronousStatsDCounter(String name, StatsDClient statsd) {
+ this.name = name;
+ this.statsd = statsd;
+ }
+
+ @Override
+ public void increment(int delta) {
+ statsd.count(name, delta, NO_TAGS);
+ }
+
+ @Override
+ public void incrementErrorCount(String cause, int delta) {
+ statsd.count(name, delta, new String[] {"cause:" + cause});
+ }
+ }
+
+ /**
+ * A {@link StatsDClient} stand-in that pays a real, serialized per-call cost instead of a true
+ * no-op: every real {@code StatsDClient} forwards {@code count} through one shared connection (a
+ * lock or an offer to a single non-blocking queue), so a true no-op would measure only
+ * virtual-dispatch overhead and understate the cost this benchmark exists to isolate. A {@code
+ * synchronized} increment of a shared counter is a reasonable stand-in for that shared
+ * serialization point without pulling in a real socket/queue implementation this benchmark
+ * doesn't need.
+ */
+ private static final class LockingStatsDClient implements StatsDClient {
+ private final Object lock = new Object();
+ private long total;
+
+ @Override
+ public void incrementCounter(String metricName, String... tags) {
+ count(metricName, 1L, tags);
+ }
+
+ @Override
+ public void count(String metricName, long delta, String... tags) {
+ synchronized (lock) {
+ total += delta;
+ }
+ }
+
+ @Override
+ public void gauge(String metricName, long value, String... tags) {}
+
+ @Override
+ public void gauge(String metricName, double value, String... tags) {}
+
+ @Override
+ public void histogram(String metricName, long value, String... tags) {}
+
+ @Override
+ public void histogram(String metricName, double value, String... tags) {}
+
+ @Override
+ public void distribution(String metricName, long value, String... tags) {}
+
+ @Override
+ public void distribution(String metricName, double value, String... tags) {}
+
+ @Override
+ public void serviceCheck(
+ String serviceCheckName, String status, String message, String... tags) {}
+
+ @Override
+ public void error(Exception error) {}
+
+ @Override
+ public int getErrorCount() {
+ return 0;
+ }
+
+ @Override
+ public void close() {}
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void accumulatorIncrement_lowContention() {
+ accumulator.inc(Metric.HITS);
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void accumulatorIncrement_highContention() {
+ accumulator.inc(Metric.HITS);
+ }
+
+ @Benchmark
+ @Threads(1)
+ public void counterIncrement_lowContention() {
+ counter.increment(1);
+ }
+
+ @Benchmark
+ @Threads(Threads.MAX)
+ public void counterIncrement_highContention() {
+ counter.increment(1);
+ }
+}
diff --git a/products/metrics/metrics-api/src/main/java/datadog/metrics/api/Accumulator.java b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/Accumulator.java
new file mode 100644
index 00000000000..4f2c618534e
--- /dev/null
+++ b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/Accumulator.java
@@ -0,0 +1,288 @@
+package datadog.metrics.api;
+
+import datadog.environment.ThreadSupport;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicLongArray;
+import javax.annotation.concurrent.ThreadSafe;
+
+/**
+ * A striped, lock-free counter primitive keyed by enum ordinal: {@code LongAdder}'s write
+ * scalability, but as one shared, thread-sharded table instead of one independent {@code LongAdder}
+ * per counter -- which avoids paying {@code LongAdder}'s per-instance striping overhead {@code
+ * E.values().length} times over.
+ *
+ * {@code
+ * enum MyCounters { FOO, BAR }
+ *
+ * Accumulator counters = Accumulator.of(MyCounters.class);
+ * counters.inc(MyCounters.FOO);
+ * counters.add(MyCounters.BAR, 5L);
+ *
+ * Accumulator.Counts drained = counters.accumulateAndReset();
+ * long foo = drained.get(MyCounters.FOO);
+ * }
+ *
+ * Each counter's own {@link #accumulateAndReset} slot is read-and-zeroed with a single atomic
+ * {@code getAndSet}, so -- like {@code Accumulator}'s previous {@code synchronized}-stripe design,
+ * and unlike {@code LongAdder#sumThenReset()} -- no individual increment can land in the gap
+ * between summing and zeroing and be silently lost. What's gone is the previous design's
+ * row-wide atomicity: {@link #inc}/{@link #add} for two different counters are no longer
+ * guaranteed to be seen together by a concurrent {@link #accumulateAndReset}. There is no {@code
+ * update}-style escape hatch for grouping several counters under one atomic operation -- callers
+ * needing that must weigh whether the guarantee was load-bearing (most call sites are logging
+ * unrelated aspects of the same event, not maintaining a cross-counter invariant a reader depends
+ * on) or bring their own coordination.
+ */
+@ThreadSafe
+public final class Accumulator> {
+ /** One full cache line of {@code long}s (64 bytes), used to pad each stripe row. */
+ private static final int CACHE_LINE_LONGS = 8;
+
+ /** Upper bound on {@link #stripeCount()}, regardless of core count. */
+ private static final int MAX_STRIPES = 64;
+
+ private final AtomicLongArray[] data;
+ private final int width;
+ private final E[] keys;
+
+ private Accumulator(AtomicLongArray[] data, int width, E[] keys) {
+ this.data = data;
+ this.width = width;
+ this.keys = keys;
+ }
+
+ /**
+ * @param enumType the enum naming each counter, e.g. {@code MyCounters.class}
+ */
+ public static > Accumulator of(Class enumType) {
+ E[] keys = enumType.getEnumConstants();
+ int width = keys.length;
+ int paddedWidth = paddedWidth(width);
+ int stripes = stripeCount();
+ AtomicLongArray[] data = new AtomicLongArray[stripes];
+ for (int i = 0; i < stripes; i++) {
+ data[i] = new AtomicLongArray(paddedWidth);
+ }
+ return new Accumulator<>(data, width, keys);
+ }
+
+ /** Increments the counter named by {@code key} in the calling thread's stripe by one. */
+ public void inc(E key) {
+ add(key, 1L);
+ }
+
+ /** Adds {@code delta} to the counter named by {@code key} in the calling thread's stripe. */
+ public void add(E key, long delta) {
+ stripeOf(data).getAndAdd(key.ordinal(), delta);
+ }
+
+ /**
+ * Combines and resets every stripe, returning the sum as a typed view. Each counter is
+ * read-and-zeroed with one atomic {@code getAndSet} -- see the class-level note on what atomicity
+ * this does and doesn't provide across different counters.
+ *
+ * @return the sum, keyed by the enum's {@code ordinal()}
+ */
+ public Counts accumulateAndReset() {
+ long[] acc = new long[width];
+ for (AtomicLongArray stripe : data) {
+ for (int i = 0; i < width; i++) {
+ acc[i] += stripe.getAndSet(i, 0L);
+ }
+ }
+ return new Counts<>(acc, keys);
+ }
+
+ /**
+ * Combines every stripe without resetting it, returning the sum as a typed view -- a live,
+ * non-destructive snapshot for a diagnostic read (e.g. {@code summary()}) that must not perturb
+ * the delta a concurrent {@link #accumulateAndReset} on a reporting cadence is about to report.
+ *
+ * @return the sum, keyed by the enum's {@code ordinal()}
+ */
+ public Counts sum() {
+ long[] acc = new long[width];
+ for (AtomicLongArray stripe : data) {
+ for (int i = 0; i < width; i++) {
+ acc[i] += stripe.get(i);
+ }
+ }
+ return new Counts<>(acc, keys);
+ }
+
+ /**
+ * A typed view over a drained snapshot, returned by {@link #accumulateAndReset} or {@link #sum}:
+ * the same enum-ordinal type checking {@link Accumulator} provides on writes, applied to the read
+ * side too.
+ */
+ public static final class Counts> {
+ private final long[] counts;
+ private final E[] keys;
+
+ private Counts(long[] counts, E[] keys) {
+ this.counts = counts;
+ this.keys = keys;
+ }
+
+ /**
+ * An all-zero {@link Counts}, sized for {@code enumType} -- for seeding a running total before
+ * any real drain has happened, without needing a scratch {@link Accumulator} just to call
+ * {@link Accumulator#sum()} on it.
+ *
+ * @param enumType the enum naming each counter, e.g. {@code MyCounters.class}
+ */
+ public static > Counts create(Class enumType) {
+ E[] keys = enumType.getEnumConstants();
+ return new Counts<>(new long[keys.length], keys);
+ }
+
+ /** The counter named by {@code key}. */
+ public long get(E key) {
+ return counts[key.ordinal()];
+ }
+
+ /**
+ * The enum constants this {@link Counts} is keyed by, in declaration order -- for a caller that
+ * wants to iterate every counter (e.g. reporting each one) without separately having to pass
+ * {@code E.values()} alongside this object. An unmodifiable view over the same backing array
+ * every {@link Counts} from the same {@link Accumulator} shares -- no copy, but a caller can't
+ * corrupt that shared array's ordering for every other snapshot the way a raw array reference
+ * would let it.
+ */
+ public List keys() {
+ return Collections.unmodifiableList(Arrays.asList(keys));
+ }
+
+ /**
+ * Adds {@code other} to this, key by key, returning a new {@link Counts} rather than mutating
+ * either input -- combines a stored running total with a fresh, non-destructive {@link #sum} to
+ * answer "what's the live total right now" without ever resetting anything.
+ */
+ public Counts plus(Counts other) {
+ long[] combined = counts.clone();
+ for (int i = 0; i < combined.length; i++) {
+ combined[i] += other.counts[i];
+ }
+ return new Counts<>(combined, keys);
+ }
+
+ /**
+ * A copy of this {@link Counts} with every entry before {@code keys()[fromIndex]} zeroed out --
+ * for a caller that partially delivered a batch downstream and wants to represent "what's left"
+ * to retry, without re-counting the part that already made it.
+ */
+ public Counts from(int fromIndex) {
+ long[] remaining = new long[counts.length];
+ System.arraycopy(counts, fromIndex, remaining, fromIndex, counts.length - fromIndex);
+ return new Counts<>(remaining, keys);
+ }
+ }
+
+ /**
+ * An opt-in, safe-by-default composition for the one recurring shape {@link Accumulator} itself
+ * deliberately doesn't try to make safe on its own: a periodic destructive drain for reporting
+ * (e.g. to statsd) alongside a separate, non-destructive live read for diagnostics (e.g. a {@code
+ * summary()} or {@code toString()}). Read independently, {@link #drain} and {@link #live} are
+ * each individually correct -- the hazard is between them: a {@link #live} call landing between a
+ * {@link #drain}'s reset and its caller publishing the delta into a running total would combine
+ * the pre-publish total with the post-drain (zeroed) {@link #sum}, silently under-reporting by
+ * the just-drained delta. {@code RunningTotal} closes that gap with one lock shared by {@link
+ * #drain} and {@link #live} -- most callers don't need this (a plain drain-only reporter, or a
+ * live-read-only diagnostic, needs no coordination at all), so it's a separate, opt-in type
+ * rather than baked into every {@link Accumulator}.
+ */
+ @ThreadSafe
+ public static final class RunningTotal> {
+ private final Accumulator accumulator;
+ private final Object lock = new Object();
+ private Counts total;
+
+ private RunningTotal(Accumulator accumulator) {
+ this.accumulator = accumulator;
+ // accumulateAndReset(), not sum() -- sum() doesn't reset, so seeding from it would double
+ // count whatever's already there once a later live() adds a fresh sum() on top of it.
+ this.total = accumulator.accumulateAndReset();
+ }
+
+ /**
+ * Wraps {@code accumulator}, seeding the running total from a drain of whatever it currently
+ * holds. That seed becomes visible through {@link #live}, but -- being folded in at
+ * construction, before any caller can observe it -- it is never returned by a later {@link
+ * #drain}. A caller that reports each {@link #drain} delta (e.g. to statsd) would therefore
+ * silently never report pre-existing counts. Always construct this from a freshly created
+ * {@code accumulator} (as every current caller does) so there is nothing to seed.
+ */
+ public static > RunningTotal of(Accumulator accumulator) {
+ return new RunningTotal<>(accumulator);
+ }
+
+ /**
+ * Atomically drains the accumulator and folds the delta into the running total -- call this
+ * right before reporting the delta to a downstream sink on a reporting cadence.
+ *
+ * @return the delta just drained
+ */
+ public Counts drain() {
+ synchronized (lock) {
+ Counts delta = accumulator.accumulateAndReset();
+ total = total.plus(delta);
+ return delta;
+ }
+ }
+
+ /**
+ * The live total: the running total as of the last {@link #drain}, folded with whatever's
+ * accumulated since -- never resets anything, safe to call at any time without perturbing a
+ * concurrent {@link #drain}.
+ */
+ public Counts live() {
+ synchronized (lock) {
+ return total.plus(accumulator.sum());
+ }
+ }
+ }
+
+ /**
+ * The calling thread's stripe: cheap masking, no allocation, no map lookup.
+ *
+ * Multiple threads can map to the same stripe (this is masking, not a bijection); each
+ * counter's own atomic slot makes that safe, just not maximally scalable under a hash collision.
+ */
+ private static AtomicLongArray stripeOf(AtomicLongArray[] data) {
+ int mask = data.length - 1;
+ int idx = (int) (ThreadSupport.threadId() & mask);
+ return data[idx];
+ }
+
+ /**
+ * A fixed, power-of-two stripe count deliberately oversized to roughly 2x {@link
+ * Runtime#availableProcessors()} (minimum 4, maximum {@link #MAX_STRIPES}). Not exposed as a
+ * per-call override: a mandatory sizing knob on every caller fails the "print test" of
+ * self-explanatory API design.
+ *
+ *
Sizing to exactly the core count leaves stripe collisions likely under real contention
+ * (birthday-paradox math: with {@code n} contending threads and {@code m} stripes, expected
+ * colliding pairs are {@code n(n-1)/(2m)}) -- and a collision costs CAS-retry/cache-line-bounce
+ * cost, the same problem {@code LongAdder}'s own {@code Cell[]} table exists to avoid. Doubling
+ * the stripe count roughly halves that collision count for a one-time, per-accumulator memory
+ * cost, at the price of a slightly more expensive (but far rarer) {@link #accumulateAndReset}
+ * drain -- the right trade given {@link #inc}/{@link #add} run on every call while {@link
+ * #accumulateAndReset} runs on a reporting cadence.
+ *
+ *
Capped at {@link #MAX_STRIPES} so a single {@link Accumulator} on a very-high-core-count
+ * host can't grow its fixed table without bound: past that many contending threads, collision
+ * math has already flattened out, so the memory cost of chasing it further isn't worth paying.
+ */
+ private static int stripeCount() {
+ int cpus = Runtime.getRuntime().availableProcessors();
+ return Math.min(MAX_STRIPES, Math.max(4, 2 * Integer.highestOneBit(Math.max(1, cpus))));
+ }
+
+ /** Rounds {@code width} up to a whole number of cache lines, plus one full trailing line. */
+ private static int paddedWidth(int width) {
+ int wholeLines = ((width + CACHE_LINE_LONGS - 1) / CACHE_LINE_LONGS) * CACHE_LINE_LONGS;
+ return wholeLines + CACHE_LINE_LONGS;
+ }
+}
diff --git a/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCountReporter.java b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCountReporter.java
new file mode 100644
index 00000000000..b149b422f4c
--- /dev/null
+++ b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCountReporter.java
@@ -0,0 +1,122 @@
+package datadog.metrics.api.statsd;
+
+import datadog.metrics.api.Accumulator;
+import java.util.List;
+import java.util.function.ToLongFunction;
+import javax.annotation.concurrent.ThreadSafe;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Owns an {@link Accumulator} plus the periodic destructive-drain-and-report cycle over it: {@link
+ * #inc}/{@link #add} write straight through to the accumulator, {@link #live} peeks it for a
+ * diagnostic read, and {@link #flush} drains it and reports the delta to a {@link StatsDClient} --
+ * holding onto whatever wasn't delivered and retrying it on the next {@link #flush}, rather than
+ * losing it.
+ *
+ *
This isn't a real transaction: statsd delivery is already best-effort over UDP, so there's no
+ * rollback on the wire and no way to recover a packet actually lost in transit. It only protects
+ * against a local exception (a bad client, a malformed tag, a config error -- not a lost packet)
+ * turning "this flush interval under-reports" into "this delta is gone forever."
+ *
+ *
The undelivered remainder is kept as {@link #pending} state on this object, not fed back into
+ * the {@link Accumulator} -- {@link Accumulator.RunningTotal#drain} already folds every drained
+ * delta into its cumulative total unconditionally, before delivery is attempted, since that total
+ * counts real events, independent of whether statsd ever got them. Re-adding the same delta to the
+ * accumulator would make a later {@link Accumulator.RunningTotal#drain} see it as new and count it
+ * a second time -- {@link #pending} carries it forward for statsd's benefit only. {@link #flush} is
+ * assumed single-threaded (it's the body of one periodic scheduled task), so no locking guards it.
+ */
+@ThreadSafe
+public final class StatsDCountReporter & StatsDCounterKey> {
+ private static final Logger log = LoggerFactory.getLogger(StatsDCountReporter.class);
+
+ private final StatsDClient statsDClient;
+ private final Accumulator accumulator;
+ private final Accumulator.RunningTotal runningTotal;
+
+ /** Undelivered from the last {@link #flush}, to retry on the next one; {@code null} if none. */
+ private Accumulator.Counts pending;
+
+ /**
+ * @param enumType the enum naming each counter, e.g. {@code MyCounters.class}
+ */
+ public static & StatsDCounterKey> StatsDCountReporter of(
+ StatsDClient statsDClient, Class enumType) {
+ return new StatsDCountReporter<>(statsDClient, Accumulator.of(enumType));
+ }
+
+ private StatsDCountReporter(StatsDClient statsDClient, Accumulator accumulator) {
+ this.statsDClient = statsDClient;
+ this.accumulator = accumulator;
+ this.runningTotal = Accumulator.RunningTotal.of(accumulator);
+ }
+
+ /** Increments the counter named by {@code key} by one. */
+ public void inc(E key) {
+ accumulator.inc(key);
+ }
+
+ /** Adds {@code delta} to the counter named by {@code key}. */
+ public void add(E key, long delta) {
+ accumulator.add(key, delta);
+ }
+
+ /**
+ * The live total -- see {@link Accumulator.RunningTotal#live}. For a diagnostic read (e.g. a
+ * {@code summary()}), never for deciding what to report: {@link #flush} owns that.
+ */
+ public Accumulator.Counts live() {
+ return runningTotal.live();
+ }
+
+ /**
+ * Drains the accumulator and reports the delta (plus anything still owed from a prior failed
+ * attempt), holding onto whatever isn't confirmed delivered this time for the next {@link
+ * #flush}.
+ */
+ public void flush() {
+ Accumulator.Counts drained = runningTotal.drain();
+ Accumulator.Counts toReport = pending == null ? drained : pending.plus(drained);
+ pending = report(toReport);
+ }
+
+ /**
+ * @return {@code null} if every counter was delivered, otherwise a {@link Accumulator.Counts}
+ * holding whatever wasn't attempted or confirmed sent
+ */
+ private Accumulator.Counts report(Accumulator.Counts counts) {
+ List keys = counts.keys();
+ for (int i = 0; i < keys.size(); i++) {
+ E key = keys.get(i);
+ long delta = counts.get(key);
+ if (delta != 0) {
+ try {
+ statsDClient.count(key.getMetricName(), delta, key.getTags());
+ } catch (RuntimeException e) {
+ log.debug(
+ "Failed to report {}, compensating {} undelivered counter(s)",
+ key,
+ keys.size() - i,
+ e);
+ return counts.from(i);
+ }
+ }
+ }
+ return null;
+ }
+
+ /**
+ * Reports every nonzero entry of {@code values}/{@code counts} directly -- no accumulator, no
+ * compensation.
+ */
+ public static & StatsDCounterKey> void report(
+ StatsDClient statsDClient, E[] values, ToLongFunction counts) {
+ for (E value : values) {
+ long delta = counts.applyAsLong(value);
+ if (delta != 0) {
+ statsDClient.count(value.getMetricName(), delta, value.getTags());
+ }
+ }
+ }
+}
diff --git a/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCounterKey.java b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCounterKey.java
new file mode 100644
index 00000000000..3aadbfc24f5
--- /dev/null
+++ b/products/metrics/metrics-api/src/main/java/datadog/metrics/api/statsd/StatsDCounterKey.java
@@ -0,0 +1,8 @@
+package datadog.metrics.api.statsd;
+
+/** A counter identity: the dogstatsd metric name and tags a batch of counts should report under. */
+public interface StatsDCounterKey {
+ String getMetricName();
+
+ String[] getTags();
+}
diff --git a/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorFootprintTest.java b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorFootprintTest.java
new file mode 100644
index 00000000000..b0b2f7b59fd
--- /dev/null
+++ b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorFootprintTest.java
@@ -0,0 +1,140 @@
+package datadog.metrics.api;
+
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assumptions.assumeFalse;
+
+import datadog.environment.JavaVirtualMachine;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.LongAdder;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.openjdk.jol.info.GraphLayout;
+
+/**
+ * Retained-footprint comparison (JOL) for {@link Accumulator} vs the alternative it actually
+ * displaces: one {@code LongAdder} per counter (there is no multi-counter {@code LongAdder} -- a
+ * caller wanting N additive counters allocates N of them, one field each, as {@code OtlpTelemetry}
+ * and {@code PayloadDispatcherImpl} do today).
+ *
+ * A freshly constructed {@code LongAdder} is nearly free -- it holds no {@code Cell[]} table
+ * until contention forces one -- so comparing fresh instances understates its real cost and
+ * flatters {@code LongAdder}. {@link Accumulator} pays its full striped array up front, at
+ * creation, sized for {@link Runtime#availableProcessors()} regardless of whether contention ever
+ * materializes. The realistic comparison is therefore not fresh-vs-fresh but contended-vs-fresh:
+ * what each actually costs once the counters they represent are hit by real concurrent writers, as
+ * they are on the telemetry paths this class targets.
+ *
+ *
{@link Accumulator}'s stripe count is fixed at creation (roughly 2x {@link
+ * Runtime#availableProcessors()}, minimum 4) and does not grow further as more contention arrives
+ * within it, while every additional concurrently-written {@code LongAdder} keeps paying its own
+ * {@code Cell[]} growth cost independently. {@code Accumulator}'s up-front cost is the more
+ * predictable one: fixed at creation, independent of runtime contention, and shared (one striped
+ * table) across however many counters the caller's enum declares, rather than paid per counter. The
+ * printed numbers below vary by run/JVM -- see the assertions for the invariants that actually
+ * matter.
+ */
+class AccumulatorFootprintTest {
+
+ enum Counters {
+ REQUESTS,
+ ERRORS,
+ RETRIES,
+ BYTES_SENT
+ }
+
+ @BeforeAll
+ static void assumeNotJ9Jvm() {
+ // JOL's GraphLayout relies on HotSpot-specific Unsafe internals and throws
+ // IllegalStateException on J9-based JVMs (IBM/Semeru) -- same guard as
+ // StringIndexFootprintTest / ScopeAndContinuationLayoutTest.
+ assumeFalse(JavaVirtualMachine.isJ9());
+ }
+
+ static long bytes(Object root) {
+ return GraphLayout.parseInstance(root).totalSize();
+ }
+
+ static LongAdder[] freshAdders() {
+ LongAdder[] adders = new LongAdder[Counters.values().length];
+ for (int i = 0; i < adders.length; i++) {
+ adders[i] = new LongAdder();
+ }
+ return adders;
+ }
+
+ @Test
+ void freshFootprint() {
+ LongAdder[] adders = freshAdders();
+ Accumulator accumulator = Accumulator.of(Counters.class);
+
+ long adderBytes = bytes((Object) adders);
+ long accumulatorBytes = bytes(accumulator);
+
+ System.out.printf(
+ "fresh: %d LongAdders = %6d bytes, Accumulator = %6d bytes%n",
+ adders.length, adderBytes, accumulatorBytes);
+ }
+
+ /**
+ * Drives real multi-threaded contention against a fresh set of {@code LongAdder}s to force their
+ * {@code Cell[]} tables to grow, then compares against {@link Accumulator}'s fixed footprint --
+ * the realistic comparison, since production callers write to these counters concurrently rather
+ * than leaving them untouched.
+ *
+ * Cell-table growth is driven by JVM-internal CAS-collision detection, not something this test
+ * controls directly, so the exact grown size can vary by run/JVM; the one invariant asserted is
+ * monotonic growth (a contended footprint can only be at least the fresh one).
+ */
+ @Test
+ void contendedFootprint() throws InterruptedException {
+ LongAdder[] adders = freshAdders();
+ long freshAdderBytes = bytes((Object) adders);
+
+ int threads = Math.min(16, Math.max(4, Runtime.getRuntime().availableProcessors()));
+ ExecutorService pool = Executors.newFixedThreadPool(threads);
+ CountDownLatch start = new CountDownLatch(1);
+ CountDownLatch done = new CountDownLatch(threads);
+ try {
+ for (int t = 0; t < threads; t++) {
+ pool.execute(
+ () -> {
+ try {
+ start.await();
+ long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(2);
+ while (System.nanoTime() < deadline) {
+ for (LongAdder adder : adders) {
+ adder.increment();
+ }
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ } finally {
+ done.countDown();
+ }
+ });
+ }
+ start.countDown();
+ assertTrue(done.await(30, TimeUnit.SECONDS));
+ } finally {
+ pool.shutdown();
+ }
+
+ long contendedAdderBytes = bytes((Object) adders);
+ Accumulator accumulator = Accumulator.of(Counters.class);
+ long accumulatorBytes = bytes(accumulator);
+
+ System.out.printf(
+ "contended: %d LongAdders = %6d bytes, Accumulator = %6d bytes%n",
+ adders.length, contendedAdderBytes, accumulatorBytes);
+
+ assertTrue(
+ contendedAdderBytes >= freshAdderBytes,
+ "contended LongAdder footprint should never shrink below the fresh footprint");
+ assertTrue(
+ accumulatorBytes < contendedAdderBytes,
+ "Accumulator's fixed footprint should be smaller than N contended LongAdders");
+ }
+}
diff --git a/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorTest.java b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorTest.java
new file mode 100644
index 00000000000..adb203fdd1c
--- /dev/null
+++ b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/AccumulatorTest.java
@@ -0,0 +1,296 @@
+package datadog.metrics.api;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicBoolean;
+import org.junit.jupiter.api.Test;
+
+class AccumulatorTest {
+
+ enum Counters {
+ FOO,
+ BAR,
+ BAZ
+ }
+
+ @Test
+ void freshAccumulatorSumsToZero() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ for (Counters c : Counters.values()) {
+ assertEquals(0L, drained.get(c));
+ }
+ }
+
+ @Test
+ void incIncrementsByOne() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+ counters.inc(Counters.FOO);
+ counters.inc(Counters.BAR);
+
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ assertEquals(2L, drained.get(Counters.FOO));
+ assertEquals(1L, drained.get(Counters.BAR));
+ assertEquals(0L, drained.get(Counters.BAZ));
+ }
+
+ @Test
+ void addAppliesArbitraryDelta() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.add(Counters.BAZ, 41L);
+ counters.add(Counters.BAZ, 1L);
+
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ assertEquals(42L, drained.get(Counters.BAZ));
+ }
+
+ @Test
+ void accumulateAndResetsSoASecondDrainIsZero() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+
+ Accumulator.Counts first = counters.accumulateAndReset();
+ assertEquals(1L, first.get(Counters.FOO));
+
+ Accumulator.Counts second = counters.accumulateAndReset();
+ for (Counters c : Counters.values()) {
+ assertEquals(0L, second.get(c));
+ }
+ }
+
+ @Test
+ void concurrentIncrementsAreNotLost() throws InterruptedException {
+ Accumulator counters = Accumulator.of(Counters.class);
+ int threadCount = 16;
+ int incrementsPerThread = 10_000;
+
+ ExecutorService pool = Executors.newFixedThreadPool(threadCount);
+ CountDownLatch start = new CountDownLatch(1);
+ CountDownLatch done = new CountDownLatch(threadCount);
+ try {
+ for (int t = 0; t < threadCount; t++) {
+ pool.execute(
+ () -> {
+ try {
+ start.await();
+ for (int i = 0; i < incrementsPerThread; i++) {
+ counters.inc(Counters.FOO);
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ } finally {
+ done.countDown();
+ }
+ });
+ }
+ start.countDown();
+ assertTrue(done.await(30, TimeUnit.SECONDS));
+ } finally {
+ pool.shutdown();
+ }
+
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ assertEquals((long) threadCount * incrementsPerThread, drained.get(Counters.FOO));
+ }
+
+ @Test
+ void concurrentAccumulateAndDuringWritesNeverExceedsWritten()
+ throws InterruptedException, ExecutionException, TimeoutException {
+ Accumulator counters = Accumulator.of(Counters.class);
+ int threadCount = 8;
+ int incrementsPerThread = 5_000;
+
+ ExecutorService pool = Executors.newFixedThreadPool(threadCount + 1);
+ CountDownLatch done = new CountDownLatch(threadCount);
+ AtomicBoolean stop = new AtomicBoolean(false);
+ long[] runningTotal = {0L};
+
+ try {
+ Future> drainer =
+ pool.submit(
+ () -> {
+ while (!stop.get()) {
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ synchronized (runningTotal) {
+ runningTotal[0] += drained.get(Counters.FOO);
+ }
+ }
+ });
+
+ for (int t = 0; t < threadCount; t++) {
+ pool.execute(
+ () -> {
+ for (int i = 0; i < incrementsPerThread; i++) {
+ counters.inc(Counters.FOO);
+ }
+ done.countDown();
+ });
+ }
+
+ assertTrue(done.await(30, TimeUnit.SECONDS));
+ stop.set(true);
+ drainer.get(30, TimeUnit.SECONDS);
+
+ Accumulator.Counts finalDrain = counters.accumulateAndReset();
+ synchronized (runningTotal) {
+ runningTotal[0] += finalDrain.get(Counters.FOO);
+ }
+
+ assertEquals((long) threadCount * incrementsPerThread, runningTotal[0]);
+ } finally {
+ pool.shutdown();
+ }
+ }
+
+ @Test
+ void sumDoesNotResetStripes() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+
+ Accumulator.Counts first = counters.sum();
+ assertEquals(1L, first.get(Counters.FOO));
+
+ // sum() didn't reset anything, so a second sum() sees the same total
+ Accumulator.Counts second = counters.sum();
+ assertEquals(1L, second.get(Counters.FOO));
+
+ // and a real drain afterwards still sees the value sum() didn't consume
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ assertEquals(1L, drained.get(Counters.FOO));
+ }
+
+ @Test
+ void sumReflectsIncrementsMadeAfterAnEarlierSum() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+ counters.sum();
+
+ counters.inc(Counters.FOO);
+ Accumulator.Counts second = counters.sum();
+ assertEquals(2L, second.get(Counters.FOO));
+ }
+
+ @Test
+ void createSeedsAnAllZeroCountsWithoutAScratchAccumulator() {
+ Accumulator.Counts zero = Accumulator.Counts.create(Counters.class);
+ assertEquals(0L, zero.get(Counters.FOO));
+ assertEquals(0L, zero.get(Counters.BAR));
+
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+
+ Accumulator.Counts live = zero.plus(counters.sum());
+ assertEquals(1L, live.get(Counters.FOO));
+ }
+
+ @Test
+ void countsExposesItsOwnKeysWithoutASeparateValuesArray() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+ counters.add(Counters.BAR, 5L);
+
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ assertEquals(Counters.values().length, drained.keys().size());
+
+ long total = 0L;
+ for (Counters c : drained.keys()) {
+ total += drained.get(c);
+ }
+ assertEquals(6L, total);
+ }
+
+ @Test
+ void keysIsUnmodifiable() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ assertThrows(UnsupportedOperationException.class, () -> drained.keys().set(0, Counters.BAZ));
+ }
+
+ @Test
+ void plusCombinesAStoredRunningTotalWithALiveSumWithoutMutatingEither() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+ counters.add(Counters.BAR, 5L);
+
+ // drain once, e.g. as if a reporting cycle already ran and stored this total
+ Accumulator.Counts storedTotal = counters.accumulateAndReset();
+
+ // more activity happens after that drain, before the next one
+ counters.inc(Counters.FOO);
+
+ Accumulator.Counts live = storedTotal.plus(counters.sum());
+ assertEquals(2L, live.get(Counters.FOO));
+ assertEquals(5L, live.get(Counters.BAR));
+
+ // neither input was mutated by combining them
+ assertEquals(1L, storedTotal.get(Counters.FOO));
+ assertEquals(1L, counters.sum().get(Counters.FOO));
+ }
+
+ @Test
+ void fromZeroesEveryEntryBeforeTheGivenIndex() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+ counters.add(Counters.BAR, 4L);
+ counters.add(Counters.BAZ, 2L);
+
+ Accumulator.Counts drained = counters.accumulateAndReset();
+ Accumulator.Counts remaining = drained.from(Counters.BAR.ordinal());
+
+ assertEquals(0L, remaining.get(Counters.FOO));
+ assertEquals(4L, remaining.get(Counters.BAR));
+ assertEquals(2L, remaining.get(Counters.BAZ));
+
+ // the original Counts is untouched by taking a remainder from it
+ assertEquals(1L, drained.get(Counters.FOO));
+ }
+
+ @Test
+ void runningTotalSeedsFromTheAccumulatorsCurrentSum() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ counters.inc(Counters.FOO);
+
+ Accumulator.RunningTotal runningTotal = Accumulator.RunningTotal.of(counters);
+ assertEquals(1L, runningTotal.live().get(Counters.FOO));
+ }
+
+ @Test
+ void runningTotalDrainReturnsTheDeltaAndFoldsItIntoTheTotal() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ Accumulator.RunningTotal runningTotal = Accumulator.RunningTotal.of(counters);
+
+ counters.inc(Counters.FOO);
+ Accumulator.Counts delta = runningTotal.drain();
+ assertEquals(1L, delta.get(Counters.FOO));
+ assertEquals(1L, runningTotal.live().get(Counters.FOO));
+
+ counters.inc(Counters.FOO);
+ Accumulator.Counts secondDelta = runningTotal.drain();
+ assertEquals(1L, secondDelta.get(Counters.FOO));
+ assertEquals(2L, runningTotal.live().get(Counters.FOO));
+ }
+
+ @Test
+ void runningTotalLiveReflectsActivitySinceTheLastDrainWithoutDraining() {
+ Accumulator counters = Accumulator.of(Counters.class);
+ Accumulator.RunningTotal runningTotal = Accumulator.RunningTotal.of(counters);
+
+ counters.inc(Counters.FOO);
+ runningTotal.drain();
+
+ counters.inc(Counters.FOO);
+ assertEquals(2L, runningTotal.live().get(Counters.FOO));
+ // live() didn't drain anything, so a real drain afterwards still sees the pending increment
+ assertEquals(1L, runningTotal.drain().get(Counters.FOO));
+ }
+}
diff --git a/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/RecordingStatsDClient.java b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/RecordingStatsDClient.java
new file mode 100644
index 00000000000..f017a134086
--- /dev/null
+++ b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/RecordingStatsDClient.java
@@ -0,0 +1,63 @@
+package datadog.metrics.api.statsd;
+
+import java.util.ArrayList;
+import java.util.List;
+
+/** Test fake that records every {@link #count} call for assertion. */
+final class RecordingStatsDClient implements StatsDClient {
+
+ static final class Count {
+ final String metricName;
+ final long delta;
+ final String[] tags;
+
+ Count(String metricName, long delta, String[] tags) {
+ this.metricName = metricName;
+ this.delta = delta;
+ this.tags = tags;
+ }
+ }
+
+ final List counts = new ArrayList<>();
+
+ @Override
+ public void incrementCounter(String metricName, String... tags) {}
+
+ @Override
+ public void count(String metricName, long delta, String... tags) {
+ counts.add(new Count(metricName, delta, tags));
+ }
+
+ @Override
+ public void gauge(String metricName, long value, String... tags) {}
+
+ @Override
+ public void gauge(String metricName, double value, String... tags) {}
+
+ @Override
+ public void histogram(String metricName, long value, String... tags) {}
+
+ @Override
+ public void histogram(String metricName, double value, String... tags) {}
+
+ @Override
+ public void distribution(String metricName, long value, String... tags) {}
+
+ @Override
+ public void distribution(String metricName, double value, String... tags) {}
+
+ @Override
+ public void serviceCheck(
+ String serviceCheckName, String status, String message, String... tags) {}
+
+ @Override
+ public void error(Exception error) {}
+
+ @Override
+ public int getErrorCount() {
+ return 0;
+ }
+
+ @Override
+ public void close() {}
+}
diff --git a/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/StatsDCountReporterTest.java b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/StatsDCountReporterTest.java
new file mode 100644
index 00000000000..9324fa62aa3
--- /dev/null
+++ b/products/metrics/metrics-api/src/test/java/datadog/metrics/api/statsd/StatsDCountReporterTest.java
@@ -0,0 +1,207 @@
+package datadog.metrics.api.statsd;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import datadog.metrics.api.Accumulator;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicBoolean;
+import org.junit.jupiter.api.Test;
+
+class StatsDCountReporterTest {
+
+ private static final String[] TAG_A = {"env:a"};
+ private static final String[] TAG_B = {"env:b"};
+
+ enum Counters implements StatsDCounterKey {
+ FOO("foo.total", TAG_A),
+ BAR("bar.total", TAG_A),
+ SHARED_A("shared.total", TAG_A),
+ SHARED_B("shared.total", TAG_B);
+
+ private final String metricName;
+ private final String[] tags;
+
+ Counters(String metricName, String[] tags) {
+ this.metricName = metricName;
+ this.tags = tags;
+ }
+
+ @Override
+ public String getMetricName() {
+ return metricName;
+ }
+
+ @Override
+ public String[] getTags() {
+ return tags;
+ }
+ }
+
+ @Test
+ void reportsNonZeroCounterWithItsOwnMetricNameAndTags() {
+ RecordingStatsDClient statsD = new RecordingStatsDClient();
+ Map deltas = new HashMap<>();
+ deltas.put(Counters.FOO, 3L);
+
+ StatsDCountReporter.report(statsD, Counters.values(), c -> deltas.getOrDefault(c, 0L));
+
+ assertEquals(1, statsD.counts.size());
+ RecordingStatsDClient.Count count = statsD.counts.get(0);
+ assertEquals("foo.total", count.metricName);
+ assertEquals(3L, count.delta);
+ assertArrayEquals(TAG_A, count.tags);
+ }
+
+ @Test
+ void skipsZeroDeltaCounters() {
+ RecordingStatsDClient statsD = new RecordingStatsDClient();
+
+ StatsDCountReporter.report(statsD, Counters.values(), c -> 0L);
+
+ assertTrue(statsD.counts.isEmpty());
+ }
+
+ @Test
+ void reportsNothingWhenEveryCounterIsZero() {
+ RecordingStatsDClient statsD = new RecordingStatsDClient();
+ Map deltas = new HashMap<>();
+
+ StatsDCountReporter.report(statsD, Counters.values(), c -> deltas.getOrDefault(c, 0L));
+
+ assertTrue(statsD.counts.isEmpty());
+ }
+
+ @Test
+ void reportsConstantsSharingAMetricNameIndependentlyByTag() {
+ RecordingStatsDClient statsD = new RecordingStatsDClient();
+ Map deltas = new HashMap<>();
+ deltas.put(Counters.SHARED_A, 5L);
+ deltas.put(Counters.SHARED_B, 7L);
+
+ StatsDCountReporter.report(statsD, Counters.values(), c -> deltas.getOrDefault(c, 0L));
+
+ assertEquals(2, statsD.counts.size());
+ RecordingStatsDClient.Count a = statsD.counts.get(0);
+ RecordingStatsDClient.Count b = statsD.counts.get(1);
+ assertEquals("shared.total", a.metricName);
+ assertEquals(5L, a.delta);
+ assertArrayEquals(TAG_A, a.tags);
+ assertEquals("shared.total", b.metricName);
+ assertEquals(7L, b.delta);
+ assertArrayEquals(TAG_B, b.tags);
+ }
+
+ @Test
+ void flushDrainsAndReportsWithoutPerturbingTheCumulativeLiveTotal() {
+ RecordingStatsDClient statsD = new RecordingStatsDClient();
+ StatsDCountReporter reporter = StatsDCountReporter.of(statsD, Counters.class);
+ reporter.inc(Counters.FOO);
+ reporter.add(Counters.BAR, 4L);
+
+ reporter.flush();
+
+ assertEquals(2, statsD.counts.size());
+ // live() is a cumulative total (real events observed), not "since the last flush" -- it
+ // doesn't reset just because a flush drained and reported the delta.
+ Accumulator.Counts live = reporter.live();
+ assertEquals(1L, live.get(Counters.FOO));
+ assertEquals(4L, live.get(Counters.BAR));
+
+ reporter.flush();
+
+ // a second flush with nothing new to report doesn't inflate the cumulative total either.
+ assertEquals(2, statsD.counts.size());
+ live = reporter.live();
+ assertEquals(1L, live.get(Counters.FOO));
+ assertEquals(4L, live.get(Counters.BAR));
+ }
+
+ @Test
+ void flushCompensatesCountersNotConfirmedDeliveredAfterAnException() {
+ List delivered = new ArrayList<>();
+ AtomicBoolean failNextBar = new AtomicBoolean(true);
+ StatsDClient failsOnBarOnce =
+ new StatsDClient() {
+ @Override
+ public void incrementCounter(String metricName, String... tags) {}
+
+ @Override
+ public void count(String metricName, long delta, String... tags) {
+ if (metricName.equals("bar.total") && failNextBar.compareAndSet(true, false)) {
+ throw new RuntimeException("boom");
+ }
+ delivered.add(new RecordingStatsDClient.Count(metricName, delta, tags));
+ }
+
+ @Override
+ public void gauge(String metricName, long value, String... tags) {}
+
+ @Override
+ public void gauge(String metricName, double value, String... tags) {}
+
+ @Override
+ public void histogram(String metricName, long value, String... tags) {}
+
+ @Override
+ public void histogram(String metricName, double value, String... tags) {}
+
+ @Override
+ public void distribution(String metricName, long value, String... tags) {}
+
+ @Override
+ public void distribution(String metricName, double value, String... tags) {}
+
+ @Override
+ public void serviceCheck(
+ String serviceCheckName, String status, String message, String... tags) {}
+
+ @Override
+ public void error(Exception error) {}
+
+ @Override
+ public int getErrorCount() {
+ return 0;
+ }
+
+ @Override
+ public void close() {}
+ };
+
+ StatsDCountReporter reporter = StatsDCountReporter.of(failsOnBarOnce, Counters.class);
+ reporter.inc(Counters.FOO);
+ reporter.add(Counters.BAR, 4L);
+ reporter.add(Counters.SHARED_A, 2L);
+
+ reporter.flush();
+
+ assertEquals(1, delivered.size());
+ assertEquals("foo.total", delivered.get(0).metricName);
+
+ // the cumulative live total counts every real event exactly once, whether or not statsd ever
+ // received it -- compensating for the failed send must not inflate this.
+ Accumulator.Counts live = reporter.live();
+ assertEquals(1L, live.get(Counters.FOO));
+ assertEquals(4L, live.get(Counters.BAR));
+ assertEquals(2L, live.get(Counters.SHARED_A));
+
+ delivered.clear();
+ reporter.flush();
+
+ assertEquals(2, delivered.size());
+ assertEquals("bar.total", delivered.get(0).metricName);
+ assertEquals(4L, delivered.get(0).delta);
+ assertEquals("shared.total", delivered.get(1).metricName);
+ assertEquals(2L, delivered.get(1).delta);
+
+ // still exactly-once in the cumulative total after the retry succeeds.
+ live = reporter.live();
+ assertEquals(1L, live.get(Counters.FOO));
+ assertEquals(4L, live.get(Counters.BAR));
+ assertEquals(2L, live.get(Counters.SHARED_A));
+ }
+}