Skip to content

Commit 7a01fd3

Browse files
committed
Fix Kafka messaging test compatibility
1 parent 8b235e0 commit 7a01fd3

7 files changed

Lines changed: 19 additions & 18 deletions

File tree

instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanKindExtractorTest.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,10 @@ class MessagingSpanKindExtractorTest {
1919

2020
@ParameterizedTest
2121
@MethodSource("spanKinds")
22-
void extractsSpanKind(MessageOperation operation, SpanKind oldKind, SpanKind stableKind) {
23-
SpanKind spanKind = MessagingSpanKindExtractor.create(operation).extract(new Object());
22+
void extractsSpanKind(MessageOperation operation, SpanKind oldKind, SpanKind kind) {
23+
SpanKind actualKind = MessagingSpanKindExtractor.create(operation).extract(new Object());
2424

25-
assertThat(spanKind).isEqualTo(emitStableMessagingSemconv() ? stableKind : oldKind);
25+
assertThat(actualKind).isEqualTo(emitStableMessagingSemconv() ? kind : oldKind);
2626
}
2727

2828
private static Stream<Arguments> spanKinds() {

instrumentation-api-incubator/src/test/java/io/opentelemetry/instrumentation/api/incubator/semconv/messaging/MessagingSpanNameExtractorTest.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@ void shouldExtractSpanName(
3434
MessageOperation operation,
3535
String operationName,
3636
String oldSpanName,
37-
String stableSpanName) {
37+
String spanName) {
3838
// given
3939
Message message = new Message();
4040

@@ -60,10 +60,10 @@ void shouldExtractSpanName(
6060
MessagingSpanNameExtractor.create(getter, operation, operationName);
6161

6262
// when
63-
String spanName = underTest.extract(message);
63+
String actualSpanName = underTest.extract(message);
6464

6565
// then
66-
assertThat(spanName).isEqualTo(emitStableMessagingSemconv() ? stableSpanName : oldSpanName);
66+
assertThat(actualSpanName).isEqualTo(emitStableMessagingSemconv() ? spanName : oldSpanName);
6767
}
6868

6969
static Stream<Arguments> spanNameParams() {

instrumentation/kafka/kafka-clients/kafka-clients-0.11/testing/src/main/java/io/opentelemetry/instrumentation/kafkaclients/common/v0_11/internal/KafkaClientBaseTest.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,6 @@
1414
import static io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions.equalTo;
1515
import static io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions.satisfies;
1616
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_BATCH_MESSAGE_COUNT;
17-
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_CLIENT_ID;
1817
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_DESTINATION_NAME;
1918
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_DESTINATION_PARTITION_ID;
2019
import static io.opentelemetry.semconv.incubating.MessagingIncubatingAttributes.MESSAGING_KAFKA_CONSUMER_GROUP;
@@ -77,8 +76,10 @@ public abstract class KafkaClientBaseTest {
7776
@RegisterExtension final AutoCleanupExtension cleanup = AutoCleanupExtension.create();
7877

7978
protected static final String SHARED_TOPIC = "shared.topic";
80-
private static final AttributeKey<String> MESSAGING_CLIENT_ID_OLD =
79+
protected static final AttributeKey<String> MESSAGING_CLIENT_ID_OLD =
8180
AttributeKey.stringKey("messaging.client_id");
81+
private static final AttributeKey<String> MESSAGING_CLIENT_ID =
82+
AttributeKey.stringKey("messaging.client.id");
8283

8384
private KafkaContainer kafka;
8485
protected Producer<Integer, String> producer;

instrumentation/kafka/kafka-clients/kafka-clients-2.6/library/src/test/java/io/opentelemetry/instrumentation/kafkaclients/v2_6/AbstractInterceptorsTest.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -164,7 +164,7 @@ private static List<AttributeAssertion> publishAttributes(boolean experimental)
164164
equalTo(MESSAGING_SYSTEM, "kafka"),
165165
equalTo(MESSAGING_DESTINATION_NAME, SHARED_TOPIC),
166166
equalTo(MESSAGING_OPERATION, "publish"),
167-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("producer")),
167+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("producer")),
168168
satisfies(
169169
stringKey("messaging.kafka.bootstrap.servers"),
170170
val -> {
@@ -181,7 +181,7 @@ private static List<AttributeAssertion> receiveAttributes() {
181181
equalTo(MESSAGING_DESTINATION_NAME, SHARED_TOPIC),
182182
equalTo(MESSAGING_OPERATION, "receive"),
183183
equalTo(MESSAGING_KAFKA_CONSUMER_GROUP, "test"),
184-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("consumer")),
184+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("consumer")),
185185
equalTo(MESSAGING_BATCH_MESSAGE_COUNT, 1));
186186
}
187187

@@ -195,7 +195,7 @@ private static List<AttributeAssertion> processAttributes(boolean experimental)
195195
satisfies(MESSAGING_DESTINATION_PARTITION_ID, AbstractStringAssert::isNotEmpty),
196196
satisfies(MESSAGING_KAFKA_MESSAGE_OFFSET, AbstractLongAssert::isNotNegative),
197197
equalTo(MESSAGING_KAFKA_CONSUMER_GROUP, "test"),
198-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("consumer")),
198+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("consumer")),
199199
satisfies(
200200
longKey("kafka.record.queue_time_ms"),
201201
val -> {

instrumentation/kafka/kafka-clients/kafka-clients-2.6/library/src/test/java/io/opentelemetry/instrumentation/kafkaclients/v2_6/InterceptorsSuppressReceiveSpansTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -51,7 +51,7 @@ void assertTraces() {
5151
equalTo(MESSAGING_SYSTEM, "kafka"),
5252
equalTo(MESSAGING_DESTINATION_NAME, SHARED_TOPIC),
5353
equalTo(MESSAGING_OPERATION, "publish"),
54-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("producer"))),
54+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("producer"))),
5555
span ->
5656
span.hasName(SHARED_TOPIC + " process")
5757
.hasKind(SpanKind.CONSUMER)
@@ -67,7 +67,7 @@ void assertTraces() {
6767
satisfies(
6868
MESSAGING_KAFKA_MESSAGE_OFFSET, AbstractLongAssert::isNotNegative),
6969
equalTo(MESSAGING_KAFKA_CONSUMER_GROUP, "test"),
70-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("consumer")),
70+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("consumer")),
7171
equalTo(stringKey("test-baggage-key-1"), "test-baggage-value-1"),
7272
equalTo(stringKey("test-baggage-key-2"), "test-baggage-value-2")),
7373
span ->

instrumentation/kafka/kafka-clients/kafka-clients-2.6/library/src/test/java/io/opentelemetry/instrumentation/kafkaclients/v2_6/WrapperSuppressReceiveSpansTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,7 @@ static List<AttributeAssertion> sendAttributes(boolean testHeaders, boolean test
7171
equalTo(MESSAGING_SYSTEM, "kafka"),
7272
equalTo(MESSAGING_DESTINATION_NAME, SHARED_TOPIC),
7373
equalTo(MESSAGING_OPERATION, "publish"),
74-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("producer")),
74+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("producer")),
7575
satisfies(MESSAGING_DESTINATION_PARTITION_ID, AbstractStringAssert::isNotEmpty),
7676
satisfies(MESSAGING_KAFKA_MESSAGE_OFFSET, AbstractLongAssert::isNotNegative)));
7777
if (testHeaders) {
@@ -98,7 +98,7 @@ static List<AttributeAssertion> processAttributes(
9898
satisfies(MESSAGING_DESTINATION_PARTITION_ID, AbstractStringAssert::isNotEmpty),
9999
satisfies(MESSAGING_KAFKA_MESSAGE_OFFSET, AbstractLongAssert::isNotNegative),
100100
equalTo(MESSAGING_KAFKA_CONSUMER_GROUP, "test"),
101-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("consumer"))));
101+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("consumer"))));
102102
if (testHeaders) {
103103
assertions.add(equalTo(headerAttributeKey("Test-Message-Header"), singletonList("test")));
104104
}

instrumentation/kafka/kafka-clients/kafka-clients-2.6/library/src/test/java/io/opentelemetry/instrumentation/kafkaclients/v2_6/WrapperTest.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -106,7 +106,7 @@ protected static List<AttributeAssertion> sendAttributes(
106106
equalTo(MESSAGING_SYSTEM, "kafka"),
107107
equalTo(MESSAGING_DESTINATION_NAME, SHARED_TOPIC),
108108
equalTo(MESSAGING_OPERATION, "publish"),
109-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("producer")),
109+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("producer")),
110110
satisfies(MESSAGING_DESTINATION_PARTITION_ID, AbstractStringAssert::isNotEmpty),
111111
satisfies(MESSAGING_KAFKA_MESSAGE_OFFSET, AbstractLongAssert::isNotNegative)));
112112
if (testHeaders) {
@@ -135,7 +135,7 @@ private static List<AttributeAssertion> processAttributes(
135135
satisfies(MESSAGING_DESTINATION_PARTITION_ID, AbstractStringAssert::isNotEmpty),
136136
satisfies(MESSAGING_KAFKA_MESSAGE_OFFSET, AbstractLongAssert::isNotNegative),
137137
equalTo(MESSAGING_KAFKA_CONSUMER_GROUP, "test"),
138-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("consumer"))));
138+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("consumer"))));
139139
if (testHeaders) {
140140
assertions.add(
141141
equalTo(
@@ -156,7 +156,7 @@ protected static List<AttributeAssertion> receiveAttributes(boolean testHeaders)
156156
equalTo(MESSAGING_DESTINATION_NAME, SHARED_TOPIC),
157157
equalTo(MESSAGING_OPERATION, "receive"),
158158
equalTo(MESSAGING_KAFKA_CONSUMER_GROUP, "test"),
159-
satisfies(MESSAGING_CLIENT_ID, val -> val.startsWith("consumer")),
159+
satisfies(MESSAGING_CLIENT_ID_OLD, val -> val.startsWith("consumer")),
160160
equalTo(MESSAGING_BATCH_MESSAGE_COUNT, 1)));
161161
if (testHeaders) {
162162
assertions.add(

0 commit comments

Comments
 (0)