From 209ea10318c277214c9cbd43157783b1f0225fe6 Mon Sep 17 00:00:00 2001 From: water <672684719@qq.com> Date: Tue, 11 Aug 2026 11:25:45 +0800 Subject: [PATCH 1/2] fix: validate every Lite consumer subscription against its bound topic --- .../apache/rocketmq/proxy/processor/ClientProcessor.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java index c73e66416da..c91b8a41ae1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java @@ -197,8 +197,10 @@ protected void validateLiteSubTopic(ProxyContext ctx, String group, Set Date: Tue, 11 Aug 2026 11:25:48 +0800 Subject: [PATCH 2/2] test: add coverage for validating every Lite subscription --- .../proxy/processor/ClientProcessorTest.java | 21 +++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ClientProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ClientProcessorTest.java index 6644341e551..7eb7da0098e 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ClientProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ClientProcessorTest.java @@ -136,6 +136,27 @@ public void testValidateLiteSubTopic_validSubList_noException() { }); } + @Test + public void testValidateLiteSubTopic_mismatchedSecondTopic_throwsException() { + String group = "group"; + String bindTopic = "topic1"; + SubscriptionData matching = new SubscriptionData(); + matching.setTopic(bindTopic); + SubscriptionData mismatching = new SubscriptionData(); + mismatching.setTopic("topic2"); + Set subList = new HashSet<>(); + subList.add(matching); + subList.add(mismatching); + + when(groupConfig.getLiteBindTopic()).thenReturn(bindTopic); + when(messagingProcessor.getSubscriptionGroupConfig(ctx, group)).thenReturn(groupConfig); + + GrpcProxyException exception = assertThrows(GrpcProxyException.class, () -> { + clientProcessor.validateLiteSubTopic(ctx, group, subList); + }); + assertTrue(exception.getMessage().contains("expected to bind topic")); + } + @Test public void testValidateLiteBindTopic_matchingTopics_noException() { String group = "group";