From 4bcdb361c96c46314aa96fee6bb6dc3db7ff8309 Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 08:17:02 -0700 Subject: [PATCH 1/2] fix(proxy): surface unexpected topic route lookup errors --- .../service/admin/DefaultAdminService.java | 8 ++++++-- .../admin/DefaultAdminServiceTest.java | 19 ++++++++++++++++++- 2 files changed, 24 insertions(+), 3 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..bdf53774af0 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 cdfc7f7fc23..55f6fcf6a99 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 @@ -35,6 +35,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; @@ -87,6 +88,22 @@ public void testCreateTopic() throws Exception { assertEquals(8, topicConfigArgumentCaptor.getValue().getReadQueueNums()); } + @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++) { @@ -100,4 +117,4 @@ private TopicRouteData createTopicRouteData(int brokerNum) { } return topicRouteData; } -} \ No newline at end of file +} From c0fac7f1b24c0dcf2c85bf83710e8a590b05c478 Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 03:50:01 -0700 Subject: [PATCH 2/2] 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 bdf53774af0..d1ac5e0d07c 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 55f6fcf6a99..4996b75a6b8 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 @@ -27,6 +27,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; @@ -96,12 +97,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) {