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 c5957d242 feat(studio): Tencent ACL boundaries, producer connections,
proxy address management, broker config preview (#2552)
c5957d242 is described below
commit c5957d242f7cc0b1121e9f8306e4ec9c0d4b44f4
Author: coder999o <[email protected]>
AuthorDate: Tue Aug 25 17:29:27 2026 +0800
feat(studio): Tencent ACL boundaries, producer connections, proxy address
management, broker config preview (#2552)
* fix: enforce Tencent ACL role boundaries
* feat: query producer connections across active groups
* feat: manage Proxy addresses from the Studio console
* feat: preview broker config updates before applying
---
.../studio/cluster/broker/ClusterController.java | 7 +
.../studio/cluster/broker/ClusterService.java | 103 +++++++++-
.../cluster/client/ProducerConnectionService.java | 4 +-
.../cluster/client/ProducerConnectionVO.java | 2 +
.../studio/cluster/client/ProducerController.java | 1 -
.../ClusterConfigPreviewVO.java} | 43 ++--
.../ProxyAddressDTO.java} | 11 +-
.../studio/cluster/proxy/ProxyController.java | 13 ++
.../rocketmq/studio/instance/acl/AclService.java | 7 +-
.../studio/instance/acl/UpdateAclRuleDTO.java | 18 +-
.../studio/instance/acl/UpdateAclUserDTO.java | 19 +-
.../provider/apache/RocketMQClientProvider.java | 29 +++
.../studio/provider/tencent/TencentAclService.java | 27 ++-
.../cluster/broker/ClusterControllerTest.java | 43 ++++
.../studio/cluster/broker/ClusterServiceTest.java | 58 ++++++
.../client/ProducerConnectionServiceTest.java | 17 +-
.../cluster/client/ProducerControllerTest.java | 18 +-
.../studio/cluster/proxy/ProxyControllerTest.java | 57 ++++++
.../studio/instance/acl/AclControllerTest.java | 74 ++++++-
.../studio/instance/acl/AclServiceTest.java | 12 +-
.../apache/RocketMQClientProviderTest.java | 54 +++++
.../provider/tencent/TencentAclServiceTest.java | 61 ++++++
web/src/api/acl.test.ts | 36 ++++
web/src/api/acl.ts | 14 +-
web/src/api/cluster.test.ts | 28 +++
web/src/api/cluster.ts | 32 +++
web/src/api/producer.test.ts | 30 +++
web/src/api/producer.ts | 4 +-
web/src/api/proxy.test.ts | 28 ++-
web/src/api/proxy.ts | 12 ++
web/src/i18n/translations.ts | 33 +++
.../pages/cluster/__tests__/ClusterPage.test.tsx | 56 ++++++
web/src/pages/cluster/index.tsx | 224 +++++++++++++++++----
web/src/pages/instance/__tests__/AclPage.test.tsx | 113 +++++++++++
web/src/pages/instance/acl.tsx | 96 ++++++---
web/src/pages/studio/Producer.tsx | 25 ++-
web/src/pages/studio/Proxy.tsx | 118 ++++++++++-
web/src/pages/studio/__tests__/Producer.test.tsx | 25 ++-
web/src/pages/studio/__tests__/Proxy.test.tsx | 58 +++++-
web/src/services/aclService.ts | 10 +-
web/src/services/clusterService.test.ts | 25 +++
web/src/services/clusterService.ts | 87 +++++++-
42 files changed, 1563 insertions(+), 169 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterController.java
index 9441f5973..6797d779d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterController.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.cluster.broker;
import org.apache.rocketmq.studio.cluster.config.ClusterConfigUpdateResultVO;
+import org.apache.rocketmq.studio.cluster.config.ClusterConfigPreviewVO;
import org.apache.rocketmq.studio.cluster.config.UpdateConfigDTO;
import org.apache.rocketmq.studio.common.domain.Result;
@@ -69,6 +70,12 @@ public class ClusterController {
return Result.ok(clusterService.updateClusterConfig(command));
}
+ @PostMapping("/config/preview")
+ public Result<ClusterConfigPreviewVO> previewClusterConfig(@Valid
@RequestBody(required = false) UpdateConfigDTO command) {
+ requireUpdateConfigCommand(command);
+ return Result.ok(clusterService.previewClusterConfig(command));
+ }
+
@PostMapping("/{clusterId}/brokers/{name}/restart")
public Result<Map<String, Object>> restartBroker(@PathVariable String
clusterId,
@PathVariable String
name) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
index e7566deb0..41c96f67d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.cluster.broker;
import org.apache.rocketmq.studio.cluster.config.BrokerConfigUpdateFailureVO;
+import org.apache.rocketmq.studio.cluster.config.ClusterConfigPreviewVO;
import org.apache.rocketmq.studio.cluster.config.ClusterConfigUpdateResultVO;
import org.apache.rocketmq.studio.cluster.config.ClusterConfigVO;
import org.apache.rocketmq.studio.cluster.config.UpdateConfigDTO;
@@ -41,7 +42,10 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
+import java.util.LinkedHashMap;
import java.util.List;
+import java.util.Map;
+import java.util.Objects;
import java.util.Properties;
@Slf4j
@@ -157,6 +161,26 @@ public class ClusterService {
requireProxy(cluster, addr);
}
+ public ClusterConfigPreviewVO previewClusterConfig(UpdateConfigDTO
command) {
+ log.info("Previewing cluster config update for: {}", command.getId());
+ requireMatchingDefaultQueueNums(command);
+ ClusterVO cluster = resolveCluster(command.getId(),
command.getInstanceId());
+ ClusterConfigVO currentConfig = copyConfig(cluster.getConfig());
+ ClusterConfigVO proposedConfig = copyConfig(currentConfig);
+ applyConfig(command, proposedConfig);
+ List<ClusterConfigPreviewVO.ConfigChangeVO> changes =
+ findConfigChanges(currentConfig, proposedConfig);
+ return ClusterConfigPreviewVO.builder()
+ .cluster(cluster)
+ .currentConfig(currentConfig)
+ .proposedConfig(proposedConfig)
+ .targetBrokers(targetBrokers(cluster))
+ .brokerProperties(buildBrokerPropertyMap(command))
+ .changes(changes)
+ .changed(!changes.isEmpty())
+ .build();
+ }
+
/**
* Attach live broker configuration (read from the first reachable master
broker via the
* admin API) to a discovered cluster. Falls back to the persisted config,
if any, when the
@@ -342,33 +366,96 @@ public class ClusterService {
private Properties buildBrokerProperties(UpdateConfigDTO command) {
Properties props = new Properties();
+ buildBrokerPropertyMap(command).forEach(props::setProperty);
+ return props;
+ }
+
+ private Map<String, String> buildBrokerPropertyMap(UpdateConfigDTO
command) {
+ Map<String, String> props = new LinkedHashMap<>();
if (command.getFlushDiskType() != null) {
- props.setProperty("flushDiskType", command.getFlushDiskType());
+ props.put("flushDiskType", command.getFlushDiskType());
}
if (command.getAutoCreateTopicEnable() != null) {
- props.setProperty("autoCreateTopicEnable",
command.getAutoCreateTopicEnable().toString());
+ props.put("autoCreateTopicEnable",
command.getAutoCreateTopicEnable().toString());
}
if (command.getAutoCreateSubscriptionGroup() != null) {
- props.setProperty("autoCreateSubscriptionGroup",
command.getAutoCreateSubscriptionGroup().toString());
+ props.put("autoCreateSubscriptionGroup",
command.getAutoCreateSubscriptionGroup().toString());
}
if (command.getMaxMessageSize() != null) {
- props.setProperty("maxMessageSize",
command.getMaxMessageSize().toString());
+ props.put("maxMessageSize",
command.getMaxMessageSize().toString());
}
if (command.getFileReservedTime() != null) {
- props.setProperty("fileReservedTime",
command.getFileReservedTime().toString());
+ props.put("fileReservedTime",
command.getFileReservedTime().toString());
}
if (command.getWriteQueueNums() != null) {
- props.setProperty("defaultTopicQueueNums",
command.getWriteQueueNums().toString());
+ props.put("defaultTopicQueueNums",
command.getWriteQueueNums().toString());
}
if (command.getReadQueueNums() != null) {
- props.setProperty("defaultTopicQueueNums",
command.getReadQueueNums().toString());
+ props.put("defaultTopicQueueNums",
command.getReadQueueNums().toString());
}
if (command.getBrokerPermission() != null) {
- props.setProperty("brokerPermission",
command.getBrokerPermission().toString());
+ props.put("brokerPermission",
command.getBrokerPermission().toString());
}
return props;
}
+ private List<ClusterConfigPreviewVO.BrokerTargetVO>
targetBrokers(ClusterVO cluster) {
+ if (cluster.getBrokers() == null) {
+ return List.of();
+ }
+ return cluster.getBrokers().stream()
+ .filter(broker -> broker.getAddr() != null &&
!broker.getAddr().isBlank())
+ .map(broker -> ClusterConfigPreviewVO.BrokerTargetVO.builder()
+ .name(broker.getName())
+ .address(broker.getAddr())
+ .build())
+ .toList();
+ }
+
+ private List<ClusterConfigPreviewVO.ConfigChangeVO> findConfigChanges(
+ ClusterConfigVO current,
+ ClusterConfigVO proposed) {
+ List<ClusterConfigPreviewVO.ConfigChangeVO> changes = new
ArrayList<>();
+ addConfigChange(changes, "flushDiskType", "flushDiskType",
+ current.getFlushDiskType(), proposed.getFlushDiskType());
+ addConfigChange(changes, "autoCreateTopicEnable",
"autoCreateTopicEnable",
+ current.isAutoCreateTopicEnable(),
proposed.isAutoCreateTopicEnable());
+ addConfigChange(changes, "autoCreateSubscriptionGroup",
"autoCreateSubscriptionGroup",
+ current.isAutoCreateSubscriptionGroup(),
proposed.isAutoCreateSubscriptionGroup());
+ addConfigChange(changes, "maxMessageSize", "maxMessageSize",
+ current.getMaxMessageSize(), proposed.getMaxMessageSize());
+ addConfigChange(changes, "fileReservedTime", "fileReservedTime",
+ current.getFileReservedTime(), proposed.getFileReservedTime());
+ addConfigChange(changes, "writeQueueNums", "defaultTopicQueueNums",
+ current.getWriteQueueNums(), proposed.getWriteQueueNums());
+ addConfigChange(changes, "readQueueNums", "defaultTopicQueueNums",
+ current.getReadQueueNums(), proposed.getReadQueueNums());
+ addConfigChange(changes, "brokerPermission", "brokerPermission",
+ current.getBrokerPermission(), proposed.getBrokerPermission());
+ return changes;
+ }
+
+ private void addConfigChange(
+ List<ClusterConfigPreviewVO.ConfigChangeVO> changes,
+ String field,
+ String brokerProperty,
+ Object currentValue,
+ Object proposedValue) {
+ if (Objects.equals(currentValue, proposedValue)) {
+ return;
+ }
+ changes.add(ClusterConfigPreviewVO.ConfigChangeVO.builder()
+ .field(field)
+ .currentValue(previewValue(currentValue))
+ .proposedValue(previewValue(proposedValue))
+ .brokerProperty(brokerProperty)
+ .build());
+ }
+
+ private String previewValue(Object value) {
+ return value == null ? null : String.valueOf(value);
+ }
+
private void requireMatchingDefaultQueueNums(UpdateConfigDTO command) {
if (command.getWriteQueueNums() != null && command.getReadQueueNums()
!= null
&&
!command.getWriteQueueNums().equals(command.getReadQueueNums())) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
index 02211b137..b1f8ca293 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
@@ -38,7 +38,7 @@ public class ProducerConnectionService {
instanceId, topic, producerGroup);
String normalizedInstanceId = requireFilter(instanceId, "instanceId");
String normalizedTopic = requireFilter(topic, "topic");
- String normalizedProducerGroup = requireFilter(producerGroup,
"producerGroup");
+ String normalizedProducerGroup =
normalizeOptionalFilter(producerGroup);
return clientProvider.findProducerConnections(normalizedInstanceId,
normalizedTopic, normalizedProducerGroup)
.stream()
.map(this::toProducerConnection)
@@ -58,6 +58,8 @@ public class ProducerConnectionService {
return ProducerConnectionVO.builder()
.clientId(connection.getClientId())
.clientAddr(connection.getAddress())
+ .topic(connection.getGroupOrTopic())
+ .producerGroup(connection.getProducerGroup())
.language(connection.getLanguage() == null ? null :
connection.getLanguage().name())
.versionDesc(connection.getVersion())
.build();
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
index c2ddee8d0..9a31a757a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
@@ -28,6 +28,8 @@ import lombok.NoArgsConstructor;
public class ProducerConnectionVO {
private String clientId;
private String clientAddr;
+ private String topic;
+ private String producerGroup;
private String language;
private String versionDesc;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
index d83db7dba..61a6ba544 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerController.java
@@ -49,7 +49,6 @@ public class ProducerController {
@RequestParam(required = false) String producerGroup) {
requireParameter(instanceId, "instanceId");
requireParameter(topic, "topic");
- requireParameter(producerGroup, "producerGroup");
return new ProducerConnectionResultVO(
producerConnectionService.listConnections(instanceId, topic,
producerGroup));
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/config/ClusterConfigPreviewVO.java
similarity index 52%
copy from
server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
copy to
server/src/main/java/org/apache/rocketmq/studio/cluster/config/ClusterConfigPreviewVO.java
index c2ddee8d0..7519a8e50 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/config/ClusterConfigPreviewVO.java
@@ -14,20 +14,39 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.cluster.client;
+package org.apache.rocketmq.studio.cluster.config;
-import lombok.AllArgsConstructor;
+import org.apache.rocketmq.studio.cluster.broker.ClusterVO;
import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
+import lombok.Value;
-@Data
+import java.util.List;
+import java.util.Map;
+
+@Value
@Builder
-@NoArgsConstructor
-@AllArgsConstructor
-public class ProducerConnectionVO {
- private String clientId;
- private String clientAddr;
- private String language;
- private String versionDesc;
+public class ClusterConfigPreviewVO {
+ ClusterVO cluster;
+ ClusterConfigVO currentConfig;
+ ClusterConfigVO proposedConfig;
+ List<BrokerTargetVO> targetBrokers;
+ Map<String, String> brokerProperties;
+ List<ConfigChangeVO> changes;
+ boolean changed;
+
+ @Value
+ @Builder
+ public static class BrokerTargetVO {
+ String name;
+ String address;
+ }
+
+ @Value
+ @Builder
+ public static class ConfigChangeVO {
+ String field;
+ String currentValue;
+ String proposedValue;
+ String brokerProperty;
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/proxy/ProxyAddressDTO.java
similarity index 82%
copy from
server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
copy to
server/src/main/java/org/apache/rocketmq/studio/cluster/proxy/ProxyAddressDTO.java
index c2ddee8d0..1b22b36ed 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/proxy/ProxyAddressDTO.java
@@ -14,8 +14,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.cluster.client;
+package org.apache.rocketmq.studio.cluster.proxy;
+import jakarta.validation.constraints.NotBlank;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
@@ -25,9 +26,7 @@ import lombok.NoArgsConstructor;
@Builder
@NoArgsConstructor
@AllArgsConstructor
-public class ProducerConnectionVO {
- private String clientId;
- private String clientAddr;
- private String language;
- private String versionDesc;
+public class ProxyAddressDTO {
+ @NotBlank(message = "addr is required")
+ private String addr;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/proxy/ProxyController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/proxy/ProxyController.java
index 423ff8b30..93b7888d6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/proxy/ProxyController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/proxy/ProxyController.java
@@ -21,6 +21,7 @@ import org.apache.rocketmq.studio.common.domain.Result;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
+import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
@@ -50,6 +51,18 @@ public class ProxyController {
return Result.ok(proxyAddressService.buildTopology());
}
+ @PostMapping("/addresses")
+ public Result<ProxyHomeVO> addProxyAddress(@Valid @RequestBody
ProxyAddressDTO request) {
+ proxyAddressService.addProxyAddr(request.getAddr());
+ return Result.ok(proxyAddressService.getHomePage());
+ }
+
+ @DeleteMapping("/addresses")
+ public Result<ProxyHomeVO> removeProxyAddress(@RequestParam String addr) {
+ proxyAddressService.removeProxyAddr(addr);
+ return Result.ok(proxyAddressService.getHomePage());
+ }
+
@PostMapping("/config/reload")
public Result<Map<String, Boolean>> reloadProxyConfig(@Valid @RequestBody
RestartProxyDTO command) {
proxyAddressService.reloadConfig(command.getClusterId(),
command.getAddr());
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclService.java
index 9bd074e60..3cf970187 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclService.java
@@ -185,11 +185,12 @@ public class AclService {
if (isTencentInstance(instanceId)) {
return tencentAclService.updateUser(instanceId,
user.toAclUserVO());
}
- if (user.getId() == null) {
+ if (!StringUtils.hasText(user.getId())) {
throw new BusinessException(400, "ACL user id is required");
}
- log.info("Updating ACL user id={}, username={}", user.getId(),
user.getUsername());
- AclUserVO existing = aclRepository.findUserById(user.getId())
+ Long userId = EntityIds.parseId(user.getId());
+ log.info("Updating ACL user id={}, username={}", userId,
user.getUsername());
+ AclUserVO existing = aclRepository.findUserById(userId)
.orElseThrow(() -> new BusinessException(404, "ACL user not
found: " + user.getId()));
if (user.getUsername() != null &&
!StringUtils.hasText(user.getUsername())) {
throw new BusinessException(400, "ACL username is required");
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/UpdateAclRuleDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/UpdateAclRuleDTO.java
index b47801a76..d9154ac02 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/UpdateAclRuleDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/UpdateAclRuleDTO.java
@@ -17,15 +17,14 @@
package org.apache.rocketmq.studio.instance.acl;
import jakarta.validation.constraints.NotBlank;
-import jakarta.validation.constraints.NotNull;
import lombok.Data;
import java.util.List;
@Data
public class UpdateAclRuleDTO {
- @NotNull(message = "id is required")
- private Long id;
+ @NotBlank(message = "id is required")
+ private String id;
@NotBlank(message = "principal is required")
private String principal;
@NotBlank(message = "resource is required")
@@ -41,7 +40,7 @@ public class UpdateAclRuleDTO {
public AclRuleVO toAclRuleVO() {
return AclRuleVO.builder()
- .id(id)
+ .id(numericIdOrNull())
.principal(principal)
.resource(resource)
.resourceType(resourceType)
@@ -52,4 +51,15 @@ public class UpdateAclRuleDTO {
.aclVersion(aclVersion)
.build();
}
+
+ private Long numericIdOrNull() {
+ if (id == null || id.isBlank()) {
+ return null;
+ }
+ try {
+ return Long.parseLong(id.trim());
+ } catch (NumberFormatException ex) {
+ return null;
+ }
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/UpdateAclUserDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/UpdateAclUserDTO.java
index 9d5d93539..b695c493e 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/UpdateAclUserDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/UpdateAclUserDTO.java
@@ -16,15 +16,15 @@
*/
package org.apache.rocketmq.studio.instance.acl;
-import jakarta.validation.constraints.NotNull;
+import jakarta.validation.constraints.NotBlank;
import lombok.Data;
import java.util.List;
@Data
public class UpdateAclUserDTO {
- @NotNull(message = "id is required")
- private Long id;
+ @NotBlank(message = "id is required")
+ private String id;
private String username;
/**
* Null when the admin flag was not part of the partial update, in which
case the existing
@@ -41,7 +41,7 @@ public class UpdateAclUserDTO {
public AclUserVO toAclUserVO() {
return AclUserVO.builder()
- .id(id)
+ .id(numericIdOrNull())
.username(username)
.admin(admin != null && admin)
.permRead(permRead)
@@ -49,4 +49,15 @@ public class UpdateAclUserDTO {
.clusters(clusters)
.build();
}
+
+ private Long numericIdOrNull() {
+ if (id == null || id.isBlank()) {
+ return null;
+ }
+ try {
+ return Long.parseLong(id.trim());
+ } catch (NumberFormatException ex) {
+ return null;
+ }
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
index d7a0d863d..237e4dce6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
@@ -130,6 +130,35 @@ public class RocketMQClientProvider implements
ClientProvider {
}
private List<ClientConnectionVO> findProducerConnections(MQAdminExt
adminExt, String topic, String producerGroup) {
+ if (producerGroup == null || producerGroup.isBlank()) {
+ return findProducerConnectionsForActiveGroups(adminExt, topic);
+ }
+ return findProducerConnectionsForGroup(adminExt, topic, producerGroup);
+ }
+
+ private List<ClientConnectionVO>
findProducerConnectionsForActiveGroups(MQAdminExt adminExt, String topic) {
+ List<String> producerGroups = findProducerGroups(adminExt, topic,
null, Integer.MAX_VALUE);
+ if (producerGroups.isEmpty()) {
+ return List.of();
+ }
+ List<ClientConnectionVO> connections = new ArrayList<>();
+ int successfulGroupQueries = 0;
+ for (String producerGroup : producerGroups) {
+ try {
+ connections.addAll(findProducerConnectionsForGroup(adminExt,
topic, producerGroup));
+ successfulGroupQueries++;
+ } catch (BusinessException e) {
+ log.warn("Failed to query producer connections for group={},
skipping", producerGroup, e);
+ }
+ }
+ if (successfulGroupQueries == 0) {
+ throw new BusinessException(502, "Failed to query producer
connections from all groups");
+ }
+ return connections;
+ }
+
+ private List<ClientConnectionVO> findProducerConnectionsForGroup(
+ MQAdminExt adminExt, String topic, String producerGroup) {
try {
ProducerConnection producerConnection =
adminExt.examineProducerConnectionInfo(producerGroup,
topic);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentAclService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentAclService.java
index f57a1e930..03b858e28 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentAclService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentAclService.java
@@ -219,8 +219,9 @@ public class TencentAclService {
public AclRuleVO createRule(String instanceId, AclRuleVO rule) {
Context context = resolve(instanceId);
String principal = requireRulePrincipal(rule);
- boolean permRead = hasAction(rule, "SUB");
- boolean permWrite = hasAction(rule, "PUB");
+ requireClusterWideAllowRule(rule);
+ boolean permRead = hasPermissionAction(rule, "SUB");
+ boolean permWrite = hasPermissionAction(rule, "PUB");
ModifyRoleRequest request = new ModifyRoleRequest();
request.setInstanceId(context.cloudInstanceId());
request.setRole(principal);
@@ -261,6 +262,28 @@ public class TencentAclService {
.anyMatch(a -> action.equalsIgnoreCase(a));
}
+ private static boolean hasPermissionAction(AclRuleVO rule, String action) {
+ return hasAction(rule, "ALL") || hasAction(rule, action);
+ }
+
+ private static void requireClusterWideAllowRule(AclRuleVO rule) {
+ if (!isOptionalValue(rule.getResourceType(), RESOURCE_TYPE)
+ || !isOptionalValue(rule.getResource(), RESOURCE)
+ || !isOptionalValue(rule.getResourcePattern(),
RESOURCE_PATTERN)
+ || !isOptionalValue(rule.getScope(), SCOPE)) {
+ throw new BusinessException(400,
+ "Tencent Cloud roles only support cluster-wide ACL rules
on resource *");
+ }
+ if (!isOptionalValue(rule.getDecision(), DECISION)) {
+ throw new BusinessException(400,
+ "Tencent Cloud roles only support ALLOW ACL rules");
+ }
+ }
+
+ private static boolean isOptionalValue(String actual, String expected) {
+ return !StringUtils.hasText(actual) ||
expected.equalsIgnoreCase(actual.trim());
+ }
+
private static AclRuleVO toRule(String principal, boolean permRead,
boolean permWrite) {
List<String> actions = new ArrayList<>();
if (permWrite) {
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterControllerTest.java
index e2cb74d60..d4704080a 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterControllerTest.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.cluster.broker;
import org.apache.rocketmq.studio.cluster.config.ClusterConfigUpdateResultVO;
+import org.apache.rocketmq.studio.cluster.config.ClusterConfigPreviewVO;
import org.apache.rocketmq.studio.cluster.config.ClusterConfigVO;
import org.apache.rocketmq.studio.cluster.config.UpdateConfigDTO;
@@ -217,6 +218,48 @@ class ClusterControllerTest {
.andExpect(jsonPath("$.data.cluster.config.readQueueNums").value(16));
}
+ @Test
+ void previewConfigShouldReturnEffectiveBrokerChanges() throws Exception {
+ ClusterVO cluster = buildCluster("cluster-1", "production-cluster",
ClusterStatus.healthy);
+
when(clusterService.previewClusterConfig(any(UpdateConfigDTO.class))).thenReturn(
+ ClusterConfigPreviewVO.builder()
+ .cluster(cluster)
+ .targetBrokers(Collections.singletonList(
+ ClusterConfigPreviewVO.BrokerTargetVO.builder()
+ .name("broker-0")
+ .address("10.0.0.1:10911")
+ .build()))
+
.brokerProperties(Collections.singletonMap("defaultTopicQueueNums", "16"))
+ .changes(Collections.singletonList(
+ ClusterConfigPreviewVO.ConfigChangeVO.builder()
+ .field("writeQueueNums")
+ .currentValue("8")
+ .proposedValue("16")
+
.brokerProperty("defaultTopicQueueNums")
+ .build()))
+ .changed(true)
+ .build());
+
+ UpdateConfigDTO command = UpdateConfigDTO.builder()
+ .id("cluster-1")
+ .writeQueueNums(16)
+ .readQueueNums(16)
+ .build();
+
+ mockMvc.perform(post("/api/clusters/config/preview")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content(objectMapper.writeValueAsString(command)))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.code").value(200))
+ .andExpect(jsonPath("$.data.changed").value(true))
+
.andExpect(jsonPath("$.data.targetBrokers[0].address").value("10.0.0.1:10911"))
+
.andExpect(jsonPath("$.data.brokerProperties.defaultTopicQueueNums").value("16"))
+
.andExpect(jsonPath("$.data.changes[0].field").value("writeQueueNums"))
+
.andExpect(jsonPath("$.data.changes[0].brokerProperty").value("defaultTopicQueueNums"));
+
+
verify(clusterService).previewClusterConfig(any(UpdateConfigDTO.class));
+ }
+
@Test
void updateConfigShouldRejectNullRequestBody() throws Exception {
mockMvc.perform(post("/api/clusters/config/update")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
index 23d850555..2dd9e7614 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.cluster.broker;
import org.apache.rocketmq.studio.cluster.config.ClusterConfigUpdateResultVO;
+import org.apache.rocketmq.studio.cluster.config.ClusterConfigPreviewVO;
import org.apache.rocketmq.studio.cluster.config.ClusterConfigVO;
import org.apache.rocketmq.studio.cluster.config.UpdateConfigDTO;
import org.apache.rocketmq.studio.cluster.nameserver.CreateNameServerDTO;
@@ -159,6 +160,63 @@ class ClusterServiceTest {
verifyNoInteractions(clusterRepository, clusterProvider);
}
+ @Test
+ void
previewConfigShouldReturnEffectiveBrokerUpdateWithoutMutatingCluster() {
+ sampleCluster.setBrokers(List.of(
+
BrokerVO.builder().name("broker-0").addr("10.0.0.1:10911").build(),
+
BrokerVO.builder().name("broker-1").addr("10.0.0.2:10911").build()));
+
when(clusterRepository.findById("cluster-1")).thenReturn(Optional.of(sampleCluster));
+ ClusterConfigVO storedConfig = sampleCluster.getConfig();
+
+ ClusterConfigPreviewVO preview =
clusterService.previewClusterConfig(UpdateConfigDTO.builder()
+ .id("cluster-1")
+ .flushDiskType("SYNC_FLUSH")
+ .autoCreateTopicEnable(false)
+ .writeQueueNums(16)
+ .build());
+
+ assertThat(preview.isChanged()).isTrue();
+ assertThat(preview.getTargetBrokers())
+ .extracting(ClusterConfigPreviewVO.BrokerTargetVO::getAddress)
+ .containsExactly("10.0.0.1:10911", "10.0.0.2:10911");
+ assertThat(preview.getBrokerProperties())
+ .containsEntry("flushDiskType", "SYNC_FLUSH")
+ .containsEntry("autoCreateTopicEnable", "false")
+ .containsEntry("defaultTopicQueueNums", "16");
+ assertThat(preview.getChanges())
+ .extracting(ClusterConfigPreviewVO.ConfigChangeVO::getField)
+ .containsExactly("flushDiskType", "autoCreateTopicEnable",
"writeQueueNums", "readQueueNums");
+
assertThat(preview.getProposedConfig().getFlushDiskType()).isEqualTo(FlushDiskType.SYNC_FLUSH);
+
assertThat(preview.getProposedConfig().getWriteQueueNums()).isEqualTo(16);
+
assertThat(preview.getProposedConfig().getReadQueueNums()).isEqualTo(16);
+ assertThat(sampleCluster.getConfig()).isSameAs(storedConfig);
+
assertThat(storedConfig.getFlushDiskType()).isEqualTo(FlushDiskType.ASYNC_FLUSH);
+ assertThat(storedConfig.getWriteQueueNums()).isEqualTo(8);
+ assertThat(storedConfig.getReadQueueNums()).isEqualTo(8);
+ verify(clusterRepository, never()).updateConfig(eq("cluster-1"),
any());
+ verifyNoInteractions(brokerConfigService, auditService);
+ }
+
+ @Test
+ void previewConfigShouldUseSelectedInstanceWithoutCallingBrokerUpdate() {
+ when(clusterProvider.refreshClusterDetail("cluster-1",
"instance-1")).thenReturn(sampleCluster);
+
+ ClusterConfigPreviewVO preview =
clusterService.previewClusterConfig(UpdateConfigDTO.builder()
+ .id("cluster-1")
+ .instanceId("instance-1")
+ .brokerPermission(4)
+ .build());
+
+ assertThat(preview.getChanges())
+ .extracting(ClusterConfigPreviewVO.ConfigChangeVO::getField)
+ .containsExactly("brokerPermission");
+ verify(clusterProvider).refreshClusterDetail("cluster-1",
"instance-1");
+ verify(brokerConfigService, never()).updateBrokerConfig(any(), any(),
any());
+ verify(brokerConfigService, never()).updateBrokerConfig(any(), any(),
any(), any());
+ verifyNoInteractions(auditService);
+ verify(clusterRepository, never()).updateConfig(eq("cluster-1"),
any());
+ }
+
@Test
void listClustersShouldReturnEmptyListWhenNoClusters() {
when(clusterProvider.discoverClusters()).thenReturn(Collections.emptyList());
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
index 1de90e871..47f569207 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
@@ -62,6 +62,8 @@ class ProducerConnectionServiceTest {
assertThat(result).hasSize(1);
assertThat(result.get(0).getClientId()).isEqualTo("producer-1");
assertThat(result.get(0).getClientAddr()).isEqualTo("10.0.0.1:38888");
+ assertThat(result.get(0).getTopic()).isEqualTo("order-topic");
+ assertThat(result.get(0).getProducerGroup()).isEqualTo("pg-order");
assertThat(result.get(0).getLanguage()).isEqualTo("Java");
assertThat(result.get(0).getVersionDesc()).isEqualTo("5.1.0");
verify(clientProvider).findProducerConnections("instance-1",
"order-topic", "pg-order");
@@ -77,12 +79,15 @@ class ProducerConnectionServiceTest {
}
@Test
- void listConnectionsShouldRejectMissingProducerGroup() {
- assertThatThrownBy(() ->
producerConnectionService.listConnections("instance-1", "order-topic", null))
- .isInstanceOf(BusinessException.class)
- .hasMessage("producerGroup is required")
- .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(400));
- verifyNoInteractions(clientProvider);
+ void listConnectionsShouldAllowMissingProducerGroupForAllGroupScan() {
+ when(clientProvider.findProducerConnections("instance-1",
"order-topic", null))
+ .thenReturn(List.of());
+
+ List<ProducerConnectionVO> result =
+ producerConnectionService.listConnections("instance-1",
"order-topic", " ");
+
+ assertThat(result).isEmpty();
+ verify(clientProvider).findProducerConnections("instance-1",
"order-topic", null);
}
@Test
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
index f2a1fe9c2..e5f862a76 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerControllerTest.java
@@ -64,6 +64,8 @@ class ProducerControllerTest {
ProducerConnectionVO connection = ProducerConnectionVO.builder()
.clientId("producer-1")
.clientAddr("10.0.0.1:38888")
+ .topic("order-topic")
+ .producerGroup("pg-order")
.language("Java")
.versionDesc("5.1.0")
.build();
@@ -78,6 +80,8 @@ class ProducerControllerTest {
.andExpect(jsonPath("$.connectionSet").isArray())
.andExpect(jsonPath("$.connectionSet[0].clientId").value("producer-1"))
.andExpect(jsonPath("$.connectionSet[0].clientAddr").value("10.0.0.1:38888"))
+
.andExpect(jsonPath("$.connectionSet[0].topic").value("order-topic"))
+
.andExpect(jsonPath("$.connectionSet[0].producerGroup").value("pg-order"))
.andExpect(jsonPath("$.connectionSet[0].language").value("Java"))
.andExpect(jsonPath("$.connectionSet[0].versionDesc").value("5.1.0"))
.andExpect(jsonPath("$.summary.totalConnections").value(1))
@@ -101,15 +105,19 @@ class ProducerControllerTest {
}
@Test
- void listConnectionsShouldRequireProducerGroup() throws Exception {
+ void listConnectionsShouldAllowMissingProducerGroup() throws Exception {
+ when(producerConnectionService.listConnections("instance-1",
"order-topic", null))
+ .thenReturn(List.of());
+
mockMvc.perform(get("/api/producer/connection")
.param("instanceId", "instance-1")
.param("topic", "order-topic"))
- .andExpect(status().isBadRequest())
- .andExpect(jsonPath("$.code").value(400))
- .andExpect(jsonPath("$.message").value("producerGroup is
required"));
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.connectionSet").isArray())
+ .andExpect(jsonPath("$.summary.totalConnections").value(0))
+
.andExpect(jsonPath("$.summary.readiness").value("UNAVAILABLE"));
- verifyNoInteractions(producerConnectionService);
+ verify(producerConnectionService).listConnections("instance-1",
"order-topic", null);
}
@Test
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/proxy/ProxyControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/proxy/ProxyControllerTest.java
index 7133a8d8b..657c1173b 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/proxy/ProxyControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/proxy/ProxyControllerTest.java
@@ -160,6 +160,63 @@ class ProxyControllerTest {
verify(proxyAddressService).buildTopology();
}
+ @Test
+ void addProxyAddressShouldReturnUpdatedAddressState() throws Exception {
+ ProxyAddressDTO request = ProxyAddressDTO.builder()
+ .addr("10.0.0.10:8081")
+ .build();
+ when(proxyAddressService.getHomePage())
+ .thenReturn(ProxyHomeVO.builder()
+ .proxyAddrList(List.of("127.0.0.1:8081",
"10.0.0.10:8081"))
+ .currentProxyAddr("127.0.0.1:8081")
+ .build());
+
+ mockMvc.perform(post("/api/proxies/addresses")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content(objectMapper.writeValueAsString(request)))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.code").value(200))
+
.andExpect(jsonPath("$.data.proxyAddrList[1]").value("10.0.0.10:8081"));
+
+ verify(proxyAddressService).addProxyAddr(eq("10.0.0.10:8081"));
+ verify(proxyAddressService).getHomePage();
+ }
+
+ @Test
+ void addProxyAddressShouldRejectBlankAddress() throws Exception {
+ ProxyAddressDTO request = ProxyAddressDTO.builder()
+ .addr(" ")
+ .build();
+
+ mockMvc.perform(post("/api/proxies/addresses")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content(objectMapper.writeValueAsString(request)))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("addr is required"));
+
+ verifyNoInteractions(proxyAddressService);
+ }
+
+ @Test
+ void removeProxyAddressShouldReturnUpdatedAddressState() throws Exception {
+ when(proxyAddressService.getHomePage())
+ .thenReturn(ProxyHomeVO.builder()
+ .proxyAddrList(List.of("127.0.0.1:8081"))
+ .currentProxyAddr("127.0.0.1:8081")
+ .build());
+
+
mockMvc.perform(org.springframework.test.web.servlet.request.MockMvcRequestBuilders
+ .delete("/api/proxies/addresses")
+ .param("addr", "10.0.0.10:8081"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.code").value(200))
+
.andExpect(jsonPath("$.data.proxyAddrList[0]").value("127.0.0.1:8081"));
+
+ verify(proxyAddressService).removeProxyAddr(eq("10.0.0.10:8081"));
+ verify(proxyAddressService).getHomePage();
+ }
+
@Test
void reloadProxyConfigShouldReturnSuccess() throws Exception {
RestartProxyDTO request = RestartProxyDTO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclControllerTest.java
index b72f8862b..65ed60ff7 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclControllerTest.java
@@ -220,6 +220,45 @@ class AclControllerTest {
.andExpect(jsonPath("$.data.decision").value("DENY"));
}
+ @Test
+ void updateRuleShouldAcceptTencentRoleNameIdentifier() throws Exception {
+ AclRuleVO updated = AclRuleVO.builder()
+ .principal("reader-role")
+ .resource("*")
+ .resourceType("Cluster")
+ .resourcePattern("LITERAL")
+ .actions(List.of("SUB"))
+ .decision("ALLOW")
+ .scope("cluster")
+ .aclVersion("1.0")
+ .build();
+ when(aclService.updateRule(any(AclRuleVO.class),
eq("tencent-rmq"))).thenReturn(updated);
+
+ mockMvc.perform(post("/api/acl/rules/update")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+ {
+ "id": "reader-role",
+ "principal": "reader-role",
+ "resource": "*",
+ "resourceType": "Cluster",
+ "resourcePattern": "LITERAL",
+ "actions": ["SUB"],
+ "decision": "ALLOW",
+ "scope": "cluster",
+ "instanceId": "tencent-rmq"
+ }
+ """))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.principal").value("reader-role"))
+ .andExpect(jsonPath("$.data.resource").value("*"));
+
+ ArgumentCaptor<AclRuleVO> captor =
ArgumentCaptor.forClass(AclRuleVO.class);
+ verify(aclService).updateRule(captor.capture(), eq("tencent-rmq"));
+ assertThat(captor.getValue().getId()).isNull();
+ assertThat(captor.getValue().getPrincipal()).isEqualTo("reader-role");
+ }
+
@Test
void updateRuleShouldReturnNotFoundForUnknownId() throws Exception {
AclRuleVO input = AclRuleVO.builder()
@@ -365,7 +404,7 @@ class AclControllerTest {
@Test
void updateUserShouldReturnMaskedUpdatedUser() throws Exception {
UpdateAclUserDTO input = new UpdateAclUserDTO();
- input.setId(1L);
+ input.setId("1");
input.setUsername("admin");
input.setAdmin(false);
AclUserVO updated = AclUserVO.builder()
@@ -389,6 +428,39 @@ class AclControllerTest {
.andExpect(jsonPath("$.data.admin").value(false));
}
+ @Test
+ void updateUserShouldAcceptTencentRoleNameIdentifier() throws Exception {
+ AclUserVO updated = AclUserVO.builder()
+ .username("reader-role")
+ .accessKey("acce****3456")
+ .secretKey("secr****7654")
+ .admin(false)
+ .permRead(true)
+ .permWrite(false)
+ .build();
+ when(aclService.updateUser(any(UpdateAclUserDTO.class),
eq("tencent-rmq"))).thenReturn(updated);
+
+ mockMvc.perform(post("/api/acl/users/update")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+ {
+ "id": "reader-role",
+ "username": "reader-role",
+ "permRead": true,
+ "permWrite": false,
+ "instanceId": "tencent-rmq"
+ }
+ """))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.username").value("reader-role"))
+ .andExpect(jsonPath("$.data.permRead").value(true));
+
+ ArgumentCaptor<UpdateAclUserDTO> captor =
ArgumentCaptor.forClass(UpdateAclUserDTO.class);
+ verify(aclService).updateUser(captor.capture(), eq("tencent-rmq"));
+ assertThat(captor.getValue().getId()).isEqualTo("reader-role");
+ assertThat(captor.getValue().getUsername()).isEqualTo("reader-role");
+ }
+
@Test
void updateUserShouldRejectMissingId() throws Exception {
mockMvc.perform(post("/api/acl/users/update")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclServiceTest.java
index b9b64d246..e1274ab03 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclServiceTest.java
@@ -474,7 +474,7 @@ class AclServiceTest {
@Test
void updateUserShouldSaveExistingUser() {
UpdateAclUserDTO input = new UpdateAclUserDTO();
- input.setId(1L);
+ input.setId("1");
input.setUsername("newuser");
input.setAdmin(true);
@@ -507,7 +507,7 @@ class AclServiceTest {
.clusters(List.of("cluster-a"))
.build();
UpdateAclUserDTO input = new UpdateAclUserDTO();
- input.setId(1L);
+ input.setId("1");
input.setUsername("renamed");
// admin intentionally left null: the existing admin flag must survive
the partial update.
@@ -526,7 +526,7 @@ class AclServiceTest {
@Test
void updateUserShouldRejectBlankUsernameWithoutSaving() {
UpdateAclUserDTO input = new UpdateAclUserDTO();
- input.setId(1L);
+ input.setId("1");
input.setUsername(" ");
when(aclRepository.findUserById(1L)).thenReturn(Optional.of(existingUser));
@@ -541,7 +541,7 @@ class AclServiceTest {
@Test
void updateUserShouldThrowWhenUserDoesNotExist() {
UpdateAclUserDTO input = new UpdateAclUserDTO();
- input.setId(999L);
+ input.setId("999");
input.setUsername("ghost");
when(aclRepository.findUserById(999L)).thenReturn(Optional.empty());
@@ -556,7 +556,7 @@ class AclServiceTest {
@Test
void updateUserShouldRejectConcurrentDeletion() {
UpdateAclUserDTO input = new UpdateAclUserDTO();
- input.setId(1L);
+ input.setId("1");
input.setUsername("renamed");
when(aclRepository.findUserById(1L)).thenReturn(Optional.of(existingUser));
when(aclRepository.replaceUser(any(AclUserVO.class))).thenReturn(Optional.empty());
@@ -640,7 +640,7 @@ class AclServiceTest {
AclUserVO listed = aclService.listUsers(null).get(0);
UpdateAclUserDTO update = new UpdateAclUserDTO();
- update.setId(listed.getId());
+ update.setId(String.valueOf(listed.getId()));
update.setUsername("orders-admin");
update.setAdmin(true);
update.setClusters(listed.getClusters());
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
index 94e2d6f74..f23be04ca 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
@@ -281,6 +281,60 @@ class RocketMQClientProviderTest {
verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
}
+ @Test
+ void producerQueryWithoutGroupScansActiveProducerGroups() throws Exception
{
+
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo("127.0.0.1:10911"));
+ when(adminExt.getAllProducerInfo("127.0.0.1:10911"))
+ .thenReturn(new ProducerTableInfo(Map.of(
+ "pg-order", List.of(producerInfo("producer-order",
"10.0.0.1:1000")),
+ "pg-payment", List.of(producerInfo("producer-payment",
"10.0.0.2:1000")))));
+ ProducerConnection orderConnection = new ProducerConnection();
+ orderConnection.setConnectionSet(new HashSet<>(List.of(
+ connection("producer-order", "10.0.0.1:1000"))));
+ ProducerConnection paymentConnection = new ProducerConnection();
+ paymentConnection.setConnectionSet(new HashSet<>(List.of(
+ connection("producer-payment", "10.0.0.2:1000"))));
+ when(adminExt.examineProducerConnectionInfo("pg-order", "TopicA"))
+ .thenReturn(orderConnection);
+ when(adminExt.examineProducerConnectionInfo("pg-payment", "TopicA"))
+ .thenReturn(paymentConnection);
+
+ List<ClientConnectionVO> connections =
provider.findProducerConnections("instance-a", "TopicA", null);
+
+ assertThat(connections)
+ .extracting(ClientConnectionVO::getProducerGroup)
+ .containsExactly("pg-order", "pg-payment");
+ assertThat(connections)
+ .extracting(ClientConnectionVO::getGroupOrTopic)
+ .containsOnly("TopicA");
+ verify(adminExt).getAllProducerInfo("127.0.0.1:10911");
+ verify(adminExt).examineProducerConnectionInfo("pg-order", "TopicA");
+ verify(adminExt).examineProducerConnectionInfo("pg-payment", "TopicA");
+ }
+
+ @Test
+ void producerQueryWithoutGroupReturnsPartialResultsWhenOneGroupFails()
throws Exception {
+
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo("127.0.0.1:10911"));
+ when(adminExt.getAllProducerInfo("127.0.0.1:10911"))
+ .thenReturn(new ProducerTableInfo(Map.of(
+ "pg-order", List.of(producerInfo("producer-order",
"10.0.0.1:1000")),
+ "pg-payment", List.of(producerInfo("producer-payment",
"10.0.0.2:1000")))));
+ ProducerConnection paymentConnection = new ProducerConnection();
+ paymentConnection.setConnectionSet(new HashSet<>(List.of(
+ connection("producer-payment", "10.0.0.2:1000"))));
+ when(adminExt.examineProducerConnectionInfo("pg-order", "TopicA"))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+ when(adminExt.examineProducerConnectionInfo("pg-payment", "TopicA"))
+ .thenReturn(paymentConnection);
+
+ List<ClientConnectionVO> connections =
provider.findProducerConnections("instance-a", "TopicA", null);
+
+ assertThat(connections).singleElement().satisfies(connection -> {
+ assertThat(connection.getClientId()).isEqualTo("producer-payment");
+ assertThat(connection.getProducerGroup()).isEqualTo("pg-payment");
+ });
+ }
+
@Test
void exactProducerQueryTranslatesAdminFailureToBadGateway() throws
Exception {
when(adminExt.examineProducerConnectionInfo("pg-order", "TopicA"))
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
index 8e431af52..4304752ff 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
@@ -23,6 +23,7 @@ import com.tencentcloudapi.trocket.v20230308.models.RoleItem;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.instance.InstanceRepository;
import org.apache.rocketmq.studio.instance.InstanceVO;
+import org.apache.rocketmq.studio.instance.acl.AclRuleVO;
import org.apache.rocketmq.studio.instance.acl.AclUserVO;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -31,6 +32,7 @@ import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
+import java.util.List;
import java.util.Optional;
import static org.assertj.core.api.Assertions.assertThat;
@@ -150,4 +152,63 @@ class TencentAclServiceTest {
.isInstanceOf(BusinessException.class)
.hasMessageContaining("ACL user not found: ghost-role");
}
+
+ @Test
+ void createRuleShouldExpandAllActionToReadAndWritePermissionsTest() throws
Exception {
+ AclRuleVO rule = AclRuleVO.builder()
+ .principal("reader-role")
+ .resource("*")
+ .resourceType("Cluster")
+ .resourcePattern("LITERAL")
+ .actions(List.of("ALL"))
+ .decision("ALLOW")
+ .scope("cluster")
+ .build();
+
+ AclRuleVO created = service.createRule(INSTANCE_ID, rule);
+
+ ArgumentCaptor<ModifyRoleRequest> requestCaptor =
ArgumentCaptor.forClass(ModifyRoleRequest.class);
+ verify(client).ModifyRole(requestCaptor.capture());
+ ModifyRoleRequest request = requestCaptor.getValue();
+ assertThat(request.getRole()).isEqualTo("reader-role");
+ assertThat(request.getPermRead()).isTrue();
+ assertThat(request.getPermWrite()).isTrue();
+ assertThat(created.getActions()).containsExactly("PUB", "SUB");
+ }
+
+ @Test
+ void createRuleShouldRejectResourceScopedTencentRulesTest() {
+ AclRuleVO rule = AclRuleVO.builder()
+ .principal("reader-role")
+ .resource("orders-*")
+ .resourceType("Topic")
+ .resourcePattern("PREFIX")
+ .actions(List.of("SUB"))
+ .decision("ALLOW")
+ .scope("cluster")
+ .build();
+
+ assertThatThrownBy(() -> service.createRule(INSTANCE_ID, rule))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Tencent Cloud roles only support cluster-wide ACL
rules on resource *")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(400));
+ }
+
+ @Test
+ void createRuleShouldRejectDenyTencentRulesTest() {
+ AclRuleVO rule = AclRuleVO.builder()
+ .principal("reader-role")
+ .resource("*")
+ .resourceType("Cluster")
+ .resourcePattern("LITERAL")
+ .actions(List.of("SUB"))
+ .decision("DENY")
+ .scope("cluster")
+ .build();
+
+ assertThatThrownBy(() -> service.createRule(INSTANCE_ID, rule))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Tencent Cloud roles only support ALLOW ACL rules")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(400));
+ }
}
diff --git a/web/src/api/acl.test.ts b/web/src/api/acl.test.ts
index abf4c357d..907d9ca85 100644
--- a/web/src/api/acl.test.ts
+++ b/web/src/api/acl.test.ts
@@ -25,6 +25,7 @@ import {
deleteAclRule,
deleteAclUser,
examineBrokerClusterAclConfig,
+ getAclUserCredentials,
listAclRules,
updateAclRule,
updateAclUser,
@@ -140,6 +141,41 @@ describe('ACL API contract', () => {
await expect(deleteAclUser(user.id)).resolves.toBeUndefined();
});
+ it('passes Tencent role names as ACL entity identifiers', async () => {
+ mock.onPost('/acl/rules/delete').reply((config) => {
+ expect(JSON.parse(config.data)).toEqual({ id: 'reader-role', instanceId:
'tencent-rmq' });
+ return [200, { code: 200 }];
+ });
+ mock.onPost('/acl/users/delete').reply((config) => {
+ expect(JSON.parse(config.data)).toEqual({ id: 'reader-role', instanceId:
'tencent-rmq' });
+ return [200, { code: 200 }];
+ });
+ mock.onGet('/acl/users/reader-role/credentials').reply((config) => {
+ expect(config.params).toEqual({ instanceId: 'tencent-rmq' });
+ return [
+ 200,
+ {
+ code: 200,
+ data: {
+ id: 'reader-role',
+ username: 'reader-role',
+ accessKey: 'ak',
+ secretKey: 'sk',
+ admin: false,
+ clusters: ['rmq-cloud'],
+ },
+ },
+ ];
+ });
+
+ await expect(deleteAclRule('reader-role',
'tencent-rmq')).resolves.toBeUndefined();
+ await expect(deleteAclUser('reader-role',
'tencent-rmq')).resolves.toBeUndefined();
+ await expect(getAclUserCredentials('reader-role',
'tencent-rmq')).resolves.toMatchObject({
+ username: 'reader-role',
+ secretKey: 'sk',
+ });
+ });
+
it('fetches cluster ACL config by clusterId', async () => {
mock.onGet('/acl/cluster-config').reply((config) => {
expect(config.params).toEqual({ clusterId: 'cluster-a' });
diff --git a/web/src/api/acl.ts b/web/src/api/acl.ts
index cf85446bd..a4ac9331d 100644
--- a/web/src/api/acl.ts
+++ b/web/src/api/acl.ts
@@ -8,8 +8,10 @@ export interface PageResult<T> {
size: number;
}
+export type AclEntityId = number | string;
+
export interface AclRule {
- id: number;
+ id: AclEntityId;
principal: string;
resource: string;
resourceType: string;
@@ -46,7 +48,7 @@ export interface AclUserPage {
}
export interface AclUser {
- id: number;
+ id: AclEntityId;
username: string;
accessKey: string;
secretKey: string;
@@ -72,7 +74,7 @@ export async function updateAclRule(data: Partial<AclRule> &
{ instanceId?: stri
return res.data.data;
}
-export async function deleteAclRule(id: number, instanceId?: string) {
+export async function deleteAclRule(id: AclEntityId, instanceId?: string) {
await client.post('/acl/rules/delete', { id, instanceId });
}
@@ -86,9 +88,9 @@ export async function pageAclUsers(params: AclUserQuery & {
page: number; pageSi
return res.data.data;
}
-export async function getAclUserCredentials(id: number, instanceId?: string) {
+export async function getAclUserCredentials(id: AclEntityId, instanceId?:
string) {
const res = await client.get<{ data: AclUser }>(
- `/acl/users/${encodeURIComponent(id)}/credentials`,
+ `/acl/users/${encodeURIComponent(String(id))}/credentials`,
{ params: { instanceId } },
);
return res.data.data;
@@ -104,7 +106,7 @@ export async function updateAclUser(data: Partial<AclUser>
& { instanceId?: stri
return res.data.data;
}
-export async function deleteAclUser(id: number, instanceId?: string) {
+export async function deleteAclUser(id: AclEntityId, instanceId?: string) {
await client.post('/acl/users/delete', { id, instanceId });
}
diff --git a/web/src/api/cluster.test.ts b/web/src/api/cluster.test.ts
index 1995f282c..21a48aa1f 100644
--- a/web/src/api/cluster.test.ts
+++ b/web/src/api/cluster.test.ts
@@ -26,6 +26,7 @@ import {
getNameServerConfigDiff,
getCluster,
listK8sCerts,
+ previewClusterConfig,
renewK8sCert,
restartBroker,
restartNameServer,
@@ -149,6 +150,33 @@ describe('K8s certificate API', () => {
);
});
+ it('previews effective broker config update changes', async () => {
+ const result = {
+ cluster: { id: 'cluster-1' },
+ currentConfig: { writeQueueNums: 8, readQueueNums: 8 },
+ proposedConfig: { writeQueueNums: 16, readQueueNums: 16 },
+ targetBrokers: [{ name: 'broker-0', address: '10.0.0.1:10911' }],
+ brokerProperties: { defaultTopicQueueNums: '16' },
+ changes: [
+ {
+ field: 'writeQueueNums',
+ currentValue: '8',
+ proposedValue: '16',
+ brokerProperty: 'defaultTopicQueueNums',
+ },
+ ],
+ changed: true,
+ };
+ mock.onPost('/clusters/config/preview').reply((config) => {
+ expect(JSON.parse(config.data)).toMatchObject({ id: 'cluster-1',
writeQueueNums: 16 });
+ return [200, { code: 200, data: result }];
+ });
+
+ await expect(previewClusterConfig({ id: 'cluster-1', writeQueueNums: 16
})).resolves.toEqual(
+ result,
+ );
+ });
+
it('sends NameServer operation payloads to their endpoints', async () => {
const target = { clusterId: 'cluster-1', addr: '127.0.0.1:9876' };
const requests = [
diff --git a/web/src/api/cluster.ts b/web/src/api/cluster.ts
index 4da2c437f..fb9869f89 100644
--- a/web/src/api/cluster.ts
+++ b/web/src/api/cluster.ts
@@ -100,6 +100,28 @@ export interface ClusterConfigUpdateResult {
failedBrokers: BrokerConfigUpdateFailure[];
}
+export interface ClusterConfigPreviewTarget {
+ name: string;
+ address: string;
+}
+
+export interface ClusterConfigPreviewChange {
+ field: string;
+ currentValue: string | null;
+ proposedValue: string | null;
+ brokerProperty: string;
+}
+
+export interface ClusterConfigPreviewResult {
+ cluster: ClusterInfo;
+ currentConfig: ClusterConfig;
+ proposedConfig: ClusterConfig;
+ targetBrokers: ClusterConfigPreviewTarget[];
+ brokerProperties: Record<string, string>;
+ changes: ClusterConfigPreviewChange[];
+ changed: boolean;
+}
+
export interface ClusterProbeResult {
connected: boolean;
namesrvAddr: string;
@@ -184,6 +206,16 @@ export async function updateClusterConfig(
return res.data.data;
}
+export async function previewClusterConfig(
+ data: { id: string; instanceId?: string } & Partial<ClusterConfig>,
+) {
+ const res = await client.post<{ data: ClusterConfigPreviewResult }>(
+ '/clusters/config/preview',
+ data,
+ );
+ return res.data.data;
+}
+
export async function restartBroker(clusterId: string, brokerName: string) {
const res = await client.post<{ data: { success: boolean; message: string }
}>(
`/clusters/${pathSegment(clusterId)}/brokers/${pathSegment(brokerName)}/restart`,
diff --git a/web/src/api/producer.test.ts b/web/src/api/producer.test.ts
index fd45e9773..9878c6744 100644
--- a/web/src/api/producer.test.ts
+++ b/web/src/api/producer.test.ts
@@ -95,6 +95,8 @@ describe('Producer API', () => {
{
clientId: 'producer-1',
clientAddr: '192.168.1.10',
+ topic: 'order-events',
+ producerGroup: 'order-producer',
language: 'JAVA',
versionDesc: '5.1.0',
},
@@ -115,10 +117,38 @@ describe('Producer API', () => {
const result = await queryProducerConnection('instance-1', 'order-events',
'order-producer');
expect(result.connectionSet).toHaveLength(2);
expect(result.connectionSet[0].clientId).toBe('producer-1');
+ expect(result.connectionSet[0].producerGroup).toBe('order-producer');
expect(result.summary.totalConnections).toBe(2);
expect(result.summary.readiness).toBe('READY');
});
+ it('queries producer connections without a producer group for all-group
scans', async () => {
+ mock.onGet('/producer/connection').reply((config) => {
+ expect(config.params.topic).toBe('order-events');
+ expect(config.params.instanceId).toBe('instance-1');
+ expect(config.params.producerGroup).toBeUndefined();
+ return [
+ 200,
+ {
+ connectionSet: [
+ {
+ clientId: 'producer-1',
+ clientAddr: '192.168.1.10',
+ topic: 'order-events',
+ producerGroup: 'pg-order',
+ language: 'JAVA',
+ versionDesc: '5.1.0',
+ },
+ ],
+ },
+ ];
+ });
+
+ const result = await queryProducerConnection('instance-1', 'order-events');
+ expect(result.connectionSet[0].producerGroup).toBe('pg-order');
+ expect(result.summary.totalConnections).toBe(1);
+ });
+
it('handles empty producer connections', async () => {
mock.onGet('/producer/connection').reply(200, { connectionSet: [] });
diff --git a/web/src/api/producer.ts b/web/src/api/producer.ts
index 14b43a228..c62ac1f12 100644
--- a/web/src/api/producer.ts
+++ b/web/src/api/producer.ts
@@ -21,6 +21,8 @@ import client from './client';
export interface ProducerConnection {
clientId: string;
clientAddr: string;
+ topic?: string;
+ producerGroup?: string;
language: string;
versionDesc: string;
}
@@ -178,7 +180,7 @@ export async function fetchProducerGroups(
export async function queryProducerConnection(
instanceId: string,
topic: string,
- producerGroup: string,
+ producerGroup?: string,
): Promise<ProducerConnectionResult> {
const res = await
client.get<ProducerConnectionResponse>('/producer/connection', {
params: { instanceId, topic, producerGroup },
diff --git a/web/src/api/proxy.test.ts b/web/src/api/proxy.test.ts
index fed361e0a..e871790d1 100644
--- a/web/src/api/proxy.test.ts
+++ b/web/src/api/proxy.test.ts
@@ -18,7 +18,7 @@
import MockAdapter from 'axios-mock-adapter';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
import client from './client';
-import { queryProxyHomePage } from './proxy';
+import { addProxyAddress, queryProxyHomePage, removeProxyAddress } from
'./proxy';
const mock = new MockAdapter(client);
@@ -54,4 +54,30 @@ describe('Proxy API', () => {
expect(result.proxyAddrList).toHaveLength(0);
expect(result.currentProxyAddr).toBe('');
});
+
+ it('adds proxy addresses through the Studio endpoint', async () => {
+ const data = {
+ proxyAddrList: ['127.0.0.1:8081', '10.0.0.10:8081'],
+ currentProxyAddr: '127.0.0.1:8081',
+ };
+ mock.onPost('/proxies/addresses').reply((config) => {
+ expect(JSON.parse(config.data)).toEqual({ addr: '10.0.0.10:8081' });
+ return [200, { code: 200, data }];
+ });
+
+ await expect(addProxyAddress('10.0.0.10:8081')).resolves.toEqual(data);
+ });
+
+ it('removes proxy addresses through the Studio endpoint', async () => {
+ const data = {
+ proxyAddrList: ['127.0.0.1:8081'],
+ currentProxyAddr: '127.0.0.1:8081',
+ };
+ mock.onDelete('/proxies/addresses').reply((config) => {
+ expect(config.params).toEqual({ addr: '10.0.0.10:8081' });
+ return [200, { code: 200, data }];
+ });
+
+ await expect(removeProxyAddress('10.0.0.10:8081')).resolves.toEqual(data);
+ });
});
diff --git a/web/src/api/proxy.ts b/web/src/api/proxy.ts
index b4d5003e9..e5b7c19a4 100644
--- a/web/src/api/proxy.ts
+++ b/web/src/api/proxy.ts
@@ -60,6 +60,18 @@ export async function getProxyTopology():
Promise<ProxyTopologyNode[]> {
return res.data.data;
}
+export async function addProxyAddress(addr: string):
Promise<ProxyHomePageData> {
+ const res = await client.post<{ data: ProxyHomePageData
}>('/proxies/addresses', { addr });
+ return res.data.data;
+}
+
+export async function removeProxyAddress(addr: string):
Promise<ProxyHomePageData> {
+ const res = await client.delete<{ data: ProxyHomePageData
}>('/proxies/addresses', {
+ params: { addr },
+ });
+ return res.data.data;
+}
+
/**
* Trigger a configuration hot-reload for a proxy.
* Uses the same DTO as restartProxy ({ clusterId, addr }).
diff --git a/web/src/i18n/translations.ts b/web/src/i18n/translations.ts
index 3321b3a9d..49dd7263c 100644
--- a/web/src/i18n/translations.ts
+++ b/web/src/i18n/translations.ts
@@ -395,6 +395,14 @@ const translations: Record<string, Record<Lang, string>> =
{
zh: '它们尚不会下发到所选 RocketMQ 实例。实例级 ACL Provider 接入完成前,请在集群侧管理实际 ACL 策略。',
en: 'They are not applied to the selected RocketMQ instance. Manage the
effective ACL policy on the cluster until instance-scoped provider support is
available.',
},
+ 'acl.tencentRoleNotice': {
+ zh: '当前 ACL 用户和规则由 Tencent Cloud Role 实时驱动',
+ en: 'Current ACL users and rules are backed by Tencent Cloud roles',
+ },
+ 'acl.tencentRoleDescription': {
+ zh: 'Tencent Cloud Role 只支持集群级允许规则,资源固定为 *,权限通过 PUB/SUB 映射为写/读角色权限。',
+ en: 'Tencent Cloud roles only support cluster-wide allow rules on resource
*, with PUB/SUB mapped to write/read role permissions.',
+ },
'acl.addRule': { zh: '添加规则', en: 'Add Rule' },
'acl.addUser': { zh: '添加用户', en: 'Add User' },
'acl.ruleTab': { zh: 'ACL 规则', en: 'ACL Rules' },
@@ -677,6 +685,16 @@ const translations: Record<string, Record<Lang, string>> =
{
zh: 'Broker 配置更新失败,请检查:{brokers}',
en: 'Broker configuration update failed. Check: {brokers}',
},
+ 'cluster.configPreview': { zh: '预览', en: 'Preview' },
+ 'cluster.configPreviewGenerated': { zh: '预览已生成', en: 'Preview generated' },
+ 'cluster.configPreviewFailed': { zh: '预览失败', en: 'Preview failed' },
+ 'cluster.configPreviewTargets': { zh: '目标 Broker', en: 'Target Brokers' },
+ 'cluster.configPreviewBrokerProperties': { zh: 'Broker 属性', en: 'Broker
Properties' },
+ 'cluster.configPreviewField': { zh: '配置项', en: 'Config Field' },
+ 'cluster.configPreviewCurrent': { zh: '当前值', en: 'Current Value' },
+ 'cluster.configPreviewProposed': { zh: '预期值', en: 'Proposed Value' },
+ 'cluster.configPreviewProperty': { zh: 'Broker 属性', en: 'Broker Property' },
+ 'cluster.configPreviewNoChanges': { zh: '无配置变更', en: 'No config changes' },
'cluster.flushDiskType': { zh: '刷盘方式', en: 'Flush Disk Type' },
'cluster.syncFlush': { zh: '同步刷盘', en: 'Sync Flush' },
'cluster.asyncFlush': { zh: '异步刷盘', en: 'Async Flush' },
@@ -1126,6 +1144,17 @@ const translations: Record<string, Record<Lang, string>>
= {
'proxy.warning': { zh: '警告', en: 'Warning' },
'proxy.clusterId': { zh: '集群 ID', en: 'Cluster ID' },
'proxy.clusterIdPlaceholder': { zh: '例:DefaultCluster', en: 'e.g.
DefaultCluster' },
+ 'proxy.address': { zh: 'Proxy 地址', en: 'Proxy Address' },
+ 'proxy.addressPlaceholder': { zh: '例:10.0.0.10:8081', en: 'e.g.
10.0.0.10:8081' },
+ 'proxy.addressRequired': { zh: '请输入 Proxy 地址', en: 'Please enter a Proxy
address' },
+ 'proxy.addAddressSuccess': { zh: 'Proxy 地址已新增', en: 'Proxy address added' },
+ 'proxy.addAddressFailed': { zh: '新增 Proxy 地址失败', en: 'Failed to add Proxy
address' },
+ 'proxy.removeAddressConfirm': {
+ zh: '确认删除 Proxy 地址 {addr}?',
+ en: 'Remove Proxy address {addr}?',
+ },
+ 'proxy.removeAddressSuccess': { zh: 'Proxy 地址已删除', en: 'Proxy address
removed' },
+ 'proxy.removeAddressFailed': { zh: '删除 Proxy 地址失败', en: 'Failed to remove
Proxy address' },
'proxy.reloadConfig': { zh: '重载配置', en: 'Reload Config' },
'proxy.reloadSuccess': { zh: '配置重载成功', en: 'Config reload succeeded' },
'proxy.reloadFailed': { zh: '配置重载失败', en: 'Config reload failed' },
@@ -1254,6 +1283,10 @@ const translations: Record<string, Record<Lang, string>>
= {
'producer.language': { zh: '语言', en: 'Language' },
'producer.selectTopic': { zh: '请选择 Topic', en: 'Select a topic' },
'producer.inputGroup': { zh: '请输入生产者组', en: 'Input producer group' },
+ 'producer.inputGroupOptional': {
+ zh: '输入生产者组,留空查询全部活跃组',
+ en: 'Input producer group, or leave empty for all active groups',
+ },
'producer.fetchTopicFailed': { zh: '获取 Topic 列表失败', en: 'Failed to fetch
topic list' },
'producer.fetchConnectionFailed': {
zh: '获取生产者连接失败',
diff --git a/web/src/pages/cluster/__tests__/ClusterPage.test.tsx
b/web/src/pages/cluster/__tests__/ClusterPage.test.tsx
index 2b1aef10e..12037bf15 100644
--- a/web/src/pages/cluster/__tests__/ClusterPage.test.tsx
+++ b/web/src/pages/cluster/__tests__/ClusterPage.test.tsx
@@ -31,6 +31,7 @@ const clusterServiceMocks = vi.hoisted(() => ({
listK8sCerts: vi.fn(),
listNameserverRegistry: vi.fn(),
listRegistryClusters: vi.fn(),
+ previewClusterConfig: vi.fn(),
restartProxy: vi.fn(),
testClusterConnection: vi.fn(),
updateClusterConfig: vi.fn(),
@@ -255,6 +256,32 @@ describe('Cluster page', () => {
]);
clusterServiceMocks.restartProxy.mockReset().mockResolvedValue(undefined);
clusterServiceMocks.testClusterConnection.mockReset();
+
clusterServiceMocks.previewClusterConfig.mockReset().mockImplementation(async
(request) => {
+ const cluster = buildCluster();
+ return {
+ cluster,
+ currentConfig: cluster.config,
+ proposedConfig: { ...cluster.config, ...request },
+ targetBrokers: cluster.brokers.map((broker) => ({
+ name: broker.name,
+ address: broker.addr,
+ })),
+ brokerProperties: {
+ ...(request.writeQueueNums != null
+ ? { defaultTopicQueueNums: String(request.writeQueueNums) }
+ : {}),
+ },
+ changes: [
+ {
+ field: 'writeQueueNums',
+ currentValue: String(cluster.config.writeQueueNums),
+ proposedValue: String(request.writeQueueNums),
+ brokerProperty: 'defaultTopicQueueNums',
+ },
+ ],
+ changed: true,
+ };
+ });
clusterServiceMocks.updateClusterConfig.mockReset().mockImplementation(async ()
=> {
const cluster = buildCluster();
return {
@@ -316,6 +343,35 @@ describe('Cluster page', () => {
expect(within(dialog).getByText('8080')).toBeInTheDocument();
});
+ it('previews broker config changes before submitting the update', async ()
=> {
+ const user = userEvent.setup();
+ renderWithProviders(<ClusterPage />);
+
+ const brokerRow = await screen.findByRole('row', { name:
/10\.101\.2\.11:10911/ });
+ await user.click(within(brokerRow).getByRole('button', { name: /配\s*置/ }));
+ const dialog = await screen.findByRole('dialog', { name: /配置 -
rocketmq-prod/ });
+ const writeQueuesInput = within(dialog).getByLabelText('写队列数');
+ await user.clear(writeQueuesInput);
+ await user.type(writeQueuesInput, '16');
+
+ await user.click(within(dialog).getByRole('button', { name: /预\s*览/ }));
+
+ await waitFor(() =>
+ expect(clusterServiceMocks.previewClusterConfig).toHaveBeenCalledWith(
+ expect.objectContaining({
+ id: 'cluster-prod',
+ instanceId: 'instance-1',
+ writeQueueNums: 16,
+ maxMessageSize: 4 * 1024 * 1024,
+ }),
+ ),
+ );
+ expect(clusterServiceMocks.updateClusterConfig).not.toHaveBeenCalled();
+ expect(within(dialog).getByText('10.101.2.11:10911')).toBeInTheDocument();
+
expect(within(dialog).getByText('defaultTopicQueueNums=16')).toBeInTheDocument();
+ expect(within(dialog).getByRole('row', { name: /写队列数/
})).toHaveTextContent('16');
+ });
+
it('keeps cluster tabs usable when address fields are missing', async () => {
const user = userEvent.setup();
const submitSearch = async (placeholder: string, value: string) => {
diff --git a/web/src/pages/cluster/index.tsx b/web/src/pages/cluster/index.tsx
index e090fd700..254882f4a 100644
--- a/web/src/pages/cluster/index.tsx
+++ b/web/src/pages/cluster/index.tsx
@@ -56,6 +56,8 @@ import type {
ProxyInfo,
NameserverRegistryEntry,
ClusterConfig,
+ ClusterConfigPreviewChange,
+ ClusterConfigPreviewResult,
ClusterInfo,
ClusterProbeResult,
} from '../../api/cluster';
@@ -66,6 +68,7 @@ import {
listK8sCerts,
listNameserverRegistry,
listRegistryClusters,
+ previewClusterConfig,
restartProxy,
testClusterConnection,
updateClusterConfig,
@@ -83,12 +86,25 @@ const REFRESH_INTERVAL_MS = 2000;
type RefreshSource = 'initial' | 'manual' | 'operation' | 'background';
type ProxyDetail = ProxyInfo & { clusterId: string; clusterName: string;
nsClusterName: string };
+type ClusterConfigFormValues = Partial<ClusterConfig> & { maxMessageSizeMB:
number };
+type ClusterConfigRequest = { id: string; instanceId?: string } &
Partial<ClusterConfig>;
const safeText = (value: string | null | undefined) => value ?? '';
const searchText = (value: string | null | undefined) =>
safeText(value).toLowerCase();
const compareText = (left: string | null | undefined, right: string | null |
undefined) =>
safeText(left).localeCompare(safeText(right));
+const CONFIG_FIELD_LABEL_KEYS: Record<string, string> = {
+ flushDiskType: 'cluster.flushDiskType',
+ autoCreateTopicEnable: 'cluster.autoCreateTopic',
+ autoCreateSubscriptionGroup: 'cluster.autoCreateSubGroup',
+ maxMessageSize: 'cluster.maxMessageSize',
+ fileReservedTime: 'cluster.fileReservedTime',
+ writeQueueNums: 'cluster.writeQueues',
+ readQueueNums: 'cluster.readQueues',
+ brokerPermission: 'cluster.brokerPermission',
+};
+
// ─── Page
─────────────────────────────────────────────────────────────────────
const ClusterPage = () => {
@@ -106,6 +122,9 @@ const ClusterPage = () => {
const [configModalOpen, setConfigModalOpen] = useState(false);
const [selectedCluster, setSelectedCluster] = useState<ClusterInfo |
null>(null);
+ const [configPreview, setConfigPreview] =
useState<ClusterConfigPreviewResult | null>(null);
+ const [configPreviewLoading, setConfigPreviewLoading] = useState(false);
+ const [configSubmitting, setConfigSubmitting] = useState(false);
const [nsRegistry, setNsRegistry] = useState<NameserverRegistryEntry[]>([]);
const [selectedProxy, setSelectedProxy] = useState<ProxyDetail | null>(null);
const [configForm] = Form.useForm();
@@ -477,6 +496,9 @@ const ClusterPage = () => {
const handleConfigOpen = (cluster: ClusterInfo) => {
const cfg: ClusterConfig = cluster.config ?? ({} as ClusterConfig);
setSelectedCluster(cluster);
+ setConfigPreview(null);
+ setConfigPreviewLoading(false);
+ setConfigSubmitting(false);
configForm.setFieldsValue({
flushDiskType: cfg.flushDiskType ?? 'ASYNC_FLUSH',
autoCreateTopicEnable: cfg.autoCreateTopicEnable ?? false,
@@ -490,6 +512,151 @@ const ClusterPage = () => {
setConfigModalOpen(true);
};
+ const buildConfigUpdateRequest = (
+ values: ClusterConfigFormValues,
+ ): ClusterConfigRequest | null => {
+ if (!selectedCluster) return null;
+ const { maxMessageSizeMB, ...configValues } = values;
+ return {
+ id: selectedCluster.id,
+ instanceId: selectedInstanceIdRef.current,
+ ...(selectedCluster.config ?? {}),
+ ...configValues,
+ maxMessageSize: maxMessageSizeMB * 1048576,
+ };
+ };
+
+ const handleConfigPreview = async () => {
+ let values: ClusterConfigFormValues;
+ try {
+ values = await configForm.validateFields();
+ } catch {
+ return;
+ }
+ const request = buildConfigUpdateRequest(values);
+ if (!request) return;
+
+ setConfigPreviewLoading(true);
+ try {
+ const preview = await previewClusterConfig(request);
+ setConfigPreview(preview);
+ message.success(t('cluster.configPreviewGenerated'));
+ } catch {
+ setConfigPreview(null);
+ message.error(t('cluster.configPreviewFailed'));
+ } finally {
+ setConfigPreviewLoading(false);
+ }
+ };
+
+ const handleConfigSubmit = async () => {
+ let values: ClusterConfigFormValues;
+ try {
+ values = await configForm.validateFields();
+ } catch {
+ return;
+ }
+ const request = buildConfigUpdateRequest(values);
+ if (!request) return;
+
+ setConfigSubmitting(true);
+ try {
+ const result = await updateClusterConfig(request);
+ if (result.status === 'SUCCESS') {
+ await requestRefresh('operation');
+ message.success(t('cluster.configUpdated'));
+ setConfigModalOpen(false);
+ setConfigPreview(null);
+ return;
+ }
+
+ const failedAddresses = result.failedBrokers.map((failure) =>
failure.address).join(', ');
+ if (result.status === 'PARTIAL') {
+ await requestRefresh('operation');
+ message.warning(t('cluster.configPartiallyUpdated', { brokers:
failedAddresses }));
+ return;
+ }
+ message.error(t('cluster.configUpdateFailed', { brokers: failedAddresses
}));
+ } catch {
+ message.error(t('cluster.configUpdateFailed', { brokers: '' }));
+ } finally {
+ setConfigSubmitting(false);
+ }
+ };
+
+ const configFieldLabel = (field: string) => {
+ const key = CONFIG_FIELD_LABEL_KEYS[field];
+ return key ? t(key) : field;
+ };
+
+ const previewValue = (value: string | null) => value ?? '-';
+
+ const renderConfigPreview = () => {
+ if (!configPreview) return null;
+ const propertyEntries = Object.entries(configPreview.brokerProperties ??
{});
+ const previewColumns: ColumnsType<ClusterConfigPreviewChange> = [
+ {
+ title: t('cluster.configPreviewField'),
+ dataIndex: 'field',
+ key: 'field',
+ render: (field: string) => configFieldLabel(field),
+ },
+ {
+ title: t('cluster.configPreviewCurrent'),
+ dataIndex: 'currentValue',
+ key: 'currentValue',
+ render: previewValue,
+ },
+ {
+ title: t('cluster.configPreviewProposed'),
+ dataIndex: 'proposedValue',
+ key: 'proposedValue',
+ render: previewValue,
+ },
+ {
+ title: t('cluster.configPreviewProperty'),
+ dataIndex: 'brokerProperty',
+ key: 'brokerProperty',
+ },
+ ];
+
+ return (
+ <Card size="small" title={t('cluster.configPreview')} style={{
marginTop: 16 }}>
+ <Descriptions size="small" column={1}>
+ <Descriptions.Item label={t('cluster.configPreviewTargets')}>
+ <Space size={[0, 4]} wrap>
+ {configPreview.targetBrokers.length > 0 ? (
+ configPreview.targetBrokers.map((broker) => (
+ <Tag key={broker.address}>{broker.address}</Tag>
+ ))
+ ) : (
+ <Text type="secondary">-</Text>
+ )}
+ </Space>
+ </Descriptions.Item>
+ <Descriptions.Item
label={t('cluster.configPreviewBrokerProperties')}>
+ <Space size={[0, 4]} wrap>
+ {propertyEntries.length > 0 ? (
+ propertyEntries.map(([key, value]) => <Tag
key={key}>{`${key}=${value}`}</Tag>)
+ ) : (
+ <Text type="secondary">-</Text>
+ )}
+ </Space>
+ </Descriptions.Item>
+ </Descriptions>
+ <Table<ClusterConfigPreviewChange>
+ columns={previewColumns}
+ dataSource={configPreview.changes}
+ rowKey="field"
+ pagination={false}
+ size="small"
+ locale={{ emptyText: t('cluster.configPreviewNoChanges') }}
+ style={{ marginTop: 12 }}
+ />
+ </Card>
+ );
+ };
+
const tabItems = [
{
key: 'nameserver',
@@ -703,48 +870,24 @@ const ClusterPage = () => {
<Modal
title={t('cluster.configTitle', { name: selectedCluster.name })}
open={configModalOpen}
- onCancel={() => setConfigModalOpen(false)}
- onOk={() => {
- configForm.validateFields().then(async (values) => {
- if (!selectedCluster) return;
- try {
- const { maxMessageSizeMB, ...configValues } = values;
- const nextConfig: ClusterConfig = {
- ...(selectedCluster.config ?? {}),
- ...configValues,
- maxMessageSize: maxMessageSizeMB * 1048576,
- };
- const result = await updateClusterConfig({
- id: selectedCluster.id,
- instanceId: selectedInstanceIdRef.current,
- ...nextConfig,
- });
- if (result.status === 'SUCCESS') {
- await requestRefresh('operation');
- message.success(t('cluster.configUpdated'));
- setConfigModalOpen(false);
- return;
- }
-
- const failedAddresses = result.failedBrokers
- .map((failure) => failure.address)
- .join(', ');
- if (result.status === 'PARTIAL') {
- await requestRefresh('operation');
- message.warning(
- t('cluster.configPartiallyUpdated', { brokers:
failedAddresses }),
- );
- return;
- }
- message.error(t('cluster.configUpdateFailed', { brokers:
failedAddresses }));
- } catch {
- message.error(t('cluster.configUpdateFailed', { brokers: ''
}));
- }
- });
+ onCancel={() => {
+ setConfigModalOpen(false);
+ setConfigPreview(null);
}}
- width={560}
+ onOk={() => void handleConfigSubmit()}
+ confirmLoading={configSubmitting}
+ width={720}
>
- <Form form={configForm} layout="vertical">
+ <Space style={{ marginBottom: 16 }}>
+ <Button
+ icon={<EyeOutlined />}
+ loading={configPreviewLoading}
+ onClick={() => void handleConfigPreview()}
+ >
+ {t('cluster.configPreview')}
+ </Button>
+ </Space>
+ <Form form={configForm} layout="vertical" onValuesChange={() =>
setConfigPreview(null)}>
<Form.Item label={t('cluster.flushDiskType')}
name="flushDiskType">
<Radio.Group>
<Radio value="SYNC_FLUSH">{t('cluster.syncFlush')}</Radio>
@@ -785,6 +928,7 @@ const ClusterPage = () => {
<InputNumber min={0} max={7} style={{ width: '100%' }} />
</Form.Item>
</Form>
+ {renderConfigPreview()}
</Modal>
)}
</div>
diff --git a/web/src/pages/instance/__tests__/AclPage.test.tsx
b/web/src/pages/instance/__tests__/AclPage.test.tsx
index 2d0f86287..cf5befec8 100644
--- a/web/src/pages/instance/__tests__/AclPage.test.tsx
+++ b/web/src/pages/instance/__tests__/AclPage.test.tsx
@@ -21,6 +21,7 @@ import userEvent from '@testing-library/user-event';
import type React from 'react';
import { MemoryRouter } from 'react-router-dom';
import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest';
+import type { AclRule, AclUser } from '../../../api/acl';
import { LangProvider } from '../../../i18n/LangContext';
import * as aclService from '../../../services/aclService';
import * as instanceService from '../../../services/instanceService';
@@ -197,6 +198,118 @@ describe('ACL page', () => {
);
});
+ it('uses Tencent role names for role-backed ACL rules and users', async ()
=> {
+ const user = userEvent.setup();
+ vi.mocked(instanceService.listInstances).mockResolvedValue([
+ {
+ id: 21,
+ name: 'tencent-rmq',
+ type: 'CLOUD',
+ endpoint: 'vpc.tencent:8080',
+ vendor: 'TENCENT',
+ cloudInstanceId: 'rmq-cloud',
+ remark: '',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ gmtCreate: '',
+ gmtModified: '',
+ },
+ ]);
+ vi.mocked(aclService.listAclRules).mockResolvedValue({
+ items: [
+ {
+ principal: 'reader-role',
+ resource: '*',
+ resourceType: 'Cluster',
+ resourcePattern: 'LITERAL',
+ actions: ['PUB'],
+ decision: 'ALLOW',
+ scope: 'cluster',
+ aclVersion: '1.0',
+ gmtCreate: '2026-07-23T00:00:00Z',
+ } as AclRule,
+ ],
+ total: 1,
+ page: 1,
+ size: 20,
+ });
+ vi.mocked(aclService.pageAclUsers).mockResolvedValue({
+ items: [
+ {
+ username: 'reader-role',
+ accessKey: 'acce****3456',
+ secretKey: 'secr****7654',
+ admin: false,
+ clusters: ['rmq-cloud'],
+ gmtCreate: '2026-07-23T00:00:00Z',
+ } as AclUser,
+ ],
+ total: 1,
+ page: 1,
+ size: 20,
+ });
+ vi.mocked(aclService.updateAclRule).mockResolvedValue({
+ id: 'reader-role',
+ principal: 'reader-role',
+ resource: '*',
+ resourceType: 'Cluster',
+ resourcePattern: 'LITERAL',
+ actions: ['PUB', 'SUB'],
+ decision: 'ALLOW',
+ scope: 'cluster',
+ aclVersion: '1.0',
+ gmtCreate: '2026-07-23T00:00:00Z',
+ });
+ vi.mocked(aclService.getAclUserCredentials).mockResolvedValue({
+ id: 'reader-role',
+ username: 'reader-role',
+ accessKey: 'full-access-key',
+ secretKey: 'full-secret-key',
+ admin: false,
+ clusters: ['rmq-cloud'],
+ gmtCreate: '2026-07-23T00:00:00Z',
+ });
+
+ renderWithProviders(<AclPage />, '/instance/tencent-rmq/acl');
+
+ await waitFor(() =>
+
expect(screen.getByTestId('acl-local-metadata-notice')).toHaveTextContent(
+ 'Tencent Cloud Role',
+ ),
+ );
+ const ruleRow = await screen.findByRole('row', { name: /reader-role/ });
+ await user.click(within(ruleRow).getByRole('button', { name: /编辑/ }));
+ const ruleDialog = await screen.findByRole('dialog');
+ expect(within(ruleDialog).getByLabelText('资源名称')).toBeDisabled();
+ await user.click(within(ruleDialog).getByLabelText('订阅 (SUB)'));
+ await user.click(within(ruleDialog).getByRole('button', { name: /保\s*存/
}));
+
+ await waitFor(() =>
expect(aclService.updateAclRule).toHaveBeenCalledTimes(1));
+ expect(aclService.updateAclRule).toHaveBeenCalledWith(
+ expect.objectContaining({
+ id: 'reader-role',
+ principal: 'reader-role',
+ resource: '*',
+ resourceType: 'Cluster',
+ resourcePattern: 'LITERAL',
+ decision: 'ALLOW',
+ scope: 'cluster',
+ instanceId: 'tencent-rmq',
+ }),
+ );
+
+ await user.click(screen.getByText('用户管理'));
+ const userRow = await screen.findByRole('row', { name: /reader-role/ });
+ const secretCell = within(userRow).getByText('••••••••••••').closest('td');
+ expect(secretCell).not.toBeNull();
+ await user.click(within(secretCell as HTMLElement).getByRole('button'));
+
+ await waitFor(() =>
+
expect(aclService.getAclUserCredentials).toHaveBeenCalledWith('reader-role',
'tencent-rmq'),
+ );
+ expect(await
within(userRow).findByText('full-secret-key')).toBeInTheDocument();
+ });
+
it('renders backend users on the user tab', async () => {
const user = userEvent.setup();
renderWithProviders(<AclPage />);
diff --git a/web/src/pages/instance/acl.tsx b/web/src/pages/instance/acl.tsx
index 1800d5a38..d06bf5416 100644
--- a/web/src/pages/instance/acl.tsx
+++ b/web/src/pages/instance/acl.tsx
@@ -67,6 +67,7 @@ import type { AclRule, AclUser, AclClusterConfig,
PlainAccessConfig } from '../.
import { useInstanceFilter } from '../../hooks/useInstanceFilter';
import { tableScrollX } from '../../utils/table';
+type AclEntityId = AclRule['id'];
type AclRuleFormValues = Pick<
AclRule,
'principal' | 'resource' | 'resourceType' | 'resourcePattern' | 'actions' |
'decision' | 'scope'
@@ -75,6 +76,7 @@ type AclUserFormValues = Pick<AclUser, 'username' | 'admin' |
'clusters'>;
const normalizeRule = (rule: AclRule): AclRule => ({
...rule,
+ id: rule.id ?? rule.principal,
principal: rule.principal ?? '',
resource: rule.resource ?? '',
resourceType: rule.resourceType ?? '',
@@ -90,6 +92,7 @@ const normalizeRule = (rule: AclRule): AclRule => ({
const normalizeUser = (user: AclUser): AclUser => ({
...user,
+ id: user.id ?? user.username,
username: user.username ?? '',
accessKey: user.accessKey ?? '',
secretKey: user.secretKey ?? '',
@@ -117,6 +120,8 @@ const AclPageContent = ({
}: AclPageContentProps) => {
const { t } = useLang();
const hasSelectedInstance = Boolean(selectedInstanceId);
+ const selectedInstance = instances.find((instance) => instance.name ===
selectedInstanceId);
+ const tencentRoleMode = selectedInstance?.vendor === 'TENCENT';
/* ─── State ─── */
const [rules, setRules] = useState<AclRule[]>([]);
@@ -153,12 +158,12 @@ const AclPageContent = ({
const [userForm] = Form.useForm();
// Secret key reveal
- const [revealedKeys, setRevealedKeys] = useState<Set<number>>(new Set());
- const [adminUpdatingIds, setAdminUpdatingIds] = useState<Set<number>>(() =>
new Set());
- const adminUpdateInFlightRef = useRef<Set<number>>(new Set());
- const revealRequestGenerationRef = useRef<Record<number, number>>({});
+ const [revealedKeys, setRevealedKeys] = useState<Set<AclEntityId>>(new
Set());
+ const [adminUpdatingIds, setAdminUpdatingIds] =
useState<Set<AclEntityId>>(() => new Set());
+ const adminUpdateInFlightRef = useRef<Set<AclEntityId>>(new Set());
+ const revealRequestGenerationRef = useRef<Record<string, number>>({});
const [credentialsByUser, setCredentialsByUser] = useState<
- Record<number, { accessKey: string; secretKey: string }>
+ Record<string, { accessKey: string; secretKey: string }>
>({});
// Cluster ACL config (examineBrokerClusterAclConfig)
@@ -256,7 +261,9 @@ const AclPageContent = ({
setEditingRule(null);
ruleForm.resetFields();
ruleForm.setFieldsValue({
- resourcePattern: 'PREFIX',
+ resourceType: tencentRoleMode ? 'Cluster' : undefined,
+ resource: tencentRoleMode ? '*' : undefined,
+ resourcePattern: tencentRoleMode ? 'LITERAL' : 'PREFIX',
actions: ['PUB'],
decision: 'ALLOW',
scope: 'cluster',
@@ -281,18 +288,28 @@ const AclPageContent = ({
const handleRuleSubmit = async () => {
try {
const values = (await ruleForm.validateFields()) as AclRuleFormValues;
+ const normalizedValues = tencentRoleMode
+ ? {
+ ...values,
+ resourceType: 'Cluster',
+ resource: '*',
+ resourcePattern: 'LITERAL',
+ decision: 'ALLOW',
+ scope: 'cluster',
+ }
+ : values;
setRuleSubmitting(true);
if (editingRule) {
await updateAclRule({
...editingRule,
- ...values,
+ ...normalizedValues,
instanceId: selectedInstanceId,
});
message.success(t('acl.ruleUpdated'));
} else {
await createAclRule({
- ...values,
- aclVersion: '2.0',
+ ...normalizedValues,
+ aclVersion: tencentRoleMode ? '1.0' : '2.0',
instanceId: selectedInstanceId,
});
setRulePage(1);
@@ -308,7 +325,7 @@ const AclPageContent = ({
}
};
- const handleDeleteRule = async (id: number) => {
+ const handleDeleteRule = async (id: AclEntityId) => {
try {
await deleteAclRule(id, selectedInstanceId);
if (rules.length === 1 && rulePage > 1) {
@@ -323,10 +340,11 @@ const AclPageContent = ({
};
/* ─── User helpers ─── */
- const toggleRevealKey = async (userId: number) => {
+ const toggleRevealKey = async (userId: AclEntityId) => {
+ const userKey = String(userId);
const revealing = !revealedKeys.has(userId);
- const revealGeneration = (revealRequestGenerationRef.current[userId] ?? 0)
+ 1;
- revealRequestGenerationRef.current[userId] = revealGeneration;
+ const revealGeneration = (revealRequestGenerationRef.current[userKey] ??
0) + 1;
+ revealRequestGenerationRef.current[userKey] = revealGeneration;
setRevealedKeys((prev) => {
const next = new Set(prev);
if (next.has(userId)) {
@@ -336,19 +354,19 @@ const AclPageContent = ({
}
return next;
});
- if (!revealing || credentialsByUser[userId]) return;
+ if (!revealing || credentialsByUser[userKey]) return;
try {
const credentials = await getAclUserCredentials(userId,
selectedInstanceId);
- if (revealRequestGenerationRef.current[userId] !== revealGeneration)
return;
+ if (revealRequestGenerationRef.current[userKey] !== revealGeneration)
return;
setCredentialsByUser((prev) => ({
...prev,
- [userId]: {
+ [userKey]: {
accessKey: credentials.accessKey,
secretKey: credentials.secretKey,
},
}));
} catch {
- if (revealRequestGenerationRef.current[userId] !== revealGeneration)
return;
+ if (revealRequestGenerationRef.current[userKey] !== revealGeneration)
return;
setRevealedKeys((prev) => {
const next = new Set(prev);
next.delete(userId);
@@ -409,7 +427,7 @@ const AclPageContent = ({
}
};
- const handleDeleteUser = async (id: number) => {
+ const handleDeleteUser = async (id: AclEntityId) => {
try {
await deleteAclUser(id, selectedInstanceId);
setUsers((prev) => prev.filter((u) => u.id !== id));
@@ -420,6 +438,7 @@ const AclPageContent = ({
};
const handleToggleAdmin = async (user: AclUser, checked: boolean) => {
+ if (tencentRoleMode) return;
if (adminUpdateInFlightRef.current.has(user.id)) return;
adminUpdateInFlightRef.current.add(user.id);
setAdminUpdatingIds((current) => new Set(current).add(user.id));
@@ -705,7 +724,7 @@ const AclPageContent = ({
sorter: (a, b) => a.accessKey.localeCompare(b.accessKey),
render: (text: string, record: AclUser) => {
const revealed = revealedKeys.has(record.id);
- const fullAccessKey = credentialsByUser[record.id]?.accessKey ?? text;
+ const fullAccessKey = credentialsByUser[String(record.id)]?.accessKey
?? text;
return (
<Space size={8}>
<Typography.Text
@@ -725,7 +744,7 @@ const AclPageContent = ({
width: 240,
render: (_: string, record: AclUser) => {
const revealed = revealedKeys.has(record.id);
- const secret = credentialsByUser[record.id]?.secretKey;
+ const secret = credentialsByUser[String(record.id)]?.secretKey;
return (
<Space size={8}>
<Typography.Text
@@ -755,7 +774,7 @@ const AclPageContent = ({
checked={val}
size="small"
loading={adminUpdatingIds.has(record.id)}
- disabled={adminUpdatingIds.has(record.id)}
+ disabled={tencentRoleMode || adminUpdatingIds.has(record.id)}
onChange={(checked) => handleToggleAdmin(record, checked)}
/>
),
@@ -935,8 +954,10 @@ const AclPageContent = ({
<InfoBanner
data-testid="acl-local-metadata-notice"
- title={t('acl.localMetadataNotice')}
- description={t('acl.localMetadataDescription')}
+ title={t(tencentRoleMode ? 'acl.tencentRoleNotice' :
'acl.localMetadataNotice')}
+ description={t(
+ tencentRoleMode ? 'acl.tencentRoleDescription' :
'acl.localMetadataDescription',
+ )}
/>
<Card variant="borderless" styles={{ body: { padding: 0 } }}>
@@ -1250,11 +1271,16 @@ const AclPageContent = ({
>
<Select
placeholder={t('acl.selectResourceType')}
- options={[
- { value: 'Topic', label: 'Topic' },
- { value: 'Group', label: 'Group' },
- { value: 'Cluster', label: 'Cluster' },
- ]}
+ disabled={tencentRoleMode}
+ options={
+ tencentRoleMode
+ ? [{ value: 'Cluster', label: 'Cluster' }]
+ : [
+ { value: 'Topic', label: 'Topic' },
+ { value: 'Group', label: 'Group' },
+ { value: 'Cluster', label: 'Cluster' },
+ ]
+ }
/>
</Form.Item>
@@ -1265,11 +1291,11 @@ const AclPageContent = ({
{ required: true, message: t('acl.inputRequired', { field:
t('acl.resourceName') }) },
]}
>
- <Input placeholder={t('acl.resourceNamePlaceholder')} />
+ <Input placeholder={t('acl.resourceNamePlaceholder')}
disabled={tencentRoleMode} />
</Form.Item>
<Form.Item name="resourcePattern" label={t('acl.matchPattern')}>
- <Radio.Group>
+ <Radio.Group disabled={tencentRoleMode}>
<Radio.Button value="LITERAL">{t('acl.literal')}</Radio.Button>
<Radio.Button value="PREFIX">{t('acl.prefix')}</Radio.Button>
</Radio.Group>
@@ -1292,7 +1318,7 @@ const AclPageContent = ({
</Form.Item>
<Form.Item name="decision" label={t('acl.decision')}>
- <Radio.Group>
+ <Radio.Group disabled={tencentRoleMode}>
<Radio.Button value="ALLOW">
<span style={{ color: '#52c41a' }}>{t('acl.allowDesc')}</span>
</Radio.Button>
@@ -1305,6 +1331,7 @@ const AclPageContent = ({
<Form.Item name="scope" label={t('acl.effectScope')}>
<Select
placeholder={t('acl.selectEffectScope')}
+ disabled={tencentRoleMode}
options={[{ value: 'cluster', label: t('acl.clusterScope') }]}
/>
</Form.Item>
@@ -1339,7 +1366,11 @@ const AclPageContent = ({
</Form.Item>
<Form.Item name="admin" label={t('acl.admin')}
valuePropName="checked">
- <Switch checkedChildren={t('common.yes')}
unCheckedChildren={t('common.no')} />
+ <Switch
+ checkedChildren={t('common.yes')}
+ unCheckedChildren={t('common.no')}
+ disabled={tencentRoleMode}
+ />
</Form.Item>
<Form.Item name="clusters" label={t('acl.associatedClusters')}>
@@ -1347,6 +1378,7 @@ const AclPageContent = ({
mode="tags"
tokenSeparators={[',']}
allowClear
+ disabled={tencentRoleMode}
options={instances.map((instance) => ({
value: instance.cloudInstanceId ?? instance.name,
label: instance.name,
diff --git a/web/src/pages/studio/Producer.tsx
b/web/src/pages/studio/Producer.tsx
index 50c514656..edd7a0e76 100644
--- a/web/src/pages/studio/Producer.tsx
+++ b/web/src/pages/studio/Producer.tsx
@@ -187,7 +187,7 @@ const ProducerPage = () => {
}
};
- const onFinish = async (values: { selectedTopic: string; producerGroup:
string }) => {
+ const onFinish = async (values: { selectedTopic: string; producerGroup?:
string }) => {
if (queryInFlightRef.current !== null) return;
if (!selectedInstanceId) {
message.error('Select an instance before querying producer
connections.');
@@ -227,6 +227,13 @@ const ProducerPage = () => {
const columns = [
{ title: 'Client ID', dataIndex: 'clientId', key: 'clientId', align:
'center' as const },
+ {
+ title: 'Producer Group',
+ dataIndex: 'producerGroup',
+ key: 'producerGroup',
+ align: 'center' as const,
+ render: (value?: string) => value || '-',
+ },
{
title: t('common.address'),
dataIndex: 'clientAddr',
@@ -276,8 +283,8 @@ const ProducerPage = () => {
const rows = connectionList.map((connection) => ({
...connection,
instanceId: selectedInstanceId ?? '',
- topic: selectedTopic,
- producerGroup,
+ topic: connection.topic ?? selectedTopic,
+ producerGroup: connection.producerGroup ?? producerGroup,
readiness: connectionSummary?.readiness,
warnings,
}));
@@ -330,14 +337,10 @@ const ProducerPage = () => {
options={topicList.map((topic) => ({ value: topic, label: topic
}))}
/>
</Form.Item>
- <Form.Item
- label="PRODUCER GROUP"
- name="producerGroup"
- rules={[{ required: true, whitespace: true, message:
t('producer.inputGroup') }]}
- >
+ <Form.Item label="PRODUCER GROUP" name="producerGroup">
<AutoComplete
allowClear
- placeholder={t('producer.inputGroup')}
+ placeholder={t('producer.inputGroupOptional')}
style={{ width: 300 }}
options={producerGroups.map((group) => ({ value: group }))}
onFocus={() => {
@@ -435,7 +438,9 @@ const ProducerPage = () => {
<Table
dataSource={connectionList}
columns={columns}
- rowKey={(record) => `${record.clientId}:${record.clientAddr}`}
+ rowKey={(record) =>
+ `${record.producerGroup ??
''}:${record.clientId}:${record.clientAddr}`
+ }
pagination={false}
bordered
size="middle"
diff --git a/web/src/pages/studio/Proxy.tsx b/web/src/pages/studio/Proxy.tsx
index c7665bd4f..dee270586 100644
--- a/web/src/pages/studio/Proxy.tsx
+++ b/web/src/pages/studio/Proxy.tsx
@@ -33,6 +33,7 @@ import {
App,
Typography,
Input,
+ Popconfirm,
} from 'antd';
import type { ColumnsType } from 'antd/es/table';
import {
@@ -42,13 +43,18 @@ import {
CheckCircle,
XCircle,
Warning,
+ Plus,
+ Trash,
} from '@phosphor-icons/react';
import PageHeader from '../../components/PageHeader';
import { useLang } from '../../i18n/LangContext';
import {
+ addProxyAddress,
getProxyTopology,
+ removeProxyAddress,
queryProxyHomePage,
reloadProxyConfig,
+ type ProxyHomePageData,
type ProxyNode,
} from '../../api/proxy';
@@ -71,6 +77,9 @@ const ProxyPage: React.FC = () => {
const [proxyNodes, setProxyNodes] = useState<ProxyNode[]>([]);
const [selectedNode, setSelectedNode] = useState<ProxyNode | null>(null);
const [configModalOpen, setConfigModalOpen] = useState(false);
+ const [newProxyAddress, setNewProxyAddress] = useState('');
+ const [addressMutationLoading, setAddressMutationLoading] = useState(false);
+ const [removingProxyAddress, setRemovingProxyAddress] = useState<string |
null>(null);
const [clusterId, setClusterId] = useState<string>(
localStorage.getItem('clusterId') || 'DefaultCluster',
);
@@ -83,12 +92,8 @@ const ProxyPage: React.FC = () => {
totalTPS: null as number | null,
});
- const loadProxyNodes = useCallback(async () => {
- const requestId = ++loadRequestId.current;
- setLoading(true);
- try {
- const { proxyAddrList, currentProxyAddr } = await queryProxyHomePage();
- if (requestId !== loadRequestId.current) return false;
+ const applyProxyHome = useCallback(
+ async ({ proxyAddrList, currentProxyAddr }: ProxyHomePageData, requestId:
number) => {
const baseNodes: ProxyNode[] = (proxyAddrList || []).map((addr) => ({
key: addr,
address: addr,
@@ -107,9 +112,7 @@ const ProxyPage: React.FC = () => {
try {
const topology = await getProxyTopology();
if (requestId !== loadRequestId.current) return false;
- const statusByAddr = new Map(
- topology.map((node) => [node.proxyAddr, node.status]),
- );
+ const statusByAddr = new Map(topology.map((node) => [node.proxyAddr,
node.status]));
nodes = baseNodes.map((node) => {
const probeStatus = statusByAddr.get(node.address);
const status: ProxyNode['status'] =
@@ -136,6 +139,17 @@ const ProxyPage: React.FC = () => {
persistProxyAddress(currentProxyAddr || proxyAddrList?.[0]);
return true;
+ },
+ [],
+ );
+
+ const loadProxyNodes = useCallback(async () => {
+ const requestId = ++loadRequestId.current;
+ setLoading(true);
+ try {
+ const home = await queryProxyHomePage();
+ if (requestId !== loadRequestId.current) return false;
+ return await applyProxyHome(home, requestId);
} catch {
if (requestId !== loadRequestId.current) return false;
message.error(t('proxy.fetchListFailed'));
@@ -145,7 +159,7 @@ const ProxyPage: React.FC = () => {
setLoading(false);
}
}
- }, [message, t]);
+ }, [applyProxyHome, message, t]);
useEffect(() => {
const requestId = loadRequestId.current;
@@ -175,6 +189,58 @@ const ProxyPage: React.FC = () => {
}
};
+ const handleAddProxyAddress = async () => {
+ const addr = newProxyAddress.trim();
+ if (!addr) {
+ message.warning(t('proxy.addressRequired'));
+ return;
+ }
+ const requestId = ++loadRequestId.current;
+ setAddressMutationLoading(true);
+ setLoading(true);
+ try {
+ const home = await addProxyAddress(addr);
+ if (requestId !== loadRequestId.current) return;
+ await applyProxyHome(home, requestId);
+ setNewProxyAddress('');
+ message.success(t('proxy.addAddressSuccess'));
+ } catch {
+ if (requestId === loadRequestId.current) {
+ message.error(t('proxy.addAddressFailed'));
+ }
+ } finally {
+ if (requestId === loadRequestId.current) {
+ setAddressMutationLoading(false);
+ setLoading(false);
+ }
+ }
+ };
+
+ const handleRemoveProxyAddress = async (addr: string) => {
+ const requestId = ++loadRequestId.current;
+ setRemovingProxyAddress(addr);
+ setLoading(true);
+ try {
+ const home = await removeProxyAddress(addr);
+ if (requestId !== loadRequestId.current) return;
+ await applyProxyHome(home, requestId);
+ if (selectedNode?.address === addr) {
+ setSelectedNode(null);
+ setConfigModalOpen(false);
+ }
+ message.success(t('proxy.removeAddressSuccess'));
+ } catch {
+ if (requestId === loadRequestId.current) {
+ message.error(t('proxy.removeAddressFailed'));
+ }
+ } finally {
+ if (requestId === loadRequestId.current) {
+ setRemovingProxyAddress(null);
+ setLoading(false);
+ }
+ }
+ };
+
const handleReloadConfig = async (node: ProxyNode) => {
try {
const result = await reloadProxyConfig(clusterId, node.address);
@@ -340,6 +406,22 @@ const ProxyPage: React.FC = () => {
onClick={() => handleReloadConfig(record)}
/>
</Tooltip>
+ <Popconfirm
+ title={t('proxy.removeAddressConfirm', { addr: record.address })}
+ okText={t('common.confirm')}
+ cancelText={t('common.cancel')}
+ disabled={removingProxyAddress === record.address}
+ onConfirm={() => void handleRemoveProxyAddress(record.address)}
+ >
+ <Button
+ type="link"
+ size="small"
+ danger
+ icon={<Trash size={14} />}
+ aria-label={t('common.delete')}
+ loading={removingProxyAddress === record.address}
+ />
+ </Popconfirm>
</Space>
),
},
@@ -354,6 +436,22 @@ const ProxyPage: React.FC = () => {
extra={
<Space>
+ <Input
+ placeholder={t('proxy.addressPlaceholder')}
+ value={newProxyAddress}
+ onChange={(e) => setNewProxyAddress(e.target.value)}
+ onPressEnter={() => void handleAddProxyAddress()}
+ style={{ width: 220 }}
+ aria-label={t('proxy.address')}
+ disabled={addressMutationLoading}
+ />
+ <Button
+ icon={<Plus size={14} />}
+ onClick={() => void handleAddProxyAddress()}
+ loading={addressMutationLoading}
+ >
+ {t('common.add')}
+ </Button>
<Input
placeholder={t('proxy.clusterIdPlaceholder')}
value={clusterId}
diff --git a/web/src/pages/studio/__tests__/Producer.test.tsx
b/web/src/pages/studio/__tests__/Producer.test.tsx
index 88255db61..31f6b7348 100644
--- a/web/src/pages/studio/__tests__/Producer.test.tsx
+++ b/web/src/pages/studio/__tests__/Producer.test.tsx
@@ -219,7 +219,19 @@ describe('ProducerPage', () => {
expect(screen.getByText('就绪')).toBeInTheDocument();
});
- it('does not query without a producer group', async () => {
+ it('queries all active producer groups when producer group is empty', async
() => {
+ vi.mocked(queryProducerConnection).mockResolvedValue(
+ producerResult([
+ {
+ clientId: 'producer-all-1',
+ clientAddr: '192.168.1.10',
+ topic: 'order-events',
+ producerGroup: 'pg-order',
+ language: 'JAVA',
+ versionDesc: '5.1.0',
+ },
+ ]),
+ );
const user = userEvent.setup();
renderWithProviders(<ProducerPage />);
@@ -231,10 +243,11 @@ describe('ProducerPage', () => {
);
await user.click(screen.getByRole('button', { name: /搜索/ }));
- expect(
- await screen.findByText('请输入生产者组', { selector:
'.ant-form-item-explain-error' }),
- ).toBeInTheDocument();
- expect(queryProducerConnection).not.toHaveBeenCalled();
+ await waitFor(() => {
+ expect(queryProducerConnection).toHaveBeenCalledWith('instance-1',
'order-events', undefined);
+ });
+ expect(await screen.findByText('producer-all-1')).toBeInTheDocument();
+ expect(screen.getByText('pg-order')).toBeInTheDocument();
});
it('renders producer connection warnings from the summary', async () => {
@@ -307,6 +320,8 @@ describe('ProducerPage', () => {
{
clientId: '@producer-a',
clientAddr: '192.168.1.10',
+ topic: 'order-events',
+ producerGroup: 'order-producer',
language: 'JAVA',
versionDesc: '5.1.0',
},
diff --git a/web/src/pages/studio/__tests__/Proxy.test.tsx
b/web/src/pages/studio/__tests__/Proxy.test.tsx
index bcaa8d8a6..9cdfbafdd 100644
--- a/web/src/pages/studio/__tests__/Proxy.test.tsx
+++ b/web/src/pages/studio/__tests__/Proxy.test.tsx
@@ -19,13 +19,20 @@ import { beforeAll, beforeEach, describe, expect, it, vi }
from 'vitest';
import { act, render, screen, waitFor } from '@testing-library/react';
import userEvent from '@testing-library/user-event';
import { App } from 'antd';
-import { queryProxyHomePage, reloadProxyConfig } from '../../../api/proxy';
+import {
+ addProxyAddress,
+ queryProxyHomePage,
+ reloadProxyConfig,
+ removeProxyAddress,
+} from '../../../api/proxy';
import { LangProvider } from '../../../i18n/LangContext';
import ProxyPage from '../Proxy';
vi.mock('../../../api/proxy', () => ({
+ addProxyAddress: vi.fn(),
queryProxyHomePage: vi.fn(),
reloadProxyConfig: vi.fn(),
+ removeProxyAddress: vi.fn(),
}));
beforeAll(() => {
@@ -71,9 +78,11 @@ describe('ProxyPage', () => {
beforeEach(() => {
vi.clearAllMocks();
vi.mocked(queryProxyHomePage).mockResolvedValue(proxyHome);
+ vi.mocked(addProxyAddress).mockResolvedValue(proxyHome);
vi.mocked(reloadProxyConfig).mockResolvedValue({
success: true,
});
+ vi.mocked(removeProxyAddress).mockResolvedValue(proxyHome);
});
it('keeps discovered nodes when browser storage is unavailable', async () =>
{
@@ -155,6 +164,53 @@ describe('ProxyPage', () => {
expect(await screen.findByText('配置重载成功')).toBeInTheDocument();
});
+ it('adds a Proxy address and applies the updated address list', async () => {
+ const user = userEvent.setup();
+ vi.mocked(addProxyAddress).mockResolvedValueOnce({
+ proxyAddrList: ['127.0.0.1:8081', '10.0.0.10:8081'],
+ currentProxyAddr: '127.0.0.1:8081',
+ });
+ renderPage();
+ await screen.findByText('127.0.0.1:8081');
+
+ await user.type(screen.getByLabelText('Proxy 地址'), '10.0.0.10:8081');
+ await user.click(screen.getByRole('button', { name: '新增' }));
+
+ await waitFor(() =>
expect(addProxyAddress).toHaveBeenCalledWith('10.0.0.10:8081'));
+ expect(await screen.findByText('10.0.0.10:8081')).toBeInTheDocument();
+ expect(await screen.findByText('Proxy 地址已新增')).toBeInTheDocument();
+ });
+
+ it('rejects an empty Proxy address before calling the API', async () => {
+ const user = userEvent.setup();
+ renderPage();
+ await screen.findByText('127.0.0.1:8081');
+
+ await user.click(screen.getByRole('button', { name: '新增' }));
+
+ expect(addProxyAddress).not.toHaveBeenCalled();
+ expect(await screen.findByText('请输入 Proxy 地址')).toBeInTheDocument();
+ });
+
+ it('removes a Proxy address and applies the updated address list', async ()
=> {
+ const user = userEvent.setup();
+ vi.mocked(queryProxyHomePage).mockResolvedValueOnce({
+ proxyAddrList: ['127.0.0.1:8081', '10.0.0.10:8081'],
+ currentProxyAddr: '127.0.0.1:8081',
+ });
+ vi.mocked(removeProxyAddress).mockResolvedValueOnce(proxyHome);
+ renderPage();
+ await screen.findByText('10.0.0.10:8081');
+
+ const deleteButtons = screen.getAllByRole('button', { name: '删除' });
+ await user.click(deleteButtons[1]);
+ await user.click(await screen.findByRole('button', { name: /确\s*认/ }));
+
+ await waitFor(() =>
expect(removeProxyAddress).toHaveBeenCalledWith('10.0.0.10:8081'));
+ expect(screen.queryByText('10.0.0.10:8081')).not.toBeInTheDocument();
+ expect(await screen.findByText('Proxy 地址已删除')).toBeInTheDocument();
+ });
+
it('keeps the latest Proxy list when an older refresh resolves last', async
() => {
const older = createDeferred<typeof proxyHome>();
const latest = createDeferred<typeof proxyHome>();
diff --git a/web/src/services/aclService.ts b/web/src/services/aclService.ts
index 580ea995e..ad1042519 100644
--- a/web/src/services/aclService.ts
+++ b/web/src/services/aclService.ts
@@ -1,6 +1,7 @@
import { isMockMode } from './dataMode';
import * as aclApi from '../api/acl';
import type {
+ AclEntityId,
AclRule,
PageResult,
AclRuleQuery,
@@ -107,7 +108,10 @@ export async function pageAclUsers(params: {
return aclApi.pageAclUsers(params);
}
-export async function getAclUserCredentials(id: number, instanceId?: string):
Promise<AclUser> {
+export async function getAclUserCredentials(
+ id: AclEntityId,
+ instanceId?: string,
+): Promise<AclUser> {
if (isMockMode()) {
const user = aclUsersState.find((u) => u.id === id);
if (!user) throw new Error(`ACL user not found: ${id}`);
@@ -155,7 +159,7 @@ export async function updateAclRule(
return aclApi.updateAclRule(data);
}
-export async function deleteAclRule(id: number, instanceId?: string):
Promise<void> {
+export async function deleteAclRule(id: AclEntityId, instanceId?: string):
Promise<void> {
if (isMockMode()) {
const idx = aclRulesState.findIndex((rule) => rule.id === id);
if (idx >= 0) aclRulesState.splice(idx, 1);
@@ -200,7 +204,7 @@ export async function updateAclUser(
return aclApi.updateAclUser(data);
}
-export async function deleteAclUser(id: number, instanceId?: string):
Promise<void> {
+export async function deleteAclUser(id: AclEntityId, instanceId?: string):
Promise<void> {
if (isMockMode()) {
const idx = aclUsersState.findIndex((user) => user.id === id);
if (idx >= 0) aclUsersState.splice(idx, 1);
diff --git a/web/src/services/clusterService.test.ts
b/web/src/services/clusterService.test.ts
index 0ef869c39..614915b77 100644
--- a/web/src/services/clusterService.test.ts
+++ b/web/src/services/clusterService.test.ts
@@ -29,6 +29,7 @@ import {
getNameServerConfigDiff,
listClusters,
listK8sCerts,
+ previewClusterConfig,
updateClusterConfig,
updateK8sCert,
updateNameServer,
@@ -89,6 +90,30 @@ describe('clusterService mock clusters', () => {
]);
});
+ it('previews mock config updates without mutating the cluster', async () => {
+ const before = await getCluster('cluster-prod');
+
+ const preview = await previewClusterConfig({
+ id: before.id,
+ writeQueueNums: before.config.writeQueueNums + 1,
+ });
+
+ expect(preview.changed).toBe(true);
+ expect(preview.brokerProperties).toMatchObject({
+ defaultTopicQueueNums: String(before.config.writeQueueNums + 1),
+ });
+ expect(preview.targetBrokers.map((broker) => broker.address)).toEqual(
+ before.brokers.map((broker) => broker.addr),
+ );
+ expect(preview.changes).toEqual([
+ expect.objectContaining({
+ field: 'writeQueueNums',
+ brokerProperty: 'defaultTopicQueueNums',
+ }),
+ ]);
+ await expect(getCluster('cluster-prod')).resolves.toMatchObject({ config:
before.config });
+ });
+
it('persists partial mock config updates without copying id into config',
async () => {
const before = await getCluster('cluster-prod');
const originalConfig = { ...before.config };
diff --git a/web/src/services/clusterService.ts
b/web/src/services/clusterService.ts
index 9ed4a2a12..c147c4240 100644
--- a/web/src/services/clusterService.ts
+++ b/web/src/services/clusterService.ts
@@ -2,6 +2,7 @@ import { isMockMode } from './dataMode';
import * as clusterApi from '../api/cluster';
import type {
ClusterConfig,
+ ClusterConfigPreviewResult,
ClusterConfigUpdateResult,
ClusterInfo,
ClusterProbeResult,
@@ -210,7 +211,8 @@ export async function updateClusterConfig(
data: { id: string; instanceId?: string } & Partial<ClusterConfig>,
) {
if (isMockMode()) {
- const { id, ...config } = data;
+ const { id } = data;
+ const config = pickClusterConfig(data);
const cluster = getMockCluster(id);
Object.assign(cluster.config, config);
return {
@@ -223,11 +225,94 @@ export async function updateClusterConfig(
return clusterApi.updateClusterConfig(data);
}
+export async function previewClusterConfig(
+ data: { id: string; instanceId?: string } & Partial<ClusterConfig>,
+): Promise<ClusterConfigPreviewResult> {
+ if (isMockMode()) {
+ const { id } = data;
+ const config = pickClusterConfig(data);
+ const cluster = getMockCluster(id);
+ const currentConfig = { ...cluster.config };
+ const proposedConfig = { ...cluster.config, ...config };
+ const changes = buildMockConfigPreviewChanges(currentConfig,
proposedConfig);
+ return {
+ cluster: copyCluster(cluster),
+ currentConfig,
+ proposedConfig,
+ targetBrokers: cluster.brokers
+ .filter((broker) => broker.addr)
+ .map((broker) => ({ name: broker.name, address: broker.addr })),
+ brokerProperties: buildMockBrokerProperties(config),
+ changes,
+ changed: changes.length > 0,
+ };
+ }
+ return clusterApi.previewClusterConfig(data);
+}
+
export async function restartBroker(clusterId: string, brokerName: string) {
if (isMockMode()) return { success: true, message: `Broker ${brokerName}
restarted (mock)` };
return clusterApi.restartBroker(clusterId, brokerName);
}
+function pickClusterConfig(config: Partial<ClusterConfig>):
Partial<ClusterConfig> {
+ const picked: Partial<ClusterConfig> = {};
+ if (config.flushDiskType !== undefined) picked.flushDiskType =
config.flushDiskType;
+ if (config.autoCreateTopicEnable !== undefined) {
+ picked.autoCreateTopicEnable = config.autoCreateTopicEnable;
+ }
+ if (config.autoCreateSubscriptionGroup !== undefined) {
+ picked.autoCreateSubscriptionGroup = config.autoCreateSubscriptionGroup;
+ }
+ if (config.maxMessageSize !== undefined) picked.maxMessageSize =
config.maxMessageSize;
+ if (config.fileReservedTime !== undefined) picked.fileReservedTime =
config.fileReservedTime;
+ if (config.writeQueueNums !== undefined) picked.writeQueueNums =
config.writeQueueNums;
+ if (config.readQueueNums !== undefined) picked.readQueueNums =
config.readQueueNums;
+ if (config.brokerPermission !== undefined) picked.brokerPermission =
config.brokerPermission;
+ return picked;
+}
+
+function buildMockBrokerProperties(config: Partial<ClusterConfig>):
Record<string, string> {
+ const props: Record<string, string> = {};
+ if (config.flushDiskType != null) props.flushDiskType =
String(config.flushDiskType);
+ if (config.autoCreateTopicEnable != null) {
+ props.autoCreateTopicEnable = String(config.autoCreateTopicEnable);
+ }
+ if (config.autoCreateSubscriptionGroup != null) {
+ props.autoCreateSubscriptionGroup =
String(config.autoCreateSubscriptionGroup);
+ }
+ if (config.maxMessageSize != null) props.maxMessageSize =
String(config.maxMessageSize);
+ if (config.fileReservedTime != null) props.fileReservedTime =
String(config.fileReservedTime);
+ if (config.writeQueueNums != null) props.defaultTopicQueueNums =
String(config.writeQueueNums);
+ if (config.readQueueNums != null) props.defaultTopicQueueNums =
String(config.readQueueNums);
+ if (config.brokerPermission != null) props.brokerPermission =
String(config.brokerPermission);
+ return props;
+}
+
+function buildMockConfigPreviewChanges(
+ currentConfig: ClusterConfig,
+ proposedConfig: ClusterConfig,
+): ClusterConfigPreviewResult['changes'] {
+ const fields: Array<{ field: keyof ClusterConfig; brokerProperty: string }>
= [
+ { field: 'flushDiskType', brokerProperty: 'flushDiskType' },
+ { field: 'autoCreateTopicEnable', brokerProperty: 'autoCreateTopicEnable'
},
+ { field: 'autoCreateSubscriptionGroup', brokerProperty:
'autoCreateSubscriptionGroup' },
+ { field: 'maxMessageSize', brokerProperty: 'maxMessageSize' },
+ { field: 'fileReservedTime', brokerProperty: 'fileReservedTime' },
+ { field: 'writeQueueNums', brokerProperty: 'defaultTopicQueueNums' },
+ { field: 'readQueueNums', brokerProperty: 'defaultTopicQueueNums' },
+ { field: 'brokerPermission', brokerProperty: 'brokerPermission' },
+ ];
+ return fields
+ .filter(({ field }) => currentConfig[field] !== proposedConfig[field])
+ .map(({ field, brokerProperty }) => ({
+ field,
+ currentValue: String(currentConfig[field]),
+ proposedValue: String(proposedConfig[field]),
+ brokerProperty,
+ }));
+}
+
function getMockCluster(clusterId: string) {
const cluster = clusters.find((item) => item.id === clusterId);
if (!cluster) throw new Error(`Cluster not found: ${clusterId}`);