Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions ci/spotbugs/exclude.xml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,10 @@

<!-- pattern sort by alpha -->

<Match>
<Bug pattern="EI_EXPOSE_REP"/>
</Match>

<Match>
<Bug pattern="EI_EXPOSE_REP2"/>
</Match>
Expand Down
2 changes: 2 additions & 0 deletions src/main/java/io/github/protocol/pulsar/InnerHttpClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ public HttpResponse<String> 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());
}
Expand All @@ -88,6 +89,7 @@ public HttpResponse<String> 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());
}
Expand Down
2 changes: 2 additions & 0 deletions src/main/java/io/github/protocol/pulsar/PulsarAdmin.java
Original file line number Diff line number Diff line change
Expand Up @@ -24,4 +24,6 @@ static PulsarAdminBuilder builder() {
}

Brokers brokers();

Tenants tenants();
}
9 changes: 9 additions & 0 deletions src/main/java/io/github/protocol/pulsar/PulsarAdminImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

}
44 changes: 44 additions & 0 deletions src/main/java/io/github/protocol/pulsar/TenantInfo.java
Original file line number Diff line number Diff line change
@@ -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<String> adminRoles;

private Set<String> allowedClusters;

}
36 changes: 36 additions & 0 deletions src/main/java/io/github/protocol/pulsar/Tenants.java
Original file line number Diff line number Diff line change
@@ -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<String> getTenants() throws PulsarAdminException;

}
107 changes: 107 additions & 0 deletions src/main/java/io/github/protocol/pulsar/TenantsImpl.java
Original file line number Diff line number Diff line change
@@ -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<String> 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<String> 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<String> 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<String> 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<String> getTenants() throws PulsarAdminException {
try {
HttpResponse<String> 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<List<String>>(){});
} catch (IOException | InterruptedException e) {
throw new PulsarAdminException(e);
}
}
}
7 changes: 6 additions & 1 deletion src/main/java/io/github/protocol/pulsar/UrlConst.java
Original file line number Diff line number Diff line change
Expand Up @@ -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 = "/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";
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
41 changes: 41 additions & 0 deletions src/test/java/io/github/protocol/pulsar/RandomUtil.java
Original file line number Diff line number Diff line change
@@ -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();
}

}
70 changes: 70 additions & 0 deletions src/test/java/io/github/protocol/pulsar/TenantsTest.java
Original file line number Diff line number Diff line change
@@ -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"));
}

}