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 b0a239ee3 fix(message): message query lifecycle integrity (#2230)
b0a239ee3 is described below

commit b0a239ee3cd2dd98bc4b5e191af8d7a238559a59
Author: shown <[email protected]>
AuthorDate: Wed Aug 19 11:19:47 2026 +0800

    fix(message): message query lifecycle integrity (#2230)
    
    * fix(message): surface topic loading failures
    
    Signed-off-by: yuluo-yx <[email protected]>
    
    * fix(message): record provider-neutral query history
    
    Signed-off-by: yuluo-yx <[email protected]>
    
    * fix(message): validate query time bounds safely
    
    Signed-off-by: yuluo-yx <[email protected]>
    
    * fix(message): clear stale query history state
    
    Signed-off-by: yuluo-yx <[email protected]>
    
    * fix(message): normalize query history filters
    
    Signed-off-by: yuluo-yx <[email protected]>
    
    ---------
    
    Signed-off-by: yuluo-yx <[email protected]>
---
 .../studio/instance/message/MessageService.java    | 41 +++++++++++--
 .../instance/message/QueryHistoryController.java   | 13 +++-
 .../provider/apache/RocketMQMessageProvider.java   | 26 --------
 .../instance/message/MessageServiceTest.java       | 50 ++++++++++++++--
 .../message/QueryHistoryControllerTest.java        | 24 ++++++++
 .../apache/RocketMQMessageProviderTest.java        | 21 +------
 web/src/components/MessageQueryHistoryDrawer.tsx   | 27 ++++++++-
 .../__tests__/MessageQueryHistoryDrawer.test.tsx   | 20 +++++++
 .../pages/instance/__tests__/MessagePage.test.tsx  | 19 +++++-
 web/src/pages/instance/message.tsx                 | 70 ++++++++++++++++------
 10 files changed, 229 insertions(+), 82 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
index e4602bbe4..b5319be5f 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
@@ -34,14 +34,17 @@ public class MessageService {
 
     private final MessageProvider messageProvider;
     private final InstanceProviderRegistry providerRegistry;
+    private final QueryHistoryService queryHistoryService;
 
     public List<MessageRecordVO> queryMessages(
             String instanceId, String topic, String msgId, String tag, String 
key, Long startTime, Long endTime) {
         validateTopicQueryWindow(topic, msgId, key, startTime, endTime);
         log.info("Querying messages: topic={}, msgId={}, tag={}, key={}", 
topic, msgId, tag, key);
-        return providerRegistry.byInstanceId(instanceId)
+        List<MessageRecordVO> result = 
providerRegistry.byInstanceId(instanceId)
                 .map(provider -> provider.queryMessages(instanceId, topic, 
msgId, tag, key, startTime, endTime))
                 .orElseGet(() -> messageProvider.queryMessages(instanceId, 
topic, msgId, tag, key, startTime, endTime));
+        recordMessageQuery(instanceId, topic, msgId, tag, key, startTime, 
endTime, result.size());
+        return result;
     }
 
     public TraceRecordVO getMessageTrace(String instanceId, String msgId, 
String topic) {
@@ -49,9 +52,33 @@ public class MessageService {
             throw new BusinessException(400, "msgId is required");
         }
         log.info("Getting message trace: msgId={}, topic={}", msgId, topic);
-        return providerRegistry.byInstanceId(instanceId)
+        TraceRecordVO result = providerRegistry.byInstanceId(instanceId)
                 .map(provider -> provider.getMessageTrace(instanceId, msgId, 
topic))
                 .orElseGet(() -> messageProvider.getMessageTrace(instanceId, 
msgId, topic));
+        recordTraceQuery(instanceId, msgId, topic, result);
+        return result;
+    }
+
+    private void recordMessageQuery(String instanceId, String topic, String 
msgId, String tag,
+                                    String key, Long startTime, Long endTime, 
int resultCount) {
+        String queryType = StringUtils.hasText(msgId) ? "MSG_ID" : 
StringUtils.hasText(key) ? "KEY" : "TOPIC";
+        try {
+            queryHistoryService.recordMessageQuery(instanceId, queryType, 
topic, msgId, tag, key,
+                    startTime, endTime, resultCount);
+        } catch (RuntimeException failure) {
+            log.warn("Failed to record message query history: {}", 
failure.getMessage());
+        }
+    }
+
+    private void recordTraceQuery(String instanceId, String msgId, String 
topic, TraceRecordVO result) {
+        int nodeCount = result == null || result.getNodes() == null ? 0 : 
result.getNodes().size();
+        int consumerCount = result == null || result.getConsumerStatus() == 
null ? 0
+                : result.getConsumerStatus().size();
+        try {
+            queryHistoryService.recordTraceQuery(instanceId, msgId, topic, 
nodeCount, consumerCount);
+        } catch (RuntimeException failure) {
+            log.warn("Failed to record trace query history: {}", 
failure.getMessage());
+        }
     }
 
     private void validateTopicQueryWindow(String topic, String msgId, String 
key, Long startTime, Long endTime) {
@@ -72,10 +99,14 @@ public class MessageService {
         }
         long end = endTime == null ? System.currentTimeMillis() : endTime;
         long start = startTime == null ? end - 60 * 60 * 1000L : startTime;
-        if (start > end) {
-            throw new BusinessException(400, "startTime must not be after 
endTime");
+        if (start < 0 || end < 0) {
+            throw new BusinessException(400, "message query timestamps must 
not be negative");
+        }
+        if (start >= end) {
+            throw new BusinessException(400, "startTime must be before 
endTime");
         }
-        if (end - start > MAX_TOPIC_QUERY_WINDOW_MILLIS) {
+        // Compare without subtracting untrusted endpoints; end - start can 
overflow long.
+        if (start < end - MAX_TOPIC_QUERY_WINDOW_MILLIS) {
             throw new BusinessException(400, "topic query time range must not 
exceed 7 days");
         }
     }
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/QueryHistoryController.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/QueryHistoryController.java
index d5ed62d2b..017e4290e 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/QueryHistoryController.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/QueryHistoryController.java
@@ -14,6 +14,7 @@ import org.springframework.web.bind.annotation.GetMapping;
 import org.springframework.web.bind.annotation.RequestMapping;
 import org.springframework.web.bind.annotation.RequestParam;
 import org.springframework.web.bind.annotation.RestController;
+import org.springframework.util.StringUtils;
 
 @RestController
 @RequestMapping("/api/query-history")
@@ -31,7 +32,8 @@ public class QueryHistoryController {
             @RequestParam(defaultValue = "20") int pageSize) {
         validatePage(page, pageSize);
         return Result.ok(queryHistoryService.listMessageQueries(
-                clusterId, queryType, search, page, pageSize));
+                normalizeFilter(clusterId), normalizeFilter(queryType), 
normalizeFilter(search),
+                page, pageSize));
     }
 
     @GetMapping("/traces")
@@ -41,12 +43,13 @@ public class QueryHistoryController {
             @RequestParam(defaultValue = "1") int page,
             @RequestParam(defaultValue = "20") int pageSize) {
         validatePage(page, pageSize);
-        return Result.ok(queryHistoryService.listTraceQueries(clusterId, 
search, page, pageSize));
+        return Result.ok(queryHistoryService.listTraceQueries(
+                normalizeFilter(clusterId), normalizeFilter(search), page, 
pageSize));
     }
 
     @GetMapping("/summary")
     public Result<QueryHistorySummaryVO> summary(@RequestParam(required = 
false) String clusterId) {
-        return Result.ok(queryHistoryService.summarize(clusterId));
+        return 
Result.ok(queryHistoryService.summarize(normalizeFilter(clusterId)));
     }
 
     private void validatePage(int page, int pageSize) {
@@ -57,4 +60,8 @@ public class QueryHistoryController {
             throw new BusinessException(400, "pageSize must be between 1 and " 
+ MAX_PAGE_SIZE);
         }
     }
+
+    private String normalizeFilter(String value) {
+        return StringUtils.hasText(value) ? value.trim() : null;
+    }
 }
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 38846bfe5..cf9f43bd3 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
@@ -34,7 +34,6 @@ import 
org.apache.rocketmq.studio.instance.message.MessageProvider;
 import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
 import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
 import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
-import org.apache.rocketmq.studio.instance.message.QueryHistoryService;
 import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
 import org.apache.rocketmq.tools.admin.api.MessageTrack;
 import org.apache.rocketmq.tools.admin.api.TrackType;
@@ -91,7 +90,6 @@ public class RocketMQMessageProvider implements 
MessageProvider {
             .thenComparing(MessageRecordVO::getMsgId, 
Comparator.nullsFirst(String::compareTo));
 
     private final RuntimeAdminClientResolver runtimeAdminClientResolver;
-    private final QueryHistoryService queryHistoryService;
 
     @Override
     public List<MessageRecordVO> queryMessages(String instanceId, String 
topic, String msgId, String tag, String key,
@@ -115,25 +113,20 @@ public class RocketMQMessageProvider implements 
MessageProvider {
         }
 
         List<MessageRecordVO> result;
-        String queryType;
         if (StringUtils.hasText(msgId)) {
-            queryType = "MSG_ID";
             result = queryByMsgId(adminExt, topic, msgId);
         } else if (StringUtils.hasText(topic) && StringUtils.hasText(key)) {
-            queryType = "KEY";
             result = queryByKey(adminExt, topic, key, tag, begin, end);
         } else if (StringUtils.hasText(topic)) {
             if (begin >= 0 && end >= 0 && end - begin > 
MAX_TOPIC_QUERY_WINDOW_MILLIS) {
                 throw new BusinessException(400, "Topic message query time 
range must not exceed 7 days");
             }
-            queryType = "TOPIC";
             result = queryByTopic(endpoint, credentialHook, topic, tag, begin, 
end, DEFAULT_TOPIC_LIMIT);
         } else {
             log.warn("queryMessages requires at least one of msgId/topic, 
returning empty list");
             return Collections.emptyList();
         }
 
-        recordMessageQuery(instanceId, queryType, topic, msgId, tag, key, 
startTime, endTime, result.size());
         return result;
     }
 
@@ -340,7 +333,6 @@ public class RocketMQMessageProvider implements 
MessageProvider {
             throw new BusinessException(502, "Failed to query message trace: " 
+ e.getMessage());
         }
 
-        recordTraceQuery(instanceId, msgId, null, nodes.size(), 
consumerStatus.size());
         return TraceRecordVO.builder()
                 .nodes(nodes)
                 .consumerStatus(consumerStatus)
@@ -619,24 +611,6 @@ public class RocketMQMessageProvider implements 
MessageProvider {
         return consumer;
     }
 
-    private void recordMessageQuery(String instanceId, String queryType, 
String topic, String msgId, String tag, String key,
-                                    Long startTime, Long endTime, int 
resultCount) {
-        try {
-            queryHistoryService.recordMessageQuery(instanceId, queryType, 
topic, msgId, tag, key,
-                    startTime, endTime, resultCount);
-        } catch (Exception e) {
-            log.warn("Failed to record message query history: {}", 
e.getMessage());
-        }
-    }
-
-    private void recordTraceQuery(String instanceId, String msgId, String 
topic, int nodeCount, int consumerCount) {
-        try {
-            queryHistoryService.recordTraceQuery(instanceId, msgId, topic, 
nodeCount, consumerCount);
-        } catch (Exception e) {
-            log.warn("Failed to record trace query history: {}", 
e.getMessage());
-        }
-    }
-
     private static TraceRecordVO emptyTrace() {
         return TraceRecordVO.builder()
                 .nodes(Collections.emptyList())
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
index 754f26476..d30b43dd5 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
@@ -12,11 +12,17 @@ package org.apache.rocketmq.studio.instance.message;
 
 import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
+import org.apache.rocketmq.studio.provider.InstanceProvider;
 import org.junit.jupiter.api.Test;
 
+import java.util.List;
+import java.util.Optional;
+
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
 
 class MessageServiceTest {
 
@@ -24,7 +30,7 @@ class MessageServiceTest {
     void rejectsKeyQueryWithoutTopicBeforeCallingProvider() {
         MessageProvider provider = mock(MessageProvider.class);
         InstanceProviderRegistry registry = 
mock(InstanceProviderRegistry.class);
-        MessageService service = new MessageService(provider, registry);
+        MessageService service = new MessageService(provider, registry, 
mock(QueryHistoryService.class));
 
         assertThatThrownBy(() -> service.queryMessages(null, null, null, null, 
"order-1", null, null))
                 .isInstanceOf(BusinessException.class)
@@ -37,7 +43,7 @@ class MessageServiceTest {
     void rejectsMessageIdQueryWithoutTopicBeforeCallingProvider() {
         MessageProvider provider = mock(MessageProvider.class);
         InstanceProviderRegistry registry = 
mock(InstanceProviderRegistry.class);
-        MessageService service = new MessageService(provider, registry);
+        MessageService service = new MessageService(provider, registry, 
mock(QueryHistoryService.class));
 
         assertThatThrownBy(() -> service.queryMessages(null, null, "msg-001", 
null, null, null, null))
                 .isInstanceOf(BusinessException.class)
@@ -50,7 +56,7 @@ class MessageServiceTest {
     void rejectsBlankMessageTraceIdBeforeCallingProvider() {
         MessageProvider provider = mock(MessageProvider.class);
         InstanceProviderRegistry registry = 
mock(InstanceProviderRegistry.class);
-        MessageService service = new MessageService(provider, registry);
+        MessageService service = new MessageService(provider, registry, 
mock(QueryHistoryService.class));
 
         assertThatThrownBy(() -> service.getMessageTrace("instance-a", "  ", 
null))
                 .isInstanceOf(BusinessException.class)
@@ -63,11 +69,11 @@ class MessageServiceTest {
     void rejectsReversedTopicQueryWindowBeforeCallingProvider() {
         MessageProvider provider = mock(MessageProvider.class);
         InstanceProviderRegistry registry = 
mock(InstanceProviderRegistry.class);
-        MessageService service = new MessageService(provider, registry);
+        MessageService service = new MessageService(provider, registry, 
mock(QueryHistoryService.class));
 
         assertThatThrownBy(() -> service.queryMessages("instance-a", "TopicA", 
null, null, null, 200L, 100L))
                 .isInstanceOf(BusinessException.class)
-                .hasMessage("startTime must not be after endTime");
+                .hasMessage("startTime must be before endTime");
 
         verifyNoInteractions(provider);
     }
@@ -76,7 +82,7 @@ class MessageServiceTest {
     void rejectsTopicQueryWindowLongerThanSevenDaysBeforeCallingProvider() {
         MessageProvider provider = mock(MessageProvider.class);
         InstanceProviderRegistry registry = 
mock(InstanceProviderRegistry.class);
-        MessageService service = new MessageService(provider, registry);
+        MessageService service = new MessageService(provider, registry, 
mock(QueryHistoryService.class));
 
         assertThatThrownBy(() -> service.queryMessages("instance-a", "TopicA", 
null, null, null, 0L,
                 8L * 24 * 60 * 60 * 1000))
@@ -85,4 +91,36 @@ class MessageServiceTest {
 
         verifyNoInteractions(provider);
     }
+
+    @Test
+    void recordsProviderNeutralMessageQueryHistory() {
+        MessageProvider fallback = mock(MessageProvider.class);
+        InstanceProvider provider = mock(InstanceProvider.class);
+        InstanceProviderRegistry registry = 
mock(InstanceProviderRegistry.class);
+        QueryHistoryService history = mock(QueryHistoryService.class);
+        MessageService service = new MessageService(fallback, registry, 
history);
+        
when(registry.byInstanceId("cloud-instance")).thenReturn(Optional.of(provider));
+        when(provider.queryMessages("cloud-instance", "orders", null, null, 
"ORDER-1", null, null))
+                
.thenReturn(List.of(MessageRecordVO.builder().msgId("msg-1").build()));
+
+        service.queryMessages("cloud-instance", "orders", null, null, 
"ORDER-1", null, null);
+
+        verify(history).recordMessageQuery("cloud-instance", "KEY", "orders", 
null, null,
+                "ORDER-1", null, null, 1);
+        verifyNoInteractions(fallback);
+    }
+
+    @Test
+    void rejectsOverflowingTopicQueryWindowBeforeCallingProvider() {
+        MessageProvider provider = mock(MessageProvider.class);
+        InstanceProviderRegistry registry = 
mock(InstanceProviderRegistry.class);
+        MessageService service = new MessageService(provider, registry, 
mock(QueryHistoryService.class));
+
+        assertThatThrownBy(() -> service.queryMessages("instance-a", "TopicA", 
null, null, null,
+                0L, Long.MAX_VALUE))
+                .isInstanceOf(BusinessException.class)
+                .hasMessage("topic query time range must not exceed 7 days");
+
+        verifyNoInteractions(provider, registry);
+    }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/QueryHistoryControllerTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/QueryHistoryControllerTest.java
index f9fb7ae36..94d3a3c6e 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/QueryHistoryControllerTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/QueryHistoryControllerTest.java
@@ -73,4 +73,28 @@ class QueryHistoryControllerTest {
                 .andExpect(jsonPath("$.data.messageQueries").value(7))
                 .andExpect(jsonPath("$.data.traceQueries").value(3));
     }
+
+    @Test
+    void normalizesOptionalHistoryFilters() throws Exception {
+        when(queryHistoryService.listMessageQueries("instance-a", "TOPIC", 
null, 1, 20))
+                .thenReturn(PageResult.of(List.of(), 0, 1, 20));
+        when(queryHistoryService.listTraceQueries("instance-a", null, 1, 20))
+                .thenReturn(PageResult.of(List.of(), 0, 1, 20));
+
+        mockMvc.perform(get("/api/query-history/messages")
+                        .param("clusterId", "  instance-a  ")
+                        .param("queryType", " TOPIC ")
+                        .param("search", "   "))
+                .andExpect(status().isOk());
+        mockMvc.perform(get("/api/query-history/traces")
+                        .param("clusterId", " instance-a ")
+                        .param("search", "\t"))
+                .andExpect(status().isOk());
+        mockMvc.perform(get("/api/query-history/summary").param("clusterId", " 
  "))
+                .andExpect(status().isOk());
+
+        verify(queryHistoryService).listMessageQueries("instance-a", "TOPIC", 
null, 1, 20);
+        verify(queryHistoryService).listTraceQueries("instance-a", null, 1, 
20);
+        verify(queryHistoryService).summarize(null);
+    }
 }
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 3d8dae308..126d72d25 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
@@ -35,7 +35,6 @@ import 
org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
 import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
 import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
 import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
-import org.apache.rocketmq.studio.instance.message.QueryHistoryService;
 import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
 import org.apache.rocketmq.tools.admin.DefaultMQAdminExtImpl;
 import org.junit.jupiter.api.BeforeEach;
@@ -82,9 +81,6 @@ class RocketMQMessageProviderTest {
     @Mock
     private RuntimeAdminClientResolver runtimeAdminClientResolver;
 
-    @Mock
-    private QueryHistoryService queryHistoryService;
-
     private RocketMQMessageProvider provider;
 
     @BeforeEach
@@ -94,7 +90,7 @@ class RocketMQMessageProviderTest {
             MqAdminExtFactory.AdminAction<Object> action = 
invocation.getArgument(1);
             return action == null ? null : action.apply(adminExt);
         });
-        provider = new RocketMQMessageProvider(runtimeAdminClientResolver, 
queryHistoryService);
+        provider = new RocketMQMessageProvider(runtimeAdminClientResolver);
     }
 
     @Test
@@ -125,8 +121,6 @@ class RocketMQMessageProviderTest {
         verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
         verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
         verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
-        verify(queryHistoryService).recordMessageQuery("instance-a", "TOPIC", 
"TopicA", null, null, null,
-                100L, 200L, 0);
     }
 
     @Test
@@ -164,8 +158,6 @@ class RocketMQMessageProviderTest {
         verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
         verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
         verify(adminExt, never()).queryMessage(anyString(), anyString(), 
anyInt(), anyLong(), anyLong());
-        verify(queryHistoryService, never()).recordMessageQuery(anyString(), 
anyString(), anyString(),
-                anyString(), anyString(), anyString(), any(), any(), anyInt());
     }
 
     @Test
@@ -179,8 +171,6 @@ class RocketMQMessageProviderTest {
         verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
         verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
         verify(adminExt, never()).queryMessage(anyString(), anyString(), 
anyInt(), anyLong(), anyLong());
-        verify(queryHistoryService, never()).recordMessageQuery(anyString(), 
anyString(), anyString(),
-                anyString(), anyString(), anyString(), any(), any(), anyInt());
     }
 
     @Test
@@ -190,9 +180,6 @@ class RocketMQMessageProviderTest {
                 .isInstanceOf(BusinessException.class)
                 .hasMessage("Topic message query time range must not exceed 7 
days")
                 .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(400));
-
-        verify(queryHistoryService, never()).recordMessageQuery(anyString(), 
anyString(), anyString(),
-                anyString(), anyString(), anyString(), any(), any(), anyInt());
     }
 
     @Test
@@ -319,8 +306,6 @@ class RocketMQMessageProviderTest {
             verify(consumer).pull(queue, "*", 20L, 32);
             verify(consumer).pull(queue, "*", 40L, 32);
         }
-        verify(queryHistoryService).recordMessageQuery("instance-a", "TOPIC", 
"TopicA", null, null, null,
-                100L, 200L, 1);
     }
 
     @Test
@@ -449,7 +434,6 @@ class RocketMQMessageProviderTest {
         TraceRecordVO record = provider.getMessageTrace("instance-a", 
"msg-123", "orders");
 
         assertThat(record.getNodes()).hasSize(2);
-        verify(queryHistoryService).recordTraceQuery(eq("instance-a"), 
eq("msg-123"), eq(null), eq(2), eq(1));
         TraceNodeVO produce = record.getNodes().get(0);
         assertThat(produce.getTitle()).isEqualTo("produce");
         assertThat(produce.getStatus()).isEqualTo("finish");
@@ -507,7 +491,6 @@ class RocketMQMessageProviderTest {
                 beginCaptor.capture(), endCaptor.capture());
         assertThat(beginCaptor.getValue()).isEqualTo(10_000_000L - 5 * 
60_000L);
         assertThat(endCaptor.getValue()).isGreaterThanOrEqualTo(10_000_000L + 
24 * 3600_000L);
-        verify(queryHistoryService).recordTraceQuery(eq("instance-a"), 
eq("msg-with-topic"), eq(null), eq(0), eq(0));
     }
 
     @Test
@@ -522,7 +505,6 @@ class RocketMQMessageProviderTest {
         verify(adminExt).queryMessage(eq("RMQ_SYS_TRACE_TOPIC"), 
eq("invalid-offset-id"), eq(64),
                 beginCaptor.capture(), endCaptor.capture());
         assertThat(endCaptor.getValue() - 
beginCaptor.getValue()).isBetween(3_660_000L, 3_670_000L);
-        verify(queryHistoryService).recordTraceQuery(eq("instance-a"), 
eq("invalid-offset-id"), eq(null), eq(0), eq(0));
     }
 
     @Test
@@ -570,7 +552,6 @@ class RocketMQMessageProviderTest {
                 .hasMessage("Failed to query message trace: broker 
unavailable")
                 .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(502));
 
-        verify(queryHistoryService, never()).recordTraceQuery(anyString(), 
anyString(), any(), anyInt(), anyInt());
     }
 
     @Test
diff --git a/web/src/components/MessageQueryHistoryDrawer.tsx 
b/web/src/components/MessageQueryHistoryDrawer.tsx
index ca5e2f5b5..846f243de 100644
--- a/web/src/components/MessageQueryHistoryDrawer.tsx
+++ b/web/src/components/MessageQueryHistoryDrawer.tsx
@@ -5,7 +5,7 @@
  * The ASF licenses this file to You under the Apache License, Version 2.0.
  */
 import { useCallback, useEffect, useRef, useState } from 'react';
-import { Alert, Drawer, Flex, Input, Statistic, Table, Tabs, Tag } from 'antd';
+import { Alert, Button, Drawer, Flex, Input, Statistic, Table, Tabs, Tag } 
from 'antd';
 import type { ColumnsType } from 'antd/es/table';
 import {
   getQueryHistorySummary,
@@ -23,7 +23,11 @@ interface Props {
 }
 
 const PAGE_SIZE = 20;
-const formatTime = (value?: string) => (value ? new 
Date(value).toLocaleString() : '-');
+const formatTime = (value?: string) => {
+  if (!value) return '-';
+  const timestamp = new Date(value);
+  return Number.isNaN(timestamp.getTime()) ? '-' : timestamp.toLocaleString();
+};
 
 const MessageQueryHistoryDrawer = ({ open, clusterId, onClose }: Props) => {
   const [tab, setTab] = useState<'messages' | 'traces'>('messages');
@@ -42,6 +46,10 @@ const MessageQueryHistoryDrawer = ({ open, clusterId, 
onClose }: Props) => {
     const id = ++requestId.current;
     setLoading(true);
     setError('');
+    setSummary(undefined);
+    setMessageRows([]);
+    setTraceRows([]);
+    setTotal(0);
     try {
       const [nextSummary, result] = await Promise.all([
         getQueryHistorySummary(clusterId),
@@ -115,7 +123,20 @@ const MessageQueryHistoryDrawer = ({ open, clusterId, 
onClose }: Props) => {
         }}
         style={{ marginBottom: 12, width: 420 }}
       />
-      {error && <Alert type="error" showIcon message={error} style={{ 
marginBottom: 12 }} />}
+      {error && (
+        <Alert
+          type="error"
+          showIcon
+          message="查询历史加载失败"
+          description={error}
+          action={
+            <Button size="small" onClick={() => void load()}>
+              重试
+            </Button>
+          }
+          style={{ marginBottom: 12 }}
+        />
+      )}
       <Tabs
         activeKey={tab}
         onChange={(key) => {
diff --git a/web/src/components/__tests__/MessageQueryHistoryDrawer.test.tsx 
b/web/src/components/__tests__/MessageQueryHistoryDrawer.test.tsx
index 5d6cfa0a5..b1fad3a60 100644
--- a/web/src/components/__tests__/MessageQueryHistoryDrawer.test.tsx
+++ b/web/src/components/__tests__/MessageQueryHistoryDrawer.test.tsx
@@ -88,4 +88,24 @@ describe('MessageQueryHistoryDrawer', () => {
     expect(await screen.findByText('msg-1')).toBeInTheDocument();
     await waitFor(() => expect(listTraceQueryHistory).toHaveBeenCalled());
   });
+
+  it('clears stale rows and offers retry when a new instance load fails', 
async () => {
+    const view = render(
+      <App>
+        <MessageQueryHistoryDrawer open clusterId="instance-a" 
onClose={vi.fn()} />
+      </App>,
+    );
+    expect(await screen.findByText('order-1')).toBeInTheDocument();
+    vi.mocked(getQueryHistorySummary).mockRejectedValueOnce(new Error('network 
unavailable'));
+
+    view.rerender(
+      <App>
+        <MessageQueryHistoryDrawer open clusterId="instance-b" 
onClose={vi.fn()} />
+      </App>,
+    );
+
+    expect(await screen.findByText('查询历史加载失败')).toBeInTheDocument();
+    expect(screen.queryByText('order-1')).not.toBeInTheDocument();
+    expect(screen.getByRole('button', { name: /重\s*试/ })).toBeEnabled();
+  });
 });
diff --git a/web/src/pages/instance/__tests__/MessagePage.test.tsx 
b/web/src/pages/instance/__tests__/MessagePage.test.tsx
index f6691da8a..b692192a3 100644
--- a/web/src/pages/instance/__tests__/MessagePage.test.tsx
+++ b/web/src/pages/instance/__tests__/MessagePage.test.tsx
@@ -118,7 +118,7 @@ describe('Message page query history', () => {
     const queryButton = screen.getByRole('button', { name: /^search查询$/ });
 
     expect(queryButton).toBeDisabled();
-    expect(queryButton).toHaveAttribute('title', '请选择 Topic');
+    await waitFor(() => expect(queryButton).toHaveAttribute('title', '请选择 
Topic'));
 
     await user.click(lastElement(screen.getAllByRole('combobox')));
     await user.click(lastElement(await screen.findAllByText('order-create')));
@@ -167,6 +167,23 @@ describe('Message page query history', () => {
     });
   });
 
+  it('surfaces Topic loading failures and retries without changing instance', 
async () => {
+    const user = userEvent.setup();
+    topicServiceMocks.listTopics
+      .mockReset()
+      .mockRejectedValueOnce(new Error('NameServer unavailable'))
+      .mockResolvedValueOnce([{ name: 'orders' }]);
+
+    renderWithProviders(<MessagePage />);
+
+    expect(await screen.findByText('Topic 列表加载失败')).toBeInTheDocument();
+    expect(screen.getByRole('button', { name: /^search查询$/ })).toBeDisabled();
+    await user.click(screen.getByRole('button', { name: /重\s*试/ }));
+    await waitFor(() => 
expect(topicServiceMocks.listTopics).toHaveBeenCalledTimes(2));
+    await user.click(lastElement(screen.getAllByRole('combobox')));
+    expect(await screen.findAllByText('orders')).not.toHaveLength(0);
+  });
+
   it('requires a topic even when a key or message ID is present', async () => {
     const user = userEvent.setup();
     renderWithProviders(<MessagePage />);
diff --git a/web/src/pages/instance/message.tsx 
b/web/src/pages/instance/message.tsx
index 0a886580b..8a3d21b21 100644
--- a/web/src/pages/instance/message.tsx
+++ b/web/src/pages/instance/message.tsx
@@ -15,7 +15,7 @@
  * limitations under the License.
  */
 
-import { useEffect, useRef, useState } from 'react';
+import { useCallback, useEffect, useRef, useState } from 'react';
 import {
   Alert,
   Card,
@@ -273,27 +273,38 @@ const MessagePageContent = ({
   const { t } = useLang();
   const [topicOptions, setTopicOptions] = useState<string[]>([]);
   const [topicError, setTopicError] = useState<string | null>(null);
+  const [topicLoading, setTopicLoading] = useState(false);
+  const topicRequestId = useRef(0);
 
-  useEffect(() => {
+  const loadTopicOptions = useCallback(async () => {
     if (!selectedInstanceId) {
+      setTopicOptions([]);
+      setTopicError(null);
+      setTopicLoading(false);
       return;
     }
-    let cancelled = false;
-    void listTopics({ instanceId: selectedInstanceId })
-      .then((nextTopics) => {
-        if (cancelled) return;
-        setTopicError(null);
-        setTopicOptions(nextTopics.map((topic) => topic.name));
-      })
-      .catch((error: unknown) => {
-        if (cancelled) return;
-        setTopicOptions([]);
-        setTopicError(error instanceof Error ? error.message : '加载 Topic 
列表失败');
-      });
+    const requestId = ++topicRequestId.current;
+    setTopicLoading(true);
+    setTopicError(null);
+    setTopicOptions([]);
+    try {
+      const nextTopics = await listTopics({ instanceId: selectedInstanceId });
+      if (requestId !== topicRequestId.current) return;
+      setTopicOptions(nextTopics.map((topic) => topic.name));
+    } catch (error: unknown) {
+      if (requestId !== topicRequestId.current) return;
+      setTopicError(error instanceof Error ? error.message : '加载 Topic 列表失败');
+    } finally {
+      if (requestId === topicRequestId.current) setTopicLoading(false);
+    }
+  }, [selectedInstanceId]);
+
+  useEffect(() => {
+    void Promise.resolve().then(loadTopicOptions);
     return () => {
-      cancelled = true;
+      topicRequestId.current += 1;
     };
-  }, [selectedInstanceId]);
+  }, [loadTopicOptions]);
   const [queryMode, setQueryMode] = useState<QueryMode>('topic');
   const [selectedTopic, setSelectedTopic] = useState<string | undefined>();
   const [dateRange, setDateRange] = useState<[Dayjs, Dayjs]>(getDefaultRange);
@@ -328,7 +339,13 @@ const MessagePageContent = ({
         ? { topic: selectedTopic, key: keyInput || undefined }
         : { topic: selectedTopic, msgId: msgIdInput || undefined };
   const queryValidationError = getQueryValidationError(queryMode, 
currentQueryParams);
-  const queryDisabledReason = !selectedInstanceId ? '请先选择实例' : 
queryValidationError;
+  const queryDisabledReason = !selectedInstanceId
+    ? '请先选择实例'
+    : topicLoading
+      ? '正在加载 Topic 列表'
+      : topicError
+        ? 'Topic 列表加载失败,请先重试'
+        : queryValidationError;
 
   /* ─── Handlers ─── */
   const handleReset = () => {
@@ -789,6 +806,8 @@ const MessagePageContent = ({
                   onChange={setSelectedTopic}
                   allowClear
                   showSearch
+                  loading={topicLoading}
+                  disabled={topicLoading || Boolean(topicError)}
                   options={topicOptions.map((t) => ({
                     value: t,
                     label: t,
@@ -816,6 +835,8 @@ const MessagePageContent = ({
                   onChange={setSelectedTopic}
                   allowClear
                   showSearch
+                  loading={topicLoading}
+                  disabled={topicLoading || Boolean(topicError)}
                   options={topicOptions.map((t) => ({
                     value: t,
                     label: t,
@@ -839,6 +860,8 @@ const MessagePageContent = ({
                   onChange={setSelectedTopic}
                   allowClear
                   showSearch
+                  loading={topicLoading}
+                  disabled={topicLoading || Boolean(topicError)}
                   options={topicOptions.map((t) => ({
                     value: t,
                     label: t,
@@ -887,7 +910,18 @@ const MessagePageContent = ({
       </Card>
 
       {topicError && (
-        <Alert showIcon type="error" message={topicError} style={{ 
marginBottom: 16 }} />
+        <Alert
+          showIcon
+          type="error"
+          message="Topic 列表加载失败"
+          description={topicError}
+          action={
+            <Button size="small" onClick={() => void loadTopicOptions()}>
+              重试
+            </Button>
+          }
+          style={{ marginBottom: 16 }}
+        />
       )}
       <MessageQueryHistoryDrawer
         open={historyDrawerOpen}

Reply via email to