Skip to content

Commit ec5eb7f

Browse files
committed
Migrate RabbitMQ messaging telemetry to v1.43
1 parent 9b5ec57 commit ec5eb7f

14 files changed

Lines changed: 580 additions & 129 deletions

File tree

instrumentation/rabbitmq-2.7/javaagent/build.gradle.kts

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,14 @@ tasks {
4444
systemProperty("metadataConfig", "otel.instrumentation.rabbitmq.experimental-span-attributes=true")
4545
}
4646

47+
val testMessagingPreview = register<Test>("testMessagingPreview") {
48+
testClassesDirs = sourceSets.test.get().output.classesDirs
49+
classpath = sourceSets.test.get().runtimeClasspath
50+
jvmArgs("-Dotel.semconv-stability.preview=messaging")
51+
systemProperty("metadataConfig", "otel.semconv-stability.opt-in=messaging")
52+
}
53+
4754
check {
48-
dependsOn(testExperimental)
55+
dependsOn(testExperimental, testMessagingPreview)
4956
}
5057
}

instrumentation/rabbitmq-2.7/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/rabbitmq/v2_7/RabbitDeliveryAttributesGetter.java

Lines changed: 10 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,9 @@
55

66
package io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7;
77

8+
import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
9+
import static io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7.RabbitInstrumenterHelper.consumerDestinationName;
10+
import static io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7.RabbitInstrumenterHelper.isGeneratedQueueName;
811
import static java.util.Collections.emptyList;
912
import static java.util.Collections.singletonList;
1013

@@ -24,11 +27,13 @@ public String getSystem(DeliveryRequest request) {
2427
@Nullable
2528
@Override
2629
public String getDestination(DeliveryRequest request) {
27-
if (request.getEnvelope() != null) {
28-
return normalizeExchangeName(request.getEnvelope().getExchange());
29-
} else {
30-
return null;
30+
if (emitStableMessagingSemconv()) {
31+
return consumerDestinationName(
32+
request.getEnvelope().getExchange(),
33+
request.getEnvelope().getRoutingKey(),
34+
request.getQueue());
3135
}
36+
return normalizeExchangeName(request.getEnvelope().getExchange());
3237
}
3338

3439
@Nullable
@@ -48,7 +53,7 @@ public boolean isTemporaryDestination(DeliveryRequest request) {
4853

4954
@Override
5055
public boolean isAnonymousDestination(DeliveryRequest request) {
51-
return false;
56+
return isGeneratedQueueName(request.getQueue());
5257
}
5358

5459
@Nullable

instrumentation/rabbitmq-2.7/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/rabbitmq/v2_7/RabbitInstrumenterHelper.java

Lines changed: 76 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,9 @@
55

66
package io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7;
77

8+
import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
89
import static io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7.RabbitSingletons.CHANNEL_AND_METHOD_CONTEXT_KEY;
10+
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_DESTINATION_ANONYMOUS;
911
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_DESTINATION_NAME;
1012
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_RABBITMQ_DESTINATION_ROUTING_KEY;
1113

@@ -17,6 +19,7 @@
1719
import io.opentelemetry.context.Context;
1820
import io.opentelemetry.instrumentation.api.incubator.config.internal.DeclarativeConfigUtil;
1921
import java.util.Map;
22+
import javax.annotation.Nullable;
2023

2124
public class RabbitInstrumenterHelper {
2225
static final AttributeKey<String> RABBITMQ_COMMAND = AttributeKey.stringKey("rabbitmq.command");
@@ -33,8 +36,19 @@ public static RabbitInstrumenterHelper helper() {
3336

3437
public void onPublish(Span span, String exchange, String routingKey) {
3538
String exchangeName = normalizeExchangeName(exchange);
36-
span.setAttribute(MESSAGING_DESTINATION_NAME, exchangeName);
37-
span.updateName(exchangeName + " publish");
39+
if (emitStableMessagingSemconv()) {
40+
String destinationName = producerDestinationName(exchange, routingKey);
41+
span.setAttribute(MESSAGING_DESTINATION_NAME, destinationName);
42+
boolean anonymousDestination =
43+
isDefaultExchange(exchange) && isGeneratedQueueName(routingKey);
44+
if (anonymousDestination) {
45+
span.setAttribute(MESSAGING_DESTINATION_ANONYMOUS, true);
46+
}
47+
span.updateName(anonymousDestination ? "publish" : "publish " + destinationName);
48+
} else {
49+
span.setAttribute(MESSAGING_DESTINATION_NAME, exchangeName);
50+
span.updateName(exchangeName + " publish");
51+
}
3852
if (routingKey != null && !routingKey.isEmpty()) {
3953
span.setAttribute(MESSAGING_RABBITMQ_DESTINATION_ROUTING_KEY, routingKey);
4054
}
@@ -58,7 +72,66 @@ public void onProps(Context context, Span span, AMQP.BasicProperties props) {
5872
}
5973

6074
private static String normalizeExchangeName(String exchange) {
61-
return exchange == null || exchange.isEmpty() ? "<default>" : exchange;
75+
return isDefaultExchange(exchange) ? "<default>" : exchange;
76+
}
77+
78+
private static boolean isDefaultExchange(@Nullable String exchange) {
79+
return exchange == null || exchange.isEmpty();
80+
}
81+
82+
static boolean isGeneratedQueueName(@Nullable String queue) {
83+
if (queue == null) {
84+
return false;
85+
}
86+
if (queue.startsWith("amq.gen-") || queue.startsWith("spring.gen-")) {
87+
return true;
88+
}
89+
return isCanonicalUuid(queue);
90+
}
91+
92+
private static boolean isCanonicalUuid(String value) {
93+
if (value.length() != 36) {
94+
return false;
95+
}
96+
for (int i = 0; i < value.length(); i++) {
97+
char ch = value.charAt(i);
98+
if (i == 8 || i == 13 || i == 18 || i == 23) {
99+
if (ch != '-') {
100+
return false;
101+
}
102+
} else if (!((ch >= '0' && ch <= '9') || (ch >= 'a' && ch <= 'f'))) {
103+
return false;
104+
}
105+
}
106+
return true;
107+
}
108+
109+
static String producerDestinationName(String exchange, String routingKey) {
110+
StringBuilder destination = new StringBuilder();
111+
appendDestinationPart(destination, exchange);
112+
appendDestinationPart(destination, routingKey);
113+
return destination.length() == 0 ? "amq.default" : destination.toString();
114+
}
115+
116+
@Nullable
117+
static String consumerDestinationName(String exchange, String routingKey, String queue) {
118+
StringBuilder destination = new StringBuilder();
119+
appendDestinationPart(destination, exchange);
120+
appendDestinationPart(destination, routingKey);
121+
if (queue != null && !queue.equals(routingKey)) {
122+
appendDestinationPart(destination, queue);
123+
}
124+
return destination.length() == 0 ? null : destination.toString();
125+
}
126+
127+
private static void appendDestinationPart(StringBuilder destination, String part) {
128+
if (part == null || part.isEmpty()) {
129+
return;
130+
}
131+
if (destination.length() != 0) {
132+
destination.append(':');
133+
}
134+
destination.append(part);
62135
}
63136

64137
public static void onCommand(Span span, Command command) {

instrumentation/rabbitmq-2.7/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/rabbitmq/v2_7/RabbitReceiveAttributesGetter.java

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,9 @@
55

66
package io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7;
77

8+
import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
9+
import static io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7.RabbitInstrumenterHelper.consumerDestinationName;
10+
import static io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7.RabbitInstrumenterHelper.isGeneratedQueueName;
811
import static java.util.Collections.emptyList;
912
import static java.util.Collections.singletonList;
1013

@@ -25,11 +28,17 @@ public String getSystem(ReceiveRequest request) {
2528
@Nullable
2629
@Override
2730
public String getDestination(ReceiveRequest request) {
28-
if (request.getResponse() != null) {
29-
return normalizeExchangeName(request.getResponse().getEnvelope().getExchange());
30-
} else {
31+
GetResponse response = request.getResponse();
32+
if (emitStableMessagingSemconv()) {
33+
return consumerDestinationName(
34+
response == null ? null : response.getEnvelope().getExchange(),
35+
response == null ? null : response.getEnvelope().getRoutingKey(),
36+
request.getQueue());
37+
}
38+
if (response == null) {
3139
return null;
3240
}
41+
return normalizeExchangeName(response.getEnvelope().getExchange());
3342
}
3443

3544
@Nullable
@@ -49,7 +58,7 @@ public boolean isTemporaryDestination(ReceiveRequest request) {
4958

5059
@Override
5160
public boolean isAnonymousDestination(ReceiveRequest request) {
52-
return false;
61+
return isGeneratedQueueName(request.getQueue());
5362
}
5463

5564
@Nullable
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,41 @@
1+
/*
2+
* Copyright The OpenTelemetry Authors
3+
* SPDX-License-Identifier: Apache-2.0
4+
*/
5+
6+
package io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7;
7+
8+
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_RABBITMQ_DESTINATION_ROUTING_KEY;
9+
10+
import com.rabbitmq.client.GetResponse;
11+
import io.opentelemetry.api.common.AttributesBuilder;
12+
import io.opentelemetry.context.Context;
13+
import io.opentelemetry.instrumentation.api.instrumenter.AttributesExtractor;
14+
import javax.annotation.Nullable;
15+
16+
class RabbitReceiveExtraAttributesExtractor
17+
implements AttributesExtractor<ReceiveRequest, GetResponse> {
18+
19+
@Override
20+
public void onStart(
21+
AttributesBuilder attributes, Context parentContext, ReceiveRequest request) {}
22+
23+
@Override
24+
public void onEnd(
25+
AttributesBuilder attributes,
26+
Context context,
27+
ReceiveRequest request,
28+
@Nullable GetResponse response,
29+
@Nullable Throwable error) {
30+
if (response == null) {
31+
response = request.getResponse();
32+
if (response == null) {
33+
return;
34+
}
35+
}
36+
String routingKey = response.getEnvelope().getRoutingKey();
37+
if (routingKey != null && !routingKey.isEmpty()) {
38+
attributes.put(MESSAGING_RABBITMQ_DESTINATION_ROUTING_KEY, routingKey);
39+
}
40+
}
41+
}

instrumentation/rabbitmq-2.7/javaagent/src/main/java/io/opentelemetry/javaagent/instrumentation/rabbitmq/v2_7/RabbitSingletons.java

Lines changed: 37 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -10,23 +10,26 @@
1010
import static io.opentelemetry.instrumentation.api.incubator.semconv.messaging.internal.MessagingExceptionEventExtractors.setMessagingProcessExceptionEventExtractor;
1111
import static io.opentelemetry.instrumentation.api.incubator.semconv.messaging.internal.MessagingExceptionEventExtractors.setMessagingReceiveExceptionEventExtractor;
1212
import static io.opentelemetry.instrumentation.api.incubator.semconv.messaging.internal.MessagingExceptionEventExtractors.setMessagingSendExceptionEventExtractor;
13+
import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
1314

1415
import com.rabbitmq.client.GetResponse;
1516
import io.opentelemetry.api.GlobalOpenTelemetry;
1617
import io.opentelemetry.context.ContextKey;
17-
import io.opentelemetry.instrumentation.api.incubator.semconv.messaging.MessageOperation;
1818
import io.opentelemetry.instrumentation.api.incubator.semconv.messaging.MessagingAttributesExtractor;
1919
import io.opentelemetry.instrumentation.api.incubator.semconv.messaging.MessagingAttributesGetter;
20+
import io.opentelemetry.instrumentation.api.incubator.semconv.messaging.MessagingOperationType;
21+
import io.opentelemetry.instrumentation.api.incubator.semconv.messaging.MessagingSpanKindExtractor;
22+
import io.opentelemetry.instrumentation.api.incubator.semconv.messaging.MessagingSpanNameExtractor;
23+
import io.opentelemetry.instrumentation.api.incubator.semconv.messaging.internal.MessagingProcessInstrumenterFactory;
2024
import io.opentelemetry.instrumentation.api.instrumenter.AttributesExtractor;
2125
import io.opentelemetry.instrumentation.api.instrumenter.Instrumenter;
2226
import io.opentelemetry.instrumentation.api.instrumenter.InstrumenterBuilder;
23-
import io.opentelemetry.instrumentation.api.instrumenter.SpanKindExtractor;
27+
import io.opentelemetry.instrumentation.api.instrumenter.SpanNameExtractor;
2428
import io.opentelemetry.instrumentation.api.internal.PropagatorBasedSpanLinksExtractor;
2529
import io.opentelemetry.instrumentation.api.semconv.network.NetworkAttributesExtractor;
2630
import io.opentelemetry.javaagent.bootstrap.internal.ExperimentalConfig;
2731
import java.util.ArrayList;
2832
import java.util.List;
29-
import javax.annotation.Nullable;
3033

3134
public class RabbitSingletons {
3235

@@ -62,8 +65,13 @@ private static Instrumenter<ChannelAndMethod, Void> createChannelInstrumenter(bo
6265
Instrumenter.<ChannelAndMethod, Void>builder(
6366
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, ChannelAndMethod::getMethod)
6467
.addAttributesExtractor(
65-
buildMessagingAttributesExtractor(
66-
new RabbitChannelAttributesGetter(), publish ? MessageOperation.PUBLISH : null))
68+
publish
69+
? buildMessagingAttributesExtractor(
70+
new RabbitChannelAttributesGetter(), MessagingOperationType.SEND)
71+
: MessagingAttributesExtractor.builder(
72+
new RabbitChannelAttributesGetter(), null)
73+
.setCapturedHeaders(ExperimentalConfig.get().getMessagingHeaders())
74+
.build())
6775
.addAttributesExtractor(
6876
NetworkAttributesExtractor.create(new RabbitChannelNetAttributesGetter()))
6977
.addContextCustomizer(
@@ -77,50 +85,62 @@ private static Instrumenter<ChannelAndMethod, Void> createChannelInstrumenter(bo
7785
}
7886

7987
private static Instrumenter<ReceiveRequest, GetResponse> createReceiveInstrumenter() {
88+
RabbitReceiveAttributesGetter getter = new RabbitReceiveAttributesGetter();
8089
List<AttributesExtractor<ReceiveRequest, GetResponse>> extractors = new ArrayList<>();
81-
extractors.add(
82-
buildMessagingAttributesExtractor(
83-
new RabbitReceiveAttributesGetter(), MessageOperation.RECEIVE));
90+
extractors.add(buildMessagingAttributesExtractor(getter, MessagingOperationType.RECEIVE));
91+
extractors.add(new RabbitReceiveExtraAttributesExtractor());
8492
extractors.add(NetworkAttributesExtractor.create(new RabbitReceiveNetAttributesGetter()));
8593
if (RabbitInstrumenterHelper.CAPTURE_EXPERIMENTAL_SPAN_ATTRIBUTES) {
8694
extractors.add(new RabbitReceiveExperimentalAttributesExtractor());
8795
}
8896

97+
SpanNameExtractor<ReceiveRequest> spanNameExtractor =
98+
emitStableMessagingSemconv()
99+
? MessagingSpanNameExtractor.create(getter, MessagingOperationType.RECEIVE)
100+
: ReceiveRequest::spanName;
89101
InstrumenterBuilder<ReceiveRequest, GetResponse> builder =
90102
Instrumenter.<ReceiveRequest, GetResponse>builder(
91-
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, ReceiveRequest::spanName)
103+
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, spanNameExtractor)
92104
.addAttributesExtractors(extractors)
93105
.setEnabled(ExperimentalConfig.get().messagingReceiveInstrumentationEnabled())
94106
.addSpanLinksExtractor(
95107
new PropagatorBasedSpanLinksExtractor<>(
96108
GlobalOpenTelemetry.getPropagators().getTextMapPropagator(),
97109
new ReceiveRequestTextMapGetter()));
98110
setMessagingReceiveExceptionEventExtractor(builder);
99-
return builder.buildInstrumenter(SpanKindExtractor.alwaysConsumer());
111+
return builder.buildInstrumenter(
112+
MessagingSpanKindExtractor.create(MessagingOperationType.RECEIVE));
100113
}
101114

102115
private static Instrumenter<DeliveryRequest, Void> createDeliverInstrumenter() {
116+
RabbitDeliveryAttributesGetter getter = new RabbitDeliveryAttributesGetter();
103117
List<AttributesExtractor<DeliveryRequest, Void>> extractors = new ArrayList<>();
104-
extractors.add(
105-
buildMessagingAttributesExtractor(
106-
new RabbitDeliveryAttributesGetter(), MessageOperation.PROCESS));
118+
extractors.add(buildMessagingAttributesExtractor(getter, MessagingOperationType.PROCESS));
107119
extractors.add(NetworkAttributesExtractor.create(new RabbitDeliveryNetAttributesGetter()));
108120
extractors.add(new RabbitDeliveryExtraAttributesExtractor());
109121
if (RabbitInstrumenterHelper.CAPTURE_EXPERIMENTAL_SPAN_ATTRIBUTES) {
110122
extractors.add(new RabbitDeliveryExperimentalAttributesExtractor());
111123
}
112124

125+
SpanNameExtractor<DeliveryRequest> spanNameExtractor =
126+
emitStableMessagingSemconv()
127+
? MessagingSpanNameExtractor.create(getter, MessagingOperationType.PROCESS)
128+
: DeliveryRequest::spanName;
113129
InstrumenterBuilder<DeliveryRequest, Void> builder =
114130
Instrumenter.<DeliveryRequest, Void>builder(
115-
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, DeliveryRequest::spanName)
131+
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, spanNameExtractor)
116132
.addAttributesExtractors(extractors);
117133
setMessagingProcessExceptionEventExtractor(builder);
118-
return builder.buildConsumerInstrumenter(new DeliveryRequestGetter());
134+
return MessagingProcessInstrumenterFactory.create(
135+
builder,
136+
GlobalOpenTelemetry.getPropagators().getTextMapPropagator(),
137+
new DeliveryRequestGetter(),
138+
false);
119139
}
120140

121141
private static <T, V> AttributesExtractor<T, V> buildMessagingAttributesExtractor(
122-
MessagingAttributesGetter<T, V> getter, @Nullable MessageOperation operation) {
123-
return MessagingAttributesExtractor.builder(getter, operation)
142+
MessagingAttributesGetter<T, V> getter, MessagingOperationType operationType) {
143+
return MessagingAttributesExtractor.builderForOperationType(getter, operationType)
124144
.setCapturedHeaders(ExperimentalConfig.get().getMessagingHeaders())
125145
.build();
126146
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
/*
2+
* Copyright The OpenTelemetry Authors
3+
* SPDX-License-Identifier: Apache-2.0
4+
*/
5+
6+
package io.opentelemetry.javaagent.instrumentation.rabbitmq.v2_7;
7+
8+
import static org.assertj.core.api.Assertions.assertThat;
9+
import static org.junit.jupiter.params.provider.Arguments.argumentSet;
10+
11+
import java.util.stream.Stream;
12+
import org.junit.jupiter.params.ParameterizedTest;
13+
import org.junit.jupiter.params.provider.Arguments;
14+
import org.junit.jupiter.params.provider.MethodSource;
15+
16+
class RabbitInstrumenterHelperTest {
17+
18+
@ParameterizedTest
19+
@MethodSource("queueNames")
20+
void identifiesGeneratedQueueNames(String queue, boolean expected) {
21+
assertThat(RabbitInstrumenterHelper.isGeneratedQueueName(queue)).isEqualTo(expected);
22+
}
23+
24+
private static Stream<Arguments> queueNames() {
25+
return Stream.of(
26+
argumentSet("null", null, false),
27+
argumentSet("named queue", "orders", false),
28+
argumentSet("RabbitMQ generated", "amq.gen-random", true),
29+
argumentSet("Spring generated", "spring.gen-random", true),
30+
argumentSet("canonical UUID", "123e4567-e89b-12d3-a456-426614174000", true),
31+
argumentSet("uppercase UUID", "123E4567-E89B-12D3-A456-426614174000", false),
32+
argumentSet("malformed UUID", "123e4567-e89b-12d3-a456-42661417400g", false));
33+
}
34+
}

0 commit comments

Comments
 (0)