From e59137bef825128e8b695f95b32e357ca412efc8 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Tue, 19 Oct 2021 22:14:01 +0800 Subject: [PATCH 1/5] Update pulsar dependencies to 2.9.0 --- .../pulsar/handlers/kop/KafkaRequestHandler.java | 8 ++------ .../handlers/kop/security/auth/SimpleAclAuthorizer.java | 7 +++---- .../pulsar/handlers/kop/utils/PulsarMessageBuilder.java | 2 +- .../handlers/kop/KafkaServiceConfigurationTest.java | 4 ++-- pom.xml | 2 +- 5 files changed, 9 insertions(+), 14 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index f49d287ae7..9a0e5c6d0f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -517,12 +517,8 @@ static CompletableFuture> expandAllowedNamespaces(Set allowe String tenant = namespace.substring(0, slash); results.add(pulsarService.getPulsarResources() .getNamespaceResources() - .getChildrenAsync(path(tenant)) - .thenAccept(children -> { - children.forEach(ns -> { - result.add(tenant + "/" + ns); - }); - })); + .listNamespacesAsync(tenant) + .thenAccept(namespaces -> namespaces.forEach(ns -> result.add(tenant + "/" + ns)))); } } return CompletableFuture diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java index e0aba7b9e8..bac2826a68 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java @@ -73,7 +73,6 @@ private CompletableFuture authorizeTopicPermission(KafkaPrincipal princ new IllegalArgumentException("Resource name must contains namespace.")); return permissionFuture; } - String policiesPath = path(namespace.toString()); String tenantName = namespace.getTenant(); isSuperUserOrTenantAdmin(tenantName, principal.getName()).whenComplete((isSuperUserOrAdmin, exception) -> { if (exception != null) { @@ -90,7 +89,7 @@ private CompletableFuture authorizeTopicPermission(KafkaPrincipal princ getPulsarService() .getPulsarResources() .getNamespaceResources() - .getAsync(policiesPath) + .getPoliciesAsync(namespace) .thenAccept(policies -> { if (!policies.isPresent()) { if (log.isDebugEnabled()) { @@ -177,7 +176,7 @@ private CompletableFuture authorizeTenantPermission(KafkaPrincipal prin getPulsarService() .getPulsarResources() .getTenantResources() - .getAsync(path(tenant)) + .getTenantAsync(tenant) .thenAccept(tenantInfo -> { permissionFuture.complete(tenantInfo.isPresent()); }).exceptionally(ex -> { @@ -234,7 +233,7 @@ private CompletableFuture isSuperUserOrTenantAdmin(String tenant, Strin if (ex != null || !isSuperUser) { pulsarService.getPulsarResources() .getTenantResources() - .getAsync(path(tenant)) + .getTenantAsync(tenant) .thenAccept(tenantInfo -> { if (!tenantInfo.isPresent()) { future.complete(false); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/PulsarMessageBuilder.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/PulsarMessageBuilder.java index fe8eb3a4d7..66624c99ac 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/PulsarMessageBuilder.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/PulsarMessageBuilder.java @@ -110,7 +110,7 @@ public PulsarMessageBuilder sequenceId(long sequenceId) { } public Message getMessage() { - return MessageImpl.create(metadata, content, SCHEMA); + return MessageImpl.create(metadata, content, SCHEMA, null); } } \ No newline at end of file diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfigurationTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfigurationTest.java index b0de3bca25..ff77bd7d43 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfigurationTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfigurationTest.java @@ -132,7 +132,7 @@ public void testConfigurationUtilsStream() throws Exception { assertNotNull(kafkaServiceConfig); assertEquals(kafkaServiceConfig.getZookeeperServers(), zkServer); assertEquals(kafkaServiceConfig.isBrokerDeleteInactiveTopicsEnabled(), true); - assertEquals(kafkaServiceConfig.getBacklogQuotaDefaultLimitGB(), 18); + assertEquals(kafkaServiceConfig.getBacklogQuotaDefaultLimitGB(), 18.0); assertEquals(kafkaServiceConfig.getClusterName(), "usc"); assertEquals(kafkaServiceConfig.getBrokerClientAuthenticationParameters(), "role:my-role"); assertEquals(kafkaServiceConfig.getBrokerServicePort().get(), new Integer(7777)); @@ -229,7 +229,7 @@ public void testAllowedNamespaces() throws Exception { NamespaceResources namespaceResources = mock(NamespaceResources.class); when(pulsarService.getPulsarResources()).thenReturn(pulsarResources); when(pulsarResources.getNamespaceResources()).thenReturn(namespaceResources); - when(namespaceResources.getChildrenAsync(any(String.class))) + when(namespaceResources.listNamespacesAsync(any(String.class))) .thenReturn(CompletableFuture.completedFuture(Arrays.asList("one", "two"))); assertEquals(KafkaRequestHandler.expandAllowedNamespaces(conf.getKopAllowedNamespaces(), "logged", pulsarService).get(), diff --git a/pom.xml b/pom.xml index b0a68954f1..529ef35203 100644 --- a/pom.xml +++ b/pom.xml @@ -46,7 +46,7 @@ 1.18.4 2.22.0 io.streamnative - 2.8.1.0 + 2.9.0-rc-202110182205 1.7.25 3.1.8 1.15.1 From bac7dace491c22baa2e864770c5c5d8bdba01dfa Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 20 Oct 2021 10:23:36 +0800 Subject: [PATCH 2/5] Fix checkstyle --- .../pulsar/handlers/kop/KafkaRequestHandler.java | 8 -------- 1 file changed, 8 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index 9a0e5c6d0f..c08a1f0447 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -28,7 +28,6 @@ import static org.apache.kafka.common.requests.CreateTopicsRequest.TopicDetails; import com.google.common.annotations.VisibleForTesting; -import com.google.common.base.Joiner; import com.google.common.collect.Lists; import com.google.common.collect.Maps; import com.google.common.collect.Sets; @@ -489,13 +488,6 @@ private static boolean isInternalTopic(final String fullTopicName) { || partitionedTopicName.endsWith("/" + TRANSACTION_STATE_TOPIC_NAME); } - private static String path(String... parts) { - StringBuilder sb = new StringBuilder(); - sb.append(POLICY_ROOT); - Joiner.on('/').appendTo(sb, parts); - return sb.toString(); - } - private CompletableFuture> expandAllowedNamespaces(Set allowedNamespaces) { String currentTenant = getCurrentTenant(kafkaConfig.getKafkaTenant()); return expandAllowedNamespaces(allowedNamespaces, currentTenant, pulsarService); From 89825ee5d079001e679318d487f57fb273ef950a Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 20 Oct 2021 10:34:10 +0800 Subject: [PATCH 3/5] Remove unused code --- .../handlers/kop/security/auth/SimpleAclAuthorizer.java | 8 -------- 1 file changed, 8 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java index bac2826a68..4d978b3373 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java @@ -36,8 +36,6 @@ @Slf4j public class SimpleAclAuthorizer implements Authorizer { - private static final String POLICY_ROOT = "/admin/policies/"; - private final PulsarService pulsarService; private final ServiceConfiguration conf; @@ -251,12 +249,6 @@ private CompletableFuture isSuperUserOrTenantAdmin(String tenant, Strin return future; } - private static String path(String... parts) { - StringBuilder sb = new StringBuilder(); - sb.append(POLICY_ROOT); - Joiner.on('/').appendTo(sb, parts); - return sb.toString(); - } @Override public CompletableFuture canAccessTenantAsync(KafkaPrincipal principal, Resource resource) { From 31936d6a3c7972f2c738661e52f06c2da6f80b2c Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 20 Oct 2021 10:39:48 +0800 Subject: [PATCH 4/5] Remove unused import --- .../pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java | 1 - 1 file changed, 1 deletion(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java index 4d978b3373..69e90e123e 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/security/auth/SimpleAclAuthorizer.java @@ -16,7 +16,6 @@ import static com.google.common.base.Preconditions.checkArgument; -import com.google.common.base.Joiner; import io.streamnative.pulsar.handlers.kop.security.KafkaPrincipal; import java.util.Map; import java.util.Set; From b1c0b7c1bfd0a32501fdffedb4e881bedb176d14 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Sun, 24 Oct 2021 11:33:02 +0800 Subject: [PATCH 5/5] Update pulsar version to 2.9.0-rc-202110221101 --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index 529ef35203..5aca3bdea9 100644 --- a/pom.xml +++ b/pom.xml @@ -46,7 +46,7 @@ 1.18.4 2.22.0 io.streamnative - 2.9.0-rc-202110182205 + 2.9.0-rc-202110221101 1.7.25 3.1.8 1.15.1