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 a004431d feat: resolve message queries by instance endpoint (#1081)
a004431d is described below

commit a004431dc9aefb10093fb8d7ca8fdf99b9b19280
Author: aias00 <[email protected]>
AuthorDate: Thu Aug 6 01:21:11 2026 -0700

    feat: resolve message queries by instance endpoint (#1081)
---
 .../studio/instance/message/MessageController.java |  3 +-
 .../studio/instance/message/MessageProvider.java   |  2 +-
 .../instance/message/MessageProviderStub.java      |  4 +--
 .../studio/instance/message/MessageService.java    |  4 +--
 .../studio/rocketmq/RocketMQMessageProvider.java   | 35 +++++++++++-----------
 .../instance/message/MessageControllerTest.java    |  5 ++--
 .../instance/message/MessageProviderStubTest.java  |  2 +-
 .../instance/message/MessageServiceTest.java       |  4 +--
 .../rocketmq/RocketMQMessageProviderTest.java      | 19 ++++++++++--
 web/src/api/message.test.ts                        |  1 +
 web/src/api/message.ts                             |  1 +
 .../pages/instance/__tests__/MessagePage.test.tsx  | 25 ++++++++++++----
 web/src/pages/instance/message.tsx                 |  2 +-
 13 files changed, 71 insertions(+), 36 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
index 674e166e..6c061a28 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
@@ -35,13 +35,14 @@ public class MessageController {
 
     @GetMapping
     public Result<List<MessageRecordVO>> queryMessages(
+            @RequestParam String instanceId,
             @RequestParam(required = false) String topic,
             @RequestParam(required = false) String msgId,
             @RequestParam(required = false) String tag,
             @RequestParam(required = false) String key,
             @RequestParam(required = false) Long startTime,
             @RequestParam(required = false) Long endTime) {
-        return Result.ok(messageService.queryMessages(topic, msgId, tag, key, 
startTime, endTime));
+        return Result.ok(messageService.queryMessages(instanceId, topic, 
msgId, tag, key, startTime, endTime));
     }
 
     @GetMapping("/{msgId}/trace")
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
index 3b3b5c0b..744d4677 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
@@ -20,7 +20,7 @@ package org.apache.rocketmq.studio.instance.message;
 import java.util.List;
 
 public interface MessageProvider {
-    List<MessageRecordVO> queryMessages(String topic, String msgId, String 
tag, String key, Long startTime,
+    List<MessageRecordVO> queryMessages(String instanceId, String topic, 
String msgId, String tag, String key, Long startTime,
                                         Long endTime);
 
     TraceRecordVO getMessageTrace(String msgId);
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
index 5c2f8f0c..f4352a1d 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
@@ -29,8 +29,8 @@ import java.util.List;
 public class MessageProviderStub implements MessageProvider {
 
     @Override
-    public List<MessageRecordVO> queryMessages(String topic, String msgId, 
String tag, String key, Long startTime,
-                                               Long endTime) {
+    public List<MessageRecordVO> queryMessages(String instanceId, String 
topic, String msgId, String tag, String key,
+                                               Long startTime, Long endTime) {
         log.warn("MessageProviderStub.queryMessages called but no real message 
provider is configured");
         throw unsupported();
     }
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 6e7b7030..549ce4f9 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
@@ -33,10 +33,10 @@ public class MessageService {
     private final MessageProvider messageProvider;
 
     public List<MessageRecordVO> queryMessages(
-            String topic, String msgId, String tag, String key, Long 
startTime, Long endTime) {
+            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 messageProvider.queryMessages(topic, msgId, tag, key, 
startTime, endTime);
+        return messageProvider.queryMessages(instanceId, topic, msgId, tag, 
key, startTime, endTime);
     }
 
     public TraceRecordVO getMessageTrace(String msgId) {
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
index 14d6899a..a08b0e81 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
@@ -25,6 +25,7 @@ import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageId;
 import org.apache.rocketmq.common.message.MessageQueue;
 import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.instance.message.ConsumerStatusVO;
 import org.apache.rocketmq.studio.instance.message.MessageProvider;
 import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
@@ -81,25 +82,27 @@ public class RocketMQMessageProvider implements 
MessageProvider {
     private static final long ONE_DAY_MILLIS = 24 * ONE_HOUR_MILLIS;
 
     private final ObjectProvider<DefaultMQAdminExt> adminExtProvider;
+    private final RuntimeAdminClientResolver runtimeAdminClientResolver;
     private final QueryHistoryService queryHistoryService;
-    private final RocketMQProperties properties;
 
     public RocketMQMessageProvider(ObjectProvider<DefaultMQAdminExt> 
adminExtProvider,
-                                   QueryHistoryService queryHistoryService,
-                                   RocketMQProperties properties) {
+                                   RuntimeAdminClientResolver 
runtimeAdminClientResolver,
+                                   QueryHistoryService queryHistoryService) {
         this.adminExtProvider = adminExtProvider;
+        this.runtimeAdminClientResolver = runtimeAdminClientResolver;
         this.queryHistoryService = queryHistoryService;
-        this.properties = properties;
     }
 
     @Override
-    public List<MessageRecordVO> queryMessages(String topic, String msgId, 
String tag, String key,
+    public List<MessageRecordVO> queryMessages(String instanceId, String 
topic, String msgId, String tag, String key,
                                                Long startTime, Long endTime) {
-        DefaultMQAdminExt adminExt = adminExtProvider.getIfAvailable();
-        if (adminExt == null) {
-            log.warn("DefaultMQAdminExt is not configured, returning empty 
message list");
-            return Collections.emptyList();
-        }
+        String endpoint = 
runtimeAdminClientResolver.resolveEndpoint(instanceId);
+        return runtimeAdminClientResolver.execute(instanceId,
+                adminExt -> queryMessages((DefaultMQAdminExt) adminExt, 
endpoint, topic, msgId, tag, key, startTime, endTime));
+    }
+
+    private List<MessageRecordVO> queryMessages(DefaultMQAdminExt adminExt, 
String endpoint, String topic, String msgId, String tag, String key,
+                                                 Long startTime, Long endTime) 
{
 
         long end = endTime != null ? endTime : System.currentTimeMillis();
         long begin = startTime != null ? startTime : end - ONE_HOUR_MILLIS;
@@ -114,7 +117,7 @@ public class RocketMQMessageProvider implements 
MessageProvider {
             result = queryByKey(adminExt, topic, key, tag, begin, end);
         } else if (StringUtils.hasText(topic)) {
             queryType = "TOPIC";
-            result = queryByTopic(topic, tag, begin, end, DEFAULT_TOPIC_LIMIT);
+            result = queryByTopic(endpoint, topic, tag, begin, end, 
DEFAULT_TOPIC_LIMIT);
         } else {
             log.warn("queryMessages requires at least one of msgId/topic, 
returning empty list");
             return Collections.emptyList();
@@ -189,8 +192,8 @@ public class RocketMQMessageProvider implements 
MessageProvider {
      * Scan a topic within a time range using a short-lived pull consumer, 
mirroring the approach
      * used by the RocketMQ dashboard for time-range topic queries.
      */
-    private List<MessageRecordVO> queryByTopic(String topic, String tag, long 
begin, long end, int limit) {
-        DefaultMQPullConsumer consumer = newPullConsumer("studio-msg-query");
+    private List<MessageRecordVO> queryByTopic(String endpoint, String topic, 
String tag, long begin, long end, int limit) {
+        DefaultMQPullConsumer consumer = newPullConsumer("studio-msg-query", 
endpoint);
         List<MessageRecordVO> result = new ArrayList<>();
         try {
             consumer.start();
@@ -477,12 +480,10 @@ public class RocketMQMessageProvider implements 
MessageProvider {
         return tag.equals(messageExt.getTags());
     }
 
-    private DefaultMQPullConsumer newPullConsumer(String groupPrefix) {
+    private DefaultMQPullConsumer newPullConsumer(String groupPrefix, String 
endpoint) {
         DefaultMQPullConsumer consumer = new DefaultMQPullConsumer(groupPrefix 
+ "-group");
         consumer.setInstanceName(ShortLivedClientName.next(groupPrefix));
-        if (StringUtils.hasText(properties.getNamesrvAddr())) {
-            consumer.setNamesrvAddr(properties.getNamesrvAddr());
-        }
+        consumer.setNamesrvAddr(endpoint);
         return consumer;
     }
 
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
index 49278190..1b1ef293 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
@@ -52,10 +52,11 @@ class MessageControllerTest {
                 .tag("created")
                 .key("order-1")
                 .build();
-        when(messageService.queryMessages(eq("orders"), isNull(), 
eq("created"), eq("order-1"),
+        when(messageService.queryMessages(eq("instance-a"), eq("orders"), 
isNull(), eq("created"), eq("order-1"),
                 eq(1784246400000L), 
eq(1784332800000L))).thenReturn(List.of(message));
 
         mockMvc.perform(get("/api/messages")
+                        .param("instanceId", "instance-a")
                         .param("topic", "orders")
                         .param("tag", "created")
                         .param("key", "order-1")
@@ -66,7 +67,7 @@ class MessageControllerTest {
                 .andExpect(jsonPath("$.data[0].msgId").value("msg-001"))
                 .andExpect(jsonPath("$.data[0].tag").value("created"));
 
-        verify(messageService).queryMessages(eq("orders"), isNull(), 
eq("created"), eq("order-1"),
+        verify(messageService).queryMessages(eq("instance-a"), eq("orders"), 
isNull(), eq("created"), eq("order-1"),
                 eq(1784246400000L), eq(1784332800000L));
     }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageProviderStubTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageProviderStubTest.java
index f33fd853..bb8d4109 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageProviderStubTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageProviderStubTest.java
@@ -28,7 +28,7 @@ class MessageProviderStubTest {
 
     @Test
     void queryMessagesShouldFailExplicitlyWhenRealProviderIsMissing() {
-        assertThatThrownBy(() -> provider.queryMessages("orders", null, null, 
null, null, null))
+        assertThatThrownBy(() -> provider.queryMessages("instance-a", 
"orders", null, null, null, null, null))
                 .isInstanceOf(BusinessException.class)
                 .hasMessage("Message query provider is not configured")
                 .extracting("code")
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 d4045388..2c8de5a3 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
@@ -24,7 +24,7 @@ class MessageServiceTest {
         MessageProvider provider = mock(MessageProvider.class);
         MessageService service = new MessageService(provider);
 
-        assertThatThrownBy(() -> service.queryMessages("TopicA", null, null, 
null, 200L, 100L))
+        assertThatThrownBy(() -> service.queryMessages("instance-a", "TopicA", 
null, null, null, 200L, 100L))
                 .isInstanceOf(BusinessException.class)
                 .hasMessage("startTime must not be after endTime");
 
@@ -36,7 +36,7 @@ class MessageServiceTest {
         MessageProvider provider = mock(MessageProvider.class);
         MessageService service = new MessageService(provider);
 
-        assertThatThrownBy(() -> service.queryMessages("TopicA", null, null, 
null, 0L,
+        assertThatThrownBy(() -> service.queryMessages("instance-a", "TopicA", 
null, null, null, 0L,
                 8L * 24 * 60 * 60 * 1000))
                 .isInstanceOf(BusinessException.class)
                 .hasMessage("topic query time range must not exceed 7 days");
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
index ac84429a..4c8e90c4 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
@@ -20,6 +20,8 @@ import org.apache.rocketmq.client.QueryResult;
 import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 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;
@@ -42,6 +44,7 @@ import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyInt;
 import static org.mockito.ArgumentMatchers.anyLong;
 import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.doNothing;
 import static org.mockito.Mockito.lenient;
 import static org.mockito.Mockito.mockConstruction;
@@ -58,6 +61,9 @@ class RocketMQMessageProviderTest {
     @Mock
     private DefaultMQAdminExt adminExt;
 
+    @Mock
+    private RuntimeAdminClientResolver runtimeAdminClientResolver;
+
     @Mock
     private QueryHistoryService queryHistoryService;
 
@@ -66,7 +72,13 @@ class RocketMQMessageProviderTest {
     @BeforeEach
     void setUp() {
         lenient().when(adminExtProvider.getIfAvailable()).thenReturn(adminExt);
-        provider = new RocketMQMessageProvider(adminExtProvider, 
queryHistoryService, new RocketMQProperties());
+        
lenient().when(runtimeAdminClientResolver.resolveEndpoint("instance-a")).thenReturn("namesrv-a:9876");
+        lenient().when(runtimeAdminClientResolver.execute(anyString(), 
any())).thenAnswer(invocation -> {
+            MqAdminExtFactory.AdminAction<Object> action = 
invocation.getArgument(1);
+            return action == null ? null : action.apply(adminExt);
+        });
+        provider = new RocketMQMessageProvider(adminExtProvider, 
runtimeAdminClientResolver, queryHistoryService,
+                new RocketMQProperties());
     }
 
     @Test
@@ -77,16 +89,19 @@ class RocketMQMessageProviderTest {
                          
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(null);
                          doNothing().when(consumer).shutdown();
                      })) {
-            List<MessageRecordVO> messages = provider.queryMessages("TopicA", 
null, null, null, 100L, 200L);
+            List<MessageRecordVO> messages = 
provider.queryMessages("instance-a", "TopicA", null, null, null, 100L, 200L);
 
             assertThat(messages).isEmpty();
             assertThat(mockedConsumers.constructed()).hasSize(1);
             DefaultMQPullConsumer consumer = 
mockedConsumers.constructed().get(0);
+            verify(consumer).setNamesrvAddr("namesrv-a:9876");
             verify(consumer).start();
             verify(consumer).fetchSubscribeMessageQueues("TopicA");
             verify(consumer, never()).pull(any(MessageQueue.class), 
anyString(), anyLong(), anyInt());
             verify(consumer).shutdown();
         }
+        verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
+        verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
         verify(queryHistoryService).recordMessageQuery(null, "TOPIC", 
"TopicA", null, null, null,
                 100L, 200L, 0);
     }
diff --git a/web/src/api/message.test.ts b/web/src/api/message.test.ts
index 75b8e697..32107a02 100644
--- a/web/src/api/message.test.ts
+++ b/web/src/api/message.test.ts
@@ -35,6 +35,7 @@ describe('message API', () => {
 
   it('sends the backend-supported query fields with epoch timestamps', async 
() => {
     const params = {
+      instanceId: 'instance-a',
       topic: 'orders',
       tag: 'created',
       key: 'order-1',
diff --git a/web/src/api/message.ts b/web/src/api/message.ts
index c21509b4..57c17138 100644
--- a/web/src/api/message.ts
+++ b/web/src/api/message.ts
@@ -35,6 +35,7 @@ export interface TraceRecord {
 }
 
 export interface MessageQuery {
+  instanceId?: string;
   topic?: string;
   tag?: string;
   key?: string;
diff --git a/web/src/pages/instance/__tests__/MessagePage.test.tsx 
b/web/src/pages/instance/__tests__/MessagePage.test.tsx
index 290a8116..1ae39b55 100644
--- a/web/src/pages/instance/__tests__/MessagePage.test.tsx
+++ b/web/src/pages/instance/__tests__/MessagePage.test.tsx
@@ -100,7 +100,10 @@ describe('Message page query history', () => {
     await user.click(screen.getByRole('button', { name: /^search查询$/ }));
 
     await waitFor(() => {
-      expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({ msgId: 
'MID-001' });
+      expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({
+        msgId: 'MID-001',
+        instanceId: '',
+      });
       expect(screen.getByRole('button', { name: /最近查询/ })).toBeEnabled();
     });
 
@@ -113,7 +116,10 @@ describe('Message page query history', () => {
 
     expect(screen.getByPlaceholderText('输入 Message 
ID')).toHaveValue('MID-001');
     await waitFor(() => {
-      expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({ msgId: 
'MID-001' });
+      expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({
+        msgId: 'MID-001',
+        instanceId: '',
+      });
     });
 
     await user.click(screen.getByRole('button', { name: /最近查询/ }));
@@ -132,7 +138,10 @@ describe('Message page query history', () => {
     await user.click(screen.getByRole('button', { name: /^search查询$/ }));
 
     await waitFor(() => {
-      expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({ msgId: 
'MID-FAILED' });
+      expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({
+        msgId: 'MID-FAILED',
+        instanceId: '',
+      });
     });
     expect(screen.getByRole('button', { name: /最近查询/ })).toBeDisabled();
     expect(localStorage).toHaveLength(0);
@@ -207,13 +216,19 @@ describe('Message page query history', () => {
     await user.click(screen.getByRole('button', { name: /最近查询/ }));
     await user.click(await screen.findByText('Topic: order-create'));
     await waitFor(() => {
-      
expect(messageServiceMocks.queryMessages).toHaveBeenLastCalledWith(topicParams);
+      expect(messageServiceMocks.queryMessages).toHaveBeenLastCalledWith({
+        ...topicParams,
+        instanceId: '',
+      });
     });
 
     await user.click(screen.getByRole('button', { name: /最近查询/ }));
     await user.click(await screen.findByText('Key: ORDER-001 · Topic: 
payment-callback'));
     await waitFor(() => {
-      
expect(messageServiceMocks.queryMessages).toHaveBeenLastCalledWith(keyParams);
+      expect(messageServiceMocks.queryMessages).toHaveBeenLastCalledWith({
+        ...keyParams,
+        instanceId: '',
+      });
       expect(screen.getByPlaceholderText('输入 Message 
Key')).toHaveValue('ORDER-001');
     });
   });
diff --git a/web/src/pages/instance/message.tsx 
b/web/src/pages/instance/message.tsx
index 5162392f..e7a59a6b 100644
--- a/web/src/pages/instance/message.tsx
+++ b/web/src/pages/instance/message.tsx
@@ -291,7 +291,7 @@ const MessagePage = () => {
     setQueryLoading(true);
     setQueryError(null);
     try {
-      const result = await queryMessages(params);
+      const result = await queryMessages({ ...params, instanceId: 
selectedInstanceId });
       if (queryGenerationRef.current !== requestGeneration) return;
       setMessages(result);
       setQueryError(null);

Reply via email to