diff --git a/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/ReplaySafeOpenTelemetry.java b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/ReplaySafeOpenTelemetry.java
index d12e51c1fd..7d05111c70 100644
--- a/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/ReplaySafeOpenTelemetry.java
+++ b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/ReplaySafeOpenTelemetry.java
@@ -4,6 +4,8 @@
import io.opentelemetry.api.OpenTelemetry;
import io.opentelemetry.api.baggage.propagation.W3CBaggagePropagator;
import io.opentelemetry.api.logs.LoggerProvider;
+import io.opentelemetry.api.metrics.Meter;
+import io.opentelemetry.api.metrics.MeterBuilder;
import io.opentelemetry.api.metrics.MeterProvider;
import io.opentelemetry.api.trace.Tracer;
import io.opentelemetry.api.trace.TracerBuilder;
@@ -19,18 +21,20 @@
import io.opentelemetry.sdk.trace.SdkTracerProviderBuilder;
import io.temporal.common.Experimental;
import io.temporal.opentelemetry.v2.internal.ReplaySafeIdGenerator;
+import io.temporal.opentelemetry.v2.internal.ReplaySafeMeter;
import io.temporal.opentelemetry.v2.internal.ReplaySafeTracer;
import java.io.Closeable;
import javax.annotation.Nonnull;
/**
* The {@link OpenTelemetry} to use for OpenTelemetry integration with Temporal. Register it with
- * {@code GlobalOpenTelemetry.set}; tracers obtained from it are replay safe inside workflows.
+ * {@code GlobalOpenTelemetry.set}; tracers and meters obtained from it are replay safe inside
+ * workflows.
*/
@Experimental
public final class ReplaySafeOpenTelemetry implements OpenTelemetry, Closeable {
private final ReplaySafeTracerProvider tracerProvider;
- private final SdkMeterProvider meterProvider; // TODO - Make the meter provider replay safe
+ private final ReplaySafeMeterProvider meterProvider;
// TODO: Make the logger provider replay safe and add logger interceptor methods for Temporal and
// OpenTelemetry loggers.
private final SdkLoggerProvider loggerProvider;
@@ -40,7 +44,7 @@ private ReplaySafeOpenTelemetry(Builder builder) {
this.tracerProvider =
new ReplaySafeTracerProvider(
builder.tracerProviderBuilder.setIdGenerator(new ReplaySafeIdGenerator()).build());
- this.meterProvider = builder.meterProviderBuilder.build();
+ this.meterProvider = new ReplaySafeMeterProvider(builder.meterProviderBuilder.build());
this.loggerProvider = builder.loggerProviderBuilder.build();
this.propagators = builder.propagators;
}
@@ -155,6 +159,54 @@ public void close() {
}
}
+ private static final class ReplaySafeMeterProvider implements MeterProvider, Closeable {
+ private final SdkMeterProvider delegate;
+
+ private ReplaySafeMeterProvider(SdkMeterProvider delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public Meter get(@Nonnull String instrumentationScopeName) {
+ return new ReplaySafeMeter(delegate.get(instrumentationScopeName));
+ }
+
+ @Override
+ public MeterBuilder meterBuilder(@Nonnull String instrumentationScopeName) {
+ return new ReplaySafeMeterBuilder(delegate.meterBuilder(instrumentationScopeName));
+ }
+
+ @Override
+ public void close() {
+ delegate.close();
+ }
+ }
+
+ private static final class ReplaySafeMeterBuilder implements MeterBuilder {
+ private final MeterBuilder delegate;
+
+ private ReplaySafeMeterBuilder(MeterBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public MeterBuilder setSchemaUrl(@Nonnull String schemaUrl) {
+ delegate.setSchemaUrl(schemaUrl);
+ return this;
+ }
+
+ @Override
+ public MeterBuilder setInstrumentationVersion(@Nonnull String instrumentationScopeVersion) {
+ delegate.setInstrumentationVersion(instrumentationScopeVersion);
+ return this;
+ }
+
+ @Override
+ public Meter build() {
+ return new ReplaySafeMeter(delegate.build());
+ }
+ }
+
private static final class ReplaySafeTracerBuilder implements TracerBuilder {
private final TracerBuilder delegate;
private final String instrumentationScopeName;
diff --git a/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/OpenTelemetrySuppression.java b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/OpenTelemetrySuppression.java
new file mode 100644
index 0000000000..e05d31543c
--- /dev/null
+++ b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/OpenTelemetrySuppression.java
@@ -0,0 +1,21 @@
+package io.temporal.opentelemetry.v2.internal;
+
+import io.temporal.workflow.unsafe.WorkflowUnsafe;
+
+/** Where replayed workflow code must not export telemetry again. */
+public final class OpenTelemetrySuppression {
+ private OpenTelemetrySuppression() {}
+
+ /**
+ * True on a workflow thread that is re-executing history. Query handlers, update validators, and
+ * side effect functions run live, at most once, even while the workflow replays, so their
+ * telemetry is kept. Await conditions are re-evaluated on every replay and so are suppressed,
+ * which is why this asks whether the calling code is subject to replay rather than whether it can
+ * mutate workflow state.
+ */
+ public static boolean shouldSuppress() {
+ return WorkflowUnsafe.isWorkflowThread()
+ && WorkflowUnsafe.isSubjectToReplay()
+ && WorkflowUnsafe.isReplaying();
+ }
+}
diff --git a/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeMeter.java b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeMeter.java
new file mode 100644
index 0000000000..80ecd9cba6
--- /dev/null
+++ b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeMeter.java
@@ -0,0 +1,670 @@
+package io.temporal.opentelemetry.v2.internal;
+
+import io.opentelemetry.api.common.Attributes;
+import io.opentelemetry.api.metrics.BatchCallback;
+import io.opentelemetry.api.metrics.DoubleCounter;
+import io.opentelemetry.api.metrics.DoubleCounterBuilder;
+import io.opentelemetry.api.metrics.DoubleGauge;
+import io.opentelemetry.api.metrics.DoubleGaugeBuilder;
+import io.opentelemetry.api.metrics.DoubleHistogram;
+import io.opentelemetry.api.metrics.DoubleHistogramBuilder;
+import io.opentelemetry.api.metrics.DoubleUpDownCounter;
+import io.opentelemetry.api.metrics.DoubleUpDownCounterBuilder;
+import io.opentelemetry.api.metrics.LongCounter;
+import io.opentelemetry.api.metrics.LongCounterBuilder;
+import io.opentelemetry.api.metrics.LongGauge;
+import io.opentelemetry.api.metrics.LongGaugeBuilder;
+import io.opentelemetry.api.metrics.LongHistogram;
+import io.opentelemetry.api.metrics.LongHistogramBuilder;
+import io.opentelemetry.api.metrics.LongUpDownCounter;
+import io.opentelemetry.api.metrics.LongUpDownCounterBuilder;
+import io.opentelemetry.api.metrics.Meter;
+import io.opentelemetry.api.metrics.ObservableDoubleCounter;
+import io.opentelemetry.api.metrics.ObservableDoubleGauge;
+import io.opentelemetry.api.metrics.ObservableDoubleMeasurement;
+import io.opentelemetry.api.metrics.ObservableDoubleUpDownCounter;
+import io.opentelemetry.api.metrics.ObservableLongCounter;
+import io.opentelemetry.api.metrics.ObservableLongGauge;
+import io.opentelemetry.api.metrics.ObservableLongMeasurement;
+import io.opentelemetry.api.metrics.ObservableLongUpDownCounter;
+import io.opentelemetry.api.metrics.ObservableMeasurement;
+import io.opentelemetry.context.Context;
+import java.util.List;
+import java.util.function.Consumer;
+
+/**
+ * Wraps a meter so the synchronous instruments it builds drop recordings made by replaying workflow
+ * code, which would otherwise be recorded again on every replay.
+ *
+ *
Observable instruments are not wrapped. Creating them from workflow code is discouraged
+ * because their callbacks run outside the workflow and the callback observes workflow state without
+ * workflow synchronization. Create observable instruments from worker or activity code instead.
+ */
+public final class ReplaySafeMeter implements Meter {
+ private final Meter delegate;
+
+ public ReplaySafeMeter(Meter delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public LongCounterBuilder counterBuilder(String name) {
+ return new ReplaySafeLongCounterBuilder(delegate.counterBuilder(name));
+ }
+
+ @Override
+ public LongUpDownCounterBuilder upDownCounterBuilder(String name) {
+ return new ReplaySafeLongUpDownCounterBuilder(delegate.upDownCounterBuilder(name));
+ }
+
+ @Override
+ public DoubleHistogramBuilder histogramBuilder(String name) {
+ return new ReplaySafeDoubleHistogramBuilder(delegate.histogramBuilder(name));
+ }
+
+ @Override
+ public DoubleGaugeBuilder gaugeBuilder(String name) {
+ return new ReplaySafeDoubleGaugeBuilder(delegate.gaugeBuilder(name));
+ }
+
+ @Override
+ public BatchCallback batchCallback(
+ Runnable callback,
+ ObservableMeasurement observableMeasurement,
+ ObservableMeasurement... additionalMeasurements) {
+ return delegate.batchCallback(callback, observableMeasurement, additionalMeasurements);
+ }
+
+ private static final class ReplaySafeLongCounterBuilder implements LongCounterBuilder {
+ private final LongCounterBuilder delegate;
+
+ ReplaySafeLongCounterBuilder(LongCounterBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public LongCounterBuilder setDescription(String description) {
+ delegate.setDescription(description);
+ return this;
+ }
+
+ @Override
+ public LongCounterBuilder setUnit(String unit) {
+ delegate.setUnit(unit);
+ return this;
+ }
+
+ @Override
+ public DoubleCounterBuilder ofDoubles() {
+ return new ReplaySafeDoubleCounterBuilder(delegate.ofDoubles());
+ }
+
+ @Override
+ public LongCounter build() {
+ return new ReplaySafeLongCounter(delegate.build());
+ }
+
+ @Override
+ public ObservableLongCounter buildWithCallback(Consumer callback) {
+ return delegate.buildWithCallback(callback);
+ }
+
+ @Override
+ public ObservableLongMeasurement buildObserver() {
+ return delegate.buildObserver();
+ }
+ }
+
+ private static final class ReplaySafeDoubleCounterBuilder implements DoubleCounterBuilder {
+ private final DoubleCounterBuilder delegate;
+
+ ReplaySafeDoubleCounterBuilder(DoubleCounterBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public DoubleCounterBuilder setDescription(String description) {
+ delegate.setDescription(description);
+ return this;
+ }
+
+ @Override
+ public DoubleCounterBuilder setUnit(String unit) {
+ delegate.setUnit(unit);
+ return this;
+ }
+
+ @Override
+ public DoubleCounter build() {
+ return new ReplaySafeDoubleCounter(delegate.build());
+ }
+
+ @Override
+ public ObservableDoubleCounter buildWithCallback(
+ Consumer callback) {
+ return delegate.buildWithCallback(callback);
+ }
+
+ @Override
+ public ObservableDoubleMeasurement buildObserver() {
+ return delegate.buildObserver();
+ }
+ }
+
+ private static final class ReplaySafeLongUpDownCounterBuilder
+ implements LongUpDownCounterBuilder {
+ private final LongUpDownCounterBuilder delegate;
+
+ ReplaySafeLongUpDownCounterBuilder(LongUpDownCounterBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public LongUpDownCounterBuilder setDescription(String description) {
+ delegate.setDescription(description);
+ return this;
+ }
+
+ @Override
+ public LongUpDownCounterBuilder setUnit(String unit) {
+ delegate.setUnit(unit);
+ return this;
+ }
+
+ @Override
+ public DoubleUpDownCounterBuilder ofDoubles() {
+ return new ReplaySafeDoubleUpDownCounterBuilder(delegate.ofDoubles());
+ }
+
+ @Override
+ public LongUpDownCounter build() {
+ return new ReplaySafeLongUpDownCounter(delegate.build());
+ }
+
+ @Override
+ public ObservableLongUpDownCounter buildWithCallback(
+ Consumer callback) {
+ return delegate.buildWithCallback(callback);
+ }
+
+ @Override
+ public ObservableLongMeasurement buildObserver() {
+ return delegate.buildObserver();
+ }
+ }
+
+ private static final class ReplaySafeDoubleUpDownCounterBuilder
+ implements DoubleUpDownCounterBuilder {
+ private final DoubleUpDownCounterBuilder delegate;
+
+ ReplaySafeDoubleUpDownCounterBuilder(DoubleUpDownCounterBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public DoubleUpDownCounterBuilder setDescription(String description) {
+ delegate.setDescription(description);
+ return this;
+ }
+
+ @Override
+ public DoubleUpDownCounterBuilder setUnit(String unit) {
+ delegate.setUnit(unit);
+ return this;
+ }
+
+ @Override
+ public DoubleUpDownCounter build() {
+ return new ReplaySafeDoubleUpDownCounter(delegate.build());
+ }
+
+ @Override
+ public ObservableDoubleUpDownCounter buildWithCallback(
+ Consumer callback) {
+ return delegate.buildWithCallback(callback);
+ }
+
+ @Override
+ public ObservableDoubleMeasurement buildObserver() {
+ return delegate.buildObserver();
+ }
+ }
+
+ private static final class ReplaySafeDoubleHistogramBuilder implements DoubleHistogramBuilder {
+ private final DoubleHistogramBuilder delegate;
+
+ ReplaySafeDoubleHistogramBuilder(DoubleHistogramBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public DoubleHistogramBuilder setDescription(String description) {
+ delegate.setDescription(description);
+ return this;
+ }
+
+ @Override
+ public DoubleHistogramBuilder setUnit(String unit) {
+ delegate.setUnit(unit);
+ return this;
+ }
+
+ @Override
+ public DoubleHistogramBuilder setExplicitBucketBoundariesAdvice(List bucketBoundaries) {
+ delegate.setExplicitBucketBoundariesAdvice(bucketBoundaries);
+ return this;
+ }
+
+ @Override
+ public LongHistogramBuilder ofLongs() {
+ return new ReplaySafeLongHistogramBuilder(delegate.ofLongs());
+ }
+
+ @Override
+ public DoubleHistogram build() {
+ return new ReplaySafeDoubleHistogram(delegate.build());
+ }
+ }
+
+ private static final class ReplaySafeLongHistogramBuilder implements LongHistogramBuilder {
+ private final LongHistogramBuilder delegate;
+
+ ReplaySafeLongHistogramBuilder(LongHistogramBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public LongHistogramBuilder setDescription(String description) {
+ delegate.setDescription(description);
+ return this;
+ }
+
+ @Override
+ public LongHistogramBuilder setUnit(String unit) {
+ delegate.setUnit(unit);
+ return this;
+ }
+
+ @Override
+ public LongHistogramBuilder setExplicitBucketBoundariesAdvice(List bucketBoundaries) {
+ delegate.setExplicitBucketBoundariesAdvice(bucketBoundaries);
+ return this;
+ }
+
+ @Override
+ public LongHistogram build() {
+ return new ReplaySafeLongHistogram(delegate.build());
+ }
+ }
+
+ private static final class ReplaySafeDoubleGaugeBuilder implements DoubleGaugeBuilder {
+ private final DoubleGaugeBuilder delegate;
+
+ ReplaySafeDoubleGaugeBuilder(DoubleGaugeBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public DoubleGaugeBuilder setDescription(String description) {
+ delegate.setDescription(description);
+ return this;
+ }
+
+ @Override
+ public DoubleGaugeBuilder setUnit(String unit) {
+ delegate.setUnit(unit);
+ return this;
+ }
+
+ @Override
+ public LongGaugeBuilder ofLongs() {
+ return new ReplaySafeLongGaugeBuilder(delegate.ofLongs());
+ }
+
+ @Override
+ public ObservableDoubleGauge buildWithCallback(Consumer callback) {
+ return delegate.buildWithCallback(callback);
+ }
+
+ @Override
+ public ObservableDoubleMeasurement buildObserver() {
+ return delegate.buildObserver();
+ }
+
+ @Override
+ public DoubleGauge build() {
+ return new ReplaySafeDoubleGauge(delegate.build());
+ }
+ }
+
+ private static final class ReplaySafeLongGaugeBuilder implements LongGaugeBuilder {
+ private final LongGaugeBuilder delegate;
+
+ ReplaySafeLongGaugeBuilder(LongGaugeBuilder delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public LongGaugeBuilder setDescription(String description) {
+ delegate.setDescription(description);
+ return this;
+ }
+
+ @Override
+ public LongGaugeBuilder setUnit(String unit) {
+ delegate.setUnit(unit);
+ return this;
+ }
+
+ @Override
+ public ObservableLongGauge buildWithCallback(Consumer callback) {
+ return delegate.buildWithCallback(callback);
+ }
+
+ @Override
+ public ObservableLongMeasurement buildObserver() {
+ return delegate.buildObserver();
+ }
+
+ @Override
+ public LongGauge build() {
+ return new ReplaySafeLongGauge(delegate.build());
+ }
+ }
+
+ private static final class ReplaySafeLongCounter implements LongCounter {
+ private final LongCounter delegate;
+
+ ReplaySafeLongCounter(LongCounter delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return delegate.isEnabled();
+ }
+
+ @Override
+ public void add(long value) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value);
+ }
+
+ @Override
+ public void add(long value, Attributes attributes) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value, attributes);
+ }
+
+ @Override
+ public void add(long value, Attributes attributes, Context context) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value, attributes, context);
+ }
+ }
+
+ private static final class ReplaySafeDoubleCounter implements DoubleCounter {
+ private final DoubleCounter delegate;
+
+ ReplaySafeDoubleCounter(DoubleCounter delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return delegate.isEnabled();
+ }
+
+ @Override
+ public void add(double value) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value);
+ }
+
+ @Override
+ public void add(double value, Attributes attributes) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value, attributes);
+ }
+
+ @Override
+ public void add(double value, Attributes attributes, Context context) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value, attributes, context);
+ }
+ }
+
+ private static final class ReplaySafeLongUpDownCounter implements LongUpDownCounter {
+ private final LongUpDownCounter delegate;
+
+ ReplaySafeLongUpDownCounter(LongUpDownCounter delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return delegate.isEnabled();
+ }
+
+ @Override
+ public void add(long value) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value);
+ }
+
+ @Override
+ public void add(long value, Attributes attributes) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value, attributes);
+ }
+
+ @Override
+ public void add(long value, Attributes attributes, Context context) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value, attributes, context);
+ }
+ }
+
+ private static final class ReplaySafeDoubleUpDownCounter implements DoubleUpDownCounter {
+ private final DoubleUpDownCounter delegate;
+
+ ReplaySafeDoubleUpDownCounter(DoubleUpDownCounter delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return delegate.isEnabled();
+ }
+
+ @Override
+ public void add(double value) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value);
+ }
+
+ @Override
+ public void add(double value, Attributes attributes) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value, attributes);
+ }
+
+ @Override
+ public void add(double value, Attributes attributes, Context context) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.add(value, attributes, context);
+ }
+ }
+
+ private static final class ReplaySafeDoubleHistogram implements DoubleHistogram {
+ private final DoubleHistogram delegate;
+
+ ReplaySafeDoubleHistogram(DoubleHistogram delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return delegate.isEnabled();
+ }
+
+ @Override
+ public void record(double value) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.record(value);
+ }
+
+ @Override
+ public void record(double value, Attributes attributes) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.record(value, attributes);
+ }
+
+ @Override
+ public void record(double value, Attributes attributes, Context context) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.record(value, attributes, context);
+ }
+ }
+
+ private static final class ReplaySafeLongHistogram implements LongHistogram {
+ private final LongHistogram delegate;
+
+ ReplaySafeLongHistogram(LongHistogram delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return delegate.isEnabled();
+ }
+
+ @Override
+ public void record(long value) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.record(value);
+ }
+
+ @Override
+ public void record(long value, Attributes attributes) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.record(value, attributes);
+ }
+
+ @Override
+ public void record(long value, Attributes attributes, Context context) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.record(value, attributes, context);
+ }
+ }
+
+ private static final class ReplaySafeDoubleGauge implements DoubleGauge {
+ private final DoubleGauge delegate;
+
+ ReplaySafeDoubleGauge(DoubleGauge delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return delegate.isEnabled();
+ }
+
+ @Override
+ public void set(double value) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.set(value);
+ }
+
+ @Override
+ public void set(double value, Attributes attributes) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.set(value, attributes);
+ }
+
+ @Override
+ public void set(double value, Attributes attributes, Context context) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.set(value, attributes, context);
+ }
+ }
+
+ private static final class ReplaySafeLongGauge implements LongGauge {
+ private final LongGauge delegate;
+
+ ReplaySafeLongGauge(LongGauge delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isEnabled() {
+ return delegate.isEnabled();
+ }
+
+ @Override
+ public void set(long value) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.set(value);
+ }
+
+ @Override
+ public void set(long value, Attributes attributes) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.set(value, attributes);
+ }
+
+ @Override
+ public void set(long value, Attributes attributes, Context context) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
+ return;
+ }
+ delegate.set(value, attributes, context);
+ }
+ }
+}
diff --git a/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeSpan.java b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeSpan.java
index 08b41b5642..937bcebbd6 100644
--- a/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeSpan.java
+++ b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeSpan.java
@@ -5,7 +5,6 @@
import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.SpanContext;
import io.opentelemetry.api.trace.StatusCode;
-import io.temporal.workflow.unsafe.WorkflowUnsafe;
import java.util.concurrent.TimeUnit;
/** Wraps a span so that replayed code does not end it, which would export a duplicate. */
@@ -18,9 +17,7 @@ public ReplaySafeSpan(Span delegate) {
@Override
public void end() {
- if (WorkflowUnsafe.isWorkflowThread()
- && WorkflowUnsafe.isSubjectToReplay()
- && WorkflowUnsafe.isReplaying()) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
return;
}
delegate.end();
@@ -28,9 +25,7 @@ public void end() {
@Override
public void end(long timestamp, TimeUnit unit) {
- if (WorkflowUnsafe.isWorkflowThread()
- && WorkflowUnsafe.isSubjectToReplay()
- && WorkflowUnsafe.isReplaying()) {
+ if (OpenTelemetrySuppression.shouldSuppress()) {
return;
}
delegate.end(timestamp, unit);
diff --git a/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeTracer.java b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeTracer.java
index daa68b8319..9bd2e7fe7a 100644
--- a/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeTracer.java
+++ b/contrib/temporal-opentelemetry-v2/src/main/java/io/temporal/opentelemetry/v2/internal/ReplaySafeTracer.java
@@ -11,7 +11,6 @@
import io.opentelemetry.context.ContextKey;
import io.opentelemetry.context.Scope;
import io.temporal.workflow.Workflow;
-import io.temporal.workflow.unsafe.WorkflowUnsafe;
import java.util.concurrent.TimeUnit;
import javax.annotation.Nonnull;
@@ -49,10 +48,7 @@ private static final class NamedStreamSpanBuilder implements SpanBuilder {
@Override
public Span startSpan() {
- if (WorkflowUnsafe.isWorkflowThread()
- && WorkflowUnsafe.isSubjectToReplay()
- && WorkflowUnsafe.isReplaying()
- && !startTimestampSet) {
+ if (OpenTelemetrySuppression.shouldSuppress() && !startTimestampSet) {
delegate.setStartTimestamp(Workflow.currentTimeMillis(), TimeUnit.MILLISECONDS);
}
try (Scope ignored = Context.current().with(TRACER_NAME, tracerName).makeCurrent()) {
diff --git a/contrib/temporal-opentelemetry-v2/src/test/java/io/temporal/opentelemetry/v2/MetricsTest.java b/contrib/temporal-opentelemetry-v2/src/test/java/io/temporal/opentelemetry/v2/MetricsTest.java
new file mode 100644
index 0000000000..0db80fa8a3
--- /dev/null
+++ b/contrib/temporal-opentelemetry-v2/src/test/java/io/temporal/opentelemetry/v2/MetricsTest.java
@@ -0,0 +1,108 @@
+package io.temporal.opentelemetry.v2;
+
+import static org.junit.Assert.assertEquals;
+
+import io.opentelemetry.api.GlobalOpenTelemetry;
+import io.opentelemetry.sdk.metrics.data.LongPointData;
+import io.temporal.activity.ActivityInterface;
+import io.temporal.activity.ActivityMethod;
+import io.temporal.activity.ActivityOptions;
+import io.temporal.api.common.v1.WorkflowExecution;
+import io.temporal.client.WorkflowClient;
+import io.temporal.client.WorkflowStub;
+import io.temporal.testing.WorkflowReplayer;
+import io.temporal.testing.internal.SDKTestWorkflowRule;
+import io.temporal.workflow.SignalMethod;
+import io.temporal.workflow.Workflow;
+import io.temporal.workflow.WorkflowInterface;
+import io.temporal.workflow.WorkflowMethod;
+import java.time.Duration;
+import java.util.Collection;
+import org.junit.Rule;
+import org.junit.Test;
+
+/**
+ * Verifies that workflow and activity metrics record during live execution and that replaying a
+ * workflow run does not emit duplicate metrics.
+ */
+public class MetricsTest extends OtelTestBase {
+ private static final String METER_NAME = "custom-metrics";
+ private static final String WORKFLOW_COUNTER = "custom_workflow_counter";
+ private static final String ACTIVITY_COUNTER = "custom_activity_counter";
+
+ @ActivityInterface
+ public interface TestActivity {
+ @ActivityMethod
+ void doActivity();
+ }
+
+ public static class TestActivityImpl implements TestActivity {
+ @Override
+ public void doActivity() {
+ GlobalOpenTelemetry.getMeter(METER_NAME).counterBuilder(ACTIVITY_COUNTER).build().add(1);
+ }
+ }
+
+ @WorkflowInterface
+ public interface TestWorkflow {
+ @WorkflowMethod
+ void run();
+
+ @SignalMethod
+ void proceed();
+ }
+
+ public static class TestWorkflowImpl implements TestWorkflow {
+ private boolean proceed;
+
+ @Override
+ public void run() {
+ GlobalOpenTelemetry.getMeter(METER_NAME).counterBuilder(WORKFLOW_COUNTER).build().add(1);
+
+ TestActivity activity =
+ Workflow.newActivityStub(
+ TestActivity.class,
+ ActivityOptions.newBuilder().setStartToCloseTimeout(Duration.ofSeconds(5)).build());
+ activity.doActivity();
+
+ Workflow.await(() -> proceed);
+ }
+
+ @Override
+ public void proceed() {
+ proceed = true;
+ }
+ }
+
+ @Rule
+ public SDKTestWorkflowRule testWorkflowRule =
+ newRuleBuilder(false)
+ .setWorkflowTypes(TestWorkflowImpl.class)
+ .setActivityImplementations(new TestActivityImpl())
+ .build();
+
+ @Test
+ public void liveExecutionRecordsAndReplayDoesNotDuplicate() throws Exception {
+ TestWorkflow workflow = testWorkflowRule.newWorkflowStub(TestWorkflow.class);
+ WorkflowExecution execution = WorkflowClient.start(workflow::run);
+
+ workflow.proceed();
+ WorkflowStub.fromTyped(workflow).getResult(Void.class);
+
+ assertEquals(1, singleLongValue(WORKFLOW_COUNTER));
+ assertEquals(1, singleLongValue(ACTIVITY_COUNTER));
+
+ // Replay the full workflow history and verify workflow metrics are suppressed on replay
+ WorkflowReplayer.replayWorkflowExecution(
+ testWorkflowRule.getExecutionHistory(execution.getWorkflowId()), TestWorkflowImpl.class);
+
+ assertEquals(1, singleLongValue(WORKFLOW_COUNTER));
+ assertEquals(1, singleLongValue(ACTIVITY_COUNTER));
+ }
+
+ private static long singleLongValue(String name) {
+ Collection points = requireMetricNamed(name).getLongSumData().getPoints();
+ assertEquals(name + " points: " + points, 1, points.size());
+ return points.iterator().next().getValue();
+ }
+}
diff --git a/contrib/temporal-opentelemetry-v2/src/test/java/io/temporal/opentelemetry/v2/OtelTestBase.java b/contrib/temporal-opentelemetry-v2/src/test/java/io/temporal/opentelemetry/v2/OtelTestBase.java
index 31276368d3..8a1d9ed6d4 100644
--- a/contrib/temporal-opentelemetry-v2/src/test/java/io/temporal/opentelemetry/v2/OtelTestBase.java
+++ b/contrib/temporal-opentelemetry-v2/src/test/java/io/temporal/opentelemetry/v2/OtelTestBase.java
@@ -6,6 +6,9 @@
import io.opentelemetry.api.GlobalOpenTelemetry;
import io.opentelemetry.api.common.AttributeKey;
+import io.opentelemetry.sdk.metrics.SdkMeterProvider;
+import io.opentelemetry.sdk.metrics.data.MetricData;
+import io.opentelemetry.sdk.testing.exporter.InMemoryMetricReader;
import io.opentelemetry.sdk.testing.exporter.InMemorySpanExporter;
import io.opentelemetry.sdk.trace.SdkTracerProvider;
import io.opentelemetry.sdk.trace.data.SpanData;
@@ -23,6 +26,7 @@
public abstract class OtelTestBase {
static final InMemorySpanExporter spanExporter = InMemorySpanExporter.create();
+ static final InMemoryMetricReader metricReader = InMemoryMetricReader.create();
private static ReplaySafeOpenTelemetry openTelemetry;
@BeforeClass
@@ -32,6 +36,7 @@ public static void registerGlobalOpenTelemetry() {
.setTracerProviderBuilder(
SdkTracerProvider.builder()
.addSpanProcessor(SimpleSpanProcessor.create(spanExporter)))
+ .setMeterProviderBuilder(SdkMeterProvider.builder().registerMetricReader(metricReader))
.build();
GlobalOpenTelemetry.set(openTelemetry);
}
@@ -75,6 +80,16 @@ static SpanData requireSpanNamed(List spans, String name) {
return null;
}
+ static MetricData requireMetricNamed(String name) {
+ for (MetricData metric : metricReader.collectAllMetrics()) {
+ if (metric.getName().equals(name)) {
+ return metric;
+ }
+ }
+ fail(name + " metric not found");
+ return null;
+ }
+
static String requireSpanAttribute(SpanData span, AttributeKey key) {
String value = span.getAttributes().get(key);
assertNotNull(key.getKey() + " attribute not found on " + span.getName(), value);