This is an automated email from the ASF dual-hosted git repository.

RongtongJin pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git


The following commit(s) were added to refs/heads/develop by this push:
     new 290d440e7d [ISSUE #11087] Validate lite.bind.topic for LiteTopic 
groups (#11088)
290d440e7d is described below

commit 290d440e7d7beefb18740968f7129694fa016537
Author: Xiao Yang <[email protected]>
AuthorDate: Mon Sep 14 10:23:13 2026 +0800

    [ISSUE #11087] Validate lite.bind.topic for LiteTopic groups (#11088)
    
    * [ISSUE #11087] Validate lite.bind.topic for LiteTopic groups
    
    * Update
---
 .../rocketmq/broker/lite/LiteMetadataUtil.java     |   5 +-
 .../rocketmq/broker/lite/LiteMetadataUtilTest.java | 112 +++++++++++++++++++++
 .../broker/processor/AdminBrokerProcessorTest.java |  73 ++++++++++++++
 .../common/SubscriptionGroupAttributes.java        |  13 ++-
 .../rocketmq/common/attribute/StringAttribute.java |  11 +-
 .../common/SubscriptionGroupAttributesTest.java    |  43 ++++++++
 6 files changed, 250 insertions(+), 7 deletions(-)

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 92aadfb6f0..5dfdd5c31c 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.Map;
 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 class LiteMetadataUtil {
         }
         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 class LiteMetadataUtil {
             
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 0000000000..ab204c9e18
--- /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<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"));
+    }
+}
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 006979ce86..a573a25211 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.GroupList;
 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.TimerMessageStore;
 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 class AdminBrokerProcessorTest {
         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 class AdminBrokerProcessorTest {
         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 3329188f8a..aaa355c3d5 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 java.util.Map;
 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.topic.TopicValidator;
 
 public class SubscriptionGroupAttributes {
 
@@ -40,7 +41,13 @@ public class SubscriptionGroupAttributes {
 
     public static final StringAttribute LITE_BIND_TOPIC_ATTRIBUTE = new 
StringAttribute(
         "lite.bind.topic",
-        true
+        true,
+        value -> {
+            TopicValidator.ValidateResult result = 
TopicValidator.validateTopic(value);
+            if (!result.isValid()) {
+                throw new RuntimeException(result.getRemark());
+            }
+        }
     );
 
     public static final EnumAttribute LITE_SUB_MODEL_ATTRIBUTE = new 
EnumAttribute(
@@ -97,4 +104,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/StringAttribute.java
 
b/common/src/main/java/org/apache/rocketmq/common/attribute/StringAttribute.java
index e66d688c78..e2a0afe6b7 100644
--- 
a/common/src/main/java/org/apache/rocketmq/common/attribute/StringAttribute.java
+++ 
b/common/src/main/java/org/apache/rocketmq/common/attribute/StringAttribute.java
@@ -17,16 +17,23 @@
 
 package org.apache.rocketmq.common.attribute;
 
-import static com.google.common.base.Preconditions.checkNotNull;
+import com.google.common.base.Preconditions;
+import java.util.function.Consumer;
 
 public class StringAttribute extends Attribute {
+    private final Consumer<String> validator;
 
     public StringAttribute(String name, boolean changeable) {
+        this(name, changeable, Preconditions::checkNotNull);
+    }
+
+    public StringAttribute(String name, boolean changeable, Consumer<String> 
validator) {
         super(name, changeable);
+        this.validator = validator;
     }
 
     @Override
     public void verify(String value) {
-        checkNotNull(value);
+        validator.accept(value);
     }
 }
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 0000000000..4dc68787b8
--- /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("");
+    }
+}

Reply via email to