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);