diff --git a/src/main/java/io/github/protocol/pulsar/MessageIdImpl.java b/src/main/java/io/github/protocol/pulsar/MessageIdImpl.java new file mode 100644 index 0000000..bf1d908 --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/MessageIdImpl.java @@ -0,0 +1,48 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package io.github.protocol.pulsar; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.EqualsAndHashCode; + +@Builder +@AllArgsConstructor +@EqualsAndHashCode +public class MessageIdImpl { + + private Long ledgerId; + + private Integer entryId; + + private Integer partitionIndex; + + @Override + public String toString() { + return new StringBuilder() + .append(ledgerId) + .append(':') + .append(entryId) + .append(':') + .append(partitionIndex) + .toString(); + } + +} diff --git a/src/main/java/io/github/protocol/pulsar/NonPersistentTopicsImpl.java b/src/main/java/io/github/protocol/pulsar/NonPersistentTopicsImpl.java new file mode 100644 index 0000000..96fd65c --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/NonPersistentTopicsImpl.java @@ -0,0 +1,34 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package io.github.protocol.pulsar; + +public class NonPersistentTopicsImpl extends PersistentTopicsImpl { + + private static final String BASE_URL_NON_PERSISTENT_DOMAIN = UrlConst.BASE_URL_V2 + "/non-persistent"; + + public NonPersistentTopicsImpl(InnerHttpClient httpClient) { + super(httpClient); + } + + @Override + public String getDomainBaseUrl() { + return BASE_URL_NON_PERSISTENT_DOMAIN; + } +} diff --git a/src/main/java/io/github/protocol/pulsar/PartitionedTopicMetadata.java b/src/main/java/io/github/protocol/pulsar/PartitionedTopicMetadata.java new file mode 100644 index 0000000..b920018 --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/PartitionedTopicMetadata.java @@ -0,0 +1,47 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package io.github.protocol.pulsar; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; +import lombok.ToString; + +import java.util.Map; + +@Getter +@Setter +@Builder +@NoArgsConstructor +@AllArgsConstructor +@ToString +@EqualsAndHashCode +public class PartitionedTopicMetadata { + + public int partitions; + + public boolean deleted; + + public Map properties; + +} diff --git a/src/main/java/io/github/protocol/pulsar/PersistentOfflineTopicStats.java b/src/main/java/io/github/protocol/pulsar/PersistentOfflineTopicStats.java new file mode 100644 index 0000000..7b56c8e --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/PersistentOfflineTopicStats.java @@ -0,0 +1,93 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package io.github.protocol.pulsar; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.Setter; +import lombok.ToString; + +import java.util.Date; +import java.util.List; +import java.util.Map; + +@Getter +@Setter +@Builder +@NoArgsConstructor +@AllArgsConstructor +@EqualsAndHashCode +@ToString +public class PersistentOfflineTopicStats { + + private Long storageSize; + + private Long totalMessages; + + private Long messageBacklog; + + private String brokerName; + + private String topicName; + + private List dataLedgerDetails; + + private Map cursorDetails; + + private Date statGeneratedAt; + + @Getter + @Setter + @Builder + @NoArgsConstructor + @AllArgsConstructor + @EqualsAndHashCode + @ToString + public static class CursorDetails { + + public Long cursorBacklog; + + public Long cursorLedgerId; + + } + + @Getter + @Setter + @Builder + @NoArgsConstructor + @AllArgsConstructor + @EqualsAndHashCode + @ToString + public static class LedgerDetails { + + public Long entries; + + public Long timestamp; + + public Long size; + + public Long ledgerId; + + } + +} diff --git a/src/main/java/io/github/protocol/pulsar/PersistentTopicsImpl.java b/src/main/java/io/github/protocol/pulsar/PersistentTopicsImpl.java new file mode 100644 index 0000000..58ef200 --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/PersistentTopicsImpl.java @@ -0,0 +1,278 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package io.github.protocol.pulsar; + +import com.fasterxml.jackson.core.type.TypeReference; + +import java.io.IOException; +import java.net.http.HttpResponse; +import java.util.List; +import java.util.Map; + +public class PersistentTopicsImpl implements Topics { + + protected final InnerHttpClient httpClient; + + private static final String BASE_URL_PERSISTENT_DOMAIN = UrlConst.BASE_URL_V2 + "/persistent"; + + public PersistentTopicsImpl(InnerHttpClient httpClient) { + this.httpClient = httpClient; + } + + public String getDomainBaseUrl() { + return BASE_URL_PERSISTENT_DOMAIN; + } + + @Override + public void createPartitionedTopic(String tenant, String namespace, String encodedTopic, int numPartitions, + boolean createLocalTopicOnly) throws PulsarAdminException { + String url = String.format("%s/%s/%s/%s%s", getDomainBaseUrl(), tenant, namespace, encodedTopic, + UrlConst.PARTITIONS); + try { + HttpResponse response = httpClient.put(url, numPartitions, "createLocalTopicOnly", + String.valueOf(createLocalTopicOnly)); + if (response.statusCode() != 204) { + throw new PulsarAdminException( + String.format("failed to create partitioned topic %s/%s/%s, status code %s, body : %s", + tenant, namespace, encodedTopic, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new RuntimeException(e); + } + } + + @Override + public void deletePartitionedTopic(String tenant, String namespace, String encodedTopic, boolean force, + boolean authoritative) throws PulsarAdminException { + String url = String.format("%s/%s/%s/%s%s", getDomainBaseUrl(), tenant, namespace, encodedTopic, + UrlConst.PARTITIONS); + try { + HttpResponse response = httpClient.delete(url, "force", String.valueOf(force), + "authoritative", String.valueOf(authoritative)); + if (response.statusCode() != 204) { + throw new PulsarAdminException( + String.format("failed to delete partitioned topic %s/%s/%s, status code %s, body : %s", + tenant, namespace, encodedTopic, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public void updatePartitionedTopic(String tenant, String namespace, String encodedTopic, + boolean updateLocalTopicOnly, boolean authoritative, + boolean force, int numPartitions) throws PulsarAdminException { + String url = String.format("%s/%s/%s/%s%s", getDomainBaseUrl(), tenant, namespace, encodedTopic, + UrlConst.PARTITIONS); + try { + HttpResponse response = httpClient.post(url, numPartitions, "updateLocalTopicOnly", + String.valueOf(updateLocalTopicOnly), + "authoritative", String.valueOf(authoritative), + "force", String.valueOf(force)); + if (response.statusCode() != 204) { + throw new PulsarAdminException( + String.format("failed to update partitioned topic %s/%s/%s, status code %s, body : %s", + tenant, namespace, encodedTopic, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public PartitionedTopicMetadata getPartitionedMetadata(String tenant, String namespace, String encodedTopic, + boolean checkAllowAutoCreation, boolean authoritative) + throws PulsarAdminException { + String url = String.format("%s/%s/%s/%s%s", getDomainBaseUrl(), + tenant, namespace, encodedTopic, UrlConst.PARTITIONS); + try { + HttpResponse response = httpClient.get(url, + "checkAllowAutoCreation", String.valueOf(checkAllowAutoCreation), + "authoritative", String.valueOf(authoritative)); + if (response.statusCode() != 200) { + throw new PulsarAdminException( + String.format("failed to update partitioned topic %s/%s/%s, status code %s, body : %s", + tenant, namespace, encodedTopic, response.statusCode(), response.body())); + } + return JacksonService.toObject(response.body(), PartitionedTopicMetadata.class); + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public void createNonPartitionedTopic(String tenant, String namespace, String encodedTopic, boolean authoritative, + Map properties) throws PulsarAdminException { + String url = String.format("%s/%s/%s/%s", getDomainBaseUrl(), tenant, namespace, encodedTopic); + try { + HttpResponse response = httpClient.put(url, properties, + "authoritative", String.valueOf(authoritative)); + if (response.statusCode() != 204) { + throw new PulsarAdminException( + String.format("failed to create non-partitioned topic %s/%s/%s, status code %s, body : %s", + tenant, namespace, encodedTopic, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public void deleteTopic(String tenant, String namespace, String encodedTopic, boolean force, boolean authoritative) + throws PulsarAdminException { + String url = String.format("%s/%s/%s/%s", getDomainBaseUrl(), tenant, namespace, encodedTopic); + try { + HttpResponse response = httpClient.delete(url, + "force", String.valueOf(force), + "authoritative", String.valueOf(authoritative)); + if (response.statusCode() != 204) { + throw new PulsarAdminException( + String.format("failed to delete non-partitioned topic %s/%s/%s, status code %s, body : %s", + tenant, namespace, encodedTopic, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public List getList(String tenant, String namespace, String bundle, boolean includeSystemTopic) + throws PulsarAdminException { + String url = String.format("%s/%s/%s", getDomainBaseUrl(), tenant, namespace); + try { + HttpResponse response; + if (bundle != null) { + response = httpClient.get(url, + "bundle", bundle, + "includeSystemTopic", String.valueOf(includeSystemTopic)); + } else { + response = httpClient.get(url, + "includeSystemTopic", String.valueOf(includeSystemTopic)); + } + if (response.statusCode() != 200) { + throw new PulsarAdminException( + String.format("failed to get list of non-partitioned-topics " + + "under namespace %s/%s, status code %s, body : %s", + tenant, namespace, response.statusCode(), response.body())); + } + return JacksonService.toRefer(response.body(), new TypeReference>() { + }); + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public List getPartitionedTopicList(String tenant, String namespace, boolean includeSystemTopic) + throws PulsarAdminException { + String url = String.format("%s/%s/%s%s", getDomainBaseUrl(), tenant, namespace, UrlConst.PARTITIONED); + try { + HttpResponse response = httpClient.get(url, "includeSystemTopic", + String.valueOf(includeSystemTopic)); + if (response.statusCode() != 200) { + throw new PulsarAdminException( + String.format("failed to get list of partitioned-topics under namespace %s/%s, " + + "status code %s, body : %s", + tenant, namespace, response.statusCode(), response.body())); + } + return JacksonService.toRefer(response.body(), new TypeReference>() { + }); + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public void createMissedPartitions(String tenant, String namespace, String encodedTopic) + throws PulsarAdminException { + String url = String.format("%s/%s/%s/%s%s", getDomainBaseUrl(), tenant, namespace, encodedTopic, + UrlConst.CREATE_MISSED_PARTITIONS); + try { + HttpResponse response = httpClient.post(url); + if (response.statusCode() != 204) { + throw new PulsarAdminException( + String.format("failed to delete non-partitioned topic %s/%s/%s, status code %s, body : %s", + tenant, namespace, encodedTopic, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public void getLastMessageId(String tenant, String namespace, String encodedTopic, boolean authoritative) + throws PulsarAdminException { + + } + + @Override + public RetentionPolicies getRetention(String tenant, String namespace, String encodedTopic, boolean isGlobal, + boolean applied, boolean authoritative) throws PulsarAdminException { + return null; + } + + @Override + public void setRetention(String tenant, String namespace, String encodedTopic, boolean authoritative, + boolean isGlobal, RetentionPolicies retention) throws PulsarAdminException { + + } + + @Override + public void removeRetention(String tenant, String namespace, String encodedTopic, boolean authoritative) + throws PulsarAdminException { + + } + + @Override + public Map getBacklogQuotaMap(String tenant, String namespace, + String encodedTopic, boolean applied, + boolean authoritative, boolean isGlobal) + throws PulsarAdminException { + return null; + } + + @Override + public void setBacklogQuota(String tenant, String namespace, String encodedTopic, boolean authoritative, + boolean isGlobal, BacklogQuotaType backlogQuotaType, BacklogQuota backlogQuota) + throws PulsarAdminException { + + } + + @Override + public void removeBacklogQuota(String tenant, String namespace, String encodedTopic, + BacklogQuotaType backlogQuotaType, boolean authoritative, boolean isGlobal) + throws PulsarAdminException { + + } + + @Override + public PersistentOfflineTopicStats getBacklog(String tenant, String namespace, String encodedTopic, + boolean authoritative) throws PulsarAdminException { + return null; + } + + @Override + public long getBacklogSizeByMessageId(String tenant, String namespace, String encodedTopic, boolean authoritative, + MessageIdImpl messageId) throws PulsarAdminException { + return 0; + } + +} diff --git a/src/main/java/io/github/protocol/pulsar/PulsarAdmin.java b/src/main/java/io/github/protocol/pulsar/PulsarAdmin.java index 3937951..c8d596b 100644 --- a/src/main/java/io/github/protocol/pulsar/PulsarAdmin.java +++ b/src/main/java/io/github/protocol/pulsar/PulsarAdmin.java @@ -28,4 +28,8 @@ static PulsarAdminBuilder builder() { Tenants tenants(); Namespaces namespaces(); + + PersistentTopicsImpl persistentTopics(); + + NonPersistentTopicsImpl nonPersistentTopics(); } diff --git a/src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java b/src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java index 0866044..915d2a4 100644 --- a/src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java +++ b/src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java @@ -26,11 +26,17 @@ public class PulsarAdminImpl implements PulsarAdmin { private final Namespaces namespaces; + private final PersistentTopicsImpl persistentTopics; + + private final NonPersistentTopicsImpl nonPersistentTopics; + PulsarAdminImpl(Configuration conf) { InnerHttpClient innerHttpClient = new InnerHttpClient(conf); this.brokers = new BrokersImpl(innerHttpClient); this.tenants = new TenantsImpl(innerHttpClient); this.namespaces = new NamespacesImpl(innerHttpClient); + this.persistentTopics = new PersistentTopicsImpl(innerHttpClient); + this.nonPersistentTopics = new NonPersistentTopicsImpl(innerHttpClient); } @Override @@ -47,4 +53,14 @@ public Tenants tenants() { public Namespaces namespaces() { return namespaces; } + + @Override + public PersistentTopicsImpl persistentTopics() { + return persistentTopics; + } + + @Override + public NonPersistentTopicsImpl nonPersistentTopics() { + return nonPersistentTopics; + } } diff --git a/src/main/java/io/github/protocol/pulsar/Topics.java b/src/main/java/io/github/protocol/pulsar/Topics.java new file mode 100644 index 0000000..cc8ef24 --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/Topics.java @@ -0,0 +1,84 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package io.github.protocol.pulsar; + +import java.util.List; +import java.util.Map; + +public interface Topics { + + void createPartitionedTopic(String tenant, String namespace, String encodedTopic, int numPartitions, + boolean createLocalTopicOnly) throws PulsarAdminException; + + void deletePartitionedTopic(String tenant, String namespace, String encodedTopic, boolean force, + boolean authoritative) throws PulsarAdminException; + + void updatePartitionedTopic(String tenant, String namespace, String encodedTopic, boolean updateLocalTopicOnly, + boolean authoritative, boolean force, int numPartitions) throws PulsarAdminException; + + PartitionedTopicMetadata getPartitionedMetadata(String tenant, String namespace, String encodedTopic, + boolean checkAllowAutoCreation, boolean authoritative) + throws PulsarAdminException; + + void createNonPartitionedTopic(String tenant, String namespace, String encodedTopic, boolean authoritative, + Map properties) throws PulsarAdminException; + + void deleteTopic(String tenant, String namespace, String encodedTopic, boolean force, boolean authoritative) + throws PulsarAdminException; + + List getList(String tenant, String namespace, String bundle, boolean includeSystemTopic) + throws PulsarAdminException; + + List getPartitionedTopicList(String tenant, String namespace, boolean includeSystemTopic) + throws PulsarAdminException; + + void createMissedPartitions(String tenant, String namespace, String encodedTopic) throws PulsarAdminException; + + void getLastMessageId(String tenant, String namespace, String encodedTopic, boolean authoritative) + throws PulsarAdminException; + + RetentionPolicies getRetention(String tenant, String namespace, String encodedTopic, + boolean isGlobal, boolean applied, boolean authoritative) + throws PulsarAdminException; + + void setRetention(String tenant, String namespace, String encodedTopic, boolean authoritative, boolean isGlobal, + RetentionPolicies retention) throws PulsarAdminException; + + void removeRetention(String tenant, String namespace, String encodedTopic, boolean authoritative) + throws PulsarAdminException; + + Map getBacklogQuotaMap(String tenant, String namespace, String encodedTopic, + boolean applied, boolean authoritative, boolean isGlobal) + throws PulsarAdminException; + + void setBacklogQuota(String tenant, String namespace, String encodedTopic, boolean authoritative, + boolean isGlobal, BacklogQuotaType backlogQuotaType, BacklogQuota backlogQuota) + throws PulsarAdminException; + + void removeBacklogQuota(String tenant, String namespace, String encodedTopic, BacklogQuotaType backlogQuotaType, + boolean authoritative, boolean isGlobal) throws PulsarAdminException; + + PersistentOfflineTopicStats getBacklog(String tenant, String namespace, String encodedTopic, boolean authoritative) + throws PulsarAdminException; + + long getBacklogSizeByMessageId(String tenant, String namespace, String encodedTopic, boolean authoritative, + MessageIdImpl messageId) throws PulsarAdminException; + +} diff --git a/src/main/java/io/github/protocol/pulsar/UrlConst.java b/src/main/java/io/github/protocol/pulsar/UrlConst.java index e282f6c..5eee3b4 100644 --- a/src/main/java/io/github/protocol/pulsar/UrlConst.java +++ b/src/main/java/io/github/protocol/pulsar/UrlConst.java @@ -32,6 +32,10 @@ public class UrlConst { public static final String BACKLOG_QUOTA = "/backlogQuota"; + public static final String BACKLOG = "/backlog"; + + public static final String BACKLOG_SIZE = "/backlogSize"; + public static final String RETENTION = "/retention"; public static final String CLEAR_BACKLOG = "/clearBacklog"; @@ -40,5 +44,11 @@ public class UrlConst { public static final String COMPACTION_THRESHOLD = "/compactionThreshold"; + public static final String PARTITIONS = "/partitions"; + + public static final String PARTITIONED = "/partitioned"; + + public static final String CREATE_MISSED_PARTITIONS = "/createMissedPartitions"; + public static final String HEALTHCHECK = BROKERS + "/health"; } diff --git a/src/test/java/io/github/protocol/pulsar/NonPersistentTopicsTest.java b/src/test/java/io/github/protocol/pulsar/NonPersistentTopicsTest.java new file mode 100644 index 0000000..00d105f --- /dev/null +++ b/src/test/java/io/github/protocol/pulsar/NonPersistentTopicsTest.java @@ -0,0 +1,95 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package io.github.protocol.pulsar; + +import io.github.embedded.pulsar.core.EmbeddedPulsarServer; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; + +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +public class NonPersistentTopicsTest { + + private static final EmbeddedPulsarServer SERVER = new EmbeddedPulsarServer(); + + private static final String CLUSTER_STANDALONE = "standalone"; + + private static final String tenant = RandomUtil.randomString(); + + private static PulsarAdmin pulsarAdmin; + + @BeforeAll + public static void setup() throws Exception { + SERVER.start(); + pulsarAdmin = PulsarAdmin.builder().port(SERVER.getWebPort()).build(); + TenantInfo initialTenantInfo = (new TenantInfo.TenantInfoBuilder()) + .adminRoles(new HashSet<>(0)) + .allowedClusters(Set.of(CLUSTER_STANDALONE)).build(); + pulsarAdmin.tenants().createTenant(tenant, initialTenantInfo); + } + + @AfterAll + public static void teardown() throws Exception { + SERVER.close(); + } + + @Test + public void partitionedTopicsTest() throws PulsarAdminException { + String namespace = RandomUtil.randomString(); + String topic = RandomUtil.randomString(); + pulsarAdmin.namespaces().createNamespace(tenant, namespace); + pulsarAdmin.nonPersistentTopics().createPartitionedTopic(tenant, namespace, topic, 2, false); + Assertions.assertEquals(List.of(String.format("non-persistent://%s/%s/%s", tenant, namespace, topic)), + pulsarAdmin.nonPersistentTopics().getPartitionedTopicList(tenant, namespace, false)); + Assertions.assertEquals(2, pulsarAdmin.nonPersistentTopics().getPartitionedMetadata(tenant, namespace, + topic, false, false).getPartitions()); + pulsarAdmin.nonPersistentTopics().updatePartitionedTopic(tenant, namespace, topic, false, false, false, 3); + Assertions.assertEquals(List.of(String.format("non-persistent://%s/%s/%s", tenant, namespace, topic)), + pulsarAdmin.nonPersistentTopics().getPartitionedTopicList(tenant, namespace, false)); + Assertions.assertEquals(3, pulsarAdmin.nonPersistentTopics().getPartitionedMetadata(tenant, namespace, + topic, false, false).getPartitions()); + pulsarAdmin.nonPersistentTopics().deletePartitionedTopic(tenant, namespace, topic, false, false); + Assertions.assertEquals(List.of(), + pulsarAdmin.nonPersistentTopics().getPartitionedTopicList(tenant, namespace, false)); + Assertions.assertEquals( + List.of(), + pulsarAdmin.nonPersistentTopics().getList(tenant, namespace, null, false)); + } + + @Test + public void nonPartitionedTopicsTest() throws PulsarAdminException { + String namespace = RandomUtil.randomString(); + String topic = RandomUtil.randomString(); + pulsarAdmin.namespaces().createNamespace(tenant, namespace); + pulsarAdmin.persistentTopics().createNonPartitionedTopic(tenant, namespace, topic, false, null); + Assertions.assertEquals( + List.of(String.format("persistent://%s/%s/%s", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + pulsarAdmin.persistentTopics().deleteTopic(tenant, namespace, topic, false, false); + Assertions.assertEquals( + List.of(), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + } + +} diff --git a/src/test/java/io/github/protocol/pulsar/PersistentTopicsTest.java b/src/test/java/io/github/protocol/pulsar/PersistentTopicsTest.java new file mode 100644 index 0000000..a2c9b7e --- /dev/null +++ b/src/test/java/io/github/protocol/pulsar/PersistentTopicsTest.java @@ -0,0 +1,126 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package io.github.protocol.pulsar; + +import io.github.embedded.pulsar.core.EmbeddedPulsarServer; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; + +import java.util.HashSet; +import java.util.List; +import java.util.Set; + +public class PersistentTopicsTest { + + private static final EmbeddedPulsarServer SERVER = new EmbeddedPulsarServer(); + + private static final String CLUSTER_STANDALONE = "standalone"; + + private static final String tenant = RandomUtil.randomString(); + + private static PulsarAdmin pulsarAdmin; + + @BeforeAll + public static void setup() throws Exception { + SERVER.start(); + pulsarAdmin = PulsarAdmin.builder().port(SERVER.getWebPort()).build(); + TenantInfo initialTenantInfo = (new TenantInfo.TenantInfoBuilder()) + .adminRoles(new HashSet<>(0)) + .allowedClusters(Set.of(CLUSTER_STANDALONE)).build(); + pulsarAdmin.tenants().createTenant(tenant, initialTenantInfo); + } + + @AfterAll + public static void teardown() throws Exception { + SERVER.close(); + } + + @Test + public void partitionedTopicsTest() throws PulsarAdminException { + String namespace = RandomUtil.randomString(); + String topic = RandomUtil.randomString(); + pulsarAdmin.namespaces().createNamespace(tenant, namespace); + pulsarAdmin.persistentTopics().createPartitionedTopic(tenant, namespace, topic, 2, false); + Assertions.assertEquals(List.of(String.format("persistent://%s/%s/%s", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getPartitionedTopicList(tenant, namespace, false)); + Assertions.assertEquals(2, pulsarAdmin.persistentTopics().getPartitionedMetadata(tenant, namespace, + topic, false, false).getPartitions()); + Assertions.assertEquals( + List.of(String.format("persistent://%s/%s/%s-partition-0", tenant, namespace, topic), + String.format("persistent://%s/%s/%s-partition-1", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + pulsarAdmin.persistentTopics().updatePartitionedTopic(tenant, namespace, topic, false, false, false, 3); + Assertions.assertEquals(List.of(String.format("persistent://%s/%s/%s", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getPartitionedTopicList(tenant, namespace, false)); + Assertions.assertEquals(3, pulsarAdmin.persistentTopics().getPartitionedMetadata(tenant, namespace, + topic, false, false).getPartitions()); + Assertions.assertEquals( + List.of(String.format("persistent://%s/%s/%s-partition-0", tenant, namespace, topic), + String.format("persistent://%s/%s/%s-partition-1", tenant, namespace, topic), + String.format("persistent://%s/%s/%s-partition-2", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + pulsarAdmin.persistentTopics().deletePartitionedTopic(tenant, namespace, topic, false, false); + Assertions.assertEquals(List.of(), + pulsarAdmin.persistentTopics().getPartitionedTopicList(tenant, namespace, false)); + Assertions.assertEquals( + List.of(), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + + } + + @Test + public void nonPartitionedTopicsTest() throws PulsarAdminException { + String namespace = RandomUtil.randomString(); + String topic = RandomUtil.randomString(); + pulsarAdmin.namespaces().createNamespace(tenant, namespace); + pulsarAdmin.persistentTopics().createNonPartitionedTopic(tenant, namespace, topic, false, null); + Assertions.assertEquals( + List.of(String.format("persistent://%s/%s/%s", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + pulsarAdmin.persistentTopics().deleteTopic(tenant, namespace, topic, false, false); + Assertions.assertEquals( + List.of(), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + } + + @Test + public void createMissedPartitionsTest() throws PulsarAdminException { + String namespace = RandomUtil.randomString(); + String topic = RandomUtil.randomString(); + pulsarAdmin.namespaces().createNamespace(tenant, namespace); + pulsarAdmin.persistentTopics().createPartitionedTopic(tenant, namespace, topic, 2, false); + Assertions.assertEquals( + List.of(String.format("persistent://%s/%s/%s-partition-0", tenant, namespace, topic), + String.format("persistent://%s/%s/%s-partition-1", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + pulsarAdmin.persistentTopics().deleteTopic(tenant, namespace, topic + "-partition-1", false, false); + Assertions.assertEquals( + List.of(String.format("persistent://%s/%s/%s-partition-0", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + pulsarAdmin.persistentTopics().createMissedPartitions(tenant, namespace, topic); + Assertions.assertEquals( + List.of(String.format("persistent://%s/%s/%s-partition-0", tenant, namespace, topic), + String.format("persistent://%s/%s/%s-partition-1", tenant, namespace, topic)), + pulsarAdmin.persistentTopics().getList(tenant, namespace, null, false)); + } + +}