From 18b6cfa685846a8cba442ee9be36d6b5fb27351a Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 2 Aug 2026 23:43:00 -0700 Subject: [PATCH 1/4] [ISSUE #10772] Skip malformed current broker data --- .../service/admin/DefaultAdminService.java | 8 +++++- .../admin/DefaultAdminServiceTest.java | 26 ++++++++++++++++++- 2 files changed, 32 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java index f3c68eab5c4..68660acbab9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java @@ -96,7 +96,13 @@ public boolean createTopicOnBroker(String topic, int wQueueNum, int rQueueNum, L Set curBrokerAddr = new HashSet<>(); if (curBrokerDataList != null) { for (BrokerData brokerData : curBrokerDataList) { - curBrokerAddr.add(brokerData.getBrokerAddrs().get(MixAll.MASTER_ID)); + if (brokerData == null || brokerData.getBrokerAddrs() == null) { + continue; + } + String addr = brokerData.getBrokerAddrs().get(MixAll.MASTER_ID); + if (addr != null) { + curBrokerAddr.add(addr); + } } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java index cdfc7f7fc23..a24302cdb34 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java @@ -17,6 +17,7 @@ package org.apache.rocketmq.proxy.service.admin; +import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.Set; @@ -87,6 +88,29 @@ public void testCreateTopic() throws Exception { assertEquals(8, topicConfigArgumentCaptor.getValue().getReadQueueNums()); } + @Test + public void testCreateTopicOnBrokerSkipsMalformedCurrentBrokerData() throws Exception { + BrokerData malformedBrokerData = new BrokerData(); + + ArgumentCaptor addrArgumentCaptor = ArgumentCaptor.forClass(String.class); + ArgumentCaptor topicConfigArgumentCaptor = ArgumentCaptor.forClass(TopicConfig.class); + doNothing().when(mqClientAPIExt) + .createTopic(addrArgumentCaptor.capture(), anyString(), topicConfigArgumentCaptor.capture(), anyLong()); + + assertTrue(defaultAdminService.createTopicOnBroker( + "createTopic", + 7, + 8, + Collections.singletonList(malformedBrokerData), + createTopicRouteData(1).getBrokerDatas(), + false, + 0 + )); + + assertEquals("127.0.0.1:10911", addrArgumentCaptor.getValue()); + assertEquals("createTopic", topicConfigArgumentCaptor.getValue().getTopicName()); + } + private TopicRouteData createTopicRouteData(int brokerNum) { TopicRouteData topicRouteData = new TopicRouteData(); for (int i = 0; i < brokerNum; i++) { @@ -100,4 +124,4 @@ private TopicRouteData createTopicRouteData(int brokerNum) { } return topicRouteData; } -} \ No newline at end of file +} From 6a5a3de9508188f1fbcbc07fb7b123aa0c102cd3 Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 04:02:55 -0700 Subject: [PATCH 2/4] test(proxy): verify malformed broker is skipped once --- .../rocketmq/proxy/service/admin/DefaultAdminServiceTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java index a24302cdb34..ffdbbd29c91 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java @@ -107,7 +107,8 @@ public void testCreateTopicOnBrokerSkipsMalformedCurrentBrokerData() throws Exce 0 )); - assertEquals("127.0.0.1:10911", addrArgumentCaptor.getValue()); + assertEquals(1, addrArgumentCaptor.getAllValues().size()); + assertEquals("127.0.0.1:10911", addrArgumentCaptor.getAllValues().get(0)); assertEquals("createTopic", topicConfigArgumentCaptor.getValue().getTopicName()); } From fb41c9a06101a44c14bc4cbb030f4ff65c95423c Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 08:17:02 -0700 Subject: [PATCH 3/4] fix(proxy): surface unexpected topic route lookup errors --- .../service/admin/DefaultAdminService.java | 8 ++++++-- .../service/admin/DefaultAdminServiceTest.java | 17 +++++++++++++++++ 2 files changed, 23 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java index 68660acbab9..fb72f3ea37c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java @@ -49,8 +49,12 @@ public boolean topicExist(String topic) { try { topicRouteData = this.getTopicRouteDataDirectlyFromNameServer(topic); topicExist = topicRouteData != null; - } catch (Throwable e) { - topicExist = false; + } catch (Exception e) { + if (TopicRouteHelper.isTopicNotExistError(e)) { + topicExist = false; + } else { + throw new IllegalStateException("get topic route " + topic + " failed", e); + } } return topicExist; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java index ffdbbd29c91..84d2e821329 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java @@ -36,6 +36,7 @@ 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.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; @@ -112,6 +113,22 @@ public void testCreateTopicOnBrokerSkipsMalformedCurrentBrokerData() throws Exce assertEquals("createTopic", topicConfigArgumentCaptor.getValue().getTopicName()); } + @Test + public void testTopicExistReturnsFalseForNotFound() throws Exception { + when(mqClientAPIExt.getTopicRouteInfoFromNameServer(eq("missingTopic"), anyLong())) + .thenThrow(new MQClientException(ResponseCode.TOPIC_NOT_EXIST, "topic not exist")); + + assertFalse(defaultAdminService.topicExist("missingTopic")); + } + + @Test(expected = IllegalStateException.class) + public void testTopicExistThrowsForUnexpectedRouteLookupFailure() throws Exception { + when(mqClientAPIExt.getTopicRouteInfoFromNameServer(eq("brokenTopic"), anyLong())) + .thenThrow(new MQClientException(ResponseCode.SYSTEM_ERROR, "namesrv unavailable")); + + defaultAdminService.topicExist("brokenTopic"); + } + private TopicRouteData createTopicRouteData(int brokerNum) { TopicRouteData topicRouteData = new TopicRouteData(); for (int i = 0; i < brokerNum; i++) { From b7652d878f8c30be1901eba2b15d9b17b531d761 Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 03:50:01 -0700 Subject: [PATCH 4/4] test(proxy): retain topic route lookup failure context --- .../proxy/service/admin/DefaultAdminService.java | 2 +- .../proxy/service/admin/DefaultAdminServiceTest.java | 12 +++++++++--- 2 files changed, 10 insertions(+), 4 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java index fb72f3ea37c..961aeabbaa2 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminService.java @@ -53,7 +53,7 @@ public boolean topicExist(String topic) { if (TopicRouteHelper.isTopicNotExistError(e)) { topicExist = false; } else { - throw new IllegalStateException("get topic route " + topic + " failed", e); + throw new IllegalStateException("get topic route for topic='" + topic + "' failed", e); } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java index 84d2e821329..b559bcbcb7e 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/admin/DefaultAdminServiceTest.java @@ -28,6 +28,7 @@ import org.apache.rocketmq.remoting.protocol.route.TopicRouteData; import org.apache.rocketmq.client.impl.mqclient.MQClientAPIExt; import org.apache.rocketmq.client.impl.mqclient.MQClientAPIFactory; +import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -121,12 +122,17 @@ public void testTopicExistReturnsFalseForNotFound() throws Exception { assertFalse(defaultAdminService.topicExist("missingTopic")); } - @Test(expected = IllegalStateException.class) + @Test public void testTopicExistThrowsForUnexpectedRouteLookupFailure() throws Exception { + MQClientException cause = new MQClientException(ResponseCode.SYSTEM_ERROR, "namesrv unavailable"); when(mqClientAPIExt.getTopicRouteInfoFromNameServer(eq("brokenTopic"), anyLong())) - .thenThrow(new MQClientException(ResponseCode.SYSTEM_ERROR, "namesrv unavailable")); + .thenThrow(cause); + + IllegalStateException exception = Assert.assertThrows(IllegalStateException.class, + () -> defaultAdminService.topicExist("brokenTopic")); - defaultAdminService.topicExist("brokenTopic"); + assertEquals("get topic route for topic='brokenTopic' failed", exception.getMessage()); + assertEquals(cause, exception.getCause()); } private TopicRouteData createTopicRouteData(int brokerNum) {