diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/AbstractInstrument.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/AbstractInstrument.java index b23855caa6d..321502a083d 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/AbstractInstrument.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/AbstractInstrument.java @@ -5,22 +5,35 @@ package io.opentelemetry.sdk.metrics; +import io.opentelemetry.context.Context; import io.opentelemetry.sdk.metrics.internal.descriptor.InstrumentDescriptor; import javax.annotation.Nullable; abstract class AbstractInstrument { private final InstrumentDescriptor descriptor; + final boolean exemplarsAlwaysOff; // All arguments cannot be null because they are checked in the abstract builder classes. - AbstractInstrument(InstrumentDescriptor descriptor) { + AbstractInstrument(InstrumentDescriptor descriptor, SdkMeter sdkMeter) { this.descriptor = descriptor; + this.exemplarsAlwaysOff = sdkMeter.isExemplarsAlwaysOff(); } final InstrumentDescriptor getDescriptor() { return descriptor; } + /** + * Returns {@link Context#current()}, or {@link Context#root()} when exemplars are known-off and + * the caller doesn't otherwise care about propagating a specific context. Used by parameterless + * record overloads on synchronous instruments to skip a thread-local lookup when the resulting + * context is only consulted by the exemplar path. + */ + final Context currentOrRootContext() { + return exemplarsAlwaysOff ? Context.root() : Context.current(); + } + @Override public boolean equals(@Nullable Object o) { if (this == o) { diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleCounter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleCounter.java index 90b4749cc99..f5cf1c75b2f 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleCounter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleCounter.java @@ -48,7 +48,7 @@ public BoundDoubleCounter bind(Attributes attributes) { @Override public void add(double value) { - add(value, Context.current()); + add(value, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleGauge.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleGauge.java index d82b3756475..36007a926b4 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleGauge.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleGauge.java @@ -47,7 +47,7 @@ public BoundDoubleGauge bind(Attributes attributes) { @Override public void set(double value) { - set(value, Context.current()); + set(value, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleHistogram.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleHistogram.java index 2a1286afbce..f8a401ff9ee 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleHistogram.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleHistogram.java @@ -48,7 +48,7 @@ public BoundDoubleHistogram bind(Attributes attributes) { @Override public void record(double value) { - record(value, Context.current()); + record(value, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleUpDownCounter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleUpDownCounter.java index b304eae6dc5..99a3972bd3c 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleUpDownCounter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkDoubleUpDownCounter.java @@ -48,7 +48,7 @@ public BoundDoubleUpDownCounter bind(Attributes attributes) { @Override public void add(double value) { - add(value, Context.current()); + add(value, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongCounter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongCounter.java index 84d8934a034..877289dde2b 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongCounter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongCounter.java @@ -47,7 +47,7 @@ public BoundLongCounter bind(Attributes attributes) { @Override public void add(long value) { - add(value, Context.current()); + add(value, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongGauge.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongGauge.java index 4f97fe8e33f..c8e4a6f45b5 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongGauge.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongGauge.java @@ -46,7 +46,7 @@ public BoundLongGauge bind(Attributes attributes) { @Override public void set(long value) { - set(value, Context.current()); + set(value, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongHistogram.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongHistogram.java index c620f58356a..a2ae99a2211 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongHistogram.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongHistogram.java @@ -48,7 +48,7 @@ public BoundLongHistogram bind(Attributes attributes) { @Override public void record(long value) { - record(value, Context.current()); + record(value, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongUpDownCounter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongUpDownCounter.java index 7abe4cf1416..bea60507fe3 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongUpDownCounter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/ExtendedSdkLongUpDownCounter.java @@ -48,7 +48,7 @@ public BoundLongUpDownCounter bind(Attributes attributes) { @Override public void add(long value) { - add(value, Context.current()); + add(value, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleCounter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleCounter.java index 72764257d2a..95e4f8a3478 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleCounter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleCounter.java @@ -28,7 +28,7 @@ class SdkDoubleCounter extends AbstractInstrument implements DoubleCounter { SdkDoubleCounter( InstrumentDescriptor descriptor, SdkMeter sdkMeter, WriteableMetricStorage storage) { - super(descriptor); + super(descriptor, sdkMeter); this.sdkMeter = sdkMeter; this.storage = storage; } @@ -48,7 +48,7 @@ public void add(double increment, Attributes attributes, Context context) { @Override public void add(double increment, Attributes attributes) { - add(increment, attributes, Context.current()); + add(increment, attributes, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleGauge.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleGauge.java index c3ee314361c..513e3cd55c8 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleGauge.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleGauge.java @@ -23,7 +23,7 @@ class SdkDoubleGauge extends AbstractInstrument implements DoubleGauge { SdkDoubleGauge( InstrumentDescriptor descriptor, SdkMeter sdkMeter, WriteableMetricStorage storage) { - super(descriptor); + super(descriptor, sdkMeter); this.sdkMeter = sdkMeter; this.storage = storage; } @@ -35,7 +35,7 @@ public boolean isEnabled() { @Override public void set(double value, Attributes attributes) { - storage.recordDouble(value, attributes, Context.current()); + storage.recordDouble(value, attributes, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleHistogram.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleHistogram.java index a4ef7ebc224..fb1e868c056 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleHistogram.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleHistogram.java @@ -28,7 +28,7 @@ class SdkDoubleHistogram extends AbstractInstrument implements DoubleHistogram { SdkDoubleHistogram( InstrumentDescriptor descriptor, SdkMeter sdkMeter, WriteableMetricStorage storage) { - super(descriptor); + super(descriptor, sdkMeter); this.sdkMeter = sdkMeter; this.storage = storage; } @@ -48,7 +48,7 @@ public void record(double value, Attributes attributes, Context context) { @Override public void record(double value, Attributes attributes) { - record(value, attributes, Context.current()); + record(value, attributes, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleUpDownCounter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleUpDownCounter.java index 0c94d89cba9..7a0603c80f9 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleUpDownCounter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkDoubleUpDownCounter.java @@ -23,7 +23,7 @@ class SdkDoubleUpDownCounter extends AbstractInstrument implements DoubleUpDownC SdkDoubleUpDownCounter( InstrumentDescriptor descriptor, SdkMeter sdkMeter, WriteableMetricStorage storage) { - super(descriptor); + super(descriptor, sdkMeter); this.sdkMeter = sdkMeter; this.storage = storage; } @@ -40,7 +40,7 @@ public void add(double increment, Attributes attributes, Context context) { @Override public void add(double increment, Attributes attributes) { - add(increment, attributes, Context.current()); + add(increment, attributes, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongCounter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongCounter.java index e03aa42b61c..cf6f1ee08fd 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongCounter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongCounter.java @@ -29,7 +29,7 @@ class SdkLongCounter extends AbstractInstrument implements LongCounter { SdkLongCounter( InstrumentDescriptor descriptor, SdkMeter sdkMeter, WriteableMetricStorage storage) { - super(descriptor); + super(descriptor, sdkMeter); this.sdkMeter = sdkMeter; this.storage = storage; } @@ -49,7 +49,7 @@ public void add(long increment, Attributes attributes, Context context) { @Override public void add(long increment, Attributes attributes) { - add(increment, attributes, Context.current()); + add(increment, attributes, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongGauge.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongGauge.java index 045d8cc5aa9..2cf3c9f3696 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongGauge.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongGauge.java @@ -22,7 +22,7 @@ class SdkLongGauge extends AbstractInstrument implements LongGauge { final WriteableMetricStorage storage; SdkLongGauge(InstrumentDescriptor descriptor, SdkMeter sdkMeter, WriteableMetricStorage storage) { - super(descriptor); + super(descriptor, sdkMeter); this.sdkMeter = sdkMeter; this.storage = storage; } @@ -34,7 +34,7 @@ public boolean isEnabled() { @Override public void set(long value, Attributes attributes) { - storage.recordLong(value, attributes, Context.current()); + storage.recordLong(value, attributes, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongHistogram.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongHistogram.java index cb4835b571f..20641857ce6 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongHistogram.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongHistogram.java @@ -29,7 +29,7 @@ class SdkLongHistogram extends AbstractInstrument implements LongHistogram { SdkLongHistogram( InstrumentDescriptor descriptor, SdkMeter sdkMeter, WriteableMetricStorage storage) { - super(descriptor); + super(descriptor, sdkMeter); this.sdkMeter = sdkMeter; this.storage = storage; } @@ -49,7 +49,7 @@ public void record(long value, Attributes attributes, Context context) { @Override public void record(long value, Attributes attributes) { - record(value, attributes, Context.current()); + record(value, attributes, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongUpDownCounter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongUpDownCounter.java index 1fe09623d9a..497f7f973fa 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongUpDownCounter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkLongUpDownCounter.java @@ -23,7 +23,7 @@ class SdkLongUpDownCounter extends AbstractInstrument implements LongUpDownCount SdkLongUpDownCounter( InstrumentDescriptor descriptor, SdkMeter sdkMeter, WriteableMetricStorage storage) { - super(descriptor); + super(descriptor, sdkMeter); this.sdkMeter = sdkMeter; this.storage = storage; } @@ -40,7 +40,7 @@ public void add(long increment, Attributes attributes, Context context) { @Override public void add(long increment, Attributes attributes) { - add(increment, attributes, Context.current()); + add(increment, attributes, currentOrRootContext()); } @Override diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkMeter.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkMeter.java index d3498408c75..2455ac66141 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkMeter.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/SdkMeter.java @@ -22,6 +22,7 @@ import io.opentelemetry.sdk.metrics.data.MetricData; import io.opentelemetry.sdk.metrics.internal.MeterConfig; import io.opentelemetry.sdk.metrics.internal.descriptor.InstrumentDescriptor; +import io.opentelemetry.sdk.metrics.internal.exemplar.AlwaysOffExemplarFilter; import io.opentelemetry.sdk.metrics.internal.export.RegisteredReader; import io.opentelemetry.sdk.metrics.internal.state.AsynchronousMetricStorage; import io.opentelemetry.sdk.metrics.internal.state.BoundStorageHandle; @@ -88,6 +89,16 @@ final class SdkMeter implements Meter { private final MeterProviderSharedState meterProviderSharedState; private final InstrumentationScopeInfo instrumentationScopeInfo; + + /** + * Returns true if the meter provider's exemplar filter samples nothing. Callers can use this to + * skip {@link io.opentelemetry.context.Context#current()} lookups on record paths that only need + * the current context to derive an exemplar span context. + */ + boolean isExemplarsAlwaysOff() { + return meterProviderSharedState.getExemplarFilter() instanceof AlwaysOffExemplarFilter; + } + private final Map readerStorageRegistries; private volatile boolean meterEnabled; diff --git a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/internal/aggregator/DoubleExplicitBucketHistogramAggregator.java b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/internal/aggregator/DoubleExplicitBucketHistogramAggregator.java index bcf9586c313..176d1bbfa6c 100644 --- a/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/internal/aggregator/DoubleExplicitBucketHistogramAggregator.java +++ b/sdk/metrics/src/main/java/io/opentelemetry/sdk/metrics/internal/aggregator/DoubleExplicitBucketHistogramAggregator.java @@ -6,7 +6,6 @@ package io.opentelemetry.sdk.metrics.internal.aggregator; import io.opentelemetry.api.common.Attributes; -import io.opentelemetry.api.internal.GuardedBy; import io.opentelemetry.context.Context; import io.opentelemetry.sdk.common.InstrumentationScopeInfo; import io.opentelemetry.sdk.common.export.MemoryMode; @@ -15,6 +14,9 @@ import io.opentelemetry.sdk.metrics.data.DoubleExemplarData; import io.opentelemetry.sdk.metrics.data.HistogramPointData; import io.opentelemetry.sdk.metrics.data.MetricData; +import io.opentelemetry.sdk.metrics.internal.concurrent.AdderUtil; +import io.opentelemetry.sdk.metrics.internal.concurrent.DoubleAdder; +import io.opentelemetry.sdk.metrics.internal.concurrent.LongAdder; import io.opentelemetry.sdk.metrics.internal.data.ImmutableHistogramData; import io.opentelemetry.sdk.metrics.internal.data.ImmutableHistogramPointData; import io.opentelemetry.sdk.metrics.internal.data.ImmutableMetricData; @@ -27,6 +29,8 @@ import java.util.Collection; import java.util.Collections; import java.util.List; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicLongFieldUpdater; import javax.annotation.Nullable; /** @@ -93,31 +97,59 @@ public MetricData toMetricData( ImmutableHistogramData.create(temporality, pointData)); } + /** + * Lock-free histogram handle inspired by prometheus/client_java. + * + *

Bucket counts and running sum use {@link LongAdder} / {@link DoubleAdder}. Min and max use + * CAS loops on volatile long bit patterns with a fast-exit for non-extremes. Total count is + * derived from bucket counts at collect. + * + *

A thread-striped {@link AtomicLong} array coordinates record and collect. Its sign bit + * signals "collect in progress". Recorders back out and spin when they observe it. This prevents + * an observation's writes from being split across collections. + * + *

Within {@link #doRecordDouble} the bucket increment is the last write. Its internal volatile + * write publishes the prior sum/min/max writes, so the collector's wait for {@code + * sum(bucketCounts) >= expected} doubles as a barrier for the whole observation. + */ static final class Handle extends AggregatorHandle { - // read-only + private static final long COLLECT_BIT = 1L << 63; + private static final long MIN_INIT_BITS = Double.doubleToRawLongBits(Double.POSITIVE_INFINITY); + private static final long MAX_INIT_BITS = Double.doubleToRawLongBits(Double.NEGATIVE_INFINITY); + private static final AtomicLongFieldUpdater MIN_BITS = + AtomicLongFieldUpdater.newUpdater(Handle.class, "minBits"); + private static final AtomicLongFieldUpdater MAX_BITS = + AtomicLongFieldUpdater.newUpdater(Handle.class, "maxBits"); + private final List boundaryList; - // read-only private final double[] boundaries; private final boolean recordMinMax; - private final Object lock = new Object(); + private final LongAdder[] bucketCounts; + private final DoubleAdder sum = AdderUtil.createDoubleAdder(); - @GuardedBy("lock") - private double sum; + // Min / max as raw double bits so they can be CAS-updated via AtomicLongFieldUpdater. Updated + // via CAS loops that fast-exit when the observation isn't a new extreme — the common + // steady-state case has no memory write. + @SuppressWarnings("UnusedVariable") + private volatile long minBits = MIN_INIT_BITS; - @GuardedBy("lock") - private double min; + @SuppressWarnings("UnusedVariable") + private volatile long maxBits = MAX_INIT_BITS; - @GuardedBy("lock") - private double max; + // Power-of-2 length so the stripe probe is a bitwise AND with stripeMask. + private final AtomicLong[] stripedStartedCounter; + private final int stripeMask; - @GuardedBy("lock") - private long count; + // Observations that have made it into buckets across all time (current + previously + // drained). Reconciles bucketCounts (drained on delta reset) with cumulativeStarted (never + // resets). Cumulative: stays 0. Delta: += drained totalCount each cycle. Collector-only. + private long cumulativeDrainedCount; - @GuardedBy("lock") - private final long[] counts; + private final long[] countsScratch; - // Used only when MemoryMode = REUSABLE_DATA + // Non-null only when MemoryMode == REUSABLE_DATA. @Nullable private final MutableHistogramPointData reusablePoint; Handle( @@ -131,13 +163,23 @@ static final class Handle extends AggregatorHandle { this.boundaryList = boundaryList; this.boundaries = boundaries; this.recordMinMax = recordMinMax; - this.counts = new long[this.boundaries.length + 1]; - this.sum = 0; - this.min = Double.MAX_VALUE; - this.max = -1; - this.count = 0; + int bucketCount = boundaries.length + 1; + this.bucketCounts = new LongAdder[bucketCount]; + for (int i = 0; i < bucketCount; i++) { + this.bucketCounts[i] = AdderUtil.createLongAdder(); + } + // Sized to NCPUS (rounded up to a power of 2 so the probe mask compiles to a bitwise AND). + // NCPUS is an upper bound on threads simultaneously executing, which bounds the useful + // stripe count for handling per-handle contention. + int stripes = roundUpToPowerOfTwo(Runtime.getRuntime().availableProcessors()); + this.stripedStartedCounter = new AtomicLong[stripes]; + for (int i = 0; i < stripes; i++) { + this.stripedStartedCounter[i] = new AtomicLong(); + } + this.stripeMask = stripes - 1; + this.countsScratch = new long[bucketCount]; if (memoryMode == MemoryMode.REUSABLE_DATA) { - this.reusablePoint = new MutableHistogramPointData(counts.length); + this.reusablePoint = new MutableHistogramPointData(bucketCount); } else { this.reusablePoint = null; } @@ -145,74 +187,163 @@ static final class Handle extends AggregatorHandle { @Override public void recordLong(long value, Attributes attributes, Context context) { - // Since there is no LongExplicitBucketHistogramAggregator and we need to support measurements - // from LongHistogram, we redirect calls from #recordLong to #recordDouble. Without this, the - // base AggregatorHandle implementation of #recordLong throws. + // There is no LongExplicitBucketHistogramAggregator. Redirect to recordDouble so + // LongHistogram measurements route through this handle. super.recordDouble((double) value, attributes, context); } @Override + @SuppressWarnings("ThreadPriorityCheck") + protected void doRecordDouble(double value) { + int bucketIndex = ExplicitBucketHistogramUtils.findBucketIndex(this.boundaries, value); + + // Reserve a pre-flip slot on our stripe via CAS. Increment only when the bit is clear at + // the time of the CAS, otherwise spin until the collector's Phase 4 clears the bit and + // retry. + AtomicLong stripe = + stripedStartedCounter[System.identityHashCode(Thread.currentThread()) & stripeMask]; + while (true) { + long current = stripe.get(); + if ((current & COLLECT_BIT) != 0) { + while ((stripe.get() & COLLECT_BIT) != 0) { + Thread.yield(); + } + continue; + } + if (stripe.compareAndSet(current, current + 1)) { + break; + } + // CAS lost the race (either the bit was just set or another recorder incremented); + // loop and reevaluate. + } + + sum.add(value); + if (recordMinMax) { + updateMin(value); + updateMax(value); + } + bucketCounts[bucketIndex].increment(); + } + + /** Fast-exits without a CAS when {@code value} is not smaller than the current min. */ + private void updateMin(double value) { + long newBits = Double.doubleToRawLongBits(value); + long cur; + do { + cur = minBits; + if (value >= Double.longBitsToDouble(cur)) { + return; + } + } while (!MIN_BITS.compareAndSet(this, cur, newBits)); + } + + /** Fast-exits without a CAS when {@code value} is not larger than the current max. */ + private void updateMax(double value) { + long newBits = Double.doubleToRawLongBits(value); + long cur; + do { + cur = maxBits; + if (value <= Double.longBitsToDouble(cur)) { + return; + } + } while (!MAX_BITS.compareAndSet(this, cur, newBits)); + } + + @Override + @SuppressWarnings("ThreadPriorityCheck") protected HistogramPointData doAggregateThenMaybeResetDoubles( long startEpochNanos, long epochNanos, Attributes attributes, List exemplars, boolean reset) { - synchronized (lock) { - HistogramPointData pointData; - if (reusablePoint == null) { - pointData = - ImmutableHistogramPointData.create( - startEpochNanos, - epochNanos, - attributes, - sum, - recordMinMax && this.count > 0, - recordMinMax ? this.min : 0, - recordMinMax && this.count > 0, - recordMinMax ? this.max : 0, - boundaryList, - PrimitiveLongList.wrap(Arrays.copyOf(counts, counts.length)), - exemplars); - } else /* REUSABLE_DATA */ { - pointData = - reusablePoint.set( - startEpochNanos, - epochNanos, - attributes, - sum, - recordMinMax && this.count > 0, - recordMinMax ? this.min : 0, - recordMinMax && this.count > 0, - recordMinMax ? this.max : 0, - boundaryList, - counts, - exemplars); - } - if (reset) { - this.sum = 0; - this.min = Double.MAX_VALUE; - this.max = -1; - this.count = 0; - Arrays.fill(this.counts, 0); - } - return pointData; + // Phase 1: set the collect bit on every stripe. Capture the pre-flip cumulative count. + long cumulativeStarted = 0; + for (AtomicLong stripe : stripedStartedCounter) { + cumulativeStarted += stripe.getAndAdd(COLLECT_BIT) & ~COLLECT_BIT; + } + + // Phase 2: wait for pre-flip recorders to publish their bucket increments. Post-flip + // recorders are spinning on the collect bit, so nothing new arrives. + while (bucketSumTotal() + cumulativeDrainedCount < cumulativeStarted) { + Thread.yield(); + } + + // Phase 3: snapshot (and reset if delta) with recorders quiescent. + long totalCount = 0; + for (int i = 0; i < bucketCounts.length; i++) { + long c = reset ? bucketCounts[i].sumThenReset() : bucketCounts[i].sum(); + countsScratch[i] = c; + totalCount += c; + } + if (reset) { + cumulativeDrainedCount += totalCount; + } + double totalSum = reset ? sum.sumThenReset() : sum.sum(); + + double snapshotMin = Double.POSITIVE_INFINITY; + double snapshotMax = Double.NEGATIVE_INFINITY; + if (recordMinMax) { + long minSnapshot = reset ? MIN_BITS.getAndSet(this, MIN_INIT_BITS) : minBits; + long maxSnapshot = reset ? MAX_BITS.getAndSet(this, MAX_INIT_BITS) : maxBits; + snapshotMin = Double.longBitsToDouble(minSnapshot); + snapshotMax = Double.longBitsToDouble(maxSnapshot); + } + + // Phase 4: clear the collect bit. addAndGet(COLLECT_BIT) toggles the sign bit off via + // two's-complement overflow. Spinning recorders resume. + for (AtomicLong stripe : stripedStartedCounter) { + stripe.addAndGet(COLLECT_BIT); } + + HistogramPointData pointData; + if (reusablePoint == null) { + pointData = + ImmutableHistogramPointData.create( + startEpochNanos, + epochNanos, + attributes, + totalSum, + recordMinMax && totalCount > 0, + recordMinMax ? snapshotMin : 0, + recordMinMax && totalCount > 0, + recordMinMax ? snapshotMax : 0, + boundaryList, + PrimitiveLongList.wrap(Arrays.copyOf(countsScratch, countsScratch.length)), + exemplars); + } else /* REUSABLE_DATA */ { + pointData = + reusablePoint.set( + startEpochNanos, + epochNanos, + attributes, + totalSum, + recordMinMax && totalCount > 0, + recordMinMax ? snapshotMin : 0, + recordMinMax && totalCount > 0, + recordMinMax ? snapshotMax : 0, + boundaryList, + countsScratch, + exemplars); + } + return pointData; } - @Override - protected void doRecordDouble(double value) { - int bucketIndex = ExplicitBucketHistogramUtils.findBucketIndex(this.boundaries, value); + private long bucketSumTotal() { + long total = 0; + for (LongAdder adder : bucketCounts) { + total += adder.sum(); + } + return total; + } - synchronized (lock) { - this.sum += value; - if (recordMinMax) { - this.min = Math.min(this.min, value); - this.max = Math.max(this.max, value); - } - this.count++; - this.counts[bucketIndex]++; + /** Smallest power of 2 >= {@code n}, with a floor of 1. */ + private static int roundUpToPowerOfTwo(int n) { + if (n <= 1) { + return 1; } + int highest = Integer.highestOneBit(n); + return highest == n ? highest : highest << 1; } } } diff --git a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/AbstractInstrumentTest.java b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/AbstractInstrumentTest.java index 659a8260c5a..dc8cb23cad8 100644 --- a/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/AbstractInstrumentTest.java +++ b/sdk/metrics/src/test/java/io/opentelemetry/sdk/metrics/AbstractInstrumentTest.java @@ -6,6 +6,7 @@ package io.opentelemetry.sdk.metrics; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; import io.opentelemetry.sdk.metrics.internal.descriptor.Advice; import io.opentelemetry.sdk.metrics.internal.descriptor.InstrumentDescriptor; @@ -22,6 +23,8 @@ class AbstractInstrumentTest { InstrumentValueType.LONG, Advice.empty()); + private static final SdkMeter SDK_METER = mock(SdkMeter.class); + @Test void getValues() { TestInstrument testInstrument = new TestInstrument(INSTRUMENT_DESCRIPTOR); @@ -37,7 +40,7 @@ void stringRepresentation() { private static final class TestInstrument extends AbstractInstrument { TestInstrument(InstrumentDescriptor descriptor) { - super(descriptor); + super(descriptor, SDK_METER); } } } diff --git a/sdk/metrics/src/testIncubating/java/io/opentelemetry/sdk/metrics/SynchronousInstrumentStressTest.java b/sdk/metrics/src/testIncubating/java/io/opentelemetry/sdk/metrics/SynchronousInstrumentStressTest.java index 29be5c455b9..3e29aac7598 100644 --- a/sdk/metrics/src/testIncubating/java/io/opentelemetry/sdk/metrics/SynchronousInstrumentStressTest.java +++ b/sdk/metrics/src/testIncubating/java/io/opentelemetry/sdk/metrics/SynchronousInstrumentStressTest.java @@ -54,9 +54,11 @@ import java.util.Collections; import java.util.List; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; import java.util.stream.LongStream; import java.util.stream.Stream; +import org.junit.jupiter.api.Timeout; import org.junit.jupiter.api.extension.RegisterExtension; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; @@ -121,35 +123,7 @@ private void stressTestOnce( MemoryMode memoryMode, InstrumentValueType instrumentValueType, boolean bound) { - // Initialize metric SDK - DefaultAggregationSelector aggregationSelector = - DefaultAggregationSelector.getDefault().with(instrumentType, aggregation); - InMemoryMetricReader reader = - InMemoryMetricReader.builder() - .setDefaultAggregationSelector(aggregationSelector) - .setAggregationTemporalitySelector(unused -> aggregationTemporality) - .setMemoryMode(memoryMode) - .build(); - SdkMeterProvider meterProvider = - SdkMeterProvider.builder().registerMetricReader(reader).build(); - cleanup.addCloseable(meterProvider); - Meter meter = meterProvider.get("test"); - List attributes = Arrays.asList(ATTR_1, ATTR_2, ATTR_3, ATTR_4); - Collections.shuffle(attributes); - // When unbound, record through `instrument`, looking up the series by attributes on each call. - // When bound, bind one instrument per series up front and record straight to those (no - // per-record attribute lookup) — the record loop below forks accordingly. - Instrument instrument = getInstrument(meter, instrumentType, instrumentValueType); - List boundInstruments = new ArrayList<>(); - if (bound) { - for (Attributes attr : attributes) { - boundInstruments.add(getBoundInstrument(meter, instrumentType, instrumentValueType, attr)); - } - } - - // Define list of measurements to record - // Later, we'll assert that the data collected matches these measurements, with no lost writes, - // partial writes, duplicate writes, etc. + // Define list of measurements to record. Shuffled 0..1999. int measurementCount = 2000; List measurements = new ArrayList<>(); for (int i = 0; i < measurementCount; i++) { @@ -157,62 +131,16 @@ private void stressTestOnce( } Collections.shuffle(measurements); - // Define recording threads + List collectedMetrics = + runStressLoop( + aggregationTemporality, + instrumentType, + aggregation, + memoryMode, + instrumentValueType, + bound, + measurements); int threadCount = 4; - List recordThreads = new ArrayList<>(); - CountDownLatch latch = new CountDownLatch(threadCount); - CountDownLatch startSignal = new CountDownLatch(1); - for (int i = 0; i < threadCount; i++) { - recordThreads.add( - new Thread( - () -> { - Uninterruptibles.awaitUninterruptibly(startSignal); - if (bound) { - for (Long measurement : measurements) { - for (BoundInstrument boundInstrument : boundInstruments) { - boundInstrument.record(measurement); - } - Thread.yield(); - } - } else { - for (Long measurement : measurements) { - for (Attributes attr : attributes) { - instrument.record(measurement, attr); - } - Thread.yield(); - } - } - latch.countDown(); - })); - } - - // Define collecting thread - // NOTE: collect makes a copy of MetricData because REUSABLE_DATA mode reuses MetricData - List collectedMetrics = new ArrayList<>(); - Thread collectThread = - new Thread( - () -> { - Uninterruptibles.awaitUninterruptibly(startSignal); - while (latch.getCount() != 0) { - Thread.yield(); - collectedMetrics.addAll( - reader.collectAllMetrics().stream() - .map(SynchronousInstrumentStressTest::copy) - .collect(toList())); - } - collectedMetrics.addAll( - reader.collectAllMetrics().stream() - .map(SynchronousInstrumentStressTest::copy) - .collect(toList())); - }); - - // Start all the threads, then release the start signal so they begin simultaneously - collectThread.start(); - recordThreads.forEach(Thread::start); - startSignal.countDown(); - - // Wait for the collect thread to end, which collects until the record threads are done - Uninterruptibles.joinUninterruptibly(collectThread); // Assert collected data is consistent with recorded measurements by independently computing the // expected aggregated value and comparing to the actual results. @@ -328,6 +256,239 @@ private void stressTestOnce( } } + /** + * Complements {@link #stressTest} by asserting that no interim collected snapshot ever + * shows a partial write — e.g., a histogram with {@code sum} disagreeing with the sum of its + * bucket counts, indicating the collector observed a recorder mid-write. + * + *

Every recorder records a uniform value of 1, so the assertions are checkable per snapshot + * without knowing which specific observations completed at snapshot time: for any histogram + * point, {@code sum == count * 1}, and the sum of bucket counts equals {@code count}. + * + *

Assertions run after the stress phase completes, so the collector loop is never slowed down. + * Scoped to histogram aggregations — the invariant is trivial for scalar aggregations + * (sum/last-value), which don't span multiple mutable fields per record. + */ + @ParameterizedTest + @MethodSource("stressTestArgs") + @Timeout(value = 10, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD) + void partialWriteStressTest( + AggregationTemporality aggregationTemporality, + InstrumentType instrumentType, + Aggregation aggregation, + MemoryMode memoryMode, + InstrumentValueType instrumentValueType, + boolean bound) { + for (int repetition = 0; repetition < STRESS_TEST_REPETITIONS; repetition++) { + partialWriteStressTestOnce( + aggregationTemporality, + instrumentType, + aggregation, + memoryMode, + instrumentValueType, + bound); + } + } + + @SuppressWarnings("ReferenceEquality") + private void partialWriteStressTestOnce( + AggregationTemporality aggregationTemporality, + InstrumentType instrumentType, + Aggregation aggregation, + MemoryMode memoryMode, + InstrumentValueType instrumentValueType, + boolean bound) { + // Uniform value of 1 so that at any interim snapshot: sum == count. + int measurementCount = 2000; + List measurements = new ArrayList<>(measurementCount); + for (int i = 0; i < measurementCount; i++) { + measurements.add(1L); + } + + List collectedMetrics = + runStressLoop( + aggregationTemporality, + instrumentType, + aggregation, + memoryMode, + instrumentValueType, + bound, + measurements); + + // Assert every collected point is internally consistent (no partial writes visible). Scalar + // aggregators can't produce partial writes structurally, but we still exercise them for + // general concurrent-path regression coverage with light per-type invariants. + if (aggregation == Aggregation.explicitBucketHistogram()) { + collectedMetrics.stream() + .flatMap(m -> m.getHistogramData().getPoints().stream()) + .forEach( + p -> { + long bucketSum = p.getCounts().stream().reduce(0L, Long::sum); + assertThat(p.getCount()).as("count == sum of bucket counts").isEqualTo(bucketSum); + assertThat(p.getSum()) + .as("sum == count (uniform value=1)") + .isEqualTo((double) p.getCount()); + if (p.getCount() > 0) { + assertThat(p.hasMin()).isTrue(); + assertThat(p.getMin()).as("min == 1 (uniform value)").isEqualTo(1.0); + assertThat(p.hasMax()).isTrue(); + assertThat(p.getMax()).as("max == 1 (uniform value)").isEqualTo(1.0); + } + }); + } else if (aggregation == Aggregation.base2ExponentialBucketHistogram()) { + collectedMetrics.stream() + .flatMap(m -> m.getExponentialHistogramData().getPoints().stream()) + .forEach( + p -> { + long positiveBucketSum = + p.getPositiveBuckets().getBucketCounts().stream().reduce(0L, Long::sum); + long negativeBucketSum = + p.getNegativeBuckets().getBucketCounts().stream().reduce(0L, Long::sum); + assertThat(p.getCount()) + .as("count == pos + neg + zero bucket counts") + .isEqualTo(positiveBucketSum + negativeBucketSum + p.getZeroCount()); + assertThat(p.getSum()) + .as("sum == count (uniform value=1)") + .isEqualTo((double) p.getCount()); + if (p.getCount() > 0) { + assertThat(p.hasMin()).isTrue(); + assertThat(p.getMin()).as("min == 1 (uniform value)").isEqualTo(1.0); + assertThat(p.hasMax()).isTrue(); + assertThat(p.getMax()).as("max == 1 (uniform value)").isEqualTo(1.0); + } + }); + } else if (aggregation == Aggregation.sum()) { + // Uniform value=1 adds only. Sum should be a non-negative whole number. + if (instrumentValueType == InstrumentValueType.DOUBLE) { + collectedMetrics.stream() + .flatMap(m -> m.getDoubleSumData().getPoints().stream()) + .forEach( + p -> { + assertThat(p.getValue()).as("sum non-negative").isGreaterThanOrEqualTo(0.0); + assertThat(p.getValue()) + .as("sum is integral") + .isEqualTo(Math.floor(p.getValue())); + }); + } else { + collectedMetrics.stream() + .flatMap(m -> m.getLongSumData().getPoints().stream()) + .forEach( + p -> assertThat(p.getValue()).as("sum non-negative").isGreaterThanOrEqualTo(0L)); + } + } else if (aggregation == Aggregation.lastValue()) { + // Every observation is 1; any non-empty snapshot's value must be 1. + if (instrumentValueType == InstrumentValueType.DOUBLE) { + collectedMetrics.stream() + .flatMap(m -> m.getDoubleGaugeData().getPoints().stream()) + .forEach(p -> assertThat(p.getValue()).as("last-value == 1").isEqualTo(1.0)); + } else { + collectedMetrics.stream() + .flatMap(m -> m.getLongGaugeData().getPoints().stream()) + .forEach(p -> assertThat(p.getValue()).as("last-value == 1").isEqualTo(1L)); + } + } else { + throw new IllegalArgumentException("Unexpected aggregation: " + aggregation); + } + } + + /** + * Runs a record/collect stress loop and returns the collected metric snapshots. Each of the 4 + * recording threads iterates {@code measurements}, and for each measurement records it against + * every attribute (or, when {@code bound}, every pre-bound instrument). The collector thread + * continuously calls {@code collectAllMetrics} until recorders complete, then does one final + * collect. Copies are taken of each collected {@link MetricData} because {@code REUSABLE_DATA} + * mode reuses instances across collects. + */ + @SuppressWarnings("ThreadPriorityCheck") + private List runStressLoop( + AggregationTemporality aggregationTemporality, + InstrumentType instrumentType, + Aggregation aggregation, + MemoryMode memoryMode, + InstrumentValueType instrumentValueType, + boolean bound, + List measurements) { + DefaultAggregationSelector aggregationSelector = + DefaultAggregationSelector.getDefault().with(instrumentType, aggregation); + InMemoryMetricReader reader = + InMemoryMetricReader.builder() + .setDefaultAggregationSelector(aggregationSelector) + .setAggregationTemporalitySelector(unused -> aggregationTemporality) + .setMemoryMode(memoryMode) + .build(); + SdkMeterProvider meterProvider = + SdkMeterProvider.builder().registerMetricReader(reader).build(); + cleanup.addCloseable(meterProvider); + Meter meter = meterProvider.get("test"); + List attributes = Arrays.asList(ATTR_1, ATTR_2, ATTR_3, ATTR_4); + Collections.shuffle(attributes); + // When unbound, record through `instrument`, looking up the series by attributes on each call. + // When bound, bind one instrument per series up front and record straight to those (no + // per-record attribute lookup) — the record loop below forks accordingly. + Instrument instrument = getInstrument(meter, instrumentType, instrumentValueType); + List boundInstruments = new ArrayList<>(); + if (bound) { + for (Attributes attr : attributes) { + boundInstruments.add(getBoundInstrument(meter, instrumentType, instrumentValueType, attr)); + } + } + + int threadCount = 4; + List recordThreads = new ArrayList<>(); + CountDownLatch latch = new CountDownLatch(threadCount); + CountDownLatch startSignal = new CountDownLatch(1); + for (int i = 0; i < threadCount; i++) { + recordThreads.add( + new Thread( + () -> { + Uninterruptibles.awaitUninterruptibly(startSignal); + if (bound) { + for (Long measurement : measurements) { + for (BoundInstrument boundInstrument : boundInstruments) { + boundInstrument.record(measurement); + } + Thread.yield(); + } + } else { + for (Long measurement : measurements) { + for (Attributes attr : attributes) { + instrument.record(measurement, attr); + } + Thread.yield(); + } + } + latch.countDown(); + })); + } + + List collectedMetrics = new ArrayList<>(); + Thread collectThread = + new Thread( + () -> { + Uninterruptibles.awaitUninterruptibly(startSignal); + while (latch.getCount() != 0) { + Thread.yield(); + collectedMetrics.addAll( + reader.collectAllMetrics().stream() + .map(SynchronousInstrumentStressTest::copy) + .collect(toList())); + } + collectedMetrics.addAll( + reader.collectAllMetrics().stream() + .map(SynchronousInstrumentStressTest::copy) + .collect(toList())); + }); + + // Start all the threads, then release the start signal so they begin simultaneously + collectThread.start(); + recordThreads.forEach(Thread::start); + startSignal.countDown(); + + // Wait for the collect thread to end, which collects until the record threads are done + Uninterruptibles.joinUninterruptibly(collectThread); + return collectedMetrics; + } + private static Stream stressTestArgs() { List argumentsList = new ArrayList<>(); for (AggregationTemporality aggregationTemporality : AggregationTemporality.values()) {