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
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,15 @@
*/
public interface BacklogQuota {

/**
* Gets quota limit in size.
* Remains for compatible
*
* @return quota limit in bytes
*/
@Deprecated
Comment thread
hezhangjian marked this conversation as resolved.
long getLimit();

/**
* Gets quota limit in size.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,23 +18,86 @@
*/
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.
* for the details: https://github.com/apache/pulsar/pull/13291
* @since 2.9.1
*/
@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;
Comment thread
hezhangjian marked this conversation as resolved.
this.limit = 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();
}
Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,15 +19,70 @@
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<Policies> simpleType = new JSONMetadataSerdeSimpleType<>(typeRef);

private final BacklogQuota.RetentionPolicy testPolicy = BacklogQuota.RetentionPolicy.consumer_backlog_eviction;

@Test
public void testV27ClientSetV28BrokerRead() throws Exception {
Policies writePolicy = new Policies();
BacklogQuotaImpl writeBacklogQuota = new BacklogQuotaImpl();
writeBacklogQuota.setLimit(1024);
writeBacklogQuota.setLimitTime(60);
writeBacklogQuota.setPolicy(testPolicy);
HashMap<BacklogQuota.BacklogQuotaType, BacklogQuota> 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 testV28ClientSetV28BrokerRead() throws Exception {
Policies writePolicy = new Policies();
BacklogQuotaImpl writeBacklogQuota = new BacklogQuotaImpl();
writeBacklogQuota.setLimitSize(1024);
writeBacklogQuota.setLimitTime(60);
writeBacklogQuota.setPolicy(testPolicy);
HashMap<BacklogQuota.BacklogQuotaType, BacklogQuota> 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 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\":{},"
Expand All @@ -41,10 +96,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(),
Expand Down