Skip to content

Commit a515ee0

Browse files
committed
Migrate RabbitMQ messaging telemetry to v1.43
1 parent 02ef2b3 commit a515ee0

10 files changed

Lines changed: 403 additions & 123 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: 50 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
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;
910
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_DESTINATION_NAME;
1011
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_RABBITMQ_DESTINATION_ROUTING_KEY;
@@ -17,6 +18,7 @@
1718
import io.opentelemetry.context.Context;
1819
import io.opentelemetry.instrumentation.api.incubator.config.internal.DeclarativeConfigUtil;
1920
import java.util.Map;
21+
import javax.annotation.Nullable;
2022

2123
public class RabbitInstrumenterHelper {
2224
static final AttributeKey<String> RABBITMQ_COMMAND = AttributeKey.stringKey("rabbitmq.command");
@@ -33,8 +35,17 @@ public static RabbitInstrumenterHelper helper() {
3335

3436
public void onPublish(Span span, String exchange, String routingKey) {
3537
String exchangeName = normalizeExchangeName(exchange);
36-
span.setAttribute(MESSAGING_DESTINATION_NAME, exchangeName);
37-
span.updateName(exchangeName + " publish");
38+
if (emitStableMessagingSemconv()) {
39+
String destinationName = producerDestinationName(exchange, routingKey);
40+
span.setAttribute(MESSAGING_DESTINATION_NAME, destinationName);
41+
span.updateName(
42+
isDefaultExchange(exchange) && isGeneratedQueueName(routingKey)
43+
? "publish"
44+
: "publish " + destinationName);
45+
} else {
46+
span.setAttribute(MESSAGING_DESTINATION_NAME, exchangeName);
47+
span.updateName(exchangeName + " publish");
48+
}
3849
if (routingKey != null && !routingKey.isEmpty()) {
3950
span.setAttribute(MESSAGING_RABBITMQ_DESTINATION_ROUTING_KEY, routingKey);
4051
}
@@ -58,7 +69,43 @@ public void onProps(Context context, Span span, AMQP.BasicProperties props) {
5869
}
5970

6071
private static String normalizeExchangeName(String exchange) {
61-
return exchange == null || exchange.isEmpty() ? "<default>" : exchange;
72+
return isDefaultExchange(exchange) ? "<default>" : exchange;
73+
}
74+
75+
private static boolean isDefaultExchange(@Nullable String exchange) {
76+
return exchange == null || exchange.isEmpty();
77+
}
78+
79+
static boolean isGeneratedQueueName(@Nullable String queue) {
80+
return queue != null && (queue.startsWith("amq.gen-") || queue.startsWith("spring.gen-"));
81+
}
82+
83+
static String producerDestinationName(String exchange, String routingKey) {
84+
StringBuilder destination = new StringBuilder();
85+
appendDestinationPart(destination, exchange);
86+
appendDestinationPart(destination, routingKey);
87+
return destination.length() == 0 ? "amq.default" : destination.toString();
88+
}
89+
90+
@Nullable
91+
static String consumerDestinationName(String exchange, String routingKey, String queue) {
92+
StringBuilder destination = new StringBuilder();
93+
appendDestinationPart(destination, exchange);
94+
appendDestinationPart(destination, routingKey);
95+
if (queue != null && !queue.equals(routingKey)) {
96+
appendDestinationPart(destination, queue);
97+
}
98+
return destination.length() == 0 ? null : destination.toString();
99+
}
100+
101+
private static void appendDestinationPart(StringBuilder destination, String part) {
102+
if (part == null || part.isEmpty()) {
103+
return;
104+
}
105+
if (destination.length() != 0) {
106+
destination.append(':');
107+
}
108+
destination.append(part);
62109
}
63110

64111
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

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

Lines changed: 36 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,61 @@ 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));
8491
extractors.add(NetworkAttributesExtractor.create(new RabbitReceiveNetAttributesGetter()));
8592
if (RabbitInstrumenterHelper.CAPTURE_EXPERIMENTAL_SPAN_ATTRIBUTES) {
8693
extractors.add(new RabbitReceiveExperimentalAttributesExtractor());
8794
}
8895

96+
SpanNameExtractor<ReceiveRequest> spanNameExtractor =
97+
emitStableMessagingSemconv()
98+
? MessagingSpanNameExtractor.create(getter, MessagingOperationType.RECEIVE)
99+
: ReceiveRequest::spanName;
89100
InstrumenterBuilder<ReceiveRequest, GetResponse> builder =
90101
Instrumenter.<ReceiveRequest, GetResponse>builder(
91-
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, ReceiveRequest::spanName)
102+
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, spanNameExtractor)
92103
.addAttributesExtractors(extractors)
93104
.setEnabled(ExperimentalConfig.get().messagingReceiveInstrumentationEnabled())
94105
.addSpanLinksExtractor(
95106
new PropagatorBasedSpanLinksExtractor<>(
96107
GlobalOpenTelemetry.getPropagators().getTextMapPropagator(),
97108
new ReceiveRequestTextMapGetter()));
98109
setMessagingReceiveExceptionEventExtractor(builder);
99-
return builder.buildInstrumenter(SpanKindExtractor.alwaysConsumer());
110+
return builder.buildInstrumenter(
111+
MessagingSpanKindExtractor.create(MessagingOperationType.RECEIVE));
100112
}
101113

102114
private static Instrumenter<DeliveryRequest, Void> createDeliverInstrumenter() {
115+
RabbitDeliveryAttributesGetter getter = new RabbitDeliveryAttributesGetter();
103116
List<AttributesExtractor<DeliveryRequest, Void>> extractors = new ArrayList<>();
104-
extractors.add(
105-
buildMessagingAttributesExtractor(
106-
new RabbitDeliveryAttributesGetter(), MessageOperation.PROCESS));
117+
extractors.add(buildMessagingAttributesExtractor(getter, MessagingOperationType.PROCESS));
107118
extractors.add(NetworkAttributesExtractor.create(new RabbitDeliveryNetAttributesGetter()));
108119
extractors.add(new RabbitDeliveryExtraAttributesExtractor());
109120
if (RabbitInstrumenterHelper.CAPTURE_EXPERIMENTAL_SPAN_ATTRIBUTES) {
110121
extractors.add(new RabbitDeliveryExperimentalAttributesExtractor());
111122
}
112123

124+
SpanNameExtractor<DeliveryRequest> spanNameExtractor =
125+
emitStableMessagingSemconv()
126+
? MessagingSpanNameExtractor.create(getter, MessagingOperationType.PROCESS)
127+
: DeliveryRequest::spanName;
113128
InstrumenterBuilder<DeliveryRequest, Void> builder =
114129
Instrumenter.<DeliveryRequest, Void>builder(
115-
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, DeliveryRequest::spanName)
130+
GlobalOpenTelemetry.get(), INSTRUMENTATION_NAME, spanNameExtractor)
116131
.addAttributesExtractors(extractors);
117132
setMessagingProcessExceptionEventExtractor(builder);
118-
return builder.buildConsumerInstrumenter(new DeliveryRequestGetter());
133+
return MessagingProcessInstrumenterFactory.create(
134+
builder,
135+
GlobalOpenTelemetry.getPropagators().getTextMapPropagator(),
136+
new DeliveryRequestGetter(),
137+
false);
119138
}
120139

121140
private static <T, V> AttributesExtractor<T, V> buildMessagingAttributesExtractor(
122-
MessagingAttributesGetter<T, V> getter, @Nullable MessageOperation operation) {
123-
return MessagingAttributesExtractor.builder(getter, operation)
141+
MessagingAttributesGetter<T, V> getter, MessagingOperationType operationType) {
142+
return MessagingAttributesExtractor.builderForOperationType(getter, operationType)
124143
.setCapturedHeaders(ExperimentalConfig.get().getMessagingHeaders())
125144
.build();
126145
}

0 commit comments

Comments
 (0)