From 12a463a3e809a8e4253f786ee7cfc2118cbaed21 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 23 Oct 2020 12:24:10 +0300 Subject: [PATCH 1/2] Validate retention policy in Admin API Fixes #8346 - prevent setting invalid retention policy where either size or time limit is set to zero while the other limit has a non-zero value. - the reason for this is that setting either size or time limit to zero will effectively disable the retention policy. - it is confusing for the end-user if it's possible to set a value for the other limit while it's ignored when the other limit has the value of zero --- .../broker/admin/impl/NamespacesBase.java | 16 ++++ .../pulsar/broker/admin/NamespacesTest.java | 77 ++++++++++++++++++- 2 files changed, 92 insertions(+), 1 deletion(-) 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 2bdd04caefd6f..df559925bca8b 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 @@ -1753,6 +1753,7 @@ protected void internalRemoveBacklogQuota(BacklogQuotaType backlogQuotaType) { } protected void internalSetRetention(RetentionPolicies retention) { + validateRetentionPolicies(retention); validateNamespacePolicyOperation(namespaceName, PolicyName.RETENTION, PolicyOperation.WRITE); validatePoliciesReadOnlyAccess(); @@ -2480,8 +2481,23 @@ private void validatePolicies(NamespaceName ns, Policies policies) { if (policies.persistence != null) { validatePersistencePolicies(policies.persistence); } + + if (policies.retention_policies != null) { + validateRetentionPolicies(policies.retention_policies); + } } + protected void validateRetentionPolicies(RetentionPolicies retention) { + checkArgument(retention.getRetentionSizeInMB() >= -1, + "Invalid retention policy: size limit must be >= -1"); + checkArgument(retention.getRetentionTimeInMinutes() >= -1, + "Invalid retention policy: time limit must be >= -1"); + checkArgument((retention.getRetentionTimeInMinutes() != 0 && retention.getRetentionSizeInMB() != 0) || + (retention.getRetentionTimeInMinutes() == 0 && retention.getRetentionSizeInMB() == 0), + "Invalid retention policy: Setting a single time or size limit to 0 is invalid when " + + "one of the limits has a non-zero value. Use the value of -1 instead of 0 to ignore a " + + "specific limit. To disable retention both limits must be set to 0."); + } protected int internalGetMaxProducersPerTopic() { validateNamespacePolicyOperation(namespaceName, PolicyName.MAX_PRODUCERS, PolicyOperation.READ); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java index d4d7f885301b4..6509ff94e0626 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/NamespacesTest.java @@ -48,6 +48,7 @@ import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import javax.ws.rs.BadRequestException; import javax.ws.rs.ClientErrorException; import javax.ws.rs.WebApplicationException; import javax.ws.rs.container.AsyncResponse; @@ -1356,4 +1357,78 @@ public void testDeletePartitionedTopicMultipleTimes() throws Exception { // Expected } } -} \ No newline at end of file + + @Test + public void testRetentionPolicyValidation() throws Exception { + String namespace = this.testTenant + "/namespace-" + System.nanoTime(); + + admin.namespaces().createNamespace(namespace, Sets.newHashSet(testLocalCluster)); + + // should pass + admin.namespaces().setRetention(namespace, new RetentionPolicies()); + admin.namespaces().setRetention(namespace, new RetentionPolicies(-1, -1)); + admin.namespaces().setRetention(namespace, new RetentionPolicies(1, 1)); + + // should not pass validation + assertInvalidRetentionPolicy(namespace, 1, 0); + assertInvalidRetentionPolicy(namespace, 0, 1); + assertInvalidRetentionPolicy(namespace, -1, 0); + assertInvalidRetentionPolicy(namespace, 0, -1); + assertInvalidRetentionPolicy(namespace, -2, 1); + assertInvalidRetentionPolicy(namespace, 1, -2); + + admin.namespaces().deleteNamespace(namespace); + } + + private void assertInvalidRetentionPolicy(String namespace, int retentionTimeInMinutes, int retentionSizeInMB) { + try { + RetentionPolicies retention = new RetentionPolicies(retentionTimeInMinutes, retentionSizeInMB); + admin.namespaces().setRetention(namespace, retention); + fail("Validation should have failed for " + retention); + } catch (PulsarAdminException e) { + assertTrue(e.getCause() instanceof BadRequestException); + assertTrue(e.getMessage().startsWith("Invalid retention policy")); + } + } + + @Test + public void testRetentionPolicyValidationAsPartOfAllPolicies() throws Exception { + Policies policies = new Policies(); + policies.replication_clusters = Sets.newHashSet(testLocalCluster); + + assertValidRetentionPolicyAsPartOfAllPolicies(policies, 0, 0); + assertValidRetentionPolicyAsPartOfAllPolicies(policies, -1, -1); + assertValidRetentionPolicyAsPartOfAllPolicies(policies, 1, 1); + + // should not pass validation + assertInvalidRetentionPolicyAsPartOfAllPolicies(policies, 1, 0); + assertInvalidRetentionPolicyAsPartOfAllPolicies(policies, 0, 1); + assertInvalidRetentionPolicyAsPartOfAllPolicies(policies, -1, 0); + assertInvalidRetentionPolicyAsPartOfAllPolicies(policies, 0, -1); + assertInvalidRetentionPolicyAsPartOfAllPolicies(policies, -2, 1); + assertInvalidRetentionPolicyAsPartOfAllPolicies(policies, 1, -2); + } + + private void assertValidRetentionPolicyAsPartOfAllPolicies(Policies policies, int retentionTimeInMinutes, + int retentionSizeInMB) throws PulsarAdminException { + String namespace = this.testTenant + "/namespace-" + System.nanoTime(); + RetentionPolicies retention = new RetentionPolicies(retentionTimeInMinutes, retentionSizeInMB); + policies.retention_policies = retention; + admin.namespaces().createNamespace(namespace, policies); + admin.namespaces().deleteNamespace(namespace); + } + + private void assertInvalidRetentionPolicyAsPartOfAllPolicies(Policies policies, int retentionTimeInMinutes, + int retentionSizeInMB) { + String namespace = this.testTenant + "/namespace-" + System.nanoTime(); + try { + RetentionPolicies retention = new RetentionPolicies(retentionTimeInMinutes, retentionSizeInMB); + policies.retention_policies = retention; + admin.namespaces().createNamespace(namespace, policies); + fail("Validation should have failed for " + retention); + } catch (PulsarAdminException e) { + assertTrue(e.getCause() instanceof BadRequestException); + assertTrue(e.getMessage().startsWith("Invalid retention policy")); + } + } +} From a2a0b146faa4f1a6f101a444e5f6cc238fd8ac18 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sat, 24 Oct 2020 21:28:35 +0300 Subject: [PATCH 2/2] Fix invalid retention policies in tests --- .../pulsar/broker/service/PersistentTopicE2ETest.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java index 5dc3d2ec7b70f..e20c18dc891ae 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicE2ETest.java @@ -744,7 +744,7 @@ public void testGcAndRetentionPolicy() throws Exception { assertNotNull(pulsar.getBrokerService().getTopicReference(topicName)); // Remove retention - admin.namespaces().setRetention("prop/ns-abc", new RetentionPolicies(0, 10)); + admin.namespaces().setRetention("prop/ns-abc", new RetentionPolicies()); Thread.sleep(300); // 2. Topic is not GCed with live connection @@ -787,7 +787,7 @@ public void testInfiniteRetentionPolicy() throws Exception { assertNotNull(pulsar.getBrokerService().getTopicReference(topicName)); // Remove retention - admin.namespaces().setRetention("prop/ns-abc", new RetentionPolicies(0, 10)); + admin.namespaces().setRetention("prop/ns-abc", new RetentionPolicies()); Thread.sleep(300); // 2. Topic is not GCed with live connection @@ -834,7 +834,7 @@ public void testServiceConfigurationRetentionPolicy() throws Exception { assertTrue(pulsar.getBrokerService().getTopicReference(topicName).isPresent()); // Remove retention - admin.namespaces().setRetention("prop/ns-abc", new RetentionPolicies(0, 10)); + admin.namespaces().setRetention("prop/ns-abc", new RetentionPolicies()); Thread.sleep(300); // 2. Topic is not GCed with live connection @@ -1684,4 +1684,4 @@ public void testWithEventTime() throws Exception { assertEquals(msg.getValue(), "test"); assertEquals(msg.getEventTime(), 5); } -} \ No newline at end of file +}