From 653c0b1d622f3fb3b9be44a0c205d3d37dcd2214 Mon Sep 17 00:00:00 2001 From: Trask Stalnaker Date: Mon, 20 Jul 2026 19:43:46 -0700 Subject: [PATCH] Add messaging v1.43 shared contract --- .../build.gradle.kts | 9 +- .../semconv/messaging/MessageOperation.java | 29 ++- .../MessagingAttributesExtractor.java | 99 ++++++++-- .../MessagingAttributesExtractorBuilder.java | 28 ++- .../messaging/MessagingAttributesGetter.java | 15 ++ .../messaging/MessagingOperationType.java | 35 ++++ .../messaging/MessagingSpanKindExtractor.java | 80 ++++++++ .../messaging/MessagingSpanNameExtractor.java | 58 +++++- .../MessagingSpanNameExtractorBuilder.java | 46 +++++ .../MessagingAttributesExtractorTest.java | 180 ++++++++++++++++-- .../MessagingSpanKindExtractorTest.java | 74 +++++++ .../MessagingSpanNameExtractorTest.java | 109 +++++++++-- .../api/internal/SemconvStability.java | 7 +- .../instrumentation/api/internal/SpanKey.java | 6 + .../api/internal/SemconvStabilityTest.java | 66 +++++++ .../v1_14/SpanKeyBridging.java | 18 ++ .../v1_14/ContextBridgeTest.java | 4 +- .../testing/AgentSpanTestingInstrumenter.java | 2 + .../SemconvMessagingStabilityUtil.java | 36 ++++ 19 files changed, 819 insertions(+), 82 deletions(-) create mode 100644 instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingOperationType.java create mode 100644 instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractor.java create mode 100644 instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorBuilder.java create mode 100644 instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractorTest.java create mode 100644 testing-common/src/main/java/io/opentelemetry/instrumentation/testing/junit/message/SemconvMessagingStabilityUtil.java 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..b9e01e99a8ee 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,40 @@ 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) { + this.operationName = requireNonNull(operationName, "operationName"); + return this; } /** @@ -47,6 +63,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/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/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..afb61067db5b --- /dev/null +++ b/instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorBuilder.java @@ -0,0 +1,46 @@ +/* + * 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) { + 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/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..dbadbf4855fc 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, String> underTest = - MessagingAttributesExtractor.create(TestGetter.INSTANCE, operation); + MessagingAttributesExtractor.builderForOperationType(TestGetter.INSTANCE, operationType) + .setOperationName(operationName) + .build(); Context context = Context.root(); @@ -78,16 +90,22 @@ void shouldExtractAllAvailableAttributes( underTest.onStart(startAttributes, context, request); AttributesBuilder endAttributes = Attributes.builder(); - underTest.onEnd(endAttributes, context, request, "42", null); + underTest.onEnd(endAttributes, context, request, "42", new IllegalStateException()); // then List, Object>> expectedEntries = new ArrayList<>(); expectedEntries.add(entry(MESSAGING_SYSTEM, "myQueue")); - expectedEntries.add(entry(MESSAGING_DESTINATION_NAME, expectedDestination)); if (temporary) { expectedEntries.add(entry(MESSAGING_DESTINATION_TEMPORARY, true)); + if (emitStableMessagingSemconv()) { + expectedEntries.add(entry(MESSAGING_DESTINATION_NAME, destination)); + expectedEntries.add(entry(MESSAGING_DESTINATION_TEMPLATE, destination)); + } else { + expectedEntries.add(entry(MESSAGING_DESTINATION_NAME, "(temporary)")); + } } else { - expectedEntries.add(entry(MESSAGING_DESTINATION_TEMPLATE, expectedDestination)); + expectedEntries.add(entry(MESSAGING_DESTINATION_NAME, destination)); + expectedEntries.add(entry(MESSAGING_DESTINATION_TEMPLATE, destination)); } if (anonymous) { expectedEntries.add(entry(MESSAGING_DESTINATION_ANONYMOUS, true)); @@ -95,22 +113,143 @@ void shouldExtractAllAvailableAttributes( expectedEntries.add(entry(MESSAGING_MESSAGE_CONVERSATION_ID, "42")); expectedEntries.add(entry(MESSAGING_MESSAGE_BODY_SIZE, 100L)); expectedEntries.add(entry(MESSAGING_MESSAGE_ENVELOPE_SIZE, 120L)); - expectedEntries.add(entry(stringKey("messaging.client_id"), "43")); - expectedEntries.add(entry(MESSAGING_OPERATION, operation.operationName())); + if (emitOldMessagingSemconv()) { + expectedEntries.add(entry(stringKey("messaging.client_id"), "43")); + expectedEntries.add(entry(MESSAGING_OPERATION, operationType.defaultOperationName())); + } + if (emitStableMessagingSemconv()) { + expectedEntries.add(entry(stringKey("messaging.client.id"), "43")); + expectedEntries.add(entry(MESSAGING_OPERATION_NAME, operationName)); + expectedEntries.add(entry(MESSAGING_OPERATION_TYPE, operationType.value())); + } @SuppressWarnings({"unchecked", "rawtypes"}) MapEntry, ?>[] expectedEntriesArr = expectedEntries.toArray(new MapEntry[0]); assertThat(startAttributes.build()).containsOnly(expectedEntriesArr); - assertThat(endAttributes.build()) - .containsOnly(entry(MESSAGING_MESSAGE_ID, "42"), entry(MESSAGING_BATCH_MESSAGE_COUNT, 2L)); + if (emitStableMessagingSemconv()) { + assertThat(endAttributes.build()) + .containsOnly( + entry(MESSAGING_MESSAGE_ID, "42"), + entry(MESSAGING_BATCH_MESSAGE_COUNT, 2L), + entry(ERROR_TYPE, IllegalStateException.class.getName())); + } else { + assertThat(endAttributes.build()) + .containsOnly( + entry(MESSAGING_MESSAGE_ID, "42"), entry(MESSAGING_BATCH_MESSAGE_COUNT, 2L)); + } } static Stream destinations() { return Stream.of( - Arguments.of(false, false, "destination", MessageOperation.RECEIVE, "destination"), - Arguments.of(true, true, null, MessageOperation.PROCESS, "(temporary)")); + argumentSet( + "create operation", + false, + false, + "destination", + MessagingOperationType.CREATE, + "create"), + argumentSet( + "regular destination", + false, + false, + "destination", + MessagingOperationType.RECEIVE, + "poll"), + argumentSet( + "temporary anonymous destination", + true, + true, + "generated-destination", + MessagingOperationType.PROCESS, + "process"), + argumentSet( + "settle operation", + false, + false, + "destination", + MessagingOperationType.SETTLE, + "settle")); + } + + @ParameterizedTest + @MethodSource("spanKeys") + void shouldReturnSpanKey(MessagingOperationType operationType, SpanKey spanKey) { + MessagingAttributesExtractor, String> underTest = + new MessagingAttributesExtractor<>( + TestGetter.INSTANCE, + operationType, + operationType.defaultOperationName(), + true, + new ArrayList<>()); + + assertThat(underTest.internalGetSpanKey()).isSameAs(spanKey); + } + + static Stream spanKeys() { + return Stream.of( + argumentSet("create", MessagingOperationType.CREATE, SpanKey.PRODUCER_CREATE), + argumentSet("send", MessagingOperationType.SEND, SpanKey.PRODUCER), + argumentSet("receive", MessagingOperationType.RECEIVE, SpanKey.CONSUMER_RECEIVE), + argumentSet("process", MessagingOperationType.PROCESS, SpanKey.CONSUMER_PROCESS), + argumentSet("settle", MessagingOperationType.SETTLE, SpanKey.CONSUMER_SETTLE)); + } + + @SuppressWarnings("deprecation") // testing deprecated API + @Test + void shouldSupportDeprecatedMessageOperation() { + AttributesExtractor, String> underTest = + MessagingAttributesExtractor.create(TestGetter.INSTANCE, MessageOperation.PUBLISH); + + AttributesBuilder attributes = Attributes.builder(); + underTest.onStart(attributes, Context.root(), singletonMap("anonymousDestination", "y")); + + assertThat(attributes.build()) + .containsOnly( + entry(MESSAGING_DESTINATION_ANONYMOUS, true), entry(MESSAGING_OPERATION, "publish")); + } + + @Test + void shouldExtractOperationNameWithoutOperationType() { + MessagingOperationType operationType = null; + AttributesExtractor, String> underTest = + MessagingAttributesExtractor.builderForOperationType(TestGetter.INSTANCE, operationType) + .setOperationName("ack") + .build(); + + AttributesBuilder attributes = Attributes.builder(); + underTest.onStart(attributes, Context.root(), emptyMap()); + + Attributes expected = + emitStableMessagingSemconv() + ? Attributes.of(MESSAGING_OPERATION_NAME, "ack") + : Attributes.empty(); + assertThat(attributes.build()).isEqualTo(expected); + } + + @Test + void shouldRequireOperationNameForStableSemconv() { + assertThatThrownBy( + () -> + MessagingAttributesExtractor.builderForOperationType(TestGetter.INSTANCE, null) + .build()) + .isInstanceOf(NullPointerException.class) + .hasMessage("operationName"); + } + + @Test + void shouldExtractErrorTypeFromResponse() { + AttributesExtractor, String> underTest = + MessagingAttributesExtractor.createForOperationType( + TestGetter.INSTANCE, MessagingOperationType.RECEIVE); + + AttributesBuilder attributes = Attributes.builder(); + underTest.onEnd(attributes, Context.root(), emptyMap(), "failure", null); + + Attributes expected = + emitStableMessagingSemconv() ? Attributes.of(ERROR_TYPE, "failure") : Attributes.empty(); + assertThat(attributes.build()).isEqualTo(expected); } @Test @@ -118,20 +257,24 @@ void shouldExtractNoAttributesIfNoneAreAvailable() { // given AttributesExtractor, String> underTest = MessagingAttributesExtractor.create(TestGetter.INSTANCE, null); + AttributesExtractor, String> builtUnderTest = + MessagingAttributesExtractor.builder(TestGetter.INSTANCE, null).build(); Context context = Context.root(); // when AttributesBuilder startAttributes = Attributes.builder(); underTest.onStart(startAttributes, context, emptyMap()); + AttributesBuilder builtStartAttributes = Attributes.builder(); + builtUnderTest.onStart(builtStartAttributes, context, emptyMap()); AttributesBuilder endAttributes = Attributes.builder(); underTest.onEnd(endAttributes, context, emptyMap(), null, null); // then - assertThat(startAttributes.build().isEmpty()).isTrue(); - - assertThat(endAttributes.build().isEmpty()).isTrue(); + assertThat(startAttributes.build()).isEmpty(); + assertThat(builtStartAttributes.build()).isEmpty(); + assertThat(endAttributes.build()).isEmpty(); } enum TestGetter implements MessagingAttributesGetter, String> { @@ -184,7 +327,7 @@ public Long getMessageEnvelopeSize(Map request) { @Override public String getMessageId(Map request, String response) { - return response; + return request.get("messageId"); } @Nullable @@ -199,5 +342,10 @@ public Long getBatchMessageCount(Map request, @Nullable String r String payloadSize = request.get("batchMessageCount"); return payloadSize == null ? null : Long.valueOf(payloadSize); } + + @Override + public String getErrorType(Map request, String response, Throwable error) { + return "failure".equals(response) ? response : null; + } } } diff --git a/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractorTest.java b/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractorTest.java new file mode 100644 index 000000000000..a7dac2e5d31d --- /dev/null +++ b/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractorTest.java @@ -0,0 +1,74 @@ +/* + * 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 org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.params.provider.Arguments.argumentSet; + +import io.opentelemetry.api.trace.SpanKind; +import java.util.stream.Stream; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +class MessagingSpanKindExtractorTest { + + @ParameterizedTest + @MethodSource("spanKinds") + void extractsSpanKind( + MessagingOperationType operationType, + boolean isSpanContextPropagated, + SpanKind oldKind, + SpanKind kind) { + SpanKind actualKind = + MessagingSpanKindExtractor.create(operationType, isSpanContextPropagated) + .extract(new Object()); + + assertThat(actualKind).isEqualTo(emitStableMessagingSemconv() ? kind : oldKind); + } + + @Test + void sendDefaultsToPropagatedSpanContext() { + SpanKind spanKind = + MessagingSpanKindExtractor.create(MessagingOperationType.SEND).extract(new Object()); + + assertThat(spanKind).isEqualTo(SpanKind.PRODUCER); + } + + @Test + void messageOperationUsesLegacySpanKind() { + SpanKind receiveKind = + MessagingSpanKindExtractor.create(MessageOperation.RECEIVE).extract(new Object()); + + assertThat(receiveKind).isEqualTo(SpanKind.CONSUMER); + } + + private static Stream spanKinds() { + return Stream.of( + argumentSet( + "create", MessagingOperationType.CREATE, true, SpanKind.PRODUCER, SpanKind.PRODUCER), + argumentSet( + "send with propagated context", + MessagingOperationType.SEND, + true, + SpanKind.PRODUCER, + SpanKind.PRODUCER), + argumentSet( + "send without propagated context", + MessagingOperationType.SEND, + false, + SpanKind.PRODUCER, + SpanKind.CLIENT), + argumentSet( + "receive", MessagingOperationType.RECEIVE, true, SpanKind.CONSUMER, SpanKind.CLIENT), + argumentSet( + "process", MessagingOperationType.PROCESS, true, SpanKind.CONSUMER, SpanKind.CONSUMER), + argumentSet( + "settle", MessagingOperationType.SETTLE, true, SpanKind.CLIENT, SpanKind.CLIENT)); + } +} diff --git a/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorTest.java b/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorTest.java index 2d658df59f87..cc9c9f8d74f8 100644 --- a/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorTest.java +++ b/instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorTest.java @@ -5,11 +5,16 @@ package io.opentelemetry.instrumentation.api.incubator.semconv.messaging; +import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv; import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.params.provider.Arguments.argumentSet; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import io.opentelemetry.instrumentation.api.instrumenter.SpanNameExtractor; import java.util.stream.Stream; +import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; @@ -22,36 +27,118 @@ class MessagingSpanNameExtractorTest { @Mock MessagingAttributesGetter getter; + @Test + void shouldKeepLegacyNameForMessageOperation() { + Message message = new Message(); + when(getter.isTemporaryDestination(message)).thenReturn(false); + when(getter.getDestination(message)).thenReturn("destination"); + + SpanNameExtractor underTest = + MessagingSpanNameExtractor.create(getter, MessageOperation.PUBLISH); + + assertThat(underTest.extract(message)).isEqualTo("destination publish"); + } + @ParameterizedTest @MethodSource("spanNameParams") void shouldExtractSpanName( boolean isTemporaryQueue, + boolean isAnonymousQueue, String destinationName, - MessageOperation operation, - String expectedSpanName) { + String destinationTemplate, + MessagingOperationType operationType, + String operationName, + String oldSpanName, + String spanName) { // given Message message = new Message(); - if (isTemporaryQueue) { - when(getter.isTemporaryDestination(message)).thenReturn(true); + if (emitStableMessagingSemconv()) { + when(getter.getDestinationTemplate(message)).thenReturn(destinationTemplate); + if (destinationTemplate == null) { + when(getter.isTemporaryDestination(message)).thenReturn(isTemporaryQueue); + if (!isTemporaryQueue) { + when(getter.isAnonymousDestination(message)).thenReturn(isAnonymousQueue); + if (!isAnonymousQueue) { + when(getter.getDestination(message)).thenReturn(destinationName); + } + } + } } else { - when(getter.getDestination(message)).thenReturn(destinationName); + when(getter.isTemporaryDestination(message)).thenReturn(isTemporaryQueue); + if (!isTemporaryQueue) { + when(getter.getDestination(message)).thenReturn(destinationName); + } } - SpanNameExtractor underTest = MessagingSpanNameExtractor.create(getter, operation); + SpanNameExtractor underTest = + MessagingSpanNameExtractor.builder(getter, operationType) + .setOperationName(operationName) + .build(); // when - String spanName = underTest.extract(message); + String actualSpanName = underTest.extract(message); // then - assertThat(spanName).isEqualTo(expectedSpanName); + assertThat(actualSpanName).isEqualTo(emitStableMessagingSemconv() ? spanName : oldSpanName); + if (emitStableMessagingSemconv() && destinationTemplate != null) { + verify(getter, never()).isTemporaryDestination(message); + verify(getter, never()).isAnonymousDestination(message); + } } static Stream spanNameParams() { return Stream.of( - Arguments.of(false, "destination", MessageOperation.PUBLISH, "destination publish"), - Arguments.of(true, null, MessageOperation.PROCESS, "(temporary) process"), - Arguments.of(false, null, MessageOperation.RECEIVE, "unknown receive")); + argumentSet( + "operation name override", + false, + false, + "destination", + null, + MessagingOperationType.SEND, + "send", + "destination publish", + "send destination"), + argumentSet( + "temporary destination", + true, + false, + "generated", + "generated-{id}", + MessagingOperationType.PROCESS, + "process", + "(temporary) process", + "process generated-{id}"), + argumentSet( + "missing destination", + false, + false, + null, + null, + MessagingOperationType.RECEIVE, + "receive", + "unknown receive", + "receive"), + argumentSet( + "destination template", + false, + false, + "customer-42", + "customer-{id}", + MessagingOperationType.SEND, + "send", + "customer-42 publish", + "send customer-{id}"), + argumentSet( + "anonymous destination", + false, + true, + "generated", + "generated-{id}", + MessagingOperationType.PROCESS, + "process", + "generated process", + "process generated-{id}")); } static class Message {} diff --git a/instrumentation-api/src/main/java/io/opentelemetry/instrumentation/api/internal/SemconvStability.java b/instrumentation-api/src/main/java/io/opentelemetry/instrumentation/api/internal/SemconvStability.java index 8efdbfe4ee0b..d4465850c5ad 100644 --- a/instrumentation-api/src/main/java/io/opentelemetry/instrumentation/api/internal/SemconvStability.java +++ b/instrumentation-api/src/main/java/io/opentelemetry/instrumentation/api/internal/SemconvStability.java @@ -169,11 +169,14 @@ private static boolean emitStable(SemconvMode mode) { return mode.version() >= 1; } - public static boolean emitOldMessagingSemconv() { + public static boolean emitOldMessagingSemconv() { // to be removed in 3.0 return emitOldMessagingSemconv; } - public static boolean emitStableMessagingSemconv() { + // Returns whether the selected v1 experimental messaging semantic conventions should be emitted. + // The method name follows the existing pattern; it does not indicate that the messaging + // conventions are stable. + public static boolean emitStableMessagingSemconv() { // to be removed in 3.0 return emitStableMessagingSemconv; } diff --git a/instrumentation-api/src/main/java/io/opentelemetry/instrumentation/api/internal/SpanKey.java b/instrumentation-api/src/main/java/io/opentelemetry/instrumentation/api/internal/SpanKey.java index 0cf43cc6162a..a7e705ad8d30 100644 --- a/instrumentation-api/src/main/java/io/opentelemetry/instrumentation/api/internal/SpanKey.java +++ b/instrumentation-api/src/main/java/io/opentelemetry/instrumentation/api/internal/SpanKey.java @@ -43,12 +43,16 @@ public final class SpanKey { private static final ContextKey DB_CLIENT_KEY = ContextKey.named("opentelemetry-traces-span-key-db-client"); + private static final ContextKey PRODUCER_CREATE_KEY = + ContextKey.named("opentelemetry-traces-span-key-producer-create"); private static final ContextKey PRODUCER_KEY = ContextKey.named("opentelemetry-traces-span-key-producer"); private static final ContextKey CONSUMER_RECEIVE_KEY = ContextKey.named("opentelemetry-traces-span-key-consumer-receive"); private static final ContextKey CONSUMER_PROCESS_KEY = ContextKey.named("opentelemetry-traces-span-key-consumer-process"); + private static final ContextKey CONSUMER_SETTLE_KEY = + ContextKey.named("opentelemetry-traces-span-key-consumer-settle"); /* Span keys */ @@ -66,9 +70,11 @@ public final class SpanKey { public static final SpanKey RPC_CLIENT = new SpanKey(RPC_CLIENT_KEY); public static final SpanKey DB_CLIENT = new SpanKey(DB_CLIENT_KEY); + public static final SpanKey PRODUCER_CREATE = new SpanKey(PRODUCER_CREATE_KEY); public static final SpanKey PRODUCER = new SpanKey(PRODUCER_KEY); public static final SpanKey CONSUMER_RECEIVE = new SpanKey(CONSUMER_RECEIVE_KEY); public static final SpanKey CONSUMER_PROCESS = new SpanKey(CONSUMER_PROCESS_KEY); + public static final SpanKey CONSUMER_SETTLE = new SpanKey(CONSUMER_SETTLE_KEY); private final ContextKey key; diff --git a/instrumentation-api/src/test/java/io/opentelemetry/instrumentation/api/internal/SemconvStabilityTest.java b/instrumentation-api/src/test/java/io/opentelemetry/instrumentation/api/internal/SemconvStabilityTest.java index 32e316457df2..bc1d672f08bd 100644 --- a/instrumentation-api/src/test/java/io/opentelemetry/instrumentation/api/internal/SemconvStabilityTest.java +++ b/instrumentation-api/src/test/java/io/opentelemetry/instrumentation/api/internal/SemconvStabilityTest.java @@ -9,6 +9,7 @@ import static java.util.Collections.emptyList; import static java.util.Collections.emptySet; import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.params.provider.Arguments.argumentSet; import io.opentelemetry.api.incubator.config.DeclarativeConfigProperties; import io.opentelemetry.common.ComponentLoader; @@ -22,6 +23,8 @@ import java.util.stream.Collectors; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; import org.junit.jupiter.params.provider.ValueSource; class SemconvStabilityTest { @@ -379,6 +382,69 @@ void previewFallbackAppliesToPreviewDomains(boolean v3Preview) { assertThat(messaging).isEqualTo(SemconvMode.V1_EXPERIMENTAL); } + @ParameterizedTest + @MethodSource("messagingSelectionModes") + void messagingSelectionMatrix( + boolean v3Preview, Set stableOptIn, Set preview, SemconvMode expectedMode) { + SemconvMode messaging = + new SemconvSelectionResolver(general(), v3Preview, stableOptIn, preview).messaging(); + + assertThat(messaging).isEqualTo(expectedMode); + } + + private static List messagingSelectionModes() { + return asList( + argumentSet("legacy default", false, noStableOptIn(), noPreview(), SemconvMode.V0_STABLE), + argumentSet( + "legacy opt-in property", + false, + stableOptIn("messaging"), + noPreview(), + SemconvMode.V1_EXPERIMENTAL), + argumentSet( + "preview property", + false, + noStableOptIn(), + preview("messaging"), + SemconvMode.V1_EXPERIMENTAL), + argumentSet( + "legacy opt-in dual emit", + false, + stableOptIn("messaging/dup"), + noPreview(), + SemconvMode.V1_EXPERIMENTAL.withDualEmit()), + argumentSet( + "preview dual emit", + false, + noStableOptIn(), + preview("messaging/dup"), + SemconvMode.V1_EXPERIMENTAL.withDualEmit()), + argumentSet( + "v3 activation remains staged", + true, + noStableOptIn(), + noPreview(), + SemconvMode.V0_STABLE), + argumentSet( + "v3 ignores legacy opt-in property", + true, + stableOptIn("messaging"), + noPreview(), + SemconvMode.V0_STABLE), + argumentSet( + "v3 ignores legacy opt-in dual emit", + true, + stableOptIn("messaging/dup"), + noPreview(), + SemconvMode.V0_STABLE), + argumentSet( + "v3 with explicit preview", + true, + noStableOptIn(), + preview("messaging"), + SemconvMode.V1_EXPERIMENTAL)); + } + @SafeVarargs private static DeclarativeConfigProperties general(Entry... entries) { Map result = new HashMap(); diff --git a/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/instrumentationapi/v1_14/SpanKeyBridging.java b/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/instrumentationapi/v1_14/SpanKeyBridging.java index 24be05a6e8d0..e588faf6bd38 100644 --- a/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/instrumentationapi/v1_14/SpanKeyBridging.java +++ b/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/instrumentationapi/v1_14/SpanKeyBridging.java @@ -51,6 +51,7 @@ public class SpanKeyBridging { application.io.opentelemetry.instrumentation.api.internal.SpanKey.DB_CLIENT, SpanKey.DB_CLIENT); + putIfPresent(map, "PRODUCER_CREATE", SpanKey.PRODUCER_CREATE); map.put( application.io.opentelemetry.instrumentation.api.internal.SpanKey.PRODUCER, SpanKey.PRODUCER); @@ -60,9 +61,26 @@ public class SpanKeyBridging { map.put( application.io.opentelemetry.instrumentation.api.internal.SpanKey.CONSUMER_PROCESS, SpanKey.CONSUMER_PROCESS); + putIfPresent(map, "CONSUMER_SETTLE", SpanKey.CONSUMER_SETTLE); return map; } + private static void putIfPresent( + Map map, + String fieldName, + SpanKey agentSpanKey) { + try { + application.io.opentelemetry.instrumentation.api.internal.SpanKey applicationSpanKey = + (application.io.opentelemetry.instrumentation.api.internal.SpanKey) + application.io.opentelemetry.instrumentation.api.internal.SpanKey.class + .getField(fieldName) + .get(null); + map.put(applicationSpanKey, agentSpanKey); + } catch (ReflectiveOperationException ignored) { + // The application may use an instrumentation-api version from before this key was added. + } + } + @Nullable public static SpanKey toAgentOrNull( application.io.opentelemetry.instrumentation.api.internal.SpanKey applicationSpanKey) { diff --git a/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/javaagent/src/test/java/io/opentelemetry/javaagent/instrumentation/instrumentationapi/v1_14/ContextBridgeTest.java b/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/javaagent/src/test/java/io/opentelemetry/javaagent/instrumentation/instrumentationapi/v1_14/ContextBridgeTest.java index 60ecc0f4201b..7322e0d9569c 100644 --- a/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/javaagent/src/test/java/io/opentelemetry/javaagent/instrumentation/instrumentationapi/v1_14/ContextBridgeTest.java +++ b/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/javaagent/src/test/java/io/opentelemetry/javaagent/instrumentation/instrumentationapi/v1_14/ContextBridgeTest.java @@ -68,9 +68,11 @@ void testSpanKeyBridge() { SpanKey.HTTP_CLIENT, SpanKey.RPC_CLIENT, SpanKey.DB_CLIENT, + SpanKey.PRODUCER_CREATE, SpanKey.PRODUCER, SpanKey.CONSUMER_RECEIVE, - SpanKey.CONSUMER_PROCESS); + SpanKey.CONSUMER_PROCESS, + SpanKey.CONSUMER_SETTLE); spanKeys.forEach( spanKey -> assertThat(spanKey.fromContextOrNull(Context.current())).isNotNull()); diff --git a/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/testing/src/main/java/io/opentelemetry/javaagent/instrumentation/testing/AgentSpanTestingInstrumenter.java b/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/testing/src/main/java/io/opentelemetry/javaagent/instrumentation/testing/AgentSpanTestingInstrumenter.java index e13499836749..74a821731111 100644 --- a/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/testing/src/main/java/io/opentelemetry/javaagent/instrumentation/testing/AgentSpanTestingInstrumenter.java +++ b/instrumentation/opentelemetry-instrumentation-api/opentelemetry-instrumentation-api-1.14/testing/src/main/java/io/opentelemetry/javaagent/instrumentation/testing/AgentSpanTestingInstrumenter.java @@ -73,9 +73,11 @@ private static SpanKey[] getAllSpanKeys() { SpanKey.HTTP_CLIENT, SpanKey.RPC_CLIENT, SpanKey.DB_CLIENT, + SpanKey.PRODUCER_CREATE, SpanKey.PRODUCER, SpanKey.CONSUMER_RECEIVE, SpanKey.CONSUMER_PROCESS, + SpanKey.CONSUMER_SETTLE, }; } diff --git a/testing-common/src/main/java/io/opentelemetry/instrumentation/testing/junit/message/SemconvMessagingStabilityUtil.java b/testing-common/src/main/java/io/opentelemetry/instrumentation/testing/junit/message/SemconvMessagingStabilityUtil.java new file mode 100644 index 000000000000..e5444956add0 --- /dev/null +++ b/testing-common/src/main/java/io/opentelemetry/instrumentation/testing/junit/message/SemconvMessagingStabilityUtil.java @@ -0,0 +1,36 @@ +/* + * Copyright The OpenTelemetry Authors + * SPDX-License-Identifier: Apache-2.0 + */ + +package io.opentelemetry.instrumentation.testing.junit.message; + +import static io.opentelemetry.api.common.AttributeKey.stringKey; +import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv; + +import io.opentelemetry.api.common.AttributeKey; +import java.util.HashMap; +import java.util.Map; + +// supports asserting on old messaging semconv, to be removed in 3.0 +public class SemconvMessagingStabilityUtil { + + private static final Map, AttributeKey> oldToNewMap = buildMap(); + + private static Map, AttributeKey> buildMap() { + Map, AttributeKey> map = new HashMap<>(); + map.put(stringKey("messaging.client_id"), stringKey("messaging.client.id")); + return map; + } + + @SuppressWarnings("unchecked") + public static AttributeKey effectiveKey(AttributeKey oldKey) { + // not testing messaging/dup + if (emitStableMessagingSemconv()) { + return (AttributeKey) oldToNewMap.getOrDefault(oldKey, oldKey); + } + return oldKey; + } + + private SemconvMessagingStabilityUtil() {} +}