Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,13 @@

import datadog.trace.bootstrap.ContextStore;
import datadog.trace.instrumentation.kafka_common.MetadataState;
import java.util.Optional;
import org.apache.kafka.clients.Metadata;

public class KafkaConsumerInstrumentationHelper {
public static String extractGroup(KafkaConsumerInfo kafkaConsumerInfo) {
if (kafkaConsumerInfo != null) {
return kafkaConsumerInfo.getConsumerGroup().get();
return kafkaConsumerInfo.getConsumerGroup().orElse(null);
}
return null;
}
Expand All @@ -16,16 +17,16 @@ public static String extractClusterId(
KafkaConsumerInfo kafkaConsumerInfo,
ContextStore<Metadata, MetadataState> metadataContextStore) {
if (kafkaConsumerInfo != null) {
Metadata metadata = kafkaConsumerInfo.getmetadata().get();
if (metadata != null) {
MetadataState state = metadataContextStore.get(metadata);
Optional<Metadata> metadata = kafkaConsumerInfo.getmetadata();
if (metadata.isPresent()) {
MetadataState state = metadataContextStore.get(metadata.get());
return state != null ? state.clusterId : null;
}
}
return null;
}

public static String extractBootstrapServers(KafkaConsumerInfo kafkaConsumerInfo) {
return kafkaConsumerInfo == null ? null : kafkaConsumerInfo.getBootstrapServers().get();
return kafkaConsumerInfo == null ? null : kafkaConsumerInfo.getBootstrapServers().orElse(null);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
package datadog.trace.instrumentation.kafka_clients38;

import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;

import datadog.trace.bootstrap.ContextStore;
import datadog.trace.instrumentation.kafka_common.MetadataState;
import org.apache.kafka.clients.Metadata;
import org.junit.jupiter.api.Test;

class KafkaConsumerInstrumentationHelperTest {

@SuppressWarnings("unchecked")
private final ContextStore<Metadata, MetadataState> metadataContextStore =
mock(ContextStore.class);

@Test
void extractGroupReturnsNullForNullKafkaConsumerInfo() {
assertNull(KafkaConsumerInstrumentationHelper.extractGroup(null));
}

@Test
void extractGroupReturnsNullWhenConsumerGroupIsNull() {
KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo(null, null, "localhost:9092");
assertNull(KafkaConsumerInstrumentationHelper.extractGroup(kafkaConsumerInfo));
}

@Test
void extractGroupReturnsConsumerGroupWhenPresent() {
KafkaConsumerInfo kafkaConsumerInfo =
new KafkaConsumerInfo("test-group", null, "localhost:9092");
assertEquals("test-group", KafkaConsumerInstrumentationHelper.extractGroup(kafkaConsumerInfo));
}

@Test
void extractBootstrapServersReturnsNullForNullKafkaConsumerInfo() {
assertNull(KafkaConsumerInstrumentationHelper.extractBootstrapServers(null));
}

@Test
void extractBootstrapServersReturnsNullWhenBootstrapServersIsNull() {
KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo("test-group", null, null);
assertNull(KafkaConsumerInstrumentationHelper.extractBootstrapServers(kafkaConsumerInfo));
}

@Test
void extractBootstrapServersReturnsValueWhenPresent() {
KafkaConsumerInfo kafkaConsumerInfo =
new KafkaConsumerInfo("test-group", null, "localhost:9092");
assertEquals(
"localhost:9092",
KafkaConsumerInstrumentationHelper.extractBootstrapServers(kafkaConsumerInfo));
}

@Test
void extractClusterIdReturnsNullForNullKafkaConsumerInfo() {
assertNull(KafkaConsumerInstrumentationHelper.extractClusterId(null, metadataContextStore));
}

@Test
void extractClusterIdReturnsNullWhenMetadataIsNull() {
KafkaConsumerInfo kafkaConsumerInfo = new KafkaConsumerInfo("test-group", "localhost:9092");
assertNull(
KafkaConsumerInstrumentationHelper.extractClusterId(
kafkaConsumerInfo, metadataContextStore));
}

@Test
void extractClusterIdReturnsNullWhenNoStateForMetadata() {
Metadata metadata = mock(Metadata.class);
KafkaConsumerInfo kafkaConsumerInfo =
new KafkaConsumerInfo("test-group", metadata, "localhost:9092");
when(metadataContextStore.get(metadata)).thenReturn(null);
assertNull(
KafkaConsumerInstrumentationHelper.extractClusterId(
kafkaConsumerInfo, metadataContextStore));
}

@Test
void extractClusterIdReturnsClusterIdWhenStatePresent() {
Metadata metadata = mock(Metadata.class);
KafkaConsumerInfo kafkaConsumerInfo =
new KafkaConsumerInfo("test-group", metadata, "localhost:9092");
MetadataState state = new MetadataState();
state.clusterId = "cluster-1";
when(metadataContextStore.get(metadata)).thenReturn(state);
assertEquals(
"cluster-1",
KafkaConsumerInstrumentationHelper.extractClusterId(
kafkaConsumerInfo, metadataContextStore));
}
}
Loading