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 b295e8f3e feat(message): support recall events in message trace (#1930)
b295e8f3e is described below

commit b295e8f3ea6d227ec126a69c2a77e62f56cd37f9
Author: majialong <[email protected]>
AuthorDate: Tue Aug 18 16:36:50 2026 +0800

    feat(message): support recall events in message trace (#1930)
---
 .../provider/apache/RocketMQMessageProvider.java   | 15 +++++++++++++++
 .../apache/RocketMQMessageProviderTest.java        | 22 ++++++++++++++++++++++
 2 files changed, 37 insertions(+)

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 3debc208f..524559972 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
@@ -397,6 +397,9 @@ public class RocketMQMessageProvider implements 
MessageProvider {
                     case "EndTransaction":
                         nodes.add(buildTransactionNode(fields));
                         break;
+                    case "Recall":
+                        nodes.add(buildRecallNode(fields));
+                        break;
                     default:
                         // SubBefore and unknown types are not surfaced as 
timeline nodes.
                         break;
@@ -457,6 +460,18 @@ public class RocketMQMessageProvider implements 
MessageProvider {
                 .build();
     }
 
+    // Recall layout (RocketMQ 5.5.0 TraceDataEncoder):
+    //               type, time, region, group, topic, msgId, isSuccess
+    private TraceNodeVO buildRecallNode(String[] f) {
+        return TraceNodeVO.builder()
+                .title("recall")
+                .timestamp(parseLong(field(f, 1)))
+                .status(parseBoolean(field(f, 6)) ? "finish" : "failed")
+                .costTime(0L)
+                .description("group=" + field(f, 3) + ", topic=" + field(f, 4))
+                .build();
+    }
+
     private List<ConsumerStatusVO> fallbackConsumerStatus(DefaultMQAdminExt 
adminExt, MessageExt message) {
         List<ConsumerStatusVO> result = new ArrayList<>();
         try {
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 2a28c8f72..c2a1a8248 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
@@ -493,6 +493,28 @@ class RocketMQMessageProviderTest {
         assertThat(address.getPort()).isEqualTo(10911);
     }
 
+    @Test
+    void getMessageTraceParsesRecallState() throws Exception {
+        String body = traceBody(traceContext("Recall", "2500", "cn", 
"producer-group", "TopicA",
+                "msg-recall", "false"));
+        MessageExt traceMessage = new MessageExt();
+        traceMessage.setBody(body.getBytes(StandardCharsets.UTF_8));
+        QueryResult queryResult = new QueryResult(0L, List.of(traceMessage));
+        when(adminExt.queryMessage(anyString(), anyString(), anyInt(), 
anyLong(), anyLong()))
+                .thenReturn(queryResult);
+
+        TraceRecordVO record = provider.getMessageTrace("instance-a", 
"msg-recall", "TopicA");
+
+        assertThat(record.getNodes()).hasSize(1);
+        TraceNodeVO recall = record.getNodes().get(0);
+        assertThat(recall.getTitle()).isEqualTo("recall");
+        assertThat(recall.getTimestamp()).isEqualTo(2500L);
+        assertThat(recall.getStatus()).isEqualTo("failed");
+        assertThat(recall.getCostTime()).isZero();
+        
assertThat(recall.getDescription()).contains("producer-group").contains("TopicA");
+        assertThat(record.getConsumerStatus()).isEmpty();
+    }
+
     @Test
     void getMessageTraceSurfacesAdminFailure() throws Exception {
         when(adminExt.queryMessage(anyString(), anyString(), anyInt(), 
anyLong(), anyLong()))

Reply via email to