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 ea55a7a9f fix(provider): preserve cloud topic pagination totals (#2369)
ea55a7a9f is described below
commit ea55a7a9f228d64d0ca7f38088959e5b36ca7b67
Author: aias00 <[email protected]>
AuthorDate: Wed Aug 19 15:05:03 2026 +0800
fix(provider): preserve cloud topic pagination totals (#2369)
Constraint: /api/topics/page must return exact total/page/size for Aliyun
and Tencent cloud instances
Rejected: raise hardcoded list caps | still truncates totals and
unreachable pages
Directive: Keep cloud topic listing aligned with provider-native pagination
and filter parameters
Confidence: high
Scope-risk: narrow
Tested: export JAVA_HOME=$(/usr/libexec/java_home -v 21) && export
PATH="$JAVA_HOME/bin:$PATH" && mvn -f
/private/tmp/rmqdashboard-fix-2355/server/pom.xml -Djava.awt.headless=true
-Dtest=AliyunInstanceProviderTest,TencentInstanceProviderTest test
Signed-off-by: liuhy <[email protected]>
---
docs/api-spec.md | 4 +-
.../provider/alibaba/AliyunInstanceProvider.java | 99 ++++++++++++------
.../provider/tencent/TencentInstanceProvider.java | 111 +++++++++++++++------
.../alibaba/AliyunInstanceProviderTest.java | 81 +++++++++++++--
.../tencent/TencentInstanceProviderTest.java | 75 +++++++++++++-
5 files changed, 291 insertions(+), 79 deletions(-)
diff --git a/docs/api-spec.md b/docs/api-spec.md
index d54d33fb5..7c7051f9a 100644
--- a/docs/api-spec.md
+++ b/docs/api-spec.md
@@ -766,9 +766,9 @@ GET
/api/topics/page?instanceId={instanceId}&clusterId={clusterId}&type={type}&s
| 字段 | 类型 | 说明 |
|------|------|------|
| `items` | `Topic[]` | 当前页数据,结构同 5.1 |
-| `total` | `number` | 总条数 |
+| `total` | `number` | 总条数;云厂商实例返回 provider 原生分页的匹配总数,不会被单次列表上限截断 |
| `page` | `number` | 当前页码 |
-| `size` | `number` | 每页条数 |
+| `size` | `number` | 每页条数(回显请求中的 `pageSize`) |
### 5.3 创建 Topic
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 390487ef1..f42381beb 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
@@ -45,8 +45,10 @@ import
com.aliyun.sdk.service.rocketmq20220801.models.ResetConsumeOffsetRequest;
import com.aliyun.sdk.service.rocketmq20220801.models.UpdateTopicRequest;
import org.springframework.util.StringUtils;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.common.util.Pagination;
import org.apache.rocketmq.studio.instance.InstanceRepository;
import org.apache.rocketmq.studio.instance.InstanceVO;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
@@ -64,6 +66,7 @@ import lombok.RequiredArgsConstructor;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import java.util.Locale;
import java.util.Set;
/**
@@ -145,42 +148,83 @@ public class AliyunInstanceProvider implements
InstanceProvider {
@Override
public List<TopicVO> listTopics(String instanceId, String type, String
search) {
Context ctx = resolve(instanceId);
- List<ListTopicsResponseBody.List> all = new ArrayList<>();
- for (int page = 1; page <= AliyunConverters.MAX_PAGES; page++) {
- ListTopicsRequest.Builder builder = ListTopicsRequest.builder()
- .instanceId(ctx.cloudInstanceId())
- .pageNumber(page)
- .pageSize(AliyunConverters.PAGE_SIZE);
- if (!!StringUtils.hasText(search)) {
- builder.filter(search);
- }
- ListTopicsRequest request = builder.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();
+ List<TopicVO> topics = new ArrayList<>();
+ for (int page = 1; ; page++) {
+ ListTopicsResponseBody.Data data = fetchTopicPage(ctx, type,
search, page, AliyunConverters.PAGE_SIZE);
List<ListTopicsResponseBody.List> list = data == null ? null :
data.getList();
if (list == null || list.isEmpty()) {
break;
}
- all.addAll(list);
- if (list.size() < AliyunConverters.PAGE_SIZE) {
+ topics.addAll(toTopics(list, instanceId));
+ if (hasFetchedAll(page, AliyunConverters.PAGE_SIZE,
data.getTotalCount())
+ || list.size() < AliyunConverters.PAGE_SIZE) {
break;
}
}
- 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);
+
+ return topics;
+ }
+
+ @Override
+ public PageResult<TopicVO> listTopicsPage(String instanceId, String type,
String search, int page, int pageSize) {
+ Context ctx = resolve(instanceId);
+ ListTopicsResponseBody.Data data = fetchTopicPage(ctx, type, search,
page, pageSize);
+ Long totalCount = data == null ? null : data.getTotalCount();
+ if (totalCount == null || totalCount < 0L) {
+ return paginate(listTopics(instanceId, type, search), page,
pageSize);
+ }
+ return PageResult.of(toTopics(data.getList(), instanceId),
boundedCount(totalCount), page, pageSize);
+ }
+
+ private ListTopicsResponseBody.Data fetchTopicPage(Context ctx, String
type, String search, int page, int pageSize) {
+ ListTopicsRequest request = buildTopicRequest(ctx.cloudInstanceId(),
type, search, page, pageSize);
+ ListTopicsResponse response = clientFactory.call(ctx.credentialId(),
ctx.regionId(),
+ client -> client.listTopics(request));
+ ListTopicsResponseBody body = response == null ? null :
response.getBody();
+ return body == null ? null : body.getData();
+ }
+
+ private ListTopicsRequest buildTopicRequest(String cloudInstanceId, String
type, String search,
+ int page, int pageSize) {
+ ListTopicsRequest.Builder builder = ListTopicsRequest.builder()
+ .instanceId(cloudInstanceId)
+ .pageNumber(page)
+ .pageSize(pageSize);
+ if (StringUtils.hasText(search)) {
+ builder.filter(search);
+ }
+ if (StringUtils.hasText(type)) {
+
builder.messageTypes(List.of(type.trim().toUpperCase(Locale.ROOT)));
+ }
+ return builder.build();
+ }
+
+ private static List<TopicVO> toTopics(List<ListTopicsResponseBody.List>
rows, String instanceId) {
+ if (rows == null || rows.isEmpty()) {
+ return List.of();
+ }
+ List<TopicVO> topics = new ArrayList<>(rows.size());
+ for (ListTopicsResponseBody.List row : rows) {
+ if (row != null) {
+ topics.add(AliyunConverters.toTopicVO(row, instanceId));
+
}
}
return topics;
}
+ private static boolean hasFetchedAll(int page, int pageSize, Long
totalCount) {
+ return totalCount != null && totalCount >= 0L && (long) page *
pageSize >= totalCount;
+ }
+
+ private static PageResult<TopicVO> paginate(List<TopicVO> topics, int
page, int pageSize) {
+ int total = topics.size();
+ long offset = Pagination.pageOffset(page, pageSize);
+ int from = (int) Math.min(offset, total);
+ int to = from + (int) Math.min(pageSize, total - from);
+ return PageResult.of(topics.subList(from, to), total, page, pageSize);
+ }
+
@Override
public TopicVO createTopic(String instanceId, TopicVO topic) {
Context ctx = resolve(instanceId);
@@ -494,13 +538,6 @@ public class AliyunInstanceProvider implements
InstanceProvider {
return new Context(instance.getCloudInstanceId(),
instance.getRegionId(), instance.getCredentialId());
}
- private static boolean matchesType(String type, TopicVO vo) {
- if (!StringUtils.hasText(type)) {
- return true;
- }
- return vo.getType() != null &&
vo.getType().name().equalsIgnoreCase(type.trim());
- }
-
private record Context(String cloudInstanceId, String regionId, Long
credentialId) {
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
index d073606b3..361b6ac7c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
@@ -39,16 +39,19 @@ import
com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListRequest;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListResponse;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicRequest;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicResponse;
+import com.tencentcloudapi.trocket.v20230308.models.Filter;
import com.tencentcloudapi.trocket.v20230308.models.ModifyTopicRequest;
import
com.tencentcloudapi.trocket.v20230308.models.ResetConsumerGroupOffsetRequest;
import com.tencentcloudapi.trocket.v20230308.models.SubscriptionData;
import com.tencentcloudapi.trocket.v20230308.models.TopicItem;
+import org.apache.rocketmq.studio.common.domain.PageResult;
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.InstanceVendor;
import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.common.util.Pagination;
import org.apache.rocketmq.studio.instance.InstanceRepository;
import org.apache.rocketmq.studio.instance.InstanceVO;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
@@ -187,39 +190,94 @@ public class TencentInstanceProvider implements
InstanceProvider {
return listTopics(instanceId, type, search, true);
}
+ @Override
+ public PageResult<TopicVO> listTopicsPage(String instanceId, String type,
String search, int page, int pageSize) {
+ Context context = resolve(instanceId);
+ DescribeTopicListResponse response = describeTopics(context, type,
search,
+ Pagination.pageOffset(page, pageSize), pageSize);
+ Long totalCount = response == null ? null : response.getTotalCount();
+ if (totalCount == null || totalCount < 0L) {
+ return paginate(listTopics(instanceId, type, search, true), page,
pageSize);
+ }
+ return PageResult.of(toTopics(response.getData(), instanceId, context,
true), totalCount, page, pageSize);
+ }
+
private List<TopicVO> listTopics(String instanceId, String type, String
search, boolean enrichTimes) {
Context context = resolve(instanceId);
List<TopicVO> topics = new ArrayList<>();
- for (int page = 0; page < MAX_PAGES; page++) {
- DescribeTopicListRequest request = new DescribeTopicListRequest();
- request.setInstanceId(context.cloudInstanceId());
- request.setOffset((long) page * PAGE_SIZE);
- request.setLimit((long) PAGE_SIZE);
- DescribeTopicListResponse response =
clientFactory.call(context.credentialId(), context.regionId(),
- client -> client.DescribeTopicList(request));
+ for (long offset = 0L; ; offset += PAGE_SIZE) {
+ DescribeTopicListResponse response = describeTopics(context, type,
search, offset, PAGE_SIZE);
TopicItem[] data = response == null ? null : response.getData();
if (data == null || data.length == 0) {
break;
}
- for (TopicItem item : data) {
- if (item == null) {
- continue;
- }
- TopicVO topic = toTopic(item, instanceId);
- if (matchesType(type, topic) && matchesSearch(search, topic)) {
- if (enrichTimes) {
- enrichTopicTimes(context, topic);
- }
- topics.add(topic);
- }
- }
- if (data.length < PAGE_SIZE) {
+ topics.addAll(toTopics(data, instanceId, context, enrichTimes));
+ if (hasFetchedAll(offset, PAGE_SIZE, response.getTotalCount()) ||
data.length < PAGE_SIZE) {
break;
}
}
return topics;
}
+ private DescribeTopicListResponse describeTopics(Context context, String
type, String search, long offset, long limit) {
+ DescribeTopicListRequest request = new DescribeTopicListRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setOffset(offset);
+ request.setLimit(limit);
+ Filter[] filters = topicFilters(type, search);
+ if (filters.length > 0) {
+ request.setFilters(filters);
+ }
+ return clientFactory.call(context.credentialId(), context.regionId(),
client -> client.DescribeTopicList(request));
+ }
+
+ private static Filter[] topicFilters(String type, String search) {
+ List<Filter> filters = new ArrayList<>(2);
+ if (StringUtils.hasText(search)) {
+ Filter filter = new Filter();
+ filter.setName("TopicName");
+ filter.setValues(new String[]{search.trim()});
+ filters.add(filter);
+ }
+ if (StringUtils.hasText(type)) {
+ Filter filter = new Filter();
+ filter.setName("TopicType");
+ filter.setValues(new
String[]{type.trim().toUpperCase(Locale.ROOT)});
+ filters.add(filter);
+ }
+ return filters.toArray(Filter[]::new);
+ }
+
+ private List<TopicVO> toTopics(TopicItem[] data, String instanceId,
Context context, boolean enrichTimes) {
+ if (data == null || data.length == 0) {
+ return List.of();
+ }
+ List<TopicVO> topics = new ArrayList<>(data.length);
+ for (TopicItem item : data) {
+ if (item == null) {
+ continue;
+ }
+ TopicVO topic = toTopic(item, instanceId);
+ if (enrichTimes) {
+ enrichTopicTimes(context, topic);
+ }
+ topics.add(topic);
+ }
+ return topics;
+ }
+
+ private static boolean hasFetchedAll(long offset, int pageSize, Long
totalCount) {
+ return totalCount != null && totalCount >= 0L && offset + pageSize >=
totalCount;
+ }
+
+ private static PageResult<TopicVO> paginate(List<TopicVO> topics, int
page, int pageSize) {
+ int total = topics.size();
+ long offset = Pagination.pageOffset(page, pageSize);
+ int from = (int) Math.min(offset, total);
+ int to = from + (int) Math.min(pageSize, total - from);
+ return PageResult.of(topics.subList(from, to), total, page, pageSize);
+ }
+
/**
* DescribeTopicList does not expose creation/update timestamps, so
resolve them per-topic
* from DescribeTopic. Kept off the cheap count path to avoid N+1 calls
for instance listings.
@@ -963,19 +1021,6 @@ public class TencentInstanceProvider implements
InstanceProvider {
return null;
}
- private static boolean matchesType(String type, TopicVO topic) {
- return !StringUtils.hasText(type)
- || topic.getType() != null &&
topic.getType().name().equalsIgnoreCase(type.trim());
- }
-
- private static boolean matchesSearch(String search, TopicVO topic) {
- if (!StringUtils.hasText(search)) {
- return true;
- }
- String needle = search.trim().toLowerCase(Locale.ROOT);
- return contains(topic.getName(), needle) ||
contains(topic.getRemark(), needle);
- }
-
private static boolean matchesSearch(String search, String value) {
if (!StringUtils.hasText(search)) {
return true;
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 357375c99..676ecf12b 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
@@ -63,6 +63,7 @@ import java.util.Map;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.function.Function;
+import java.util.stream.IntStream;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -110,13 +111,17 @@ class AliyunInstanceProviderTest {
void listTopicsShouldMapMessageTypeAndFilterTest() {
stubInstance();
stubCallThrough();
- ListTopicsResponse response = topicsResponse(
- topicRow("topic-normal", "NORMAL"),
- null,
- topicRow("topic-fifo", "FIFO"),
- topicRow("topic-mystery", "MYSTERY"));
- when(asyncClient.listTopics(any(ListTopicsRequest.class)))
- .thenReturn(CompletableFuture.completedFuture(response));
+
when(asyncClient.listTopics(any(ListTopicsRequest.class))).thenAnswer(invocation
-> {
+ ListTopicsRequest request = invocation.getArgument(0);
+ if (request.getMessageTypes() != null &&
request.getMessageTypes().contains("FIFO")) {
+ return
CompletableFuture.completedFuture(topicsResponse(topicRow("topic-fifo",
"FIFO")));
+ }
+ return CompletableFuture.completedFuture(topicsResponse(
+ topicRow("topic-normal", "NORMAL"),
+ null,
+ topicRow("topic-fifo", "FIFO"),
+ topicRow("topic-mystery", "MYSTERY")));
+ });
List<TopicVO> all = provider.listTopics(STUDIO_INSTANCE_ID, null,
null);
@@ -135,6 +140,57 @@ class AliyunInstanceProviderTest {
assertThat(fifos.get(0).getType()).isEqualTo(TopicType.FIFO);
}
+ @Test
+ void listTopicsShouldTraversePastLegacyFivePageCapTest() {
+ stubInstance();
+ stubCallThrough();
+
when(asyncClient.listTopics(any(ListTopicsRequest.class))).thenAnswer(invocation
-> {
+ ListTopicsRequest request = invocation.getArgument(0);
+ int pageNumber = request.getPageNumber();
+ if (pageNumber <= 5) {
+ return CompletableFuture.completedFuture(topicsResponse(501L,
pageNumber, AliyunConverters.PAGE_SIZE,
+ IntStream.range(0, AliyunConverters.PAGE_SIZE)
+ .mapToObj(index -> topicRow("topic-" +
((pageNumber - 1) * AliyunConverters.PAGE_SIZE + index),
+ "NORMAL"))
+ .toArray(ListTopicsResponseBody.List[]::new)));
+ }
+ return CompletableFuture.completedFuture(topicsResponse(501L,
pageNumber, AliyunConverters.PAGE_SIZE,
+ topicRow("topic-500", "NORMAL")));
+ });
+
+ List<TopicVO> topics = provider.listTopics(STUDIO_INSTANCE_ID, null,
null);
+
+ assertThat(topics).hasSize(501);
+ ArgumentCaptor<ListTopicsRequest> captor =
ArgumentCaptor.forClass(ListTopicsRequest.class);
+ verify(asyncClient, times(6)).listTopics(captor.capture());
+
assertThat(captor.getAllValues()).extracting(ListTopicsRequest::getPageNumber)
+ .containsExactly(1, 2, 3, 4, 5, 6);
+ }
+
+ @Test
+ void listTopicsPageShouldUseAliyunNativePaginationAndFiltersTest() {
+ stubInstance();
+ stubCallThrough();
+
when(asyncClient.listTopics(any(ListTopicsRequest.class))).thenReturn(CompletableFuture.completedFuture(
+ topicsResponse(321L, 3L, 20L,
+ topicRow("orders-fifo-40", "FIFO"),
+ topicRow("orders-fifo-41", "FIFO"))));
+
+ var page = provider.listTopicsPage(STUDIO_INSTANCE_ID, "fifo",
"orders", 3, 20);
+
+ assertThat(page.getTotal()).isEqualTo(321);
+ assertThat(page.getPage()).isEqualTo(3);
+ assertThat(page.getSize()).isEqualTo(20);
+ assertThat(page.getItems()).extracting(TopicVO::getName)
+ .containsExactly("orders-fifo-40", "orders-fifo-41");
+ ArgumentCaptor<ListTopicsRequest> captor =
ArgumentCaptor.forClass(ListTopicsRequest.class);
+ verify(asyncClient).listTopics(captor.capture());
+ assertThat(captor.getValue().getPageNumber()).isEqualTo(3);
+ assertThat(captor.getValue().getPageSize()).isEqualTo(20);
+ assertThat(captor.getValue().getFilter()).isEqualTo("orders");
+
assertThat(captor.getValue().getMessageTypes()).containsExactly("FIFO");
+ }
+
@Test
void listConsumerGroupsShouldMapGroupIdTest() {
stubInstance();
@@ -506,14 +562,19 @@ class AliyunInstanceProviderTest {
}
private static ListTopicsResponse
topicsResponse(ListTopicsResponseBody.List... rows) {
+ return topicsResponse((long) rows.length, 1L, 100L, rows);
+ }
+
+ private static ListTopicsResponse topicsResponse(long totalCount, long
pageNumber, long pageSize,
+ ListTopicsResponseBody.List... rows) {
return ListTopicsResponse.create().toBuilder()
.statusCode(200)
.body(ListTopicsResponseBody.builder()
.data(ListTopicsResponseBody.Data.builder()
.list(java.util.Arrays.asList(rows))
- .pageNumber(1L)
- .pageSize(100L)
- .totalCount((long) rows.length)
+ .pageNumber(pageNumber)
+ .pageSize(pageSize)
+ .totalCount(totalCount)
.build())
.build())
.build();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
index 3f63461e3..640ea6895 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
@@ -36,6 +36,7 @@ import
com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListRequest;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListResponse;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicRequest;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicResponse;
+import com.tencentcloudapi.trocket.v20230308.models.Filter;
import com.tencentcloudapi.trocket.v20230308.models.ModifyTopicRequest;
import
com.tencentcloudapi.trocket.v20230308.models.ResetConsumerGroupOffsetRequest;
import com.tencentcloudapi.trocket.v20230308.models.SubscriptionData;
@@ -166,9 +167,17 @@ class TencentInstanceProviderTest {
void listTopicsShouldMapAndFilterAndEnrichTimesTest() throws Exception {
TopicItem normal = topicItem("orders", "NORMAL", 8L);
TopicItem fifo = topicItem("orders-fifo", "FIFO", 4L);
- DescribeTopicListResponse response = new DescribeTopicListResponse();
- response.setData(new TopicItem[]{normal, fifo});
- when(client.DescribeTopicList(any())).thenReturn(response);
+ when(client.DescribeTopicList(any())).thenAnswer(invocation -> {
+ DescribeTopicListRequest request = invocation.getArgument(0);
+ DescribeTopicListResponse response = new
DescribeTopicListResponse();
+ Filter[] filters = request.getFilters();
+ if (filters != null && filters.length == 2) {
+ response.setData(new TopicItem[]{fifo});
+ } else {
+ response.setData(new TopicItem[]{normal, fifo});
+ }
+ return response;
+ });
DescribeTopicResponse detail = new DescribeTopicResponse();
detail.setCreatedTime(1600000000000L);
detail.setLastUpdateTime(1600000100000L);
@@ -188,6 +197,66 @@ class TencentInstanceProviderTest {
assertThat(topics.get(0).getGmtModified())
.isEqualTo(java.time.Instant.ofEpochMilli(1600000100000L)
.atZone(java.time.ZoneId.systemDefault()).toLocalDateTime());
+ ArgumentCaptor<DescribeTopicListRequest> captor =
+ ArgumentCaptor.forClass(DescribeTopicListRequest.class);
+ verify(client).DescribeTopicList(captor.capture());
+ assertThat(captor.getValue().getFilters()).extracting(Filter::getName)
+ .containsExactly("TopicName", "TopicType");
+
assertThat(captor.getValue().getFilters()[0].getValues()).containsExactly("fifo");
+
assertThat(captor.getValue().getFilters()[1].getValues()).containsExactly("FIFO");
+ }
+
+ @Test
+ void listTopicsPageShouldUseTencentNativePaginationAndFiltersTest() throws
Exception {
+ TopicItem item = topicItem("orders-fifo-10000", "FIFO", 8L);
+ DescribeTopicListResponse response = new DescribeTopicListResponse();
+ response.setTotalCount(10001L);
+ response.setData(new TopicItem[]{item});
+ when(client.DescribeTopicList(any())).thenReturn(response);
+ DescribeTopicResponse detail = new DescribeTopicResponse();
+ detail.setCreatedTime(1700000000000L);
+ detail.setLastUpdateTime(1700000100000L);
+ when(client.DescribeTopic(any())).thenReturn(detail);
+
+ var page = provider.listTopicsPage(STUDIO_INSTANCE_ID, "fifo",
"orders", 101, 100);
+
+ assertThat(page.getTotal()).isEqualTo(10001L);
+ assertThat(page.getPage()).isEqualTo(101);
+ assertThat(page.getSize()).isEqualTo(100);
+
assertThat(page.getItems()).extracting(TopicVO::getName).containsExactly("orders-fifo-10000");
+ ArgumentCaptor<DescribeTopicListRequest> captor =
+ ArgumentCaptor.forClass(DescribeTopicListRequest.class);
+ verify(client).DescribeTopicList(captor.capture());
+ assertThat(captor.getValue().getOffset()).isEqualTo(10000L);
+ assertThat(captor.getValue().getLimit()).isEqualTo(100L);
+ assertThat(captor.getValue().getFilters()).extracting(Filter::getName)
+ .containsExactly("TopicName", "TopicType");
+ }
+
+ @Test
+ void
countTopicsShouldFallBackToCompleteListingPastLegacyTenThousandCapTest() throws
Exception {
+ when(client.DescribeTopicList(any())).thenAnswer(invocation -> {
+ DescribeTopicListRequest request = invocation.getArgument(0);
+ DescribeTopicListResponse response = new
DescribeTopicListResponse();
+ if (request.getLimit() == 1L && request.getOffset() == 0L) {
+ response.setData(new TopicItem[]{topicItem("seed", "NORMAL",
8L)});
+ return response;
+ }
+ int start = request.getOffset().intValue();
+ int count = Math.min(request.getLimit().intValue(), Math.max(10001
- start, 0));
+ response.setTotalCount(10001L);
+ response.setData(IntStream.range(0, count)
+ .mapToObj(index -> topicItem("topic-" + (start + index),
"NORMAL", 8L))
+ .toArray(TopicItem[]::new));
+ return response;
+ });
+
+ assertThat(provider.countTopics(STUDIO_INSTANCE_ID)).isEqualTo(10001);
+ ArgumentCaptor<DescribeTopicListRequest> captor =
+ ArgumentCaptor.forClass(DescribeTopicListRequest.class);
+ verify(client, times(102)).DescribeTopicList(captor.capture());
+ assertThat(captor.getAllValues().get(0).getLimit()).isEqualTo(1L);
+
assertThat(captor.getAllValues().get(101).getOffset()).isEqualTo(10000L);
}
@Test