Skip to content
Open
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 @@ -22,6 +22,7 @@
import java.util.Set;
import java.util.concurrent.ConcurrentMap;
import java.util.stream.Collectors;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.broker.BrokerController;
import org.apache.rocketmq.common.TopicConfig;
import org.apache.rocketmq.common.attribute.TopicMessageType;
Expand Down Expand Up @@ -52,7 +53,7 @@ public static boolean isLiteGroupType(String group, BrokerController brokerContr
}
SubscriptionGroupConfig groupConfig =
brokerController.getSubscriptionGroupManager().findSubscriptionGroupConfig(group);
return null != groupConfig && groupConfig.getLiteBindTopic() != null;
return null != groupConfig && StringUtils.isNotBlank(groupConfig.getLiteBindTopic());
}

public static String getLiteBindTopic(String group, BrokerController brokerController) {
Expand Down Expand Up @@ -135,7 +136,7 @@ public static Map<String, Set<String>> getSubscriberGroupMap(BrokerController br
brokerController.getSubscriptionGroupManager().getSubscriptionGroupTable();

return groupTable.entrySet().stream()
.filter(entry -> entry.getValue().getLiteBindTopic() != null)
.filter(entry -> StringUtils.isNotBlank(entry.getValue().getLiteBindTopic()))
.collect(Collectors.groupingBy(
entry -> entry.getValue().getLiteBindTopic(),
Collectors.mapping(Map.Entry::getKey, Collectors.toSet())
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
/*
* 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.rocketmq.broker.lite;

import java.util.Collections;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import org.apache.rocketmq.broker.BrokerController;
import org.apache.rocketmq.broker.subscription.SubscriptionGroupManager;
import org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mock;
import org.mockito.junit.MockitoJUnitRunner;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import static org.mockito.Mockito.when;

@RunWith(MockitoJUnitRunner.class)
public class LiteMetadataUtilTest {

@Mock
private BrokerController brokerController;

@Mock
private SubscriptionGroupManager subscriptionGroupManager;

@Before
public void setUp() {
when(brokerController.getSubscriptionGroupManager()).thenReturn(subscriptionGroupManager);
}

@Test
public void testIsLiteGroupTypeTreatsEmptyBindTopicAsNonLite() {
SubscriptionGroupConfig emptyBindGroup = new SubscriptionGroupConfig();
emptyBindGroup.setGroupName("emptyBindGroup");
emptyBindGroup.setLiteBindTopic("");

SubscriptionGroupConfig blankBindGroup = new SubscriptionGroupConfig();
blankBindGroup.setGroupName("blankBindGroup");
blankBindGroup.setLiteBindTopic(" ");

SubscriptionGroupConfig liteGroup = new SubscriptionGroupConfig();
liteGroup.setGroupName("liteGroup");
liteGroup.setLiteBindTopic("parentTopic");

when(subscriptionGroupManager.findSubscriptionGroupConfig("normalGroup"))
.thenReturn(new SubscriptionGroupConfig());
when(subscriptionGroupManager.findSubscriptionGroupConfig("emptyBindGroup"))
.thenReturn(emptyBindGroup);
when(subscriptionGroupManager.findSubscriptionGroupConfig("blankBindGroup"))
.thenReturn(blankBindGroup);
when(subscriptionGroupManager.findSubscriptionGroupConfig("liteGroup"))
.thenReturn(liteGroup);

assertFalse(LiteMetadataUtil.isLiteGroupType("missingGroup", brokerController));
assertFalse(LiteMetadataUtil.isLiteGroupType("normalGroup", brokerController));
assertFalse(LiteMetadataUtil.isLiteGroupType("emptyBindGroup", brokerController));
assertFalse(LiteMetadataUtil.isLiteGroupType("blankBindGroup", brokerController));
assertTrue(LiteMetadataUtil.isLiteGroupType("liteGroup", brokerController));
}

@Test
public void testGetSubscriberGroupMapSkipsEmptyBindTopic() {
ConcurrentMap<String, SubscriptionGroupConfig> groupTable = new ConcurrentHashMap<>();
groupTable.put("normalGroup", new SubscriptionGroupConfig());

SubscriptionGroupConfig emptyBindGroup = new SubscriptionGroupConfig();
emptyBindGroup.setGroupName("emptyBindGroup");
emptyBindGroup.setLiteBindTopic("");
groupTable.put("emptyBindGroup", emptyBindGroup);

SubscriptionGroupConfig blankBindGroup = new SubscriptionGroupConfig();
blankBindGroup.setGroupName("blankBindGroup");
blankBindGroup.setLiteBindTopic(" ");
groupTable.put("blankBindGroup", blankBindGroup);

SubscriptionGroupConfig liteGroup = new SubscriptionGroupConfig();
liteGroup.setGroupName("liteGroup");
liteGroup.setLiteBindTopic("parentTopic");
groupTable.put("liteGroup", liteGroup);

when(subscriptionGroupManager.getSubscriptionGroupTable()).thenReturn(groupTable);

Map<String, Set<String>> result = LiteMetadataUtil.getSubscriberGroupMap(brokerController);

assertFalse(result.containsKey(""));
assertFalse(result.containsKey(" "));
assertFalse(result.containsKey(null));
assertEquals(Collections.singleton("liteGroup"), result.get("parentTopic"));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@
import org.apache.rocketmq.remoting.protocol.body.HARuntimeInfo;
import org.apache.rocketmq.remoting.protocol.body.LockBatchRequestBody;
import org.apache.rocketmq.remoting.protocol.body.QueryCorrectionOffsetBody;
import org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupList;
import org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
import org.apache.rocketmq.remoting.protocol.body.TopicConfigSerializeWrapper;
import org.apache.rocketmq.remoting.protocol.body.UnlockBatchRequestBody;
Expand Down Expand Up @@ -146,6 +147,7 @@
import org.apache.rocketmq.store.timer.TimerMetrics;
import org.apache.rocketmq.store.util.LibC;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
Expand Down Expand Up @@ -659,6 +661,41 @@ public void testDeleteSubscriptionGroupListNoCleanOffset() throws Exception {
assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
}

@Test
public void testDeleteSubscriptionGroupWithEmptyLiteBindTopicDoesNotCleanOffset() throws Exception {
String groupName = "GID-EMPTY-LITE-BIND";
SubscriptionGroupConfig groupConfig = new SubscriptionGroupConfig();
groupConfig.setGroupName(groupName);
groupConfig.setLiteBindTopic("");
brokerController.getSubscriptionGroupManager().getSubscriptionGroupTable().put(groupName, groupConfig);
brokerController.setConsumerOffsetManager(consumerOffsetManager);

RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.DELETE_SUBSCRIPTIONGROUP, null);
request.addExtField("groupName", groupName);
request.addExtField("cleanOffset", "false");
RemotingCommand response = adminBrokerProcessor.processRequest(handlerContext, request);

assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
verify(consumerOffsetManager, never()).removeOffset(groupName);
}

@Test
public void testDeleteSubscriptionGroupListWithEmptyLiteBindTopicDoesNotCleanOffset() throws Exception {
brokerController.getBrokerConfig().setBatchDeleteSubscriptionGroupMaxRate(0);
String groupName = "GID-BATCH-EMPTY-LITE-BIND";
SubscriptionGroupConfig groupConfig = new SubscriptionGroupConfig();
groupConfig.setGroupName(groupName);
groupConfig.setLiteBindTopic("");
brokerController.getSubscriptionGroupManager().getSubscriptionGroupTable().put(groupName, groupConfig);
brokerController.setConsumerOffsetManager(consumerOffsetManager);

RemotingCommand request = buildDeleteSubscriptionGroupListRequest(Collections.singletonList(groupName), false);
RemotingCommand response = adminBrokerProcessor.processRequest(handlerContext, request);

assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
verify(consumerOffsetManager, never()).removeOffset(groupName);
}

@Test
public void testDeleteTopicListWithPopRetryTopics() throws Exception {
// When clearRetryTopicWhenDeleteTopic=true, POP retry topics should be collected and deleted
Expand Down Expand Up @@ -1077,6 +1114,42 @@ public void testUpdateAndCreateSubscriptionGroup() throws RemotingCommandExcepti
assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
}

@Test
public void testUpdateAndCreateSubscriptionGroupRejectsEmptyLiteBindTopic() {
String groupName = "GID-EMPTY-LITE-BIND";
SubscriptionGroupConfig subscriptionGroupConfig = new SubscriptionGroupConfig();
subscriptionGroupConfig.setGroupName(groupName);
subscriptionGroupConfig.setAttributes(ImmutableMap.of("+lite.bind.topic", ""));

RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.UPDATE_AND_CREATE_SUBSCRIPTIONGROUP, null);
request.setBody(JSON.toJSON(subscriptionGroupConfig).toString().getBytes(StandardCharsets.UTF_8));

RuntimeException exception = Assert.assertThrows(RuntimeException.class,
() -> adminBrokerProcessor.processRequest(handlerContext, request));

assertThat(exception).hasMessageContaining("The specified topic is blank");
assertFalse(brokerController.getSubscriptionGroupManager().getSubscriptionGroupTable().containsKey(groupName));
}

@Test
public void testUpdateAndCreateSubscriptionGroupListRejectsEmptyLiteBindTopic() {
String groupName = "GID-LIST-EMPTY-LITE-BIND";
SubscriptionGroupConfig subscriptionGroupConfig = new SubscriptionGroupConfig();
subscriptionGroupConfig.setGroupName(groupName);
subscriptionGroupConfig.setAttributes(ImmutableMap.of("+lite.bind.topic", ""));

SubscriptionGroupList subscriptionGroupList =
new SubscriptionGroupList(Collections.singletonList(subscriptionGroupConfig));
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.UPDATE_AND_CREATE_SUBSCRIPTIONGROUP_LIST, null);
request.setBody(subscriptionGroupList.encode());

RuntimeException exception = Assert.assertThrows(RuntimeException.class,
() -> adminBrokerProcessor.processRequest(handlerContext, request));

assertThat(exception).hasMessageContaining("The specified topic is blank");
assertFalse(brokerController.getSubscriptionGroupManager().getSubscriptionGroupTable().containsKey(groupName));
}

@Test
public void testGetAllSubscriptionGroupInRocksdb() throws Exception {
initRocksdbSubscriptionManager();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,10 @@
import org.apache.rocketmq.common.attribute.Attribute;
import org.apache.rocketmq.common.attribute.BooleanAttribute;
import org.apache.rocketmq.common.attribute.EnumAttribute;
import org.apache.rocketmq.common.attribute.LiteSubModel;
import org.apache.rocketmq.common.attribute.LongRangeAttribute;
import org.apache.rocketmq.common.attribute.StringAttribute;
import org.apache.rocketmq.common.attribute.LiteSubModel;
import org.apache.rocketmq.common.attribute.TopicNameAttribute;

public class SubscriptionGroupAttributes {

Expand All @@ -38,7 +39,7 @@ public class SubscriptionGroupAttributes {
100
);

public static final StringAttribute LITE_BIND_TOPIC_ATTRIBUTE = new StringAttribute(
public static final StringAttribute LITE_BIND_TOPIC_ATTRIBUTE = new TopicNameAttribute(
"lite.bind.topic",
true
);
Expand Down Expand Up @@ -97,4 +98,4 @@ public class SubscriptionGroupAttributes {
ALL.put(LITE_SUB_CLIENT_MAX_EVENT_COUNT_ATTRIBUTE.getName(), LITE_SUB_CLIENT_MAX_EVENT_COUNT_ATTRIBUTE);
ALL.put(LITE_SUB_WILDCARD_ATTRIBUTE.getName(), LITE_SUB_WILDCARD_ATTRIBUTE);
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
/*
* 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.rocketmq.common.attribute;

import org.apache.rocketmq.common.topic.TopicValidator;

public class TopicNameAttribute extends StringAttribute {

public TopicNameAttribute(String name, boolean changeable) {
super(name, changeable);
}

@Override
public void verify(String value) {
TopicValidator.ValidateResult result = TopicValidator.validateTopic(value);
if (!result.isValid()) {
throw new RuntimeException(result.getRemark());
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
/*
* 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.rocketmq.common;

import org.junit.Assert;
import org.junit.Test;

public class SubscriptionGroupAttributesTest {

@Test
public void testLiteBindTopicAttributeValidatesTopicName() {
SubscriptionGroupAttributes.LITE_BIND_TOPIC_ATTRIBUTE.verify("parentTopic");
SubscriptionGroupAttributes.LITE_BIND_TOPIC_ATTRIBUTE.verify("parent_topic");

Assert.assertThrows(RuntimeException.class,
() -> SubscriptionGroupAttributes.LITE_BIND_TOPIC_ATTRIBUTE.verify(null));
Assert.assertThrows(RuntimeException.class,
() -> SubscriptionGroupAttributes.LITE_BIND_TOPIC_ATTRIBUTE.verify(""));
Assert.assertThrows(RuntimeException.class,
() -> SubscriptionGroupAttributes.LITE_BIND_TOPIC_ATTRIBUTE.verify(" "));
Assert.assertThrows(RuntimeException.class,
() -> SubscriptionGroupAttributes.LITE_BIND_TOPIC_ATTRIBUTE.verify("parent topic"));
}

@Test
public void testLiteSubWildcardAttributeStillAllowsEmptyValue() {
SubscriptionGroupAttributes.LITE_SUB_WILDCARD_ATTRIBUTE.verify("");
}
}
Loading