From 73d4db300bfaa2db9c3bc7ec18ecb25561d3492f Mon Sep 17 00:00:00 2001 From: yx9o Date: Wed, 9 Sep 2026 12:04:27 +0800 Subject: [PATCH] [ISSUE #11087] Validate lite.bind.topic for LiteTopic groups --- .../broker/lite/LiteMetadataUtil.java | 5 +- .../broker/lite/LiteMetadataUtilTest.java | 112 ++++++++++++++++++ .../processor/AdminBrokerProcessorTest.java | 73 ++++++++++++ .../common/SubscriptionGroupAttributes.java | 7 +- .../common/attribute/TopicNameAttribute.java | 34 ++++++ .../SubscriptionGroupAttributesTest.java | 43 +++++++ 6 files changed, 269 insertions(+), 5 deletions(-) create mode 100644 broker/src/test/java/org/apache/rocketmq/broker/lite/LiteMetadataUtilTest.java create mode 100644 common/src/main/java/org/apache/rocketmq/common/attribute/TopicNameAttribute.java create mode 100644 common/src/test/java/org/apache/rocketmq/common/SubscriptionGroupAttributesTest.java diff --git a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteMetadataUtil.java b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteMetadataUtil.java index 92aadfb6f0f..5dfdd5c31cf 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteMetadataUtil.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteMetadataUtil.java @@ -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; @@ -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) { @@ -135,7 +136,7 @@ public static Map> 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()) diff --git a/broker/src/test/java/org/apache/rocketmq/broker/lite/LiteMetadataUtilTest.java b/broker/src/test/java/org/apache/rocketmq/broker/lite/LiteMetadataUtilTest.java new file mode 100644 index 00000000000..ab204c9e185 --- /dev/null +++ b/broker/src/test/java/org/apache/rocketmq/broker/lite/LiteMetadataUtilTest.java @@ -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 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> result = LiteMetadataUtil.getSubscriberGroupMap(brokerController); + + assertFalse(result.containsKey("")); + assertFalse(result.containsKey(" ")); + assertFalse(result.containsKey(null)); + assertEquals(Collections.singleton("liteGroup"), result.get("parentTopic")); + } +} diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java index 006979ce865..a573a25211c 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessorTest.java @@ -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; @@ -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; @@ -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 @@ -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(); diff --git a/common/src/main/java/org/apache/rocketmq/common/SubscriptionGroupAttributes.java b/common/src/main/java/org/apache/rocketmq/common/SubscriptionGroupAttributes.java index 3329188f8aa..f50a62fa657 100644 --- a/common/src/main/java/org/apache/rocketmq/common/SubscriptionGroupAttributes.java +++ b/common/src/main/java/org/apache/rocketmq/common/SubscriptionGroupAttributes.java @@ -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 { @@ -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 ); @@ -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); } -} \ No newline at end of file +} diff --git a/common/src/main/java/org/apache/rocketmq/common/attribute/TopicNameAttribute.java b/common/src/main/java/org/apache/rocketmq/common/attribute/TopicNameAttribute.java new file mode 100644 index 00000000000..c619db3b70b --- /dev/null +++ b/common/src/main/java/org/apache/rocketmq/common/attribute/TopicNameAttribute.java @@ -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()); + } + } +} diff --git a/common/src/test/java/org/apache/rocketmq/common/SubscriptionGroupAttributesTest.java b/common/src/test/java/org/apache/rocketmq/common/SubscriptionGroupAttributesTest.java new file mode 100644 index 00000000000..4dc68787b83 --- /dev/null +++ b/common/src/test/java/org/apache/rocketmq/common/SubscriptionGroupAttributesTest.java @@ -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(""); + } +}