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 4ad376b18 fix(aliyun): preserve trace consumer status (#2239)
4ad376b18 is described below
commit 4ad376b18aa6c04393418eeb1e9681d3bdd13137
Author: shown <[email protected]>
AuthorDate: Wed Aug 19 11:21:48 2026 +0800
fix(aliyun): preserve trace consumer status (#2239)
* fix(aliyun): preserve trace consumer status
Signed-off-by: yuluo-yx <[email protected]>
* fix(aliyun): tolerate null collection rows
Signed-off-by: yuluo-yx <[email protected]>
---------
Signed-off-by: yuluo-yx <[email protected]>
---
.../studio/provider/alibaba/AliyunConverters.java | 47 ++++++++++++++---
.../provider/alibaba/AliyunInstanceProvider.java | 15 ++++++
.../alibaba/AliyunInstanceProviderTest.java | 60 +++++++++++++++++++---
3 files changed, 107 insertions(+), 15 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
index 82f6628d4..cfe1f2194 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConverters.java
@@ -28,11 +28,13 @@ import
com.aliyun.sdk.service.rocketmq20220801.models.ListRegionsResponseBody;
import
com.aliyun.sdk.service.rocketmq20220801.models.ListTopicSubscriptionsResponseBody;
import com.aliyun.sdk.service.rocketmq20220801.models.ListTopicsResponseBody;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
+import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
+import org.apache.rocketmq.studio.instance.message.ConsumerStatusVO;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
@@ -233,6 +235,7 @@ final class AliyunConverters {
static TraceRecordVO toTraceRecord(GetTraceResponseBody.Data data) {
List<TraceNodeVO> nodes = new ArrayList<>();
+ List<ConsumerStatusVO> consumerStatuses = new ArrayList<>();
if (data.getProducerInfo() != null &&
data.getProducerInfo().getRecords() != null) {
for (GetTraceResponseBody.ProducerInfoRecords record :
data.getProducerInfo().getRecords()) {
if (record == null) {
@@ -264,30 +267,60 @@ final class AliyunConverters {
continue;
}
if (consumerInfo.getRecords() == null ||
consumerInfo.getRecords().isEmpty()) {
+ String status = consumerInfo.getConsumeStatus();
nodes.add(TraceNodeVO.builder()
.title("Consumer " +
consumerInfo.getConsumerGroupId())
- .status(consumerInfo.getConsumeStatus())
+ .status(status)
.build());
+ consumerStatuses.add(consumerStatus(
+ consumerInfo.getConsumerGroupId(), status, 0L));
continue;
}
for (GetTraceResponseBody.Records record :
consumerInfo.getRecords()) {
if (record == null) {
continue;
}
- String operateTime = null;
- if (record.getOperations() != null &&
!record.getOperations().isEmpty()) {
- operateTime =
record.getOperations().get(0).getOperateTime();
- }
+ String operateTime =
firstOperateTime(record.getOperations());
+ long consumeTime = parseTimeMillis(operateTime);
nodes.add(TraceNodeVO.builder()
.title("Consumer " +
consumerInfo.getConsumerGroupId())
- .timestamp(parseTimeMillis(operateTime))
+ .timestamp(consumeTime)
.status(record.getConsumeStatus())
.description(joinParts(", ",
record.getClientHost(), record.getUserName()))
.build());
+ consumerStatuses.add(consumerStatus(
+ consumerInfo.getConsumerGroupId(),
record.getConsumeStatus(), consumeTime));
}
}
}
- return TraceRecordVO.builder().nodes(nodes).build();
+ return TraceRecordVO.builder()
+ .nodes(nodes)
+ .consumerStatus(consumerStatuses)
+ .build();
+ }
+
+ private static String
firstOperateTime(List<GetTraceResponseBody.RecordsOperations> operations) {
+ if (operations == null) {
+ return null;
+ }
+ return operations.stream()
+ .filter(operation -> operation != null &&
operation.getOperateTime() != null)
+ .map(GetTraceResponseBody.RecordsOperations::getOperateTime)
+ .findFirst()
+ .orElse(null);
+ }
+
+ private static ConsumerStatusVO consumerStatus(String group, String
rawStatus, long consumeTime) {
+ String normalized = rawStatus == null ? "" :
rawStatus.toUpperCase(Locale.ROOT);
+ DeliveryStatus status = normalized.contains("SUCCESS") ||
normalized.contains("OK")
+ ? DeliveryStatus.success
+ : normalized.contains("FAIL") ? DeliveryStatus.failed :
DeliveryStatus.pending;
+ return ConsumerStatusVO.builder()
+ .group(group)
+ .deliveryStatus(status)
+ .consumeTime(consumeTime)
+ .retryCount(0)
+ .build();
}
static java.time.LocalDateTime parseDateTime(String value) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
index c2a95721d..390487ef1 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
@@ -170,6 +170,9 @@ public class AliyunInstanceProvider implements
InstanceProvider {
}
List<TopicVO> topics = new ArrayList<>();
for (ListTopicsResponseBody.List item : all) {
+ if (item == null) {
+ continue;
+ }
TopicVO vo = AliyunConverters.toTopicVO(item, instanceId);
if (matchesType(type, vo)) {
topics.add(vo);
@@ -242,6 +245,9 @@ public class AliyunInstanceProvider implements
InstanceProvider {
return consumers;
}
for (ListTopicSubscriptionsResponseBody.Data item : data) {
+ if (item == null) {
+ continue;
+ }
consumers.add(AliyunConverters.toTopicConsumerVO(item));
}
return consumers;
@@ -275,6 +281,9 @@ public class AliyunInstanceProvider implements
InstanceProvider {
}
List<ConsumerGroupVO> groups = new ArrayList<>();
for (ListConsumerGroupsResponseBody.List item : all) {
+ if (item == null) {
+ continue;
+ }
groups.add(AliyunConverters.toConsumerGroupVO(item, instanceId));
}
return groups;
@@ -371,6 +380,9 @@ public class AliyunInstanceProvider implements
InstanceProvider {
return subscriptions;
}
for (ListConsumerGroupSubscriptionsResponseBody.Data item : data) {
+ if (item == null) {
+ continue;
+ }
subscriptions.add(AliyunConverters.toSubscriptionEntry(item));
}
return subscriptions;
@@ -430,6 +442,9 @@ public class AliyunInstanceProvider implements
InstanceProvider {
break;
}
for (ListMessagesResponseBody.List item : list) {
+ if (item == null) {
+ continue;
+ }
MessageRecordVO vo = AliyunConverters.toMessageRecord(item);
if (!StringUtils.hasText(tag) || tag.equals(vo.getTag())) {
records.add(vo);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
index 2d46f3e91..357375c99 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
@@ -112,6 +112,7 @@ class AliyunInstanceProviderTest {
stubCallThrough();
ListTopicsResponse response = topicsResponse(
topicRow("topic-normal", "NORMAL"),
+ null,
topicRow("topic-fifo", "FIFO"),
topicRow("topic-mystery", "MYSTERY"));
when(asyncClient.listTopics(any(ListTopicsRequest.class)))
@@ -142,12 +143,13 @@ class AliyunInstanceProviderTest {
.statusCode(200)
.body(ListConsumerGroupsResponseBody.builder()
.data(ListConsumerGroupsResponseBody.Data.builder()
-
.list(List.of(ListConsumerGroupsResponseBody.List.builder()
- .consumerGroupId("GID_test")
- .messageModel("Clustering")
- .status("RUNNING")
- .remark("test group")
- .build()))
+ .list(java.util.Arrays.asList(null,
+
ListConsumerGroupsResponseBody.List.builder()
+ .consumerGroupId("GID_test")
+ .messageModel("Clustering")
+ .status("RUNNING")
+ .remark("test group")
+ .build()))
.pageNumber(1L)
.pageSize(100L)
.totalCount(1L)
@@ -209,7 +211,8 @@ class AliyunInstanceProviderTest {
.statusCode(200)
.body(ListMessagesResponseBody.builder()
.data(ListMessagesResponseBody.Data.builder()
- .list(List.of(
+ .list(java.util.Arrays.asList(
+ null,
ListMessagesResponseBody.List.builder()
.messageId("msg-1")
.topicName("topic-a")
@@ -369,6 +372,47 @@ class AliyunInstanceProviderTest {
TraceNodeVO consumer = trace.getNodes().get(2);
assertThat(consumer.getTitle()).isEqualTo("Consumer GID_test");
assertThat(consumer.getStatus()).isEqualTo("CONSUME_OK");
+ assertThat(trace.getConsumerStatus()).singleElement().satisfies(status
-> {
+ assertThat(status.getGroup()).isEqualTo("GID_test");
+ assertThat(status.getDeliveryStatus().name()).isEqualTo("success");
+ assertThat(status.getConsumeTime())
+ .isEqualTo(AliyunConverters.parseTimeMillis("2023-03-22
12:17:10"));
+ });
+ }
+
+ @Test
+ void getMessageTraceShouldSkipNullConsumerOperations() {
+ stubInstance();
+ stubCallThrough();
+ GetTraceResponse response = GetTraceResponse.create().toBuilder()
+ .statusCode(200)
+ .body(GetTraceResponseBody.builder()
+ .data(GetTraceResponseBody.Data.builder()
+
.consumerInfos(List.of(GetTraceResponseBody.ConsumerInfos.builder()
+ .consumerGroupId("GID_test")
+
.records(List.of(GetTraceResponseBody.Records.builder()
+
.consumeStatus("CONSUME_FAILED")
+
.operations(java.util.Arrays.asList(null,
+
GetTraceResponseBody.RecordsOperations.builder()
+
.operateTime("2023-03-22 12:17:10")
+ .build()))
+ .build()))
+ .build()))
+ .build())
+ .build())
+ .build();
+
when(asyncClient.getTrace(any())).thenReturn(CompletableFuture.completedFuture(response));
+
+ TraceRecordVO trace = provider.getMessageTrace(STUDIO_INSTANCE_ID,
"msg-1", "orders");
+
+ assertThat(trace.getNodes()).singleElement()
+ .extracting(TraceNodeVO::getTimestamp)
+ .isEqualTo(AliyunConverters.parseTimeMillis("2023-03-22
12:17:10"));
+ assertThat(trace.getConsumerStatus()).singleElement().satisfies(status
-> {
+ assertThat(status.getDeliveryStatus().name()).isEqualTo("failed");
+ assertThat(status.getConsumeTime())
+ .isEqualTo(AliyunConverters.parseTimeMillis("2023-03-22
12:17:10"));
+ });
}
@Test
@@ -466,7 +510,7 @@ class AliyunInstanceProviderTest {
.statusCode(200)
.body(ListTopicsResponseBody.builder()
.data(ListTopicsResponseBody.Data.builder()
- .list(List.of(rows))
+ .list(java.util.Arrays.asList(rows))
.pageNumber(1L)
.pageSize(100L)
.totalCount((long) rows.length)