From 66e2651a45b7c4a199a94519627c8a1f042e7419 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 8 Jun 2021 02:52:31 +0800 Subject: [PATCH 1/3] Upgrade pulsar to 2.8.0-rc-202106071430 --- .../handlers/kop/KafkaProtocolHandler.java | 20 ++++++----- .../handlers/kop/utils/MetadataUtils.java | 6 +++- .../handlers/kop/utils/MetadataUtilsTest.java | 11 +++--- pom.xml | 2 +- .../kop/DifferentNamespaceTestBase.java | 6 ++-- .../handlers/kop/DistributedClusterTest.java | 18 ---------- .../kop/EntryPublishTimeKafkaFormatTest.java | 2 +- .../handlers/kop/EntryPublishTimeTest.java | 34 ------------------- .../kop/InnerTopicProtectionTest.java | 26 +++----------- .../pulsar/handlers/kop/KafkaApisTest.java | 18 ---------- .../handlers/kop/KafkaIntegrationTest.java | 25 -------------- .../kop/KafkaMessageOrderTestBase.java | 24 ------------- .../handlers/kop/KafkaProducerStatsTest.java | 28 ++------------- .../handlers/kop/KafkaRequestHandlerTest.java | 28 +++------------ .../handlers/kop/KafkaRequestTypeTest.java | 33 ------------------ .../handlers/kop/KafkaSSLChannelTest.java | 33 ------------------ .../KafkaSSLChannelWithClientAuthTest.java | 33 ------------------ .../kop/KafkaTopicConsumerManagerTest.java | 8 ++--- .../kop/KopProtocolHandlerTestBase.java | 7 +++- .../handlers/kop/MetricsProviderTest.java | 33 ------------------ .../pulsar/handlers/kop/MultiLedgerTest.java | 33 ------------------ .../handlers/kop/PublishRateLimitTest.java | 33 ------------------ .../kop/PulsarAuthEnabledTestBase.java | 5 ++- .../handlers/kop/SaslPlainTestBase.java | 5 ++- .../pulsar/handlers/kop/TransactionTest.java | 33 ------------------ .../group/GroupMetadataManagerTest.java | 33 ------------------ 26 files changed, 58 insertions(+), 479 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index 78fb06af9f..ea1d98f2a9 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -369,10 +369,12 @@ public void initGroupCoordinator(BrokerService service) throws Exception { .build(); PulsarAdmin pulsarAdmin = service.pulsar().getAdminClient(); - ClusterData clusterData = new ClusterData(service.getPulsar().getWebServiceAddress(), - service.getPulsar().getWebServiceAddressTls(), - service.getPulsar().getBrokerServiceUrl(), - service.getPulsar().getBrokerServiceUrlTls()); + final ClusterData clusterData = ClusterData.builder() + .serviceUrl(brokerService.getPulsar().getWebServiceAddress()) + .serviceUrlTls(brokerService.getPulsar().getWebServiceAddressTls()) + .brokerServiceUrl(brokerService.getPulsar().getBrokerServiceUrl()) + .brokerServiceUrlTls(brokerService.getPulsar().getBrokerServiceUrlTls()) + .build(); MetadataUtils.createOffsetMetadataIfMissing(pulsarAdmin, clusterData, kafkaConfig); @@ -404,10 +406,12 @@ public void initTransactionCoordinator() throws Exception { .build(); PulsarAdmin pulsarAdmin = brokerService.getPulsar().getAdminClient(); - ClusterData clusterData = new ClusterData(brokerService.getPulsar().getWebServiceAddress(), - brokerService.getPulsar().getWebServiceAddressTls(), - brokerService.getPulsar().getBrokerServiceUrl(), - brokerService.getPulsar().getBrokerServiceUrlTls()); + final ClusterData clusterData = ClusterData.builder() + .serviceUrl(brokerService.getPulsar().getWebServiceAddress()) + .serviceUrlTls(brokerService.getPulsar().getWebServiceAddressTls()) + .brokerServiceUrl(brokerService.getPulsar().getBrokerServiceUrl()) + .brokerServiceUrlTls(brokerService.getPulsar().getBrokerServiceUrlTls()) + .build(); MetadataUtils.createTxnMetadataIfMissing(pulsarAdmin, clusterData, kafkaConfig); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtils.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtils.java index 0c66cd5b69..0670adf9b0 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtils.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtils.java @@ -15,6 +15,7 @@ import com.google.common.collect.Sets; import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; +import java.util.Collections; import java.util.List; import java.util.Set; import lombok.extern.slf4j.Slf4j; @@ -114,7 +115,10 @@ private static void createKafkaMetadataIfMissing(PulsarAdmin pulsarAdmin, if (!tenants.getTenants().contains(kafkaMetadataTenant)) { log.info("Tenant: {} does not exist, creating it ...", kafkaMetadataTenant); tenants.createTenant(kafkaMetadataTenant, - new TenantInfo(Sets.newHashSet(conf.getSuperUserRoles()), Sets.newHashSet(cluster))); + TenantInfo.builder() + .adminRoles(conf.getSuperUserRoles()) + .allowedClusters(Collections.singleton(cluster)) + .build()); } else { TenantInfo kafkaMetadataTenantInfo = tenants.getTenantInfo(kafkaMetadataTenant); Set allowedClusters = kafkaMetadataTenantInfo.getAllowedClusters(); diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java index 0969e4de5d..661ff5e2dd 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java @@ -26,6 +26,7 @@ import com.google.common.collect.Sets; import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Set; import org.apache.pulsar.client.admin.Clusters; @@ -48,7 +49,7 @@ public class MetadataUtilsTest { public void testCreateKafkaMetadataIfMissing() throws Exception { KopTopic.initialize("public/default"); KafkaServiceConfiguration conf = new KafkaServiceConfiguration(); - ClusterData clusterData = new ClusterData(); + ClusterData clusterData = ClusterData.builder().build(); conf.setClusterName("test"); conf.setKafkaMetadataTenant("public"); conf.setKafkaMetadataNamespace("default"); @@ -82,7 +83,7 @@ public void testCreateKafkaMetadataIfMissing() throws Exception { doReturn(mockNamespaces).when(mockPulsarAdmin).namespaces(); doReturn(mockTopics).when(mockPulsarAdmin).topics(); - TenantInfo partialTenant = new TenantInfo(); + TenantInfo partialTenant = TenantInfo.builder().build(); doReturn(partialTenant).when(mockTenants).getTenantInfo(eq(conf.getKafkaMetadataTenant())); MetadataUtils.createOffsetMetadataIfMissing(mockPulsarAdmin, clusterData, conf); @@ -120,8 +121,10 @@ public void testCreateKafkaMetadataIfMissing() throws Exception { doReturn(Lists.newArrayList("public")).when(mockTenants).getTenants(); - partialTenant = new TenantInfo(Sets.newHashSet(conf.getSuperUserRoles()), - Sets.newHashSet("other-cluster")); + partialTenant = TenantInfo.builder() + .adminRoles(conf.getSuperUserRoles()) + .allowedClusters(Collections.singleton("other-cluster")) + .build(); doReturn(partialTenant).when(mockTenants).getTenantInfo(eq(conf.getKafkaMetadataTenant())); doReturn(Lists.newArrayList("test")).when(mockNamespaces).getNamespaces("public"); diff --git a/pom.xml b/pom.xml index 6e29d44baa..a416aac072 100644 --- a/pom.xml +++ b/pom.xml @@ -48,7 +48,7 @@ 1.18.4 2.22.0 io.streamnative - 2.8.0-rc-202105291235 + 2.8.0-rc-202106071430 1.7.25 3.1.8 1.15.1 diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DifferentNamespaceTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DifferentNamespaceTestBase.java index 59d5dfbe23..808af46794 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DifferentNamespaceTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DifferentNamespaceTestBase.java @@ -16,7 +16,6 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; -import com.google.common.collect.Sets; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import java.time.Duration; import java.util.ArrayList; @@ -73,7 +72,10 @@ protected void setup() throws Exception { super.internalSetup(); admin.tenants().createTenant(ANOTHER_TENANT, - new TenantInfo(Sets.newHashSet("admin_user"), Sets.newHashSet(super.configClusterName))); + TenantInfo.builder() + .adminRoles(Collections.singleton("admin_user")) + .allowedClusters(Collections.singleton(configClusterName)) + .build()); admin.namespaces().createNamespace(ANOTHER_TENANT + "/" + ANOTHER_NAMESPACE); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DistributedClusterTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DistributedClusterTest.java index 16a2c4705c..487d6dc7ef 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DistributedClusterTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/DistributedClusterTest.java @@ -36,9 +36,7 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.TopicPartition; import org.apache.pulsar.broker.PulsarService; -import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.junit.Assert; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -157,22 +155,6 @@ protected void stopBroker() throws Exception { public void setup() throws Exception { super.internalSetup(); - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } if (!admin.namespaces().getNamespaces("public").contains("public/default")) { admin.namespaces().createNamespace("public/default"); admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeKafkaFormatTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeKafkaFormatTest.java index 6df5ff5c67..2c28af395d 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeKafkaFormatTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeKafkaFormatTest.java @@ -49,7 +49,7 @@ public class EntryPublishTimeKafkaFormatTest extends EntryPublishTimeTest { private static final Logger log = LoggerFactory.getLogger(EntryPublishTimeKafkaFormatTest.class); - public EntryPublishTimeKafkaFormatTest(String format) { + public EntryPublishTimeKafkaFormatTest() { super("kafka"); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeTest.java index 4fcc5cb6ff..acb30b6669 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeTest.java @@ -16,7 +16,6 @@ import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; -import com.google.common.collect.Sets; import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; @@ -25,9 +24,6 @@ import java.net.InetSocketAddress; import java.net.SocketAddress; import org.apache.pulsar.broker.protocol.ProtocolHandler; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.annotations.AfterMethod; @@ -51,36 +47,6 @@ public EntryPublishTimeTest(String format) { @Override protected void setup() throws Exception { super.internalSetup(); - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(-1, -1)); - } - - log.info("created namespaces, init handler"); ProtocolHandler handler = pulsar.getProtocolHandlers().protocol("kafka"); GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler).getGroupCoordinator(); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/InnerTopicProtectionTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/InnerTopicProtectionTest.java index 2280b2bd21..92e13ead54 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/InnerTopicProtectionTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/InnerTopicProtectionTest.java @@ -26,12 +26,10 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import org.apache.pulsar.client.admin.PulsarAdminException; -import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; -import org.testng.annotations.AfterMethod; -import org.testng.annotations.BeforeMethod; +import org.testng.annotations.AfterClass; +import org.testng.annotations.BeforeClass; import org.testng.annotations.Test; /** @@ -97,28 +95,12 @@ protected void resetConfig() { brokerPort, brokerWebservicePort, kafkaBrokerPort); } - @BeforeMethod + @BeforeClass @Override protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } if (!admin.namespaces().getNamespaces("public").contains("public/default")) { admin.namespaces().createNamespace("public/default"); admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); @@ -133,7 +115,7 @@ protected void setup() throws Exception { } } - @AfterMethod + @AfterClass @Override protected void cleanup() throws Exception { super.internalCleanup(); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java index b6dfbfbcce..9f68e7fe57 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java @@ -71,9 +71,7 @@ import org.apache.kafka.common.requests.RequestHeader; import org.apache.kafka.common.serialization.StringSerializer; import org.apache.pulsar.broker.protocol.ProtocolHandler; -import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Ignore; @@ -102,22 +100,6 @@ protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } if (!admin.namespaces().getNamespaces("public").contains("public/default")) { admin.namespaces().createNamespace("public/default"); admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaIntegrationTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaIntegrationTest.java index f1c48181d1..1b0097331c 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaIntegrationTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaIntegrationTest.java @@ -38,9 +38,7 @@ import java.util.Optional; import java.util.concurrent.TimeUnit; import lombok.extern.slf4j.Slf4j; -import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testcontainers.Testcontainers; import org.testcontainers.containers.GenericContainer; import org.testcontainers.containers.output.WaitingConsumer; @@ -161,29 +159,6 @@ protected void setup() throws Exception { + "," + SSL_PREFIX + "127.0.0.1:" + kafkaBrokerPortTls); super.internalSetup(); - - if (!this.admin.clusters().getClusters().contains(this.configClusterName)) { - // so that clients can test short names - this.admin.clusters().createCluster(this.configClusterName, - new ClusterData("http://127.0.0.1:" + this.brokerWebservicePort)); - } else { - this.admin.clusters().updateCluster(this.configClusterName, - new ClusterData("http://127.0.0.1:" + this.brokerWebservicePort)); - } - - if (!this.admin.tenants().getTenants().contains("public")) { - this.admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - this.admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!this.admin.namespaces().getNamespaces("public").contains("public/default")) { - this.admin.namespaces().createNamespace("public/default"); - this.admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - this.admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } if (!this.admin.namespaces().getNamespaces("public").contains("public/__kafka")) { this.admin.namespaces().createNamespace("public/__kafka"); this.admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaMessageOrderTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaMessageOrderTestBase.java index 908f638363..6458339b92 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaMessageOrderTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaMessageOrderTestBase.java @@ -38,9 +38,7 @@ import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.impl.BatchMessageIdImpl; import org.apache.pulsar.client.impl.TopicMessageIdImpl; -import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; import org.testng.annotations.DataProvider; @@ -69,28 +67,6 @@ protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { admin.namespaces().createNamespace("public/__kafka"); admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaProducerStatsTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaProducerStatsTest.java index c9eb70e291..96b63e54da 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaProducerStatsTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaProducerStatsTest.java @@ -26,9 +26,7 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.pulsar.broker.service.Producer; import org.apache.pulsar.broker.service.persistent.PersistentTopic; -import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -57,28 +55,6 @@ protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { admin.namespaces().createNamespace("public/__kafka"); admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); @@ -123,9 +99,9 @@ public void testKafkaProducePulsarMetrics(int partitionNumber, boolean isBatch) } } - long msgInCounter = admin.topics().getPartitionedStats(pulsarTopicName, false).msgInCounter; + long msgInCounter = admin.topics().getPartitionedStats(pulsarTopicName, false).getMsgInCounter(); assertEquals(msgInCounter, totalMsgs); - long bytesInCounter = admin.topics().getPartitionedStats(pulsarTopicName, false).bytesInCounter; + long bytesInCounter = admin.topics().getPartitionedStats(pulsarTopicName, false).getBytesInCounter(); assertNotEquals(bytesInCounter, 0); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java index 3717032e99..a94d92f72e 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java @@ -89,7 +89,6 @@ import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; import org.apache.pulsar.common.naming.TopicName; -import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; @@ -118,28 +117,6 @@ protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { admin.namespaces().createNamespace("public/__kafka"); admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); @@ -148,7 +125,10 @@ protected void setup() throws Exception { } admin.tenants().createTenant("my-tenant", - new TenantInfo(Sets.newHashSet(), Sets.newHashSet(super.configClusterName))); + TenantInfo.builder() + .adminRoles(Collections.emptySet()) + .allowedClusters(Collections.singleton(configClusterName)) + .build()); admin.namespaces().createNamespace("my-tenant/my-ns"); log.info("created namespaces, init handler"); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestTypeTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestTypeTest.java index 3edd034501..0b6c8c9a63 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestTypeTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestTypeTest.java @@ -21,7 +21,6 @@ import static org.testng.Assert.assertTrue; import com.google.common.collect.Lists; -import com.google.common.collect.Sets; import java.time.Duration; import java.util.Base64; import java.util.Collections; @@ -46,9 +45,6 @@ import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerBuilder; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.common.util.FutureUtil; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -93,35 +89,6 @@ public static Object[][] partitionsAndBatch() { protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(-1, -1)); - } } @AfterMethod diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelTest.java index cc21436c05..ff2685bb82 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelTest.java @@ -15,7 +15,6 @@ import static java.nio.charset.StandardCharsets.UTF_8; -import com.google.common.collect.Sets; import java.io.Closeable; import java.util.Properties; import javax.net.ssl.HostnameVerifier; @@ -27,9 +26,6 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.kafka.common.serialization.StringSerializer; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Factory; @@ -86,35 +82,6 @@ protected void setup() throws Exception { sslSetUpForBroker(); super.internalSetup(); log.info("success internal setup"); - - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(-1, -1)); - } } @AfterMethod diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelWithClientAuthTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelWithClientAuthTest.java index b4a372ac46..f22b7c275d 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelWithClientAuthTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaSSLChannelWithClientAuthTest.java @@ -15,7 +15,6 @@ import static java.nio.charset.StandardCharsets.UTF_8; -import com.google.common.collect.Sets; import java.io.Closeable; import java.util.Properties; import javax.net.ssl.HostnameVerifier; @@ -27,9 +26,6 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.kafka.common.serialization.StringSerializer; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Factory; @@ -88,35 +84,6 @@ protected void setup() throws Exception { sslSetUpForBroker(); super.internalSetup(); log.info("success internal setup"); - - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(-1, -1)); - } } @AfterMethod diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java index 24454f6880..f174d331e0 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java @@ -290,14 +290,14 @@ private void verifyBacklogAndNumCursor(PersistentTopic persistentTopic, TopicStats topicStats = persistentTopic.getStats(true, true); log.info(" dump topicStats for topic : {}, storageSize: {}, backlogSize: {}, expected: {}", persistentTopic.getName(), - topicStats.storageSize, topicStats.backlogSize, expectedBacklog); + topicStats.getStorageSize(), topicStats.getBacklogSize(), expectedBacklog); - topicStats.subscriptions.forEach((subname, substats) -> { + topicStats.getSubscriptions().forEach((subname, substats) -> { log.debug(" dump sub: subname - {}, activeConsumerName {}, " + "consumers {}, msgBacklog {}, unackedMessages {}.", subname, - substats.activeConsumerName, substats.consumers, - substats.msgBacklog, substats.unackedMessages); + substats.getActiveConsumerName(), substats.getConsumers(), + substats.getMsgBacklog(), substats.getUnackedMessages()); }); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java index 4acff67e7f..482eca6b38 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java @@ -227,7 +227,12 @@ protected final void init() throws Exception { String brokerServiceUrl = "pulsar://" + this.conf.getAdvertisedAddress() + ":" + brokerPort; String brokerServiceUrlTls = null; // TLS not supported at this time - ClusterData clusterData = new ClusterData(serviceUrl, serviceUrlTls, brokerServiceUrl, null); + final ClusterData clusterData = ClusterData.builder() + .serviceUrl(serviceUrl) + .serviceUrlTls(serviceUrlTls) + .brokerServiceUrl(brokerServiceUrl) + .brokerServiceUrlTls(brokerServiceUrlTls) + .build(); mockZooKeeper = createMockZooKeeper(configClusterName, serviceUrl, serviceUrlTls, brokerServiceUrl, brokerServiceUrlTls); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MetricsProviderTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MetricsProviderTest.java index 2f7ebf6f6f..b76457ca04 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MetricsProviderTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MetricsProviderTest.java @@ -13,7 +13,6 @@ */ package io.streamnative.pulsar.handlers.kop; -import com.google.common.collect.Sets; import java.io.BufferedReader; import java.io.InputStream; import java.io.InputStreamReader; @@ -33,9 +32,6 @@ import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.TopicPartition; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -52,35 +48,6 @@ public class MetricsProviderTest extends KopProtocolHandlerTestBase{ protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(-1, -1)); - } } @AfterMethod diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MultiLedgerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MultiLedgerTest.java index 4ee785fd7a..e9dee8cf6a 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MultiLedgerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/MultiLedgerTest.java @@ -19,7 +19,6 @@ import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; -import com.google.common.collect.Sets; import java.time.Duration; import java.util.Base64; import java.util.List; @@ -37,9 +36,6 @@ import org.apache.pulsar.client.api.SubscriptionInitialPosition; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.TopicMessageIdImpl; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Factory; @@ -76,35 +72,6 @@ protected void resetConfig() { protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(-1, -1)); - } } @AfterMethod diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/PublishRateLimitTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/PublishRateLimitTest.java index 07e1855779..e6c3dda1c9 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/PublishRateLimitTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/PublishRateLimitTest.java @@ -14,7 +14,6 @@ package io.streamnative.pulsar.handlers.kop; -import com.google.common.collect.Sets; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.producer.ProducerRecord; @@ -22,9 +21,6 @@ import org.apache.pulsar.broker.service.PublishRateLimiter; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.impl.ProducerImpl; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -50,35 +46,6 @@ protected void resetConfig() { protected void setup() throws Exception { super.internalSetup(); log.info("success internal setup"); - - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(-1, -1)); - } } @AfterMethod diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/PulsarAuthEnabledTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/PulsarAuthEnabledTestBase.java index 3ad4931711..ae288c4fc8 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/PulsarAuthEnabledTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/PulsarAuthEnabledTestBase.java @@ -99,7 +99,10 @@ protected void setup() throws Exception { super.internalSetup(); admin.tenants().createTenant(TENANT, - new TenantInfo(Sets.newHashSet(ADMIN_USER), Sets.newHashSet(super.configClusterName))); + TenantInfo.builder() + .adminRoles(Collections.singleton(ADMIN_USER)) + .allowedClusters(Collections.singleton(configClusterName)) + .build()); admin.namespaces().createNamespace(TENANT + "/" + NAMESPACE); admin.namespaces() .setNamespaceReplicationClusters(TENANT + "/" + NAMESPACE, Sets.newHashSet(super.configClusterName)); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java index b03683bffe..398a340ef6 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java @@ -106,7 +106,10 @@ protected void setup() throws Exception { super.internalSetup(); admin.tenants().createTenant(TENANT, - new TenantInfo(Sets.newHashSet(ADMIN_USER), Sets.newHashSet(super.configClusterName))); + TenantInfo.builder() + .adminRoles(Collections.singleton(ADMIN_USER)) + .allowedClusters(Collections.singleton(configClusterName)) + .build()); admin.namespaces().createNamespace(TENANT + "/" + NAMESPACE); admin.topics().createPartitionedTopic(TOPIC, 1); admin diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/TransactionTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/TransactionTest.java index 2c6256e61a..812e42a1b8 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/TransactionTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/TransactionTest.java @@ -13,7 +13,6 @@ */ package io.streamnative.pulsar.handlers.kop; -import com.google.common.collect.Sets; import java.time.Duration; import java.time.temporal.ChronoUnit; import java.util.ArrayList; @@ -40,9 +39,6 @@ import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; @@ -60,35 +56,6 @@ protected void setup() throws Exception { this.conf.setEnableTransactionCoordinator(true); super.internalSetup(); log.info("success internal setup"); - - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(-1, -1)); - } } @AfterClass diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManagerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManagerTest.java index 1b38711182..b1d5a25441 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManagerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/coordinator/group/GroupMetadataManagerTest.java @@ -84,9 +84,6 @@ import org.apache.pulsar.client.api.ReaderBuilder; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionInitialPosition; -import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.RetentionPolicies; -import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -120,36 +117,6 @@ public void setup() throws Exception { .numThreads(1) .build(); - - if (!admin.clusters().getClusters().contains(configClusterName)) { - // so that clients can test short names - admin.clusters().createCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } else { - admin.clusters().updateCluster(configClusterName, - new ClusterData("http://127.0.0.1:" + brokerWebservicePort)); - } - - if (!admin.tenants().getTenants().contains("public")) { - admin.tenants().createTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } else { - admin.tenants().updateTenant("public", - new TenantInfo(Sets.newHashSet("appid1", "appid2"), Sets.newHashSet("test"))); - } - if (!admin.namespaces().getNamespaces("public").contains("public/default")) { - admin.namespaces().createNamespace("public/default"); - admin.namespaces().setNamespaceReplicationClusters("public/default", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/default", - new RetentionPolicies(60, 1000)); - } - if (!admin.namespaces().getNamespaces("public").contains("public/__kafka")) { - admin.namespaces().createNamespace("public/__kafka"); - admin.namespaces().setNamespaceReplicationClusters("public/__kafka", Sets.newHashSet("test")); - admin.namespaces().setRetention("public/__kafka", - new RetentionPolicies(20, 100)); - } - //groupMetadataManager = kafkaService.getGroupCoordinator().getGroupManager(); ProtocolHandler handler = pulsar.getProtocolHandlers().protocol("kafka"); groupMetadataManager = ((KafkaProtocolHandler) handler).getGroupCoordinator().getGroupManager(); } From 5aaf6afcc5c07c255b7f2c553f5970e795060a54 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 8 Jun 2021 10:53:39 +0800 Subject: [PATCH 2/3] Fix testCreateKafkaMetadataIfMissing --- .../pulsar/handlers/kop/utils/MetadataUtilsTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java index 661ff5e2dd..6686e62ef4 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java @@ -123,7 +123,7 @@ public void testCreateKafkaMetadataIfMissing() throws Exception { partialTenant = TenantInfo.builder() .adminRoles(conf.getSuperUserRoles()) - .allowedClusters(Collections.singleton("other-cluster")) + .allowedClusters(Sets.newHashSet("other-cluster")) .build(); doReturn(partialTenant).when(mockTenants).getTenantInfo(eq(conf.getKafkaMetadataTenant())); From 0cd8e578fb89f1ea616733eb0e7bf599d6bdb964 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 8 Jun 2021 10:53:57 +0800 Subject: [PATCH 3/3] Fix checkstyle --- .../pulsar/handlers/kop/utils/MetadataUtilsTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java index 6686e62ef4..ced1c331cd 100644 --- a/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java +++ b/kafka-impl/src/test/java/io/streamnative/pulsar/handlers/kop/utils/MetadataUtilsTest.java @@ -26,7 +26,6 @@ import com.google.common.collect.Sets; import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import java.util.ArrayList; -import java.util.Collections; import java.util.List; import java.util.Set; import org.apache.pulsar.client.admin.Clusters;