Skip to content
98 changes: 3 additions & 95 deletions core/src/main/scala/kafka/server/KafkaApis.scala
Original file line number Diff line number Diff line change
Expand Up @@ -27,16 +27,13 @@ import kafka.server.share.SharePartitionManager
import kafka.utils.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
Expand Down Expand Up @@ -2520,45 +2517,11 @@ 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 {
sendResponse(Some(processLegacyAlterConfigsRequest(request, remaining)))
}
}

def processLegacyAlterConfigsRequest(
originalRequest: RequestChannel.Request,
data: AlterConfigsRequestData
): AlterConfigsResponseData = {
val zkSupport = metadataSupport.requireZkOrThrow(KafkaApis.shouldAlwaysForward(originalRequest))
Comment thread
TaiJuWu marked this conversation as resolved.
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 = {
Expand Down Expand Up @@ -2647,31 +2610,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))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

For my understanding (and here I am working under the assumption that these requests ought to be authorised on the KRaft controller now), should this be removed as well?

@TaiJuWu TaiJuWu Jan 10, 2025

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

when we set broker config (only for LOG4J), we also need to do authorization locally so we need to keep this one.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Got it, thank you!

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) =>
Expand All @@ -2684,49 +2628,13 @@ 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 {
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))
Comment thread
TaiJuWu marked this conversation as resolved.
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)
}
new IncrementalAlterConfigsResponse(0, (authorizedResult ++ unauthorizedResult).asJava).data()
}

def handleDescribeConfigsRequest(request: RequestChannel.Request): Unit = {
Expand Down
Loading