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

lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git


The following commit(s) were added to refs/heads/rocketmq-studio by this push:
     new d7843a3f7 fix(server): return empty results when the topic or DLQ is 
absent on the broker (#4146)
d7843a3f7 is described below

commit d7843a3f7357d8487153f2cc33adf28eba430698
Author: cyberslack_lee <[email protected]>
AuthorDate: Tue Sep 15 19:48:17 2026 +0800

    fix(server): return empty results when the topic or DLQ is absent on the 
broker (#4146)
    
    getQueueOffsets and queryByTopic turned a broker-side TOPIC_NOT_EXIST into 
a 502,
    and the DLQ scan only recognised a missing topic by matching the client's 
message
    text. A topic that does not exist on the broker is an expected state rather 
than a
    gateway error, so degrade both to an empty result. The retry-topic guard 
keeps its
    own branch ahead of the new one so its dedicated warning is not swallowed.
    
    The DLQ check now goes through the shared MqResponseCodes response-code 
helper
    instead of adding a second, DLQ-local classification; the existing text 
match is
    kept as a fallback because some client paths surface the condition with no
    response code at all.
    
    Signed-off-by: enkilee <[email protected]>
---
 .../provider/apache/RocketMQDLQProvider.java       |  9 ++++++
 .../provider/apache/RocketMQMessageProvider.java   |  8 ++++++
 .../provider/apache/RocketMQDLQProviderTest.java   | 32 ++++++++++++++++++++++
 .../apache/RocketMQMessageProviderTest.java        | 19 +++++++++++++
 4 files changed, 68 insertions(+)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
index cdb4b09fd..f3b32cad3 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
@@ -28,6 +28,7 @@ import org.apache.rocketmq.common.message.MessageConst;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageQueue;
 import org.apache.rocketmq.common.topic.TopicValidator;
+import org.apache.rocketmq.remoting.protocol.ResponseCode;
 import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
 import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
 import org.apache.rocketmq.remoting.protocol.body.TopicList;
@@ -35,6 +36,7 @@ import 
org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.common.domain.PageResult;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.common.util.MessagePropertyDisplay;
+import org.apache.rocketmq.studio.common.util.MqResponseCodes;
 import org.apache.rocketmq.studio.common.util.Pagination;
 import org.apache.rocketmq.studio.common.util.SystemTopicFilter;
 import org.apache.rocketmq.studio.instance.dlq.DLQExcelExportResultVO;
@@ -543,8 +545,15 @@ public class RocketMQDLQProvider implements DLQProvider {
      * True when the failure only means "the {@code %DLQ%} topic does not 
exist / has no route or
      * message queue yet" — the expected state for a consumer group that has 
never dead-lettered a
      * message. Any other cause is a real scan failure and must not be 
silently degraded to empty.
+     *
+     * <p>Checked by response code first (shared with the message provider 
through
+     * {@link MqResponseCodes}), then by message text: some client paths 
surface the condition
+     * without a response code, so the text match stays as the fallback rather 
than being replaced.
      */
     private static boolean isDlqTopicMissing(Throwable e) {
+        if (MqResponseCodes.hasResponseCode(e, ResponseCode.TOPIC_NOT_EXIST, 
ResponseCode.NO_MESSAGE)) {
+            return true;
+        }
         Throwable cause = e;
         while (cause != null) {
             String message = cause.getMessage();
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
index b22ebee38..ad98ae062 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
@@ -278,6 +278,10 @@ public class RocketMQMessageProvider implements 
MessageProvider {
                 result.sort(Comparator.comparing(QueueOffsetVO::getBrokerName)
                         .thenComparingInt(QueueOffsetVO::getQueueId));
             } catch (Exception e) {
+                if (MqResponseCodes.hasResponseCode(e, 
ResponseCode.TOPIC_NOT_EXIST)) {
+                    log.info("getQueueOffsets(topic={}) matched nothing ({}), 
returning empty list", topic, e.getMessage());
+                    return Collections.emptyList();
+                }
                 log.warn("getQueueOffsets(topic={}) failed: {}", topic, 
e.getMessage());
                 throw new BusinessException(502, "Failed to get queue offsets: 
" + e.getMessage());
             }
@@ -396,6 +400,10 @@ public class RocketMQMessageProvider implements 
MessageProvider {
                             + "group-matched pull consumer; returning empty. 
cause={}", topic, e.getMessage());
                     return Collections.emptyList();
                 }
+                if (MqResponseCodes.hasResponseCode(e, 
ResponseCode.TOPIC_NOT_EXIST, ResponseCode.NO_MESSAGE)) {
+                    log.info("queryByTopic(topic={}) matched nothing ({}), 
returning empty list", topic, e.getMessage());
+                    return Collections.emptyList();
+                }
                 log.warn("queryByTopic(topic={}) failed: {}", topic, 
e.getMessage());
                 throw new BusinessException(502, "Failed to query messages by 
topic: " + e.getMessage());
             }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
index 5530bfca2..48986b883 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
@@ -29,6 +29,7 @@ import org.apache.rocketmq.common.message.MessageConst;
 import org.apache.rocketmq.common.message.MessageDecoder;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.remoting.protocol.ResponseCode;
 import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
 import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
 import org.apache.rocketmq.remoting.protocol.body.TopicList;
@@ -502,6 +503,37 @@ class RocketMQDLQProviderTest {
         verify(runtimeAdminClientResolver, 
never()).executeProducer(anyString(), any());
     }
 
+    @Test
+    void listMessagesDegradesToEmptyWhenClientReportsNoMessageTest() throws 
Exception {
+        String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+        // NO_MESSAGE carries neither "can not find message queue" nor "no 
topic route info", so
+        // only the response-code check recognises it as a missing DLQ topic.
+        when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic))
+                .thenThrow(new MQClientException(ResponseCode.NO_MESSAGE,
+                        "query message by key finished, but no message."));
+
+        PageResult<DLQMessageVO> page = provider.listMessages("instance-a", 
"group-a", 100L, 200L, 1, 20);
+
+        assertThat(page.getTotal()).isZero();
+        assertThat(page.getItems()).isEmpty();
+        verify(pullConsumer, never()).pull(any(MessageQueue.class), 
anyString(), anyLong(), anyInt());
+    }
+
+    @Test
+    void resendMessagesThrowsNotFoundWhenClientReportsNoMessageTest() throws 
Exception {
+        String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+        when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic))
+                .thenThrow(new MQClientException(ResponseCode.NO_MESSAGE,
+                        "query message by key finished, but no message."));
+
+        assertThatThrownBy(() -> provider.resendMessages("instance-a", 
"group-a", 100L, 200L, null))
+                .isInstanceOf(BusinessException.class)
+                .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(404));
+        verify(auditService).record(eq("RESEND_DLQ"), eq("DLQ"), 
eq("group-a"), isNull(),
+                contains("dlqTopicMissing=true"), eq("NOT_FOUND"));
+        verify(runtimeAdminClientResolver, 
never()).executeProducer(anyString(), any());
+    }
+
     @Test
     void resendMessagesUsesPooledClientsForScanAndResend() throws Exception {
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
index da307f9b1..50e9c0ba5 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
@@ -498,6 +498,16 @@ class RocketMQMessageProviderTest {
         assertThat(messages).isEmpty();
     }
 
+    @Test
+    void queryByTopicReturnsEmptyListWhenTopicNotExistTest() throws Exception {
+        when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+                .thenThrow(new MQClientException(ResponseCode.TOPIC_NOT_EXIST,
+                        "No topic route info in name server for the topic: 
TopicA"));
+
+        assertThat(provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, 100L, 
200L)).isEmpty();
+    }
+
     @Test
     @Timeout(value = 1, unit = TimeUnit.SECONDS)
     void queryByTopicStopsWhenPullOffsetDoesNotAdvance() throws Exception {
@@ -1012,6 +1022,15 @@ class RocketMQMessageProviderTest {
                 .containsExactly(0, 1);
     }
 
+    @Test
+    void getQueueOffsetsReturnsEmptyListWhenTopicNotExistTest() throws 
Exception {
+        when(adminExt.examineTopicRouteInfo("TopicA"))
+                .thenThrow(new MQClientException(ResponseCode.TOPIC_NOT_EXIST,
+                        "No topic route info in name server for the topic: 
TopicA"));
+
+        assertThat(provider.getQueueOffsets("instance-a", "TopicA")).isEmpty();
+    }
+
     private MQClientAPIImpl mockOffsetLookupClient() {
         DefaultMQAdminExtImpl adminExtImpl = mock(DefaultMQAdminExtImpl.class);
         MQClientInstance clientInstance = mock(MQClientInstance.class);

Reply via email to