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..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); @@ -517,12 +509,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..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; @@ -36,8 +35,6 @@ @Slf4j public class SimpleAclAuthorizer implements Authorizer { - private static final String POLICY_ROOT = "/admin/policies/"; - private final PulsarService pulsarService; private final ServiceConfiguration conf; @@ -73,7 +70,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 +86,7 @@ private CompletableFuture authorizeTopicPermission(KafkaPrincipal princ getPulsarService() .getPulsarResources() .getNamespaceResources() - .getAsync(policiesPath) + .getPoliciesAsync(namespace) .thenAccept(policies -> { if (!policies.isPresent()) { if (log.isDebugEnabled()) { @@ -177,7 +173,7 @@ private CompletableFuture authorizeTenantPermission(KafkaPrincipal prin getPulsarService() .getPulsarResources() .getTenantResources() - .getAsync(path(tenant)) + .getTenantAsync(tenant) .thenAccept(tenantInfo -> { permissionFuture.complete(tenantInfo.isPresent()); }).exceptionally(ex -> { @@ -234,7 +230,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); @@ -252,12 +248,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) { 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..5aca3bdea9 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-202110221101 1.7.25 3.1.8 1.15.1