diff --git a/instrumentation-api-incubator/build.gradle.kts b/instrumentation-api-incubator/build.gradle.kts
index 1a59ad1364a8..5d87535b7a2b 100644
--- a/instrumentation-api-incubator/build.gradle.kts
+++ b/instrumentation-api-incubator/build.gradle.kts
@@ -92,6 +92,7 @@ tasks {
testClassesDirs = sourceSets.test.get().output.classesDirs
classpath = sourceSets.test.get().runtimeClasspath
jvmArgs("-Dotel.semconv-stability.opt-in=database,code,service.peer,rpc")
+ jvmArgs("-Dotel.semconv-stability.preview=messaging")
inputs.dir(jflexOutputDir)
}
@@ -99,6 +100,7 @@ tasks {
testClassesDirs = sourceSets.test.get().output.classesDirs
classpath = sourceSets.test.get().runtimeClasspath
jvmArgs("-Dotel.semconv-stability.opt-in=database/dup,code/dup,service.peer/dup,rpc/dup")
+ jvmArgs("-Dotel.semconv-stability.preview=messaging/dup")
inputs.dir(jflexOutputDir)
}
@@ -117,6 +119,11 @@ tasks {
}
check {
- dependsOn(testStableSemconv, testBothSemconv, testExceptionSignalLogs, testExceptionSignalLogsDup)
+ dependsOn(
+ testStableSemconv,
+ testBothSemconv,
+ testExceptionSignalLogs,
+ testExceptionSignalLogsDup,
+ )
}
}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessageOperation.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessageOperation.java
index e7ef5c479ed0..c1062bd0d056 100644
--- a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessageOperation.java
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessageOperation.java
@@ -5,24 +5,19 @@
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
-import java.util.Locale;
-
-/**
- * Represents type of operations
- * that may be used in a messaging system.
- */
+/** Represents an operation that may be used in a messaging system. */
public enum MessageOperation {
- PUBLISH,
- RECEIVE,
- PROCESS;
+ PUBLISH(MessagingOperationType.SEND),
+ RECEIVE(MessagingOperationType.RECEIVE),
+ PROCESS(MessagingOperationType.PROCESS);
+
+ private final MessagingOperationType operationType;
+
+ MessageOperation(MessagingOperationType operationType) {
+ this.operationType = operationType;
+ }
- /**
- * Returns the operation name as defined in the
- * specification.
- */
- String operationName() {
- return name().toLowerCase(Locale.ROOT);
+ MessagingOperationType type() {
+ return operationType;
}
}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractor.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractor.java
index 3ad89017a047..3a680fa8f625 100644
--- a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractor.java
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractor.java
@@ -5,6 +5,10 @@
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitOldMessagingSemconv;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
+import static io.opentelemetry.semconv.ErrorAttributes.ERROR_TYPE;
+
import io.opentelemetry.api.common.AttributeKey;
import io.opentelemetry.api.common.AttributesBuilder;
import io.opentelemetry.context.Context;
@@ -17,7 +21,7 @@
/**
* Extractor of messaging
+ * href="https://github.com/open-telemetry/semantic-conventions/blob/v1.43.0/docs/messaging/messaging-spans.md">messaging
* attributes.
*
*
This class delegates to a type-specific {@link MessagingAttributesGetter} for individual
@@ -29,8 +33,10 @@ public final class MessagingAttributesExtractor
// copied from MessagingIncubatingAttributes
private static final AttributeKey MESSAGING_BATCH_MESSAGE_COUNT =
AttributeKey.longKey("messaging.batch.message_count");
- private static final AttributeKey MESSAGING_CLIENT_ID =
+ private static final AttributeKey MESSAGING_CLIENT_ID_OLD =
AttributeKey.stringKey("messaging.client_id");
+ private static final AttributeKey MESSAGING_CLIENT_ID =
+ AttributeKey.stringKey("messaging.client.id");
private static final AttributeKey MESSAGING_DESTINATION_ANONYMOUS =
AttributeKey.booleanKey("messaging.destination.anonymous");
private static final AttributeKey MESSAGING_DESTINATION_NAME =
@@ -51,49 +57,78 @@ public final class MessagingAttributesExtractor
AttributeKey.stringKey("messaging.message.id");
private static final AttributeKey MESSAGING_OPERATION =
AttributeKey.stringKey("messaging.operation");
+ private static final AttributeKey MESSAGING_OPERATION_NAME =
+ AttributeKey.stringKey("messaging.operation.name");
+ private static final AttributeKey MESSAGING_OPERATION_TYPE =
+ AttributeKey.stringKey("messaging.operation.type");
private static final AttributeKey MESSAGING_SYSTEM =
AttributeKey.stringKey("messaging.system");
static final String TEMP_DESTINATION_NAME = "(temporary)";
- /**
- * Creates the messaging attributes extractor for the given {@link MessageOperation operation}
- * with default configuration.
- */
+ /** Creates the messaging attributes extractor for the given operation type. */
+ public static AttributesExtractor createForOperationType(
+ MessagingAttributesGetter getter, MessagingOperationType operationType) {
+ return builderForOperationType(getter, operationType).build();
+ }
+
+ /** Creates the messaging attributes extractor for the given operation. */
public static AttributesExtractor create(
- MessagingAttributesGetter getter, MessageOperation operation) {
+ MessagingAttributesGetter getter, @Nullable MessageOperation operation) {
return builder(getter, operation).build();
}
/**
- * Returns a new {@link MessagingAttributesExtractorBuilder} for the given {@link MessageOperation
- * operation} that can be used to configure the messaging attributes extractor.
+ * Returns a new {@link MessagingAttributesExtractorBuilder} configured for the given operation
+ * type.
*/
+ public static
+ MessagingAttributesExtractorBuilder builderForOperationType(
+ MessagingAttributesGetter getter,
+ @Nullable MessagingOperationType operationType) {
+ return new MessagingAttributesExtractorBuilder<>(getter, operationType, true);
+ }
+
+ /** Returns a new messaging attributes extractor builder for the given operation. */
public static MessagingAttributesExtractorBuilder builder(
- MessagingAttributesGetter getter, MessageOperation operation) {
- return new MessagingAttributesExtractorBuilder<>(getter, operation);
+ MessagingAttributesGetter getter, @Nullable MessageOperation operation) {
+ return new MessagingAttributesExtractorBuilder<>(
+ getter, operation == null ? null : operation.type(), false);
}
private final MessagingAttributesGetter getter;
- private final MessageOperation operation;
+ @Nullable private final MessagingOperationType operationType;
+ @Nullable private final String operationName;
+ private final boolean supportsStableSemconv;
private final List capturedHeaders;
MessagingAttributesExtractor(
MessagingAttributesGetter getter,
- MessageOperation operation,
+ @Nullable MessagingOperationType operationType,
+ @Nullable String operationName,
+ boolean supportsStableSemconv,
List capturedHeaders) {
this.getter = getter;
- this.operation = operation;
+ this.operationType = operationType;
+ this.operationName = operationName;
+ this.supportsStableSemconv = supportsStableSemconv;
this.capturedHeaders = new ArrayList<>(capturedHeaders);
}
@Override
public void onStart(AttributesBuilder attributes, Context parentContext, REQUEST request) {
+ boolean emitOldSemconv = !supportsStableSemconv || emitOldMessagingSemconv();
+ boolean emitStableSemconv = supportsStableSemconv && emitStableMessagingSemconv();
attributes.put(MESSAGING_SYSTEM, getter.getSystem(request));
boolean isTemporaryDestination = getter.isTemporaryDestination(request);
if (isTemporaryDestination) {
attributes.put(MESSAGING_DESTINATION_TEMPORARY, true);
- attributes.put(MESSAGING_DESTINATION_NAME, TEMP_DESTINATION_NAME);
+ if (emitStableSemconv) {
+ attributes.put(MESSAGING_DESTINATION_NAME, getter.getDestination(request));
+ attributes.put(MESSAGING_DESTINATION_TEMPLATE, getter.getDestinationTemplate(request));
+ } else {
+ attributes.put(MESSAGING_DESTINATION_NAME, TEMP_DESTINATION_NAME);
+ }
} else {
attributes.put(MESSAGING_DESTINATION_NAME, getter.getDestination(request));
attributes.put(MESSAGING_DESTINATION_TEMPLATE, getter.getDestinationTemplate(request));
@@ -106,9 +141,20 @@ public void onStart(AttributesBuilder attributes, Context parentContext, REQUEST
attributes.put(MESSAGING_MESSAGE_CONVERSATION_ID, getter.getConversationId(request));
attributes.put(MESSAGING_MESSAGE_BODY_SIZE, getter.getMessageBodySize(request));
attributes.put(MESSAGING_MESSAGE_ENVELOPE_SIZE, getter.getMessageEnvelopeSize(request));
- attributes.put(MESSAGING_CLIENT_ID, getter.getClientId(request));
- if (operation != null) {
- attributes.put(MESSAGING_OPERATION, operation.operationName());
+ if (emitOldSemconv) {
+ attributes.put(MESSAGING_CLIENT_ID_OLD, getter.getClientId(request));
+ }
+ if (emitStableSemconv) {
+ attributes.put(MESSAGING_CLIENT_ID, getter.getClientId(request));
+ }
+ if (emitOldSemconv && operationType != null) {
+ attributes.put(MESSAGING_OPERATION, operationType.defaultOperationName());
+ }
+ if (emitStableSemconv) {
+ attributes.put(MESSAGING_OPERATION_NAME, operationName);
+ if (operationType != null) {
+ attributes.put(MESSAGING_OPERATION_TYPE, operationType.value());
+ }
}
}
@@ -121,6 +167,13 @@ public void onEnd(
@Nullable Throwable error) {
attributes.put(MESSAGING_MESSAGE_ID, getter.getMessageId(request, response));
attributes.put(MESSAGING_BATCH_MESSAGE_COUNT, getter.getBatchMessageCount(request, response));
+ if (supportsStableSemconv && emitStableMessagingSemconv()) {
+ String errorType = getter.getErrorType(request, response, error);
+ if (errorType == null && error != null) {
+ errorType = error.getClass().getName();
+ }
+ attributes.put(ERROR_TYPE, errorType);
+ }
for (String name : capturedHeaders) {
List values = getter.getMessageHeader(request, name);
@@ -137,17 +190,21 @@ public void onEnd(
@Override
@Nullable
public SpanKey internalGetSpanKey() {
- if (operation == null) {
+ if (operationType == null) {
return null;
}
- switch (operation) {
- case PUBLISH:
+ switch (operationType) {
+ case CREATE:
+ return SpanKey.PRODUCER_CREATE;
+ case SEND:
return SpanKey.PRODUCER;
case RECEIVE:
return SpanKey.CONSUMER_RECEIVE;
case PROCESS:
return SpanKey.CONSUMER_PROCESS;
+ case SETTLE:
+ return SpanKey.CONSUMER_SETTLE;
}
throw new IllegalStateException("Can't possibly happen");
}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorBuilder.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorBuilder.java
index 61f6c6916549..925047dabecf 100644
--- a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorBuilder.java
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorBuilder.java
@@ -6,24 +6,43 @@
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
import static java.util.Collections.emptyList;
+import static java.util.Objects.requireNonNull;
import com.google.errorprone.annotations.CanIgnoreReturnValue;
import io.opentelemetry.instrumentation.api.instrumenter.AttributesExtractor;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
+import javax.annotation.Nullable;
/** A builder of {@link MessagingAttributesExtractor}. */
public final class MessagingAttributesExtractorBuilder {
final MessagingAttributesGetter getter;
- final MessageOperation operation;
+ @Nullable private final MessagingOperationType operationType;
+ @Nullable private String operationName;
+ private final boolean supportsStableSemconv;
List capturedHeaders = emptyList();
MessagingAttributesExtractorBuilder(
- MessagingAttributesGetter getter, MessageOperation operation) {
+ MessagingAttributesGetter getter,
+ @Nullable MessagingOperationType operationType,
+ boolean supportsStableSemconv) {
this.getter = getter;
- this.operation = operation;
+ this.operationType = operationType;
+ this.operationName = operationType == null ? null : operationType.defaultOperationName();
+ this.supportsStableSemconv = supportsStableSemconv;
+ }
+
+ /** Configures the system-specific operation name emitted as {@code messaging.operation.name}. */
+ @CanIgnoreReturnValue
+ public MessagingAttributesExtractorBuilder setOperationName(
+ String operationName) {
+ if (!supportsStableSemconv) {
+ throw new IllegalStateException("Operation name is not configurable for legacy builders");
+ }
+ this.operationName = requireNonNull(operationName, "operationName");
+ return this;
}
/**
@@ -47,6 +66,10 @@ public MessagingAttributesExtractorBuilder setCapturedHeaders
* MessagingAttributesExtractorBuilder}.
*/
public AttributesExtractor build() {
- return new MessagingAttributesExtractor<>(getter, operation, capturedHeaders);
+ if (supportsStableSemconv) {
+ requireNonNull(operationName, "operationName");
+ }
+ return new MessagingAttributesExtractor<>(
+ getter, operationType, operationName, supportsStableSemconv, capturedHeaders);
}
}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesGetter.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesGetter.java
index ae0b29471cf0..eaea16aaf74f 100644
--- a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesGetter.java
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesGetter.java
@@ -55,6 +55,21 @@ default String getDestinationPartitionId(REQUEST request) {
return null;
}
+ /**
+ * Returns a description of a class of error the operation ended with.
+ *
+ * If this method returns {@code null}, the exception class name (if any) will be used as error
+ * type.
+ *
+ *
The cardinality of the error type should be low. The instrumentations implementing this
+ * method are recommended to document the custom values they support.
+ */
+ @Nullable
+ default String getErrorType(
+ REQUEST request, @Nullable RESPONSE response, @Nullable Throwable error) {
+ return null;
+ }
+
/**
* Extracts all values of header named {@code name} from the request, or an empty list if there
* were none.
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingConsumerMetrics.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingConsumerMetrics.java
index 87a98b106d74..76a62a828e84 100644
--- a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingConsumerMetrics.java
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingConsumerMetrics.java
@@ -5,6 +5,9 @@
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitOldMessagingSemconv;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
+import static io.opentelemetry.semconv.ErrorAttributes.ERROR_TYPE;
import static java.util.concurrent.TimeUnit.SECONDS;
import static java.util.logging.Level.FINE;
@@ -23,10 +26,11 @@
import io.opentelemetry.instrumentation.api.instrumenter.OperationMetrics;
import io.opentelemetry.instrumentation.api.internal.OperationMetricsUtil;
import java.util.logging.Logger;
+import javax.annotation.Nullable;
/**
* {@link OperationListener} which keeps track of Consumer
+ * href="https://github.com/open-telemetry/semantic-conventions/blob/v1.43.0/docs/messaging/messaging-metrics.md#consumer-metrics">consumer
* metrics.
*/
public final class MessagingConsumerMetrics implements OperationListener {
@@ -35,39 +39,70 @@ public final class MessagingConsumerMetrics implements OperationListener {
// copied from MessagingIncubatingAttributes
private static final AttributeKey MESSAGING_BATCH_MESSAGE_COUNT =
AttributeKey.longKey("messaging.batch.message_count");
+ private static final AttributeKey MESSAGING_OPERATION =
+ AttributeKey.stringKey("messaging.operation");
+ private static final AttributeKey MESSAGING_OPERATION_TYPE =
+ AttributeKey.stringKey("messaging.operation.type");
private static final ContextKey MESSAGING_CONSUMER_METRICS_STATE =
ContextKey.named("messaging-consumer-metrics-state");
private static final Logger logger = Logger.getLogger(MessagingConsumerMetrics.class.getName());
- private final DoubleHistogram receiveDurationHistogram;
- private final LongCounter receiveMessageCount;
+ private final boolean supportsStableSemconv;
+ private final boolean consumedMessagesOnly;
+ private final boolean enabled;
+ @Nullable private final DoubleHistogram receiveDurationHistogram;
+ @Nullable private final LongCounter receiveMessageCount;
+ @Nullable private final DoubleHistogram clientOperationDurationHistogram;
+ @Nullable private final LongCounter consumedMessagesCounter;
- private MessagingConsumerMetrics(Meter meter) {
- DoubleHistogramBuilder durationBuilder =
- meter
- .histogramBuilder("messaging.receive.duration")
- .setDescription("Measures the duration of receive operation.")
- .setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
- .setUnit("s");
- MessagingMetricsAdvice.applyReceiveDurationAdvice(durationBuilder);
- receiveDurationHistogram = durationBuilder.build();
+ private MessagingConsumerMetrics(Meter meter, boolean supportsStableSemconv) {
+ this(meter, supportsStableSemconv, false);
+ }
- LongCounterBuilder longCounterBuilder =
- meter
- .counterBuilder("messaging.receive.messages")
- .setDescription("Measures the number of received messages.")
- .setUnit("{message}");
- MessagingMetricsAdvice.applyReceiveMessagesAdvice(longCounterBuilder);
- receiveMessageCount = longCounterBuilder.build();
+ private MessagingConsumerMetrics(
+ Meter meter, boolean supportsStableSemconv, boolean consumedMessagesOnly) {
+ this.supportsStableSemconv = supportsStableSemconv;
+ this.consumedMessagesOnly = consumedMessagesOnly;
+ boolean emitOldSemconv = !supportsStableSemconv || emitOldMessagingSemconv();
+ boolean emitStableSemconv = supportsStableSemconv && emitStableMessagingSemconv();
+ receiveDurationHistogram =
+ !consumedMessagesOnly && emitOldSemconv ? buildReceiveDuration(meter) : null;
+ receiveMessageCount =
+ !consumedMessagesOnly && emitOldSemconv ? buildReceiveMessages(meter) : null;
+ clientOperationDurationHistogram =
+ !consumedMessagesOnly && emitStableSemconv ? buildClientOperationDuration(meter) : null;
+ consumedMessagesCounter = emitStableSemconv ? buildConsumedMessages(meter) : null;
+ enabled =
+ receiveDurationHistogram != null
+ || receiveMessageCount != null
+ || clientOperationDurationHistogram != null
+ || consumedMessagesCounter != null;
}
+ /** Returns metrics for extractors configured with {@link MessageOperation}. */
public static OperationMetrics get() {
- return OperationMetricsUtil.create("messaging consumer", MessagingConsumerMetrics::new);
+ return OperationMetricsUtil.create(
+ "messaging consumer", meter -> new MessagingConsumerMetrics(meter, false));
+ }
+
+ /** Returns metrics for extractors configured with {@link MessagingOperationType}. */
+ public static OperationMetrics getForOperationType() {
+ return OperationMetricsUtil.create(
+ "messaging consumer", meter -> new MessagingConsumerMetrics(meter, true));
+ }
+
+ /** Returns only the stable consumed-messages metric for a delivered message. */
+ public static OperationMetrics getConsumedMessages() {
+ return OperationMetricsUtil.create(
+ "messaging consumed messages", meter -> new MessagingConsumerMetrics(meter, true, true));
}
@Override
@CanIgnoreReturnValue
public Context onStart(Context context, Attributes startAttributes, long startNanos) {
+ if (!enabled) {
+ return context;
+ }
return context.with(
MESSAGING_CONSUMER_METRICS_STATE,
new AutoValue_MessagingConsumerMetrics_State(startAttributes, startNanos));
@@ -75,6 +110,9 @@ public Context onStart(Context context, Attributes startAttributes, long startNa
@Override
public void onEnd(Context context, Attributes endAttributes, long endNanos) {
+ if (!enabled) {
+ return;
+ }
MessagingConsumerMetrics.State state = context.get(MESSAGING_CONSUMER_METRICS_STATE);
if (state == null) {
logger.log(
@@ -85,21 +123,104 @@ public void onEnd(Context context, Attributes endAttributes, long endNanos) {
}
Attributes attributes = state.startAttributes().toBuilder().putAll(endAttributes).build();
- receiveDurationHistogram.record(
- (endNanos - state.startTimeNanos()) / NANOS_PER_S, attributes, context);
+ double duration = (endNanos - state.startTimeNanos()) / NANOS_PER_S;
+ String operationType = attributes.get(MESSAGING_OPERATION_TYPE);
+ boolean recordsLegacyReceive =
+ !supportsStableSemconv
+ || "receive".equals(attributes.get(MESSAGING_OPERATION))
+ || MessagingOperationType.RECEIVE.value().equals(operationType);
+ if (receiveDurationHistogram != null && recordsLegacyReceive) {
+ receiveDurationHistogram.record(duration, attributes, context);
+ }
+ // Metric view attribute advice can only select keys statically. The concrete destination name
+ // must be omitted when a template is available or the destination is temporary or anonymous,
+ // so this conditional requirement must be enforced before recording.
+ Attributes filteredAttributes =
+ clientOperationDurationHistogram != null || consumedMessagesCounter != null
+ ? MessagingMetricsAdvice.filterAttributes(attributes)
+ : attributes;
+ if (clientOperationDurationHistogram != null
+ && !MessagingOperationType.PROCESS.value().equals(operationType)) {
+ clientOperationDurationHistogram.record(duration, filteredAttributes, context);
+ }
- long receiveMessagesCount = getReceiveMessagesCount(state.startAttributes(), endAttributes);
- receiveMessageCount.add(receiveMessagesCount, attributes, context);
+ Long batchMessageCount = getBatchMessageCount(state.startAttributes(), endAttributes);
+ if (receiveMessageCount != null && recordsLegacyReceive) {
+ long receiveMessagesCount = batchMessageCount == null ? 1 : batchMessageCount;
+ if (!supportsStableSemconv || receiveMessagesCount > 0) {
+ receiveMessageCount.add(receiveMessagesCount, attributes, context);
+ }
+ }
+ if (consumedMessagesCounter != null
+ && (consumedMessagesOnly || MessagingOperationType.RECEIVE.value().equals(operationType))) {
+ long consumedMessagesCount =
+ getConsumedMessagesCount(attributes, batchMessageCount, consumedMessagesOnly);
+ if (consumedMessagesCount > 0) {
+ consumedMessagesCounter.add(consumedMessagesCount, filteredAttributes, context);
+ }
+ }
}
- private static long getReceiveMessagesCount(Attributes... attributesList) {
+ @Nullable
+ private static Long getBatchMessageCount(Attributes... attributesList) {
for (Attributes attributes : attributesList) {
Long value = attributes.get(MESSAGING_BATCH_MESSAGE_COUNT);
if (value != null) {
return value;
}
}
- return 1;
+ return null;
+ }
+
+ private static long getConsumedMessagesCount(
+ Attributes attributes, @Nullable Long batchMessageCount, boolean consumedMessagesOnly) {
+ if (batchMessageCount != null) {
+ return batchMessageCount;
+ }
+ return consumedMessagesOnly || attributes.get(ERROR_TYPE) == null ? 1 : 0;
+ }
+
+ private static DoubleHistogram buildReceiveDuration(Meter meter) {
+ DoubleHistogramBuilder builder =
+ meter
+ .histogramBuilder("messaging.receive.duration")
+ .setDescription("Measures the duration of receive operation.")
+ .setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
+ .setUnit("s");
+ MessagingMetricsAdvice.applyOldDurationAdvice(builder);
+ return builder.build();
+ }
+
+ private static LongCounter buildReceiveMessages(Meter meter) {
+ LongCounterBuilder builder =
+ meter
+ .counterBuilder("messaging.receive.messages")
+ .setDescription("Measures the number of received messages.")
+ .setUnit("{message}");
+ MessagingMetricsAdvice.applyOldMessagesAdvice(builder);
+ return builder.build();
+ }
+
+ private static DoubleHistogram buildClientOperationDuration(Meter meter) {
+ DoubleHistogramBuilder builder =
+ meter
+ .histogramBuilder("messaging.client.operation.duration")
+ .setDescription(
+ "Duration of messaging operation initiated by a producer or consumer client.")
+ .setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
+ .setUnit("s");
+ MessagingMetricsAdvice.applyClientOperationDurationAdvice(builder);
+ return builder.build();
+ }
+
+ private static LongCounter buildConsumedMessages(Meter meter) {
+ LongCounterBuilder builder =
+ meter
+ .counterBuilder("messaging.client.consumed.messages")
+ .setDescription("Number of messages that were delivered to the application.")
+ .setUnit("{message}");
+ MessagingMetricsAdvice.applyConsumedMessagesAdvice(builder);
+ return builder.build();
}
@AutoValue
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingMetricsAdvice.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingMetricsAdvice.java
index ddc4115ea7f4..defc418fc406 100644
--- a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingMetricsAdvice.java
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingMetricsAdvice.java
@@ -12,10 +12,12 @@
import static java.util.Collections.unmodifiableList;
import io.opentelemetry.api.common.AttributeKey;
+import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.incubator.metrics.ExtendedDoubleHistogramBuilder;
import io.opentelemetry.api.incubator.metrics.ExtendedLongCounterBuilder;
import io.opentelemetry.api.metrics.DoubleHistogramBuilder;
import io.opentelemetry.api.metrics.LongCounterBuilder;
+import java.util.ArrayList;
import java.util.List;
final class MessagingMetricsAdvice {
@@ -28,14 +30,26 @@ final class MessagingMetricsAdvice {
AttributeKey.stringKey("messaging.system");
private static final AttributeKey MESSAGING_DESTINATION_NAME =
AttributeKey.stringKey("messaging.destination.name");
+ private static final AttributeKey MESSAGING_DESTINATION_ANONYMOUS =
+ AttributeKey.booleanKey("messaging.destination.anonymous");
+ private static final AttributeKey MESSAGING_DESTINATION_TEMPORARY =
+ AttributeKey.booleanKey("messaging.destination.temporary");
private static final AttributeKey MESSAGING_OPERATION =
AttributeKey.stringKey("messaging.operation");
+ private static final AttributeKey MESSAGING_OPERATION_NAME =
+ AttributeKey.stringKey("messaging.operation.name");
+ private static final AttributeKey MESSAGING_OPERATION_TYPE =
+ AttributeKey.stringKey("messaging.operation.type");
+ private static final AttributeKey MESSAGING_CONSUMER_GROUP_NAME =
+ AttributeKey.stringKey("messaging.consumer.group.name");
+ private static final AttributeKey MESSAGING_DESTINATION_SUBSCRIPTION_NAME =
+ AttributeKey.stringKey("messaging.destination.subscription.name");
private static final AttributeKey MESSAGING_DESTINATION_PARTITION_ID =
AttributeKey.stringKey("messaging.destination.partition.id");
private static final AttributeKey MESSAGING_DESTINATION_TEMPLATE =
AttributeKey.stringKey("messaging.destination.template");
- private static final List> MESSAGING_ATTRIBUTES =
+ private static final List> OLD_ATTRIBUTES =
asList(
MESSAGING_SYSTEM,
MESSAGING_DESTINATION_NAME,
@@ -46,25 +60,88 @@ final class MessagingMetricsAdvice {
SERVER_PORT,
SERVER_ADDRESS);
- static void applyPublishDurationAdvice(DoubleHistogramBuilder builder) {
+ private static final List> CLIENT_OPERATION_DURATION_ATTRIBUTES =
+ buildAttributes(true, true, true);
+ private static final List> SENT_MESSAGES_ATTRIBUTES =
+ buildAttributes(false, false, true);
+ private static final List> CONSUMED_MESSAGES_ATTRIBUTES =
+ buildAttributes(true, false, true);
+ private static final List> PROCESS_DURATION_ATTRIBUTES =
+ buildAttributes(true, false, true);
+
+ private static List> buildAttributes(
+ boolean includeConsumerAttributes, boolean includeOperationType, boolean includeErrorType) {
+ List> attributes = new ArrayList<>();
+ attributes.add(MESSAGING_OPERATION_NAME);
+ attributes.add(MESSAGING_SYSTEM);
+ if (includeErrorType) {
+ attributes.add(ERROR_TYPE);
+ }
+ if (includeConsumerAttributes) {
+ attributes.add(MESSAGING_CONSUMER_GROUP_NAME);
+ }
+ attributes.add(MESSAGING_DESTINATION_NAME);
+ if (includeConsumerAttributes) {
+ attributes.add(MESSAGING_DESTINATION_SUBSCRIPTION_NAME);
+ }
+ attributes.add(MESSAGING_DESTINATION_TEMPLATE);
+ if (includeOperationType) {
+ attributes.add(MESSAGING_OPERATION_TYPE);
+ }
+ attributes.add(MESSAGING_DESTINATION_PARTITION_ID);
+ return unmodifiableList(attributes);
+ }
+
+ static Attributes filterAttributes(Attributes attributes) {
+ if (attributes.get(MESSAGING_DESTINATION_TEMPLATE) == null
+ && !Boolean.TRUE.equals(attributes.get(MESSAGING_DESTINATION_ANONYMOUS))
+ && !Boolean.TRUE.equals(attributes.get(MESSAGING_DESTINATION_TEMPORARY))) {
+ return attributes;
+ }
+ return attributes.toBuilder().remove(MESSAGING_DESTINATION_NAME).build();
+ }
+
+ static void applyOldDurationAdvice(DoubleHistogramBuilder builder) {
if (!(builder instanceof ExtendedDoubleHistogramBuilder)) {
return;
}
- ((ExtendedDoubleHistogramBuilder) builder).setAttributesAdvice(MESSAGING_ATTRIBUTES);
+ ((ExtendedDoubleHistogramBuilder) builder).setAttributesAdvice(OLD_ATTRIBUTES);
}
- static void applyReceiveDurationAdvice(DoubleHistogramBuilder builder) {
+ static void applyClientOperationDurationAdvice(DoubleHistogramBuilder builder) {
if (!(builder instanceof ExtendedDoubleHistogramBuilder)) {
return;
}
- ((ExtendedDoubleHistogramBuilder) builder).setAttributesAdvice(MESSAGING_ATTRIBUTES);
+ ((ExtendedDoubleHistogramBuilder) builder)
+ .setAttributesAdvice(CLIENT_OPERATION_DURATION_ATTRIBUTES);
+ }
+
+ static void applyProcessDurationAdvice(DoubleHistogramBuilder builder) {
+ if (!(builder instanceof ExtendedDoubleHistogramBuilder)) {
+ return;
+ }
+ ((ExtendedDoubleHistogramBuilder) builder).setAttributesAdvice(PROCESS_DURATION_ATTRIBUTES);
+ }
+
+ static void applyOldMessagesAdvice(LongCounterBuilder builder) {
+ if (!(builder instanceof ExtendedLongCounterBuilder)) {
+ return;
+ }
+ ((ExtendedLongCounterBuilder) builder).setAttributesAdvice(OLD_ATTRIBUTES);
+ }
+
+ static void applySentMessagesAdvice(LongCounterBuilder builder) {
+ if (!(builder instanceof ExtendedLongCounterBuilder)) {
+ return;
+ }
+ ((ExtendedLongCounterBuilder) builder).setAttributesAdvice(SENT_MESSAGES_ATTRIBUTES);
}
- static void applyReceiveMessagesAdvice(LongCounterBuilder builder) {
+ static void applyConsumedMessagesAdvice(LongCounterBuilder builder) {
if (!(builder instanceof ExtendedLongCounterBuilder)) {
return;
}
- ((ExtendedLongCounterBuilder) builder).setAttributesAdvice(MESSAGING_ATTRIBUTES);
+ ((ExtendedLongCounterBuilder) builder).setAttributesAdvice(CONSUMED_MESSAGES_ATTRIBUTES);
}
private MessagingMetricsAdvice() {}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingOperationType.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingOperationType.java
new file mode 100644
index 000000000000..53b90807068b
--- /dev/null
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingOperationType.java
@@ -0,0 +1,35 @@
+/*
+ * Copyright The OpenTelemetry Authors
+ * SPDX-License-Identifier: Apache-2.0
+ */
+
+package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
+
+/**
+ * Represents a messaging
+ * operation type.
+ */
+public enum MessagingOperationType {
+ CREATE("create", "create"),
+ SEND("send", "publish"),
+ RECEIVE("receive", "receive"),
+ PROCESS("process", "process"),
+ SETTLE("settle", "settle");
+
+ private final String value;
+ private final String defaultOperationName;
+
+ MessagingOperationType(String value, String defaultOperationName) {
+ this.value = value;
+ this.defaultOperationName = defaultOperationName;
+ }
+
+ String value() {
+ return value;
+ }
+
+ String defaultOperationName() {
+ return defaultOperationName;
+ }
+}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingProcessMetrics.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingProcessMetrics.java
new file mode 100644
index 000000000000..287741bea625
--- /dev/null
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingProcessMetrics.java
@@ -0,0 +1,97 @@
+/*
+ * Copyright The OpenTelemetry Authors
+ * SPDX-License-Identifier: Apache-2.0
+ */
+
+package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
+
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
+import static java.util.concurrent.TimeUnit.SECONDS;
+import static java.util.logging.Level.FINE;
+
+import com.google.auto.value.AutoValue;
+import com.google.errorprone.annotations.CanIgnoreReturnValue;
+import io.opentelemetry.api.common.Attributes;
+import io.opentelemetry.api.metrics.DoubleHistogram;
+import io.opentelemetry.api.metrics.DoubleHistogramBuilder;
+import io.opentelemetry.api.metrics.Meter;
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.ContextKey;
+import io.opentelemetry.instrumentation.api.instrumenter.OperationListener;
+import io.opentelemetry.instrumentation.api.instrumenter.OperationMetrics;
+import io.opentelemetry.instrumentation.api.internal.OperationMetricsUtil;
+import java.util.logging.Logger;
+import javax.annotation.Nullable;
+
+/**
+ * {@link OperationListener} which keeps track of message
+ * processing metrics.
+ */
+public final class MessagingProcessMetrics implements OperationListener {
+ private static final double NANOS_PER_S = SECONDS.toNanos(1);
+
+ private static final ContextKey MESSAGING_PROCESS_METRICS_STATE =
+ ContextKey.named("messaging-process-metrics-state");
+ private static final Logger logger = Logger.getLogger(MessagingProcessMetrics.class.getName());
+
+ @Nullable private final DoubleHistogram processDurationHistogram;
+
+ private MessagingProcessMetrics(Meter meter) {
+ processDurationHistogram = emitStableMessagingSemconv() ? buildProcessDuration(meter) : null;
+ }
+
+ public static OperationMetrics get() {
+ return OperationMetricsUtil.create("messaging process", MessagingProcessMetrics::new);
+ }
+
+ @Override
+ @CanIgnoreReturnValue
+ public Context onStart(Context context, Attributes startAttributes, long startNanos) {
+ if (processDurationHistogram == null) {
+ return context;
+ }
+ return context.with(
+ MESSAGING_PROCESS_METRICS_STATE,
+ new AutoValue_MessagingProcessMetrics_State(startAttributes, startNanos));
+ }
+
+ @Override
+ public void onEnd(Context context, Attributes endAttributes, long endNanos) {
+ if (processDurationHistogram == null) {
+ return;
+ }
+ State state = context.get(MESSAGING_PROCESS_METRICS_STATE);
+ if (state == null) {
+ logger.log(
+ FINE,
+ "No state present when ending context {0}. Cannot record messaging process metrics.",
+ context);
+ return;
+ }
+
+ Attributes attributes = state.startAttributes().toBuilder().putAll(endAttributes).build();
+ processDurationHistogram.record(
+ (endNanos - state.startTimeNanos()) / NANOS_PER_S,
+ MessagingMetricsAdvice.filterAttributes(attributes),
+ context);
+ }
+
+ private static DoubleHistogram buildProcessDuration(Meter meter) {
+ DoubleHistogramBuilder builder =
+ meter
+ .histogramBuilder("messaging.process.duration")
+ .setDescription("Duration of processing operation.")
+ .setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
+ .setUnit("s");
+ MessagingMetricsAdvice.applyProcessDurationAdvice(builder);
+ return builder.build();
+ }
+
+ @AutoValue
+ abstract static class State {
+ abstract Attributes startAttributes();
+
+ abstract long startTimeNanos();
+ }
+}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingProducerMetrics.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingProducerMetrics.java
index ee79fd0c224e..cd022e15a96a 100644
--- a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingProducerMetrics.java
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingProducerMetrics.java
@@ -5,14 +5,19 @@
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitOldMessagingSemconv;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
import static java.util.concurrent.TimeUnit.SECONDS;
import static java.util.logging.Level.FINE;
import com.google.auto.value.AutoValue;
import com.google.errorprone.annotations.CanIgnoreReturnValue;
+import io.opentelemetry.api.common.AttributeKey;
import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.metrics.DoubleHistogram;
import io.opentelemetry.api.metrics.DoubleHistogramBuilder;
+import io.opentelemetry.api.metrics.LongCounter;
+import io.opentelemetry.api.metrics.LongCounterBuilder;
import io.opentelemetry.api.metrics.Meter;
import io.opentelemetry.context.Context;
import io.opentelemetry.context.ContextKey;
@@ -20,34 +25,52 @@
import io.opentelemetry.instrumentation.api.instrumenter.OperationMetrics;
import io.opentelemetry.instrumentation.api.internal.OperationMetricsUtil;
import java.util.logging.Logger;
+import javax.annotation.Nullable;
/**
* {@link OperationListener} which keeps track of Producer
+ * href="https://github.com/open-telemetry/semantic-conventions/blob/v1.43.0/docs/messaging/messaging-metrics.md#producer-metrics">producer
* metrics.
*/
public final class MessagingProducerMetrics implements OperationListener {
private static final double NANOS_PER_S = SECONDS.toNanos(1);
+ // copied from MessagingIncubatingAttributes
+ private static final AttributeKey MESSAGING_BATCH_MESSAGE_COUNT =
+ AttributeKey.longKey("messaging.batch.message_count");
+ private static final AttributeKey MESSAGING_OPERATION =
+ AttributeKey.stringKey("messaging.operation");
+ private static final AttributeKey MESSAGING_OPERATION_TYPE =
+ AttributeKey.stringKey("messaging.operation.type");
private static final ContextKey MESSAGING_PRODUCER_METRICS_STATE =
ContextKey.named("messaging-producer-metrics-state");
private static final Logger logger = Logger.getLogger(MessagingProducerMetrics.class.getName());
- private final DoubleHistogram publishDurationHistogram;
+ private final boolean supportsStableSemconv;
+ @Nullable private final DoubleHistogram publishDurationHistogram;
+ @Nullable private final DoubleHistogram clientOperationDurationHistogram;
+ @Nullable private final LongCounter sentMessagesCounter;
- private MessagingProducerMetrics(Meter meter) {
- DoubleHistogramBuilder durationBuilder =
- meter
- .histogramBuilder("messaging.publish.duration")
- .setDescription("Measures the duration of publish operation.")
- .setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
- .setUnit("s");
- MessagingMetricsAdvice.applyPublishDurationAdvice(durationBuilder);
- publishDurationHistogram = durationBuilder.build();
+ private MessagingProducerMetrics(Meter meter, boolean supportsStableSemconv) {
+ this.supportsStableSemconv = supportsStableSemconv;
+ boolean emitOldSemconv = !supportsStableSemconv || emitOldMessagingSemconv();
+ boolean emitStableSemconv = supportsStableSemconv && emitStableMessagingSemconv();
+ publishDurationHistogram = emitOldSemconv ? buildPublishDuration(meter) : null;
+ clientOperationDurationHistogram =
+ emitStableSemconv ? buildClientOperationDuration(meter) : null;
+ sentMessagesCounter = emitStableSemconv ? buildSentMessages(meter) : null;
}
+ /** Returns metrics for extractors configured with {@link MessageOperation}. */
public static OperationMetrics get() {
- return OperationMetricsUtil.create("messaging producer", MessagingProducerMetrics::new);
+ return OperationMetricsUtil.create(
+ "messaging producer", meter -> new MessagingProducerMetrics(meter, false));
+ }
+
+ /** Returns metrics for extractors configured with {@link MessagingOperationType}. */
+ public static OperationMetrics getForOperationType() {
+ return OperationMetricsUtil.create(
+ "messaging producer", meter -> new MessagingProducerMetrics(meter, true));
}
@Override
@@ -70,9 +93,64 @@ public void onEnd(Context context, Attributes endAttributes, long endNanos) {
}
Attributes attributes = state.startAttributes().toBuilder().putAll(endAttributes).build();
+ double duration = (endNanos - state.startTimeNanos()) / NANOS_PER_S;
- publishDurationHistogram.record(
- (endNanos - state.startTimeNanos()) / NANOS_PER_S, attributes, context);
+ if (publishDurationHistogram != null
+ && (!supportsStableSemconv
+ || "publish".equals(attributes.get(MESSAGING_OPERATION))
+ || MessagingOperationType.SEND
+ .value()
+ .equals(attributes.get(MESSAGING_OPERATION_TYPE)))) {
+ publishDurationHistogram.record(duration, attributes, context);
+ }
+ Attributes filteredAttributes =
+ clientOperationDurationHistogram != null || sentMessagesCounter != null
+ ? MessagingMetricsAdvice.filterAttributes(attributes)
+ : attributes;
+ if (clientOperationDurationHistogram != null) {
+ clientOperationDurationHistogram.record(duration, filteredAttributes, context);
+ }
+ if (sentMessagesCounter != null
+ && MessagingOperationType.SEND.value().equals(attributes.get(MESSAGING_OPERATION_TYPE))) {
+ Long batchMessageCount = attributes.get(MESSAGING_BATCH_MESSAGE_COUNT);
+ long sentMessagesCount = batchMessageCount == null ? 1 : batchMessageCount;
+ if (sentMessagesCount > 0) {
+ sentMessagesCounter.add(sentMessagesCount, filteredAttributes, context);
+ }
+ }
+ }
+
+ private static DoubleHistogram buildPublishDuration(Meter meter) {
+ DoubleHistogramBuilder builder =
+ meter
+ .histogramBuilder("messaging.publish.duration")
+ .setDescription("Measures the duration of publish operation.")
+ .setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
+ .setUnit("s");
+ MessagingMetricsAdvice.applyOldDurationAdvice(builder);
+ return builder.build();
+ }
+
+ private static DoubleHistogram buildClientOperationDuration(Meter meter) {
+ DoubleHistogramBuilder builder =
+ meter
+ .histogramBuilder("messaging.client.operation.duration")
+ .setDescription(
+ "Duration of messaging operation initiated by a producer or consumer client.")
+ .setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
+ .setUnit("s");
+ MessagingMetricsAdvice.applyClientOperationDurationAdvice(builder);
+ return builder.build();
+ }
+
+ private static LongCounter buildSentMessages(Meter meter) {
+ LongCounterBuilder builder =
+ meter
+ .counterBuilder("messaging.client.sent.messages")
+ .setDescription("Number of messages producer attempted to send to the broker.")
+ .setUnit("{message}");
+ MessagingMetricsAdvice.applySentMessagesAdvice(builder);
+ return builder.build();
}
@AutoValue
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractor.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractor.java
new file mode 100644
index 000000000000..4a7008c1420a
--- /dev/null
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractor.java
@@ -0,0 +1,80 @@
+/*
+ * Copyright The OpenTelemetry Authors
+ * SPDX-License-Identifier: Apache-2.0
+ */
+
+package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
+
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
+
+import io.opentelemetry.api.trace.SpanKind;
+import io.opentelemetry.instrumentation.api.instrumenter.SpanKindExtractor;
+
+/** Selects messaging span kinds according to the configured semantic convention version. */
+public final class MessagingSpanKindExtractor {
+
+ /**
+ * Returns a span kind extractor following the v1.43
+ * messaging span kind conventions.
+ */
+ public static SpanKindExtractor create(MessagingOperationType operationType) {
+ return create(operationType, true);
+ }
+
+ /**
+ * Returns a span kind extractor following the v1.43
+ * messaging span kind conventions.
+ *
+ * @param isSpanContextPropagated whether the context of a {@link MessagingOperationType#SEND}
+ * span is propagated as the message creation context; ignored for other operation types
+ */
+ public static SpanKindExtractor create(
+ MessagingOperationType operationType, boolean isSpanContextPropagated) {
+ SpanKind spanKind;
+ switch (operationType) {
+ case CREATE:
+ spanKind = SpanKind.PRODUCER;
+ break;
+ case SEND:
+ spanKind =
+ emitStableMessagingSemconv() && !isSpanContextPropagated
+ ? SpanKind.CLIENT
+ : SpanKind.PRODUCER;
+ break;
+ case RECEIVE:
+ spanKind = emitStableMessagingSemconv() ? SpanKind.CLIENT : SpanKind.CONSUMER;
+ break;
+ case PROCESS:
+ spanKind = SpanKind.CONSUMER;
+ break;
+ case SETTLE:
+ spanKind = SpanKind.CLIENT;
+ break;
+ default:
+ throw new IllegalStateException("Can't possibly happen");
+ }
+ SpanKind result = spanKind;
+ return request -> result;
+ }
+
+ /** Returns a span kind extractor for the given operation. */
+ public static SpanKindExtractor create(MessageOperation operation) {
+ SpanKind spanKind;
+ switch (operation) {
+ case PUBLISH:
+ spanKind = SpanKind.PRODUCER;
+ break;
+ case RECEIVE:
+ case PROCESS:
+ spanKind = SpanKind.CONSUMER;
+ break;
+ default:
+ throw new IllegalStateException("Can't possibly happen");
+ }
+ return request -> spanKind;
+ }
+
+ private MessagingSpanKindExtractor() {}
+}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractor.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractor.java
index 624d3b8e00bf..a6d6aa679747 100644
--- a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractor.java
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractor.java
@@ -5,35 +5,75 @@
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
+
import io.opentelemetry.instrumentation.api.instrumenter.SpanNameExtractor;
public final class MessagingSpanNameExtractor implements SpanNameExtractor {
/**
* Returns a {@link SpanNameExtractor} that constructs the span name according to
- * messaging semantic conventions: {@code }.
+ * href="https://github.com/open-telemetry/semantic-conventions/blob/v1.43.0/docs/messaging/messaging-spans.md#span-name">
+ * messaging semantic conventions.
*
* @see MessagingAttributesGetter#getDestination(Object) used to extract {@code }.
- * @see MessageOperation used to extract {@code }.
+ * @see MessagingOperationType used to extract {@code }.
*/
+ public static SpanNameExtractor create(
+ MessagingAttributesGetter getter, MessagingOperationType operationType) {
+ return builder(getter, operationType).build();
+ }
+
+ /** Returns a messaging span name extractor for the given operation. */
public static SpanNameExtractor create(
MessagingAttributesGetter getter, MessageOperation operation) {
- return new MessagingSpanNameExtractor<>(getter, operation);
+ return builder(getter, operation).build();
}
- private final MessagingAttributesGetter getter;
- private final MessageOperation operation;
+ /**
+ * Returns a new {@link MessagingSpanNameExtractorBuilder} that can be used to configure the
+ * messaging span name extractor.
+ */
+ public static MessagingSpanNameExtractorBuilder builder(
+ MessagingAttributesGetter getter, MessagingOperationType operationType) {
+ return new MessagingSpanNameExtractorBuilder<>(getter, operationType, true);
+ }
- private MessagingSpanNameExtractor(
+ /** Returns a messaging span name extractor builder for the given operation. */
+ public static MessagingSpanNameExtractorBuilder builder(
MessagingAttributesGetter getter, MessageOperation operation) {
+ return new MessagingSpanNameExtractorBuilder<>(getter, operation.type(), false);
+ }
+
+ private final MessagingAttributesGetter getter;
+ private final MessagingOperationType operationType;
+ private final String operationName;
+ private final boolean supportsStableSemconv;
+
+ MessagingSpanNameExtractor(
+ MessagingAttributesGetter getter,
+ MessagingOperationType operationType,
+ String operationName,
+ boolean supportsStableSemconv) {
this.getter = getter;
- this.operation = operation;
+ this.operationType = operationType;
+ this.operationName = operationName;
+ this.supportsStableSemconv = supportsStableSemconv;
}
@Override
public String extract(REQUEST request) {
+ if (supportsStableSemconv && emitStableMessagingSemconv()) {
+ String destinationName = getter.getDestinationTemplate(request);
+ if (destinationName == null
+ && !getter.isTemporaryDestination(request)
+ && !getter.isAnonymousDestination(request)) {
+ destinationName = getter.getDestination(request);
+ }
+ return destinationName == null ? operationName : operationName + " " + destinationName;
+ }
+
String destinationName =
getter.isTemporaryDestination(request)
? MessagingAttributesExtractor.TEMP_DESTINATION_NAME
@@ -42,6 +82,6 @@ public String extract(REQUEST request) {
destinationName = "unknown";
}
- return destinationName + " " + operation.operationName();
+ return destinationName + " " + operationType.defaultOperationName();
}
}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorBuilder.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorBuilder.java
new file mode 100644
index 000000000000..665ceaa58888
--- /dev/null
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorBuilder.java
@@ -0,0 +1,49 @@
+/*
+ * Copyright The OpenTelemetry Authors
+ * SPDX-License-Identifier: Apache-2.0
+ */
+
+package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
+
+import static java.util.Objects.requireNonNull;
+
+import com.google.errorprone.annotations.CanIgnoreReturnValue;
+import io.opentelemetry.instrumentation.api.instrumenter.SpanNameExtractor;
+
+/** A builder of {@link MessagingSpanNameExtractor}. */
+public final class MessagingSpanNameExtractorBuilder {
+
+ private final MessagingAttributesGetter getter;
+ private final MessagingOperationType operationType;
+ private final boolean supportsStableSemconv;
+ private String operationName;
+
+ MessagingSpanNameExtractorBuilder(
+ MessagingAttributesGetter getter,
+ MessagingOperationType operationType,
+ boolean supportsStableSemconv) {
+ this.getter = getter;
+ this.operationType = requireNonNull(operationType, "operationType");
+ this.supportsStableSemconv = supportsStableSemconv;
+ this.operationName = operationType.defaultOperationName();
+ }
+
+ /** Configures the system-specific operation name used in the v1.43 messaging span name. */
+ @CanIgnoreReturnValue
+ public MessagingSpanNameExtractorBuilder setOperationName(String operationName) {
+ if (!supportsStableSemconv) {
+ throw new IllegalStateException("Operation name is not configurable for legacy builders");
+ }
+ this.operationName = requireNonNull(operationName, "operationName");
+ return this;
+ }
+
+ /**
+ * Returns a new {@link MessagingSpanNameExtractor} with the settings of this {@link
+ * MessagingSpanNameExtractorBuilder}.
+ */
+ public SpanNameExtractor build() {
+ return new MessagingSpanNameExtractor<>(
+ getter, operationType, operationName, supportsStableSemconv);
+ }
+}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/internal/MessagingProcessContextCustomizer.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/internal/MessagingProcessContextCustomizer.java
new file mode 100644
index 000000000000..03b45ddf5c86
--- /dev/null
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/internal/MessagingProcessContextCustomizer.java
@@ -0,0 +1,48 @@
+/*
+ * Copyright The OpenTelemetry Authors
+ * SPDX-License-Identifier: Apache-2.0
+ */
+
+package io.opentelemetry.instrumentation.api.incubator.semconv.messaging.internal;
+
+import io.opentelemetry.api.common.Attributes;
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.context.Context;
+import io.opentelemetry.context.propagation.TextMapGetter;
+import io.opentelemetry.context.propagation.TextMapPropagator;
+import io.opentelemetry.instrumentation.api.instrumenter.ContextCustomizer;
+import java.util.function.BiFunction;
+
+/**
+ * This class is internal and is hence not for public use. Its APIs are unstable and can change at
+ * any time.
+ */
+public class MessagingProcessContextCustomizer implements ContextCustomizer {
+
+ public static ContextCustomizer create(
+ TextMapPropagator propagator, TextMapGetter getter) {
+ return create((parentContext, request) -> propagator.extract(parentContext, request, getter));
+ }
+
+ public static ContextCustomizer create(
+ BiFunction producerContextExtractor) {
+ return new MessagingProcessContextCustomizer<>(producerContextExtractor);
+ }
+
+ private final BiFunction producerContextExtractor;
+
+ private MessagingProcessContextCustomizer(
+ BiFunction producerContextExtractor) {
+ this.producerContextExtractor = producerContextExtractor;
+ }
+
+ @Override
+ public Context onStart(Context parentContext, REQUEST request, Attributes startAttributes) {
+ Context extractedContext = producerContextExtractor.apply(parentContext, request);
+ Span parentSpan = Span.fromContext(parentContext);
+ if (!parentSpan.getSpanContext().isValid()) {
+ return extractedContext;
+ }
+ return extractedContext.with(parentSpan);
+ }
+}
diff --git a/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/internal/MessagingProcessInstrumenterFactory.java b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/internal/MessagingProcessInstrumenterFactory.java
new file mode 100644
index 000000000000..726e5ae92194
--- /dev/null
+++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/internal/MessagingProcessInstrumenterFactory.java
@@ -0,0 +1,55 @@
+/*
+ * Copyright The OpenTelemetry Authors
+ * SPDX-License-Identifier: Apache-2.0
+ */
+
+package io.opentelemetry.instrumentation.api.incubator.semconv.messaging.internal;
+
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
+
+import io.opentelemetry.api.trace.Span;
+import io.opentelemetry.api.trace.SpanContext;
+import io.opentelemetry.context.propagation.TextMapGetter;
+import io.opentelemetry.context.propagation.TextMapPropagator;
+import io.opentelemetry.instrumentation.api.instrumenter.Instrumenter;
+import io.opentelemetry.instrumentation.api.instrumenter.InstrumenterBuilder;
+import io.opentelemetry.instrumentation.api.instrumenter.SpanKindExtractor;
+import io.opentelemetry.instrumentation.api.internal.PropagatorBasedSpanLinksExtractor;
+
+/**
+ * This class is internal and is hence not for public use. Its APIs are unstable and can change at
+ * any time.
+ */
+public class MessagingProcessInstrumenterFactory {
+
+ public static Instrumenter create(
+ InstrumenterBuilder builder,
+ TextMapPropagator propagator,
+ TextMapGetter getter,
+ boolean receiveInstrumentationEnabled) {
+ if (emitStableMessagingSemconv()) {
+ builder.addSpanLinksExtractor(
+ (spanLinks, parentContext, request) -> {
+ SpanContext parentSpanContext = Span.fromContext(parentContext).getSpanContext();
+ SpanContext producerSpanContext =
+ Span.fromContext(propagator.extract(parentContext, request, getter))
+ .getSpanContext();
+ if (parentSpanContext.isValid()
+ && producerSpanContext.isValid()
+ && (!producerSpanContext.getTraceId().equals(parentSpanContext.getTraceId())
+ || !producerSpanContext.getSpanId().equals(parentSpanContext.getSpanId()))) {
+ spanLinks.addLink(producerSpanContext);
+ }
+ });
+ builder.addContextCustomizer(MessagingProcessContextCustomizer.create(propagator, getter));
+ return builder.buildInstrumenter(SpanKindExtractor.alwaysConsumer());
+ }
+ if (receiveInstrumentationEnabled) {
+ builder.addSpanLinksExtractor(new PropagatorBasedSpanLinksExtractor<>(propagator, getter));
+ return builder.buildInstrumenter(SpanKindExtractor.alwaysConsumer());
+ }
+ return builder.buildConsumerInstrumenter(getter);
+ }
+
+ private MessagingProcessInstrumenterFactory() {}
+}
diff --git a/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorTest.java b/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorTest.java
index a8d56f7e4f98..2a7959df2b08 100644
--- a/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorTest.java
+++ b/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorTest.java
@@ -6,7 +6,10 @@
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
import static io.opentelemetry.api.common.AttributeKey.stringKey;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitOldMessagingSemconv;
+import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
import static io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions.assertThat;
+import static io.opentelemetry.semconv.ErrorAttributes.ERROR_TYPE;
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_BATCH_MESSAGE_COUNT;
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_DESTINATION_ANONYMOUS;
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_DESTINATION_NAME;
@@ -17,15 +20,21 @@
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_MESSAGE_ENVELOPE_SIZE;
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_MESSAGE_ID;
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_OPERATION;
+import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_OPERATION_NAME;
+import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_OPERATION_TYPE;
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_SYSTEM;
import static java.util.Collections.emptyMap;
+import static java.util.Collections.singletonMap;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.assertj.core.api.Assertions.entry;
+import static org.junit.jupiter.params.provider.Arguments.argumentSet;
import io.opentelemetry.api.common.AttributeKey;
import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.common.AttributesBuilder;
import io.opentelemetry.context.Context;
import io.opentelemetry.instrumentation.api.instrumenter.AttributesExtractor;
+import io.opentelemetry.instrumentation.api.internal.SpanKey;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
@@ -47,8 +56,8 @@ void shouldExtractAllAvailableAttributes(
boolean temporary,
boolean anonymous,
String destination,
- MessageOperation operation,
- String expectedDestination) {
+ MessagingOperationType operationType,
+ String operationName) {
// given
Map request = new HashMap<>();
request.put("system", "myQueue");
@@ -63,13 +72,16 @@ void shouldExtractAllAvailableAttributes(
}
request.put("url", "http://broker/topic");
request.put("conversationId", "42");
+ request.put("messageId", "42");
request.put("bodySize", "100");
request.put("envelopeSize", "120");
request.put("clientId", "43");
request.put("batchMessageCount", "2");
AttributesExtractor