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 0920783d [ISSUE #1608] Count Aliyun resources from API totals (#1610)
0920783d is described below
commit 0920783dfc922b6ff10dd9b9c0af843ca3228bd4
Author: youngkermit8-coder <[email protected]>
AuthorDate: Tue Aug 11 20:30:32 2026 +0800
[ISSUE #1608] Count Aliyun resources from API totals (#1610)
Signed-off-by: youngkermit8-coder <[email protected]>
---
.../provider/alibaba/AliyunInstanceProvider.java | 35 ++++++-
.../alibaba/AliyunInstanceProviderTest.java | 109 ++++++++++++++++++---
2 files changed, 127 insertions(+), 17 deletions(-)
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 927e3bd3..447249c0 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
@@ -77,6 +77,7 @@ public class AliyunInstanceProvider implements
InstanceProvider {
private static final String FIXED_RETRY_POLICY = "FixedRetryPolicy";
private static final int DEFAULT_MAX_RETRY_TIMES = 16;
private static final int DEFAULT_FIXED_RETRY_INTERVAL_SECONDS = 10;
+ private static final int COUNT_PAGE_SIZE = 1;
private static final String RESET_TYPE_SPECIFIED_TIME = "SPECIFIED_TIME";
private static final String RESET_TYPE_LATEST_OFFSET = "LATEST_OFFSET";
@@ -90,12 +91,42 @@ public class AliyunInstanceProvider implements
InstanceProvider {
@Override
public int countTopics(String instanceId) {
- return listTopics(instanceId, null, null).size();
+ Context ctx = resolve(instanceId);
+ ListTopicsRequest request = ListTopicsRequest.builder()
+ .instanceId(ctx.cloudInstanceId())
+ .pageNumber(1)
+ .pageSize(COUNT_PAGE_SIZE)
+ .build();
+ ListTopicsResponse response = clientFactory.call(ctx.credentialId(),
ctx.regionId(),
+ client -> client.listTopics(request));
+ ListTopicsResponseBody body = response == null ? null :
response.getBody();
+ ListTopicsResponseBody.Data data = body == null ? null :
body.getData();
+ Long totalCount = data == null ? null : data.getTotalCount();
+ return totalCount == null || totalCount < 0
+ ? listTopics(instanceId, null, null).size()
+ : boundedCount(totalCount);
}
@Override
public int countGroups(String instanceId) {
- return listConsumerGroups(instanceId, null).size();
+ Context ctx = resolve(instanceId);
+ ListConsumerGroupsRequest request = ListConsumerGroupsRequest.builder()
+ .instanceId(ctx.cloudInstanceId())
+ .pageNumber(1)
+ .pageSize(COUNT_PAGE_SIZE)
+ .build();
+ ListConsumerGroupsResponse response =
clientFactory.call(ctx.credentialId(), ctx.regionId(),
+ client -> client.listConsumerGroups(request));
+ ListConsumerGroupsResponseBody body = response == null ? null :
response.getBody();
+ ListConsumerGroupsResponseBody.Data data = body == null ? null :
body.getData();
+ Long totalCount = data == null ? null : data.getTotalCount();
+ return totalCount == null || totalCount < 0
+ ? listConsumerGroups(instanceId, null).size()
+ : boundedCount(totalCount);
+ }
+
+ private int boundedCount(long totalCount) {
+ return totalCount > Integer.MAX_VALUE ? Integer.MAX_VALUE : (int)
totalCount;
}
@Override
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 645d9537..0d8ef275 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
@@ -25,6 +25,7 @@ import
com.aliyun.sdk.service.rocketmq20220801.models.GetConsumerGroupLagRespons
import
com.aliyun.sdk.service.rocketmq20220801.models.GetConsumerGroupLagResponseBody;
import com.aliyun.sdk.service.rocketmq20220801.models.GetTraceResponse;
import com.aliyun.sdk.service.rocketmq20220801.models.GetTraceResponseBody;
+import
com.aliyun.sdk.service.rocketmq20220801.models.ListConsumerGroupsRequest;
import
com.aliyun.sdk.service.rocketmq20220801.models.ListConsumerGroupsResponse;
import
com.aliyun.sdk.service.rocketmq20220801.models.ListConsumerGroupsResponseBody;
import com.aliyun.sdk.service.rocketmq20220801.models.ListMessagesResponse;
@@ -66,6 +67,7 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -427,42 +429,119 @@ class AliyunInstanceProviderTest {
}
@Test
- void countTopicsShouldReuseListTopicsSizeTest() {
+ void countTopicsShouldUseTotalCountWithoutFetchingEveryTopicTest() {
stubInstance();
stubCallThrough();
+ ListTopicsResponse response = ListTopicsResponse.create().toBuilder()
+ .statusCode(200)
+ .body(ListTopicsResponseBody.builder()
+ .data(ListTopicsResponseBody.Data.builder()
+ .list(List.of(topicRow("topic-a", "NORMAL")))
+ .pageNumber(1L)
+ .pageSize(1L)
+ .totalCount(321L)
+ .build())
+ .build())
+ .build();
when(asyncClient.listTopics(any(ListTopicsRequest.class)))
- .thenReturn(CompletableFuture.completedFuture(topicsResponse(
- topicRow("topic-a", "NORMAL"),
- topicRow("topic-b", "FIFO"))));
+ .thenReturn(CompletableFuture.completedFuture(response));
- assertThat(provider.countTopics(STUDIO_INSTANCE_ID)).isEqualTo(2);
+ assertThat(provider.countTopics(STUDIO_INSTANCE_ID)).isEqualTo(321);
+ ArgumentCaptor<ListTopicsRequest> captor =
ArgumentCaptor.forClass(ListTopicsRequest.class);
+ verify(asyncClient).listTopics(captor.capture());
+ assertThat(captor.getValue().getPageNumber()).isEqualTo(1L);
+ assertThat(captor.getValue().getPageSize()).isEqualTo(1L);
}
@Test
- void countGroupsShouldReuseListConsumerGroupsSizeTest() {
+ void countGroupsShouldUseTotalCountWithoutFetchingEveryGroupTest() {
stubInstance();
stubCallThrough();
ListConsumerGroupsResponse response =
ListConsumerGroupsResponse.create().toBuilder()
.statusCode(200)
.body(ListConsumerGroupsResponseBody.builder()
.data(ListConsumerGroupsResponseBody.Data.builder()
- .list(List.of(
-
ListConsumerGroupsResponseBody.List.builder()
- .consumerGroupId("GID_one")
- .build(),
-
ListConsumerGroupsResponseBody.List.builder()
- .consumerGroupId("GID_two")
- .build()))
+
.list(List.of(ListConsumerGroupsResponseBody.List.builder()
+ .consumerGroupId("GID_one")
+ .build()))
.pageNumber(1L)
- .pageSize(100L)
- .totalCount(2L)
+ .pageSize(1L)
+ .totalCount(654L)
.build())
.build())
.build();
when(asyncClient.listConsumerGroups(any()))
.thenReturn(CompletableFuture.completedFuture(response));
+ assertThat(provider.countGroups(STUDIO_INSTANCE_ID)).isEqualTo(654);
+ ArgumentCaptor<ListConsumerGroupsRequest> captor =
+ ArgumentCaptor.forClass(ListConsumerGroupsRequest.class);
+ verify(asyncClient).listConsumerGroups(captor.capture());
+ assertThat(captor.getValue().getPageNumber()).isEqualTo(1L);
+ assertThat(captor.getValue().getPageSize()).isEqualTo(1L);
+ }
+
+ @Test
+ void countTopicsShouldFallBackToFullListingWhenTotalCountIsMissingTest() {
+ stubInstance();
+ stubCallThrough();
+ ListTopicsResponse missingTotal =
ListTopicsResponse.create().toBuilder()
+ .statusCode(200)
+ .body(ListTopicsResponseBody.builder()
+ .data(ListTopicsResponseBody.Data.builder()
+ .list(List.of(topicRow("topic-a", "NORMAL")))
+ .pageNumber(1L)
+ .pageSize(1L)
+ .build())
+ .build())
+ .build();
+ when(asyncClient.listTopics(any(ListTopicsRequest.class)))
+ .thenReturn(CompletableFuture.completedFuture(missingTotal))
+ .thenReturn(CompletableFuture.completedFuture(topicsResponse(
+ topicRow("topic-a", "NORMAL"),
+ topicRow("topic-b", "FIFO"))));
+
+ assertThat(provider.countTopics(STUDIO_INSTANCE_ID)).isEqualTo(2);
+ ArgumentCaptor<ListTopicsRequest> captor =
ArgumentCaptor.forClass(ListTopicsRequest.class);
+ verify(asyncClient, times(2)).listTopics(captor.capture());
+
assertThat(captor.getAllValues()).extracting(ListTopicsRequest::getPageSize)
+ .containsExactly(1, AliyunConverters.PAGE_SIZE);
+ }
+
+ @Test
+ void countGroupsShouldFallBackToFullListingWhenTotalCountIsMissingTest() {
+ stubInstance();
+ stubCallThrough();
+ ListConsumerGroupsResponse missingTotal = groupsResponse(null,
"GID_one");
+ ListConsumerGroupsResponse completeListing = groupsResponse(2L,
"GID_one", "GID_two");
+ when(asyncClient.listConsumerGroups(any()))
+ .thenReturn(CompletableFuture.completedFuture(missingTotal))
+
.thenReturn(CompletableFuture.completedFuture(completeListing));
+
assertThat(provider.countGroups(STUDIO_INSTANCE_ID)).isEqualTo(2);
+ ArgumentCaptor<ListConsumerGroupsRequest> captor =
+ ArgumentCaptor.forClass(ListConsumerGroupsRequest.class);
+ verify(asyncClient, times(2)).listConsumerGroups(captor.capture());
+
assertThat(captor.getAllValues()).extracting(ListConsumerGroupsRequest::getPageSize)
+ .containsExactly(1, AliyunConverters.PAGE_SIZE);
+ }
+
+ private static ListConsumerGroupsResponse groupsResponse(Long totalCount,
String... groupIds) {
+ return ListConsumerGroupsResponse.create().toBuilder()
+ .statusCode(200)
+ .body(ListConsumerGroupsResponseBody.builder()
+ .data(ListConsumerGroupsResponseBody.Data.builder()
+ .list(java.util.Arrays.stream(groupIds)
+ .map(groupId ->
ListConsumerGroupsResponseBody.List.builder()
+ .consumerGroupId(groupId)
+ .build())
+ .toList())
+ .pageNumber(1L)
+ .pageSize((long) AliyunConverters.PAGE_SIZE)
+ .totalCount(totalCount)
+ .build())
+ .build())
+ .build();
}
@Test