diff --git a/clients/src/main/java/org/apache/kafka/server/authorizer/Authorizer.java b/clients/src/main/java/org/apache/kafka/server/authorizer/Authorizer.java index 45bd6d939928d..1865e7e6ad4a4 100644 --- a/clients/src/main/java/org/apache/kafka/server/authorizer/Authorizer.java +++ b/clients/src/main/java/org/apache/kafka/server/authorizer/Authorizer.java @@ -114,6 +114,8 @@ public interface Authorizer extends Configurable, Closeable { * This is an asynchronous API that enables the caller to avoid blocking during the update. Implementations of this * API can return completed futures using {@link java.util.concurrent.CompletableFuture#completedFuture(Object)} * to process the update synchronously on the request thread. + *

+ * Refer to the authorizer implementation docs for details on concurrent update guarantees. * * @param requestContext Request context if the ACL is being deleted by a broker to handle * a client request to delete ACLs. This may be null if ACLs are deleted directly in ZooKeeper diff --git a/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala b/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala index 5f2be9053515d..af6d01e97b620 100644 --- a/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala +++ b/core/src/main/scala/kafka/security/authorizer/AclAuthorizer.scala @@ -216,6 +216,22 @@ class AclAuthorizer extends Authorizer with Logging { results.toList.map(CompletableFuture.completedFuture[AclCreateResult]).asJava } + /** + * + * Concurrent updates: + *

+ */ override def deleteAcls(requestContext: AuthorizableRequestContext, aclBindingFilters: util.List[AclBindingFilter]): util.List[_ <: CompletionStage[AclDeleteResult]] = { val deletedBindings = new mutable.HashMap[AclBinding, Int]() @@ -542,12 +558,16 @@ class AclAuthorizer extends Authorizer with Logging { } } + private[authorizer] def processAclChangeNotification(resource: ResourcePattern): Unit = { + lock synchronized { + val versionedAcls = getAclsFromZk(resource) + updateCache(resource, versionedAcls) + } + } + object AclChangedNotificationHandler extends AclChangeNotificationHandler { override def processNotification(resource: ResourcePattern): Unit = { - lock synchronized { - val versionedAcls = getAclsFromZk(resource) - updateCache(resource, versionedAcls) - } + processAclChangeNotification(resource) } } } 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 9403870a1caff..89d2095e6edbb 100644 --- a/core/src/test/scala/unit/kafka/security/authorizer/AclAuthorizerTest.scala +++ b/core/src/test/scala/unit/kafka/security/authorizer/AclAuthorizerTest.scala @@ -20,7 +20,7 @@ import java.io.File import java.net.InetAddress import java.nio.charset.StandardCharsets.UTF_8 import java.nio.file.Files -import java.util.UUID +import java.util.{Collections, UUID} import java.util.concurrent.{Executors, Semaphore, TimeUnit} import kafka.Kafka @@ -912,6 +912,82 @@ class AclAuthorizerTest extends ZooKeeperTestHarness { }) } + @Test + def testCreateDeleteTiming(): Unit = { + val literalResource = new ResourcePattern(TOPIC, "foo-" + UUID.randomUUID(), LITERAL) + val prefixedResource = new ResourcePattern(TOPIC, "bar-", PREFIXED) + val wildcardResource = new ResourcePattern(TOPIC, "*", LITERAL) + val ace = new AccessControlEntry(principal.toString, WildcardHost, READ, ALLOW) + val updateSemaphore = new Semaphore(1) + + def createAcl(createAuthorizer: AclAuthorizer, resource: ResourcePattern): AclBinding = { + val acl = new AclBinding(resource, ace) + createAuthorizer.createAcls(requestContext, Collections.singletonList(acl)).asScala + .foreach(_.toCompletableFuture.get(15, TimeUnit.SECONDS)) + acl + } + + def deleteAcl(deleteAuthorizer: AclAuthorizer, + resource: ResourcePattern, + deletePatternType: PatternType): List[AclBinding] = { + + val filter = new AclBindingFilter( + new ResourcePatternFilter(resource.resourceType(), resource.name(), deletePatternType), + AccessControlEntryFilter.ANY) + deleteAuthorizer.deleteAcls(requestContext, Collections.singletonList(filter)).asScala + .map(_.toCompletableFuture.get(15, TimeUnit.SECONDS)) + .flatMap(_.aclBindingDeleteResults.asScala) + .map(_.aclBinding) + .toList + } + + def listAcls(authorizer: AclAuthorizer): List[AclBinding] = { + authorizer.acls(AclBindingFilter.ANY).asScala.toList + } + + def verifyCreateDeleteAcl(deleteAuthorizer: AclAuthorizer, + resource: ResourcePattern, + deletePatternType: PatternType): Unit = { + updateSemaphore.acquire() + assertEquals(List.empty, listAcls(deleteAuthorizer)) + val acl = createAcl(aclAuthorizer, resource) + val deleted = deleteAcl(deleteAuthorizer, resource, deletePatternType) + if (deletePatternType != PatternType.MATCH) { + assertEquals(List(acl), deleted) + } else { + assertEquals(List.empty[AclBinding], deleted) + } + updateSemaphore.release() + if (deletePatternType == PatternType.MATCH) { + TestUtils.waitUntilTrue(() => listAcls(deleteAuthorizer).nonEmpty, "ACL not propagated") + assertEquals(List(acl), deleteAcl(deleteAuthorizer, resource, deletePatternType)) + } + TestUtils.waitUntilTrue(() => listAcls(deleteAuthorizer).isEmpty, "ACL delete not propagated") + } + + val deleteAuthorizer = new AclAuthorizer { + override def processAclChangeNotification(resource: ResourcePattern): Unit = { + updateSemaphore.acquire() + try { + super.processAclChangeNotification(resource) + } finally { + updateSemaphore.release() + } + } + } + + try { + deleteAuthorizer.configure(config.originals) + List(literalResource, prefixedResource, wildcardResource).foreach { resource => + verifyCreateDeleteAcl(deleteAuthorizer, resource, resource.patternType()) + verifyCreateDeleteAcl(deleteAuthorizer, resource, PatternType.ANY) + verifyCreateDeleteAcl(deleteAuthorizer, resource, PatternType.MATCH) + } + } finally { + deleteAuthorizer.close() + } + } + private def givenAuthorizerWithProtocolVersion(protocolVersion: Option[ApiVersion]): Unit = { aclAuthorizer.close()