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);