Skip to content

Commit adccb3b

Browse files
committed
Pin messaging semantic conventions to v1.43 behind preview
1 parent 18f13c2 commit adccb3b

22 files changed

Lines changed: 1191 additions & 158 deletions

File tree

instrumentation-api-incubator/build.gradle.kts

Lines changed: 17 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -92,13 +92,23 @@ tasks {
9292
testClassesDirs = sourceSets.test.get().output.classesDirs
9393
classpath = sourceSets.test.get().runtimeClasspath
9494
jvmArgs("-Dotel.semconv-stability.opt-in=database,code,service.peer,rpc")
95+
jvmArgs("-Dotel.semconv-stability.preview=messaging")
9596
inputs.dir(jflexOutputDir)
9697
}
9798

9899
val testBothSemconv = register<Test>("testBothSemconv") {
99100
testClassesDirs = sourceSets.test.get().output.classesDirs
100101
classpath = sourceSets.test.get().runtimeClasspath
101102
jvmArgs("-Dotel.semconv-stability.opt-in=database/dup,code/dup,service.peer/dup,rpc/dup")
103+
jvmArgs("-Dotel.semconv-stability.preview=messaging/dup")
104+
inputs.dir(jflexOutputDir)
105+
}
106+
107+
val testV3Preview = register<Test>("testV3Preview") {
108+
testClassesDirs = sourceSets.test.get().output.classesDirs
109+
classpath = sourceSets.test.get().runtimeClasspath
110+
jvmArgs("-Dotel.instrumentation.common.v3-preview=true")
111+
jvmArgs("-Dotel.semconv-stability.preview=messaging")
102112
inputs.dir(jflexOutputDir)
103113
}
104114

@@ -117,6 +127,12 @@ tasks {
117127
}
118128

119129
check {
120-
dependsOn(testStableSemconv, testBothSemconv, testExceptionSignalLogs, testExceptionSignalLogsDup)
130+
dependsOn(
131+
testStableSemconv,
132+
testBothSemconv,
133+
testV3Preview,
134+
testExceptionSignalLogs,
135+
testExceptionSignalLogsDup,
136+
)
121137
}
122138
}

instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/db/SqlQueryAnalyzer.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,9 @@ private static SqlQuery analyzeWithSummaryImpl(String query, SqlDialect dialect)
7878

7979
// visible for tests
8080
static boolean isCached(String query, SqlDialect dialect) {
81-
return sqlToQueryCache.get(CacheKey.create(query, dialect)) != null;
81+
Cache<CacheKey, SqlQuery> cache =
82+
SemconvStability.v3Preview() ? sqlToQueryCacheWithSummary : sqlToQueryCache;
83+
return cache.get(CacheKey.create(query, dialect)) != null;
8284
}
8385

8486
@AutoValue

instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessageOperation.java

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@
99

1010
/**
1111
* Represents type of <a
12-
* href="https://github.com/open-telemetry/semantic-conventions/blob/main/docs/messaging/messaging-spans.md#operation-names">operations</a>
12+
* href="https://github.com/open-telemetry/semantic-conventions/blob/v1.43.0/docs/messaging/messaging-spans.md#operation-types">operations</a>
1313
* that may be used in a messaging system.
1414
*/
1515
public enum MessageOperation {
@@ -18,9 +18,8 @@ public enum MessageOperation {
1818
PROCESS;
1919

2020
/**
21-
* Returns the operation name as defined in <a
22-
* href="https://github.com/open-telemetry/semantic-conventions/blob/main/docs/messaging/messaging-spans.md#operation-names">the
23-
* specification</a>.
21+
* Returns the legacy operation name. The v1.43 operation name defaults to this value unless an
22+
* instrumentation supplies a system-specific override.
2423
*/
2524
String operationName() {
2625
return name().toLowerCase(Locale.ROOT);

instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractor.java

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

66
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
77

8+
import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitOldMessagingSemconv;
9+
import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
10+
import static io.opentelemetry.semconv.ErrorAttributes.ERROR_TYPE;
11+
812
import io.opentelemetry.api.common.AttributeKey;
913
import io.opentelemetry.api.common.AttributesBuilder;
1014
import io.opentelemetry.context.Context;
@@ -17,7 +21,7 @@
1721

1822
/**
1923
* Extractor of <a
20-
* href="https://github.com/open-telemetry/semantic-conventions/blob/main/docs/messaging/messaging-spans.md">messaging
24+
* href="https://github.com/open-telemetry/semantic-conventions/blob/v1.43.0/docs/messaging/messaging-spans.md">messaging
2125
* attributes</a>.
2226
*
2327
* <p>This class delegates to a type-specific {@link MessagingAttributesGetter} for individual
@@ -29,8 +33,10 @@ public final class MessagingAttributesExtractor<REQUEST, RESPONSE>
2933
// copied from MessagingIncubatingAttributes
3034
private static final AttributeKey<Long> MESSAGING_BATCH_MESSAGE_COUNT =
3135
AttributeKey.longKey("messaging.batch.message_count");
32-
private static final AttributeKey<String> MESSAGING_CLIENT_ID =
36+
private static final AttributeKey<String> MESSAGING_CLIENT_ID_OLD =
3337
AttributeKey.stringKey("messaging.client_id");
38+
private static final AttributeKey<String> MESSAGING_CLIENT_ID =
39+
AttributeKey.stringKey("messaging.client.id");
3440
private static final AttributeKey<Boolean> MESSAGING_DESTINATION_ANONYMOUS =
3541
AttributeKey.booleanKey("messaging.destination.anonymous");
3642
private static final AttributeKey<String> MESSAGING_DESTINATION_NAME =
@@ -51,6 +57,10 @@ public final class MessagingAttributesExtractor<REQUEST, RESPONSE>
5157
AttributeKey.stringKey("messaging.message.id");
5258
private static final AttributeKey<String> MESSAGING_OPERATION =
5359
AttributeKey.stringKey("messaging.operation");
60+
private static final AttributeKey<String> MESSAGING_OPERATION_NAME =
61+
AttributeKey.stringKey("messaging.operation.name");
62+
private static final AttributeKey<String> MESSAGING_OPERATION_TYPE =
63+
AttributeKey.stringKey("messaging.operation.type");
5464
private static final AttributeKey<String> MESSAGING_SYSTEM =
5565
AttributeKey.stringKey("messaging.system");
5666

@@ -61,26 +71,51 @@ public final class MessagingAttributesExtractor<REQUEST, RESPONSE>
6171
* with default configuration.
6272
*/
6373
public static <REQUEST, RESPONSE> AttributesExtractor<REQUEST, RESPONSE> create(
64-
MessagingAttributesGetter<REQUEST, RESPONSE> getter, MessageOperation operation) {
74+
MessagingAttributesGetter<REQUEST, RESPONSE> getter, @Nullable MessageOperation operation) {
6575
return builder(getter, operation).build();
6676
}
6777

78+
/**
79+
* Creates the messaging attributes extractor with a system-specific v1.43 operation name.
80+
*
81+
* <p>The {@code operationName} is emitted as {@code messaging.operation.name}. The legacy {@code
82+
* messaging.operation} value remains derived from {@code operation}.
83+
*/
84+
public static <REQUEST, RESPONSE> AttributesExtractor<REQUEST, RESPONSE> create(
85+
MessagingAttributesGetter<REQUEST, RESPONSE> getter,
86+
MessageOperation operation,
87+
String operationName) {
88+
return builder(getter, operation, operationName).build();
89+
}
90+
6891
/**
6992
* Returns a new {@link MessagingAttributesExtractorBuilder} for the given {@link MessageOperation
7093
* operation} that can be used to configure the messaging attributes extractor.
7194
*/
7295
public static <REQUEST, RESPONSE> MessagingAttributesExtractorBuilder<REQUEST, RESPONSE> builder(
73-
MessagingAttributesGetter<REQUEST, RESPONSE> getter, MessageOperation operation) {
96+
MessagingAttributesGetter<REQUEST, RESPONSE> getter, @Nullable MessageOperation operation) {
7497
return new MessagingAttributesExtractorBuilder<>(getter, operation);
7598
}
7699

100+
/**
101+
* Returns a new {@link MessagingAttributesExtractorBuilder} with a system-specific v1.43
102+
* operation name.
103+
*/
104+
public static <REQUEST, RESPONSE> MessagingAttributesExtractorBuilder<REQUEST, RESPONSE> builder(
105+
MessagingAttributesGetter<REQUEST, RESPONSE> getter,
106+
MessageOperation operation,
107+
String operationName) {
108+
return new MessagingAttributesExtractorBuilder<>(
109+
getter, MessagingOperation.create(operation, operationName));
110+
}
111+
77112
private final MessagingAttributesGetter<REQUEST, RESPONSE> getter;
78-
private final MessageOperation operation;
113+
@Nullable private final MessagingOperation operation;
79114
private final List<String> capturedHeaders;
80115

81116
MessagingAttributesExtractor(
82117
MessagingAttributesGetter<REQUEST, RESPONSE> getter,
83-
MessageOperation operation,
118+
@Nullable MessagingOperation operation,
84119
List<String> capturedHeaders) {
85120
this.getter = getter;
86121
this.operation = operation;
@@ -93,7 +128,12 @@ public void onStart(AttributesBuilder attributes, Context parentContext, REQUEST
93128
boolean isTemporaryDestination = getter.isTemporaryDestination(request);
94129
if (isTemporaryDestination) {
95130
attributes.put(MESSAGING_DESTINATION_TEMPORARY, true);
96-
attributes.put(MESSAGING_DESTINATION_NAME, TEMP_DESTINATION_NAME);
131+
if (emitStableMessagingSemconv()) {
132+
attributes.put(MESSAGING_DESTINATION_NAME, getter.getDestination(request));
133+
attributes.put(MESSAGING_DESTINATION_TEMPLATE, getter.getDestinationTemplate(request));
134+
} else {
135+
attributes.put(MESSAGING_DESTINATION_NAME, TEMP_DESTINATION_NAME);
136+
}
97137
} else {
98138
attributes.put(MESSAGING_DESTINATION_NAME, getter.getDestination(request));
99139
attributes.put(MESSAGING_DESTINATION_TEMPLATE, getter.getDestinationTemplate(request));
@@ -106,9 +146,20 @@ public void onStart(AttributesBuilder attributes, Context parentContext, REQUEST
106146
attributes.put(MESSAGING_MESSAGE_CONVERSATION_ID, getter.getConversationId(request));
107147
attributes.put(MESSAGING_MESSAGE_BODY_SIZE, getter.getMessageBodySize(request));
108148
attributes.put(MESSAGING_MESSAGE_ENVELOPE_SIZE, getter.getMessageEnvelopeSize(request));
109-
attributes.put(MESSAGING_CLIENT_ID, getter.getClientId(request));
149+
if (emitOldMessagingSemconv()) {
150+
attributes.put(MESSAGING_CLIENT_ID_OLD, getter.getClientId(request));
151+
}
152+
if (emitStableMessagingSemconv()) {
153+
attributes.put(MESSAGING_CLIENT_ID, getter.getClientId(request));
154+
}
110155
if (operation != null) {
111-
attributes.put(MESSAGING_OPERATION, operation.operationName());
156+
if (emitOldMessagingSemconv()) {
157+
attributes.put(MESSAGING_OPERATION, operation.operation().operationName());
158+
}
159+
if (emitStableMessagingSemconv()) {
160+
attributes.put(MESSAGING_OPERATION_NAME, operation.name());
161+
attributes.put(MESSAGING_OPERATION_TYPE, operation.type());
162+
}
112163
}
113164
}
114165

@@ -121,6 +172,9 @@ public void onEnd(
121172
@Nullable Throwable error) {
122173
attributes.put(MESSAGING_MESSAGE_ID, getter.getMessageId(request, response));
123174
attributes.put(MESSAGING_BATCH_MESSAGE_COUNT, getter.getBatchMessageCount(request, response));
175+
if (emitStableMessagingSemconv() && error != null) {
176+
attributes.put(ERROR_TYPE, error.getClass().getName());
177+
}
124178

125179
for (String name : capturedHeaders) {
126180
List<String> values = getter.getMessageHeader(request, name);
@@ -141,7 +195,7 @@ public SpanKey internalGetSpanKey() {
141195
return null;
142196
}
143197

144-
switch (operation) {
198+
switch (operation.operation()) {
145199
case PUBLISH:
146200
return SpanKey.PRODUCER;
147201
case RECEIVE:

instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingAttributesExtractorBuilder.java

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,16 +12,22 @@
1212
import java.util.ArrayList;
1313
import java.util.Collection;
1414
import java.util.List;
15+
import javax.annotation.Nullable;
1516

1617
/** A builder of {@link MessagingAttributesExtractor}. */
1718
public final class MessagingAttributesExtractorBuilder<REQUEST, RESPONSE> {
1819

1920
final MessagingAttributesGetter<REQUEST, RESPONSE> getter;
20-
final MessageOperation operation;
21+
@Nullable final MessagingOperation operation;
2122
List<String> capturedHeaders = emptyList();
2223

2324
MessagingAttributesExtractorBuilder(
24-
MessagingAttributesGetter<REQUEST, RESPONSE> getter, MessageOperation operation) {
25+
MessagingAttributesGetter<REQUEST, RESPONSE> getter, @Nullable MessageOperation operation) {
26+
this(getter, MessagingOperation.createNullable(operation));
27+
}
28+
29+
MessagingAttributesExtractorBuilder(
30+
MessagingAttributesGetter<REQUEST, RESPONSE> getter, @Nullable MessagingOperation operation) {
2531
this.getter = getter;
2632
this.operation = operation;
2733
}

instrumentation-api-incubator/src/main/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingConsumerMetrics.java

Lines changed: 73 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@
55

66
package io.opentelemetry.instrumentation.api.incubator.semconv.messaging;
77

8+
import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitOldMessagingSemconv;
9+
import static io.opentelemetry.instrumentation.api.internal.SemconvStability.emitStableMessagingSemconv;
810
import static java.util.concurrent.TimeUnit.SECONDS;
911
import static java.util.logging.Level.FINE;
1012

@@ -23,10 +25,11 @@
2325
import io.opentelemetry.instrumentation.api.instrumenter.OperationMetrics;
2426
import io.opentelemetry.instrumentation.api.internal.OperationMetricsUtil;
2527
import java.util.logging.Logger;
28+
import javax.annotation.Nullable;
2629

2730
/**
2831
* {@link OperationListener} which keeps track of <a
29-
* href="https://github.com/open-telemetry/semantic-conventions/blob/v1.26.0/docs/messaging/messaging-metrics.md#consumer-metrics">Consumer
32+
* href="https://github.com/open-telemetry/semantic-conventions/blob/v1.43.0/docs/messaging/messaging-metrics.md#consumer-metrics">consumer
3033
* metrics</a>.
3134
*/
3235
public final class MessagingConsumerMetrics implements OperationListener {
@@ -39,26 +42,17 @@ public final class MessagingConsumerMetrics implements OperationListener {
3942
ContextKey.named("messaging-consumer-metrics-state");
4043
private static final Logger logger = Logger.getLogger(MessagingConsumerMetrics.class.getName());
4144

42-
private final DoubleHistogram receiveDurationHistogram;
43-
private final LongCounter receiveMessageCount;
45+
@Nullable private final DoubleHistogram receiveDurationHistogram;
46+
@Nullable private final LongCounter receiveMessageCount;
47+
@Nullable private final DoubleHistogram clientOperationDurationHistogram;
48+
@Nullable private final LongCounter consumedMessagesCounter;
4449

4550
private MessagingConsumerMetrics(Meter meter) {
46-
DoubleHistogramBuilder durationBuilder =
47-
meter
48-
.histogramBuilder("messaging.receive.duration")
49-
.setDescription("Measures the duration of receive operation.")
50-
.setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
51-
.setUnit("s");
52-
MessagingMetricsAdvice.applyReceiveDurationAdvice(durationBuilder);
53-
receiveDurationHistogram = durationBuilder.build();
54-
55-
LongCounterBuilder longCounterBuilder =
56-
meter
57-
.counterBuilder("messaging.receive.messages")
58-
.setDescription("Measures the number of received messages.")
59-
.setUnit("{message}");
60-
MessagingMetricsAdvice.applyReceiveMessagesAdvice(longCounterBuilder);
61-
receiveMessageCount = longCounterBuilder.build();
51+
receiveDurationHistogram = emitOldMessagingSemconv() ? buildReceiveDuration(meter) : null;
52+
receiveMessageCount = emitOldMessagingSemconv() ? buildReceiveMessages(meter) : null;
53+
clientOperationDurationHistogram =
54+
emitStableMessagingSemconv() ? buildClientOperationDuration(meter) : null;
55+
consumedMessagesCounter = emitStableMessagingSemconv() ? buildConsumedMessages(meter) : null;
6256
}
6357

6458
public static OperationMetrics get() {
@@ -85,11 +79,25 @@ public void onEnd(Context context, Attributes endAttributes, long endNanos) {
8579
}
8680

8781
Attributes attributes = state.startAttributes().toBuilder().putAll(endAttributes).build();
88-
receiveDurationHistogram.record(
89-
(endNanos - state.startTimeNanos()) / NANOS_PER_S, attributes, context);
82+
double duration = (endNanos - state.startTimeNanos()) / NANOS_PER_S;
83+
if (receiveDurationHistogram != null) {
84+
receiveDurationHistogram.record(duration, attributes, context);
85+
}
86+
Attributes filteredAttributes =
87+
clientOperationDurationHistogram != null || consumedMessagesCounter != null
88+
? MessagingMetricsAdvice.filterAttributes(attributes)
89+
: attributes;
90+
if (clientOperationDurationHistogram != null) {
91+
clientOperationDurationHistogram.record(duration, filteredAttributes, context);
92+
}
9093

9194
long receiveMessagesCount = getReceiveMessagesCount(state.startAttributes(), endAttributes);
92-
receiveMessageCount.add(receiveMessagesCount, attributes, context);
95+
if (receiveMessageCount != null) {
96+
receiveMessageCount.add(receiveMessagesCount, attributes, context);
97+
}
98+
if (consumedMessagesCounter != null) {
99+
consumedMessagesCounter.add(receiveMessagesCount, filteredAttributes, context);
100+
}
93101
}
94102

95103
private static long getReceiveMessagesCount(Attributes... attributesList) {
@@ -102,6 +110,49 @@ private static long getReceiveMessagesCount(Attributes... attributesList) {
102110
return 1;
103111
}
104112

113+
private static DoubleHistogram buildReceiveDuration(Meter meter) {
114+
DoubleHistogramBuilder builder =
115+
meter
116+
.histogramBuilder("messaging.receive.duration")
117+
.setDescription("Measures the duration of receive operation.")
118+
.setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
119+
.setUnit("s");
120+
MessagingMetricsAdvice.applyOldDurationAdvice(builder);
121+
return builder.build();
122+
}
123+
124+
private static LongCounter buildReceiveMessages(Meter meter) {
125+
LongCounterBuilder builder =
126+
meter
127+
.counterBuilder("messaging.receive.messages")
128+
.setDescription("Measures the number of received messages.")
129+
.setUnit("{message}");
130+
MessagingMetricsAdvice.applyOldMessagesAdvice(builder);
131+
return builder.build();
132+
}
133+
134+
private static DoubleHistogram buildClientOperationDuration(Meter meter) {
135+
DoubleHistogramBuilder builder =
136+
meter
137+
.histogramBuilder("messaging.client.operation.duration")
138+
.setDescription(
139+
"Duration of messaging operation initiated by a producer or consumer client.")
140+
.setExplicitBucketBoundariesAdvice(MessagingMetricsAdvice.DURATION_SECONDS_BUCKETS)
141+
.setUnit("s");
142+
MessagingMetricsAdvice.applyClientOperationDurationAdvice(builder);
143+
return builder.build();
144+
}
145+
146+
private static LongCounter buildConsumedMessages(Meter meter) {
147+
LongCounterBuilder builder =
148+
meter
149+
.counterBuilder("messaging.client.consumed.messages")
150+
.setDescription("Number of messages that were delivered to the application.")
151+
.setUnit("{message}");
152+
MessagingMetricsAdvice.applyConsumedMessagesAdvice(builder);
153+
return builder.build();
154+
}
155+
105156
@AutoValue
106157
abstract static class State {
107158

0 commit comments

Comments
 (0)