From d141b8755d8551980ea0925e2631cbebaa16a170 Mon Sep 17 00:00:00 2001 From: Zixuan Liu Date: Fri, 11 Feb 2022 23:14:47 +0800 Subject: [PATCH] [Broker] Fix NPE of internalExpireMessagesByTimestamp Signed-off-by: Zixuan Liu --- .../admin/impl/PersistentTopicsBase.java | 12 +++- .../admin/AdminApiSubscriptionTest.java | 72 +++++++++++++++++++ .../auth/MockedPulsarServiceBaseTest.java | 20 ++++++ 3 files changed, 102 insertions(+), 2 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiSubscriptionTest.java diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 8890669588d9f..372ecdf7705c8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -3483,11 +3483,19 @@ private CompletableFuture internalExpireMessagesByTimestampForSinglePartit String remoteCluster = PersistentReplicator.getRemoteCluster(subName); PersistentReplicator repl = (PersistentReplicator) topic .getPersistentReplicator(remoteCluster); - checkNotNull(repl); + if (repl == null) { + resultFuture.completeExceptionally( + new RestException(Status.NOT_FOUND, "Replicator not found")); + return; + } issued = repl.expireMessages(expireTimeInSeconds); } else { PersistentSubscription sub = topic.getSubscription(subName); - checkNotNull(sub); + if (sub == null) { + resultFuture.completeExceptionally( + new RestException(Status.NOT_FOUND, "Subscription not found")); + return; + } issued = sub.expireMessages(expireTimeInSeconds); } if (issued) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiSubscriptionTest.java new file mode 100644 index 0000000000000..6f38ccd8b5cf5 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiSubscriptionTest.java @@ -0,0 +1,72 @@ +/** + * 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 org.apache.pulsar.broker.admin; + +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.expectThrows; +import javax.ws.rs.core.Response; +import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest; +import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.client.api.MessageId; +import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; + +@Slf4j +@Test(groups = "broker-admin") +public class AdminApiSubscriptionTest extends MockedPulsarServiceBaseTest { + @BeforeMethod + @Override + public void setup() throws Exception { + super.internalSetup(); + super.setupDefaultTenantAndNamespace(); + } + + @AfterMethod(alwaysRun = true) + @Override + public void cleanup() throws Exception { + super.internalCleanup(); + } + + @Test + public void testExpireNonExistTopic() throws Exception { + String topic = "test-expire-messages-topic"; + String subscriptionName = "test-expire-messages-sub"; + admin.topics().createSubscription(topic, subscriptionName, MessageId.latest); + assertEquals(expectThrows(PulsarAdminException.class, + () -> admin.topics().expireMessages(topic, subscriptionName, 1)).getStatusCode(), + Response.Status.CONFLICT.getStatusCode()); + assertEquals(expectThrows(PulsarAdminException.class, + () -> admin.topics().expireMessagesForAllSubscriptions(topic, 1)).getStatusCode(), + Response.Status.CONFLICT.getStatusCode()); + } + + @Test + public void TestExpireNonExistTopicAndNonExistSub() { + String topic = "test-expire-messages-topic"; + String subscriptionName = "test-expire-messages-sub"; + assertEquals(expectThrows(PulsarAdminException.class, + () -> admin.topics().expireMessages(topic, subscriptionName, 1)).getStatusCode(), + Response.Status.NOT_FOUND.getStatusCode()); + assertEquals(expectThrows(PulsarAdminException.class, + () -> admin.topics().expireMessagesForAllSubscriptions(topic, 1)).getStatusCode(), + Response.Status.NOT_FOUND.getStatusCode()); + } +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java index 7c1cfd01d541d..24f0465da6e42 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java @@ -58,6 +58,7 @@ import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.common.policies.data.ClusterData; +import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.common.policies.data.TenantInfoImpl; import org.apache.pulsar.metadata.api.MetadataStoreException; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; @@ -479,5 +480,24 @@ protected static ServiceConfiguration getDefaultConf() { return configuration; } + protected void setupDefaultTenantAndNamespace() throws Exception { + final String tenant = "public"; + final String namespace = tenant + "/default"; + + if (!admin.clusters().getClusters().contains(configClusterName)) { + admin.clusters().createCluster(configClusterName, + ClusterData.builder().serviceUrl(pulsar.getWebServiceAddress()).build()); + } + + if (!admin.tenants().getTenants().contains(tenant)) { + admin.tenants().createTenant(tenant, TenantInfo.builder().allowedClusters( + Sets.newHashSet(configClusterName)).build()); + } + + if (!admin.namespaces().getNamespaces(tenant).contains(namespace)) { + admin.namespaces().createNamespace(namespace); + } + } + private static final Logger log = LoggerFactory.getLogger(MockedPulsarServiceBaseTest.class); }