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 c7c8e8612 feat(inventory): bounded producer-group discovery and paged
consumer-group inventory (#2331)
c7c8e8612 is described below
commit c7c8e8612a8578fcc754b78cdf4968958b63fa9a
Author: aias00 <[email protected]>
AuthorDate: Wed Aug 19 15:16:55 2026 +0800
feat(inventory): bounded producer-group discovery and paged consumer-group
inventory (#2331)
* Keep producer-group discovery bounded to selector queries
Constraint: Studio still needs free-form producer group entry while
selector suggestions stay lightweight and deterministic.
Rejected: reuse findConnections | it builds a full producer connection
inventory just to render selector options.
Directive: keep selector discovery separate from detailed producer
connection lookup.
Confidence: medium
Scope-risk: narrow
Tested: export
JAVA_HOME=/Users/aias/Library/Java/JavaVirtualMachines/openjdk-21.0.2/Contents/Home;
export PATH="$JAVA_HOME/bin:$PATH"; mvn -q
-Dtest=ProducerControllerTest,ProducerConnectionServiceTest,ClientProviderStubTest,RocketMQClientProviderTest
test
Tested: export
JAVA_HOME=/Users/aias/Library/Java/JavaVirtualMachines/openjdk-21.0.2/Contents/Home;
export PATH="$JAVA_HOME/bin:$PATH"; mvn -q -DskipTests compile
Tested: npm exec eslint -- src/pages/studio/Producer.tsx
src/pages/studio/__tests__/Producer.test.tsx src/api/producer.ts
src/api/producer.test.ts
Tested: npm test -- --run src/api/producer.test.ts
src/pages/studio/__tests__/Producer.test.tsx
Tested: npm run build
* Preserve consumer-group inventory pagination semantics across Studio
Constraint: keep the legacy /api/groups list contract intact while adding a
paged metadata inventory path for the web consumer page
Rejected: replace /api/groups with PageResult | would break existing
callers including AI/tooling paths
Directive: keep PageResult pagination 1-based and preserve clusterId/search
filtering before slicing
Confidence: high
Scope-risk: narrow
Tested:
JAVA_HOME=/Users/aias/Library/Java/JavaVirtualMachines/openjdk-21.0.2/Contents/Home
mvn -f server/pom.xml -Dtest=ConsumerGroupControllerTest,MetadataServiceTest
test
Tested: npm test -- --run src/api/consumerGroups.test.ts
src/services/consumerService.test.ts
src/pages/instance/__tests__/ConsumerPage.test.tsx
Tested: npm run build
---
.../studio/cluster/client/ClientProvider.java | 2 +
.../studio/cluster/client/ClientProviderStub.java | 8 ++
.../cluster/client/ProducerConnectionService.java | 32 +++--
.../studio/cluster/client/ProducerController.java | 8 +-
.../instance/group/ConsumerGroupController.java | 11 ++
.../studio/instance/topic/MetadataService.java | 23 ++++
.../provider/apache/RocketMQClientProvider.java | 50 +++++++
.../cluster/client/ClientProviderStubTest.java | 8 ++
.../client/ProducerConnectionServiceTest.java | 27 ++--
.../cluster/client/ProducerControllerTest.java | 10 +-
.../group/ConsumerGroupControllerTest.java | 36 ++++++
.../studio/instance/topic/MetadataServiceTest.java | 51 ++++++++
.../apache/RocketMQClientProviderTest.java | 33 +++++
web/src/api/consumerGroups.test.ts | 18 +++
web/src/api/metadata.ts | 17 +++
web/src/api/producer.test.ts | 19 ++-
web/src/api/producer.ts | 18 ++-
.../pages/instance/__tests__/ConsumerPage.test.tsx | 143 +++++++++++++--------
web/src/pages/instance/consumer.tsx | 52 +++++---
web/src/pages/studio/Producer.tsx | 59 +++++----
web/src/pages/studio/__tests__/Producer.test.tsx | 34 ++++-
web/src/services/consumerService.test.ts | 27 ++++
web/src/services/consumerService.ts | 38 +++++-
23 files changed, 592 insertions(+), 132 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
index 6a02e68cd..b04736009 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
@@ -27,5 +27,7 @@ public interface ClientProvider {
throw new BusinessException(501, "Client connection provider does not
support nameserver lookup");
}
+ List<String> findProducerGroups(String instanceId, String topic, String
query, int limit);
+
List<ClientConnectionVO> findProducerConnections(String instanceId, String
topic, String producerGroup);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStub.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStub.java
index c83d2b00e..9eec263de 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStub.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStub.java
@@ -32,6 +32,14 @@ public class ClientProviderStub implements ClientProvider {
throw new BusinessException(501, "Client connection provider is not
configured");
}
+ @Override
+ public List<String> findProducerGroups(String instanceId, String topic,
String query, int limit) {
+ log.warn("ClientProviderStub.findProducerGroups called without a real
client provider. "
+ + "instanceId={}, topic={}, query={}, limit={}",
+ instanceId, topic, query, limit);
+ throw new BusinessException(501, "Client connection provider is not
configured");
+ }
+
@Override
public List<ClientConnectionVO> findProducerConnections(String instanceId,
String topic, String producerGroup) {
log.warn("ClientProviderStub.findProducerConnections called without a
real client provider. "
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
index bf32351ce..02211b137 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
@@ -16,7 +16,6 @@
*/
package org.apache.rocketmq.studio.cluster.client;
-import org.apache.rocketmq.studio.common.domain.enums.ClientType;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
@@ -29,6 +28,9 @@ import java.util.List;
@RequiredArgsConstructor
public class ProducerConnectionService {
+ private static final int DEFAULT_PRODUCER_GROUP_SELECTOR_LIMIT = 20;
+ private static final int MAX_PRODUCER_GROUP_SELECTOR_LIMIT = 100;
+
private final ClientProvider clientProvider;
public List<ProducerConnectionVO> listConnections(String instanceId,
String topic, String producerGroup) {
@@ -43,15 +45,13 @@ public class ProducerConnectionService {
.toList();
}
- public List<String> listProducerGroups(String instanceId) {
+ public List<String> listProducerGroups(String instanceId, String topic,
String query, Integer limit) {
String normalizedInstanceId = requireFilter(instanceId, "instanceId");
- return clientProvider.findConnections(normalizedInstanceId, null,
ClientType.Producer.name()).stream()
- .map(ClientConnectionVO::getProducerGroup)
- .filter(this::hasText)
- .map(String::trim)
- .distinct()
- .sorted()
- .toList();
+ return clientProvider.findProducerGroups(
+ normalizedInstanceId,
+ normalizeOptionalFilter(topic),
+ normalizeOptionalFilter(query),
+ normalizeSelectorLimit(limit));
}
private ProducerConnectionVO toProducerConnection(ClientConnectionVO
connection) {
@@ -67,6 +67,20 @@ public class ProducerConnectionService {
return value != null && !value.trim().isEmpty();
}
+ private String normalizeOptionalFilter(String value) {
+ return hasText(value) ? value.trim() : null;
+ }
+
+ private int normalizeSelectorLimit(Integer limit) {
+ if (limit == null) {
+ return DEFAULT_PRODUCER_GROUP_SELECTOR_LIMIT;
+ }
+ if (limit < 1) {
+ return 1;
+ }
+ return Math.min(limit, MAX_PRODUCER_GROUP_SELECTOR_LIMIT);
+ }
+
private String requireFilter(String value, String fieldName) {
if (!hasText(value)) {
throw new BusinessException(400, fieldName + " is required");
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
index 163d4319e..d83db7dba 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
@@ -34,8 +34,12 @@ public class ProducerController {
private final ProducerConnectionService producerConnectionService;
@GetMapping("/groups")
- public Result<List<String>> listProducerGroups(@RequestParam String
instanceId) {
- return
Result.ok(producerConnectionService.listProducerGroups(instanceId));
+ public Result<List<String>> listProducerGroups(
+ @RequestParam String instanceId,
+ @RequestParam(required = false) String topic,
+ @RequestParam(required = false) String query,
+ @RequestParam(required = false) Integer limit) {
+ return
Result.ok(producerConnectionService.listProducerGroups(instanceId, topic,
query, limit));
}
@GetMapping("/connection")
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
index eb5d8d0bc..bf7988d31 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.studio.instance.group;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.instance.topic.MetadataService;
import org.apache.rocketmq.studio.common.domain.Result;
@@ -48,6 +49,16 @@ public class ConsumerGroupController {
return Result.ok(metadataService.listConsumerGroups(instanceId,
clusterId, search));
}
+ @GetMapping("/page")
+ public Result<PageResult<ConsumerGroupVO>> listConsumerGroupsPage(
+ @RequestParam(required = false) String instanceId,
+ @RequestParam(required = false) String clusterId,
+ @RequestParam(required = false) String search,
+ @RequestParam(defaultValue = "1") int page,
+ @RequestParam(defaultValue = "20") int pageSize) {
+ return Result.ok(metadataService.listConsumerGroupsPage(instanceId,
clusterId, search, page, pageSize));
+ }
+
@GetMapping("/{name}")
public Result<ConsumerGroupVO> getConsumerGroup(
@PathVariable String name,
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
index 4d05c81ab..2b61fb0b4 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
@@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.instance.topic;
import org.apache.rocketmq.studio.provider.apache.AdminClient;
import org.apache.rocketmq.studio.provider.apache.MetadataProvider;
+import org.apache.rocketmq.studio.common.util.Pagination;
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;
@@ -39,6 +40,8 @@ import java.util.List;
@RequiredArgsConstructor
public class MetadataService {
+ private static final int MAX_PAGE_SIZE = 100;
+
private final MetadataProvider metadataProvider;
private final AdminClient adminClient;
private final InstanceProviderRegistry providerRegistry;
@@ -179,6 +182,17 @@ public class MetadataService {
return resolve(instanceId).listConsumerGroups(instanceId,
normalizeFilter(search));
}
+ public PageResult<ConsumerGroupVO> listConsumerGroupsPage(String
instanceId, String clusterId, String search,
+ int page, int
pageSize) {
+ validatePagination(page, pageSize);
+
+ List<ConsumerGroupVO> groups = listConsumerGroups(instanceId,
clusterId, search);
+ int total = groups.size();
+ int from = (int) Math.min(Pagination.pageOffset(page, pageSize),
total);
+ int to = Math.min(from + pageSize, total);
+ return PageResult.of(groups.subList(from, to), total, page, pageSize);
+ }
+
public ConsumerGroupVO getConsumerGroup(String name) {
return getConsumerGroup(null, name);
@@ -255,6 +269,15 @@ public class MetadataService {
return !StringUtils.hasText(value) ? null : value.trim();
}
+ private void validatePagination(int page, int pageSize) {
+ if (page < 1) {
+ throw new BusinessException(400, "page must be greater than zero");
+ }
+ if (pageSize < 1 || pageSize > MAX_PAGE_SIZE) {
+ throw new BusinessException(400, "pageSize must be between 1 and "
+ MAX_PAGE_SIZE);
+ }
+ }
+
private void requireTopic(TopicVO topic) {
if (topic == null) {
throw new BusinessException(400, "Topic request is required");
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
index df2acd67f..d7a0d863d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
@@ -44,6 +44,8 @@ import org.springframework.context.annotation.Primary;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.LinkedHashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
@@ -95,6 +97,38 @@ public class RocketMQClientProvider implements
ClientProvider {
adminExt -> findProducerConnections(adminExt, topic,
producerGroup));
}
+ @Override
+ public List<String> findProducerGroups(String instanceId, String topic,
String query, int limit) {
+ return runtimeAdminClientResolver.execute(instanceId,
+ adminExt -> findProducerGroups(adminExt, topic, query, limit));
+ }
+
+ private List<String> findProducerGroups(MQAdminExt adminExt, String topic,
String query, int limit) {
+ BrokerTopology topology = discoverBrokerTopology(adminExt, null,
"producer group selector");
+ if (topology.brokerAddresses().isEmpty()) {
+ return List.of();
+ }
+ String normalizedQuery = query == null ? null :
query.toLowerCase(Locale.ROOT);
+ LinkedHashSet<String> groups = new LinkedHashSet<>();
+ int successfulBrokers = 0;
+ for (String brokerAddress : topology.brokerAddresses()) {
+ try {
+ ProducerTableInfo producerTable =
adminExt.getAllProducerInfo(brokerAddress);
+ successfulBrokers++;
+ collectProducerGroups(groups, producerTable, normalizedQuery);
+ } catch (Exception e) {
+ log.warn("Failed to fetch producer groups from broker={},
skipping", brokerAddress, e);
+ }
+ }
+ if (successfulBrokers == 0) {
+ throw new BusinessException(502, "Failed to query producer groups
from all brokers");
+ }
+ return groups.stream()
+ .sorted(Comparator.naturalOrder())
+ .limit(limit)
+ .toList();
+ }
+
private List<ClientConnectionVO> findProducerConnections(MQAdminExt
adminExt, String topic, String producerGroup) {
try {
ProducerConnection producerConnection =
@@ -201,6 +235,22 @@ public class RocketMQClientProvider implements
ClientProvider {
});
}
+ private void collectProducerGroups(
+ LinkedHashSet<String> groups,
+ ProducerTableInfo producerTable,
+ String normalizedQuery) {
+ if (producerTable == null || producerTable.getData() == null) {
+ return;
+ }
+ producerTable.getData().keySet().stream()
+ .filter(Objects::nonNull)
+ .map(String::trim)
+ .filter(group -> !group.isEmpty())
+ .filter(group -> normalizedQuery == null
+ ||
group.toLowerCase(Locale.ROOT).contains(normalizedQuery))
+ .forEach(groups::add);
+ }
+
private ClientConnectionVO toConnectionVO(
ProducerInfo producerInfo, String producerGroup, String clusterId)
{
String remoteIp = producerInfo.getRemoteIP();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStubTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStubTest.java
index 6308f775b..28663699a 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStubTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStubTest.java
@@ -41,6 +41,14 @@ class ClientProviderStubTest {
.satisfies(ex -> assertThatBusinessExceptionCode(ex, 501));
}
+ @Test
+ void findProducerGroupsShouldFailWhenRealProviderIsMissing() {
+ assertThatThrownBy(() -> provider.findProducerGroups("instance-1",
"order-topic", "pg", 20))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Client connection provider is not configured")
+ .satisfies(ex -> assertThatBusinessExceptionCode(ex, 501));
+ }
+
private void assertThatBusinessExceptionCode(Throwable ex, int code) {
org.assertj.core.api.Assertions.assertThat(((BusinessException)
ex).getCode()).isEqualTo(code);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
index 933d82690..1de90e871 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
@@ -97,18 +97,23 @@ class ProducerConnectionServiceTest {
}
@Test
- void listProducerGroupsShouldReturnSortedUniqueActiveGroups() {
- when(clientProvider.findConnections("instance-1", null,
ClientType.Producer.name()))
- .thenReturn(List.of(
- ClientConnectionVO.builder().producerGroup("
pg-payment ").build(),
-
ClientConnectionVO.builder().producerGroup("pg-order").build(),
-
ClientConnectionVO.builder().producerGroup("pg-payment").build(),
- ClientConnectionVO.builder().producerGroup("
").build(),
- ClientConnectionVO.builder().build()));
-
- assertThat(producerConnectionService.listProducerGroups("instance-1"))
+ void
listProducerGroupsShouldDelegateSelectorDiscoveryWithNormalizedFilters() {
+ when(clientProvider.findProducerGroups("instance-1", "order-topic",
"pg", 100))
+ .thenReturn(List.of("pg-order", "pg-payment"));
+
+ assertThat(producerConnectionService.listProducerGroups(" instance-1
", " order-topic ", " pg ", 1000))
.containsExactly("pg-order", "pg-payment");
- verify(clientProvider).findConnections("instance-1", null,
ClientType.Producer.name());
+ verify(clientProvider).findProducerGroups("instance-1", "order-topic",
"pg", 100);
+ }
+
+ @Test
+ void listProducerGroupsShouldApplyDefaultSelectorLimit() {
+ when(clientProvider.findProducerGroups("instance-1", null, null, 20))
+ .thenReturn(List.of("pg-order"));
+
+ assertThat(producerConnectionService.listProducerGroups("instance-1",
" ", " ", null))
+ .containsExactly("pg-order");
+ verify(clientProvider).findProducerGroups("instance-1", null, null,
20);
}
@Test
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
index 5e824875e..f2a1fe9c2 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
@@ -44,15 +44,19 @@ class ProducerControllerTest {
@Test
void listProducerGroupsShouldReturnSuggestions() throws Exception {
- when(producerConnectionService.listProducerGroups("instance-1"))
+ when(producerConnectionService.listProducerGroups("instance-1",
"order-topic", "pg", 20))
.thenReturn(List.of("pg-order", "pg-payment"));
- mockMvc.perform(get("/api/producer/groups").param("instanceId",
"instance-1"))
+ mockMvc.perform(get("/api/producer/groups")
+ .param("instanceId", "instance-1")
+ .param("topic", "order-topic")
+ .param("query", "pg")
+ .param("limit", "20"))
.andExpect(status().isOk())
.andExpect(jsonPath("$.data[0]").value("pg-order"))
.andExpect(jsonPath("$.data[1]").value("pg-payment"));
- verify(producerConnectionService).listProducerGroups("instance-1");
+ verify(producerConnectionService).listProducerGroups("instance-1",
"order-topic", "pg", 20);
}
@Test
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
index abcb2b05b..205a6d7cc 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
@@ -18,6 +18,7 @@
package org.apache.rocketmq.studio.instance.group;
import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.instance.topic.MetadataService;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
@@ -63,6 +64,41 @@ class ConsumerGroupControllerTest {
@MockBean
private ConsumerDiagnosticsService consumerDiagnosticsService;
+ @Test
+ void listConsumerGroupsShouldPassQueryParams() throws Exception {
+ when(metadataService.listConsumerGroups("instance-a", "cluster-a",
"orders"))
+ .thenReturn(List.of());
+
+ mockMvc.perform(get("/api/groups")
+ .param("instanceId", "instance-a")
+ .param("clusterId", "cluster-a")
+ .param("search", "orders"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data").isArray());
+
+ verify(metadataService).listConsumerGroups("instance-a", "cluster-a",
"orders");
+ }
+
+ @Test
+ void listConsumerGroupsPageShouldPassSelectedInstanceFiltersAndPaging()
throws Exception {
+ PageResult<ConsumerGroupVO> page = PageResult.of(List.of(), 3, 2, 20);
+ when(metadataService.listConsumerGroupsPage("instance-a", "cluster-a",
"orders", 2, 20))
+ .thenReturn(page);
+
+ mockMvc.perform(get("/api/groups/page")
+ .param("instanceId", "instance-a")
+ .param("clusterId", "cluster-a")
+ .param("search", "orders")
+ .param("page", "2")
+ .param("pageSize", "20"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.total").value(3))
+ .andExpect(jsonPath("$.data.page").value(2))
+ .andExpect(jsonPath("$.data.size").value(20));
+
+ verify(metadataService).listConsumerGroupsPage("instance-a",
"cluster-a", "orders", 2, 20);
+ }
+
@Test
void createConsumerGroupShouldPassValidatedRequest() throws Exception {
Map<String, Object> body = Map.of(
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
index 3da8f87d8..b02659653 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.instance.topic;
+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.instance.group.ConsumerGroupVO;
@@ -267,4 +268,54 @@ class MetadataServiceTest {
verify(metadataProvider).listConsumerGroups("cluster-1", "order");
}
+ @Test
+ void listConsumerGroupsPageShouldPaginateFromOneBasedIndexes() {
+ ConsumerGroupVO first = new ConsumerGroupVO();
+ first.setName("cg-a");
+ ConsumerGroupVO second = new ConsumerGroupVO();
+ second.setName("cg-b");
+ ConsumerGroupVO third = new ConsumerGroupVO();
+ third.setName("cg-c");
+ when(apacheProvider.listConsumerGroups("instance-a", "order"))
+ .thenReturn(List.of(first, second, third));
+
+ PageResult<ConsumerGroupVO> result =
+ metadataService.listConsumerGroupsPage("instance-a", null,
"order", 2, 2);
+
+ assertThat(result.getItems()).containsExactly(third);
+ assertThat(result.getTotal()).isEqualTo(3);
+ assertThat(result.getPage()).isEqualTo(2);
+ assertThat(result.getSize()).isEqualTo(2);
+ verify(apacheProvider).listConsumerGroups("instance-a", "order");
+ }
+
+ @Test
+ void
listConsumerGroupsPageShouldReturnEmptyItemsWhenPageStartsPastFilteredTotal() {
+ ConsumerGroupVO first = new ConsumerGroupVO();
+ first.setName("cg-a");
+ when(metadataProvider.listConsumerGroups("cluster-1",
"order")).thenReturn(List.of(first));
+
+ PageResult<ConsumerGroupVO> result =
+ metadataService.listConsumerGroupsPage(null, "cluster-1",
"order", 2, 1);
+
+ assertThat(result.getItems()).isEmpty();
+ assertThat(result.getTotal()).isEqualTo(1);
+ assertThat(result.getPage()).isEqualTo(2);
+ assertThat(result.getSize()).isEqualTo(1);
+ verify(metadataProvider).listConsumerGroups("cluster-1", "order");
+ }
+
+ @Test
+ void listConsumerGroupsPageShouldRejectInvalidPaginationBounds() {
+ assertThatThrownBy(() -> metadataService.listConsumerGroupsPage(null,
"cluster-1", null, 0, 20))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("page must be greater than zero");
+ assertThatThrownBy(() -> metadataService.listConsumerGroupsPage(null,
"cluster-1", null, 1, 0))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("pageSize must be between 1 and 100");
+ assertThatThrownBy(() -> metadataService.listConsumerGroupsPage(null,
"cluster-1", null, 1, 101))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("pageSize must be between 1 and 100");
+ }
+
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
index c1775d3f0..94e2d6f74 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
@@ -229,6 +229,39 @@ class RocketMQClientProviderTest {
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
}
+ @Test
+ void producerGroupSelectorReturnsSortedUniqueBoundedMatches() throws
Exception {
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo(
+ "127.0.0.1:10911", "127.0.0.2:10911"));
+ when(adminExt.getAllProducerInfo("127.0.0.1:10911"))
+ .thenReturn(new ProducerTableInfo(Map.of(
+ " pg-payment ",
List.of(producerInfo("producer-payment", "10.0.0.2:1000")),
+ "pg-order", List.of(producerInfo("producer-order",
"10.0.0.1:1000")))));
+ when(adminExt.getAllProducerInfo("127.0.0.2:10911"))
+ .thenReturn(new ProducerTableInfo(Map.of(
+ "pg-order", List.of(producerInfo("producer-order-2",
"10.0.0.3:1000")),
+ "pg-shipment",
List.of(producerInfo("producer-shipment", "10.0.0.4:1000")),
+ " ", List.of(producerInfo("ignored",
"10.0.0.5:1000")))));
+
+ List<String> groups = provider.findProducerGroups("instance-a",
"TopicA", "pg", 2);
+
+ assertThat(groups).containsExactly("pg-order", "pg-payment");
+ verify(adminExt, never()).examineProducerConnectionInfo(anyString(),
anyString());
+ }
+
+ @Test
+ void producerGroupSelectorFailsWhenEveryBrokerQueryFails() throws
Exception {
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo(
+ "127.0.0.1:10911", "127.0.0.2:10911"));
+ when(adminExt.getAllProducerInfo(anyString()))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() -> provider.findProducerGroups("instance-a",
"TopicA", "pg", 20))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to query producer groups from all brokers")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
@Test
void exactProducerQueryPassesNonBlankGroupToAdminApi() throws Exception {
ProducerConnection producerConnection = new ProducerConnection();
diff --git a/web/src/api/consumerGroups.test.ts
b/web/src/api/consumerGroups.test.ts
index ae6615b38..8ad59ece1 100644
--- a/web/src/api/consumerGroups.test.ts
+++ b/web/src/api/consumerGroups.test.ts
@@ -22,6 +22,7 @@ import {
getConsumerGroup,
getConsumerProgress,
getConsumerSubscriptions,
+ listConsumerGroupPage,
listConsumerGroups,
deleteConsumerGroup,
resetConsumerOffset,
@@ -66,6 +67,23 @@ describe('consumer groups API contract', () => {
await expect(listConsumerGroups(params)).resolves.toEqual([group]);
});
+ it('uses the paged inventory query fields supported by the backend', async
() => {
+ const params = {
+ instanceId: 'instance-1',
+ clusterId: 'cluster-a',
+ search: 'orders',
+ page: 2,
+ pageSize: 10,
+ };
+ const page = { items: [group], total: 11, page: 2, size: 10 };
+ mock.onGet('/groups/page').reply((config) => {
+ expect(config.params).toEqual(params);
+ return [200, { code: 200, data: page }];
+ });
+
+ await expect(listConsumerGroupPage(params)).resolves.toEqual(page);
+ });
+
it('encodes consumer group names and passes instance context for runtime
queries', async () => {
const groupName = '%RETRY%cg-order';
const instanceId = 'instance-1';
diff --git a/web/src/api/metadata.ts b/web/src/api/metadata.ts
index 3e5d77bf5..953ae9ccc 100644
--- a/web/src/api/metadata.ts
+++ b/web/src/api/metadata.ts
@@ -54,6 +54,13 @@ export interface TopicConsumerPage {
pageSize: number;
}
+export interface PageResult<T> {
+ items: T[];
+ total: number;
+ page: number;
+ size: number;
+}
+
// ─── Consumer Group (matches mock/consumers.ts) ─────────────────
export interface ConsumerGroup {
name: string;
@@ -126,6 +133,11 @@ export interface ConsumerGroupQuery {
search?: string;
}
+export interface ConsumerGroupPageQuery extends ConsumerGroupQuery {
+ page?: number;
+ pageSize?: number;
+}
+
export interface ResetConsumerOffsetRequest {
name: string;
timestamp: number;
@@ -219,6 +231,11 @@ export async function listConsumerGroups(params?:
ConsumerGroupQuery) {
return res.data.data;
}
+export async function listConsumerGroupPage(params?: ConsumerGroupPageQuery) {
+ const res = await client.get<{ data: PageResult<ConsumerGroup>
}>('/groups/page', { params });
+ return res.data.data;
+}
+
export async function getConsumerGroup(name: string, instanceId?: string) {
const res = await client.get<{ data: ConsumerGroupDetail }>(
`/groups/${encodeURIComponent(name)}`,
diff --git a/web/src/api/producer.test.ts b/web/src/api/producer.test.ts
index 54cbfabce..fd45e9773 100644
--- a/web/src/api/producer.test.ts
+++ b/web/src/api/producer.test.ts
@@ -71,12 +71,23 @@ describe('Producer API', () => {
});
it('fetches active producer group suggestions', async () => {
- mock.onGet('/producer/groups').reply(200, {
- code: 200,
- data: ['pg-order', 'pg-payment'],
+ mock.onGet('/producer/groups').reply((config) => {
+ expect(config.params.instanceId).toBe('instance-1');
+ expect(config.params.topic).toBe('order-events');
+ expect(config.params.query).toBe('pg');
+ expect(config.params.limit).toBe(20);
+ return [
+ 200,
+ {
+ code: 200,
+ data: ['pg-order', 'pg-payment'],
+ },
+ ];
});
- await
expect(fetchProducerGroups('instance-1')).resolves.toEqual(['pg-order',
'pg-payment']);
+ await expect(
+ fetchProducerGroups('instance-1', { topic: 'order-events', query: 'pg',
limit: 20 }),
+ ).resolves.toEqual(['pg-order', 'pg-payment']);
});
it('queries producer connections by topic and group', async () => {
diff --git a/web/src/api/producer.ts b/web/src/api/producer.ts
index 26b458853..14b43a228 100644
--- a/web/src/api/producer.ts
+++ b/web/src/api/producer.ts
@@ -155,8 +155,22 @@ export async function fetchTopicList(instanceId: string):
Promise<string[]> {
}
/** Fetch active producer groups for query suggestions */
-export async function fetchProducerGroups(instanceId: string):
Promise<string[]> {
- const res = await client.get<{ data?: string[] }>('/producer/groups', {
params: { instanceId } });
+export async function fetchProducerGroups(
+ instanceId: string,
+ options: {
+ topic?: string;
+ query?: string;
+ limit?: number;
+ } = {},
+): Promise<string[]> {
+ const res = await client.get<{ data?: string[] }>('/producer/groups', {
+ params: {
+ instanceId,
+ topic: options.topic,
+ query: options.query,
+ limit: options.limit,
+ },
+ });
return res.data.data ?? [];
}
diff --git a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
index 454b044ba..68f908204 100644
--- a/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
+++ b/web/src/pages/instance/__tests__/ConsumerPage.test.tsx
@@ -35,7 +35,7 @@ vi.mock('../../../services/consumerService', () => ({
getConsumerProgress: vi.fn(),
getConsumerStack: vi.fn(),
getConsumerSubscriptions: vi.fn(),
- listConsumerGroups: vi.fn(),
+ listConsumerGroupPage: vi.fn(),
resetConsumerOffset: vi.fn(),
}));
const instanceServiceMocks = vi.hoisted(() => ({ listInstances: vi.fn() }));
@@ -84,6 +84,17 @@ const group: ConsumerGroup = {
instances: [],
};
+const groupPage = (
+ items: ConsumerGroup[],
+ overrides: Partial<{ total: number; page: number; size: number }> = {},
+) => ({
+ items,
+ total: items.length,
+ page: 1,
+ size: 20,
+ ...overrides,
+});
+
const renderWithProviders = (ui: React.ReactElement, initialEntry =
'/instance/consumer') =>
render(
<App>
@@ -109,7 +120,7 @@ describe('Consumer page', () => {
gmtModified: '2026-01-01T00:00:00Z',
},
]);
- vi.mocked(consumerService.listConsumerGroups).mockResolvedValue([group]);
+
vi.mocked(consumerService.listConsumerGroupPage).mockResolvedValue(groupPage([group]));
vi.mocked(consumerService.createConsumerGroup).mockImplementation(
async (data: Partial<ConsumerGroup>) =>
({
@@ -180,7 +191,13 @@ describe('Consumer page', () => {
expect(await screen.findByText('remote-cg')).toBeInTheDocument();
expect(screen.getByText('Push')).toBeInTheDocument();
- expect(consumerService.listConsumerGroups).toHaveBeenCalledTimes(1);
+ expect(consumerService.listConsumerGroupPage).toHaveBeenCalledTimes(1);
+ expect(consumerService.listConsumerGroupPage).toHaveBeenCalledWith({
+ instanceId: 'instance-1',
+ page: 1,
+ pageSize: 20,
+ search: undefined,
+ });
});
it('downloads the currently filtered consumer groups when exporting', async
() => {
@@ -191,7 +208,7 @@ describe('Consumer page', () => {
exportedBlob = blob as Blob;
return 'blob:consumer-group-export';
});
- vi.mocked(consumerService.listConsumerGroups).mockResolvedValue([
+ const exportGroups = [
{
...group,
name: 'orders-cg',
@@ -204,7 +221,17 @@ describe('Consumer page', () => {
namespace: '=formula-risk',
subscribedTopics: ['users-topic'],
},
- ]);
+ ];
+ vi.mocked(consumerService.listConsumerGroupPage).mockImplementation(async
(params) => {
+ const filtered = params?.search
+ ? exportGroups.filter((item) => item.name.includes(params.search ??
''))
+ : exportGroups;
+ return groupPage(filtered, {
+ total: filtered.length,
+ page: params?.page ?? 1,
+ size: params?.pageSize ?? 20,
+ });
+ });
renderWithProviders(<ConsumerPage />);
expect(await screen.findByText('orders-cg')).toBeInTheDocument();
@@ -261,9 +288,9 @@ describe('Consumer page', () => {
gmtModified: '2026-07-23T00:00:00Z',
},
]);
- vi.mocked(consumerService.listConsumerGroups).mockResolvedValue([
- { ...group, instanceId: 'instance-a' },
- ]);
+ vi.mocked(consumerService.listConsumerGroupPage).mockResolvedValue(
+ groupPage([{ ...group, instanceId: 'instance-a' }]),
+ );
const user = userEvent.setup();
renderWithProviders(<ConsumerPage />);
@@ -345,9 +372,12 @@ describe('Consumer page', () => {
gmtModified: '2026-07-23T00:00:00Z',
},
]);
- vi.mocked(consumerService.listConsumerGroups).mockImplementation(async
(params) => [
- { ...group, instanceId: params?.instanceId ?? 'instance-a' },
- ]);
+ vi.mocked(consumerService.listConsumerGroupPage).mockImplementation(async
(params) =>
+ groupPage([{ ...group, instanceId: params?.instanceId ?? 'instance-a'
}], {
+ page: params?.page ?? 1,
+ size: params?.pageSize ?? 20,
+ }),
+ );
const user = userEvent.setup();
renderWithProviders(<ConsumerPage />, '/instance/instance-a/consumer');
@@ -364,7 +394,12 @@ describe('Consumer page', () => {
await screen.findByText('instance-b', { selector:
'.ant-select-item-option-content' }),
);
await waitFor(() =>
- expect(consumerService.listConsumerGroups).toHaveBeenCalledWith({
instanceId: 'instance-b' }),
+ expect(consumerService.listConsumerGroupPage).toHaveBeenCalledWith({
+ instanceId: 'instance-b',
+ page: 1,
+ pageSize: 20,
+ search: undefined,
+ }),
);
expect(screen.queryByRole('dialog')).not.toBeInTheDocument();
await user.click(await screen.findByRole('button', { name: /详情/ }));
@@ -421,29 +456,31 @@ describe('Consumer page', () => {
resolveSecond = resolve;
}),
);
- vi.mocked(consumerService.listConsumerGroups).mockResolvedValue([
- {
- ...group,
- instances: [
- {
- clientId: 'client-1',
- protocol: 'Remoting',
- address: '10.0.0.1:1',
- subscribedTopics: [],
- lastHeartbeat: '',
- topicLag: {},
- },
- {
- clientId: 'client-2',
- protocol: 'Remoting',
- address: '10.0.0.2:2',
- subscribedTopics: [],
- lastHeartbeat: '',
- topicLag: {},
- },
- ],
- },
- ]);
+ vi.mocked(consumerService.listConsumerGroupPage).mockResolvedValue(
+ groupPage([
+ {
+ ...group,
+ instances: [
+ {
+ clientId: 'client-1',
+ protocol: 'Remoting',
+ address: '10.0.0.1:1',
+ subscribedTopics: [],
+ lastHeartbeat: '',
+ topicLag: {},
+ },
+ {
+ clientId: 'client-2',
+ protocol: 'Remoting',
+ address: '10.0.0.2:2',
+ subscribedTopics: [],
+ lastHeartbeat: '',
+ topicLag: {},
+ },
+ ],
+ },
+ ]),
+ );
const user = userEvent.setup();
renderWithProviders(<ConsumerPage />);
@@ -461,21 +498,23 @@ describe('Consumer page', () => {
});
it('loads a consumer client stack trace from the selected instance', async
() => {
- vi.mocked(consumerService.listConsumerGroups).mockResolvedValue([
- {
- ...group,
- instances: [
- {
- clientId: 'client-1',
- protocol: 'Remoting',
- address: '10.0.0.1:39210',
- subscribedTopics: ['remote-topic'],
- lastHeartbeat: '2026-07-23T00:00:00Z',
- topicLag: {},
- },
- ],
- },
- ]);
+ vi.mocked(consumerService.listConsumerGroupPage).mockResolvedValue(
+ groupPage([
+ {
+ ...group,
+ instances: [
+ {
+ clientId: 'client-1',
+ protocol: 'Remoting',
+ address: '10.0.0.1:39210',
+ subscribedTopics: ['remote-topic'],
+ lastHeartbeat: '2026-07-23T00:00:00Z',
+ topicLag: {},
+ },
+ ],
+ },
+ ]),
+ );
const user = userEvent.setup();
renderWithProviders(<ConsumerPage />);
@@ -583,7 +622,7 @@ describe('Consumer page', () => {
});
it('keeps per-row state when consumer group CSV import partially fails',
async () => {
- vi.mocked(consumerService.listConsumerGroups).mockResolvedValue([]);
+
vi.mocked(consumerService.listConsumerGroupPage).mockResolvedValue(groupPage([]));
vi.mocked(consumerService.createConsumerGroup).mockImplementation(
async (data: Partial<ConsumerGroup>) => {
if (data.name === 'cg-fail') throw new Error('broker rejected group');
@@ -656,7 +695,7 @@ describe('Consumer page', () => {
renderWithProviders(<ConsumerPage />);
expect(await screen.findByText('选择实例')).toBeInTheDocument();
- expect(consumerService.listConsumerGroups).not.toHaveBeenCalled();
+ expect(consumerService.listConsumerGroupPage).not.toHaveBeenCalled();
expect(document.querySelector('.ant-spin-spinning')).toBeNull();
expect(screen.getByRole('button', { name: /导入/ })).toBeDisabled();
expect(screen.getByRole('button', { name: '创建 Group' })).toBeDisabled();
diff --git a/web/src/pages/instance/consumer.tsx
b/web/src/pages/instance/consumer.tsx
index ef7c038e6..8e27f3335 100644
--- a/web/src/pages/instance/consumer.tsx
+++ b/web/src/pages/instance/consumer.tsx
@@ -77,7 +77,7 @@ import {
getConsumerProgress,
getConsumerStack,
getConsumerSubscriptions,
- listConsumerGroups,
+ listConsumerGroupPage,
resetConsumerOffset,
} from '../../services/consumerService';
import { useInstanceFilter } from '../../hooks/useInstanceFilter';
@@ -176,11 +176,14 @@ const ConsumerPageContent = ({
selectedInstance?.vendor === 'ALIYUN' || selectedInstance?.vendor ===
'TENCENT';
const hasSelectedInstance = Boolean(selectedInstanceId);
const [groups, setGroups] = useState<ConsumerGroup[]>([]);
+ const [totalGroups, setTotalGroups] = useState(0);
const [loading, setLoading] = useState(hasSelectedInstance);
const [submitting, setSubmitting] = useState(false);
const [resetSubmitting, setResetSubmitting] = useState(false);
const [selectedRowKeys, setSelectedRowKeys] = useState<React.Key[]>([]);
const [search, setSearch] = useState('');
+ const [page, setPage] = useState(1);
+ const [pageSize, setPageSize] = useState(20);
const [modeFilter, setModeFilter] = useState<string>('ALL');
const [sortKey, setSortKey] = useState<string>('name_asc');
const [modalOpen, setModalOpen] = useState(false);
@@ -225,9 +228,17 @@ const ConsumerPageContent = ({
const requestId = ++groupRequestIdRef.current;
const timer = window.setTimeout(() => {
setLoading(true);
- void listConsumerGroups({ instanceId: selectedInstanceId })
- .then((nextGroups) => {
- if (requestId === groupRequestIdRef.current) setGroups(nextGroups);
+ void listConsumerGroupPage({
+ instanceId: selectedInstanceId,
+ search: search.trim() || undefined,
+ page,
+ pageSize,
+ })
+ .then((result) => {
+ if (requestId === groupRequestIdRef.current) {
+ setGroups(result.items);
+ setTotalGroups(result.total);
+ }
})
.catch(() => {
if (requestId === groupRequestIdRef.current)
message.error(t('consumer.fetchListFailed'));
@@ -239,7 +250,7 @@ const ConsumerPageContent = ({
return () => {
window.clearTimeout(timer);
};
- }, [t, selectedInstanceId]);
+ }, [t, selectedInstanceId, search, page, pageSize]);
const loadSubscriptions = useCallback(
async (groupName: string, force = false) => {
@@ -279,13 +290,7 @@ const ConsumerPageContent = ({
/* ─── Filtered & sorted data ─── */
const filtered = useMemo(() => {
- let data = groups.filter(
- (g) => g.name.includes(search) || (g.subscribedTopics ?? []).some((t) =>
t.includes(search)),
- );
-
- if (selectedInstanceId) {
- data = data.filter((g) => g.instanceId === selectedInstanceId);
- }
+ let data = groups;
if (modeFilter !== 'ALL') {
data = data.filter((g) => g.subscriptionMode === modeFilter);
@@ -298,7 +303,7 @@ const ConsumerPageContent = ({
}
return data;
- }, [groups, search, modeFilter, sortKey, selectedInstanceId]);
+ }, [groups, modeFilter, sortKey]);
/* ─── Open detail modal ─── */
const openModal = (group: ConsumerGroup) => {
@@ -802,7 +807,7 @@ const ConsumerPageContent = ({
{/* ─── Header ─── */}
<PageHeader
title={t('group.title')}
- subtitle={`管理消费者组订阅关系与消费进度,共 ${groups.length} 个 Group`}
+ subtitle={`管理消费者组订阅关系与消费进度,共 ${totalGroups} 个 Group`}
/>
{/* ─── Filter Bar ─── */}
@@ -818,8 +823,14 @@ const ConsumerPageContent = ({
placeholder="搜索 Group 名称或 Topic"
allowClear
value={search}
- onChange={(e) => setSearch(e.target.value)}
- onSearch={setSearch}
+ onChange={(e) => {
+ setSearch(e.target.value);
+ setPage(1);
+ }}
+ onSearch={(value) => {
+ setSearch(value);
+ setPage(1);
+ }}
style={{ width: 320 }}
prefix={<MagnifyingGlass size={14} color="#9CA3AF" />}
/>
@@ -933,9 +944,16 @@ const ConsumerPageContent = ({
onChange: (keys) => setSelectedRowKeys(keys),
}}
pagination={{
- pageSize: 20,
+ current: page,
+ pageSize,
+ total: totalGroups,
showSizeChanger: true,
showTotal: (total) => `共 ${total} 个 Group`,
+ pageSizeOptions: [10, 20, 50, 100],
+ onChange: (nextPage, nextPageSize) => {
+ setPage(nextPage);
+ setPageSize(nextPageSize);
+ },
}}
size="small"
expandable={{
diff --git a/web/src/pages/studio/Producer.tsx
b/web/src/pages/studio/Producer.tsx
index 9ac04261a..50c514656 100644
--- a/web/src/pages/studio/Producer.tsx
+++ b/web/src/pages/studio/Producer.tsx
@@ -70,6 +70,8 @@ const PRODUCER_CONNECTION_EXPORT_COLUMNS:
CsvColumn<ProducerConnectionExportRow>
{ header: 'Version', value: (connection) => connection.versionDesc },
];
+const PRODUCER_GROUP_SELECTOR_LIMIT = 20;
+
const ProducerPage = () => {
const [form] = Form.useForm();
const [topicList, setTopicList] = useState<string[]>([]);
@@ -85,8 +87,12 @@ const ProducerPage = () => {
const { message } = App.useApp();
const fetchTopicFailedMessage = t('producer.fetchTopicFailed');
const queryRequestIdRef = useRef(0);
+
const queryInFlightRef = useRef<number | null>(null);
+ const producerGroupRequestIdRef = useRef(0);
+ const selectedTopic = Form.useWatch('selectedTopic', form);
+
useEffect(() => {
let cancelled = false;
@@ -153,31 +159,33 @@ const ProducerPage = () => {
};
}, [fetchTopicFailedMessage, form, message, selectedInstanceId]);
- useEffect(() => {
- let cancelled = false;
+ const handleTopicChange = () => {
+ producerGroupRequestIdRef.current += 1;
+ setProducerGroups([]);
+ form.setFieldValue('producerGroup', undefined);
+ };
- if (!selectedInstanceId) {
- return () => {
- cancelled = true;
- };
+ const loadProducerGroups = async (query = '') => {
+ if (!selectedInstanceId || !selectedTopic) {
+ setProducerGroups([]);
+ return;
}
-
- void fetchProducerGroups(selectedInstanceId)
- .then((groups) => {
- if (!cancelled) {
- setProducerGroups(groups);
- }
- })
- .catch(() => {
- if (!cancelled) {
- setProducerGroups([]);
- }
+ const requestId = ++producerGroupRequestIdRef.current;
+ try {
+ const groups = await fetchProducerGroups(selectedInstanceId, {
+ topic: selectedTopic,
+ query,
+ limit: PRODUCER_GROUP_SELECTOR_LIMIT,
});
-
- return () => {
- cancelled = true;
- };
- }, [selectedInstanceId]);
+ if (requestId === producerGroupRequestIdRef.current) {
+ setProducerGroups(groups);
+ }
+ } catch {
+ if (requestId === producerGroupRequestIdRef.current) {
+ setProducerGroups([]);
+ }
+ }
+ };
const onFinish = async (values: { selectedTopic: string; producerGroup:
string }) => {
if (queryInFlightRef.current !== null) return;
@@ -318,6 +326,7 @@ const ProducerPage = () => {
placeholder={t('producer.selectTopic')}
style={{ width: 300 }}
optionFilterProp="label"
+ onChange={handleTopicChange}
options={topicList.map((topic) => ({ value: topic, label: topic
}))}
/>
</Form.Item>
@@ -331,6 +340,12 @@ const ProducerPage = () => {
placeholder={t('producer.inputGroup')}
style={{ width: 300 }}
options={producerGroups.map((group) => ({ value: group }))}
+ onFocus={() => {
+ void loadProducerGroups();
+ }}
+ onSearch={(value) => {
+ void loadProducerGroups(value);
+ }}
filterOption={(inputValue, option) =>
option?.value.toLowerCase().includes(inputValue.toLowerCase())
?? false
}
diff --git a/web/src/pages/studio/__tests__/Producer.test.tsx
b/web/src/pages/studio/__tests__/Producer.test.tsx
index 1167ad269..88255db61 100644
--- a/web/src/pages/studio/__tests__/Producer.test.tsx
+++ b/web/src/pages/studio/__tests__/Producer.test.tsx
@@ -111,6 +111,7 @@ describe('ProducerPage', () => {
await waitFor(() => {
expect(fetchTopicList).toHaveBeenCalledWith('instance-1');
});
+ expect(fetchProducerGroups).not.toHaveBeenCalled();
});
it('uses an Apache instance rather than a cloud instance for producer
diagnostics', async () => {
@@ -159,14 +160,27 @@ describe('ProducerPage', () => {
expect(await screen.findByRole('option', { name: 'payment-events'
})).toBeInTheDocument();
});
- it('suggests active producer groups while keeping free-form input', async ()
=> {
+ it('discovers producer groups from selector interactions after topic
selection', async () => {
const user = userEvent.setup();
renderWithProviders(<ProducerPage />);
- await waitFor(() => expect(fetchProducerGroups).toHaveBeenCalledTimes(1));
- const groupInput = screen.getAllByRole('combobox')[2];
+ await waitFor(() => expect(fetchTopicList).toHaveBeenCalledTimes(1));
+ const [, topicSelect, groupInput] = screen.getAllByRole('combobox');
+ fireEvent.mouseDown(topicSelect.parentElement!);
+ await user.click(
+ await screen.findByText('order-events', { selector:
'.ant-select-item-option-content' }),
+ );
+
+ await user.click(groupInput);
await user.type(groupInput, 'payment');
+ await waitFor(() => {
+ expect(fetchProducerGroups).toHaveBeenLastCalledWith('instance-1', {
+ topic: 'order-events',
+ query: 'payment',
+ limit: 20,
+ });
+ });
expect(await screen.findByRole('option', { name: 'pg-payment'
})).toBeInTheDocument();
});
@@ -349,6 +363,7 @@ describe('ProducerPage', () => {
await user.click(
await screen.findByText('order-events', { selector:
'.ant-select-item-option-content' }),
);
+ await user.click(groupInput);
await user.type(groupInput, 'manual-producer');
await user.click(screen.getByRole('button', { name: /搜索/ }));
@@ -406,6 +421,7 @@ describe('ProducerPage', () => {
await user.click(
await screen.findByText('order-events', { selector:
'.ant-select-item-option-content' }),
);
+ await user.click(groupInput);
await user.type(groupInput, 'order-producer');
await user.click(screen.getByRole('button', { name: /搜索/ }));
expect(await screen.findByText('producer-1')).toBeInTheDocument();
@@ -484,4 +500,16 @@ describe('ProducerPage', () => {
resolveQuery?.(producerResult([]));
await waitFor(() => expect(search).not.toHaveClass('ant-btn-loading'));
});
+
+ it('does not discover producer groups before a topic is selected', async ()
=> {
+ const user = userEvent.setup();
+ renderWithProviders(<ProducerPage />);
+
+ await waitFor(() => expect(fetchTopicList).toHaveBeenCalledTimes(1));
+ const groupInput = screen.getAllByRole('combobox')[2];
+ await user.click(groupInput);
+ await user.type(groupInput, 'order');
+
+ expect(fetchProducerGroups).not.toHaveBeenCalled();
+ });
});
diff --git a/web/src/services/consumerService.test.ts
b/web/src/services/consumerService.test.ts
index e5d71e0ee..b40872f1d 100644
--- a/web/src/services/consumerService.test.ts
+++ b/web/src/services/consumerService.test.ts
@@ -22,6 +22,7 @@ import {
getConsumerProgress,
getConsumerStack,
getConsumerSubscriptions,
+ listConsumerGroupPage,
listConsumerGroups,
} from './consumerService';
@@ -75,6 +76,32 @@ describe('consumer service mock data', () => {
expect(blankSearchGroups).toHaveLength(allGroups.length);
});
+ it('returns paged mock consumer groups with the filtered total', async () =>
{
+ const page = await listConsumerGroupPage({
+ search: 'cg-order-notify',
+ page: 1,
+ pageSize: 1,
+ });
+
+ expect(page.items.map((group) => group.name)).toEqual(['cg-order-notify']);
+ expect(page.total).toBe(1);
+ expect(page.page).toBe(1);
+ expect(page.size).toBe(1);
+ });
+
+ it('returns an empty page when the one-based offset starts past the filtered
total', async () => {
+ const page = await listConsumerGroupPage({
+ search: 'cg-order-notify',
+ page: 2,
+ pageSize: 1,
+ });
+
+ expect(page.items).toEqual([]);
+ expect(page.total).toBe(1);
+ expect(page.page).toBe(2);
+ expect(page.size).toBe(1);
+ });
+
it('returns copied progress and subscription rows', async () => {
const firstProgress = await getConsumerProgress('cg-order-notify');
const firstSubscriptions = await
getConsumerSubscriptions('cg-order-notify');
diff --git a/web/src/services/consumerService.ts
b/web/src/services/consumerService.ts
index 23f347194..7ef632a68 100644
--- a/web/src/services/consumerService.ts
+++ b/web/src/services/consumerService.ts
@@ -2,9 +2,11 @@ import { isMockMode } from './dataMode';
import * as metadataApi from '../api/metadata';
import type {
ConsumerGroup,
+ ConsumerGroupPageQuery,
ConsumerGroupQuery,
ConsumerGroupDetail,
ConsumerStackTrace,
+ PageResult,
QueueProgress,
ResetConsumerOffsetRequest,
SubscriptionEntry,
@@ -45,19 +47,41 @@ const normalizeConsumerGroup = <T extends
ConsumerGroup>(group: T): T => ({
instances: group.instances ?? [],
});
+function filterConsumerGroups(params?: ConsumerGroupQuery): ConsumerGroup[] {
+ let result = [...consumerGroupsState];
+ if (params?.clusterId) result = result.filter((group) => group.clusterId ===
params.clusterId);
+ if (params?.search) {
+ const kw = params.search.trim().toLowerCase();
+ if (kw) result = result.filter((group) =>
group.name.toLowerCase().includes(kw));
+ }
+ return result;
+}
+
export async function listConsumerGroups(params?: ConsumerGroupQuery):
Promise<ConsumerGroup[]> {
if (isMockMode()) {
- let result = [...consumerGroupsState];
- if (params?.clusterId) result = result.filter((group) => group.clusterId
=== params.clusterId);
- if (params?.search) {
- const kw = params.search.trim().toLowerCase();
- if (kw) result = result.filter((g) => g.name.toLowerCase().includes(kw));
- }
- return result.map(copyConsumerGroup);
+ return filterConsumerGroups(params).map(copyConsumerGroup);
}
return (await
metadataApi.listConsumerGroups(params)).map(normalizeConsumerGroup);
}
+export async function listConsumerGroupPage(
+ params: ConsumerGroupPageQuery = {},
+): Promise<PageResult<ConsumerGroup>> {
+ if (isMockMode()) {
+ const page = params.page ?? 1;
+ const pageSize = params.pageSize ?? 20;
+ const groups = filterConsumerGroups(params);
+ const from = Math.min((page - 1) * pageSize, groups.length);
+ return {
+ items: groups.slice(from, from + pageSize).map(copyConsumerGroup),
+ total: groups.length,
+ page,
+ size: pageSize,
+ };
+ }
+ return metadataApi.listConsumerGroupPage(params);
+}
+
export async function getConsumerProgress(
name: string,
instanceId?: string,