diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/CreateAclsResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/CreateAclsResult.java index 6e69554635efc..efe1e08a9e407 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/CreateAclsResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/CreateAclsResult.java @@ -33,7 +33,8 @@ public class CreateAclsResult { private final Map> futures; - CreateAclsResult(Map> futures) { + // Visible for testing + public CreateAclsResult(Map> futures) { this.futures = futures; } diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java index b97fc7bc577ee..315ee6e414f89 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java @@ -20,10 +20,12 @@ import java.util.Locale; import java.util.Map.Entry; +import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicReference; import org.apache.kafka.clients.admin.Admin; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.IsolationLevel; +import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.config.ConfigValue; import org.apache.kafka.common.errors.SecurityDisabledException; import org.apache.kafka.connect.connector.Task; @@ -74,6 +76,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import static org.apache.kafka.common.utils.Utils.sleep; import static org.apache.kafka.connect.mirror.MirrorSourceConfig.SYNC_TOPIC_ACLS_ENABLED; /** Replicate data, configuration, and ACLs between clusters. @@ -104,6 +107,7 @@ public class MirrorSourceConnector extends SourceConnector { private Admin offsetSyncsAdminClient; private volatile boolean useIncrementalAlterConfigs; private AtomicBoolean noAclAuthorizer = new AtomicBoolean(false); + private Set knownTopicAclBindings = Collections.emptySet(); public MirrorSourceConnector() { // nop @@ -581,13 +585,30 @@ void incrementalAlterConfigs(Map topicConfigs) { })); } - private void updateTopicAcls(List bindings) { - log.trace("Syncing {} topic ACL bindings.", bindings.size()); - targetAdminClient.createAcls(bindings).values().forEach((k, v) -> v.whenComplete((x, e) -> { - if (e != null) { - log.warn("Could not sync ACL of topic {}.", k.pattern().name(), e); + // Visible for testing + int updateTopicAcls(List bindings) { + Set addBindings = new HashSet<>(bindings); + Set failedBindings = new HashSet<>(); + addBindings.removeAll(knownTopicAclBindings); + int newBindCount = addBindings.size(); + if (!addBindings.isEmpty()) { + log.info("Syncing new found {} topic ACL bindings.", newBindCount); + Map> futureMap = targetAdminClient.createAcls(addBindings).values(); + futureMap.forEach((k, v) -> v.whenComplete((x, e) -> { + if (e != null) { + log.warn("Could not sync ACL of topic {}.", k.pattern().name(), e); + failedBindings.add(k); + } + })); + while (!futureMap.values().stream().allMatch(Future::isDone)) { + sleep(1000); } - })); + knownTopicAclBindings = new HashSet<>(bindings); + knownTopicAclBindings.removeAll(failedBindings); + } else { + log.debug("Not found new topic Acl info, skip sync!"); + } + return newBindCount - failedBindings.size(); } private static Stream expandTopicDescription(TopicDescription description) { @@ -672,4 +693,9 @@ boolean isCycle(String topic) { String formatRemoteTopic(String topic) { return replicationPolicy.formatRemoteTopic(sourceAndTarget.source(), topic); } + + // Visible for testing + Set knownTopicAclBindings() { + return knownTopicAclBindings; + } } diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java index 5e626679e9df7..793a32dc88caa 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java @@ -18,6 +18,7 @@ import org.apache.kafka.clients.admin.AlterConfigOp; import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.CreateAclsResult; import org.apache.kafka.clients.admin.DescribeAclsResult; import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.TopicPartition; @@ -27,8 +28,10 @@ import org.apache.kafka.common.acl.AclPermissionType; import org.apache.kafka.common.config.ConfigResource; import org.apache.kafka.common.config.ConfigValue; +import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.errors.SecurityDisabledException; +import org.apache.kafka.common.internals.KafkaFutureImpl; import org.apache.kafka.common.resource.PatternType; import org.apache.kafka.common.resource.ResourcePattern; import org.apache.kafka.common.resource.ResourceType; @@ -41,9 +44,11 @@ import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.source.ExactlyOnceSupport; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentMatchers; import static org.apache.kafka.clients.admin.AdminClientTestUtils.alterConfigsResult; import static org.apache.kafka.clients.consumer.ConsumerConfig.ISOLATION_LEVEL_CONFIG; +import static org.apache.kafka.common.utils.Utils.sleep; import static org.apache.kafka.connect.mirror.MirrorConnectorConfig.CONSUMER_CLIENT_PREFIX; import static org.apache.kafka.connect.mirror.MirrorConnectorConfig.SOURCE_PREFIX; import static org.apache.kafka.connect.mirror.MirrorSourceConfig.OFFSET_LAG_MAX; @@ -55,6 +60,7 @@ import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.mockito.ArgumentMatchers.anySet; import static org.mockito.ArgumentMatchers.isA; import static org.mockito.Mockito.any; import static org.mockito.Mockito.doAnswer; @@ -74,10 +80,14 @@ import java.util.Collection; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.Set; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.stream.Collectors; public class MirrorSourceConnectorTest { @@ -683,4 +693,65 @@ private Optional validateProperty(String name, Map assertNotNull(result, "Connector should not have record null config value for '" + name + "' property"); return Optional.of(result); } + + @Test + public void testUpdateIncrementTopicAcls() { + Admin sourceAdmin = mock(Admin.class); + Admin targetAdmin = mock(Admin.class); + MirrorSourceConnector connector = new MirrorSourceConnector(sourceAdmin, targetAdmin); + + List filteredBindings = new ArrayList<>(); + AclBinding binding1 = mock(AclBinding.class); + AclBinding binding2 = mock(AclBinding.class); + filteredBindings.add(binding1); + filteredBindings.add(binding2); + doReturn(mock(CreateAclsResult.class)).when(targetAdmin).createAcls(anySet()); + + // First topic acl info update when starting `syncTopicAcls` thread + int newAddCount = connector.updateTopicAcls(filteredBindings); + assertEquals(new HashSet<>(filteredBindings), connector.knownTopicAclBindings()); + assertEquals(filteredBindings.size(), newAddCount); + + List newAddBindings = new ArrayList<>(); + AclBinding binding3 = mock(AclBinding.class); + AclBinding binding4 = mock(AclBinding.class); + newAddBindings.add(binding3); + newAddBindings.add(binding4); + filteredBindings.addAll(newAddBindings); + + // The next increment topic acl info update + newAddCount = connector.updateTopicAcls(filteredBindings); + assertEquals(new HashSet<>(filteredBindings), connector.knownTopicAclBindings()); + assertEquals(newAddBindings.size(), newAddCount); + + // The next increment topic acl info update, contains failed create + List newAddFailedBindings = new ArrayList<>(); + AclBinding binding5 = mock(AclBinding.class); + AclBinding binding6 = mock(AclBinding.class); + newAddFailedBindings.add(binding5); + newAddFailedBindings.add(binding6); + filteredBindings.addAll(newAddFailedBindings); + + Map> futures = new HashMap<>(); + KafkaFutureImpl futureForBinding5 = new KafkaFutureImpl<>(); + KafkaFutureImpl futureForBinding6 = new KafkaFutureImpl<>(); + ExecutorService singleThread = Executors.newSingleThreadExecutor(); + // Delayed completion of `createAclRequest` for simulating actual scenarios + singleThread.submit(() -> { + sleep(5000); + futureForBinding5.complete(null); + futureForBinding6.completeExceptionally(new ApiException("mock create acl failure.")); + }); + futures.put(binding5, futureForBinding5); + futures.put(binding6, futureForBinding6); + CreateAclsResult mockCreateAclsResult = new CreateAclsResult(new HashMap<>(futures)); + doReturn(new ResourcePattern(ResourceType.TOPIC, "topic6", PatternType.LITERAL)).when(binding6).pattern(); + doReturn(mockCreateAclsResult).when(targetAdmin).createAcls(ArgumentMatchers.eq(new HashSet<>(newAddFailedBindings))); + + int newAddSuccessCount = connector.updateTopicAcls(filteredBindings); + Set bindingToVerify = new HashSet<>(filteredBindings); + bindingToVerify.remove(binding6); + assertEquals(bindingToVerify, connector.knownTopicAclBindings()); + assertEquals(1, newAddSuccessCount); + } }