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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
110 changes: 110 additions & 0 deletions contrib/temporal-opentelemetry-v2/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
# OpenTelemetry v2 integration for the Temporal Java SDK

Module `io.temporal:temporal-opentelemetry-v2` provides replay-safe
OpenTelemetry tracing, metrics, and logs for Temporal.

## Setup

Use the same version as the rest of your Temporal Java SDK dependencies:

```groovy
implementation 'io.temporal:temporal-opentelemetry-v2:<temporal-java-sdk-version>'
// Add the exporters you use, such as this OTLP exporter.
implementation 'io.opentelemetry:opentelemetry-exporter-otlp'
```

Create a replay-safe OpenTelemetry instance, register it as the global, and
attach the plugin to your service stubs:

```java
import io.opentelemetry.api.GlobalOpenTelemetry;
import io.opentelemetry.exporter.otlp.trace.OtlpGrpcSpanExporter;
import io.opentelemetry.sdk.trace.SdkTracerProvider;
import io.opentelemetry.sdk.trace.export.BatchSpanProcessor;
import io.temporal.client.WorkflowClient;
import io.temporal.opentelemetry.v2.OpenTelemetryPlugin;
import io.temporal.opentelemetry.v2.ReplaySafeOpenTelemetry;
import io.temporal.serviceclient.WorkflowServiceStubs;
import io.temporal.serviceclient.WorkflowServiceStubsOptions;
import io.temporal.worker.WorkerFactory;

ReplaySafeOpenTelemetry openTelemetry =
ReplaySafeOpenTelemetry.newBuilder()
.setTracerProviderBuilder(

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Wire metric and log exporters in the setup sample

This sample configures only setTracerProviderBuilder, and ReplaySafeOpenTelemetry.Builder defaults the meter and logger providers to bare SdkMeterProvider.builder() and SdkLoggerProvider.builder() with no reader or processor. A user who follows this Setup section and then the Metrics and Logs sections verbatim records into providers that export nothing, with no hint why. Please add .setMeterProviderBuilder(SdkMeterProvider.builder().registerMetricReader(PeriodicMetricReader.builder(OtlpGrpcMetricExporter.builder().build()).build())) and .setLoggerProviderBuilder(SdkLoggerProvider.builder().addLogRecordProcessor(BatchLogRecordProcessor.builder(OtlpGrpcLogRecordExporter.builder().build()).build())) to the sample, or one sentence in each section saying an exporter must be configured on the corresponding provider.

SdkTracerProvider.builder()
.addSpanProcessor(
BatchSpanProcessor.builder(OtlpGrpcSpanExporter.builder().build()).build()))
.build();
GlobalOpenTelemetry.set(openTelemetry);

WorkflowServiceStubs service =
WorkflowServiceStubs.newServiceStubs(
WorkflowServiceStubsOptions.newBuilder()
.setPlugins(OpenTelemetryPlugin.newBuilder().build())
.build());

WorkflowClient client = WorkflowClient.newInstance(service);
WorkerFactory factory = WorkerFactory.newInstance(client);
```

Plugins configured on `WorkflowServiceStubsOptions` propagate to clients and
workers created from those stubs. It can also be configured directly on
`WorkflowClientOptions` or `WorkerFactoryOptions`.

## Tracing

The plugin propagates application trace context through Temporal headers.
Application spans remain connected across clients, workflows, activities, and
Nexus operations.

Set `OpenTelemetryPlugin.Builder.setAddTemporalSpans(true)` to emit spans for
operations such as `StartWorkflow`, `RunWorkflow`, `RunActivity`, and
`ContinueAsNew`.

Create spans in workflow, client, and activity code with the standard OpenTelemetry API:

```java
import io.opentelemetry.api.GlobalOpenTelemetry;
import io.opentelemetry.api.trace.Span;
import io.opentelemetry.context.Scope;

Span span =
GlobalOpenTelemetry.getTracer("my-workflows")
.spanBuilder("my-span")
.startSpan();
try (Scope ignored = span.makeCurrent()) {
activity.doWork();
} finally {
span.end();
}
```

## Metrics

Create synchronous metrics in workflow, client, and activity code with the standard OpenTelemetry API:

```java
import io.opentelemetry.api.GlobalOpenTelemetry;
import io.opentelemetry.api.metrics.LongCounter;

LongCounter counter =
GlobalOpenTelemetry.getMeter("my-workflows")
.counterBuilder("workflow.items.processed")
.build();
counter.add(1);
```

## Logs

Create log records in workflow, client, and activity code with the standard OpenTelemetry API:

```java
import io.opentelemetry.api.GlobalOpenTelemetry;

GlobalOpenTelemetry.get()
.getLogsBridge()
.get("my-workflows")
.logRecordBuilder()
.setBody("workflow step completed")
.emit();
```
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@
import io.opentelemetry.api.GlobalOpenTelemetry;
import io.opentelemetry.api.OpenTelemetry;
import io.opentelemetry.api.baggage.propagation.W3CBaggagePropagator;
import io.opentelemetry.api.logs.Logger;
import io.opentelemetry.api.logs.LoggerBuilder;
import io.opentelemetry.api.logs.LoggerProvider;
import io.opentelemetry.api.metrics.Meter;
import io.opentelemetry.api.metrics.MeterBuilder;
Expand All @@ -21,31 +23,30 @@
import io.opentelemetry.sdk.trace.SdkTracerProviderBuilder;
import io.temporal.common.Experimental;
import io.temporal.opentelemetry.v2.internal.ReplaySafeIdGenerator;
import io.temporal.opentelemetry.v2.internal.ReplaySafeLogger;
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 and meters obtained from it are replay safe inside
* workflows.
* {@code GlobalOpenTelemetry.set}; tracers, meters, and loggers obtained from it are replay safe
* inside workflows.
*/
@Experimental
public final class ReplaySafeOpenTelemetry implements OpenTelemetry, Closeable {
private final ReplaySafeTracerProvider tracerProvider;
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;
private final ReplaySafeLoggerProvider loggerProvider;
private final ContextPropagators propagators;

private ReplaySafeOpenTelemetry(Builder builder) {
this.tracerProvider =
new ReplaySafeTracerProvider(
builder.tracerProviderBuilder.setIdGenerator(new ReplaySafeIdGenerator()).build());
this.meterProvider = new ReplaySafeMeterProvider(builder.meterProviderBuilder.build());
this.loggerProvider = builder.loggerProviderBuilder.build();
this.loggerProvider = new ReplaySafeLoggerProvider(builder.loggerProviderBuilder.build());
Comment thread
patbeqo marked this conversation as resolved.
this.propagators = builder.propagators;
}

Expand Down Expand Up @@ -207,6 +208,54 @@ public Meter build() {
}
}

private static final class ReplaySafeLoggerProvider implements LoggerProvider, Closeable {
private final SdkLoggerProvider delegate;

private ReplaySafeLoggerProvider(SdkLoggerProvider delegate) {
this.delegate = delegate;
}

@Override
public Logger get(@Nonnull String instrumentationScopeName) {
return new ReplaySafeLogger(delegate.get(instrumentationScopeName));
}

@Override
public LoggerBuilder loggerBuilder(@Nonnull String instrumentationScopeName) {
return new ReplaySafeLoggerBuilder(delegate.loggerBuilder(instrumentationScopeName));
}

@Override
public void close() {
delegate.close();
}
}

private static final class ReplaySafeLoggerBuilder implements LoggerBuilder {
private final LoggerBuilder delegate;

private ReplaySafeLoggerBuilder(LoggerBuilder delegate) {
this.delegate = delegate;
}

@Override
public LoggerBuilder setSchemaUrl(@Nonnull String schemaUrl) {
delegate.setSchemaUrl(schemaUrl);
return this;
}

@Override
public LoggerBuilder setInstrumentationVersion(@Nonnull String instrumentationScopeVersion) {
delegate.setInstrumentationVersion(instrumentationScopeVersion);
return this;
}
Comment thread
patbeqo marked this conversation as resolved.

@Override
public Logger build() {
return new ReplaySafeLogger(delegate.build());
}
}

private static final class ReplaySafeTracerBuilder implements TracerBuilder {
private final TracerBuilder delegate;
private final String instrumentationScopeName;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
package io.temporal.opentelemetry.v2.internal;

import io.opentelemetry.api.common.AttributeKey;
import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.common.Value;
import io.opentelemetry.api.logs.LogRecordBuilder;
import io.opentelemetry.api.logs.Logger;
import io.opentelemetry.api.logs.Severity;
import io.opentelemetry.context.Context;
import java.time.Instant;
import java.util.concurrent.TimeUnit;

/**
* Wraps a logger so the records it builds are not emitted by replaying workflow code, which would
* otherwise emit them again on every replay.
*/
public final class ReplaySafeLogger implements Logger {
private final Logger delegate;

public ReplaySafeLogger(Logger delegate) {
this.delegate = delegate;
}

@Override
public LogRecordBuilder logRecordBuilder() {
return new ReplaySafeLogRecordBuilder(delegate.logRecordBuilder());
}

@Override
public boolean isEnabled(Severity severity, Context context) {
return !OpenTelemetrySuppression.shouldSuppress() && delegate.isEnabled(severity, context);
}

@Override
public boolean isEnabled(Severity severity) {
return !OpenTelemetrySuppression.shouldSuppress() && delegate.isEnabled(severity);
}

private static final class ReplaySafeLogRecordBuilder implements LogRecordBuilder {
private final LogRecordBuilder delegate;

ReplaySafeLogRecordBuilder(LogRecordBuilder delegate) {
this.delegate = delegate;
}

@Override
public void emit() {
if (OpenTelemetrySuppression.shouldSuppress()) {
return;
}
delegate.emit();
}

@Override
public LogRecordBuilder setTimestamp(long timestamp, TimeUnit unit) {
delegate.setTimestamp(timestamp, unit);
return this;
}

@Override
public LogRecordBuilder setTimestamp(Instant instant) {
delegate.setTimestamp(instant);
return this;
}

@Override
public LogRecordBuilder setObservedTimestamp(long timestamp, TimeUnit unit) {
delegate.setObservedTimestamp(timestamp, unit);
return this;
}

@Override
public LogRecordBuilder setObservedTimestamp(Instant instant) {
delegate.setObservedTimestamp(instant);
return this;
}

@Override
public LogRecordBuilder setContext(Context context) {
delegate.setContext(context);
return this;
}

@Override
public LogRecordBuilder setSeverity(Severity severity) {
delegate.setSeverity(severity);
return this;
}

@Override
public LogRecordBuilder setSeverityText(String severityText) {
delegate.setSeverityText(severityText);
return this;
}

@Override
public LogRecordBuilder setBody(String body) {
delegate.setBody(body);
return this;
}

@Override
public LogRecordBuilder setBody(Value<?> body) {
delegate.setBody(body);
return this;
}

@Override
public LogRecordBuilder setAllAttributes(Attributes attributes) {
delegate.setAllAttributes(attributes);
return this;
}

@Override
public <T> LogRecordBuilder setAttribute(AttributeKey<T> key, T value) {
delegate.setAttribute(key, value);
return this;
}

@Override
public LogRecordBuilder setAttribute(String key, String value) {
delegate.setAttribute(key, value);
return this;
}

@Override
public LogRecordBuilder setAttribute(String key, long value) {
delegate.setAttribute(key, value);
return this;
}

@Override
public LogRecordBuilder setAttribute(String key, double value) {
delegate.setAttribute(key, value);
return this;
}

@Override
public LogRecordBuilder setAttribute(String key, boolean value) {
delegate.setAttribute(key, value);
return this;
}

@Override
public LogRecordBuilder setAttribute(String key, int value) {
delegate.setAttribute(key, value);
return this;
}

@Override
public LogRecordBuilder setAttribute(String key, Value<?> value) {
delegate.setAttribute(key, value);
return this;
}

@Override
public LogRecordBuilder setEventName(String eventName) {
delegate.setEventName(eventName);
return this;
}

@Override
public LogRecordBuilder setException(Throwable throwable) {
delegate.setException(throwable);
return this;
}
}
}
Loading
Loading