From 12e02d201311fa0cd3f9f5bcc856ffd592a432b3 Mon Sep 17 00:00:00 2001 From: xiang Date: Mon, 13 Mar 2023 21:41:06 +0800 Subject: [PATCH] pulsar tenant rest api --- ci/spotbugs/exclude.xml | 4 + .../protocol/pulsar/InnerHttpClient.java | 2 + .../github/protocol/pulsar/PulsarAdmin.java | 2 + .../protocol/pulsar/PulsarAdminImpl.java | 9 ++ .../io/github/protocol/pulsar/TenantInfo.java | 44 +++++++ .../io/github/protocol/pulsar/Tenants.java | 36 ++++++ .../github/protocol/pulsar/TenantsImpl.java | 107 ++++++++++++++++++ .../io/github/protocol/pulsar/UrlConst.java | 7 +- .../{BrokerTest.java => BrokersTest.java} | 2 +- .../io/github/protocol/pulsar/RandomUtil.java | 41 +++++++ .../github/protocol/pulsar/TenantsTest.java | 70 ++++++++++++ 11 files changed, 322 insertions(+), 2 deletions(-) create mode 100644 src/main/java/io/github/protocol/pulsar/TenantInfo.java create mode 100644 src/main/java/io/github/protocol/pulsar/Tenants.java create mode 100644 src/main/java/io/github/protocol/pulsar/TenantsImpl.java rename src/test/java/io/github/protocol/pulsar/{BrokerTest.java => BrokersTest.java} (98%) create mode 100644 src/test/java/io/github/protocol/pulsar/RandomUtil.java create mode 100644 src/test/java/io/github/protocol/pulsar/TenantsTest.java diff --git a/ci/spotbugs/exclude.xml b/ci/spotbugs/exclude.xml index b3aaea5..7ef8216 100644 --- a/ci/spotbugs/exclude.xml +++ b/ci/spotbugs/exclude.xml @@ -27,6 +27,10 @@ + + + + diff --git a/src/main/java/io/github/protocol/pulsar/InnerHttpClient.java b/src/main/java/io/github/protocol/pulsar/InnerHttpClient.java index f40ada3..bf6c954 100644 --- a/src/main/java/io/github/protocol/pulsar/InnerHttpClient.java +++ b/src/main/java/io/github/protocol/pulsar/InnerHttpClient.java @@ -66,6 +66,7 @@ public HttpResponse post(String url, String body, String... params) HttpRequest request = HttpRequest.newBuilder() .uri(getUri(url, params)) .POST(HttpRequest.BodyPublishers.ofString(body)) + .setHeader("Content-Type", "application/json") .build(); return client.send(request, HttpResponse.BodyHandlers.ofString()); } @@ -88,6 +89,7 @@ public HttpResponse put(String url, String body, String... params) HttpRequest request = HttpRequest.newBuilder() .uri(getUri(url, params)) .PUT(HttpRequest.BodyPublishers.ofString(body)) + .setHeader("Content-Type", "application/json") .build(); return client.send(request, HttpResponse.BodyHandlers.ofString()); } diff --git a/src/main/java/io/github/protocol/pulsar/PulsarAdmin.java b/src/main/java/io/github/protocol/pulsar/PulsarAdmin.java index 3f2ded4..3c30bcf 100644 --- a/src/main/java/io/github/protocol/pulsar/PulsarAdmin.java +++ b/src/main/java/io/github/protocol/pulsar/PulsarAdmin.java @@ -24,4 +24,6 @@ static PulsarAdminBuilder builder() { } Brokers brokers(); + + Tenants tenants(); } diff --git a/src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java b/src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java index fe7be9f..8bad83e 100644 --- a/src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java +++ b/src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java @@ -22,13 +22,22 @@ public class PulsarAdminImpl implements PulsarAdmin { private final Brokers brokers; + private final Tenants tenants; + PulsarAdminImpl(Configuration conf) { InnerHttpClient innerHttpClient = new InnerHttpClient(conf); this.brokers = new BrokersImpl(innerHttpClient); + this.tenants = new TenantsImpl(innerHttpClient); } @Override public Brokers brokers() { return brokers; } + + @Override + public Tenants tenants() { + return tenants; + } + } diff --git a/src/main/java/io/github/protocol/pulsar/TenantInfo.java b/src/main/java/io/github/protocol/pulsar/TenantInfo.java new file mode 100644 index 0000000..c399030 --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/TenantInfo.java @@ -0,0 +1,44 @@ +/* + * 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 java.util.Set; + +@Getter +@Setter +@Builder +@NoArgsConstructor +@AllArgsConstructor +@EqualsAndHashCode +public class TenantInfo { + + private Set adminRoles; + + private Set allowedClusters; + +} diff --git a/src/main/java/io/github/protocol/pulsar/Tenants.java b/src/main/java/io/github/protocol/pulsar/Tenants.java new file mode 100644 index 0000000..846ff17 --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/Tenants.java @@ -0,0 +1,36 @@ +/* + * 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; + +public interface Tenants { + + void createTenant(String tenant, TenantInfo tenantInfo) throws PulsarAdminException; + + void deleteTenant(String tenant, boolean force) throws PulsarAdminException; + + void updateTenant(String tenant, TenantInfo tenantInfo) throws PulsarAdminException; + + TenantInfo getTenantAdmin(String tenant) throws PulsarAdminException; + + List getTenants() throws PulsarAdminException; + +} diff --git a/src/main/java/io/github/protocol/pulsar/TenantsImpl.java b/src/main/java/io/github/protocol/pulsar/TenantsImpl.java new file mode 100644 index 0000000..e7000f2 --- /dev/null +++ b/src/main/java/io/github/protocol/pulsar/TenantsImpl.java @@ -0,0 +1,107 @@ +/* + * 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; + +public class TenantsImpl implements Tenants { + + private final InnerHttpClient httpClient; + + public TenantsImpl(InnerHttpClient httpClient) { + this.httpClient = httpClient; + } + + @Override + public void createTenant(String tenant, TenantInfo tenantInfo) throws PulsarAdminException { + try { + HttpResponse response = httpClient.put( + String.format("%s/%s", UrlConst.TENANTS, tenant), tenantInfo); + if (response.statusCode() != 204){ + throw new PulsarAdminException(String.format("failed to create tenant %s, status code %s, body : %s", + tenant, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public void deleteTenant(String tenant, boolean force) throws PulsarAdminException { + try { + HttpResponse response = + httpClient.delete(String.format("%s/%s", UrlConst.TENANTS, tenant), "force", String.valueOf(force)); + if (response.statusCode() != 204){ + throw new PulsarAdminException(String.format("failed to create tenant %s, status code %s, body : %s", + tenant, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public void updateTenant(String tenant, TenantInfo tenantInfo) throws PulsarAdminException { + try { + HttpResponse response = + httpClient.post( + String.format("%s/%s", UrlConst.TENANTS, tenant), tenantInfo); + if (response.statusCode() != 204){ + throw new PulsarAdminException(String.format("failed to update tenant %s, status code %s, body : %s", + tenant, response.statusCode(), response.body())); + } + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public TenantInfo getTenantAdmin(String tenant) throws PulsarAdminException { + try { + HttpResponse response = httpClient.get( + String.format("%s/%s", UrlConst.TENANTS, tenant)); + if (response.statusCode() != 200){ + throw new PulsarAdminException(String.format("failed to get tenant %s, status code %s, body : %s", + tenant, response.statusCode(), response.body())); + } + return JacksonService.toObject(response.body(), TenantInfo.class); + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } + + @Override + public List getTenants() throws PulsarAdminException { + try { + HttpResponse response = httpClient.get(UrlConst.TENANTS); + if (response.statusCode() != 200){ + throw new PulsarAdminException(String.format("failed to get list of tenant, status code %s, body : %s", + response.statusCode(), response.body())); + } + return JacksonService.toList(response.body(), new TypeReference>(){}); + } catch (IOException | InterruptedException e) { + throw new PulsarAdminException(e); + } + } +} diff --git a/src/main/java/io/github/protocol/pulsar/UrlConst.java b/src/main/java/io/github/protocol/pulsar/UrlConst.java index edeeb4d..4f15bb6 100644 --- a/src/main/java/io/github/protocol/pulsar/UrlConst.java +++ b/src/main/java/io/github/protocol/pulsar/UrlConst.java @@ -19,7 +19,12 @@ package io.github.protocol.pulsar; public class UrlConst { - public static final String BROKERS = "/admin/v2/brokers"; + + public static final String BASE_URL_V2 = "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/admin/v2"; + + public static final String BROKERS = BASE_URL_V2 + "/brokers"; + + public static final String TENANTS = BASE_URL_V2 + "/tenants"; public static final String HEALTHCHECK = BROKERS + "/health"; } diff --git a/src/test/java/io/github/protocol/pulsar/BrokerTest.java b/src/test/java/io/github/protocol/pulsar/BrokersTest.java similarity index 98% rename from src/test/java/io/github/protocol/pulsar/BrokerTest.java rename to src/test/java/io/github/protocol/pulsar/BrokersTest.java index 0ee0b53..25a2163 100644 --- a/src/test/java/io/github/protocol/pulsar/BrokerTest.java +++ b/src/test/java/io/github/protocol/pulsar/BrokersTest.java @@ -23,7 +23,7 @@ import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; -public class BrokerTest { +public class BrokersTest { private static final EmbeddedPulsarServer SERVER = new EmbeddedPulsarServer(); @BeforeAll diff --git a/src/test/java/io/github/protocol/pulsar/RandomUtil.java b/src/test/java/io/github/protocol/pulsar/RandomUtil.java new file mode 100644 index 0000000..0c8dd38 --- /dev/null +++ b/src/test/java/io/github/protocol/pulsar/RandomUtil.java @@ -0,0 +1,41 @@ +/* + * 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.Random; +import java.util.UUID; + +public class RandomUtil { + + public static Random random = new Random(); + + public static String randomString() { + return UUID.randomUUID().toString(); + } + + public static int randomInt() { + return random.nextInt(); + } + + public static long randomLong() { + return random.nextLong(); + } + +} diff --git a/src/test/java/io/github/protocol/pulsar/TenantsTest.java b/src/test/java/io/github/protocol/pulsar/TenantsTest.java new file mode 100644 index 0000000..634cb41 --- /dev/null +++ b/src/test/java/io/github/protocol/pulsar/TenantsTest.java @@ -0,0 +1,70 @@ +/* + * 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 TenantsTest { + + protected static final EmbeddedPulsarServer SERVER = new EmbeddedPulsarServer(); + + protected static final String CLUSTER_STANDALONE = "standalone"; + + protected static PulsarAdmin pulsarAdmin; + + @BeforeAll + public static void setup() throws Exception { + SERVER.start(); + pulsarAdmin = PulsarAdmin.builder().port(SERVER.getWebPort()).build(); + } + + @AfterAll + public static void teardown() throws Exception { + SERVER.close(); + } + + @Test + public void tenantsTest() throws PulsarAdminException { + String tenantName = RandomUtil.randomString(); + TenantInfo initialTenantInfo = (new TenantInfo.TenantInfoBuilder()) + .adminRoles(new HashSet<>(0)) + .allowedClusters(Set.of(CLUSTER_STANDALONE)).build(); + TenantInfo updatedTenantInfo = (new TenantInfo.TenantInfoBuilder()) + .adminRoles(Set.of("test")) + .allowedClusters(Set.of("global")) + .build(); + pulsarAdmin.tenants().createTenant(tenantName, initialTenantInfo); + Assertions.assertEquals(pulsarAdmin.tenants().getTenants(), List.of(tenantName, "public", "pulsar")); + Assertions.assertEquals(pulsarAdmin.tenants().getTenantAdmin(tenantName), initialTenantInfo); + pulsarAdmin.tenants().updateTenant(tenantName, updatedTenantInfo); + Assertions.assertEquals(pulsarAdmin.tenants().getTenantAdmin(tenantName), updatedTenantInfo); + pulsarAdmin.tenants().deleteTenant(tenantName, false); + Assertions.assertEquals(pulsarAdmin.tenants().getTenants(), List.of("public", "pulsar")); + } + +}