From 52c247f777146ee1aa3cdd821138ca2743dc0de2 Mon Sep 17 00:00:00 2001 From: Aman Singh Date: Wed, 29 Jun 2022 11:26:09 +0530 Subject: [PATCH 1/3] santize and desantize resource name in acls --- core/src/main/scala/kafka/zk/KafkaZkClient.scala | 4 ++-- core/src/main/scala/kafka/zk/ZkData.scala | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/zk/KafkaZkClient.scala b/core/src/main/scala/kafka/zk/KafkaZkClient.scala index fa7ce00882aee..ff7f2ccc6035c 100644 --- a/core/src/main/scala/kafka/zk/KafkaZkClient.scala +++ b/core/src/main/scala/kafka/zk/KafkaZkClient.scala @@ -32,7 +32,7 @@ import kafka.zookeeper._ import org.apache.kafka.common.errors.ControllerMovedException import org.apache.kafka.common.resource.{PatternType, ResourcePattern, ResourceType} import org.apache.kafka.common.security.token.delegation.{DelegationToken, TokenInformation} -import org.apache.kafka.common.utils.{Time, Utils} +import org.apache.kafka.common.utils.{Sanitizer, Time, Utils} import org.apache.kafka.common.{KafkaException, TopicPartition, Uuid} import org.apache.zookeeper.KeeperException.{Code, NodeExistsException} import org.apache.zookeeper.OpResult.{CreateResult, ErrorResult, SetDataResult} @@ -1307,7 +1307,7 @@ class KafkaZkClient private[zk] (zooKeeperClient: ZooKeeperClient, isSecure: Boo * @return list of resource names */ def getResourceNames(patternType: PatternType, resourceType: ResourceType): Seq[String] = { - getChildren(ZkAclStore(patternType).path(resourceType)) + getChildren(ZkAclStore(patternType).path(resourceType)).map(resourceName => Sanitizer.desanitize(resourceName)) } /** diff --git a/core/src/main/scala/kafka/zk/ZkData.scala b/core/src/main/scala/kafka/zk/ZkData.scala index 7006a21f94bfb..b55e23647717b 100644 --- a/core/src/main/scala/kafka/zk/ZkData.scala +++ b/core/src/main/scala/kafka/zk/ZkData.scala @@ -38,7 +38,7 @@ import org.apache.kafka.common.network.ListenerName import org.apache.kafka.common.resource.{PatternType, ResourcePattern, ResourceType} import org.apache.kafka.common.security.auth.SecurityProtocol import org.apache.kafka.common.security.token.delegation.{DelegationToken, TokenInformation} -import org.apache.kafka.common.utils.{SecurityUtils, Time} +import org.apache.kafka.common.utils.{Sanitizer, SecurityUtils, Time} import org.apache.kafka.common.{KafkaException, TopicPartition, Uuid} import org.apache.kafka.metadata.LeaderRecoveryState import org.apache.kafka.server.common.{MetadataVersion, ProducerIdsBlock} @@ -599,7 +599,7 @@ sealed trait ZkAclStore { def path(resourceType: ResourceType): String = s"$aclPath/${SecurityUtils.resourceTypeName(resourceType)}" - def path(resourceType: ResourceType, resourceName: String): String = s"$aclPath/${SecurityUtils.resourceTypeName(resourceType)}/$resourceName" + def path(resourceType: ResourceType, resourceName: String): String = s"$aclPath/${SecurityUtils.resourceTypeName(resourceType)}/${Sanitizer.sanitize(resourceName)}" def changeStore: ZkAclChangeStore } From 531fb845487fe6df4f739409c1dc3a440d39b38c Mon Sep 17 00:00:00 2001 From: Aman Singh Date: Wed, 6 Jul 2022 15:26:00 +0530 Subject: [PATCH 2/3] address review comment --- .../security/authorizer/AclAuthorizer.scala | 7 +++++ .../main/scala/kafka/zk/KafkaZkClient.scala | 4 +-- core/src/main/scala/kafka/zk/ZkData.scala | 4 +-- .../authorizer/AclAuthorizerTest.scala | 29 +++++++++++-------- 4 files changed, 28 insertions(+), 16 deletions(-) diff --git a/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala b/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala index 0c2b6f619f458..4dc62231a4e44 100644 --- a/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala +++ b/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala @@ -122,6 +122,10 @@ object AclAuthorizer { if (aclBinding.isUnknown) throw new IllegalArgumentException("ACL binding contains unknown elements") } + + private def inValidAclBindingResourceName(resourceName: String): Boolean = { + resourceName.contains("/") + } } class AclAuthorizer extends Authorizer with Logging { @@ -209,6 +213,9 @@ class AclAuthorizer extends Authorizer with Logging { throw new UnsupportedVersionException(s"Adding ACLs on prefixed resource patterns requires " + s"${KafkaConfig.InterBrokerProtocolVersionProp} of $IBP_2_0_IV1 or greater") } + if (inValidAclBindingResourceName(aclBinding.pattern().name())) { + throw new IllegalArgumentException(s"ACL binding contains invalid resource name: ${aclBinding.pattern().name()}") + } validateAclBinding(aclBinding) true } catch { diff --git a/core/src/main/scala/kafka/zk/KafkaZkClient.scala b/core/src/main/scala/kafka/zk/KafkaZkClient.scala index ff7f2ccc6035c..fa7ce00882aee 100644 --- a/core/src/main/scala/kafka/zk/KafkaZkClient.scala +++ b/core/src/main/scala/kafka/zk/KafkaZkClient.scala @@ -32,7 +32,7 @@ import kafka.zookeeper._ import org.apache.kafka.common.errors.ControllerMovedException import org.apache.kafka.common.resource.{PatternType, ResourcePattern, ResourceType} import org.apache.kafka.common.security.token.delegation.{DelegationToken, TokenInformation} -import org.apache.kafka.common.utils.{Sanitizer, Time, Utils} +import org.apache.kafka.common.utils.{Time, Utils} import org.apache.kafka.common.{KafkaException, TopicPartition, Uuid} import org.apache.zookeeper.KeeperException.{Code, NodeExistsException} import org.apache.zookeeper.OpResult.{CreateResult, ErrorResult, SetDataResult} @@ -1307,7 +1307,7 @@ class KafkaZkClient private[zk] (zooKeeperClient: ZooKeeperClient, isSecure: Boo * @return list of resource names */ def getResourceNames(patternType: PatternType, resourceType: ResourceType): Seq[String] = { - getChildren(ZkAclStore(patternType).path(resourceType)).map(resourceName => Sanitizer.desanitize(resourceName)) + getChildren(ZkAclStore(patternType).path(resourceType)) } /** diff --git a/core/src/main/scala/kafka/zk/ZkData.scala b/core/src/main/scala/kafka/zk/ZkData.scala index b55e23647717b..7006a21f94bfb 100644 --- a/core/src/main/scala/kafka/zk/ZkData.scala +++ b/core/src/main/scala/kafka/zk/ZkData.scala @@ -38,7 +38,7 @@ import org.apache.kafka.common.network.ListenerName import org.apache.kafka.common.resource.{PatternType, ResourcePattern, ResourceType} import org.apache.kafka.common.security.auth.SecurityProtocol import org.apache.kafka.common.security.token.delegation.{DelegationToken, TokenInformation} -import org.apache.kafka.common.utils.{Sanitizer, SecurityUtils, Time} +import org.apache.kafka.common.utils.{SecurityUtils, Time} import org.apache.kafka.common.{KafkaException, TopicPartition, Uuid} import org.apache.kafka.metadata.LeaderRecoveryState import org.apache.kafka.server.common.{MetadataVersion, ProducerIdsBlock} @@ -599,7 +599,7 @@ sealed trait ZkAclStore { def path(resourceType: ResourceType): String = s"$aclPath/${SecurityUtils.resourceTypeName(resourceType)}" - def path(resourceType: ResourceType, resourceName: String): String = s"$aclPath/${SecurityUtils.resourceTypeName(resourceType)}/${Sanitizer.sanitize(resourceName)}" + def path(resourceType: ResourceType, resourceName: String): String = s"$aclPath/${SecurityUtils.resourceTypeName(resourceType)}/$resourceName" def changeStore: ZkAclChangeStore } diff --git a/core/src/test/scala/unit/kafka/security/authorizer/AclAuthorizerTest.scala b/core/src/test/scala/unit/kafka/security/authorizer/AclAuthorizerTest.scala index ce7bca25d12fc..3be34921423dc 100644 --- a/core/src/test/scala/unit/kafka/security/authorizer/AclAuthorizerTest.scala +++ b/core/src/test/scala/unit/kafka/security/authorizer/AclAuthorizerTest.scala @@ -16,40 +16,39 @@ */ package kafka.security.authorizer -import java.io.File -import java.net.InetAddress -import java.nio.charset.StandardCharsets.UTF_8 -import java.nio.file.Files -import java.util.{Collections, UUID} -import java.util.concurrent.{Executors, Semaphore, TimeUnit} - import kafka.Kafka import kafka.security.authorizer.AclEntry.{WildcardHost, WildcardPrincipalString} import kafka.server.{KafkaConfig, QuorumTestHarness} import kafka.utils.TestUtils import kafka.zk.ZkAclStore import kafka.zookeeper.{GetChildrenRequest, GetDataRequest, ZooKeeperClient} -import org.apache.kafka.common.acl._ import org.apache.kafka.common.acl.AclOperation._ import org.apache.kafka.common.acl.AclPermissionType.{ALLOW, DENY} +import org.apache.kafka.common.acl._ import org.apache.kafka.common.errors.{ApiException, UnsupportedVersionException} import org.apache.kafka.common.requests.RequestContext -import org.apache.kafka.common.resource.{PatternType, ResourcePattern, ResourcePatternFilter, ResourceType} +import org.apache.kafka.common.resource.PatternType.{LITERAL, MATCH, PREFIXED} import org.apache.kafka.common.resource.Resource.CLUSTER_NAME import org.apache.kafka.common.resource.ResourcePattern.WILDCARD_RESOURCE import org.apache.kafka.common.resource.ResourceType._ -import org.apache.kafka.common.resource.PatternType.{LITERAL, MATCH, PREFIXED} +import org.apache.kafka.common.resource.{PatternType, ResourcePattern, ResourcePatternFilter, ResourceType} import org.apache.kafka.common.security.auth.KafkaPrincipal -import org.apache.kafka.server.authorizer._ import org.apache.kafka.common.utils.{Time, SecurityUtils => JSecurityUtils} +import org.apache.kafka.server.authorizer._ import org.apache.kafka.server.common.MetadataVersion import org.apache.kafka.server.common.MetadataVersion.{IBP_2_0_IV0, IBP_2_0_IV1} import org.apache.zookeeper.client.ZKClientConfig import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.{AfterEach, BeforeEach, Test, TestInfo} -import scala.jdk.CollectionConverters._ +import java.io.File +import java.net.InetAddress +import java.nio.charset.StandardCharsets.UTF_8 +import java.nio.file.Files +import java.util.concurrent.{Executors, Semaphore, TimeUnit} +import java.util.{Collections, UUID} import scala.collection.mutable +import scala.jdk.CollectionConverters._ class AclAuthorizerTest extends QuorumTestHarness with BaseAuthorizerTest { @@ -722,6 +721,12 @@ class AclAuthorizerTest extends QuorumTestHarness with BaseAuthorizerTest { assertTrue(e.getCause.isInstanceOf[UnsupportedVersionException], s"Unexpected exception $e") } + @Test + def testCreateAclWithInvalidResourceName(): Unit = { + assertThrows(classOf[ApiException], + () => addAcls(aclAuthorizer, Set(allowReadAcl), new ResourcePattern(TOPIC, "test/1", LITERAL))) + } + @Test def testWritesExtendedAclChangeEventIfInterBrokerProtocolNotSet(): Unit = { givenAuthorizerWithProtocolVersion(Option.empty) From 0b0cb0e1352029b6901b819ce9d40f8cae085456 Mon Sep 17 00:00:00 2001 From: Aman Singh Date: Thu, 7 Jul 2022 19:22:56 +0530 Subject: [PATCH 3/3] address review comment --- .../scala/kafka/security/authorizer/AclAuthorizer.scala | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) diff --git a/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala b/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala index 4dc62231a4e44..1de9a27402cb3 100644 --- a/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala +++ b/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala @@ -121,10 +121,8 @@ object AclAuthorizer { private def validateAclBinding(aclBinding: AclBinding): Unit = { if (aclBinding.isUnknown) throw new IllegalArgumentException("ACL binding contains unknown elements") - } - - private def inValidAclBindingResourceName(resourceName: String): Boolean = { - resourceName.contains("/") + if (aclBinding.pattern().name().contains("/")) + throw new IllegalArgumentException(s"ACL binding contains invalid resource name: ${aclBinding.pattern().name()}") } } @@ -213,9 +211,6 @@ class AclAuthorizer extends Authorizer with Logging { throw new UnsupportedVersionException(s"Adding ACLs on prefixed resource patterns requires " + s"${KafkaConfig.InterBrokerProtocolVersionProp} of $IBP_2_0_IV1 or greater") } - if (inValidAclBindingResourceName(aclBinding.pattern().name())) { - throw new IllegalArgumentException(s"ACL binding contains invalid resource name: ${aclBinding.pattern().name()}") - } validateAclBinding(aclBinding) true } catch {