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