From c81537a4fe90084a8418fbc1d9aac74137eb2b08 Mon Sep 17 00:00:00 2001 From: Shoothzj Date: Wed, 15 Dec 2021 08:41:02 +0800 Subject: [PATCH 1/2] Fix incompatibility of GetBacklogQuota --- .../common/policies/data/BacklogQuota.java | 9 +++ .../policies/data/impl/BacklogQuotaImpl.java | 74 +++++++++++++++++-- .../policies/data/BacklogQuotaMixIn.java | 26 ------- .../common/util/ObjectMapperFactory.java | 2 - .../common/util/ObjectMapperFactoryTest.java | 27 ------- .../BacklogQuotaCompatibilityTest.java | 54 +++++++++++++- 6 files changed, 126 insertions(+), 66 deletions(-) delete mode 100644 pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/BacklogQuotaMixIn.java diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/BacklogQuota.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/BacklogQuota.java index d4b5c4bba1c5b..4604710c3a68c 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/BacklogQuota.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/BacklogQuota.java @@ -28,6 +28,15 @@ */ public interface BacklogQuota { + /** + * Gets quota limit in size. + * Remains for compatible + * + * @return quota limit in bytes + */ + @Deprecated + long getLimit(); + /** * Gets quota limit in size. * diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/BacklogQuotaImpl.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/BacklogQuotaImpl.java index 591e8b8c95a8a..3971f3ddfcc0e 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/BacklogQuotaImpl.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/BacklogQuotaImpl.java @@ -18,23 +18,83 @@ */ package org.apache.pulsar.common.policies.data.impl; -import lombok.AllArgsConstructor; -import lombok.Data; +import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; +import lombok.ToString; import org.apache.pulsar.common.policies.data.BacklogQuota; -@Data -@AllArgsConstructor +@ToString +@EqualsAndHashCode @NoArgsConstructor public class BacklogQuotaImpl implements BacklogQuota { public static final long BYTES_IN_GIGABYTE = 1024 * 1024 * 1024; - // backlog quota by size in byte - private long limitSize; - // backlog quota by time in second + /** + * backlog quota by size in byte, remains for compatible. + */ + @Deprecated + private Long limit; + + /** + * backlog quota by size in byte. + */ + private Long limitSize; + + /** + * backlog quota by time in second. + */ private int limitTime; private RetentionPolicy policy; + public BacklogQuotaImpl(long limitSize, int limitTime, RetentionPolicy policy) { + this.limitSize = limitSize; + this.limitTime = limitTime; + this.policy = policy; + } + + @Deprecated + public long getLimit() { + if (limitSize == null) { + // the limitSize and limit can't be both null + return limit; + } + return limitSize; + } + + @Deprecated + public void setLimit(long limit) { + this.limit = limit; + this.limitSize = limit; + } + + public long getLimitSize() { + if (limitSize == null) { + // the limitSize and limit can't be both null + return limit; + } + return limitSize; + } + + public void setLimitSize(long limitSize) { + this.limitSize = limitSize; + } + + public int getLimitTime() { + return limitTime; + } + + public void setLimitTime(int limitTime) { + this.limitTime = limitTime; + } + + public RetentionPolicy getPolicy() { + return policy; + } + + public void setPolicy(RetentionPolicy policy) { + this.policy = policy; + } + public static BacklogQuotaImplBuilder builder() { return new BacklogQuotaImplBuilder(); } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/BacklogQuotaMixIn.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/BacklogQuotaMixIn.java deleted file mode 100644 index a1562400f0818..0000000000000 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/BacklogQuotaMixIn.java +++ /dev/null @@ -1,26 +0,0 @@ -/** - * 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.common.policies.data; - -import com.fasterxml.jackson.annotation.JsonAlias; - -public abstract class BacklogQuotaMixIn { - @JsonAlias("limit") - private long limitSize; -} diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/util/ObjectMapperFactory.java b/pulsar-common/src/main/java/org/apache/pulsar/common/util/ObjectMapperFactory.java index 94e1b7af4a278..ef2e4894a721f 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/util/ObjectMapperFactory.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/util/ObjectMapperFactory.java @@ -37,7 +37,6 @@ import org.apache.pulsar.common.policies.data.AutoSubscriptionCreationOverride; import org.apache.pulsar.common.policies.data.AutoTopicCreationOverride; import org.apache.pulsar.common.policies.data.BacklogQuota; -import org.apache.pulsar.common.policies.data.BacklogQuotaMixIn; import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BookieInfo; import org.apache.pulsar.common.policies.data.BookiesClusterInfo; @@ -192,7 +191,6 @@ private static void setAnnotationsModule(ObjectMapper mapper) { resolver.addMapping(AutoSubscriptionCreationOverride.class, AutoSubscriptionCreationOverrideImpl.class); // we use MixIn class to add jackson annotations - mapper.addMixIn(BacklogQuotaImpl.class, BacklogQuotaMixIn.class); mapper.addMixIn(ResourceQuota.class, ResourceQuotaMixIn.class); mapper.addMixIn(FunctionConfig.class, JsonIgnorePropertiesMixIn.class); mapper.addMixIn(FunctionState.class, JsonIgnorePropertiesMixIn.class); diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/util/ObjectMapperFactoryTest.java b/pulsar-common/src/test/java/org/apache/pulsar/common/util/ObjectMapperFactoryTest.java index 61f54699ba239..466585d03d515 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/util/ObjectMapperFactoryTest.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/util/ObjectMapperFactoryTest.java @@ -19,39 +19,12 @@ package org.apache.pulsar.common.util; import com.fasterxml.jackson.databind.ObjectMapper; -import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.ResourceQuota; import org.apache.pulsar.common.stats.Metrics; import org.testng.Assert; import org.testng.annotations.Test; public class ObjectMapperFactoryTest { - @Test - public void testBacklogQuotaMixIn() { - ObjectMapper objectMapper = ObjectMapperFactory.getThreadLocal(); - String json = "{\"limit\":10,\"limitTime\":0,\"policy\":\"producer_request_hold\"}"; - try { - BacklogQuota backlogQuota = objectMapper.readValue(json, BacklogQuota.class); - Assert.assertEquals(backlogQuota.getLimitSize(), 10); - Assert.assertEquals(backlogQuota.getLimitTime(), 0); - Assert.assertEquals(backlogQuota.getPolicy(), BacklogQuota.RetentionPolicy.producer_request_hold); - } catch (Exception ex) { - Assert.fail("shouldn't have thrown exception", ex); - } - - try { - String expectJson = "{\"limitSize\":10,\"limitTime\":0,\"policy\":\"producer_request_hold\"}"; - BacklogQuota backlogQuota = BacklogQuota.builder() - .limitSize(10) - .limitTime(0) - .retentionPolicy(BacklogQuota.RetentionPolicy.producer_request_hold) - .build(); - String writeJson = objectMapper.writeValueAsString(backlogQuota); - Assert.assertEquals(expectJson, writeJson); - } catch (Exception ex) { - Assert.fail("shouldn't have thrown exception", ex); - } - } @Test public void testResourceQuotaMixIn() { diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BacklogQuotaCompatibilityTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BacklogQuotaCompatibilityTest.java index 06765e307ef29..daf1b009f1940 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BacklogQuotaCompatibilityTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BacklogQuotaCompatibilityTest.java @@ -19,15 +19,64 @@ package org.apache.pulsar.metadata; import static org.testng.Assert.assertEquals; + +import com.fasterxml.jackson.databind.JavaType; import com.fasterxml.jackson.databind.type.TypeFactory; import java.io.IOException; +import java.util.HashMap; + import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.Policies; +import org.apache.pulsar.common.policies.data.impl.BacklogQuotaImpl; import org.apache.pulsar.metadata.cache.impl.JSONMetadataSerdeSimpleType; +import org.testng.Assert; import org.testng.annotations.Test; public class BacklogQuotaCompatibilityTest { + private final JavaType typeRef = TypeFactory.defaultInstance().constructSimpleType(Policies.class, null); + + private final JSONMetadataSerdeSimpleType simpleType = new JSONMetadataSerdeSimpleType<>(typeRef); + + private final BacklogQuota.RetentionPolicy testPolicy = BacklogQuota.RetentionPolicy.consumer_backlog_eviction; + + @Test + public void testV27SetV28Read() throws Exception { + Policies writePolicy = new Policies(); + BacklogQuotaImpl writeBacklogQuota = new BacklogQuotaImpl(); + writeBacklogQuota.setLimit(1024); + writeBacklogQuota.setLimitTime(60); + writeBacklogQuota.setPolicy(testPolicy); + HashMap quotaHashMap = new HashMap<>(); + quotaHashMap.put(BacklogQuota.BacklogQuotaType.destination_storage, writeBacklogQuota); + writePolicy.backlog_quota_map = quotaHashMap; + byte[] serialize = simpleType.serialize("/path", writePolicy); + Policies policies = simpleType.deserialize("/path", serialize, null); + BacklogQuota readBacklogQuota = policies.backlog_quota_map.get(BacklogQuota.BacklogQuotaType.destination_storage); + Assert.assertEquals(readBacklogQuota.getLimitSize(), 1024); + Assert.assertEquals(readBacklogQuota.getLimitTime(), 60); + Assert.assertEquals(readBacklogQuota.getPolicy(), testPolicy); + } + + @Test + public void testV28SetV27Read() throws Exception { + JSONMetadataSerdeSimpleType simpleType = new JSONMetadataSerdeSimpleType<>(typeRef); + Policies writePolicy = new Policies(); + BacklogQuotaImpl writeBacklogQuota = new BacklogQuotaImpl(); + writeBacklogQuota.setLimitSize(1024); + writeBacklogQuota.setLimitTime(60); + writeBacklogQuota.setPolicy(testPolicy); + HashMap quotaHashMap = new HashMap<>(); + quotaHashMap.put(BacklogQuota.BacklogQuotaType.destination_storage, writeBacklogQuota); + writePolicy.backlog_quota_map = quotaHashMap; + byte[] serialize = simpleType.serialize("/path", writePolicy); + Policies policies = simpleType.deserialize("/path", serialize, null); + BacklogQuota readBacklogQuota = policies.backlog_quota_map.get(BacklogQuota.BacklogQuotaType.destination_storage); + Assert.assertEquals(readBacklogQuota.getLimit(), 1024); + Assert.assertEquals(readBacklogQuota.getLimitTime(), 60); + Assert.assertEquals(readBacklogQuota.getPolicy(), testPolicy); + } + @Test public void testBackwardCompatibility() throws IOException { String oldPolicyStr = "{\"auth_policies\":{\"namespace_auth\":{},\"destination_auth\":{}," @@ -41,10 +90,7 @@ public void testBackwardCompatibility() throws IOException { + "\"schema_auto_update_compatibility_strategy\":\"Full\",\"schema_compatibility_strategy\":" + "\"UNDEFINED\",\"is_allow_auto_update_schema\":true,\"schema_validation_enforced\":false," + "\"subscription_types_enabled\":[]}\n"; - - JSONMetadataSerdeSimpleType jsonMetadataSerdeSimpleType = new JSONMetadataSerdeSimpleType( - TypeFactory.defaultInstance().constructSimpleType(Policies.class, null)); - Policies policies = (Policies) jsonMetadataSerdeSimpleType.deserialize(null, oldPolicyStr.getBytes(), null); + Policies policies = simpleType.deserialize(null, oldPolicyStr.getBytes(), null); assertEquals(policies.backlog_quota_map.get(BacklogQuota.BacklogQuotaType.destination_storage).getLimitSize(), 1001); assertEquals(policies.backlog_quota_map.get(BacklogQuota.BacklogQuotaType.destination_storage).getLimitTime(), From 4560e801d378bd8f94b39f1e9a7a6ef3a7efa61c Mon Sep 17 00:00:00 2001 From: Shoothzj Date: Wed, 15 Dec 2021 15:21:34 +0800 Subject: [PATCH 2/2] fix incompatibilty set --- .../common/policies/data/impl/BacklogQuotaImpl.java | 3 +++ .../metadata/BacklogQuotaCompatibilityTest.java | 12 +++++++++--- 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/BacklogQuotaImpl.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/BacklogQuotaImpl.java index 3971f3ddfcc0e..3d97fa06426f9 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/BacklogQuotaImpl.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/policies/data/impl/BacklogQuotaImpl.java @@ -31,6 +31,8 @@ public class BacklogQuotaImpl implements BacklogQuota { /** * backlog quota by size in byte, remains for compatible. + * for the details: https://github.com/apache/pulsar/pull/13291 + * @since 2.9.1 */ @Deprecated private Long limit; @@ -77,6 +79,7 @@ public long getLimitSize() { public void setLimitSize(long limitSize) { this.limitSize = limitSize; + this.limit = limitSize; } public int getLimitTime() { diff --git a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BacklogQuotaCompatibilityTest.java b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BacklogQuotaCompatibilityTest.java index daf1b009f1940..659825c74858a 100644 --- a/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BacklogQuotaCompatibilityTest.java +++ b/pulsar-metadata/src/test/java/org/apache/pulsar/metadata/BacklogQuotaCompatibilityTest.java @@ -41,7 +41,7 @@ public class BacklogQuotaCompatibilityTest { private final BacklogQuota.RetentionPolicy testPolicy = BacklogQuota.RetentionPolicy.consumer_backlog_eviction; @Test - public void testV27SetV28Read() throws Exception { + public void testV27ClientSetV28BrokerRead() throws Exception { Policies writePolicy = new Policies(); BacklogQuotaImpl writeBacklogQuota = new BacklogQuotaImpl(); writeBacklogQuota.setLimit(1024); @@ -59,8 +59,7 @@ public void testV27SetV28Read() throws Exception { } @Test - public void testV28SetV27Read() throws Exception { - JSONMetadataSerdeSimpleType simpleType = new JSONMetadataSerdeSimpleType<>(typeRef); + public void testV28ClientSetV28BrokerRead() throws Exception { Policies writePolicy = new Policies(); BacklogQuotaImpl writeBacklogQuota = new BacklogQuotaImpl(); writeBacklogQuota.setLimitSize(1024); @@ -77,6 +76,13 @@ public void testV28SetV27Read() throws Exception { Assert.assertEquals(readBacklogQuota.getPolicy(), testPolicy); } + @Test + public void testV28ClientSetV27BrokerRead() { + BacklogQuotaImpl writeBacklogQuota = new BacklogQuotaImpl(); + writeBacklogQuota.setLimitSize(1024); + Assert.assertEquals(1024, writeBacklogQuota.getLimit()); + } + @Test public void testBackwardCompatibility() throws IOException { String oldPolicyStr = "{\"auth_policies\":{\"namespace_auth\":{},\"destination_auth\":{},"