From 0bc97107dabd82e5f9eee3f2e8e1365abd349a01 Mon Sep 17 00:00:00 2001 From: rajan Date: Mon, 7 Feb 2022 17:49:18 -0800 Subject: [PATCH 1/3] [pulsar-broker] Fix: return non-null persistence policy for namespace --- .../apache/pulsar/broker/admin/AdminResource.java | 8 ++++++++ .../pulsar/broker/admin/impl/NamespacesBase.java | 2 +- .../broker/admin/impl/PersistentTopicsBase.java | 8 +------- .../apache/pulsar/broker/admin/AdminApiTest.java | 13 +++++++++++++ 4 files changed, 23 insertions(+), 8 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 9022208be89f3..5f592ee9278fe 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -305,6 +305,14 @@ protected Policies getNamespacePolicies(NamespaceName namespaceName) { } + protected PersistencePolicies getOrDefaultPersistencePolicy(PersistencePolicies policies) { + return policies != null ? policies + : new PersistencePolicies(pulsar().getConfiguration().getManagedLedgerDefaultEnsembleSize(), + pulsar().getConfiguration().getManagedLedgerDefaultWriteQuorum(), + pulsar().getConfiguration().getManagedLedgerDefaultAckQuorum(), + pulsar().getConfiguration().getManagedLedgerDefaultMarkDeleteRateLimit()); + } + protected CompletableFuture getNamespacePoliciesAsync(NamespaceName namespaceName) { return namespaceResources().getPoliciesAsync(namespaceName).thenCompose(policies -> { if (policies.isPresent()) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index 1fc20d9b9d2fa..fadff788dc637 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -1519,7 +1519,7 @@ protected PersistencePolicies internalGetPersistence() { validateNamespacePolicyOperation(namespaceName, PolicyName.PERSISTENCE, PolicyOperation.READ); Policies policies = getNamespacePolicies(namespaceName); - return policies.persistence; + return getOrDefaultPersistencePolicy(policies.persistence); } protected void internalClearNamespaceBacklog(AsyncResponse asyncResponse, boolean authoritative) { 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 69e9bd52c67bc..07f0f0140e112 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 @@ -3094,13 +3094,7 @@ protected CompletableFuture internalGetPersistence(boolean if (applied) { PersistencePolicies namespacePolicy = getNamespacePolicies(namespaceName) .persistence; - return namespacePolicy == null - ? new PersistencePolicies( - pulsar().getConfiguration().getManagedLedgerDefaultEnsembleSize(), - pulsar().getConfiguration().getManagedLedgerDefaultWriteQuorum(), - pulsar().getConfiguration().getManagedLedgerDefaultAckQuorum(), - pulsar().getConfiguration().getManagedLedgerDefaultMarkDeleteRateLimit()) - : namespacePolicy; + return getOrDefaultPersistencePolicy(namespacePolicy); } return null; })); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java index 1266a49146e56..ce3525b258d42 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java @@ -24,6 +24,7 @@ import static org.mockito.Mockito.verify; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; @@ -3119,4 +3120,16 @@ public void testPeekEncryptedMessages() throws Exception { assertEquals(peekedMessages.get(i).getData(), receivedMessages.get(i).getData()); } } + + @Test + public void testGetPersistenceAPI() throws Exception { + String namespace = "prop-xyz/test/persistence"; + admin.namespaces().createNamespace(namespace); + PersistencePolicies persistencePolicies = admin.namespaces().getPersistence(namespace); + assertNotNull(persistencePolicies); + assertEquals(persistencePolicies.getBookkeeperEnsemble(), conf.getManagedLedgerDefaultEnsembleSize()); + assertEquals(persistencePolicies.getBookkeeperWriteQuorum(), conf.getManagedLedgerDefaultWriteQuorum()); + assertEquals(persistencePolicies.getBookkeeperAckQuorum(), conf.getManagedLedgerDefaultAckQuorum()); + } + } From 805e1fb074764e1a6b6a7e2d835d7a72f38cd73e Mon Sep 17 00:00:00 2001 From: rajan Date: Mon, 7 Feb 2022 23:52:51 -0800 Subject: [PATCH 2/3] fix test --- .../test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index 1f31ae34d052b..36188a571ae28 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -390,7 +390,7 @@ public void testSetPersistencePolicies() throws Exception { final String namespace = "prop-xyz/ns2"; admin.namespaces().createNamespace(namespace, Sets.newHashSet("test")); - assertNull(admin.namespaces().getPersistence(namespace)); + assertNotNull(admin.namespaces().getPersistence(namespace)); admin.namespaces().setPersistence(namespace, new PersistencePolicies(3, 3, 3, 10.0)); assertEquals(admin.namespaces().getPersistence(namespace), new PersistencePolicies(3, 3, 3, 10.0)); From b9dbab0dd223dc170c66f1c4ba40a578a00df128 Mon Sep 17 00:00:00 2001 From: rajan Date: Wed, 9 Feb 2022 17:17:08 -0800 Subject: [PATCH 3/3] fix test --- .../test/java/org/apache/pulsar/broker/admin/AdminApiTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java index ce3525b258d42..f9d5d0cc8a2df 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java @@ -749,7 +749,7 @@ public void namespaces() throws Exception { policies.auth_policies.getNamespaceAuthentication().remove("my-role"); assertEquals(admin.namespaces().getPolicies("prop-xyz/ns1"), policies); - assertEquals(admin.namespaces().getPersistence("prop-xyz/ns1"), null); + assertEquals(admin.namespaces().getPersistence("prop-xyz/ns1"), new PersistencePolicies(2, 2, 2, 1.0)); admin.namespaces().setPersistence("prop-xyz/ns1", new PersistencePolicies(3, 2, 1, 10.0)); assertEquals(admin.namespaces().getPersistence("prop-xyz/ns1"), new PersistencePolicies(3, 2, 1, 10.0));