From 7643d68ffb5b54c225be6865110e1822e4584cc0 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Wed, 8 Jan 2025 12:45:35 +0000 Subject: [PATCH 01/13] remove alterConfig and IncrementalAlterConfig --- .../main/scala/kafka/server/KafkaApis.scala | 96 +-------- .../unit/kafka/server/KafkaApisTest.scala | 198 +++--------------- 2 files changed, 38 insertions(+), 256 deletions(-) diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index c1e8807de9d68..5f5b568614b07 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -27,16 +27,13 @@ import kafka.server.share.SharePartitionManager import kafka.utils.{CoreUtils, Logging} import org.apache.kafka.admin.AdminUtils import org.apache.kafka.clients.CommonClientConfigs -import org.apache.kafka.clients.admin.AlterConfigOp.OpType -import org.apache.kafka.clients.admin.{AlterConfigOp, ConfigEntry, EndpointType} +import org.apache.kafka.clients.admin.EndpointType import org.apache.kafka.common.acl.AclOperation import org.apache.kafka.common.acl.AclOperation._ -import org.apache.kafka.common.config.ConfigResource import org.apache.kafka.common.errors._ import org.apache.kafka.common.internals.Topic.{GROUP_METADATA_TOPIC_NAME, SHARE_GROUP_STATE_TOPIC_NAME, TRANSACTION_STATE_TOPIC_NAME, isInternal} import org.apache.kafka.common.internals.{FatalExitError, Topic} import org.apache.kafka.common.message.AddPartitionsToTxnResponseData.{AddPartitionsToTxnResult, AddPartitionsToTxnResultCollection} -import org.apache.kafka.common.message.AlterConfigsResponseData.AlterConfigsResourceResponse import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.{ReassignablePartitionResponse, ReassignableTopicResponse} import org.apache.kafka.common.message.CreatePartitionsResponseData.CreatePartitionsTopicResult import org.apache.kafka.common.message.CreateTopicsRequestData.CreatableTopic @@ -2686,42 +2683,10 @@ class KafkaApis(val requestChannel: RequestChannel, new AlterConfigsRequest(remaining, request.header.apiVersion()), response => sendResponse(response.map(_.data()))) } else { - sendResponse(Some(processLegacyAlterConfigsRequest(request, remaining))) + throw KafkaApis.shouldAlwaysForward(request) } } - def processLegacyAlterConfigsRequest( - originalRequest: RequestChannel.Request, - data: AlterConfigsRequestData - ): AlterConfigsResponseData = { - val zkSupport = metadataSupport.requireZkOrThrow(KafkaApis.shouldAlwaysForward(originalRequest)) - val alterConfigsRequest = new AlterConfigsRequest(data, originalRequest.header.apiVersion()) - val (authorizedResources, unauthorizedResources) = alterConfigsRequest.configs.asScala.toMap.partition { case (resource, _) => - resource.`type` match { - case ConfigResource.Type.BROKER_LOGGER => - throw new InvalidRequestException(s"AlterConfigs is deprecated and does not support the resource type ${ConfigResource.Type.BROKER_LOGGER}") - case ConfigResource.Type.BROKER | ConfigResource.Type.CLIENT_METRICS => - authHelper.authorize(originalRequest.context, ALTER_CONFIGS, CLUSTER, CLUSTER_NAME) - case ConfigResource.Type.TOPIC => - authHelper.authorize(originalRequest.context, ALTER_CONFIGS, TOPIC, resource.name) - case rt => throw new InvalidRequestException(s"Unexpected resource type $rt") - } - } - val authorizedResult = zkSupport.adminManager.alterConfigs(authorizedResources, alterConfigsRequest.validateOnly) - val unauthorizedResult = unauthorizedResources.keys.map { resource => - resource -> configsAuthorizationApiError(resource) - } - val response = new AlterConfigsResponseData() - (authorizedResult ++ unauthorizedResult).foreach { case (resource, error) => - response.responses().add(new AlterConfigsResourceResponse() - .setErrorCode(error.error.code) - .setErrorMessage(error.message) - .setResourceName(resource.name) - .setResourceType(resource.`type`.id)) - } - response - } - def handleAlterPartitionReassignmentsRequest(request: RequestChannel.Request): Unit = { val zkSupport = metadataSupport.requireZkOrThrow(KafkaApis.shouldAlwaysForward(request)) authHelper.authorizeClusterOperation(request, ALTER) @@ -2808,31 +2773,12 @@ class KafkaApis(val requestChannel: RequestChannel, zkSupport.controller.listPartitionReassignments(partitionsOpt, sendResponseCallback) } - private def configsAuthorizationApiError(resource: ConfigResource): ApiError = { - val error = resource.`type` match { - case ConfigResource.Type.BROKER | ConfigResource.Type.BROKER_LOGGER => Errors.CLUSTER_AUTHORIZATION_FAILED - case ConfigResource.Type.TOPIC => Errors.TOPIC_AUTHORIZATION_FAILED - case ConfigResource.Type.GROUP => Errors.GROUP_AUTHORIZATION_FAILED - case rt => throw new InvalidRequestException(s"Unexpected resource type $rt for resource ${resource.name}") - } - new ApiError(error, null) - } - def handleIncrementalAlterConfigsRequest(request: RequestChannel.Request): Unit = { val original = request.body[IncrementalAlterConfigsRequest] val preprocessingResponses = configManager.preprocess(original.data(), (rType, rName) => authHelper.authorize(request.context, ALTER_CONFIGS, rType, rName)) val remaining = ConfigAdminManager.copyWithoutPreprocessed(original.data(), preprocessingResponses) - // Before deciding whether to forward or handle locally, a ZK broker needs to check if - // the active controller is ZK or KRaft. If the controller is KRaft, we need to forward. - // If the controller is ZK, we need to process the request locally. - val isKRaftController = metadataSupport match { - case ZkSupport(_, _, _, _, metadataCache, _) => - metadataCache.getControllerId.exists(_.isInstanceOf[KRaftCachedControllerId]) - case RaftSupport(_, _) => true - } - def sendResponse(secondPart: Option[ApiMessage]): Unit = { secondPart match { case Some(result: IncrementalAlterConfigsResponseData) => @@ -2845,49 +2791,15 @@ class KafkaApis(val requestChannel: RequestChannel, } } - // Forwarding has not happened yet, so handle both ZK and KRaft cases here if (remaining.resources().isEmpty) { sendResponse(Some(new IncrementalAlterConfigsResponseData())) - } else if ((!request.isForwarded) && metadataSupport.canForward() && isKRaftController) { + } else if ((!request.isForwarded) && metadataSupport.canForward()) { metadataSupport.forwardingManager.get.forwardRequest(request, new IncrementalAlterConfigsRequest(remaining, request.header.apiVersion()), response => sendResponse(response.map(_.data()))) } else { - sendResponse(Some(processIncrementalAlterConfigsRequest(request, remaining))) - } - } - - def processIncrementalAlterConfigsRequest( - originalRequest: RequestChannel.Request, - data: IncrementalAlterConfigsRequestData - ): IncrementalAlterConfigsResponseData = { - val zkSupport = metadataSupport.requireZkOrThrow(KafkaApis.shouldAlwaysForward(originalRequest)) - val configs = data.resources.iterator.asScala.map { alterConfigResource => - val configResource = new ConfigResource(ConfigResource.Type.forId(alterConfigResource.resourceType), - alterConfigResource.resourceName) - configResource -> alterConfigResource.configs.iterator.asScala.map { - alterConfig => new AlterConfigOp(new ConfigEntry(alterConfig.name, alterConfig.value), - OpType.forId(alterConfig.configOperation)) - }.toBuffer - }.toMap - - val (authorizedResources, unauthorizedResources) = configs.partition { case (resource, _) => - resource.`type` match { - case ConfigResource.Type.BROKER | ConfigResource.Type.BROKER_LOGGER | ConfigResource.Type.CLIENT_METRICS => - authHelper.authorize(originalRequest.context, ALTER_CONFIGS, CLUSTER, CLUSTER_NAME) - case ConfigResource.Type.TOPIC => - authHelper.authorize(originalRequest.context, ALTER_CONFIGS, TOPIC, resource.name) - case ConfigResource.Type.GROUP => - authHelper.authorize(originalRequest.context, ALTER_CONFIGS, GROUP, resource.name) - case rt => throw new InvalidRequestException(s"Unexpected resource type $rt") - } - } - - val authorizedResult = zkSupport.adminManager.incrementalAlterConfigs(authorizedResources, data.validateOnly) - val unauthorizedResult = unauthorizedResources.keys.map { resource => - resource -> configsAuthorizationApiError(resource) + throw KafkaApis.shouldAlwaysForward(request) } - new IncrementalAlterConfigsResponse(0, (authorizedResult ++ unauthorizedResult).asJava).data() } def handleDescribeConfigsRequest(request: RequestChannel.Request): Unit = { diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 229b1a18a19a7..b349c57bbc87c 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -98,7 +98,7 @@ import org.apache.kafka.server.util.{FutureUtils, MockTime} import org.apache.kafka.storage.internals.log.{AppendOrigin, LogConfig} import org.apache.kafka.storage.log.metrics.BrokerTopicStats import org.junit.jupiter.api.Assertions._ -import org.junit.jupiter.api.{AfterEach, Test} +import org.junit.jupiter.api.{AfterEach, Test, Disabled} import org.junit.jupiter.params.ParameterizedTest import org.junit.jupiter.params.provider.{CsvSource, EnumSource, ValueSource} import org.mockito.ArgumentMatchers._ @@ -324,79 +324,6 @@ class KafkaApisTest extends Logging { assertEquals(propValue, describeConfigsResponseData.value) } - @Test - def testEnvelopeRequestHandlingAsController(): Unit = { - testEnvelopeRequestWithAlterConfig( - alterConfigHandler = () => ApiError.NONE, - expectedError = Errors.NONE - ) - } - - @Test - def testEnvelopeRequestWithAlterConfigUnhandledError(): Unit = { - testEnvelopeRequestWithAlterConfig( - alterConfigHandler = () => throw new IllegalStateException(), - expectedError = Errors.UNKNOWN_SERVER_ERROR - ) - } - - private def testEnvelopeRequestWithAlterConfig( - alterConfigHandler: () => ApiError, - expectedError: Errors - ): Unit = { - val authorizer: Authorizer = mock(classOf[Authorizer]) - - authorizeResource(authorizer, AclOperation.CLUSTER_ACTION, ResourceType.CLUSTER, Resource.CLUSTER_NAME, AuthorizationResult.ALLOWED) - - val operation = AclOperation.ALTER_CONFIGS - val resourceName = "topic-1" - val requestHeader = new RequestHeader(ApiKeys.ALTER_CONFIGS, ApiKeys.ALTER_CONFIGS.latestVersion, - clientId, 0) - - when(controller.isActive).thenReturn(true) - - authorizeResource(authorizer, operation, ResourceType.TOPIC, resourceName, AuthorizationResult.ALLOWED) - - val configResource = new ConfigResource(ConfigResource.Type.TOPIC, resourceName) - when(adminManager.alterConfigs(any(), ArgumentMatchers.eq(false))) - .thenAnswer(_ => { - Map(configResource -> alterConfigHandler.apply()) - }) - - val configs = Map( - configResource -> new AlterConfigsRequest.Config( - Seq(new AlterConfigsRequest.ConfigEntry("foo", "bar")).asJava)) - val alterConfigsRequest = new AlterConfigsRequest.Builder(configs.asJava, false).build(requestHeader.apiVersion) - - val startTimeNanos = time.nanoseconds() - val queueDurationNanos = 5 * 1000 * 1000 - val request = TestUtils.buildEnvelopeRequest( - alterConfigsRequest, kafkaPrincipalSerde, requestChannelMetrics, startTimeNanos, startTimeNanos + queueDurationNanos) - - val capturedResponse: ArgumentCaptor[AlterConfigsResponse] = ArgumentCaptor.forClass(classOf[AlterConfigsResponse]) - val capturedRequest: ArgumentCaptor[RequestChannel.Request] = ArgumentCaptor.forClass(classOf[RequestChannel.Request]) - kafkaApis = createKafkaApis(authorizer = Some(authorizer), enableForwarding = true) - kafkaApis.handle(request, RequestLocal.withThreadConfinedCaching) - - verify(requestChannel).sendResponse( - capturedRequest.capture(), - capturedResponse.capture(), - any() - ) - assertEquals(Some(request), capturedRequest.getValue.envelope) - // the dequeue time of forwarded request should equals to envelop request - assertEquals(request.requestDequeueTimeNanos, capturedRequest.getValue.requestDequeueTimeNanos) - val innerResponse = capturedResponse.getValue - val responseMap = innerResponse.data.responses().asScala.map { resourceResponse => - resourceResponse.resourceName -> Errors.forCode(resourceResponse.errorCode) - }.toMap - - assertEquals(Map(resourceName -> expectedError), responseMap) - - verify(controller).isActive - verify(adminManager).alterConfigs(any(), ArgumentMatchers.eq(false)) - } - @Test def testInvalidEnvelopeRequestWithNonForwardableAPI(): Unit = { val requestHeader = new RequestHeader(ApiKeys.LEAVE_GROUP, ApiKeys.LEAVE_GROUP.latestVersion, @@ -485,44 +412,6 @@ class KafkaApisTest extends Logging { } } - @Test - def testAlterConfigsWithAuthorizer(): Unit = { - val authorizer: Authorizer = mock(classOf[Authorizer]) - - val authorizedTopic = "authorized-topic" - val unauthorizedTopic = "unauthorized-topic" - val (authorizedResource, unauthorizedResource) = - createConfigsWithAuthorization(authorizer, authorizedTopic, unauthorizedTopic) - - val configs = Map( - authorizedResource -> new AlterConfigsRequest.Config( - Seq(new AlterConfigsRequest.ConfigEntry("foo", "bar")).asJava), - unauthorizedResource -> new AlterConfigsRequest.Config( - Seq(new AlterConfigsRequest.ConfigEntry("foo-1", "bar-1")).asJava) - ) - - val topicHeader = new RequestHeader(ApiKeys.ALTER_CONFIGS, ApiKeys.ALTER_CONFIGS.latestVersion, - clientId, 0) - - val alterConfigsRequest = new AlterConfigsRequest.Builder(configs.asJava, false) - .build(topicHeader.apiVersion) - val request = buildRequest(alterConfigsRequest) - - when(controller.isActive).thenReturn(false) - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), - any[Long])).thenReturn(0) - when(adminManager.alterConfigs(any(), ArgumentMatchers.eq(false))) - .thenReturn(Map(authorizedResource -> ApiError.NONE)) - kafkaApis = createKafkaApis(authorizer = Some(authorizer)) - kafkaApis.handleAlterConfigsRequest(request) - - val response = verifyNoThrottling[AlterConfigsResponse](request) - verifyAlterConfigResult(response, Map(authorizedTopic -> Errors.NONE, - unauthorizedTopic -> Errors.TOPIC_AUTHORIZATION_FAILED)) - verify(authorizer, times(2)).authorize(any(), any()) - verify(adminManager).alterConfigs(any(), anyBoolean()) - } - @Test def testElectLeadersForwarding(): Unit = { val requestBuilder = new ElectLeadersRequest.Builder(ElectionType.PREFERRED, null, 30000) @@ -547,17 +436,12 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, fromPrivilegedListener = true, requestHeader = Option(requestHeader)) - when(controller.isActive).thenReturn(true) when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), any[Long])).thenReturn(0) - when(adminManager.incrementalAlterConfigs(any(), ArgumentMatchers.eq(false))) - .thenReturn(Map(resource -> ApiError.NONE)) - createKafkaApis(authorizer = Some(authorizer)).handleIncrementalAlterConfigsRequest(request) - val response = verifyNoThrottling[IncrementalAlterConfigsResponse](request) - verifyIncrementalAlterConfigResult(response, Map(consumerGroupId -> Errors.NONE)) - verify(authorizer, times(1)).authorize(any(), any()) - verify(adminManager).incrementalAlterConfigs(any(), anyBoolean()) + metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) + createKafkaApis(authorizer = Some(authorizer), raftSupport = true).handleIncrementalAlterConfigsRequest(request) + testNewBodyForwardableApi(request) } @Test @@ -630,51 +514,31 @@ class KafkaApisTest extends Logging { val configs = Map(authorizedResource -> new AlterConfigsRequest.Config(configEntries)) val requestHeader = new RequestHeader(ApiKeys.ALTER_CONFIGS, ApiKeys.ALTER_CONFIGS.latestVersion, clientId, 0) - val request = buildRequest( - new AlterConfigsRequest.Builder(configs.asJava, false).build(requestHeader.apiVersion)) + val apiRequest = new AlterConfigsRequest.Builder(configs.asJava, false).build(requestHeader.apiVersion) + val request = buildRequest(apiRequest) - when(controller.isActive).thenReturn(false) - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), - any[Long])).thenReturn(0) - when(adminManager.alterConfigs(any(), ArgumentMatchers.eq(false))) - .thenReturn(Map(authorizedResource -> ApiError.NONE)) - kafkaApis = createKafkaApis(authorizer = Some(authorizer)) + metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) + kafkaApis = createKafkaApis(raftSupport = true) kafkaApis.handleAlterConfigsRequest(request) - val response = verifyNoThrottling[AlterConfigsResponse](request) - verifyAlterConfigResult(response, Map(subscriptionName -> Errors.NONE)) - verify(authorizer, times(1)).authorize(any(), any()) - verify(adminManager).alterConfigs(any(), anyBoolean()) + testNewBodyForwardableApi(request) } @Test def testIncrementalClientMetricAlterConfigs(): Unit = { - val authorizer: Authorizer = mock(classOf[Authorizer]) - val subscriptionName = "client_metric_subscription_1" val resource = new ConfigResource(ConfigResource.Type.CLIENT_METRICS, subscriptionName) - authorizeResource(authorizer, AclOperation.ALTER_CONFIGS, ResourceType.CLUSTER, - Resource.CLUSTER_NAME, AuthorizationResult.ALLOWED) - val requestHeader = new RequestHeader(ApiKeys.INCREMENTAL_ALTER_CONFIGS, ApiKeys.INCREMENTAL_ALTER_CONFIGS.latestVersion, clientId, 0) val incrementalAlterConfigsRequest = getIncrementalAlterConfigRequestBuilder( Seq(resource), "metrics", "foo.bar").build(requestHeader.apiVersion) - val request = buildRequest(incrementalAlterConfigsRequest, - fromPrivilegedListener = true, requestHeader = Option(requestHeader)) + val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) - when(controller.isActive).thenReturn(true) - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), - any[Long])).thenReturn(0) - when(adminManager.incrementalAlterConfigs(any(), ArgumentMatchers.eq(false))) - .thenReturn(Map(resource -> ApiError.NONE)) - kafkaApis = createKafkaApis(authorizer = Some(authorizer)) + metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) + kafkaApis = createKafkaApis(raftSupport = true) kafkaApis.handleIncrementalAlterConfigsRequest(request) - val response = verifyNoThrottling[IncrementalAlterConfigsResponse](request) - verifyIncrementalAlterConfigResult(response, Map(subscriptionName -> Errors.NONE )) - verify(authorizer, times(1)).authorize(any(), any()) - verify(adminManager).incrementalAlterConfigs(any(), anyBoolean()) + testNewBodyForwardableApi(request) } private def getIncrementalAlterConfigRequestBuilder(configResources: Seq[ConfigResource], @@ -761,6 +625,21 @@ class KafkaApisTest extends Logging { ) } + private def testNewBodyForwardableApi(request: RequestChannel.Request + ): Unit = { + val forwardCallback: ArgumentCaptor[Option[AbstractResponse] => Unit] = ArgumentCaptor.forClass(classOf[Option[AbstractResponse] => Unit]) + val newBody: ArgumentCaptor[AbstractRequest] = ArgumentCaptor.forClass(classOf[AbstractRequest]) + + verify(forwardingManager).forwardRequest( + ArgumentMatchers.eq(request), + newBody.capture(), + forwardCallback.capture() + ) + + assertNotNull(request.buffer, "The buffer was unexpectedly deallocated") + assertNotNull(newBody, "The newBody field should not be null") + } + private def testForwardableApi(apiKey: ApiKeys, requestBuilder: AbstractRequest.Builder[_ <: AbstractRequest]): Unit = { kafkaApis = createKafkaApis(enableForwarding = true) testForwardableApi(kafkaApis = kafkaApis, @@ -830,15 +709,6 @@ class KafkaApisTest extends Logging { .thenReturn(Seq(result).asJava) } - private def verifyAlterConfigResult(response: AlterConfigsResponse, - expectedResults: Map[String, Errors]): Unit = { - val responseMap = response.data.responses().asScala.map { resourceResponse => - resourceResponse.resourceName -> Errors.forCode(resourceResponse.errorCode) - }.toMap - - assertEquals(expectedResults, responseMap) - } - private def createConfigsWithAuthorization(authorizer: Authorizer, authorizedTopic: String, unauthorizedTopic: String): (ConfigResource, ConfigResource) = { @@ -850,6 +720,7 @@ class KafkaApisTest extends Logging { (authorizedResource, unauthorizedResource) } + @Disabled @Test def testIncrementalAlterConfigsWithAuthorizer(): Unit = { val authorizer: Authorizer = mock(classOf[Authorizer]) @@ -863,15 +734,13 @@ class KafkaApisTest extends Logging { val incrementalAlterConfigsRequest = getIncrementalAlterConfigRequestBuilder(Seq(authorizedResource, unauthorizedResource)) .build(requestHeader.apiVersion) - val request = buildRequest(incrementalAlterConfigsRequest, - fromPrivilegedListener = true, requestHeader = Option(requestHeader)) + val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) - when(controller.isActive).thenReturn(true) when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), any[Long])).thenReturn(0) - when(adminManager.incrementalAlterConfigs(any(), ArgumentMatchers.eq(false))) - .thenReturn(Map(authorizedResource -> ApiError.NONE)) - kafkaApis = createKafkaApis(authorizer = Some(authorizer)) + + metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) + kafkaApis = createKafkaApis(authorizer = Some(authorizer), raftSupport = true) kafkaApis.handleIncrementalAlterConfigsRequest(request) val capturedResponse = verifyNoThrottling[IncrementalAlterConfigsResponse](request) @@ -894,6 +763,7 @@ class KafkaApisTest extends Logging { new IncrementalAlterConfigsRequest.Builder(resourceMap, false) } + private def verifyIncrementalAlterConfigResult(response: IncrementalAlterConfigsResponse, expectedResults: Map[String, Errors]): Unit = { val responseMap = response.data.responses.asScala.map { resourceResponse => From 14ade5cae5aeb0ac0697c9a4f1d3510c1d5155c3 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Wed, 8 Jan 2025 18:30:40 +0000 Subject: [PATCH 02/13] remove incrementalAlterConfigWithAuthorizer --- .../unit/kafka/server/KafkaApisTest.scala | 83 +------------------ 1 file changed, 2 insertions(+), 81 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index b349c57bbc87c..d133b7111094c 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -98,7 +98,7 @@ import org.apache.kafka.server.util.{FutureUtils, MockTime} import org.apache.kafka.storage.internals.log.{AppendOrigin, LogConfig} import org.apache.kafka.storage.log.metrics.BrokerTopicStats import org.junit.jupiter.api.Assertions._ -import org.junit.jupiter.api.{AfterEach, Test, Disabled} +import org.junit.jupiter.api.{AfterEach, Test} import org.junit.jupiter.params.ParameterizedTest import org.junit.jupiter.params.provider.{CsvSource, EnumSource, ValueSource} import org.mockito.ArgumentMatchers._ @@ -433,11 +433,7 @@ class KafkaApisTest extends Logging { val incrementalAlterConfigsRequest = getIncrementalAlterConfigRequestBuilder( Seq(resource), "consumer.session.timeout.ms", "45000").build(requestHeader.apiVersion) - val request = buildRequest(incrementalAlterConfigsRequest, - fromPrivilegedListener = true, requestHeader = Option(requestHeader)) - - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), - any[Long])).thenReturn(0) + val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) createKafkaApis(authorizer = Some(authorizer), raftSupport = true).handleIncrementalAlterConfigsRequest(request) @@ -709,69 +705,6 @@ class KafkaApisTest extends Logging { .thenReturn(Seq(result).asJava) } - private def createConfigsWithAuthorization(authorizer: Authorizer, - authorizedTopic: String, - unauthorizedTopic: String): (ConfigResource, ConfigResource) = { - val authorizedResource = new ConfigResource(ConfigResource.Type.TOPIC, authorizedTopic) - - val unauthorizedResource = new ConfigResource(ConfigResource.Type.TOPIC, unauthorizedTopic) - - createTopicAuthorization(authorizer, AclOperation.ALTER_CONFIGS, authorizedTopic, unauthorizedTopic) - (authorizedResource, unauthorizedResource) - } - - @Disabled - @Test - def testIncrementalAlterConfigsWithAuthorizer(): Unit = { - val authorizer: Authorizer = mock(classOf[Authorizer]) - - val authorizedTopic = "authorized-topic" - val unauthorizedTopic = "unauthorized-topic" - val (authorizedResource, unauthorizedResource) = - createConfigsWithAuthorization(authorizer, authorizedTopic, unauthorizedTopic) - - val requestHeader = new RequestHeader(ApiKeys.INCREMENTAL_ALTER_CONFIGS, ApiKeys.INCREMENTAL_ALTER_CONFIGS.latestVersion, clientId, 0) - - val incrementalAlterConfigsRequest = getIncrementalAlterConfigRequestBuilder(Seq(authorizedResource, unauthorizedResource)) - .build(requestHeader.apiVersion) - val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) - - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), - any[Long])).thenReturn(0) - - metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - kafkaApis = createKafkaApis(authorizer = Some(authorizer), raftSupport = true) - kafkaApis.handleIncrementalAlterConfigsRequest(request) - - val capturedResponse = verifyNoThrottling[IncrementalAlterConfigsResponse](request) - verifyIncrementalAlterConfigResult(capturedResponse, Map( - authorizedTopic -> Errors.NONE, - unauthorizedTopic -> Errors.TOPIC_AUTHORIZATION_FAILED - )) - - verify(authorizer, times(2)).authorize(any(), any()) - verify(adminManager).incrementalAlterConfigs(any(), anyBoolean()) - } - - private def getIncrementalAlterConfigRequestBuilder(configResources: Seq[ConfigResource]): IncrementalAlterConfigsRequest.Builder = { - val resourceMap = configResources.map(configResource => { - configResource -> Set( - new AlterConfigOp(new ConfigEntry("foo", "bar"), - OpType.forId(configResource.`type`.id))).asJavaCollection - }).toMap.asJava - - new IncrementalAlterConfigsRequest.Builder(resourceMap, false) - } - - - private def verifyIncrementalAlterConfigResult(response: IncrementalAlterConfigsResponse, - expectedResults: Map[String, Errors]): Unit = { - val responseMap = response.data.responses.asScala.map { resourceResponse => - resourceResponse.resourceName -> Errors.forCode(resourceResponse.errorCode) - }.toMap - assertEquals(expectedResults, responseMap) - } - @Test def testAlterClientQuotasWithAuthorizer(): Unit = { val authorizer: Authorizer = mock(classOf[Authorizer]) @@ -1011,18 +944,6 @@ class KafkaApisTest extends Logging { assertEquals(Some(Errors.TOPIC_AUTHORIZATION_FAILED), results.find(_.name == "bar").map(result => Errors.forCode(result.errorCode))) } - private def createTopicAuthorization(authorizer: Authorizer, - operation: AclOperation, - authorizedTopic: String, - unauthorizedTopic: String, - logIfAllowed: Boolean = true, - logIfDenied: Boolean = true): Unit = { - authorizeResource(authorizer, operation, ResourceType.TOPIC, - authorizedTopic, AuthorizationResult.ALLOWED, logIfAllowed, logIfDenied) - authorizeResource(authorizer, operation, ResourceType.TOPIC, - unauthorizedTopic, AuthorizationResult.DENIED, logIfAllowed, logIfDenied) - } - private def createCombinedTopicAuthorization(authorizer: Authorizer, operation: AclOperation, authorizedTopic: String, From 616bc610d075b87a6ac273d052eae9ce23bf466c Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Fri, 10 Jan 2025 16:37:28 +0000 Subject: [PATCH 03/13] address comments --- core/src/main/scala/kafka/server/KafkaApis.scala | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 334eb33a9131c..dd94afdfe72fa 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -2517,12 +2517,10 @@ class KafkaApis(val requestChannel: RequestChannel, } if (remaining.resources().isEmpty) { sendResponse(Some(new AlterConfigsResponseData())) - } else if ((!request.isForwarded) && metadataSupport.canForward()) { + } else { metadataSupport.forwardingManager.get.forwardRequest(request, new AlterConfigsRequest(remaining, request.header.apiVersion()), response => sendResponse(response.map(_.data()))) - } else { - throw KafkaApis.shouldAlwaysForward(request) } } From da9e6675046aa4ca0b25dfc790675b4562275c90 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Fri, 10 Jan 2025 18:09:32 +0000 Subject: [PATCH 04/13] address comments --- core/src/main/scala/kafka/server/KafkaApis.scala | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index dd94afdfe72fa..6cd1fec4abfdf 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -2630,12 +2630,10 @@ class KafkaApis(val requestChannel: RequestChannel, if (remaining.resources().isEmpty) { sendResponse(Some(new IncrementalAlterConfigsResponseData())) - } else if ((!request.isForwarded) && metadataSupport.canForward()) { + } else { metadataSupport.forwardingManager.get.forwardRequest(request, new IncrementalAlterConfigsRequest(remaining, request.header.apiVersion()), response => sendResponse(response.map(_.data()))) - } else { - throw KafkaApis.shouldAlwaysForward(request) } } From 02805f3ef75a05d80c3b8dbb6fcefc9c6f7aea7a Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Sat, 11 Jan 2025 18:07:48 +0000 Subject: [PATCH 05/13] authorize --- .../unit/kafka/server/KafkaApisTest.scala | 29 +++++++++++++++++++ 1 file changed, 29 insertions(+) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 6745b22239126..1d17b3be3915b 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -615,6 +615,35 @@ class KafkaApisTest extends Logging { .thenReturn(Seq(result).asJava) } + @Test + def testIncrementalAlterConfigsWithAuthorizer(): Unit = { + val authorizer: Authorizer = mock(classOf[Authorizer]) + + val authorizedResource = new ConfigResource(ConfigResource.Type.BROKER_LOGGER, "authorizedResource") + val unauthorizedResource = new ConfigResource(ConfigResource.Type.BROKER_LOGGER, "unauthorizedResource") + + val requestHeader = new RequestHeader(ApiKeys.INCREMENTAL_ALTER_CONFIGS, ApiKeys.INCREMENTAL_ALTER_CONFIGS.latestVersion, clientId, 0) + + val incrementalAlterConfigsRequest = getIncrementalAlterConfigRequestBuilder(Seq(authorizedResource, unauthorizedResource)) + .build(requestHeader.apiVersion) + val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) + + kafkaApis = createKafkaApis(authorizer = Some(authorizer)) + kafkaApis.handleIncrementalAlterConfigsRequest(request) + + verify(authorizer, times(2)).authorize(any(), any()) + } + + private def getIncrementalAlterConfigRequestBuilder(configResources: Seq[ConfigResource]): IncrementalAlterConfigsRequest.Builder = { + val resourceMap = configResources.map(configResource => { + configResource -> Set( + new AlterConfigOp(new ConfigEntry("foo", "bar"), + OpType.SET)).asJavaCollection + }).toMap.asJava + + new IncrementalAlterConfigsRequest.Builder(resourceMap, false) + } + @Test def testAlterClientQuotasWithAuthorizer(): Unit = { val authorizer: Authorizer = mock(classOf[Authorizer]) From 15193792f4bf22d9cbb1ccacc013b455d164d54c Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Sat, 11 Jan 2025 18:19:00 +0000 Subject: [PATCH 06/13] fix verify --- .../unit/kafka/server/KafkaApisTest.scala | 33 +++++++++---------- 1 file changed, 15 insertions(+), 18 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 1d17b3be3915b..aa9fd93f4a13d 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -347,7 +347,11 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) createKafkaApis(authorizer = Some(authorizer), raftSupport = true).handleIncrementalAlterConfigsRequest(request) - testNewBodyForwardableApi(request) + verify(forwardingManager, times(1)).forwardRequest( + any(), + any(), + any() + ) } @Test @@ -426,7 +430,11 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) kafkaApis = createKafkaApis(raftSupport = true) kafkaApis.handleAlterConfigsRequest(request) - testNewBodyForwardableApi(request) + verify(forwardingManager, times(1)).forwardRequest( + any(), + any(), + any() + ) } @Test @@ -444,7 +452,11 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) kafkaApis = createKafkaApis(raftSupport = true) kafkaApis.handleIncrementalAlterConfigsRequest(request) - testNewBodyForwardableApi(request) + verify(forwardingManager, times(1)).forwardRequest( + any(), + any(), + any() + ) } private def getIncrementalAlterConfigRequestBuilder(configResources: Seq[ConfigResource], @@ -531,21 +543,6 @@ class KafkaApisTest extends Logging { ) } - private def testNewBodyForwardableApi(request: RequestChannel.Request - ): Unit = { - val forwardCallback: ArgumentCaptor[Option[AbstractResponse] => Unit] = ArgumentCaptor.forClass(classOf[Option[AbstractResponse] => Unit]) - val newBody: ArgumentCaptor[AbstractRequest] = ArgumentCaptor.forClass(classOf[AbstractRequest]) - - verify(forwardingManager).forwardRequest( - ArgumentMatchers.eq(request), - newBody.capture(), - forwardCallback.capture() - ) - - assertNotNull(request.buffer, "The buffer was unexpectedly deallocated") - assertNotNull(newBody, "The newBody field should not be null") - } - private def testForwardableApi(apiKey: ApiKeys, requestBuilder: AbstractRequest.Builder[_ <: AbstractRequest]): Unit = { kafkaApis = createKafkaApis(enableForwarding = true) testForwardableApi(kafkaApis = kafkaApis, From bbb41aaa7f2d4c0bdd51cdaf37ba72f01b357784 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Sat, 11 Jan 2025 18:55:52 +0000 Subject: [PATCH 07/13] forward + authorization --- .../scala/unit/kafka/server/KafkaApisTest.scala | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index aa9fd93f4a13d..d72560504fcb7 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -616,19 +616,24 @@ class KafkaApisTest extends Logging { def testIncrementalAlterConfigsWithAuthorizer(): Unit = { val authorizer: Authorizer = mock(classOf[Authorizer]) - val authorizedResource = new ConfigResource(ConfigResource.Type.BROKER_LOGGER, "authorizedResource") - val unauthorizedResource = new ConfigResource(ConfigResource.Type.BROKER_LOGGER, "unauthorizedResource") + val localResource = new ConfigResource(ConfigResource.Type.BROKER_LOGGER, "localResource") + val forwardedResource = new ConfigResource(ConfigResource.Type.GROUP, "forwardedResource") val requestHeader = new RequestHeader(ApiKeys.INCREMENTAL_ALTER_CONFIGS, ApiKeys.INCREMENTAL_ALTER_CONFIGS.latestVersion, clientId, 0) - val incrementalAlterConfigsRequest = getIncrementalAlterConfigRequestBuilder(Seq(authorizedResource, unauthorizedResource)) + val incrementalAlterConfigsRequest = getIncrementalAlterConfigRequestBuilder(Seq(localResource, forwardedResource)) .build(requestHeader.apiVersion) val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) kafkaApis = createKafkaApis(authorizer = Some(authorizer)) kafkaApis.handleIncrementalAlterConfigsRequest(request) - verify(authorizer, times(2)).authorize(any(), any()) + verify(authorizer, times(1)).authorize(any(), any()) + verify(forwardingManager, times(1)).forwardRequest( + any(), + any(), + any() + ) } private def getIncrementalAlterConfigRequestBuilder(configResources: Seq[ConfigResource]): IncrementalAlterConfigsRequest.Builder = { From 7e8213b68eae9ef3afc85f698d27c073f9956237 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Sat, 11 Jan 2025 18:59:24 +0000 Subject: [PATCH 08/13] fix zk issue --- core/src/test/scala/unit/kafka/server/KafkaApisTest.scala | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index d72560504fcb7..2bb26a5e712ae 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -625,7 +625,8 @@ class KafkaApisTest extends Logging { .build(requestHeader.apiVersion) val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) - kafkaApis = createKafkaApis(authorizer = Some(authorizer)) + metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) + kafkaApis = createKafkaApis(authorizer = Some(authorizer), raftSupport = true) kafkaApis.handleIncrementalAlterConfigsRequest(request) verify(authorizer, times(1)).authorize(any(), any()) From c9bf263612354b42d62612c3e31297027a8c6fb4 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Tue, 14 Jan 2025 23:06:19 +0000 Subject: [PATCH 09/13] cleanup zk related in kafkaApisTest --- .../unit/kafka/server/KafkaApisTest.scala | 394 +++++------------- 1 file changed, 93 insertions(+), 301 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index d4cad343c0492..44604b9b4fbad 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -17,8 +17,7 @@ package kafka.server -import kafka.cluster.{Broker, Partition} -import kafka.controller.KafkaController +import kafka.cluster.Partition import kafka.coordinator.transaction.{InitProducerIdResult, TransactionCoordinator} import kafka.log.UnifiedLog import kafka.network.RequestChannel @@ -26,7 +25,6 @@ import kafka.server.QuotaFactory.QuotaManagers import kafka.server.metadata.{ConfigRepository, KRaftMetadataCache, MockConfigRepository, ZkMetadataCache} import kafka.server.share.SharePartitionManager import kafka.utils.{CoreUtils, Log4jController, Logging, TestUtils} -import kafka.zk.KafkaZkClient import org.apache.kafka.clients.admin.AlterConfigOp.OpType import org.apache.kafka.clients.admin.{AlterConfigOp, ConfigEntry} import org.apache.kafka.common._ @@ -82,7 +80,7 @@ import org.apache.kafka.security.authorizer.AclEntry import org.apache.kafka.server.{BrokerFeatures, ClientMetricsManager} import org.apache.kafka.server.authorizer.{Action, AuthorizationResult, Authorizer} import org.apache.kafka.server.common.{FeatureVersion, FinalizedFeatures, GroupVersion, KRaftVersion, MetadataVersion, RequestLocal, TransactionVersion} -import org.apache.kafka.server.config.{ConfigType, KRaftConfigs, ReplicationConfigs, ServerConfigs, ServerLogConfigs} +import org.apache.kafka.server.config.{KRaftConfigs, ReplicationConfigs, ServerConfigs, ServerLogConfigs} import org.apache.kafka.server.metrics.ClientMetricsTestUtils import org.apache.kafka.server.share.{CachedSharePartition, ErroneousAndValidPartitionData} import org.apache.kafka.server.quota.ThrottleCallback @@ -119,9 +117,7 @@ class KafkaApisTest extends Logging { private val replicaManager: ReplicaManager = mock(classOf[ReplicaManager]) private val groupCoordinator: GroupCoordinator = mock(classOf[GroupCoordinator]) private val shareCoordinator: ShareCoordinator = mock(classOf[ShareCoordinator]) - private val adminManager: ZkAdminManager = mock(classOf[ZkAdminManager]) private val txnCoordinator: TransactionCoordinator = mock(classOf[TransactionCoordinator]) - private val controller: KafkaController = mock(classOf[KafkaController]) private val forwardingManager: ForwardingManager = mock(classOf[ForwardingManager]) private val autoTopicCreationManager: AutoTopicCreationManager = mock(classOf[AutoTopicCreationManager]) @@ -129,12 +125,10 @@ class KafkaApisTest extends Logging { override def serialize(principal: KafkaPrincipal): Array[Byte] = Utils.utf8(principal.toString) override def deserialize(bytes: Array[Byte]): KafkaPrincipal = SecurityUtils.parseKafkaPrincipal(Utils.utf8(bytes)) } - private val zkClient: KafkaZkClient = mock(classOf[KafkaZkClient]) private val metrics = new Metrics() private val brokerId = 1 // KRaft tests should override this with a KRaftMetadataCache - private var metadataCache: MetadataCache = MetadataCache.zkMetadataCache(brokerId, MetadataVersion.latestTesting()) - private val brokerEpochManager: ZkBrokerEpochManager = new ZkBrokerEpochManager(metadataCache, controller, None) + private var metadataCache: MetadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) private val clientQuotaManager: ClientQuotaManager = mock(classOf[ClientQuotaManager]) private val clientRequestQuotaManager: ClientRequestQuotaManager = mock(classOf[ClientRequestQuotaManager]) private val clientControllerQuotaManager: ControllerMutationQuotaManager = mock(classOf[ControllerMutationQuotaManager]) @@ -162,59 +156,37 @@ class KafkaApisTest extends Logging { def createKafkaApis(interBrokerProtocolVersion: MetadataVersion = MetadataVersion.latestTesting, authorizer: Option[Authorizer] = None, - enableForwarding: Boolean = false, configRepository: ConfigRepository = new MockConfigRepository(), - raftSupport: Boolean = false, overrideProperties: Map[String, String] = Map.empty, featureVersions: Seq[FeatureVersion] = Seq.empty): KafkaApis = { - val properties = if (raftSupport) { - val properties = TestUtils.createBrokerConfig(brokerId) - properties.put(KRaftConfigs.NODE_ID_CONFIG, brokerId.toString) - properties.put(KRaftConfigs.PROCESS_ROLES_CONFIG, "broker") - val voterId = brokerId + 1 - properties.put(QuorumConfig.QUORUM_VOTERS_CONFIG, s"$voterId@localhost:9093") - properties - } else { - TestUtils.createBrokerConfig(brokerId) - } + + val properties = TestUtils.createBrokerConfig(brokerId) + properties.put(KRaftConfigs.NODE_ID_CONFIG, brokerId.toString) + properties.put(KRaftConfigs.PROCESS_ROLES_CONFIG, "broker") + val voterId = brokerId + 1 + properties.put(QuorumConfig.QUORUM_VOTERS_CONFIG, s"$voterId@localhost:9093") + overrideProperties.foreach( p => properties.put(p._1, p._2)) TestUtils.setIbpVersion(properties, interBrokerProtocolVersion) val config = new KafkaConfig(properties) - val forwardingManagerOpt = if (enableForwarding) - Some(this.forwardingManager) - else - None - - val metadataSupport = if (raftSupport) { - // it will be up to the test to replace the default ZkMetadataCache implementation - // with a KRaftMetadataCache instance - metadataCache match { + val metadataSupport = metadataCache match { case cache: KRaftMetadataCache => RaftSupport(forwardingManager, cache) case _ => throw new IllegalStateException("Test must set an instance of KRaftMetadataCache") } - } else { - metadataCache match { - case zkMetadataCache: ZkMetadataCache => - ZkSupport(adminManager, controller, zkClient, forwardingManagerOpt, zkMetadataCache, brokerEpochManager) - case _ => throw new IllegalStateException("Test must set an instance of ZkMetadataCache") - } - } - val listenerType = if (raftSupport) ListenerType.BROKER else ListenerType.ZK_BROKER - val enabledApis = if (enableForwarding) { - ApiKeys.apisForListener(listenerType).asScala ++ Set(ApiKeys.ENVELOPE) - } else { - ApiKeys.apisForListener(listenerType).asScala.toSet - } + + val listenerType = ListenerType.BROKER + val enabledApis = ApiKeys.apisForListener(listenerType).asScala ++ Set(ApiKeys.ENVELOPE) + val apiVersionManager = new SimpleApiVersionManager( listenerType, enabledApis, BrokerFeatures.defaultSupportedFeatures(true), true, - () => new FinalizedFeatures(MetadataVersion.latestTesting(), Collections.emptyMap[String, java.lang.Short], 0, raftSupport)) + () => new FinalizedFeatures(MetadataVersion.latestTesting(), Collections.emptyMap[String, java.lang.Short], 0, true)) - val clientMetricsManagerOpt = if (raftSupport) Some(clientMetricsManager) else None + val clientMetricsManagerOpt = Some(clientMetricsManager) when(groupCoordinator.isNewGroupCoordinator).thenReturn(config.isNewGroupCoordinatorEnabled) setupFeatures(featureVersions) @@ -344,7 +316,7 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - createKafkaApis(authorizer = Some(authorizer), raftSupport = true).handleIncrementalAlterConfigsRequest(request) + createKafkaApis(authorizer = Some(authorizer)).handleIncrementalAlterConfigsRequest(request) verify(forwardingManager, times(1)).forwardRequest( any(), any(), @@ -422,7 +394,7 @@ class KafkaApisTest extends Logging { val request = buildRequest(apiRequest) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handleAlterConfigsRequest(request) verify(forwardingManager, times(1)).forwardRequest( any(), @@ -444,7 +416,7 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handleIncrementalAlterConfigsRequest(request) verify(forwardingManager, times(1)).forwardRequest( any(), @@ -518,7 +490,7 @@ class KafkaApisTest extends Logging { val requestData = DescribeQuorumRequest.singletonRequest(KafkaRaftServer.MetadataPartition) val requestBuilder = new DescribeQuorumRequest.Builder(requestData) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() testForwardableApi(kafkaApis = kafkaApis, ApiKeys.DESCRIBE_QUORUM, requestBuilder @@ -530,7 +502,7 @@ class KafkaApisTest extends Logging { requestBuilder: AbstractRequest.Builder[_ <: AbstractRequest] ): Unit = { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(enableForwarding = true, raftSupport = true) + kafkaApis = createKafkaApis() testForwardableApi(kafkaApis = kafkaApis, apiKey, requestBuilder @@ -548,13 +520,6 @@ class KafkaApisTest extends Logging { val apiRequest = requestBuilder.build(topicHeader.apiVersion) val request = buildRequest(apiRequest) - if (kafkaApis.metadataSupport.isInstanceOf[ZkSupport]) { - // The controller check only makes sense for ZK clusters. For KRaft, - // controller requests are handled on a separate listener, so there - // is no choice but to forward them. - when(controller.isActive).thenReturn(false) - } - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), any[Long])).thenReturn(0) val forwardCallback: ArgumentCaptor[Option[AbstractResponse] => Unit] = ArgumentCaptor.forClass(classOf[Option[AbstractResponse] => Unit]) @@ -573,9 +538,6 @@ class KafkaApisTest extends Logging { val capturedResponse = verifyNoThrottling[AbstractResponse](request) assertEquals(expectedResponse.data, capturedResponse.data) - if (kafkaApis.metadataSupport.isInstanceOf[ZkSupport]) { - verify(controller).isActive - } } private def authorizeResource(authorizer: Authorizer, @@ -612,7 +574,7 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - kafkaApis = createKafkaApis(authorizer = Some(authorizer), raftSupport = true) + kafkaApis = createKafkaApis(authorizer = Some(authorizer)) kafkaApis.handleIncrementalAlterConfigsRequest(request) verify(authorizer, times(1)).authorize(any(), any()) @@ -652,7 +614,7 @@ class KafkaApisTest extends Logging { val requestBuilder = new CreateTopicsRequest.Builder(requestData).build() val request = buildRequest(requestBuilder) - kafkaApis = createKafkaApis(enableForwarding = true, raftSupport = true) + kafkaApis = createKafkaApis() val forwardCallback: ArgumentCaptor[Option[AbstractResponse] => Unit] = ArgumentCaptor.forClass(classOf[Option[AbstractResponse] => Unit]) @@ -904,8 +866,7 @@ class KafkaApisTest extends Logging { any[Long])).thenReturn(0) val capturedRequest = verifyTopicCreation(topicName, enableAutoTopicCreation, isInternal, request) - kafkaApis = createKafkaApis(authorizer = Some(authorizer), enableForwarding = enableAutoTopicCreation, - overrideProperties = topicConfigOverride) + kafkaApis = createKafkaApis(authorizer = Some(authorizer), overrideProperties = topicConfigOverride) kafkaApis.handleTopicMetadataRequest(request) val response = verifyNoThrottling[MetadataResponse](request) @@ -3612,175 +3573,6 @@ class KafkaApisTest extends Logging { assertEquals(Set(0), response.brokers.asScala.map(_.id).toSet) } - - /** - * Metadata request to fetch all topics should not result in the followings: - * 1) Auto topic creation - * 2) UNKNOWN_TOPIC_OR_PARTITION - * - * This case is testing the case that a topic is being deleted from MetadataCache right after - * authorization but before checking in MetadataCache. - */ - @Test - def testGetAllTopicMetadataShouldNotCreateTopicOrReturnUnknownTopicPartition(): Unit = { - // Setup: authorizer authorizes 2 topics, but one got deleted in metadata cache - metadataCache = mock(classOf[ZkMetadataCache]) - when(metadataCache.getAliveBrokerNodes(any())).thenReturn(List(new Node(brokerId,"localhost", 0))) - when(metadataCache.getControllerId).thenReturn(None) - - // 2 topics returned for authorization in during handle - val topicsReturnedFromMetadataCacheForAuthorization = Set("remaining-topic", "later-deleted-topic") - when(metadataCache.getAllTopics()).thenReturn(topicsReturnedFromMetadataCacheForAuthorization) - // 1 topic is deleted from metadata right at the time between authorization and the next getTopicMetadata() call - when(metadataCache.getTopicMetadata( - ArgumentMatchers.eq(topicsReturnedFromMetadataCacheForAuthorization), - any[ListenerName], - anyBoolean, - anyBoolean - )).thenReturn(Seq( - new MetadataResponseTopic() - .setErrorCode(Errors.NONE.code) - .setName("remaining-topic") - .setIsInternal(false) - )) - - - var createTopicIsCalled: Boolean = false - // Specific mock on zkClient for this use case - // Expect it's never called to do auto topic creation - when(zkClient.setOrCreateEntityConfigs( - ArgumentMatchers.eq(ConfigType.TOPIC), - anyString, - any[Properties] - )).thenAnswer(_ => { - createTopicIsCalled = true - }) - // No need to use - when(zkClient.getAllBrokersInCluster) - .thenReturn(Seq(new Broker( - brokerId, "localhost", 9902, - ListenerName.forSecurityProtocol(SecurityProtocol.PLAINTEXT), SecurityProtocol.PLAINTEXT - ))) - - - val (requestListener, _) = updateMetadataCacheWithInconsistentListeners() - val response = sendMetadataRequestWithInconsistentListeners(requestListener) - - assertFalse(createTopicIsCalled) - val responseTopics = response.topicMetadata().asScala.map { metadata => metadata.topic() } - assertEquals(List("remaining-topic"), responseTopics) - assertTrue(response.topicsByError(Errors.UNKNOWN_TOPIC_OR_PARTITION).isEmpty) - } - - @Test - def testUnauthorizedTopicMetadataRequest(): Unit = { - // 1. Set up broker information - val plaintextListener = ListenerName.forSecurityProtocol(SecurityProtocol.PLAINTEXT) - val broker = new UpdateMetadataBroker() - .setId(0) - .setRack("rack") - .setEndpoints(Seq( - new UpdateMetadataEndpoint() - .setHost("broker0") - .setPort(9092) - .setSecurityProtocol(SecurityProtocol.PLAINTEXT.id) - .setListener(plaintextListener.value) - ).asJava) - - // 2. Set up authorizer - val authorizer: Authorizer = mock(classOf[Authorizer]) - val unauthorizedTopic = "unauthorized-topic" - val authorizedTopic = "authorized-topic" - - val expectedActions = Seq( - new Action(AclOperation.DESCRIBE, new ResourcePattern(ResourceType.TOPIC, unauthorizedTopic, PatternType.LITERAL), 1, true, true), - new Action(AclOperation.DESCRIBE, new ResourcePattern(ResourceType.TOPIC, authorizedTopic, PatternType.LITERAL), 1, true, true) - ) - - when(authorizer.authorize(any[RequestContext], argThat((t: java.util.List[Action]) => t.containsAll(expectedActions.asJava)))) - .thenAnswer { invocation => - val actions = invocation.getArgument(1).asInstanceOf[util.List[Action]].asScala - actions.map { action => - if (action.resourcePattern().name().equals(authorizedTopic)) - AuthorizationResult.ALLOWED - else - AuthorizationResult.DENIED - }.asJava - } - - // 3. Set up MetadataCache - val authorizedTopicId = Uuid.randomUuid() - val unauthorizedTopicId = Uuid.randomUuid() - - val topicIds = new util.HashMap[String, Uuid]() - topicIds.put(authorizedTopic, authorizedTopicId) - topicIds.put(unauthorizedTopic, unauthorizedTopicId) - - def createDummyPartitionStates(topic: String) = { - new UpdateMetadataPartitionState() - .setTopicName(topic) - .setPartitionIndex(0) - .setControllerEpoch(0) - .setLeader(0) - .setLeaderEpoch(0) - .setReplicas(Collections.singletonList(0)) - .setZkVersion(0) - .setIsr(Collections.singletonList(0)) - } - - // Send UpdateMetadataReq to update MetadataCache - val partitionStates = Seq(unauthorizedTopic, authorizedTopic).map(createDummyPartitionStates) - - val updateMetadataRequest = new UpdateMetadataRequest.Builder(ApiKeys.UPDATE_METADATA.latestVersion, 0, - 0, 0, partitionStates.asJava, Seq(broker).asJava, topicIds).build() - metadataCache.asInstanceOf[ZkMetadataCache].updateMetadata(correlationId = 0, updateMetadataRequest) - - // 4. Send TopicMetadataReq using topicId - val metadataReqByTopicId = new MetadataRequest.Builder(util.Arrays.asList(authorizedTopicId, unauthorizedTopicId)).build() - val repByTopicId = buildRequest(metadataReqByTopicId, plaintextListener) - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), - any[Long])).thenReturn(0) - kafkaApis = createKafkaApis(authorizer = Some(authorizer)) - kafkaApis.handleTopicMetadataRequest(repByTopicId) - val metadataByTopicIdResp = verifyNoThrottling[MetadataResponse](repByTopicId) - - val metadataByTopicId = metadataByTopicIdResp.data().topics().asScala.groupBy(_.topicId()).map(kv => (kv._1, kv._2.head)) - - metadataByTopicId.foreach { case (topicId, metadataResponseTopic) => - if (topicId == unauthorizedTopicId) { - // Return an TOPIC_AUTHORIZATION_FAILED on unauthorized error regardless of leaking the existence of topic id - assertEquals(Errors.TOPIC_AUTHORIZATION_FAILED.code(), metadataResponseTopic.errorCode()) - // Do not return topic information on unauthorized error - assertNull(metadataResponseTopic.name()) - } else { - assertEquals(Errors.NONE.code(), metadataResponseTopic.errorCode()) - assertEquals(authorizedTopic, metadataResponseTopic.name()) - } - } - kafkaApis.close() - - // 4. Send TopicMetadataReq using topic name - reset(clientRequestQuotaManager, requestChannel) - val metadataReqByTopicName = new MetadataRequest.Builder(util.Arrays.asList(authorizedTopic, unauthorizedTopic), false).build() - val repByTopicName = buildRequest(metadataReqByTopicName, plaintextListener) - kafkaApis = createKafkaApis(authorizer = Some(authorizer)) - kafkaApis.handleTopicMetadataRequest(repByTopicName) - val metadataByTopicNameResp = verifyNoThrottling[MetadataResponse](repByTopicName) - - val metadataByTopicName = metadataByTopicNameResp.data().topics().asScala.groupBy(_.name()).map(kv => (kv._1, kv._2.head)) - - metadataByTopicName.foreach { case (topicName, metadataResponseTopic) => - if (topicName == unauthorizedTopic) { - assertEquals(Errors.TOPIC_AUTHORIZATION_FAILED.code(), metadataResponseTopic.errorCode()) - // Do not return topic Id on unauthorized error - assertEquals(Uuid.ZERO_UUID, metadataResponseTopic.topicId()) - } else { - assertEquals(Errors.NONE.code(), metadataResponseTopic.errorCode()) - assertEquals(authorizedTopicId, metadataResponseTopic.topicId()) - } - } - } - /** * Verifies that sending a fetch request with version 9 works correctly when * ReplicaManager.getLogConfig returns None. @@ -4023,7 +3815,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) val responseData = response.data() @@ -4106,7 +3898,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) var response = verifyNoThrottling[ShareFetchResponse](request) var responseData = response.data() @@ -4209,7 +4001,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) var response = verifyNoThrottling[ShareFetchResponse](request) var responseData = response.data() @@ -4292,7 +4084,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) val responseData = response.data() @@ -4369,7 +4161,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) val responseData = response.data() @@ -4433,7 +4225,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) val responseData = response.data() @@ -4489,7 +4281,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) val responseData = response.data() @@ -4560,7 +4352,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) var response = verifyNoThrottling[ShareFetchResponse](request) var responseData = response.data() @@ -4653,7 +4445,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) var response = verifyNoThrottling[ShareFetchResponse](request) var responseData = response.data() @@ -4800,7 +4592,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) var response = verifyNoThrottling[ShareFetchResponse](request) var responseData = response.data() @@ -5133,7 +4925,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) var response = verifyNoThrottling[ShareFetchResponse](request) var responseData = response.data() @@ -5498,7 +5290,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val fetchResult: Map[TopicIdPartition, ShareFetchResponseData.PartitionData] = kafkaApis.handleFetchFromShareFetchRequest( request, @@ -5644,7 +5436,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val fetchResult: Map[TopicIdPartition, ShareFetchResponseData.PartitionData] = kafkaApis.handleFetchFromShareFetchRequest( request, @@ -5787,7 +5579,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val fetchResult: Map[TopicIdPartition, ShareFetchResponseData.PartitionData] = kafkaApis.handleFetchFromShareFetchRequest( request, @@ -5956,7 +5748,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val fetchResult: Map[TopicIdPartition, ShareFetchResponseData.PartitionData] = kafkaApis.handleFetchFromShareFetchRequest( request, @@ -6129,7 +5921,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) var response = verifyNoThrottling[ShareFetchResponse](request) var responseData = response.data() @@ -6216,7 +6008,7 @@ class KafkaApisTest extends Logging { GroupCoordinatorConfig.NEW_GROUP_COORDINATOR_ENABLE_CONFIG -> "false", ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) @@ -6259,7 +6051,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "false"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) @@ -6312,7 +6104,7 @@ class KafkaApisTest extends Logging { ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), authorizer = Option(authorizer), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) @@ -6375,7 +6167,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareFetchRequest(request) val response = verifyNoThrottling[ShareFetchResponse](request) val responseData = response.data() @@ -6442,7 +6234,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) val responseData = response.data() @@ -6490,7 +6282,7 @@ class KafkaApisTest extends Logging { GroupCoordinatorConfig.NEW_GROUP_COORDINATOR_ENABLE_CONFIG -> "false", ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) @@ -6532,7 +6324,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "false"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) @@ -6584,7 +6376,7 @@ class KafkaApisTest extends Logging { ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), authorizer = Option(authorizer), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) @@ -6635,7 +6427,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) @@ -6686,7 +6478,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) @@ -6735,7 +6527,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) @@ -6810,7 +6602,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) @@ -6873,7 +6665,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) val responseData = response.data() @@ -6940,7 +6732,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) val responseData = response.data() @@ -7008,7 +6800,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) kafkaApis.handleShareAcknowledgeRequest(request) val response = verifyNoThrottling[ShareAcknowledgeResponse](request) val responseData = response.data() @@ -7094,7 +6886,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val acknowledgeBatches = kafkaApis.getAcknowledgeBatchesFromShareFetchRequest(shareFetchRequest, topicNames, erroneous) assertEquals(4, acknowledgeBatches.size) @@ -7163,7 +6955,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val acknowledgeBatches = kafkaApis.getAcknowledgeBatchesFromShareFetchRequest(shareFetchRequest, topicIdNames, erroneous) val erroneousTopicIdPartitions = kafkaApis.validateAcknowledgementBatches(acknowledgeBatches, erroneous) @@ -7236,7 +7028,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val acknowledgeBatches = kafkaApis.getAcknowledgeBatchesFromShareAcknowledgeRequest(shareAcknowledgeRequest, topicNames, erroneous) assertEquals(3, acknowledgeBatches.size) @@ -7303,7 +7095,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val acknowledgeBatches = kafkaApis.getAcknowledgeBatchesFromShareAcknowledgeRequest(shareAcknowledgeRequest, topicIdNames, erroneous) val erroneousTopicIdPartitions = kafkaApis.validateAcknowledgementBatches(acknowledgeBatches, erroneous) @@ -7375,7 +7167,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val ackResult = kafkaApis.handleAcknowledgements( acknowledgementData, erroneous, @@ -7454,7 +7246,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val ackResult = kafkaApis.handleAcknowledgements( acknowledgementData, erroneous, @@ -7534,7 +7326,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val ackResult = kafkaApis.handleAcknowledgements( acknowledgementData, erroneous, @@ -7608,7 +7400,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val ackResult = kafkaApis.handleAcknowledgements( acknowledgementData, erroneous, @@ -7705,7 +7497,7 @@ class KafkaApisTest extends Logging { overrideProperties = Map( ServerConfigs.UNSTABLE_API_VERSIONS_ENABLE_CONFIG -> "true", ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true) + ) val response = kafkaApis.processShareAcknowledgeResponse(responseAcknowledgeData, request) val responseData = response.data() val topicResponses = responseData.responses() @@ -9790,14 +9582,14 @@ class KafkaApisTest extends Logging { @Test def testRaftShouldAlwaysForwardCreateAcls(): Unit = { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() verifyShouldAlwaysForwardErrorMessage(kafkaApis.handleCreateAcls) } @Test def testRaftShouldAlwaysForwardDeleteAcls(): Unit = { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() verifyShouldAlwaysForwardErrorMessage(kafkaApis.handleDeleteAcls) } @@ -9807,7 +9599,7 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), any[Long])).thenReturn(0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handleAlterConfigsRequest(request) val response = verifyNoThrottling[AlterConfigsResponse](request) assertEquals(new AlterConfigsResponseData(), response.data()) @@ -9827,7 +9619,7 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), any[Long])).thenReturn(0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handleAlterConfigsRequest(request) val response = verifyNoThrottling[AlterConfigsResponse](request) assertEquals(new AlterConfigsResponseData().setResponses(asList( @@ -9845,7 +9637,7 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), any[Long])).thenReturn(0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handleIncrementalAlterConfigsRequest(request) val response = verifyNoThrottling[IncrementalAlterConfigsResponse](request) assertEquals(new IncrementalAlterConfigsResponseData(), response.data()) @@ -9865,7 +9657,7 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), any[Long])).thenReturn(0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handleIncrementalAlterConfigsRequest(request) val response = verifyNoThrottling[IncrementalAlterConfigsResponse](request) assertEquals(new IncrementalAlterConfigsResponseData().setResponses(asList( @@ -9883,7 +9675,7 @@ class KafkaApisTest extends Logging { val requestChannelRequest = buildRequest(new ConsumerGroupHeartbeatRequest.Builder(consumerGroupHeartbeatRequest).build()) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) val expectedHeartbeatResponse = new ConsumerGroupHeartbeatResponseData() @@ -9907,7 +9699,7 @@ class KafkaApisTest extends Logging { )).thenReturn(future) kafkaApis = createKafkaApis( featureVersions = Seq(GroupVersion.GV_1), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -9934,7 +9726,7 @@ class KafkaApisTest extends Logging { )).thenReturn(future) kafkaApis = createKafkaApis( featureVersions = Seq(GroupVersion.GV_1), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -9957,7 +9749,7 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( authorizer = Some(authorizer), featureVersions = Seq(GroupVersion.GV_1), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -9983,7 +9775,7 @@ class KafkaApisTest extends Logging { )).thenReturn(future) kafkaApis = createKafkaApis( featureVersions = Seq(GroupVersion.GV_1), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10026,7 +9818,7 @@ class KafkaApisTest extends Logging { val expectedResponse = new ConsumerGroupDescribeResponseData() expectedResponse.groups.add(expectedDescribedGroup) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_1) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) val response = verifyNoThrottling[ConsumerGroupDescribeResponse](requestChannelRequest) @@ -10054,7 +9846,7 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( authorizer = Some(authorizer), featureVersions = Seq(GroupVersion.GV_1), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10077,7 +9869,7 @@ class KafkaApisTest extends Logging { )).thenReturn(future) kafkaApis = createKafkaApis( featureVersions = Seq(GroupVersion.GV_1), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10097,7 +9889,7 @@ class KafkaApisTest extends Logging { new GetTelemetrySubscriptionsResponseData())) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(request, RequestLocal.noCaching) val response = verifyNoThrottling[GetTelemetrySubscriptionsResponse](request) @@ -10116,7 +9908,7 @@ class KafkaApisTest extends Logging { any[RequestContext]())).thenThrow(new RuntimeException("test")) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(request, RequestLocal.noCaching) val response = verifyNoThrottling[GetTelemetrySubscriptionsResponse](request) @@ -10134,7 +9926,7 @@ class KafkaApisTest extends Logging { .thenReturn(new PushTelemetryResponse(new PushTelemetryResponseData())) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(request, RequestLocal.noCaching) val response = verifyNoThrottling[PushTelemetryResponse](request) @@ -10151,7 +9943,7 @@ class KafkaApisTest extends Logging { .thenThrow(new RuntimeException("test")) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(request, RequestLocal.noCaching) val response = verifyNoThrottling[PushTelemetryResponse](request) @@ -10168,7 +9960,7 @@ class KafkaApisTest extends Logging { resources.add("test1") resources.add("test2") when(clientMetricsManager.listClientMetricsResources).thenReturn(resources.asJava) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(request, RequestLocal.noCaching) val response = verifyNoThrottling[ListClientMetricsResourcesResponse](request) val expectedResponse = new ListClientMetricsResourcesResponseData().setClientMetricsResources( @@ -10183,7 +9975,7 @@ class KafkaApisTest extends Logging { val resources = new mutable.HashSet[String] when(clientMetricsManager.listClientMetricsResources).thenReturn(resources.asJava) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(request, RequestLocal.noCaching) val response = verifyNoThrottling[ListClientMetricsResourcesResponse](request) val expectedResponse = new ListClientMetricsResourcesResponseData() @@ -10196,7 +9988,7 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) when(clientMetricsManager.listClientMetricsResources).thenThrow(new RuntimeException("test")) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(request, RequestLocal.noCaching) val response = verifyNoThrottling[ListClientMetricsResourcesResponse](request) @@ -10210,7 +10002,7 @@ class KafkaApisTest extends Logging { val requestChannelRequest = buildRequest(new ShareGroupHeartbeatRequest.Builder(shareGroupHeartbeatRequest, true).build()) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) - kafkaApis = createKafkaApis(raftSupport = true) + kafkaApis = createKafkaApis() kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) val expectedHeartbeatResponse = new ShareGroupHeartbeatResponseData() @@ -10233,7 +10025,7 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) kafkaApis = createKafkaApis( overrideProperties = Map(ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10258,7 +10050,7 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( overrideProperties = Map(ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), authorizer = Some(authorizer), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10280,7 +10072,7 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) kafkaApis = createKafkaApis( overrideProperties = Map(ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10583,7 +10375,7 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( overrideProperties = configOverrides, authorizer = Option(authorizer), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10612,7 +10404,7 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( overrideProperties = configOverrides, authorizer = Option(authorizer), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching()) @@ -10642,7 +10434,7 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( overrideProperties = configOverrides, authorizer = Option(authorizer), - raftSupport = true + ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching()) From 4a2abc40952bdff951904336912b95255eb0c472 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Tue, 14 Jan 2025 23:28:24 +0000 Subject: [PATCH 10/13] fix describe testcase --- .../unit/kafka/server/KafkaApisTest.scala | 49 ++----------------- 1 file changed, 4 insertions(+), 45 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 44604b9b4fbad..03a476e2cd4ce 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -22,7 +22,7 @@ import kafka.coordinator.transaction.{InitProducerIdResult, TransactionCoordinat import kafka.log.UnifiedLog import kafka.network.RequestChannel import kafka.server.QuotaFactory.QuotaManagers -import kafka.server.metadata.{ConfigRepository, KRaftMetadataCache, MockConfigRepository, ZkMetadataCache} +import kafka.server.metadata.{ConfigRepository, KRaftMetadataCache, MockConfigRepository} import kafka.server.share.SharePartitionManager import kafka.utils.{CoreUtils, Log4jController, Logging, TestUtils} import org.apache.kafka.clients.admin.AlterConfigOp.OpType @@ -127,7 +127,6 @@ class KafkaApisTest extends Logging { } private val metrics = new Metrics() private val brokerId = 1 - // KRaft tests should override this with a KRaftMetadataCache private var metadataCache: MetadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) private val clientQuotaManager: ClientQuotaManager = mock(classOf[ClientQuotaManager]) private val clientRequestQuotaManager: ClientRequestQuotaManager = mock(classOf[ClientRequestQuotaManager]) @@ -262,19 +261,14 @@ class KafkaApisTest extends Logging { topicConfigs.put(propName, propValue) when(configRepository.topicConfig(resourceName)).thenReturn(topicConfigs) - metadataCache = mock(classOf[ZkMetadataCache]) - when(metadataCache.contains(resourceName)).thenReturn(true) - val describeConfigsRequest = new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() .setIncludeSynonyms(true) .setResources(List(new DescribeConfigsRequestData.DescribeConfigsResource() .setResourceName(resourceName) .setResourceType(ConfigResource.Type.TOPIC.id)).asJava)) .build(requestHeader.apiVersion) - val request = buildRequest(describeConfigsRequest, - requestHeader = Option(requestHeader)) - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), - any[Long])).thenReturn(0) + val request = buildRequest(describeConfigsRequest, requestHeader = Option(requestHeader)) + kafkaApis = createKafkaApis(authorizer = Some(authorizer), configRepository = configRepository) kafkaApis.handleDescribeConfigsRequest(request) @@ -285,11 +279,6 @@ class KafkaApisTest extends Logging { val describeConfigsResult = results.get(0) assertEquals(ConfigResource.Type.TOPIC.id, describeConfigsResult.resourceType) assertEquals(resourceName, describeConfigsResult.resourceName) - val configs = describeConfigsResult.configs.asScala.filter(_.name == propName) - assertEquals(1, configs.length) - val describeConfigsResponseData = configs.head - assertEquals(propName, describeConfigsResponseData.name) - assertEquals(propValue, describeConfigsResponseData.value) } @Test @@ -456,9 +445,6 @@ class KafkaApisTest extends Logging { val cmConfigs = ClientMetricsTestUtils.defaultProperties when(configRepository.config(resource)).thenReturn(cmConfigs) - metadataCache = mock(classOf[ZkMetadataCache]) - when(metadataCache.contains(subscriptionName)).thenReturn(true) - val describeConfigsRequest = new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() .setIncludeSynonyms(true) .setResources(List(new DescribeConfigsRequestData.DescribeConfigsResource() @@ -467,8 +453,7 @@ class KafkaApisTest extends Logging { .build(requestHeader.apiVersion) val request = buildRequest(describeConfigsRequest, requestHeader = Option(requestHeader)) - when(clientRequestQuotaManager.maybeRecordAndGetThrottleTimeMs(any[RequestChannel.Request](), - any[Long])).thenReturn(0) + kafkaApis = createKafkaApis(authorizer = Some(authorizer), configRepository = configRepository) kafkaApis.handleDescribeConfigsRequest(request) @@ -2203,7 +2188,6 @@ class KafkaApisTest extends Logging { @Test def testProduceResponseMetadataLookupErrorOnNotLeaderOrFollower(): Unit = { val topic = "topic" - metadataCache = mock(classOf[ZkMetadataCache]) for (version <- 10 to ApiKeys.PRODUCE.latestVersion) { @@ -8898,30 +8882,6 @@ class KafkaApisTest extends Logging { @Test def testDescribeClusterRequest(): Unit = { val plaintextListener = ListenerName.forSecurityProtocol(SecurityProtocol.PLAINTEXT) - val brokers = Seq( - new UpdateMetadataBroker() - .setId(0) - .setRack("rack") - .setEndpoints(Seq( - new UpdateMetadataEndpoint() - .setHost("broker0") - .setPort(9092) - .setSecurityProtocol(SecurityProtocol.PLAINTEXT.id) - .setListener(plaintextListener.value) - ).asJava), - new UpdateMetadataBroker() - .setId(1) - .setRack("rack") - .setEndpoints(Seq( - new UpdateMetadataEndpoint() - .setHost("broker1") - .setPort(9092) - .setSecurityProtocol(SecurityProtocol.PLAINTEXT.id) - .setListener(plaintextListener.value)).asJava) - ) - val updateMetadataRequest = new UpdateMetadataRequest.Builder(ApiKeys.UPDATE_METADATA.latestVersion, 0, - 0, 0, Seq.empty[UpdateMetadataPartitionState].asJava, brokers.asJava, Collections.emptyMap()).build() - MetadataCacheTest.updateCache(metadataCache, updateMetadataRequest) val describeClusterRequest = new DescribeClusterRequest.Builder(new DescribeClusterRequestData() .setIncludeClusterAuthorizedOperations(true)).build() @@ -8932,7 +8892,6 @@ class KafkaApisTest extends Logging { val describeClusterResponse = verifyNoThrottling[DescribeClusterResponse](request) - assertEquals(metadataCache.getControllerId.get.id, describeClusterResponse.data.controllerId) assertEquals(clusterId, describeClusterResponse.data.clusterId) assertEquals(8096, describeClusterResponse.data.clusterAuthorizedOperations) assertEquals(metadataCache.getAliveBrokerNodes(plaintextListener).toSet, From 047b385bc250b764a55ca475a7d824bc0ad80bf9 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Tue, 14 Jan 2025 23:47:40 +0000 Subject: [PATCH 11/13] fix testProduceResponseMetadataLookupErrorOnNotLeaderOrFollower --- core/src/test/scala/unit/kafka/server/KafkaApisTest.scala | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 03a476e2cd4ce..3f8c48d09ac10 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -176,7 +176,7 @@ class KafkaApisTest extends Logging { val listenerType = ListenerType.BROKER - val enabledApis = ApiKeys.apisForListener(listenerType).asScala ++ Set(ApiKeys.ENVELOPE) + val enabledApis = ApiKeys.apisForListener(listenerType).asScala val apiVersionManager = new SimpleApiVersionManager( listenerType, @@ -2189,6 +2189,7 @@ class KafkaApisTest extends Logging { def testProduceResponseMetadataLookupErrorOnNotLeaderOrFollower(): Unit = { val topic = "topic" + metadataCache = mock(classOf[KRaftMetadataCache]) for (version <- 10 to ApiKeys.PRODUCE.latestVersion) { reset(replicaManager, clientQuotaManager, clientRequestQuotaManager, requestChannel, txnCoordinator) From 6f3965c0e06c05e5ac5fdcb71c8400248cfda82b Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Wed, 15 Jan 2025 11:59:49 +0000 Subject: [PATCH 12/13] fix merge --- .../unit/kafka/server/KafkaApisTest.scala | 20 ++++--------------- 1 file changed, 4 insertions(+), 16 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 988d8367bed03..099add1e92eec 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -313,11 +313,8 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) -<<<<<<< HEAD + createKafkaApis(authorizer = Some(authorizer)).handleIncrementalAlterConfigsRequest(request) -======= - createKafkaApis(authorizer = Some(authorizer), raftSupport = true).handleIncrementalAlterConfigsRequest(request) ->>>>>>> trunk verify(forwardingManager, times(1)).forwardRequest( any(), any(), @@ -395,11 +392,8 @@ class KafkaApisTest extends Logging { val request = buildRequest(apiRequest) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) -<<<<<<< HEAD + kafkaApis = createKafkaApis() -======= - kafkaApis = createKafkaApis(raftSupport = true) ->>>>>>> trunk kafkaApis.handleAlterConfigsRequest(request) verify(forwardingManager, times(1)).forwardRequest( any(), @@ -421,11 +415,8 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) -<<<<<<< HEAD + kafkaApis = createKafkaApis() -======= - kafkaApis = createKafkaApis(raftSupport = true) ->>>>>>> trunk kafkaApis.handleIncrementalAlterConfigsRequest(request) verify(forwardingManager, times(1)).forwardRequest( any(), @@ -579,11 +570,8 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) -<<<<<<< HEAD + kafkaApis = createKafkaApis(authorizer = Some(authorizer)) -======= - kafkaApis = createKafkaApis(authorizer = Some(authorizer), raftSupport = true) ->>>>>>> trunk kafkaApis.handleIncrementalAlterConfigsRequest(request) verify(authorizer, times(1)).authorize(any(), any()) From dcc2a38bd43a73b5bbbca6872326277806378e56 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Wed, 15 Jan 2025 12:29:32 +0000 Subject: [PATCH 13/13] revert some change --- .../unit/kafka/server/KafkaApisTest.scala | 45 +++++++++++++------ 1 file changed, 31 insertions(+), 14 deletions(-) diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index 099add1e92eec..d778c7859821a 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -261,6 +261,9 @@ class KafkaApisTest extends Logging { topicConfigs.put(propName, propValue) when(configRepository.topicConfig(resourceName)).thenReturn(topicConfigs) + metadataCache = mock(classOf[KRaftMetadataCache]) + when(metadataCache.contains(resourceName)).thenReturn(true) + val describeConfigsRequest = new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() .setIncludeSynonyms(true) .setResources(List(new DescribeConfigsRequestData.DescribeConfigsResource() @@ -279,14 +282,11 @@ class KafkaApisTest extends Logging { val describeConfigsResult = results.get(0) assertEquals(ConfigResource.Type.TOPIC.id, describeConfigsResult.resourceType) assertEquals(resourceName, describeConfigsResult.resourceName) -<<<<<<< HEAD -======= val configs = describeConfigsResult.configs.asScala.filter(_.name == propName) assertEquals(1, configs.length) val describeConfigsResponseData = configs.head assertEquals(propName, describeConfigsResponseData.name) assertEquals(propValue, describeConfigsResponseData.value) ->>>>>>> trunk } @Test @@ -313,7 +313,6 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - createKafkaApis(authorizer = Some(authorizer)).handleIncrementalAlterConfigsRequest(request) verify(forwardingManager, times(1)).forwardRequest( any(), @@ -392,7 +391,6 @@ class KafkaApisTest extends Logging { val request = buildRequest(apiRequest) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - kafkaApis = createKafkaApis() kafkaApis.handleAlterConfigsRequest(request) verify(forwardingManager, times(1)).forwardRequest( @@ -415,7 +413,6 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - kafkaApis = createKafkaApis() kafkaApis.handleIncrementalAlterConfigsRequest(request) verify(forwardingManager, times(1)).forwardRequest( @@ -456,6 +453,9 @@ class KafkaApisTest extends Logging { val cmConfigs = ClientMetricsTestUtils.defaultProperties when(configRepository.config(resource)).thenReturn(cmConfigs) + metadataCache = mock(classOf[KRaftMetadataCache]) + when(metadataCache.contains(subscriptionName)).thenReturn(true) + val describeConfigsRequest = new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() .setIncludeSynonyms(true) .setResources(List(new DescribeConfigsRequestData.DescribeConfigsResource() @@ -570,7 +570,6 @@ class KafkaApisTest extends Logging { val request = buildRequest(incrementalAlterConfigsRequest, requestHeader = Option(requestHeader)) metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.LATEST_PRODUCTION) - kafkaApis = createKafkaApis(authorizer = Some(authorizer)) kafkaApis.handleIncrementalAlterConfigsRequest(request) @@ -2200,8 +2199,8 @@ class KafkaApisTest extends Logging { @Test def testProduceResponseMetadataLookupErrorOnNotLeaderOrFollower(): Unit = { val topic = "topic" - metadataCache = mock(classOf[KRaftMetadataCache]) + for (version <- 10 to ApiKeys.PRODUCE.latestVersion) { reset(replicaManager, clientQuotaManager, clientRequestQuotaManager, requestChannel, txnCoordinator) @@ -8895,6 +8894,30 @@ class KafkaApisTest extends Logging { @Test def testDescribeClusterRequest(): Unit = { val plaintextListener = ListenerName.forSecurityProtocol(SecurityProtocol.PLAINTEXT) + val brokers = Seq( + new UpdateMetadataBroker() + .setId(0) + .setRack("rack") + .setEndpoints(Seq( + new UpdateMetadataEndpoint() + .setHost("broker0") + .setPort(9092) + .setSecurityProtocol(SecurityProtocol.PLAINTEXT.id) + .setListener(plaintextListener.value) + ).asJava), + new UpdateMetadataBroker() + .setId(1) + .setRack("rack") + .setEndpoints(Seq( + new UpdateMetadataEndpoint() + .setHost("broker1") + .setPort(9092) + .setSecurityProtocol(SecurityProtocol.PLAINTEXT.id) + .setListener(plaintextListener.value)).asJava) + ) + val updateMetadataRequest = new UpdateMetadataRequest.Builder(ApiKeys.UPDATE_METADATA.latestVersion, 0, + 0, 0, Seq.empty[UpdateMetadataPartitionState].asJava, brokers.asJava, Collections.emptyMap()).build() + MetadataCacheTest.updateCache(metadataCache, updateMetadataRequest) val describeClusterRequest = new DescribeClusterRequest.Builder(new DescribeClusterRequestData() .setIncludeClusterAuthorizedOperations(true)).build() @@ -9969,7 +9992,6 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) kafkaApis = createKafkaApis( overrideProperties = Map(ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -9994,7 +10016,6 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( overrideProperties = Map(ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), authorizer = Some(authorizer), - ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10016,7 +10037,6 @@ class KafkaApisTest extends Logging { metadataCache = MetadataCache.kRaftMetadataCache(brokerId, () => KRaftVersion.KRAFT_VERSION_0) kafkaApis = createKafkaApis( overrideProperties = Map(ShareGroupConfig.SHARE_GROUP_ENABLE_CONFIG -> "true"), - ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10319,7 +10339,6 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( overrideProperties = configOverrides, authorizer = Option(authorizer), - ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching) @@ -10348,7 +10367,6 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( overrideProperties = configOverrides, authorizer = Option(authorizer), - ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching()) @@ -10378,7 +10396,6 @@ class KafkaApisTest extends Logging { kafkaApis = createKafkaApis( overrideProperties = configOverrides, authorizer = Option(authorizer), - ) kafkaApis.handle(requestChannelRequest, RequestLocal.noCaching())