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 d09609e9 fix(message): use decoded physical offset for message lookup 
(#1911)
d09609e9 is described below

commit d09609e9aa35fa1461c4ffad9edfa7d7ff85a5f1
Author: majialong <[email protected]>
AuthorDate: Fri Aug 14 11:02:26 2026 +0800

    fix(message): use decoded physical offset for message lookup (#1911)
    
    Decode the broker address and physical offset embedded in the offset
    msgId and query the broker directly with the topic hint, fixing trace
    queries that could not locate the original message.
---
 .../provider/apache/RocketMQMessageProvider.java   | 15 +++---
 .../apache/RocketMQMessageProviderTest.java        | 59 ++++++++++++++++++++++
 2 files changed, 66 insertions(+), 8 deletions(-)

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 1432e823..3debc208 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
@@ -79,6 +79,7 @@ public class RocketMQMessageProvider implements 
MessageProvider {
     private static final int MAX_BINARY_BODY_DISPLAY_BYTES = 48 * 1024;
     private static final int MAX_PROPERTIES = 64;
     private static final int MAX_PROPERTY_VALUE_CHARS = 1024;
+    private static final long VIEW_MESSAGE_TIMEOUT_MILLIS = 3000L;
     private static final long ONE_HOUR_MILLIS = 3600_000L;
     private static final long ONE_DAY_MILLIS = 24 * ONE_HOUR_MILLIS;
     private static final int MAX_CONSECUTIVE_OFFSET_ILLEGAL = 3;
@@ -138,7 +139,7 @@ public class RocketMQMessageProvider implements 
MessageProvider {
             }
         }
         if (messageExt == null) {
-            messageExt = viewMessageByOffsetId(adminExt, msgId);
+            messageExt = viewMessageByOffsetId(adminExt, topic, msgId);
         }
         if (messageExt == null) {
             return Collections.emptyList();
@@ -147,10 +148,10 @@ public class RocketMQMessageProvider implements 
MessageProvider {
     }
 
     /**
-     * Locate a message purely by its offset msgId by decoding the broker 
address embedded in the
-     * id and querying that broker directly. Used when no topic hint is 
available.
+     * Locate a message by decoding the broker address and physical offset 
embedded in its offset
+     * msgId, then querying that broker directly.
      */
-    private MessageExt viewMessageByOffsetId(DefaultMQAdminExt adminExt, 
String msgId) {
+    private MessageExt viewMessageByOffsetId(DefaultMQAdminExt adminExt, 
String topic, String msgId) {
         try {
             MessageId messageId = MessageDecoder.decodeMessageId(msgId);
             SocketAddress address = messageId.getAddress();
@@ -162,7 +163,7 @@ public class RocketMQMessageProvider implements 
MessageProvider {
             return adminExt.getDefaultMQAdminExtImpl()
                     .getMqClientInstance()
                     .getMQClientAPIImpl()
-                    .viewMessage(brokerAddr, msgId, 3000L, 3000L);
+                    .viewMessage(brokerAddr, topic, messageId.getOffset(), 
VIEW_MESSAGE_TIMEOUT_MILLIS);
         } catch (Exception e) {
             log.warn("viewMessage by decoded offset id failed for msgId={}: 
{}", msgId, e.getMessage());
             return null;
@@ -348,9 +349,7 @@ public class RocketMQMessageProvider implements 
MessageProvider {
             }
         }
         try {
-            // No topic hint is available in the trace flow, so locate the 
message purely
-            // by its offset msgId.
-            MessageExt messageExt = viewMessageByOffsetId(adminExt, msgId);
+            MessageExt messageExt = viewMessageByOffsetId(adminExt, topic, 
msgId);
             if (messageExt != null) {
                 return messageExt.getStoreTimestamp();
             }
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 033241dc..2a28c8f7 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
@@ -20,6 +20,8 @@ import org.apache.rocketmq.client.QueryResult;
 import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
 import org.apache.rocketmq.client.consumer.PullResult;
 import org.apache.rocketmq.client.consumer.PullStatus;
+import org.apache.rocketmq.client.impl.MQClientAPIImpl;
+import org.apache.rocketmq.client.impl.factory.MQClientInstance;
 import org.apache.rocketmq.client.trace.TraceConstants;
 import org.apache.rocketmq.common.message.MessageDecoder;
 import org.apache.rocketmq.common.message.MessageExt;
@@ -34,6 +36,7 @@ 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;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.Timeout;
@@ -61,6 +64,7 @@ import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.doNothing;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.lenient;
+import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.times;
@@ -158,6 +162,25 @@ class RocketMQMessageProviderTest {
                 .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(502));
     }
 
+    @Test
+    void queryByMsgIdUsesDecodedPhysicalOffsetForFallback() throws Exception {
+        String msgId = "AC1E0A6400002A9F0000000001A3F2B1";
+        MQClientAPIImpl clientApi = mockOffsetLookupClient();
+        MessageExt message = new MessageExt();
+        message.setMsgId(msgId);
+        message.setTopic("TopicA");
+        when(adminExt.viewMessage("TopicA", msgId))
+                .thenThrow(new IllegalStateException("primary lookup failed"));
+        when(clientApi.viewMessage("172.30.10.100:10911", "TopicA", 27521713L, 
3000L))
+                .thenReturn(message);
+
+        List<MessageRecordVO> result = provider.queryMessages(
+                "instance-a", "TopicA", msgId, null, null, 100L, 200L);
+
+        
assertThat(result).singleElement().extracting(MessageRecordVO::getMsgId).isEqualTo(msgId);
+        verify(clientApi).viewMessage("172.30.10.100:10911", "TopicA", 
27521713L, 3000L);
+    }
+
     @Test
     void queryByTopicSurfacesPullConsumerFailure() throws Exception {
         try (MockedConstruction<DefaultMQPullConsumer> ignored =
@@ -483,6 +506,42 @@ class RocketMQMessageProviderTest {
         verify(queryHistoryService, never()).recordTraceQuery(anyString(), 
anyString(), any(), anyInt(), anyInt());
     }
 
+    @Test
+    void getMessageTraceUsesDecodedOffsetToResolveQueryWindow() throws 
Exception {
+        String msgId = "AC1E0A6400002A9F0000000001A3F2B1";
+        long storeTimestamp = System.currentTimeMillis() - 2 * 60 * 60 * 1000L;
+        MQClientAPIImpl clientApi = mockOffsetLookupClient();
+        MessageExt message = new MessageExt();
+        message.setMsgId(msgId);
+        message.setTopic("TopicA");
+        message.setStoreTimestamp(storeTimestamp);
+        when(clientApi.viewMessage("172.30.10.100:10911", "TopicA", 27521713L, 
3000L))
+                .thenReturn(message);
+        when(adminExt.queryMessage(
+                "RMQ_SYS_TRACE_TOPIC", msgId, 64, storeTimestamp - 5 * 60_000L,
+                storeTimestamp + 24 * 60 * 60 * 1000L))
+                .thenReturn(new QueryResult(0L, List.of()));
+
+        TraceRecordVO result = provider.getMessageTrace("instance-a", msgId, 
"TopicA");
+
+        assertThat(result.getNodes()).isEmpty();
+        assertThat(result.getConsumerStatus()).isEmpty();
+        verify(clientApi).viewMessage("172.30.10.100:10911", "TopicA", 
27521713L, 3000L);
+        verify(adminExt).queryMessage(
+                "RMQ_SYS_TRACE_TOPIC", msgId, 64, storeTimestamp - 5 * 60_000L,
+                storeTimestamp + 24 * 60 * 60 * 1000L);
+    }
+
+    private MQClientAPIImpl mockOffsetLookupClient() {
+        DefaultMQAdminExtImpl adminExtImpl = mock(DefaultMQAdminExtImpl.class);
+        MQClientInstance clientInstance = mock(MQClientInstance.class);
+        MQClientAPIImpl clientApi = mock(MQClientAPIImpl.class);
+        when(adminExt.getDefaultMQAdminExtImpl()).thenReturn(adminExtImpl);
+        when(adminExtImpl.getMqClientInstance()).thenReturn(clientInstance);
+        when(clientInstance.getMQClientAPIImpl()).thenReturn(clientApi);
+        return clientApi;
+    }
+
     private static String traceContext(String... fields) {
         return String.join(String.valueOf(TraceConstants.CONTENT_SPLITOR), 
fields);
     }

Reply via email to