Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 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
52 changes: 46 additions & 6 deletions core/src/main/scala/kafka/server/KafkaApis.scala
Original file line number Diff line number Diff line change
Expand Up @@ -3747,16 +3747,56 @@ class KafkaApis(val requestChannel: RequestChannel,
CompletableFuture.completedFuture[Unit](())
}

// Just a place holder for now.
def handleGetTelemetrySubscriptionsRequest(request: RequestChannel.Request): Unit = {
requestHelper.sendMaybeThrottle(request, request.body[GetTelemetrySubscriptionsRequest].getErrorResponse(Errors.UNSUPPORTED_VERSION.exception))
CompletableFuture.completedFuture[Unit](())
val subscriptionRequest = request.body[GetTelemetrySubscriptionsRequest]

clientMetricsManager match {
case Some(metricsManager) =>
try {
if (metricsManager.isTelemetryReceiverConfigured) {
requestHelper.sendResponseMaybeThrottle(request, requestThrottleMs =>
Comment thread
apoorvmittal10 marked this conversation as resolved.
Outdated
metricsManager.processGetTelemetrySubscriptionRequest(subscriptionRequest, request.context))
} else {
info("Received get telemetry client request for metrics receiver, but no metrics receiver configured")
requestHelper.sendResponseMaybeThrottle(request, requestThrottleMs =>
subscriptionRequest.getErrorResponse(requestThrottleMs, Errors.UNSUPPORTED_VERSION.exception))
}
} catch {
case _: Exception =>
requestHelper.sendResponseMaybeThrottle(request, requestThrottleMs =>
subscriptionRequest.getErrorResponse(requestThrottleMs, Errors.INVALID_REQUEST.exception))
}
case None =>
info("Received get telemetry client request for zookeeper based cluster")
requestHelper.sendResponseMaybeThrottle(request, requestThrottleMs =>
subscriptionRequest.getErrorResponse(requestThrottleMs, Errors.UNSUPPORTED_VERSION.exception))
}
}

// Just a place holder for now.
def handlePushTelemetryRequest(request: RequestChannel.Request): Unit = {
requestHelper.sendMaybeThrottle(request, request.body[PushTelemetryRequest].getErrorResponse(Errors.UNSUPPORTED_VERSION.exception))
CompletableFuture.completedFuture[Unit](())
val pushTelemetryRequest = request.body[PushTelemetryRequest]

clientMetricsManager match {
case Some(metricsManager) =>
try {
if (metricsManager.isTelemetryReceiverConfigured) {
requestHelper.sendResponseMaybeThrottle(request, requestThrottleMs =>
metricsManager.processPushTelemetryRequest(pushTelemetryRequest, request.context))
} else {
info("Received push telemetry client request, but no metrics receiver configured")
requestHelper.sendResponseMaybeThrottle(request, requestThrottleMs =>
pushTelemetryRequest.getErrorResponse(requestThrottleMs, Errors.UNSUPPORTED_VERSION.exception))
}
} catch {
case _: Exception =>
requestHelper.sendResponseMaybeThrottle(request, requestThrottleMs =>
pushTelemetryRequest.getErrorResponse(requestThrottleMs, Errors.INVALID_REQUEST.exception))
}
case None =>
info("Received push telemetry client request for zookeeper based cluster")
requestHelper.sendResponseMaybeThrottle(request, requestThrottleMs =>
pushTelemetryRequest.getErrorResponse(requestThrottleMs, Errors.UNSUPPORTED_VERSION.exception))
}
}

private def updateRecordConversionStats(request: RequestChannel.Request,
Expand Down
93 changes: 81 additions & 12 deletions core/src/test/scala/unit/kafka/server/KafkaApisTest.scala
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ import org.apache.kafka.common.message.CreatePartitionsRequestData.CreatePartiti
import org.apache.kafka.common.message.CreateTopicsResponseData.CreatableTopicResult
import org.apache.kafka.common.message.OffsetDeleteResponseData.{OffsetDeleteResponsePartition, OffsetDeleteResponsePartitionCollection, OffsetDeleteResponseTopic, OffsetDeleteResponseTopicCollection}
import org.apache.kafka.coordinator.group.GroupCoordinator
import org.apache.kafka.server.ClientMetricsManager
import org.apache.kafka.server.common.{Features, MetadataVersion}
import org.apache.kafka.server.common.MetadataVersion.{IBP_0_10_2_IV0, IBP_2_2_IV1}
import org.apache.kafka.server.metrics.ClientMetricsTestUtils
Expand Down Expand Up @@ -129,6 +130,7 @@ class KafkaApisTest {
private val quotas = QuotaManagers(clientQuotaManager, clientQuotaManager, clientRequestQuotaManager,
clientControllerQuotaManager, replicaQuotaManager, replicaQuotaManager, replicaQuotaManager, None)
private val fetchManager: FetchManager = mock(classOf[FetchManager])
private val clientMetricsManager: ClientMetricsManager = mock(classOf[ClientMetricsManager])
private val brokerTopicStats = new BrokerTopicStats
private val clusterId = "clusterId"
private val time = new MockTime
Expand Down Expand Up @@ -196,6 +198,8 @@ class KafkaApisTest {
false,
() => new Features(MetadataVersion.latest(), Collections.emptyMap[String, java.lang.Short], 0, raftSupport))

val clientMetricsManagerOpt = if (raftSupport) Some(clientMetricsManager) else None

new KafkaApis(
requestChannel = requestChannel,
metadataSupport = metadataSupport,
Expand All @@ -216,7 +220,7 @@ class KafkaApisTest {
time = time,
tokenManager = null,
apiVersionManager = apiVersionManager,
clientMetricsManager = null)
clientMetricsManager = clientMetricsManagerOpt)
}

@Test
Expand Down Expand Up @@ -6760,18 +6764,54 @@ class KafkaApisTest {
}

@Test
def testGetTelemetrySubscriptionsUnsupportedVersionForKRaftClusters(): Unit = {
val data = new GetTelemetrySubscriptionsRequestData()
def testGetTelemetrySubscriptions(): Unit = {
val request = buildRequest(new GetTelemetrySubscriptionsRequest.Builder(
new GetTelemetrySubscriptionsRequestData(), true).build())

when(clientMetricsManager.isTelemetryReceiverConfigured).thenReturn(true)
when(clientMetricsManager.processGetTelemetrySubscriptionRequest(any[GetTelemetrySubscriptionsRequest](),
any[RequestContext]())).thenReturn(new GetTelemetrySubscriptionsResponse(
new GetTelemetrySubscriptionsResponseData()))

metadataCache = MetadataCache.kRaftMetadataCache(brokerId)
createKafkaApis(raftSupport = true).handle(request, RequestLocal.NoCaching)

val response = verifyNoThrottling[GetTelemetrySubscriptionsResponse](request)

val request = buildRequest(new GetTelemetrySubscriptionsRequest.Builder(data, true).build())
val errorCode = Errors.UNSUPPORTED_VERSION.code
val expectedResponse = new GetTelemetrySubscriptionsResponseData()
expectedResponse.setErrorCode(errorCode)
assertEquals(expectedResponse, response.data)
}

@Test
def testGetTelemetrySubscriptionsNoMetricsPlugin(): Unit = {
val request = buildRequest(new GetTelemetrySubscriptionsRequest.Builder(
new GetTelemetrySubscriptionsRequestData(), true).build())

when(clientMetricsManager.isTelemetryReceiverConfigured).thenReturn(false)
metadataCache = MetadataCache.kRaftMetadataCache(brokerId)
createKafkaApis(raftSupport = true).handle(request, RequestLocal.NoCaching)

val response = verifyNoThrottling[GetTelemetrySubscriptionsResponse](request)

val expectedResponse = new GetTelemetrySubscriptionsResponseData().setErrorCode(Errors.UNSUPPORTED_VERSION.code)
assertEquals(expectedResponse, response.data)
}

@Test
def testGetTelemetrySubscriptionsWithException(): Unit = {
val request = buildRequest(new GetTelemetrySubscriptionsRequest.Builder(
new GetTelemetrySubscriptionsRequestData(), true).build())

when(clientMetricsManager.isTelemetryReceiverConfigured).thenReturn(true)
when(clientMetricsManager.processGetTelemetrySubscriptionRequest(any[GetTelemetrySubscriptionsRequest](),
any[RequestContext]())).thenThrow(new RuntimeException("test"))

metadataCache = MetadataCache.kRaftMetadataCache(brokerId)
createKafkaApis(raftSupport = true).handle(request, RequestLocal.NoCaching)

val response = verifyNoThrottling[GetTelemetrySubscriptionsResponse](request)

val expectedResponse = new GetTelemetrySubscriptionsResponseData().setErrorCode(Errors.INVALID_REQUEST.code)
assertEquals(expectedResponse, response.data)
}

Expand All @@ -6787,18 +6827,47 @@ class KafkaApisTest {
}

@Test
def testPushTelemetryUnsupportedVersionForKRaftClusters(): Unit = {
val data = new PushTelemetryRequestData()
def testPushTelemetry(): Unit = {
val request = buildRequest(new PushTelemetryRequest.Builder(new PushTelemetryRequestData(), true).build())

val request = buildRequest(new PushTelemetryRequest.Builder(data, true).build())
val errorCode = Errors.UNSUPPORTED_VERSION.code
val expectedResponse = new PushTelemetryResponseData()
expectedResponse.setErrorCode(errorCode)
when(clientMetricsManager.isTelemetryReceiverConfigured).thenReturn(true)
when(clientMetricsManager.processPushTelemetryRequest(any[PushTelemetryRequest](), any[RequestContext]()))
.thenReturn(new PushTelemetryResponse(new PushTelemetryResponseData()))

metadataCache = MetadataCache.kRaftMetadataCache(brokerId)
createKafkaApis(raftSupport = true).handle(request, RequestLocal.NoCaching)
val response = verifyNoThrottling[PushTelemetryResponse](request)

val expectedResponse = new PushTelemetryResponseData().setErrorCode(Errors.NONE.code)
assertEquals(expectedResponse, response.data)
}

@Test
def testPushTelemetryNoMetricsPlugin(): Unit = {
val request = buildRequest(new PushTelemetryRequest.Builder(new PushTelemetryRequestData(), true).build())

when(clientMetricsManager.isTelemetryReceiverConfigured).thenReturn(false)
metadataCache = MetadataCache.kRaftMetadataCache(brokerId)
createKafkaApis(raftSupport = true).handle(request, RequestLocal.NoCaching)
val response = verifyNoThrottling[PushTelemetryResponse](request)

val expectedResponse = new PushTelemetryResponseData().setErrorCode(Errors.UNSUPPORTED_VERSION.code)
assertEquals(expectedResponse, response.data)
}

@Test
def testPushTelemetryWithException(): Unit = {
val request = buildRequest(new PushTelemetryRequest.Builder(new PushTelemetryRequestData(), true).build())

when(clientMetricsManager.isTelemetryReceiverConfigured()).thenReturn(true)
when(clientMetricsManager.processPushTelemetryRequest(any[PushTelemetryRequest](), any[RequestContext]()))
.thenThrow(new RuntimeException("test"))

metadataCache = MetadataCache.kRaftMetadataCache(brokerId)
createKafkaApis(raftSupport = true).handle(request, RequestLocal.NoCaching)
val response = verifyNoThrottling[PushTelemetryResponse](request)

val expectedResponse = new PushTelemetryResponseData().setErrorCode(Errors.INVALID_REQUEST.code)
assertEquals(expectedResponse, response.data)
}
}