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 7443227f1 feat(instance): capabilities, proxy access types and
paginated topic metadata inventory (#2307)
7443227f1 is described below
commit 7443227f1b7456b9ab8b8a63876ebb56e28e87af
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 18 21:13:28 2026 +0800
feat(instance): capabilities, proxy access types and paginated topic
metadata inventory (#2307)
* feat: distinguish Proxy deployment access types
Split the generic PROXY access type into PROXY_LOCAL (Proxy co-located
with the Broker process) and PROXY_CLUSTER (standalone Proxy cluster),
keeping PROXY as a legacy compatibility value: list filters expand the
legacy value to every Proxy subtype, Apache create/update requests
normalize it to PROXY_CLUSTER, and the dashboard reports the matching
cluster type. The instance and topic pages surface the explicit
deployment guidance for each subtype.
* feat(instance): add instance capabilities and paginated topic metadata
inventory
* fix(web): align DLQ group API test with paged envelope
---
deploy/mysql/upgrade-demo-instance.sql | 8 +-
docs/api-spec.md | 69 ++++++++++++--
.../studio/common/domain/enums/InstanceType.java | 14 ++-
.../InstanceCapabilitiesVO.java} | 18 +++-
.../studio/instance/InstanceCapabilityService.java | 50 ++++++++++
.../studio/instance/InstanceController.java | 8 ++
.../rocketmq/studio/instance/InstanceService.java | 3 +-
.../instance/MybatisPlusInstanceRepository.java | 15 ++-
.../studio/instance/topic/MetadataService.java | 18 ++++
.../studio/instance/topic/TopicController.java | 12 +++
.../InstanceCapability.java} | 15 ++-
.../rocketmq/studio/provider/InstanceProvider.java | 16 ++++
.../provider/alibaba/AliyunInstanceProvider.java | 12 +++
.../provider/apache/ApacheInstanceProvider.java | 20 ++++
.../studio/provider/apache/MetadataProvider.java | 17 ++++
.../provider/apache/RocketMQDashboardProvider.java | 8 +-
.../provider/apache/RocketMQMetadataProvider.java | 18 ++++
.../provider/tencent/TencentInstanceProvider.java | 12 +++
server/src/main/resources/db/schema.sql | 2 +-
.../studio/auth/AuthCorsIntegrationTest.java | 4 +
.../instance/InstanceCapabilityServiceTest.java | 104 +++++++++++++++++++++
.../studio/instance/InstanceControllerTest.java | 25 +++++
.../studio/instance/InstanceServiceTest.java | 18 +++-
.../MybatisPlusInstanceRepositoryTest.java | 32 +++++++
.../alibaba/AliyunInstanceProviderTest.java | 10 ++
.../apache/ApacheInstanceProviderTest.java | 12 +++
.../apache/RocketMQDashboardProviderTest.java | 23 +++++
.../tencent/TencentInstanceProviderTest.java | 12 ++-
web/src/api/dlq.test.ts | 11 ++-
web/src/api/instance.test.ts | 31 +++++-
web/src/api/instance.ts | 33 ++++++-
web/src/api/metadata.ts | 12 +++
web/src/layouts/MainLayout.test.tsx | 95 ++++++++++++++++++-
web/src/layouts/MainLayout.tsx | 73 +++++++++++++--
web/src/mock/instances.ts | 6 +-
web/src/pages/cluster/index.tsx | 3 +-
.../pages/home/__tests__/DashboardPage.test.tsx | 44 +++++++++
web/src/pages/home/dashboard.tsx | 4 +-
.../pages/instance/__tests__/InstancePage.test.tsx | 36 ++++++-
.../pages/instance/__tests__/TopicPage.test.tsx | 42 ++++++---
web/src/pages/instance/index.tsx | 27 ++++--
web/src/pages/instance/topic.tsx | 51 +++++++---
web/src/pages/studio/BrokerCluster.tsx | 7 +-
web/src/pages/studio/Producer.tsx | 11 ++-
web/src/pages/studio/__tests__/Producer.test.tsx | 33 +++++++
web/src/services/instanceService.test.ts | 33 ++++++-
web/src/services/instanceService.ts | 47 +++++++++-
web/src/services/topicService.ts | 41 ++++++--
48 files changed, 1105 insertions(+), 110 deletions(-)
diff --git a/deploy/mysql/upgrade-demo-instance.sql
b/deploy/mysql/upgrade-demo-instance.sql
index 1dc8e5857..493bdba7e 100644
--- a/deploy/mysql/upgrade-demo-instance.sql
+++ b/deploy/mysql/upgrade-demo-instance.sql
@@ -15,7 +15,7 @@ CREATE TABLE IF NOT EXISTS rmq_instance (
id VARCHAR(64) PRIMARY KEY,
name VARCHAR(128) NOT NULL,
remark VARCHAR(255),
- type VARCHAR(32) NOT NULL COMMENT 'PROXY/DIRECT',
+ type VARCHAR(32) NOT NULL COMMENT 'PROXY/PROXY_LOCAL/PROXY_CLUSTER/DIRECT',
endpoint VARCHAR(512) NOT NULL,
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
@@ -40,9 +40,9 @@ PREPARE stmt FROM @sql; EXECUTE stmt; DEALLOCATE PREPARE stmt;
INSERT IGNORE INTO rmq_instance (id, name, remark, type, endpoint) VALUES
('instance-direct-1', 'instance-direct-1', '直连实例 1,交易核心链路(NameServer 直连)',
'DIRECT', '10.0.1.11:9876'),
('instance-direct-2', 'instance-direct-2', '直连实例 2,风控与审计链路(NameServer 直连)',
'DIRECT', '10.0.1.12:9876'),
- ('instance-proxy-1', 'instance-proxy-1', 'Proxy 实例 1,电商交易主链路', 'PROXY',
'10.0.2.21:8080'),
- ('instance-proxy-2', 'instance-proxy-2', 'Proxy 实例 2,营销与会员链路', 'PROXY',
'10.0.2.22:8080'),
- ('instance-proxy-3', 'instance-proxy-3', 'Proxy 实例 3,物流与大数据链路', 'PROXY',
'10.0.2.23:8080');
+ ('instance-proxy-1', 'instance-proxy-1', 'Proxy 实例 1,电商交易主链路',
'PROXY_CLUSTER', '10.0.2.21:8080'),
+ ('instance-proxy-2', 'instance-proxy-2', 'Proxy 实例 2,营销与会员链路',
'PROXY_CLUSTER', '10.0.2.22:8080'),
+ ('instance-proxy-3', 'instance-proxy-3', 'Proxy 实例 3,物流与大数据链路',
'PROXY_CLUSTER', '10.0.2.23:8080');
-- 4. 旧种子数据回填 instance_id(旧部署里这 9 个 topic、8 个 group 已存在,INSERT IGNORE 不会更新它们)
UPDATE rmq_topic SET instance_id = 'instance-proxy-1'
diff --git a/docs/api-spec.md b/docs/api-spec.md
index aa632e915..d54d33fb5 100644
--- a/docs/api-spec.md
+++ b/docs/api-spec.md
@@ -121,6 +121,8 @@
| 78 | GET | `/api/metrics/grafana/dashboards/:uid` | Grafana 看板 JSON 模型 |
| 79 | GET | `/api/metrics/grafana/dashboards/:uid/export` | 导出单个 Grafana 看板
JSON |
| 80 | GET | `/api/metrics/grafana/dashboards/export` | 打包导出全部 Grafana 看板 |
+| 81 | GET | `/api/instances/:instanceId/capabilities` | 实例能力契约 |
+| 82 | GET | `/api/topics/page` | Topic 分页列表 |
## 通用响应格式
@@ -296,7 +298,7 @@ GET /api/instances?type={type}&search={keyword}
| 参数 | 类型 | 必填 | 说明 |
|------|------|------|------|
-| `type` | `string` | 否 | 按类型过滤: `PROXY` / `DIRECT` |
+| `type` | `string` | 否 | 按类型过滤: `PROXY`(全部 Proxy)/ `PROXY_LOCAL` /
`PROXY_CLUSTER` / `DIRECT` |
| `search` | `string` | 否 | 按名称或地址搜索 |
**Response `data`:** `Instance[]`
@@ -306,7 +308,7 @@ GET /api/instances?type={type}&search={keyword}
| `id` | `string` | 实例 ID |
| `name` | `string` | 实例名称 |
| `remark` | `string` | 备注 |
-| `type` | `string` | 接入类型: `PROXY` / `DIRECT` |
+| `type` | `string` | 接入类型: `PROXY`(兼容值)/ `PROXY_LOCAL` / `PROXY_CLUSTER` /
`DIRECT` |
| `endpoint` | `string` | 接入地址 |
| `topicCount` | `number` | Topic 数量 |
| `consumerGroupCount` | `number` | 消费组数量 |
@@ -324,7 +326,7 @@ POST /api/instances/create
| 字段 | 类型 | 必填 | 说明 |
|------|------|------|------|
| `name` | `string` | 是 | 实例名称 |
-| `type` | `string` | 是 | `PROXY` / `DIRECT` |
+| `type` | `string` | 是 | Apache 实例使用 `PROXY_LOCAL` / `PROXY_CLUSTER` /
`DIRECT`;旧 `PROXY` 请求归一为 `PROXY_CLUSTER` |
| `endpoint` | `string` | 是 | 接入地址 |
**Response `data`:** `Instance`
@@ -341,7 +343,7 @@ POST /api/instances/update
|------|------|------|------|
| `id` | `string` | 是 | 实例 ID |
| `name` | `string` | 是 | 实例名称 |
-| `type` | `string` | 是 | `PROXY` / `DIRECT` |
+| `type` | `string` | 是 | `PROXY_LOCAL` / `PROXY_CLUSTER` / `DIRECT`;旧 `PROXY`
请求归一为 `PROXY_CLUSTER` |
| `endpoint` | `string` | 是 | 接入地址 |
**Response `data`:** `Instance`
@@ -360,6 +362,27 @@ POST /api/instances/delete
**Response `data`:** `null`
+### 3.5 获取实例能力契约
+
+```
+GET /api/instances/{instanceId}/capabilities
+```
+
+**Path Parameters:**
+
+| 参数 | 类型 | 必填 | 说明 |
+|------|------|------|------|
+| `instanceId` | `string` | 是 | 实例 ID(全局唯一字符串) |
+
+**Response `data`:** `InstanceCapabilities`
+
+| 字段 | 类型 | 说明 |
+|------|------|------|
+| `instanceId` | `string` | 实例 ID |
+| `vendor` | `string` | 厂商: `APACHE` / `ALIYUN` / `TENCENT` |
+| `accessType` | `string` | 接入类型: `PROXY_LOCAL` / `PROXY_CLUSTER` / `DIRECT` 等
|
+| `capabilities` | `string[]` | 能力列表: `TOPIC_MANAGEMENT` /
`CONSUMER_GROUP_MANAGEMENT` / `MESSAGE_QUERY` / `MESSAGE_TRACE` /
`ACL_MANAGEMENT` / `DLQ_MANAGEMENT` |
+
---
## 4. 集群管理 Cluster / NameServer / Proxy
@@ -721,7 +744,33 @@ GET
/api/topics?clusterId={clusterId}&type={type}&search={keyword}
| `createdAt` | `string` | 创建时间 (ISO 8601) |
| `updatedAt` | `string` | 更新时间 (ISO 8601) |
-### 5.2 创建 Topic
+### 5.2 分页获取 Topic 列表
+
+```
+GET
/api/topics/page?instanceId={instanceId}&clusterId={clusterId}&type={type}&search={keyword}&page={page}&pageSize={pageSize}
+```
+
+**Query Parameters:**
+
+| 参数 | 类型 | 必填 | 说明 |
+|------|------|------|------|
+| `instanceId` | `string` | 否 | 实例 ID(全局唯一字符串) |
+| `clusterId` | `string` | 否 | 按集群过滤 |
+| `type` | `string` | 否 | 按类型过滤 |
+| `search` | `string` | 否 | 按名称搜索 |
+| `page` | `number` | 否 | 页码,默认 `1` |
+| `pageSize` | `number` | 否 | 每页条数,默认 `20`,最大 `100` |
+
+**Response `data`:** `PageResult<Topic>`
+
+| 字段 | 类型 | 说明 |
+|------|------|------|
+| `items` | `Topic[]` | 当前页数据,结构同 5.1 |
+| `total` | `number` | 总条数 |
+| `page` | `number` | 当前页码 |
+| `size` | `number` | 每页条数 |
+
+### 5.3 创建 Topic
```
POST /api/topics/create
@@ -742,7 +791,7 @@ POST /api/topics/create
**Response `data`:** `Topic`
-### 5.3 更新 Topic
+### 5.4 更新 Topic
```
POST /api/topics/update
@@ -763,7 +812,7 @@ POST /api/topics/update
**Response `data`:** `Topic`
-### 5.4 删除 Topic
+### 5.5 删除 Topic
```
POST /api/topics/delete
@@ -777,7 +826,7 @@ POST /api/topics/delete
**Response `data`:** `null`
-### 5.5 获取 Topic 路由信息
+### 5.6 获取 Topic 路由信息
```
GET /api/topics/:name/routes
@@ -793,7 +842,7 @@ GET /api/topics/:name/routes
| `readQueues` | `number` | 读队列数 |
| `perm` | `string` | 权限 |
-### 5.6 获取 Topic 消费者列表
+### 5.7 获取 Topic 消费者列表
```
GET /api/topics/:name/consumers
@@ -809,7 +858,7 @@ GET /api/topics/:name/consumers
| `consumeTps` | `number` | 消费 TPS |
| `diffTotal` | `number` | 堆积消息数 |
-### 5.7 发送消息到 Topic
+### 5.8 发送消息到 Topic
```
POST /api/topics/send
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
b/server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
index c529a6d40..4afd4c880 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
@@ -18,5 +18,17 @@
package org.apache.rocketmq.studio.common.domain.enums;
public enum InstanceType {
- PROXY, DIRECT
+ /** Legacy generic Proxy value retained for persisted and cloud-managed
instances. */
+ PROXY,
+ PROXY_LOCAL,
+ PROXY_CLUSTER,
+ DIRECT;
+
+ public boolean isProxy() {
+ return this != DIRECT;
+ }
+
+ public InstanceType normalizeApacheType() {
+ return this == PROXY ? PROXY_CLUSTER : this;
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceCapabilitiesVO.java
similarity index 56%
copy from
server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
copy to
server/src/main/java/org/apache/rocketmq/studio/instance/InstanceCapabilitiesVO.java
index c529a6d40..09dd28b06 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceCapabilitiesVO.java
@@ -14,9 +14,21 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+package org.apache.rocketmq.studio.instance;
-package org.apache.rocketmq.studio.common.domain.enums;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
-public enum InstanceType {
- PROXY, DIRECT
+import java.util.List;
+
+/**
+ * Capability contract for a single instance. {@code instanceId} is the
canonical,
+ * globally unique instance identifier (the instance name), never the numeric
primary key.
+ */
+public record InstanceCapabilitiesVO(
+ String instanceId,
+ InstanceVendor vendor,
+ InstanceType accessType,
+ List<InstanceCapability> capabilities) {
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceCapabilityService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceCapabilityService.java
new file mode 100644
index 000000000..7cc1ca045
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceCapabilityService.java
@@ -0,0 +1,50 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance;
+
+import lombok.RequiredArgsConstructor;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
+import org.apache.rocketmq.studio.provider.InstanceProvider;
+import org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
+import org.springframework.stereotype.Service;
+
+import java.util.Comparator;
+import java.util.List;
+
+/**
+ * Resolves the capability contract for an instance by delegating to the
vendor provider.
+ */
+@Service
+@RequiredArgsConstructor
+public class InstanceCapabilityService {
+
+ private final InstanceRepository instanceRepository;
+ private final InstanceProviderRegistry providerRegistry;
+
+ public InstanceCapabilitiesVO getCapabilities(Long instanceId) {
+ InstanceVO instance = instanceRepository.findById(instanceId)
+ .orElseThrow(() -> new BusinessException(404, "Instance not
found: " + instanceId));
+ InstanceVendor vendor = instance.getVendor() == null ?
InstanceVendor.APACHE : instance.getVendor();
+ InstanceProvider provider = providerRegistry.forVendor(vendor);
+ List<InstanceCapability> capabilities =
provider.capabilities().stream()
+ .sorted(Comparator.comparingInt(InstanceCapability::ordinal))
+ .toList();
+ return new InstanceCapabilitiesVO(instance.getName(), vendor,
instance.getType(), capabilities);
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceController.java
index 1e8e123a6..d3dbb0ed2 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceController.java
@@ -23,6 +23,7 @@ import
org.apache.rocketmq.studio.common.domain.enums.InstanceType;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
@@ -37,6 +38,7 @@ import java.util.List;
public class InstanceController {
private final InstanceService instanceService;
+ private final InstanceCapabilityService instanceCapabilityService;
@GetMapping
public Result<List<InstanceVO>> listInstances(
@@ -45,6 +47,12 @@ public class InstanceController {
return Result.ok(instanceService.listInstances(type, search));
}
+ @GetMapping("/{instanceId}/capabilities")
+ public Result<InstanceCapabilitiesVO> getCapabilities(@PathVariable String
instanceId) {
+ return Result.ok(instanceCapabilityService.getCapabilities(
+ instanceService.resolveInstanceId(instanceId)));
+ }
+
@PostMapping("/create")
public Result<InstanceVO> createInstance(@Valid @RequestBody
CreateInstanceDTO request) {
return
Result.ok(instanceService.createInstance(request.toInstanceVO()));
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
index 3298b3f1c..e3d6002d9 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
@@ -155,6 +155,7 @@ public class InstanceService {
if (instance.getType() == null) {
throw new BusinessException(400, "InstanceVO type is required");
}
+ instance.setType(instance.getType().normalizeApacheType());
}
/**
@@ -270,7 +271,7 @@ public class InstanceService {
}
if (!cloudInstance) {
if (instance.getType() != null) {
- updated.setType(instance.getType());
+ updated.setType(instance.getType().normalizeApacheType());
}
if (instance.getEndpoint() != null) {
updated.setEndpoint(requireValidEndpoint(instance.getEndpoint()));
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
index b0c3efe85..0d49534ae 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
@@ -33,6 +33,7 @@ import
org.springframework.transaction.annotation.Transactional;
import java.util.List;
import java.util.Optional;
+import java.util.stream.Stream;
@RequiredArgsConstructor
@Repository
@@ -54,7 +55,7 @@ public class MybatisPlusInstanceRepository implements
InstanceRepository {
public List<InstanceVO> findByType(InstanceType type) {
return instanceMapper.selectList(
new QueryWrapper<RmqInstance>()
- .eq("type", type.name())
+ .in("type", typeNamesForFilter(type))
.orderByAsc("id")).stream()
.map(this::toVO)
.toList();
@@ -76,7 +77,7 @@ public class MybatisPlusInstanceRepository implements
InstanceRepository {
public List<InstanceVO> findByTypeAndSearch(InstanceType type, String
keyword) {
return instanceMapper.selectList(
new QueryWrapper<RmqInstance>()
- .eq("type", type.name())
+ .in("type", typeNamesForFilter(type))
.and(w -> w.like("name", keyword)
.or().like("endpoint", keyword)
.or().like("remark", keyword))
@@ -85,6 +86,16 @@ public class MybatisPlusInstanceRepository implements
InstanceRepository {
.toList();
}
+ private List<String> typeNamesForFilter(InstanceType type) {
+ if (type == InstanceType.PROXY) {
+ return Stream.of(InstanceType.values())
+ .filter(InstanceType::isProxy)
+ .map(Enum::name)
+ .toList();
+ }
+ return List.of(type.name());
+ }
+
@Override
public Optional<InstanceVO> findById(Long id) {
if (id == null) {
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 44e47406f..4d05c81ab 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.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.springframework.util.StringUtils;
@@ -75,6 +76,23 @@ public class MetadataService {
return resolve(instanceId).listTopics(instanceId,
normalizeFilter(type), normalizeFilter(search));
}
+ public PageResult<TopicVO> listTopicsPage(String instanceId, String
clusterId, String type,
+ String search, int page, int pageSize) {
+ if (page < 1) {
+ throw new BusinessException(400, "page must be greater than zero");
+ }
+ if (pageSize < 1 || pageSize > 100) {
+ throw new BusinessException(400, "pageSize must be between 1 and
100");
+ }
+ instanceId = normalizeInstanceId(instanceId);
+ if (!StringUtils.hasText(instanceId) &&
StringUtils.hasText(clusterId)) {
+ return metadataProvider.listTopicsPage(normalizeFilter(clusterId),
+ normalizeFilter(type), normalizeFilter(search), page,
pageSize);
+ }
+ return resolve(instanceId).listTopicsPage(instanceId,
normalizeFilter(type),
+ normalizeFilter(search), page, pageSize);
+ }
+
public TopicVO createTopic(TopicVO topic) {
requireTopic(topic);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
index ad04e281c..eab11d963 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.instance.topic;
import org.apache.rocketmq.studio.common.domain.Result;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
@@ -47,6 +48,17 @@ public class TopicController {
return Result.ok(metadataService.listTopics(instanceId, clusterId,
type, search));
}
+ @GetMapping("/page")
+ public Result<PageResult<TopicVO>> listTopicsPage(
+ @RequestParam(required = false) String instanceId,
+ @RequestParam(required = false) String clusterId,
+ @RequestParam(required = false) String type,
+ @RequestParam(required = false) String search,
+ @RequestParam(defaultValue = "1") int page,
+ @RequestParam(defaultValue = "20") int pageSize) {
+ return Result.ok(metadataService.listTopicsPage(instanceId, clusterId,
type, search, page, pageSize));
+ }
+
@PostMapping("/create")
public Result<TopicVO> createTopic(@Valid @RequestBody(required = false)
CreateTopicDTO topic) {
requireCreateTopicRequest(topic);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceCapability.java
similarity index 73%
copy from
server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
copy to
server/src/main/java/org/apache/rocketmq/studio/provider/InstanceCapability.java
index c529a6d40..30fed36d2 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/enums/InstanceType.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceCapability.java
@@ -14,9 +14,16 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+package org.apache.rocketmq.studio.provider;
-package org.apache.rocketmq.studio.common.domain.enums;
-
-public enum InstanceType {
- PROXY, DIRECT
+/**
+ * Stable, user-facing capabilities implemented by an instance provider.
+ */
+public enum InstanceCapability {
+ TOPIC_MANAGEMENT,
+ CONSUMER_GROUP_MANAGEMENT,
+ MESSAGE_QUERY,
+ MESSAGE_TRACE,
+ ACL_MANAGEMENT,
+ DLQ_MANAGEMENT
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
index 43e64627d..13ee1133a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.studio.provider;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.util.Pagination;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
@@ -28,6 +29,7 @@ import
org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import java.util.List;
+import java.util.Set;
/**
* Unified instance-scoped operations SPI. Every method takes the Studio
instance id as its
@@ -38,12 +40,26 @@ public interface InstanceProvider {
InstanceVendor vendor();
+ default Set<InstanceCapability> capabilities() {
+ return Set.of();
+ }
+
int countTopics(String instanceId);
int countGroups(String instanceId);
List<TopicVO> listTopics(String instanceId, String type, String search);
+ default PageResult<TopicVO> listTopicsPage(String instanceId, String type,
String search,
+ int page, int pageSize) {
+ List<TopicVO> topics = listTopics(instanceId, type, search);
+ 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);
+ }
+
TopicVO createTopic(String instanceId, TopicVO topic);
TopicVO updateTopic(String instanceId, TopicVO 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 dc1812664..c2a95721d 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
@@ -57,12 +57,14 @@ import
org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.springframework.stereotype.Component;
import lombok.RequiredArgsConstructor;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import java.util.Set;
/**
* Aliyun RocketMQ 5.x implementation of the instance-scoped operations SPI,
backed by the
@@ -90,6 +92,16 @@ public class AliyunInstanceProvider implements
InstanceProvider {
return InstanceVendor.ALIYUN;
}
+ @Override
+ public Set<InstanceCapability> capabilities() {
+ return Set.of(
+ InstanceCapability.TOPIC_MANAGEMENT,
+ InstanceCapability.CONSUMER_GROUP_MANAGEMENT,
+ InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.MESSAGE_TRACE,
+ InstanceCapability.ACL_MANAGEMENT);
+ }
+
@Override
public int countTopics(String instanceId) {
Context ctx = resolve(instanceId);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
index 8eae9b305..6c4147ecb 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.studio.provider.apache;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.instance.InstanceRepository;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
@@ -28,11 +29,13 @@ import
org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.springframework.stereotype.Component;
import lombok.RequiredArgsConstructor;
import java.util.List;
import java.util.Objects;
+import java.util.Set;
/**
* Open-source Apache RocketMQ implementation: pure delegation to the existing
admin-client
@@ -52,6 +55,17 @@ public class ApacheInstanceProvider implements
InstanceProvider {
return InstanceVendor.APACHE;
}
+ @Override
+ public Set<InstanceCapability> capabilities() {
+ return Set.of(
+ InstanceCapability.TOPIC_MANAGEMENT,
+ InstanceCapability.CONSUMER_GROUP_MANAGEMENT,
+ InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.MESSAGE_TRACE,
+ InstanceCapability.ACL_MANAGEMENT,
+ InstanceCapability.DLQ_MANAGEMENT);
+ }
+
@Override
public int countTopics(String instanceId) {
// The instance_id foreign key is numeric; resolve the external
identifier first.
@@ -74,6 +88,12 @@ public class ApacheInstanceProvider implements
InstanceProvider {
.toList();
}
+ @Override
+ public PageResult<TopicVO> listTopicsPage(String instanceId, String type,
String search,
+ int page, int pageSize) {
+ return metadataProvider.listTopicsPage(instanceId, null, type, search,
page, pageSize);
+ }
+
@Override
public TopicVO createTopic(String instanceId, TopicVO topic) {
return adminClient.createTopic(topic);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
index b565055dd..181267c8f 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.studio.provider.apache;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.util.Pagination;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerPageVO;
@@ -29,6 +30,22 @@ import java.util.List;
public interface MetadataProvider {
List<TopicVO> listTopics(String clusterId, String type, String search);
+
+ default PageResult<TopicVO> listTopicsPage(String clusterId, String type,
String search,
+ int page, int pageSize) {
+ List<TopicVO> topics = listTopics(clusterId, type, search);
+ 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);
+ }
+
+ default PageResult<TopicVO> listTopicsPage(String instanceId, String
clusterId, String type,
+ String search, int page, int pageSize) {
+ return listTopicsPage(clusterId, type, search, page, pageSize);
+ }
+
List<ConsumerGroupVO> listConsumerGroups(String clusterId, String search);
default List<ConsumerGroupVO> listConsumerGroups(String instanceId, String
clusterId, String search) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProvider.java
index 2ba491506..dc94c2f75 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProvider.java
@@ -34,7 +34,6 @@ import
org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.ClusterStatus;
import org.apache.rocketmq.studio.common.domain.enums.ClusterType;
-import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
import org.apache.rocketmq.studio.instance.InstanceVO;
import org.apache.rocketmq.studio.ops.dashboard.ClusterOverviewVO;
import org.apache.rocketmq.studio.ops.dashboard.DashboardDataVO;
@@ -320,8 +319,11 @@ public class RocketMQDashboardProvider implements
DashboardProvider {
}
private ClusterType clusterTypeFor(InstanceVO instance) {
- return instance.getType() == InstanceType.DIRECT
- ? ClusterType.V4_DIRECT : ClusterType.V5_PROXY_CLUSTER;
+ return switch (instance.getType()) {
+ case DIRECT -> ClusterType.V4_DIRECT;
+ case PROXY_LOCAL -> ClusterType.V5_PROXY_LOCAL;
+ case PROXY, PROXY_CLUSTER -> ClusterType.V5_PROXY_CLUSTER;
+ };
}
private DashboardDataVO emptyDashboard() {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
index 4882beeed..1b6986c2c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.provider.apache;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.common.MixAll;
import org.apache.rocketmq.common.TopicConfig;
@@ -33,6 +34,7 @@ import
org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.tools.admin.MQAdminExt;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
import org.apache.rocketmq.studio.common.domain.enums.SubscriptionMode;
@@ -130,6 +132,22 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
return result;
}
+ @Override
+ public PageResult<TopicVO> listTopicsPage(String instanceId, String
clusterId, String type,
+ String search, int page, int pageSize) {
+ LambdaQueryWrapper<RmqTopic> query = new LambdaQueryWrapper<RmqTopic>()
+ .eq(StringUtils.hasText(instanceId), RmqTopic::getInstanceId,
instanceId)
+ .eq(StringUtils.hasText(clusterId), RmqTopic::getClusterId,
clusterId)
+ .eq(StringUtils.hasText(type), RmqTopic::getTopicType, type)
+ .like(StringUtils.hasText(search), RmqTopic::getName, search)
+ .notLikeRight(RmqTopic::getName, "RMQ_SYS_")
+ .notLikeRight(RmqTopic::getName, "rmq_sys_")
+ .orderByAsc(RmqTopic::getName, RmqTopic::getId);
+ Page<RmqTopic> result = topicMapper.selectPage(new Page<>(page,
pageSize), query);
+ return
PageResult.of(result.getRecords().stream().map(this::toTopicVO).toList(),
+ result.getTotal(), page, pageSize);
+ }
+
private TopicVO toTopicVO(RmqTopic entity) {
TopicVO vo = new TopicVO();
vo.setId(entity.getId());
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 7d2529d80..d073606b3 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
@@ -61,6 +61,7 @@ import
org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import lombok.RequiredArgsConstructor;
@@ -79,6 +80,7 @@ import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
+import java.util.Set;
import java.util.Map;
import java.util.UUID;
@@ -123,6 +125,16 @@ public class TencentInstanceProvider implements
InstanceProvider {
return InstanceVendor.TENCENT;
}
+ @Override
+ public Set<InstanceCapability> capabilities() {
+ return Set.of(
+ InstanceCapability.TOPIC_MANAGEMENT,
+ InstanceCapability.CONSUMER_GROUP_MANAGEMENT,
+ InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.MESSAGE_TRACE,
+ InstanceCapability.ACL_MANAGEMENT);
+ }
+
@Override
public int countTopics(String instanceId) {
Context context = resolve(instanceId);
diff --git a/server/src/main/resources/db/schema.sql
b/server/src/main/resources/db/schema.sql
index 0a04de3bf..9132df72a 100644
--- a/server/src/main/resources/db/schema.sql
+++ b/server/src/main/resources/db/schema.sql
@@ -63,7 +63,7 @@ CREATE TABLE IF NOT EXISTS rmq_instance (
`gmt_modified` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE
CURRENT_TIMESTAMP COMMENT '修改时间',
name VARCHAR(128) NOT NULL,
remark VARCHAR(255),
- type VARCHAR(32) NOT NULL COMMENT 'PROXY/DIRECT',
+ type VARCHAR(32) NOT NULL COMMENT 'PROXY/PROXY_LOCAL/PROXY_CLUSTER/DIRECT',
endpoint VARCHAR(512) NOT NULL,
vendor VARCHAR(32),
cloud_instance_id VARCHAR(128),
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/auth/AuthCorsIntegrationTest.java
b/server/src/test/java/org/apache/rocketmq/studio/auth/AuthCorsIntegrationTest.java
index 844b3432e..b5fe66234 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/auth/AuthCorsIntegrationTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/auth/AuthCorsIntegrationTest.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.auth;
import org.apache.rocketmq.studio.common.config.CorsConfig;
import org.apache.rocketmq.studio.instance.InstanceController;
+import org.apache.rocketmq.studio.instance.InstanceCapabilityService;
import org.apache.rocketmq.studio.instance.InstanceService;
import org.apache.rocketmq.studio.settings.SettingsRepository;
import org.junit.jupiter.api.BeforeEach;
@@ -60,6 +61,9 @@ class AuthCorsIntegrationTest {
@MockBean
private InstanceService instanceService;
+ @MockBean
+ private InstanceCapabilityService instanceCapabilityService;
+
@MockBean
private AuthProperties authProperties;
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceCapabilityServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceCapabilityServiceTest.java
new file mode 100644
index 000000000..428580186
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceCapabilityServiceTest.java
@@ -0,0 +1,104 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance;
+
+import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
+import org.apache.rocketmq.studio.provider.InstanceProvider;
+import org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import java.util.Optional;
+import java.util.Set;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.mockito.Mockito.when;
+
+@ExtendWith(MockitoExtension.class)
+class InstanceCapabilityServiceTest {
+
+ @Mock
+ private InstanceRepository instanceRepository;
+ @Mock
+ private InstanceProviderRegistry providerRegistry;
+ @Mock
+ private InstanceProvider instanceProvider;
+
+ private InstanceCapabilityService service;
+
+ @BeforeEach
+ void setUp() {
+ service = new InstanceCapabilityService(instanceRepository,
providerRegistry);
+ }
+
+ @Test
+ void getCapabilitiesShouldReturnProviderCapabilitiesInStableOrderTest() {
+ InstanceVO instance = InstanceVO.builder()
+ .name("cloud-1")
+ .vendor(InstanceVendor.ALIYUN)
+ .type(InstanceType.PROXY)
+ .build();
+ instance.setId(1L);
+
when(instanceRepository.findById(1L)).thenReturn(Optional.of(instance));
+
when(providerRegistry.forVendor(InstanceVendor.ALIYUN)).thenReturn(instanceProvider);
+ when(instanceProvider.capabilities()).thenReturn(Set.of(
+ InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.TOPIC_MANAGEMENT));
+
+ InstanceCapabilitiesVO result = service.getCapabilities(1L);
+
+ assertThat(result.instanceId()).isEqualTo("cloud-1");
+ assertThat(result.vendor()).isEqualTo(InstanceVendor.ALIYUN);
+ assertThat(result.accessType()).isEqualTo(InstanceType.PROXY);
+ assertThat(result.capabilities()).containsExactly(
+ InstanceCapability.TOPIC_MANAGEMENT,
+ InstanceCapability.MESSAGE_QUERY);
+ }
+
+ @Test
+ void getCapabilitiesShouldDefaultLegacyNullVendorToApacheTest() {
+ InstanceVO instance = InstanceVO.builder()
+ .name("direct-1")
+ .type(InstanceType.DIRECT)
+ .build();
+ instance.setId(2L);
+
when(instanceRepository.findById(2L)).thenReturn(Optional.of(instance));
+
when(providerRegistry.forVendor(InstanceVendor.APACHE)).thenReturn(instanceProvider);
+
when(instanceProvider.capabilities()).thenReturn(Set.of(InstanceCapability.DLQ_MANAGEMENT));
+
+ InstanceCapabilitiesVO result = service.getCapabilities(2L);
+
+ assertThat(result.instanceId()).isEqualTo("direct-1");
+ assertThat(result.vendor()).isEqualTo(InstanceVendor.APACHE);
+ }
+
+ @Test
+ void getCapabilitiesShouldRejectUnknownInstanceTest() {
+ when(instanceRepository.findById(999L)).thenReturn(Optional.empty());
+
+ assertThatThrownBy(() -> service.getCapabilities(999L))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(404));
+ }
+}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceControllerTest.java
index 8a9c8507c..2c87ce067 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceControllerTest.java
@@ -19,7 +19,9 @@ package org.apache.rocketmq.studio.instance;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import
org.springframework.boot.test.autoconfigure.web.servlet.AutoConfigureMockMvc;
@@ -55,6 +57,9 @@ class InstanceControllerTest {
@MockBean
private InstanceService instanceService;
+ @MockBean
+ private InstanceCapabilityService instanceCapabilityService;
+
@Autowired
private ObjectMapper objectMapper;
@@ -103,6 +108,26 @@ class InstanceControllerTest {
verify(instanceService).listInstances(isNull(), eq("prod"));
}
+ @Test
+ void getCapabilitiesShouldResolveStringInstanceIdAndReturnContractTest()
throws Exception {
+ when(instanceService.resolveInstanceId("instance-1")).thenReturn(1L);
+ when(instanceCapabilityService.getCapabilities(1L)).thenReturn(new
InstanceCapabilitiesVO(
+ "instance-1",
+ InstanceVendor.APACHE,
+ InstanceType.DIRECT,
+ List.of(InstanceCapability.TOPIC_MANAGEMENT,
InstanceCapability.DLQ_MANAGEMENT)));
+
+ mockMvc.perform(get("/api/instances/instance-1/capabilities"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.instanceId").value("instance-1"))
+ .andExpect(jsonPath("$.data.vendor").value("APACHE"))
+ .andExpect(jsonPath("$.data.accessType").value("DIRECT"))
+
.andExpect(jsonPath("$.data.capabilities[0]").value("TOPIC_MANAGEMENT"))
+
.andExpect(jsonPath("$.data.capabilities[1]").value("DLQ_MANAGEMENT"));
+
+ verify(instanceService).resolveInstanceId("instance-1");
+ }
+
@Test
void createInstanceShouldReturnCreatedInstance() throws Exception {
InstanceVO input = InstanceVO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
index ea23aefda..2e6f93ce4 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
@@ -290,9 +290,10 @@ class InstanceServiceTest {
assertThat(result.getGmtCreate()).isNotNull();
assertThat(result.getGmtModified()).isNotNull();
assertThat(result.getName()).isEqualTo("new-instance");
+ assertThat(result.getType()).isEqualTo(InstanceType.PROXY_CLUSTER);
verify(instanceRepository).save(any(InstanceVO.class));
verify(operationAuditService).record(eq("CREATE_INSTANCE"),
eq("INSTANCE"), eq("1"), eq(null),
- argThat(detail -> detail.equals("name=new-instance,
vendor=APACHE, type=PROXY")
+ argThat(detail -> detail.equals("name=new-instance,
vendor=APACHE, type=PROXY_CLUSTER")
&& !detail.contains("10.0.1.1:8080")),
eq("SUCCESS"), eq(null));
}
@@ -798,9 +799,24 @@ class InstanceServiceTest {
InstanceVO created = instanceService.createInstance(instance);
assertThat(created.getVendor()).isEqualTo(InstanceVendor.APACHE);
+ assertThat(created.getType()).isEqualTo(InstanceType.PROXY_CLUSTER);
verifyNoInteractions(cloudCredentialRepository, providerRegistry);
}
+ @Test
+ void createApacheInstanceShouldPreserveExplicitProxyLocalTypeTest() {
+ InstanceVO instance = InstanceVO.builder()
+ .name("local-proxy")
+ .endpoint("broker-proxy:8080")
+ .type(InstanceType.PROXY_LOCAL)
+ .build();
+
when(instanceRepository.save(any(InstanceVO.class))).thenAnswer(invocation ->
invocation.getArgument(0));
+
+ InstanceVO created = instanceService.createInstance(instance);
+
+ assertThat(created.getType()).isEqualTo(InstanceType.PROXY_LOCAL);
+ }
+
@Test
void createApacheInstanceShouldRequireTypeTest() {
InstanceVO instance = InstanceVO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
index 1dde42bd1..727689dc2 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
@@ -80,6 +80,38 @@ class MybatisPlusInstanceRepositoryTest {
verifyNoInteractions(topicMapper, groupMapper);
}
+ @Test
+ void findByTypeShouldTreatLegacyProxyAsAllProxyAccessTypes() {
+ when(instanceMapper.selectList(any(QueryWrapper.class)))
+ .thenReturn(List.of(
+ entity(3L, "legacy", InstanceType.PROXY),
+ entity(4L, "local", InstanceType.PROXY_LOCAL),
+ entity(5L, "cluster", InstanceType.PROXY_CLUSTER)));
+
+ List<InstanceVO> result = repository.findByType(InstanceType.PROXY);
+
+ assertThat(result).extracting(InstanceVO::getType)
+ .containsExactly(InstanceType.PROXY, InstanceType.PROXY_LOCAL,
InstanceType.PROXY_CLUSTER);
+ ArgumentCaptor<QueryWrapper<RmqInstance>> query =
ArgumentCaptor.forClass(QueryWrapper.class);
+ verify(instanceMapper).selectList(query.capture());
+ assertThat(query.getValue().getSqlSegment()).contains("type IN");
+ assertThat(query.getValue().getParamNameValuePairs().values())
+ .contains(InstanceType.PROXY.name(),
InstanceType.PROXY_LOCAL.name(), InstanceType.PROXY_CLUSTER.name());
+ }
+
+ @Test
+ void findByTypeShouldFilterExplicitProxyLocalTypeOnly() {
+ when(instanceMapper.selectList(any(QueryWrapper.class)))
+ .thenReturn(List.of(entity(4L, "local",
InstanceType.PROXY_LOCAL)));
+
+
assertThat(repository.findByType(InstanceType.PROXY_LOCAL)).singleElement()
+ .extracting(InstanceVO::getType)
+ .isEqualTo(InstanceType.PROXY_LOCAL);
+ ArgumentCaptor<QueryWrapper<RmqInstance>> query =
ArgumentCaptor.forClass(QueryWrapper.class);
+ verify(instanceMapper).selectList(query.capture());
+ assertThat(query.getValue().getSqlSegment()).contains("type IN");
+ }
+
@Test
void constructorShouldNotSeedDemoInstances() {
when(instanceMapper.selectList(any(QueryWrapper.class))).thenReturn(List.of());
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 a14744602..2d46f3e91 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
@@ -48,6 +48,7 @@ import
org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@@ -96,6 +97,15 @@ class AliyunInstanceProviderTest {
provider = new AliyunInstanceProvider(clientFactory,
instanceRepository);
}
+ @Test
+ void capabilitiesShouldExcludeUnsupportedDlqOperationsTest() {
+ assertThat(provider.capabilities())
+ .contains(InstanceCapability.TOPIC_MANAGEMENT,
+ InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.ACL_MANAGEMENT)
+ .doesNotContain(InstanceCapability.DLQ_MANAGEMENT);
+ }
+
@Test
void listTopicsShouldMapMessageTypeAndFilterTest() {
stubInstance();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
index 48fd02cd1..b99657915 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
@@ -20,6 +20,7 @@ import
org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.instance.InstanceRepository;
import org.apache.rocketmq.studio.instance.InstanceVO;
import org.apache.rocketmq.studio.instance.message.MessageProvider;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
@@ -50,6 +51,17 @@ class ApacheInstanceProviderTest {
@InjectMocks
private ApacheInstanceProvider provider;
+ @Test
+ void capabilitiesShouldIncludeApacheOnlyOperationsTest() {
+ assertThat(provider.capabilities()).contains(
+ InstanceCapability.TOPIC_MANAGEMENT,
+ InstanceCapability.CONSUMER_GROUP_MANAGEMENT,
+ InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.MESSAGE_TRACE,
+ InstanceCapability.ACL_MANAGEMENT,
+ InstanceCapability.DLQ_MANAGEMENT);
+ }
+
@Test
void vendorShouldBeApacheTest() {
assertThat(provider.vendor()).isEqualTo(InstanceVendor.APACHE);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProviderTest.java
index 035d79fcd..b26499d69 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProviderTest.java
@@ -378,6 +378,29 @@ class RocketMQDashboardProviderTest {
verify(resolver).execute(eq(instance), any());
}
+ @Test
+ void dashboardShouldReportSelectedProxyLocalInstanceType() throws
Exception {
+ DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+ RuntimeAdminClientResolver resolver =
mock(RuntimeAdminClientResolver.class);
+ InstanceVO instance = InstanceVO.builder()
+ .type(InstanceType.PROXY_LOCAL)
+ .endpoint("local-proxy:8080")
+ .build();
+ instance.setId(1L);
+ when(resolver.resolveInstance("instance-local")).thenReturn(instance);
+ when(resolver.execute(eq(instance), any())).thenAnswer(invocation ->
+
invocation.<MqAdminExtFactory.AdminAction<DashboardDataVO>>getArgument(1).apply(adminExt));
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
+ when(adminExt.fetchAllTopicList()).thenReturn(topicList());
+
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtimeStats());
+
+ DashboardDataVO dashboard = newProvider(adminExt,
resolver).getDashboardData("instance-local");
+
+ assertThat(dashboard.getClusters()).singleElement()
+ .extracting(cluster -> cluster.getType())
+ .isEqualTo(ClusterType.V5_PROXY_LOCAL);
+ }
+
@Test
void dashboardShouldMarkClusterWarningWhenBrokerRuntimeStatsFail() throws
Exception {
DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
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 dfc2cf4c5..3f63461e3 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
@@ -56,6 +56,7 @@ import
org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
+import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@@ -100,7 +101,7 @@ class TencentInstanceProviderTest {
@BeforeEach
void setUp() {
provider = new TencentInstanceProvider(clientFactory,
instanceRepository);
-
when(instanceRepository.findByIdentifier(STUDIO_INSTANCE_ID)).thenReturn(Optional.of(InstanceVO.builder()
+
lenient().when(instanceRepository.findByIdentifier(STUDIO_INSTANCE_ID)).thenReturn(Optional.of(InstanceVO.builder()
.name("tencent-prod")
.vendor(InstanceVendor.TENCENT)
.cloudInstanceId(CLOUD_INSTANCE_ID)
@@ -113,6 +114,15 @@ class TencentInstanceProviderTest {
});
}
+ @Test
+ void capabilitiesShouldExcludeUnsupportedDlqOperationsTest() {
+ assertThat(provider.capabilities())
+ .contains(InstanceCapability.TOPIC_MANAGEMENT,
+ InstanceCapability.MESSAGE_QUERY,
+ InstanceCapability.ACL_MANAGEMENT)
+ .doesNotContain(InstanceCapability.DLQ_MANAGEMENT);
+ }
+
@Test
void countTopicsShouldClampOversizedTotals() throws Exception {
DescribeTopicListResponse response = new DescribeTopicListResponse();
diff --git a/web/src/api/dlq.test.ts b/web/src/api/dlq.test.ts
index c77e97b66..796063a5a 100644
--- a/web/src/api/dlq.test.ts
+++ b/web/src/api/dlq.test.ts
@@ -43,11 +43,14 @@ describe('DLQ API', () => {
});
it('loads and unwraps DLQ groups', async () => {
- mock
- .onGet('/dlq', { params: { instanceId: 'instance-1' } })
- .reply(200, { code: 200, data: [group] });
+ const pageData = { items: [group], total: 1, page: 1, size: 20 };
+ mock.onGet('/dlq').reply((config) =>
+ config.params?.instanceId === 'instance-1'
+ ? [200, { code: 200, data: pageData }]
+ : [404, {}],
+ );
- await expect(listDLQGroups('instance-1')).resolves.toEqual([group]);
+ await expect(listDLQGroups('instance-1')).resolves.toEqual(pageData);
});
it('sends epoch milliseconds for the resend time range', async () => {
diff --git a/web/src/api/instance.test.ts b/web/src/api/instance.test.ts
index 152a05185..7998d6dc7 100644
--- a/web/src/api/instance.test.ts
+++ b/web/src/api/instance.test.ts
@@ -18,7 +18,14 @@
import MockAdapter from 'axios-mock-adapter';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import client from './client';
-import { createInstance, deleteInstance, listInstances, updateInstance } from
'./instance';
+import {
+ createInstance,
+ deleteInstance,
+ getInstanceCapabilities,
+ listInstances,
+ supportsApacheRuntime,
+ updateInstance,
+} from './instance';
const mock = new MockAdapter(client);
const instance = {
@@ -86,4 +93,26 @@ describe('instance API', () => {
await expect(deleteInstance(instance.name)).resolves.toBeUndefined();
});
+
+ it('loads the capability contract for an encoded instance id', async () => {
+ const capabilities = {
+ instanceId: 'instance/proxy',
+ vendor: 'APACHE' as const,
+ accessType: 'PROXY' as const,
+ capabilities: ['TOPIC_MANAGEMENT', 'DLQ_MANAGEMENT'] as const,
+ };
+ mock.onGet('/instances/instance%2Fproxy/capabilities').reply(200, {
+ code: 200,
+ data: capabilities,
+ });
+
+ await
expect(getInstanceCapabilities('instance/proxy')).resolves.toEqual(capabilities);
+ });
+
+ it('identifies instances supported by Apache MQAdmin runtime APIs', () => {
+ expect(supportsApacheRuntime({ vendor: 'APACHE' })).toBe(true);
+ expect(supportsApacheRuntime({})).toBe(true);
+ expect(supportsApacheRuntime({ vendor: 'ALIYUN' })).toBe(false);
+ expect(supportsApacheRuntime({ vendor: 'TENCENT' })).toBe(false);
+ });
});
diff --git a/web/src/api/instance.ts b/web/src/api/instance.ts
index e04482d2f..d61e03c1c 100644
--- a/web/src/api/instance.ts
+++ b/web/src/api/instance.ts
@@ -19,12 +19,20 @@ import client from './client';
// ─── Types ──────────────────────────────────────────────────────
export type InstanceVendor = 'APACHE' | 'ALIYUN' | 'TENCENT';
+export type InstanceType = 'PROXY' | 'PROXY_LOCAL' | 'PROXY_CLUSTER' |
'DIRECT';
+export type InstanceCapability =
+ | 'TOPIC_MANAGEMENT'
+ | 'CONSUMER_GROUP_MANAGEMENT'
+ | 'MESSAGE_QUERY'
+ | 'MESSAGE_TRACE'
+ | 'ACL_MANAGEMENT'
+ | 'DLQ_MANAGEMENT';
export interface Instance {
id: number;
name: string;
remark: string | null;
- type: 'PROXY' | 'DIRECT';
+ type: InstanceType;
endpoint: string;
vendor?: InstanceVendor;
cloudInstanceId?: string;
@@ -40,7 +48,7 @@ export interface Instance {
export interface CreateInstanceRequest {
name?: string;
- type?: 'PROXY' | 'DIRECT';
+ type?: InstanceType;
endpoint?: string;
remark?: string;
vendor?: InstanceVendor;
@@ -53,7 +61,7 @@ export interface CreateInstanceRequest {
export interface UpdateInstanceRequest {
instanceId: string;
name?: string;
- type?: 'PROXY' | 'DIRECT';
+ type?: InstanceType;
endpoint?: string;
remark?: string;
adminCredentialRef?: string;
@@ -64,6 +72,18 @@ export interface InstanceQuery {
search?: string;
}
+/** Whether an instance can use Apache MQAdmin-backed runtime diagnostics. */
+export function supportsApacheRuntime(instance: Pick<Instance, 'vendor'>):
boolean {
+ return instance.vendor === undefined || instance.vendor === 'APACHE';
+}
+
+export interface InstanceCapabilities {
+ instanceId: string;
+ vendor: InstanceVendor;
+ accessType: Instance['type'];
+ capabilities: InstanceCapability[];
+}
+
// ─── Instance CRUD ──────────────────────────────────────────────
export async function listInstances(query: InstanceQuery = {}) {
const search = query.search?.trim();
@@ -75,6 +95,13 @@ export async function listInstances(query: InstanceQuery =
{}) {
return res.data.data;
}
+export async function getInstanceCapabilities(instanceId: string) {
+ const res = await client.get<{ data: InstanceCapabilities }>(
+ `/instances/${encodeURIComponent(instanceId)}/capabilities`,
+ );
+ return res.data.data;
+}
+
export async function createInstance(data: CreateInstanceRequest) {
const res = await client.post<{ data: Instance }>('/instances/create', data);
return res.data.data;
diff --git a/web/src/api/metadata.ts b/web/src/api/metadata.ts
index bcaaad51f..3e5d77bf5 100644
--- a/web/src/api/metadata.ts
+++ b/web/src/api/metadata.ts
@@ -138,6 +138,18 @@ export async function listTopics(params?: TopicQuery) {
return res.data.data;
}
+export interface TopicPage {
+ items: Topic[];
+ total: number;
+ page: number;
+ size: number;
+}
+
+export async function listTopicsPage(params?: TopicQuery & { page?: number;
pageSize?: number }) {
+ const res = await client.get<{ data: TopicPage }>('/topics/page', { params
});
+ return res.data.data;
+}
+
export async function createTopic(data: Partial<Topic>) {
const res = await client.post<{ data: Topic }>('/topics/create', data);
return res.data.data;
diff --git a/web/src/layouts/MainLayout.test.tsx
b/web/src/layouts/MainLayout.test.tsx
index 22c868281..3a9b3b7dc 100644
--- a/web/src/layouts/MainLayout.test.tsx
+++ b/web/src/layouts/MainLayout.test.tsx
@@ -24,7 +24,12 @@ import useAuthStore from '../stores/authStore';
import { ThemeProvider } from '../theme/ThemeProvider';
import MainLayout from './MainLayout';
+const instanceServiceMocks = vi.hoisted(() => ({
+ getInstanceCapabilities: vi.fn(),
+}));
+
vi.mock('../api/auth', () => ({ logout: vi.fn() }));
+vi.mock('../services/instanceService', () => instanceServiceMocks);
vi.mock('antd', async () => {
const React = await import('react');
@@ -38,12 +43,27 @@ vi.mock('antd', async () => {
return {
Layout,
- Menu: () => null,
+ Menu: ({
+ items,
+ }: {
+ items?: Array<{ key: string; label: React.ReactNode; children?:
unknown[] }>;
+ }) => {
+ const renderItems = (entries: typeof items): React.ReactNode[] =>
+ (entries ?? []).flatMap((item) => [
+ React.createElement('span', { key: `${item.key}-label` },
item.label),
+ ...renderItems(item.children as typeof items),
+ ]);
+ return React.createElement('nav', null, renderItems(items));
+ },
Breadcrumb: ({ items }: { items: Array<{ key: string; title:
React.ReactNode }> }) =>
React.createElement(
- 'nav',
+ 'div',
null,
- items.map((item) => React.createElement(React.Fragment, { key:
item.key }, item.title)),
+ React.createElement(
+ 'nav',
+ null,
+ items.map((item) => React.createElement(React.Fragment, { key:
item.key }, item.title)),
+ ),
),
Avatar: () => React.createElement('span', null, 'avatar'),
Dropdown: ({ children, menu }: { children?: React.ReactNode; menu:
DropdownMenu }) =>
@@ -74,6 +94,19 @@ describe('MainLayout authentication navigation', () => {
localStorage.clear();
vi.mocked(logout).mockReset().mockResolvedValue(undefined);
useAuthStore.getState().login('admin', 7, true);
+
instanceServiceMocks.getInstanceCapabilities.mockReset().mockResolvedValue({
+ instanceId: 'apache-1',
+ vendor: 'APACHE',
+ accessType: 'DIRECT',
+ capabilities: [
+ 'TOPIC_MANAGEMENT',
+ 'CONSUMER_GROUP_MANAGEMENT',
+ 'MESSAGE_QUERY',
+ 'MESSAGE_TRACE',
+ 'ACL_MANAGEMENT',
+ 'DLQ_MANAGEMENT',
+ ],
+ });
});
it('replaces a protected route with the login page after logout', async ()
=> {
@@ -137,4 +170,60 @@ describe('MainLayout authentication navigation', () => {
fireEvent.click(screen.getByRole('button', { name: '切换到英语' }));
expect(screen.getByRole('button', { name: 'Switch to Chinese'
})).toBeInTheDocument();
});
+
+ it('hides unsupported instance navigation after capabilities load', async ()
=> {
+ instanceServiceMocks.getInstanceCapabilities.mockResolvedValue({
+ instanceId: 'cloud-1',
+ vendor: 'ALIYUN',
+ accessType: 'PROXY',
+ capabilities: ['TOPIC_MANAGEMENT'],
+ });
+
+ render(
+ <LangProvider>
+ <ThemeProvider>
+ <MemoryRouter initialEntries={['/instance/cloud-1/topic']}>
+ <Routes>
+ <Route path="/" element={<MainLayout />}>
+ <Route path="instance/:instanceId/topic" element={<div>cloud
topic</div>} />
+ </Route>
+ </Routes>
+ </MemoryRouter>
+ </ThemeProvider>
+ </LangProvider>,
+ );
+
+ await waitFor(() =>
+
expect(instanceServiceMocks.getInstanceCapabilities).toHaveBeenCalledWith('cloud-1'),
+ );
+ await waitFor(() => expect(screen.getByText('Topic
管理')).toBeInTheDocument());
+ expect(screen.queryByText('死信队列')).not.toBeInTheDocument();
+ expect(screen.queryByText('ACL 管理')).not.toBeInTheDocument();
+ expect(screen.queryByText('Group 管理')).not.toBeInTheDocument();
+ expect(screen.queryByText('消息查询')).not.toBeInTheDocument();
+ });
+
+ it('keeps navigation available when capability discovery fails', async () =>
{
+ instanceServiceMocks.getInstanceCapabilities.mockRejectedValue(new
Error('unavailable'));
+
+ render(
+ <LangProvider>
+ <ThemeProvider>
+ <MemoryRouter initialEntries={['/instance/apache-1/topic']}>
+ <Routes>
+ <Route path="/" element={<MainLayout />}>
+ <Route path="instance/:instanceId/topic" element={<div>apache
topic</div>} />
+ </Route>
+ </Routes>
+ </MemoryRouter>
+ </ThemeProvider>
+ </LangProvider>,
+ );
+
+ await waitFor(() =>
+
expect(instanceServiceMocks.getInstanceCapabilities).toHaveBeenCalledWith('apache-1'),
+ );
+ expect(screen.getByText('ACL 管理')).toBeInTheDocument();
+ expect(screen.getByText('死信队列')).toBeInTheDocument();
+ });
});
diff --git a/web/src/layouts/MainLayout.tsx b/web/src/layouts/MainLayout.tsx
index e79a09107..f03c548e0 100644
--- a/web/src/layouts/MainLayout.tsx
+++ b/web/src/layouts/MainLayout.tsx
@@ -51,11 +51,20 @@ import {
type NavigationSearchEntry,
} from './navigationSearch';
import { useDataModeStore } from '../stores/dataModeStore';
+import { getInstanceCapabilities } from '../services/instanceService';
+import type { InstanceCapability } from '../api/instance';
const { Sider, Content } = Layout;
const iconSize = 18;
+function hasInstanceCapability(
+ capabilities: Set<InstanceCapability> | null,
+ capability: InstanceCapability,
+) {
+ return capabilities === null || capabilities.has(capability);
+}
+
const MainLayout = () => {
const navigate = useNavigate();
const location = useLocation();
@@ -69,6 +78,10 @@ const MainLayout = () => {
const admin = useAuthStore((state) => state.admin);
const useMock = useDataModeStore((state) => state.useMock);
const toggleDataMode = useDataModeStore((state) => state.toggle);
+ const [capabilityState, setCapabilityState] = useState<{
+ instanceId: string;
+ capabilities: Set<InstanceCapability>;
+ } | null>(null);
// Pages fetch on mount, so reload to re-request everything from the new
data source.
const handleDataModeToggle = () => {
@@ -112,6 +125,40 @@ const MainLayout = () => {
() =>
location.pathname.match(/^\/instance\/[^/]+\/(topic|consumer|message|acl|dlq)$/),
[location.pathname],
);
+ const selectedInstanceId = useMemo(() => {
+ const match = location.pathname.match(/^\/instance\/([^/]+)\//);
+ if (!match) return null;
+ try {
+ const value = decodeURIComponent(match[1]);
+ return value || null;
+ } catch {
+ return null;
+ }
+ }, [location.pathname]);
+
+ useEffect(() => {
+ if (!selectedInstanceId) return;
+ let active = true;
+ void getInstanceCapabilities(selectedInstanceId)
+ .then((result) => {
+ if (active) {
+ setCapabilityState({
+ instanceId: selectedInstanceId,
+ capabilities: new Set(result.capabilities),
+ });
+ }
+ })
+ .catch(() => {
+ // Preserve existing navigation when capability discovery is
unavailable.
+ });
+ return () => {
+ active = false;
+ };
+ }, [selectedInstanceId]);
+
+ const instanceCapabilities =
+ capabilityState?.instanceId === selectedInstanceId ?
capabilityState.capabilities : null;
+
const selectedMenuKey = instanceScopedMatch
? `/instance/${instanceScopedMatch[1]}`
: location.pathname;
@@ -125,15 +172,21 @@ const MainLayout = () => {
label: t('nav.instance'),
children: [
{ key: '/instance', icon: <Database size={16} />, label:
t('nav.instanceList') },
- { key: '/instance/topic', icon: <ListDashes size={16} />, label:
t('nav.topic') },
- { key: '/instance/consumer', icon: <ChatCircleText size={16} />,
label: t('nav.group') },
- { key: '/instance/acl', icon: <Key size={16} />, label: t('nav.acl')
},
- {
- key: '/instance/message',
- icon: <MagnifyingGlass size={16} />,
- label: t('nav.message'),
- },
- { key: '/instance/dlq', icon: <TrashSimple size={16} />, label:
t('nav.dlq') },
+ ...(hasInstanceCapability(instanceCapabilities, 'TOPIC_MANAGEMENT')
+ ? [{ key: '/instance/topic', icon: <ListDashes size={16} />,
label: t('nav.topic') }]
+ : []),
+ ...(hasInstanceCapability(instanceCapabilities,
'CONSUMER_GROUP_MANAGEMENT')
+ ? [{ key: '/instance/consumer', icon: <ChatCircleText size={16}
/>, label: t('nav.group') }]
+ : []),
+ ...(hasInstanceCapability(instanceCapabilities, 'ACL_MANAGEMENT')
+ ? [{ key: '/instance/acl', icon: <Key size={16} />, label:
t('nav.acl') }]
+ : []),
+ ...(hasInstanceCapability(instanceCapabilities, 'MESSAGE_QUERY')
+ ? [{ key: '/instance/message', icon: <MagnifyingGlass size={16}
/>, label: t('nav.message') }]
+ : []),
+ ...(hasInstanceCapability(instanceCapabilities, 'DLQ_MANAGEMENT')
+ ? [{ key: '/instance/dlq', icon: <TrashSimple size={16} />, label:
t('nav.dlq') }]
+ : []),
{ key: '/instance/alerts', icon: <Warning size={16} />, label:
t('nav.alertRuleAssets') },
],
},
@@ -166,7 +219,7 @@ const MainLayout = () => {
label: t('nav.settings'),
},
],
- [t],
+ [t, instanceCapabilities],
);
const breadcrumbMap: Record<string, string> = useMemo(
diff --git a/web/src/mock/instances.ts b/web/src/mock/instances.ts
index bfb3dfebe..caf315d4f 100644
--- a/web/src/mock/instances.ts
+++ b/web/src/mock/instances.ts
@@ -44,7 +44,7 @@ export const mockInstances: Instance[] = [
id: 3,
name: 'instance-proxy-1',
remark: 'Proxy 实例 1,电商交易主链路',
- type: 'PROXY',
+ type: 'PROXY_CLUSTER',
endpoint: '10.0.2.21:8080',
topicCount: 8,
consumerGroupCount: 6,
@@ -55,7 +55,7 @@ export const mockInstances: Instance[] = [
id: 4,
name: 'instance-proxy-2',
remark: 'Proxy 实例 2,营销与会员链路',
- type: 'PROXY',
+ type: 'PROXY_CLUSTER',
endpoint: '10.0.2.22:8080',
topicCount: 6,
consumerGroupCount: 5,
@@ -66,7 +66,7 @@ export const mockInstances: Instance[] = [
id: 5,
name: 'instance-proxy-3',
remark: 'Proxy 实例 3,物流与大数据链路',
- type: 'PROXY',
+ type: 'PROXY_LOCAL',
endpoint: '10.0.2.23:8080',
topicCount: 6,
consumerGroupCount: 6,
diff --git a/web/src/pages/cluster/index.tsx b/web/src/pages/cluster/index.tsx
index d047366a0..1541cf277 100644
--- a/web/src/pages/cluster/index.tsx
+++ b/web/src/pages/cluster/index.tsx
@@ -72,6 +72,7 @@ import {
updateNameserverRegistry,
} from '../../services/clusterService';
import { listInstances } from '../../services/instanceService';
+import { supportsApacheRuntime } from '../../api/instance';
import { isMockMode } from '../../services/dataMode';
const { Text } = Typography;
@@ -317,7 +318,7 @@ const ClusterPage = () => {
void listInstances()
.then((nextInstances) => {
if (cancelled) return;
- const apacheInstances = nextInstances.filter((instance) =>
instance.vendor === 'APACHE');
+ const apacheInstances = nextInstances.filter(supportsApacheRuntime);
const initialInstanceId = apacheInstances.some(
(instance) => instance.name === requestedInstanceId,
)
diff --git a/web/src/pages/home/__tests__/DashboardPage.test.tsx
b/web/src/pages/home/__tests__/DashboardPage.test.tsx
index bb68d3094..c63045c15 100644
--- a/web/src/pages/home/__tests__/DashboardPage.test.tsx
+++ b/web/src/pages/home/__tests__/DashboardPage.test.tsx
@@ -191,4 +191,48 @@ describe('DashboardPage', () => {
expect(screen.getByTestId('location')).toHaveTextContent('/cluster?instanceId=instance-b');
});
+
+ it('does not offer cloud instances for MQAdmin runtime diagnostics', async
() => {
+ vi.mocked(instanceService.listInstances).mockResolvedValue([
+ {
+ id: 1,
+ name: 'apache-instance',
+ endpoint: 'apache:9876',
+ type: 'DIRECT',
+ vendor: 'APACHE',
+ remark: '',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ gmtCreate: '',
+ gmtModified: '',
+ },
+ {
+ id: 2,
+ name: 'cloud-instance',
+ endpoint: 'cloud:9876',
+ type: 'DIRECT',
+ vendor: 'ALIYUN',
+ remark: '',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ gmtCreate: '',
+ gmtModified: '',
+ },
+ ]);
+
vi.mocked(dashboardService.getDashboard).mockResolvedValue(dashboard('apache-cluster'));
+ const user = userEvent.setup();
+ renderWithProviders(<DashboardPage />);
+
+ await screen.findByText('apache-cluster');
+ await user.click(screen.getByRole('combobox', { name: 'Dashboard instance'
}));
+
+ expect(
+ await screen.findByText('apache-instance', {
+ selector: '.ant-select-item-option-content',
+ }),
+ ).toBeInTheDocument();
+ expect(
+ screen.queryByText('cloud-instance', { selector:
'.ant-select-item-option-content' }),
+ ).not.toBeInTheDocument();
+ });
});
diff --git a/web/src/pages/home/dashboard.tsx b/web/src/pages/home/dashboard.tsx
index bcbc2aefc..2965e26e6 100644
--- a/web/src/pages/home/dashboard.tsx
+++ b/web/src/pages/home/dashboard.tsx
@@ -23,8 +23,8 @@ import MetricsExplorer from
'../../components/MetricsExplorer';
import { CLUSTER_TYPE_MAP } from '../../constants/theme';
import { getDashboard } from '../../services/dashboardService';
import type { DashboardData } from '../../api/metrics';
+import { supportsApacheRuntime, type Instance } from '../../api/instance';
import { listInstances } from '../../services/instanceService';
-import type { Instance } from '../../api/instance';
import { useLang } from '../../i18n/LangContext';
const { Text } = Typography;
@@ -68,7 +68,7 @@ const DashboardPage = () => {
let cancelled = false;
void listInstances()
.then((nextInstances) => {
- if (!cancelled) setInstances(nextInstances);
+ if (!cancelled)
setInstances(nextInstances.filter(supportsApacheRuntime));
})
.catch(() => {
if (!cancelled) setInstances([]);
diff --git a/web/src/pages/instance/__tests__/InstancePage.test.tsx
b/web/src/pages/instance/__tests__/InstancePage.test.tsx
index 1fbfb12b8..0558fd42e 100644
--- a/web/src/pages/instance/__tests__/InstancePage.test.tsx
+++ b/web/src/pages/instance/__tests__/InstancePage.test.tsx
@@ -227,7 +227,7 @@ describe('InstancePage', () => {
await user.type(within(dialog).getByLabelText('实例 ID'), 'new-proxy');
const createTypeSelect = within(dialog).getByRole('combobox');
fireEvent.mouseDown(createTypeSelect.parentElement!);
- const proxyOptions = await screen.findAllByText('Proxy 模式', {
+ const proxyOptions = await screen.findAllByText('Proxy Cluster 模式', {
selector: '.ant-select-item-option-content',
});
await user.click(proxyOptions[proxyOptions.length - 1]);
@@ -260,7 +260,7 @@ describe('InstancePage', () => {
await user.type(within(dialog).getByLabelText('实例 ID'), 'new-proxy');
const createTypeSelect = within(dialog).getByRole('combobox');
fireEvent.mouseDown(createTypeSelect.parentElement!);
- const proxyOptions = await screen.findAllByText('Proxy 模式', {
+ const proxyOptions = await screen.findAllByText('Proxy Cluster 模式', {
selector: '.ant-select-item-option-content',
});
await user.click(proxyOptions[proxyOptions.length - 1]);
@@ -270,13 +270,43 @@ describe('InstancePage', () => {
await waitFor(() =>
expect(instanceService.createInstance).toHaveBeenCalledWith({
name: 'new-proxy',
- type: 'PROXY',
+ type: 'PROXY_CLUSTER',
endpoint: 'proxy-new:8080',
}),
);
expect(instanceService.listInstances).toHaveBeenLastCalledWith({ type:
'DIRECT' });
});
+ it('creates an Apache instance with an explicit Proxy Local deployment
type', async () => {
+ const user = userEvent.setup();
+ vi.mocked(instanceService.createInstance).mockResolvedValue(
+ instance(3, 'local-proxy', 'PROXY_LOCAL'),
+ );
+ renderPage();
+
+ expect(await screen.findByText('production-proxy')).toBeInTheDocument();
+ await user.click(screen.getByRole('button', { name: /添加实例/ }));
+ const dialog = await screen.findByRole('dialog');
+ await user.type(within(dialog).getByLabelText('实例 ID'), 'local-proxy');
+ const createTypeSelect = within(dialog).getByRole('combobox');
+ fireEvent.mouseDown(createTypeSelect.parentElement!);
+ await user.click(
+ await screen.findByText('Proxy Local 模式', {
+ selector: '.ant-select-item-option-content',
+ }),
+ );
+ await user.type(within(dialog).getByLabelText('接入地址'),
'broker-proxy:8080');
+ await user.click(within(dialog).getByRole('button', { name: /连\s*接/ }));
+
+ await waitFor(() =>
+ expect(instanceService.createInstance).toHaveBeenCalledWith({
+ name: 'local-proxy',
+ type: 'PROXY_LOCAL',
+ endpoint: 'broker-proxy:8080',
+ }),
+ );
+ });
+
it('updates instance type and endpoint through the edit dialog', async () =>
{
const user = userEvent.setup();
vi.mocked(instanceService.updateInstance).mockResolvedValue(
diff --git a/web/src/pages/instance/__tests__/TopicPage.test.tsx
b/web/src/pages/instance/__tests__/TopicPage.test.tsx
index a20eb6c5f..d4562a259 100644
--- a/web/src/pages/instance/__tests__/TopicPage.test.tsx
+++ b/web/src/pages/instance/__tests__/TopicPage.test.tsx
@@ -32,6 +32,7 @@ const topicServiceMocks = vi.hoisted(() => ({
getTopicConsumerPage: vi.fn(),
getTopicRoutes: vi.fn(),
listTopics: vi.fn(),
+ listTopicsPage: vi.fn(),
sendTopicMessage: vi.fn(),
}));
@@ -40,6 +41,14 @@ const instanceServiceMocks = vi.hoisted(() => ({
}));
vi.mock('../../../services/topicService', () => topicServiceMocks);
+
+const mockTopicsList = (items: Topic[]) =>
+ topicServiceMocks.listTopicsPage.mockResolvedValue({
+ items,
+ total: items.length,
+ page: 1,
+ size: 20,
+ });
vi.mock('../../../services/instanceService', () => instanceServiceMocks);
beforeAll(() => {
@@ -121,7 +130,7 @@ const getTableBody = () => {
describe('TopicPage', () => {
beforeEach(() => {
- topicServiceMocks.listTopics.mockResolvedValue(buildTopics(25));
+ mockTopicsList(buildTopics(25));
topicServiceMocks.batchDeleteTopics.mockResolvedValue({ deleted: [],
failed: [] });
topicServiceMocks.createTopic.mockImplementation(async (data:
Partial<Topic>) => ({
...buildTopics(1)[0],
@@ -161,6 +170,17 @@ describe('TopicPage', () => {
vi.clearAllMocks();
});
+ it('shows the selected Proxy deployment type explicitly', async () => {
+ instanceServiceMocks.listInstances.mockResolvedValue([
+ { ...selectedInstance, type: 'PROXY_LOCAL' },
+ ]);
+
+ renderWithProviders('/instance/instance-proxy-1/topic');
+
+ expect(await screen.findByText(/Proxy Local 模式/)).toBeInTheDocument();
+ expect(screen.getByText(/与 Broker 同进程部署的 Proxy 地址/)).toBeInTheDocument();
+ });
+
it('ignores duplicate Topic creates while the first request is pending',
async () => {
topicServiceMocks.createTopic.mockImplementation(() => new Promise(() =>
{}));
const user = userEvent.setup();
@@ -187,7 +207,7 @@ describe('TopicPage', () => {
exportedBlob = blob as Blob;
return 'blob:topic-export';
});
- topicServiceMocks.listTopics.mockResolvedValue([
+ mockTopicsList([
{
...buildTopics(1)[0],
name: 'orders-topic',
@@ -237,7 +257,7 @@ describe('TopicPage', () => {
gmtModified: '2026-01-01T00:00:00Z',
},
]);
- topicServiceMocks.listTopics.mockResolvedValue(
+ mockTopicsList(
buildTopics(25).map((topic) => ({ ...topic, instanceId: 'instance-a' })),
);
renderWithProviders('/instance/instance-a/topic');
@@ -281,7 +301,7 @@ describe('TopicPage', () => {
type: 'DIRECT',
},
]);
- topicServiceMocks.listTopics.mockResolvedValue([topic]);
+ mockTopicsList([topic]);
renderWithProviders('/instance/instance-a/topic');
expect(await screen.findByText('topic-01')).toBeInTheDocument();
@@ -302,7 +322,7 @@ describe('TopicPage', () => {
it('keeps failed topics selected after a partially successful batch
deletion', async () => {
const user = userEvent.setup();
- topicServiceMocks.listTopics.mockResolvedValue(buildTopics(3));
+ mockTopicsList(buildTopics(3));
topicServiceMocks.batchDeleteTopics.mockResolvedValue({
deleted: ['topic-01', 'topic-03'],
failed: ['topic-02'],
@@ -325,7 +345,7 @@ describe('TopicPage', () => {
it('filters topics by the instance from the route and shows its endpoint',
async () => {
const base = buildTopics(1)[0];
- topicServiceMocks.listTopics.mockResolvedValue([
+ mockTopicsList([
{ ...base, name: 'topic-a', instanceId: 'instance-proxy-1' },
{ ...base, name: 'topic-b', instanceId: 'instance-proxy-2' },
]);
@@ -363,7 +383,7 @@ describe('TopicPage', () => {
it('imports valid topic CSV rows through the create service with the
selected instance', async () => {
const user = userEvent.setup();
- topicServiceMocks.listTopics.mockResolvedValue([]);
+ mockTopicsList([]);
instanceServiceMocks.listInstances.mockResolvedValue([selectedInstance]);
renderWithProviders('/instance/instance-proxy-1/topic');
@@ -394,7 +414,7 @@ describe('TopicPage', () => {
it('does not call createTopic when imported topic CSV is invalid or
duplicated', async () => {
const user = userEvent.setup();
instanceServiceMocks.listInstances.mockResolvedValue([selectedInstance]);
- topicServiceMocks.listTopics.mockResolvedValue([
+ mockTopicsList([
{ ...buildTopics(1)[0], instanceId: 'instance-proxy-1' },
]);
renderWithProviders('/instance/instance-proxy-1/topic');
@@ -417,7 +437,7 @@ describe('TopicPage', () => {
it('imports valid topic rows while skipping duplicate rows', async () => {
const user = userEvent.setup();
- topicServiceMocks.listTopics.mockResolvedValue([]);
+ mockTopicsList([]);
instanceServiceMocks.listInstances.mockResolvedValue([selectedInstance]);
renderWithProviders('/instance/instance-proxy-1/topic');
@@ -452,7 +472,7 @@ describe('TopicPage', () => {
renderWithProviders();
expect(await screen.findByText('共 0 个 Topic')).toBeInTheDocument();
- expect(topicServiceMocks.listTopics).not.toHaveBeenCalled();
+ expect(topicServiceMocks.listTopicsPage).not.toHaveBeenCalled();
await waitFor(() =>
expect(document.querySelector('.ant-spin-spinning')).toBeNull());
expect(screen.getByRole('button', { name: /导入/ })).toBeDisabled();
expect(screen.getByRole('button', { name: /创建 Topic/ })).toBeDisabled();
@@ -460,7 +480,7 @@ describe('TopicPage', () => {
it('renders unavailable Topic consumer metrics distinctly from zero', async
() => {
const user = userEvent.setup();
- topicServiceMocks.listTopics.mockResolvedValue([buildTopics(1)[0]]);
+ mockTopicsList([buildTopics(1)[0]]);
topicServiceMocks.getTopicConsumerPage.mockResolvedValue({
items: [
{
diff --git a/web/src/pages/instance/index.tsx b/web/src/pages/instance/index.tsx
index 2b776afbd..63fbd7643 100644
--- a/web/src/pages/instance/index.tsx
+++ b/web/src/pages/instance/index.tsx
@@ -65,7 +65,9 @@ const DEFAULT_CLOUD_REGION_IDS:
Partial<Record<InstanceVendor, string>> = {
/* ─── Helpers ─── */
const typeLabel: Record<string, { text: string; color: string }> = {
- PROXY: { text: 'Proxy 模式', color: 'blue' },
+ PROXY: { text: 'Proxy 模式(部署形态未标明)', color: 'blue' },
+ PROXY_LOCAL: { text: 'Proxy Local 模式', color: 'cyan' },
+ PROXY_CLUSTER: { text: 'Proxy Cluster 模式', color: 'blue' },
DIRECT: { text: 'Direct 模式', color: 'orange' },
};
@@ -109,7 +111,7 @@ const InstancePage = () => {
const [addModalOpen, setAddModalOpen] = useState(false);
const [vendor, setVendor] = useState<InstanceVendor>(DEFAULT_VENDOR);
const [addForm] = Form.useForm();
- const addInstanceType = Form.useWatch<'PROXY' | 'DIRECT' |
undefined>('type', addForm);
+ const addInstanceType = Form.useWatch<Instance['type'] | undefined>('type',
addForm);
const addCredentialId = Form.useWatch<number | undefined>('credentialId',
addForm);
const addRegionId = Form.useWatch<string | undefined>('regionId', addForm);
const [credentials, setCredentials] = useState<CloudCredential[]>([]);
@@ -121,7 +123,7 @@ const InstancePage = () => {
const [editModalOpen, setEditModalOpen] = useState(false);
const [editingInstance, setEditingInstance] = useState<Instance |
null>(null);
const [editForm] = Form.useForm();
- const editInstanceType = Form.useWatch<'PROXY' | 'DIRECT' |
undefined>('type', editForm);
+ const editInstanceType = Form.useWatch<Instance['type'] | undefined>('type',
editForm);
const [submitting, setSubmitting] = useState(false);
const requestIdRef = useRef(0);
const mutationInFlightRef = useRef(false);
@@ -567,7 +569,9 @@ const InstancePage = () => {
style={{ width: 140 }}
options={[
{ value: 'ALL', label: '全部架构' },
- { value: 'PROXY', label: 'Proxy 模式' },
+ { value: 'PROXY', label: '全部 Proxy 模式' },
+ { value: 'PROXY_LOCAL', label: 'Proxy Local 模式' },
+ { value: 'PROXY_CLUSTER', label: 'Proxy Cluster 模式' },
{ value: 'DIRECT', label: 'Direct 模式' },
]}
/>
@@ -740,7 +744,8 @@ const InstancePage = () => {
<Select
placeholder="选择接入方式"
options={[
- { value: 'PROXY', label: 'Proxy 模式' },
+ { value: 'PROXY_LOCAL', label: 'Proxy Local 模式' },
+ { value: 'PROXY_CLUSTER', label: 'Proxy Cluster 模式' },
{ value: 'DIRECT', label: 'Direct 模式' },
]}
/>
@@ -759,9 +764,11 @@ const InstancePage = () => {
extra={
addInstanceType === 'DIRECT'
? 'Direct 模式请填写 NameServer SLB 地址(K8s 场景下一般为 NameServer
Service 地址,如 namesrv.mq.svc:9876)'
- : addInstanceType === 'PROXY'
- ? 'Proxy 模式请填写 Proxy SLB 内网地址(如 proxy.mq.svc:8080)'
- : '请先选择接入方式'
+ : addInstanceType === 'PROXY_LOCAL'
+ ? 'Proxy Local 模式请填写与 Broker 同进程部署的 Proxy 接入地址(如
broker-proxy.mq.svc:8080)'
+ : addInstanceType === 'PROXY_CLUSTER'
+ ? 'Proxy Cluster 模式请填写独立 Proxy 集群的 SLB 内网地址(如
proxy.mq.svc:8080)'
+ : '请先选择接入方式'
}
>
<Input
@@ -811,7 +818,9 @@ const InstancePage = () => {
>
<Select
options={[
- { value: 'PROXY', label: 'Proxy 模式' },
+ { value: 'PROXY', label: 'Proxy 模式(部署形态未标明)' },
+ { value: 'PROXY_LOCAL', label: 'Proxy Local 模式' },
+ { value: 'PROXY_CLUSTER', label: 'Proxy Cluster 模式' },
{ value: 'DIRECT', label: 'Direct 模式' },
]}
/>
diff --git a/web/src/pages/instance/topic.tsx b/web/src/pages/instance/topic.tsx
index 63b8b29a7..b6eb6575b 100644
--- a/web/src/pages/instance/topic.tsx
+++ b/web/src/pages/instance/topic.tsx
@@ -66,10 +66,11 @@ import {
deleteTopic,
getTopicConsumerPage,
getTopicRoutes,
- listTopics,
+ listTopicsPage,
sendTopicMessage,
} from '../../services/topicService';
import { useInstanceFilter } from '../../hooks/useInstanceFilter';
+import type { Instance } from '../../api/instance';
import {
parseCsvTable,
validateTopicCsvImport,
@@ -79,6 +80,24 @@ import { buildCsv, downloadCsv, type CsvColumn } from
'../../utils/download';
const { Text } = Typography;
+const INSTANCE_ACCESS_LABEL: Record<Instance['type'], string> = {
+ PROXY: 'Proxy 模式(部署形态未标明)',
+ PROXY_LOCAL: 'Proxy Local 模式',
+ PROXY_CLUSTER: 'Proxy Cluster 模式',
+ DIRECT: 'Direct 模式',
+};
+
+const INSTANCE_ACCESS_DESCRIPTION: Record<Instance['type'], string> = {
+ PROXY:
+ '接入点为 Proxy 地址,但该存量实例未标明 Local/Cluster 部署形态。若客户端环境无法解析该地址,请自行配置 DNS
解析或在客户端 hosts 中映射。',
+ PROXY_LOCAL:
+ '接入点为与 Broker 同进程部署的 Proxy 地址。若客户端环境无法解析该地址,请自行配置 DNS 解析或在客户端 hosts 中映射。',
+ PROXY_CLUSTER:
+ '接入点为独立 Proxy 集群的 SLB 内网地址。若客户端环境无法解析该地址,请自行配置 DNS 解析或在客户端 hosts 中映射。',
+ DIRECT:
+ '接入点为 NameServer SLB 地址(K8s 场景下一般为 NameServer Service 地址),Direct
模式客户端通过该地址发现 Broker。若客户端环境无法解析该地址,请自行配置 DNS 解析或在客户端 hosts 中映射。',
+};
+
// ─── Cluster name lookup ───────────────────────────────────────────
const CLUSTER_NAME_MAP: Record<string, { name: string; type: string }> = {
'rmq-cn-v5-prod-01': { name: 'rmq-cn-v5-prod-01', type: 'V5_PROXY_CLUSTER' },
@@ -274,6 +293,7 @@ const TopicPage = () => {
// ─── State ─────────────────────────────────────────────────────
const [topics, setTopics] = useState<Topic[]>([]);
+ const [totalTopics, setTotalTopics] = useState(0);
const [loading, setLoading] = useState(false);
const [routesByTopic, setRoutesByTopic] = useState<Record<string,
BrokerRoute[]>>({});
const [consumersByTopic, setConsumersByTopic] = useState<Record<string,
TopicConsumerPage>>({});
@@ -319,6 +339,7 @@ const TopicPage = () => {
topicRequestIdRef.current += 1;
const resetTimer = window.setTimeout(() => {
setTopics([]);
+ setTotalTopics(0);
setSelectedRowKeys([]);
setLoading(false);
}, 0);
@@ -329,9 +350,18 @@ const TopicPage = () => {
const requestId = ++topicRequestIdRef.current;
const timer = window.setTimeout(() => {
setLoading(true);
- void listTopics({ instanceId: selectedInstanceId })
- .then((nextTopics) => {
- if (requestId === topicRequestIdRef.current) setTopics(nextTopics);
+ void listTopicsPage({
+ instanceId: selectedInstanceId,
+ type: typeFilter || undefined,
+ search: searchText.trim() || undefined,
+ page: tablePage,
+ pageSize: tablePageSize,
+ })
+ .then((result) => {
+ if (requestId === topicRequestIdRef.current) {
+ setTopics(result.items);
+ setTotalTopics(result.total);
+ }
})
.catch(() => {
if (requestId === topicRequestIdRef.current)
@@ -345,7 +375,7 @@ const TopicPage = () => {
return () => {
window.clearTimeout(timer);
};
- }, [selectedInstanceId]);
+ }, [selectedInstanceId, typeFilter, searchText, tablePage, tablePageSize]);
// ─── Filtered data ─────────────────────────────────────────────
const filteredTopics = useMemo(
@@ -361,7 +391,7 @@ const TopicPage = () => {
[topics, selectedInstanceId, searchText, typeFilter],
);
- const maxTablePage = Math.max(1, Math.ceil(filteredTopics.length /
tablePageSize));
+ const maxTablePage = Math.max(1, Math.ceil(totalTopics / tablePageSize));
const currentTablePage = Math.min(tablePage, maxTablePage);
const resetTablePage = () => {
@@ -945,7 +975,7 @@ const TopicPage = () => {
return (
<div style={{ padding: 24 }}>
{/* ── Header ────────────────────────────────────────────── */}
- <PageHeader title={t('topic.title')} subtitle={`共
${filteredTopics.length} 个 Topic`} />
+ <PageHeader title={t('topic.title')} subtitle={`共 ${totalTopics} 个
Topic`} />
{/* ── Current instance banner ───────────────────────────── */}
{selectedInstance && (
@@ -957,7 +987,7 @@ const TopicPage = () => {
</span>
<span>
<span style={{ color: '#8c8c8c', marginRight: 6 }}>接入模式</span>
- <span>{selectedInstance.type === 'DIRECT' ? 'Direct 模式' : 'Proxy
模式'}</span>
+ <span>{INSTANCE_ACCESS_LABEL[selectedInstance.type]}</span>
</span>
{selectedInstance.vendor === 'ALIYUN' && (
<span>
@@ -979,9 +1009,7 @@ const TopicPage = () => {
</span>
</Flex>
<div style={{ marginTop: 10, fontSize: 14, lineHeight: 1.6, color:
'#8c8c8c' }}>
- {selectedInstance.type === 'DIRECT'
- ? '接入点为 NameServer SLB 地址(K8s 场景下一般为 NameServer Service
地址),Direct 模式客户端通过该地址发现 Broker。若客户端环境无法解析该地址,请自行配置 DNS 解析或在客户端 hosts 中映射。'
- : '接入点为 Proxy SLB 内网地址,gRPC/Remoting
客户端直接连接该地址收发消息。若客户端环境无法解析该地址,请自行配置 DNS 解析或在客户端 hosts 中映射。'}
+ {INSTANCE_ACCESS_DESCRIPTION[selectedInstance.type]}
</div>
</InfoBanner>
)}
@@ -1149,6 +1177,7 @@ const TopicPage = () => {
pagination={{
current: currentTablePage,
pageSize: tablePageSize,
+ total: totalTopics,
showSizeChanger: true,
showTotal: (t) => `共 ${t} 条`,
onChange: (page, pageSize) => {
diff --git a/web/src/pages/studio/BrokerCluster.tsx
b/web/src/pages/studio/BrokerCluster.tsx
index 7cbea26a3..0356f4ab0 100644
--- a/web/src/pages/studio/BrokerCluster.tsx
+++ b/web/src/pages/studio/BrokerCluster.tsx
@@ -22,8 +22,8 @@ import { useLang } from '../../i18n/LangContext';
import { listClusters } from '../../services/clusterService';
import { isMockMode } from '../../services/dataMode';
import type { ClusterInfo } from '../../api/cluster';
+import { supportsApacheRuntime, type Instance } from '../../api/instance';
import { listInstances } from '../../services/instanceService';
-import type { Instance } from '../../api/instance';
import { useVisiblePolling } from '../../hooks/useVisiblePolling';
// ─── Types ──────────────────────────────────────────────────────
@@ -192,8 +192,9 @@ const BrokerClusterPage = () => {
void listInstances()
.then((nextInstances) => {
if (!active) return;
- setInstances(nextInstances);
- setSelectedInstanceId(nextInstances[0]?.name);
+ const apacheInstances = nextInstances.filter(supportsApacheRuntime);
+ setInstances(apacheInstances);
+ setSelectedInstanceId(apacheInstances[0]?.name);
})
.catch(() => {
if (!active) return;
diff --git a/web/src/pages/studio/Producer.tsx
b/web/src/pages/studio/Producer.tsx
index 5a6b661de..68cff6813 100644
--- a/web/src/pages/studio/Producer.tsx
+++ b/web/src/pages/studio/Producer.tsx
@@ -40,7 +40,7 @@ import {
type ProducerConnectionWarning,
type ProducerReadiness,
} from '../../api/producer';
-import type { Instance } from '../../api/instance';
+import { supportsApacheRuntime, type Instance } from '../../api/instance';
import { listInstances } from '../../services/instanceService';
import { buildCsv, downloadCsv, type CsvColumn } from '../../utils/download';
@@ -92,8 +92,13 @@ const ProducerPage = () => {
void listInstances()
.then((nextInstances) => {
if (cancelled) return;
- setInstances(nextInstances);
- setSelectedInstanceId((current) => current ?? nextInstances[0]?.name);
+ const apacheInstances = nextInstances.filter(supportsApacheRuntime);
+ setInstances(apacheInstances);
+ setSelectedInstanceId((current) =>
+ apacheInstances.some((instance) => instance.name === current)
+ ? current
+ : apacheInstances[0]?.name,
+ );
})
.catch(() => {
if (!cancelled) {
diff --git a/web/src/pages/studio/__tests__/Producer.test.tsx
b/web/src/pages/studio/__tests__/Producer.test.tsx
index ae70c2209..3fabf951f 100644
--- a/web/src/pages/studio/__tests__/Producer.test.tsx
+++ b/web/src/pages/studio/__tests__/Producer.test.tsx
@@ -113,6 +113,39 @@ describe('ProducerPage', () => {
});
});
+ it('uses an Apache instance rather than a cloud instance for producer
diagnostics', async () => {
+ vi.mocked(listInstances).mockResolvedValue([
+ {
+ id: 1,
+ name: 'cloud-instance',
+ remark: '',
+ type: 'DIRECT',
+ endpoint: 'cloud:9876',
+ vendor: 'ALIYUN',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ gmtCreate: '2026-08-01T00:00:00',
+ gmtModified: '2026-08-01T00:00:00',
+ },
+ {
+ id: 2,
+ name: 'apache-instance',
+ remark: '',
+ type: 'DIRECT',
+ endpoint: 'apache:9876',
+ vendor: 'APACHE',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ gmtCreate: '2026-08-01T00:00:00',
+ gmtModified: '2026-08-01T00:00:00',
+ },
+ ]);
+ renderWithProviders(<ProducerPage />);
+
+ await waitFor(() =>
expect(fetchTopicList).toHaveBeenCalledWith('apache-instance'));
+ expect(screen.queryByText('cloud-instance')).not.toBeInTheDocument();
+ });
+
it('renders topic options loaded from the API', async () => {
const user = userEvent.setup();
renderWithProviders(<ProducerPage />);
diff --git a/web/src/services/instanceService.test.ts
b/web/src/services/instanceService.test.ts
index aeb68c571..49798c903 100644
--- a/web/src/services/instanceService.test.ts
+++ b/web/src/services/instanceService.test.ts
@@ -22,7 +22,13 @@ vi.mock('../config', () => ({
API_BASE_URL: '/api',
}));
-import { createInstance, deleteInstance, listInstances, updateInstance } from
'./instanceService';
+import {
+ createInstance,
+ deleteInstance,
+ getInstanceCapabilities,
+ listInstances,
+ updateInstance,
+} from './instanceService';
describe('instanceService mock instances', () => {
it('returns defensive copies from list reads', async () => {
@@ -49,6 +55,16 @@ describe('instanceService mock instances', () => {
const combined = await listInstances({ type: 'DIRECT', search:
'instance-direct-2' });
expect(combined.map((instance) => instance.id)).toEqual([2]);
+
+ const allProxy = await listInstances({ type: 'PROXY' });
+ expect(allProxy.map((instance) => instance.type)).toEqual([
+ 'PROXY_CLUSTER',
+ 'PROXY_CLUSTER',
+ 'PROXY_LOCAL',
+ ]);
+ await expect(listInstances({ type: 'PROXY_LOCAL' })).resolves.toEqual([
+ expect.objectContaining({ type: 'PROXY_LOCAL' }),
+ ]);
});
it('does not expose created or updated store records by reference', async ()
=> {
@@ -66,6 +82,7 @@ describe('instanceService mock instances', () => {
const storedCreated = afterCreate.find((instance) => instance.id ===
created.id);
expect(storedCreated).toMatchObject({
name: 'rocketmq-copy-test',
+ type: 'PROXY_CLUSTER',
remark: 'created',
});
@@ -89,4 +106,18 @@ describe('instanceService mock instances', () => {
await expect(listInstances()).resolves.toEqual(before);
});
+
+ it('returns provider-specific mock capabilities without sharing mutable
arrays', async () => {
+ const first = await getInstanceCapabilities('instance-direct-1');
+ expect(first.instanceId).toBe('instance-direct-1');
+ expect(first.vendor).toBe('APACHE');
+ first.capabilities.length = 0;
+
+ const second = await getInstanceCapabilities('instance-direct-1');
+
+ expect(second.capabilities).toContain('DLQ_MANAGEMENT');
+ await expect(getInstanceCapabilities('missing-instance')).rejects.toThrow(
+ 'Instance not found: missing-instance',
+ );
+ });
});
diff --git a/web/src/services/instanceService.ts
b/web/src/services/instanceService.ts
index 593a97b35..bdfff2349 100644
--- a/web/src/services/instanceService.ts
+++ b/web/src/services/instanceService.ts
@@ -5,6 +5,7 @@ import type {
CreateInstanceRequest,
InstanceQuery,
UpdateInstanceRequest,
+ InstanceCapabilities,
} from '../api/instance';
import { mockInstances } from '../mock/instances';
@@ -14,11 +15,33 @@ function copyInstance(instance: Instance): Instance {
return { ...instance };
}
+function matchesType(instance: Instance, type?: Instance['type']) {
+ if (!type) return true;
+ return type === 'PROXY' ? instance.type !== 'DIRECT' : instance.type ===
type;
+}
+
+const APACHE_CAPABILITIES: InstanceCapabilities['capabilities'] = [
+ 'TOPIC_MANAGEMENT',
+ 'CONSUMER_GROUP_MANAGEMENT',
+ 'MESSAGE_QUERY',
+ 'MESSAGE_TRACE',
+ 'ACL_MANAGEMENT',
+ 'DLQ_MANAGEMENT',
+];
+
+const CLOUD_CAPABILITIES: InstanceCapabilities['capabilities'] = [
+ 'TOPIC_MANAGEMENT',
+ 'CONSUMER_GROUP_MANAGEMENT',
+ 'MESSAGE_QUERY',
+ 'MESSAGE_TRACE',
+ 'ACL_MANAGEMENT',
+];
+
export async function listInstances(query: InstanceQuery = {}):
Promise<Instance[]> {
if (isMockMode()) {
const search = query.search?.trim().toLowerCase();
return mockInstances
- .filter((instance) => !query.type || instance.type === query.type)
+ .filter((instance) => matchesType(instance, query.type))
.filter(
(instance) =>
!search ||
@@ -31,13 +54,33 @@ export async function listInstances(query: InstanceQuery =
{}): Promise<Instance
return instanceApi.listInstances(query);
}
+export async function getInstanceCapabilities(instanceId: string):
Promise<InstanceCapabilities> {
+ if (!isMockMode()) {
+ return instanceApi.getInstanceCapabilities(instanceId);
+ }
+ const instance = mockInstances.find((candidate) => candidate.name ===
instanceId);
+ if (!instance) throw new Error(`Instance not found: ${instanceId}`);
+ const vendor = instance.vendor ?? 'APACHE';
+ return {
+ instanceId: instance.name,
+ vendor,
+ accessType: instance.type,
+ capabilities: [...(vendor === 'APACHE' ? APACHE_CAPABILITIES :
CLOUD_CAPABILITIES)],
+ };
+}
+
export async function createInstance(data: CreateInstanceRequest):
Promise<Instance> {
if (isMockMode()) {
+ const cloudManaged = data.vendor === 'ALIYUN' || data.vendor === 'TENCENT';
const instance: Instance = {
id: Date.now(),
...data,
name: data.name || '',
- type: data.type || 'PROXY',
+ type: cloudManaged
+ ? 'PROXY'
+ : data.type === 'PROXY'
+ ? 'PROXY_CLUSTER'
+ : data.type || 'PROXY_CLUSTER',
endpoint: data.endpoint || '',
vendor: data.vendor || 'APACHE',
remark: data.remark || '',
diff --git a/web/src/services/topicService.ts b/web/src/services/topicService.ts
index fd7013430..9218e030d 100644
--- a/web/src/services/topicService.ts
+++ b/web/src/services/topicService.ts
@@ -3,6 +3,7 @@ import * as metadataApi from '../api/metadata';
import type {
Topic,
TopicQuery,
+ TopicPage,
BrokerRoute,
ConsumerGroupInfo,
TopicConsumerPage,
@@ -16,21 +17,43 @@ const cloneRoutes = (routes: BrokerRoute[]): BrokerRoute[]
=> routes.map((route)
const cloneConsumers = (consumers: ConsumerGroupInfo[]): ConsumerGroupInfo[] =>
consumers.map((consumer) => ({ ...consumer }));
+function filterMockTopics(params?: TopicQuery): Topic[] {
+ let result = [...mockTopics];
+ if (params?.search) {
+ const keyword = params.search.trim().toLowerCase();
+ if (keyword) result = result.filter((topic) =>
topic.name.toLowerCase().includes(keyword));
+ }
+ if (params?.type) result = result.filter((t) => t.type === params.type);
+ if (params?.clusterId) result = result.filter((t) => t.clusterId ===
params.clusterId);
+ if (params?.instanceId) result = result.filter((t) => t.instanceId ===
params.instanceId);
+ return (result as unknown as Topic[]).map(cloneTopic);
+}
+
export async function listTopics(params?: TopicQuery): Promise<Topic[]> {
if (isMockMode()) {
- let result = [...mockTopics];
- if (params?.search) {
- const keyword = params.search.trim().toLowerCase();
- if (keyword) result = result.filter((topic) =>
topic.name.toLowerCase().includes(keyword));
- }
- if (params?.type) result = result.filter((t) => t.type === params.type);
- if (params?.clusterId) result = result.filter((t) => t.clusterId ===
params.clusterId);
- if (params?.instanceId) result = result.filter((t) => t.instanceId ===
params.instanceId);
- return (result as unknown as Topic[]).map(cloneTopic);
+ return filterMockTopics(params);
}
return metadataApi.listTopics(params);
}
+export async function listTopicsPage(
+ params?: TopicQuery & { page?: number; pageSize?: number },
+): Promise<TopicPage> {
+ if (isMockMode()) {
+ const filtered = filterMockTopics(params);
+ const page = Math.max(params?.page ?? 1, 1);
+ const pageSize = Math.min(Math.max(params?.pageSize ?? 20, 1), 100);
+ const from = Math.min((page - 1) * pageSize, filtered.length);
+ return {
+ items: filtered.slice(from, from + pageSize),
+ total: filtered.length,
+ page,
+ size: pageSize,
+ };
+ }
+ return metadataApi.listTopicsPage(params);
+}
+
export async function createTopic(data: Partial<Topic>): Promise<Topic> {
if (isMockMode()) {
const duplicate = mockTopics.some(