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 6f4ef279 feat: add instance-scoped Apache ACL 2.0 read foundation
(#1454)
6f4ef279 is described below
commit 6f4ef279e0c6c2b44415bae1ac7704eca6c77f1f
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 11 20:02:37 2026 +0800
feat: add instance-scoped Apache ACL 2.0 read foundation (#1454)
* feat: expose instance ACL capabilities
* feat: resolve authenticated admin clients by instance
* feat: read Apache ACL policies by instance
* fix: persist Apache ACL admin credential references
---
deploy/mysql/upgrade-cloud-vendor.sql | 20 ++++-
.../studio/cluster/broker/MqAdminExtFactory.java | 57 +++++++++++---
.../studio/cluster/broker/MqAdminProperties.java | 19 +++++
.../cluster/broker/RuntimeAdminClientResolver.java | 23 +++++-
.../studio/instance/CreateInstanceDTO.java | 3 +
.../rocketmq/studio/instance/InstanceService.java | 26 +++++-
.../rocketmq/studio/instance/InstanceVO.java | 1 +
.../instance/MybatisPlusInstanceRepository.java | 2 +
.../studio/instance/UpdateInstanceDTO.java | 3 +
.../AclCapabilitiesVO.java} | 40 +++-------
.../studio/instance/acl/AclController.java | 14 ++++
.../rocketmq/studio/instance/acl/AclService.java | 15 ++++
.../studio/instance/acl/ApacheAclReadService.java | 80 +++++++++++++++++++
.../studio/instance/acl/RemoteAclPolicyVO.java | 52 ++++++++++++
.../RemoteAclReadResult.java} | 46 ++++-------
.../studio/persistence/entity/RmqInstance.java | 2 +
server/src/main/resources/db/schema.sql | 1 +
.../rocketmq/studio/StudioApplicationTest.java | 30 +++++++
.../cluster/broker/MqAdminExtFactoryTest.java | 37 +++++++++
.../broker/RuntimeAdminClientResolverTest.java | 92 ++++++++++++++++++++--
.../studio/instance/InstanceServiceTest.java | 27 +++++++
.../MybatisPlusInstanceRepositoryTest.java | 8 +-
.../studio/instance/acl/AclControllerTest.java | 32 ++++++++
.../studio/instance/acl/AclServiceTest.java | 36 +++++++++
.../instance/acl/ApacheAclReadServiceTest.java | 61 ++++++++++++++
web/src/api/instance.ts | 3 +
web/src/pages/instance/index.tsx | 27 ++++++-
27 files changed, 675 insertions(+), 82 deletions(-)
diff --git a/deploy/mysql/upgrade-cloud-vendor.sql
b/deploy/mysql/upgrade-cloud-vendor.sql
index 26684247..ca7ac15e 100644
--- a/deploy/mysql/upgrade-cloud-vendor.sql
+++ b/deploy/mysql/upgrade-cloud-vendor.sql
@@ -72,10 +72,26 @@ PREPARE region_id_statement FROM @region_id_sql;
EXECUTE region_id_statement;
DEALLOCATE PREPARE region_id_statement;
--- 5. 存量实例兜底回填
+-- 5. rmq_instance.admin_credential_ref. The reference is non-secret; actual
Apache admin
+-- credentials remain external Spring configuration and are never stored in
this database.
+SET @admin_credential_ref_column_exists := (
+ SELECT COUNT(*)
+ FROM information_schema.columns
+ WHERE table_schema = @schema_name
+ AND table_name = 'rmq_instance'
+ AND column_name = 'admin_credential_ref'
+);
+SET @admin_credential_ref_sql := IF(@admin_credential_ref_column_exists = 0,
+ "ALTER TABLE rmq_instance ADD COLUMN admin_credential_ref VARCHAR(128)
COMMENT 'External Apache admin credential reference; no secret material' AFTER
credential_id",
+ 'SELECT 1');
+PREPARE admin_credential_ref_statement FROM @admin_credential_ref_sql;
+EXECUTE admin_credential_ref_statement;
+DEALLOCATE PREPARE admin_credential_ref_statement;
+
+-- 6. 存量实例兜底回填
UPDATE rmq_instance SET vendor = 'APACHE' WHERE vendor IS NULL OR vendor = '';
--- 6. 云厂商凭据表(secret_key 为 base64 编码,禁止明文;access_key 明文用于唯一键与打码展示)
+-- 7. 云厂商凭据表(secret_key 为 base64 编码,禁止明文;access_key 明文用于唯一键与打码展示)
CREATE TABLE IF NOT EXISTS rmq_cloud_credential (
id VARCHAR(64) PRIMARY KEY,
name VARCHAR(128) NOT NULL COMMENT '凭据显示名',
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactory.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactory.java
index c08ad75d..2bffa05e 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactory.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactory.java
@@ -27,6 +27,7 @@ import org.springframework.stereotype.Component;
import java.util.Arrays;
import java.util.Map;
+import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;
@@ -35,12 +36,12 @@ import java.util.stream.Collectors;
* Central lifecycle owner for real {@link DefaultMQAdminExt} connections.
*
* <p>This is the single real-network entry point shared by every live
cluster/metadata provider.
- * Admin clients are created lazily, started once and cached per NameServer
address so subsequent
- * calls reuse the established connection. All cached clients are shut down on
context destruction.
+ * Admin clients are created lazily, started once and cached per NameServer
address and credential
+ * reference so subsequent calls reuse the established connection without
crossing identities. All
+ * cached clients are shut down on context destruction.
*
- * <p>The {@link RPCHook} parameter is reserved for the ACL / authentication
work (AUTH-01); it is
- * currently always {@code null} but is threaded through so credentials can be
injected later
- * without changing this contract.
+ * <p>The {@link RPCHook} parameter carries Remoting authentication when a
selected instance has
+ * an externally configured admin credential.
*/
@Slf4j
@Component
@@ -49,7 +50,7 @@ public class MqAdminExtFactory {
/** Default admin RPC timeout in milliseconds. */
private static final long DEFAULT_TIMEOUT_MILLIS = 5000L;
- private final Map<String, DefaultMQAdminExt> cache = new
ConcurrentHashMap<>();
+ private final Map<AdminClientCacheKey, DefaultMQAdminExt> cache = new
ConcurrentHashMap<>();
private final AtomicInteger instanceCounter = new AtomicInteger();
private volatile boolean closed = false;
@@ -64,6 +65,19 @@ public class MqAdminExtFactory {
* @throws BusinessException if the connection cannot be established or
the action fails
*/
public <T> T execute(String namesrvAddr, RPCHook rpcHook, AdminAction<T>
action) {
+ String authenticationIdentity = rpcHook == null ? "anonymous"
+ : "custom-hook-" +
Integer.toUnsignedString(System.identityHashCode(rpcHook));
+ return execute(namesrvAddr, rpcHook, authenticationIdentity, action);
+ }
+
+ /**
+ * Runs an action with a cache entry isolated by a non-secret
authentication identity.
+ *
+ * <p>The identity must be a stable reference, never an access key or
secret key. Callers that
+ * do not use a configured credential should use the three-argument
overload.
+ */
+ public <T> T execute(String namesrvAddr, RPCHook rpcHook, String
authenticationIdentity,
+ AdminAction<T> action) {
if (namesrvAddr == null || namesrvAddr.isBlank()) {
throw new BusinessException(400, "NameServer address is required");
}
@@ -74,8 +88,10 @@ public class MqAdminExtFactory {
if (normalizedNamesrvAddr.isEmpty()) {
throw new BusinessException(400, "NameServer address is required");
}
- DefaultMQAdminExt admin = cache.computeIfAbsent(normalizedNamesrvAddr,
- addr -> createAndStart(addr, rpcHook));
+ AdminClientCacheKey cacheKey = new
AdminClientCacheKey(normalizedNamesrvAddr,
+ normalizeAuthenticationIdentity(authenticationIdentity));
+ DefaultMQAdminExt admin = cache.computeIfAbsent(cacheKey,
+ key -> createAndStart(key.namesrvAddr(), rpcHook));
try {
return action.apply(admin);
} catch (BusinessException ex) {
@@ -99,11 +115,14 @@ public class MqAdminExtFactory {
if (normalizedNamesrvAddr.isEmpty()) {
return;
}
- DefaultMQAdminExt admin = cache.remove(normalizedNamesrvAddr);
- if (admin != null) {
- safeShutdown(admin);
- log.info("Released RocketMQ admin client for namesrv {}",
normalizedNamesrvAddr);
- }
+ cache.entrySet().removeIf(entry -> {
+ if (!entry.getKey().namesrvAddr().equals(normalizedNamesrvAddr)) {
+ return false;
+ }
+ safeShutdown(entry.getValue());
+ return true;
+ });
+ log.info("Released RocketMQ admin clients for namesrv {}",
normalizedNamesrvAddr);
}
private DefaultMQAdminExt createAndStart(String namesrvAddr, RPCHook
rpcHook) {
@@ -143,6 +162,11 @@ public class MqAdminExtFactory {
.collect(Collectors.joining(";"));
}
+ private String normalizeAuthenticationIdentity(String
authenticationIdentity) {
+ return authenticationIdentity == null ||
authenticationIdentity.isBlank()
+ ? "anonymous" : authenticationIdentity.trim();
+ }
+
private void safeShutdown(MQAdminExt admin) {
try {
admin.shutdown();
@@ -173,4 +197,11 @@ public class MqAdminExtFactory {
public interface AdminAction<T> {
T apply(MQAdminExt admin) throws Exception;
}
+
+ private record AdminClientCacheKey(String namesrvAddr, String
authenticationIdentity) {
+ private AdminClientCacheKey {
+ Objects.requireNonNull(namesrvAddr, "namesrvAddr");
+ Objects.requireNonNull(authenticationIdentity,
"authenticationIdentity");
+ }
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminProperties.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminProperties.java
index fe5ece97..0e879553 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminProperties.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminProperties.java
@@ -21,6 +21,9 @@ import lombok.Setter;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
+import java.util.LinkedHashMap;
+import java.util.Map;
+
/**
* Configuration for the default RocketMQ cluster the studio connects to.
*
@@ -35,4 +38,20 @@ public class MqAdminProperties {
/** NameServer address list, e.g. {@code host1:9876;host2:9876}; blank
disables discovery. */
private String namesrvAddr;
+
+ /**
+ * Externally supplied credentials for Apache RocketMQ admin calls.
+ *
+ * <p>Instances persist only a reference into this map. Access keys and
secret keys must be
+ * supplied through externalized Spring configuration, such as environment
variables or a
+ * Kubernetes Secret, and must never be persisted in the Studio database.
+ */
+ private Map<String, Credential> credentials = new LinkedHashMap<>();
+
+ @Getter
+ @Setter
+ public static class Credential {
+ private String accessKey;
+ private String secretKey;
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
index 2bd87db7..d42bb010 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
@@ -7,6 +7,9 @@
package org.apache.rocketmq.studio.cluster.broker;
import lombok.RequiredArgsConstructor;
+import org.apache.rocketmq.acl.common.AclClientRPCHook;
+import org.apache.rocketmq.acl.common.SessionCredentials;
+import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.instance.InstanceRepository;
@@ -20,6 +23,7 @@ public class RuntimeAdminClientResolver {
private final InstanceRepository instanceRepository;
private final MqAdminExtFactory adminFactory;
+ private final MqAdminProperties adminProperties;
public InstanceVO resolveInstance(String instanceId) {
if (!StringUtils.hasText(instanceId)) {
@@ -46,7 +50,24 @@ public class RuntimeAdminClientResolver {
if (instance == null || !StringUtils.hasText(instance.getEndpoint())) {
throw new BusinessException(400, "Instance endpoint is required");
}
- return adminFactory.execute(instance.getEndpoint().trim(), null,
action);
+ String credentialRef =
StringUtils.hasText(instance.getAdminCredentialRef())
+ ? instance.getAdminCredentialRef().trim() : null;
+ return adminFactory.execute(instance.getEndpoint().trim(),
resolveCredential(credentialRef),
+ credentialRef, action);
+ }
+
+ private RPCHook resolveCredential(String credentialRef) {
+ if (!StringUtils.hasText(credentialRef)) {
+ return null;
+ }
+ MqAdminProperties.Credential credential =
adminProperties.getCredentials().get(credentialRef);
+ if (credential == null ||
!StringUtils.hasText(credential.getAccessKey())
+ || !StringUtils.hasText(credential.getSecretKey())) {
+ throw new BusinessException(422,
+ "Admin credential reference is not configured: " +
credentialRef);
+ }
+ return new AclClientRPCHook(new SessionCredentials(
+ credential.getAccessKey().trim(), credential.getSecretKey()));
}
private InstanceVO requireApacheInstance(InstanceVO instance) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/CreateInstanceDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/CreateInstanceDTO.java
index f338dd16..efb73287 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/CreateInstanceDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/CreateInstanceDTO.java
@@ -37,6 +37,8 @@ public class CreateInstanceDTO {
private String credentialId;
+ private String adminCredentialRef;
+
private String regionId;
public InstanceVO toInstanceVO() {
@@ -48,6 +50,7 @@ public class CreateInstanceDTO {
.vendor(vendor)
.cloudInstanceId(cloudInstanceId)
.credentialId(credentialId)
+ .adminCredentialRef(adminCredentialRef)
.regionId(regionId)
.build();
return vo;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
index f0050857..29b6f569 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
@@ -36,6 +36,7 @@ import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Locale;
+import java.util.Objects;
import java.util.UUID;
@Slf4j
@@ -108,6 +109,7 @@ public class InstanceService {
instance.setVendor(InstanceVendor.APACHE);
instance.setName(requireInstanceName(instance.getName()));
instance.setEndpoint(requireValidEndpoint(instance.getEndpoint()));
+
instance.setAdminCredentialRef(normalizeCredentialRef(instance.getAdminCredentialRef()));
if (instance.getType() == null) {
throw new BusinessException(400, "InstanceVO type is required");
}
@@ -188,6 +190,10 @@ public class InstanceService {
return name.trim();
}
+ private String normalizeCredentialRef(String credentialRef) {
+ return StringUtils.hasText(credentialRef) ? credentialRef.trim() :
null;
+ }
+
public InstanceVO updateInstance(InstanceVO instance) {
requireInstance(instance);
log.info("Updating instance: {}", instance.getId());
@@ -219,10 +225,13 @@ public class InstanceService {
if (instance.getRemark() != null) {
updated.setRemark(instance.getRemark());
}
+ if (!cloudInstance && instance.getAdminCredentialRef() != null) {
+
updated.setAdminCredentialRef(normalizeCredentialRef(instance.getAdminCredentialRef()));
+ }
updated.setUpdatedAt(LocalDateTime.now());
InstanceVO saved = instanceRepository.save(updated);
- releaseApacheEndpointIfUnused(existing, saved.getEndpoint());
+ releaseApacheClientIfChanged(existing, saved);
recordAudit("UPDATE_INSTANCE", "INSTANCE", saved.getId(), null,
instanceAuditDetail(saved));
return saved;
@@ -273,6 +282,7 @@ public class InstanceService {
.vendor(instance.getVendor() == null ? InstanceVendor.APACHE :
instance.getVendor())
.cloudInstanceId(instance.getCloudInstanceId())
.credentialId(instance.getCredentialId())
+ .adminCredentialRef(instance.getAdminCredentialRef())
.regionId(instance.getRegionId())
.topicCount(instance.getTopicCount())
.consumerGroupCount(instance.getConsumerGroupCount())
@@ -291,6 +301,20 @@ public class InstanceService {
releaseEndpointIfUnused(existing.getEndpoint(), currentEndpoint,
existing.getId());
}
+ private void releaseApacheClientIfChanged(InstanceVO existing, InstanceVO
saved) {
+ InstanceVendor vendor = existing.getVendor() == null ?
InstanceVendor.APACHE : existing.getVendor();
+ if (vendor != InstanceVendor.APACHE) {
+ return;
+ }
+ if
(!Objects.equals(normalizeCredentialRef(existing.getAdminCredentialRef()),
+ normalizeCredentialRef(saved.getAdminCredentialRef()))
+ && Objects.equals(normalizeEndpoint(existing.getEndpoint()),
normalizeEndpoint(saved.getEndpoint()))) {
+ adminFactory.release(existing.getEndpoint());
+ return;
+ }
+ releaseEndpointIfUnused(existing.getEndpoint(), saved.getEndpoint(),
existing.getId());
+ }
+
private void releaseEndpointIfUnused(String previousEndpoint, String
currentEndpoint, String excludedInstanceId) {
String oldEndpoint = normalizeEndpoint(previousEndpoint);
if (oldEndpoint == null ||
oldEndpoint.equals(normalizeEndpoint(currentEndpoint))) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceVO.java
index 3e8ed673..dceab86b 100644
--- a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceVO.java
+++ b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceVO.java
@@ -39,6 +39,7 @@ public class InstanceVO extends BaseEntity {
private InstanceVendor vendor;
private String cloudInstanceId;
private String credentialId;
+ private String adminCredentialRef;
private String regionId;
private int topicCount;
private int consumerGroupCount;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
index b0b83661..4a393e68 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
@@ -141,6 +141,7 @@ public class MybatisPlusInstanceRepository implements
InstanceRepository {
.vendor(parseVendor(entity.getId(), entity.getVendor()))
.cloudInstanceId(entity.getCloudInstanceId())
.credentialId(entity.getCredentialId())
+ .adminCredentialRef(entity.getAdminCredentialRef())
.regionId(entity.getRegionId())
.build();
vo.setId(entity.getId());
@@ -180,6 +181,7 @@ public class MybatisPlusInstanceRepository implements
InstanceRepository {
entity.setVendor(vo.getVendor() == null ? InstanceVendor.APACHE.name()
: vo.getVendor().name());
entity.setCloudInstanceId(vo.getCloudInstanceId());
entity.setCredentialId(vo.getCredentialId());
+ entity.setAdminCredentialRef(vo.getAdminCredentialRef());
entity.setRegionId(vo.getRegionId());
entity.setCreatedAt(vo.getCreatedAt());
entity.setUpdatedAt(vo.getUpdatedAt());
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
index 6b94f840..971e15d3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
@@ -34,12 +34,15 @@ public class UpdateInstanceDTO {
private String remark;
+ private String adminCredentialRef;
+
public InstanceVO toInstanceVO() {
InstanceVO vo = InstanceVO.builder()
.name(name)
.type(type)
.endpoint(endpoint)
.remark(remark)
+ .adminCredentialRef(adminCredentialRef)
.build();
vo.setId(id);
return vo;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclCapabilitiesVO.java
similarity index 58%
copy from
server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
copy to
server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclCapabilitiesVO.java
index 6b94f840..0b808804 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclCapabilitiesVO.java
@@ -14,34 +14,20 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.instance;
-import jakarta.validation.constraints.NotBlank;
-import lombok.Data;
-import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
-
-@Data
-public class UpdateInstanceDTO {
-
- @NotBlank(message = "instance id is required")
- private String id;
-
- private String name;
+package org.apache.rocketmq.studio.instance.acl;
- private InstanceType type;
-
- private String endpoint;
-
- private String remark;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
- public InstanceVO toInstanceVO() {
- InstanceVO vo = InstanceVO.builder()
- .name(name)
- .type(type)
- .endpoint(endpoint)
- .remark(remark)
- .build();
- vo.setId(id);
- return vo;
- }
+/**
+ * Describes whether ACL operations are backed by the selected instance.
+ */
+public record AclCapabilitiesVO(
+ String instanceId,
+ InstanceVendor vendor,
+ InstanceType instanceType,
+ String stateSource,
+ boolean remoteReadSupported,
+ boolean remoteWriteSupported) {
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclController.java
index a636415b..36bc8e9a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclController.java
@@ -37,6 +37,20 @@ import java.util.List;
public class AclController {
private final AclService aclService;
+ private final ApacheAclReadService apacheAclReadService;
+
+ @GetMapping("/remote/rules")
+ public Result<RemoteAclReadResult> listRemoteRules(
+ @RequestParam String instanceId,
+ @RequestParam(required = false) String subject,
+ @RequestParam(required = false) String resource) {
+ return Result.ok(apacheAclReadService.listRules(instanceId, subject,
resource));
+ }
+
+ @GetMapping("/capabilities")
+ public Result<AclCapabilitiesVO> capabilities(@RequestParam String
instanceId) {
+ return Result.ok(aclService.capabilities(instanceId));
+ }
@GetMapping("/rules")
public Result<List<AclRuleVO>> listRules(
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 237259ac..ce21c352 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
@@ -19,9 +19,12 @@ package org.apache.rocketmq.studio.instance.acl;
import org.springframework.util.StringUtils;
import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.util.CredentialUtils;
import org.apache.rocketmq.studio.audit.OperationAuditService;
import org.apache.rocketmq.studio.model.Acl2PolicyContext;
+import org.apache.rocketmq.studio.instance.InstanceRepository;
+import org.apache.rocketmq.studio.instance.InstanceVO;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@@ -37,6 +40,18 @@ public class AclService {
private final AclRepository aclRepository;
private final OperationAuditService operationAuditService;
+ private final InstanceRepository instanceRepository;
+
+ public AclCapabilitiesVO capabilities(String instanceId) {
+ if (!StringUtils.hasText(instanceId)) {
+ throw new BusinessException(400, "instanceId is required");
+ }
+ InstanceVO instance = instanceRepository.findById(instanceId)
+ .orElseThrow(() -> new BusinessException(404, "Instance not
found: " + instanceId));
+ boolean apacheInstance = instance.getVendor() == null ||
instance.getVendor() == InstanceVendor.APACHE;
+ return new AclCapabilitiesVO(instance.getId(), instance.getVendor(),
instance.getType(),
+ apacheInstance ? "APACHE_ACL2" : "STUDIO_LOCAL",
apacheInstance, false);
+ }
public List<AclRuleVO> listRules(String clusterId, String principal) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/ApacheAclReadService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/ApacheAclReadService.java
new file mode 100644
index 00000000..9edc9add
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/ApacheAclReadService.java
@@ -0,0 +1,80 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance.acl;
+
+import lombok.RequiredArgsConstructor;
+import org.apache.rocketmq.common.MixAll;
+import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
+import org.apache.rocketmq.remoting.protocol.route.BrokerData;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
+import org.springframework.stereotype.Service;
+
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+@Service
+@RequiredArgsConstructor
+public class ApacheAclReadService {
+ private final RuntimeAdminClientResolver runtimeAdminClientResolver;
+
+ public RemoteAclReadResult listRules(String instanceId, String subject,
String resource) {
+ return runtimeAdminClientResolver.execute(instanceId, admin -> {
+ Map<String, List<RemoteAclPolicyVO>> policies = new
LinkedHashMap<>();
+ Map<String, String> failures = new LinkedHashMap<>();
+ ClusterInfo clusterInfo = admin.examineBrokerClusterInfo();
+ if (clusterInfo == null || clusterInfo.getBrokerAddrTable() ==
null) {
+ return result(policies, failures);
+ }
+ for (BrokerData broker :
clusterInfo.getBrokerAddrTable().values()) {
+ String address = masterAddress(broker);
+ if (address == null || policies.containsKey(address) ||
failures.containsKey(address)) {
+ continue;
+ }
+ try {
+ policies.put(address, admin.listAcl(address, subject,
resource).stream()
+ .map(RemoteAclPolicyVO::from)
+ .toList());
+ } catch (Exception ex) {
+ failures.put(address, rootMessage(ex));
+ }
+ }
+ return result(policies, failures);
+ });
+ }
+
+ private RemoteAclReadResult result(Map<String, List<RemoteAclPolicyVO>>
policies,
+ Map<String, String> failures) {
+ return RemoteAclReadResult.builder().source("APACHE_ACL2")
+
.policiesByBroker(Map.copyOf(policies)).failuresByBroker(Map.copyOf(failures)).build();
+ }
+
+ private String masterAddress(BrokerData broker) {
+ if (broker == null || broker.getBrokerAddrs() == null) {
+ return null;
+ }
+ return broker.getBrokerAddrs().get(MixAll.MASTER_ID);
+ }
+
+ private String rootMessage(Exception ex) {
+ Throwable cause = ex;
+ while (cause.getCause() != null && cause.getCause() != cause) {
+ cause = cause.getCause();
+ }
+ return cause.getMessage() == null ? cause.getClass().getSimpleName() :
cause.getMessage();
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/RemoteAclPolicyVO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/RemoteAclPolicyVO.java
new file mode 100644
index 00000000..a7f25c22
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/RemoteAclPolicyVO.java
@@ -0,0 +1,52 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance.acl;
+
+import org.apache.rocketmq.remoting.protocol.body.AclInfo;
+
+import java.util.List;
+
+/** REST representation of an Apache ACL 2.0 policy. */
+public record RemoteAclPolicyVO(String subject, List<PolicyGroupVO> policies) {
+
+ public static RemoteAclPolicyVO from(AclInfo policy) {
+ List<PolicyGroupVO> groups = policy.getPolicies() == null ? List.of()
: policy.getPolicies().stream()
+ .map(PolicyGroupVO::from)
+ .toList();
+ return new RemoteAclPolicyVO(policy.getSubject(), groups);
+ }
+
+ public record PolicyGroupVO(String policyType, List<PolicyEntryVO>
entries) {
+ private static PolicyGroupVO from(AclInfo.PolicyInfo policy) {
+ List<PolicyEntryVO> policyEntries = policy.getEntries() == null ?
List.of() : policy.getEntries().stream()
+ .map(PolicyEntryVO::from)
+ .toList();
+ return new PolicyGroupVO(policy.getPolicyType(), policyEntries);
+ }
+ }
+
+ public record PolicyEntryVO(String resource, List<String> actions,
List<String> sourceIps, String decision) {
+ private static PolicyEntryVO from(AclInfo.PolicyEntryInfo entry) {
+ return new PolicyEntryVO(entry.getResource(),
listOrEmpty(entry.getActions()),
+ listOrEmpty(entry.getSourceIps()), entry.getDecision());
+ }
+
+ private static List<String> listOrEmpty(List<String> values) {
+ return values == null ? List.of() : List.copyOf(values);
+ }
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/RemoteAclReadResult.java
similarity index 54%
copy from
server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
copy to
server/src/main/java/org/apache/rocketmq/studio/instance/acl/RemoteAclReadResult.java
index 6b94f840..883207f0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/UpdateInstanceDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/RemoteAclReadResult.java
@@ -14,34 +14,22 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.instance;
-
-import jakarta.validation.constraints.NotBlank;
-import lombok.Data;
-import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
-
-@Data
-public class UpdateInstanceDTO {
-
- @NotBlank(message = "instance id is required")
- private String id;
-
- private String name;
-
- private InstanceType type;
-
- private String endpoint;
-
- private String remark;
-
- public InstanceVO toInstanceVO() {
- InstanceVO vo = InstanceVO.builder()
- .name(name)
- .type(type)
- .endpoint(endpoint)
- .remark(remark)
- .build();
- vo.setId(id);
- return vo;
+package org.apache.rocketmq.studio.instance.acl;
+
+import lombok.Builder;
+import lombok.Value;
+import java.util.List;
+import java.util.Map;
+
+/** Provider-backed ACL 2.0 read result, retaining per-Broker provenance. */
+@Value
+@Builder
+public class RemoteAclReadResult {
+ String source;
+ Map<String, List<RemoteAclPolicyVO>> policiesByBroker;
+ Map<String, String> failuresByBroker;
+
+ public boolean isPartial() {
+ return !failuresByBroker.isEmpty() && !policiesByBroker.isEmpty();
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqInstance.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqInstance.java
index c172e113..b37ec656 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqInstance.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqInstance.java
@@ -44,6 +44,8 @@ public class RmqInstance {
private String credentialId;
+ private String adminCredentialRef;
+
private String regionId;
private LocalDateTime createdAt;
diff --git a/server/src/main/resources/db/schema.sql
b/server/src/main/resources/db/schema.sql
index 82194c75..34d2ccf6 100644
--- a/server/src/main/resources/db/schema.sql
+++ b/server/src/main/resources/db/schema.sql
@@ -27,6 +27,7 @@ CREATE TABLE IF NOT EXISTS rmq_instance (
vendor VARCHAR(32),
cloud_instance_id VARCHAR(128),
credential_id VARCHAR(64),
+ admin_credential_ref VARCHAR(128) COMMENT 'External Apache admin credential
reference; no secret material',
region_id VARCHAR(128),
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/StudioApplicationTest.java
b/server/src/test/java/org/apache/rocketmq/studio/StudioApplicationTest.java
index 125d65be..a31a2f48 100644
--- a/server/src/test/java/org/apache/rocketmq/studio/StudioApplicationTest.java
+++ b/server/src/test/java/org/apache/rocketmq/studio/StudioApplicationTest.java
@@ -18,6 +18,10 @@ package org.apache.rocketmq.studio;
import org.apache.rocketmq.studio.ops.ai.tool.ToolCatalog;
import org.apache.rocketmq.studio.ops.ai.tool.ToolGatewayService;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
+import org.apache.rocketmq.studio.instance.InstanceVO;
+import org.apache.rocketmq.studio.instance.MybatisPlusInstanceRepository;
import org.apache.rocketmq.studio.persistence.mapper.RmqInstanceMapper;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
@@ -43,6 +47,9 @@ class StudioApplicationTest {
@Autowired
private RmqInstanceMapper instanceMapper;
+ @Autowired
+ private MybatisPlusInstanceRepository instanceRepository;
+
@Autowired
private MockMvc mockMvc;
@@ -56,4 +63,27 @@ class StudioApplicationTest {
.andExpect(status().isOk())
.andExpect(jsonPath("$.data").isEmpty());
}
+
+ @Test
+ void instanceRepositoryPersistsAdminCredentialReference() {
+ InstanceVO instance = InstanceVO.builder()
+ .name("acl-admin-ref-roundtrip")
+ .type(InstanceType.DIRECT)
+ .vendor(InstanceVendor.APACHE)
+ .endpoint("127.0.0.1:9876")
+ .adminCredentialRef("production-admin")
+ .build();
+ instance.setId("acl-admin-ref-roundtrip");
+
+ instanceRepository.save(instance);
+
+
assertThat(instanceMapper.selectById(instance.getId()).getAdminCredentialRef())
+ .isEqualTo("production-admin");
+ assertThat(instanceRepository.findById(instance.getId()))
+ .get()
+ .extracting(InstanceVO::getAdminCredentialRef)
+ .isEqualTo("production-admin");
+
+ instanceRepository.deleteById(instance.getId());
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactoryTest.java
index 88b60484..31965e12 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactoryTest.java
@@ -109,6 +109,43 @@ class MqAdminExtFactoryTest {
verify(admin).setNamesrvAddr("10.0.0.1:9876;10.0.0.2:9876");
}
+ @Test
+ void executeShouldNotShareClientsAcrossCredentialReferences() throws
Exception {
+ DefaultMQAdminExt admin = mock(DefaultMQAdminExt.class);
+ RecordingFactory factory = new RecordingFactory(admin);
+
+ factory.execute("10.0.0.1:9876", mock(RPCHook.class), "credential-a",
ignored -> null);
+ factory.execute("10.0.0.1:9876", mock(RPCHook.class), "credential-b",
ignored -> null);
+
+ assertThat(factory.created.get()).isEqualTo(2);
+ verify(admin, times(2)).start();
+ }
+
+ @Test
+ void executeShouldNotShareClientsAcrossLegacyHookInstances() throws
Exception {
+ DefaultMQAdminExt admin = mock(DefaultMQAdminExt.class);
+ RecordingFactory factory = new RecordingFactory(admin);
+
+ factory.execute("10.0.0.1:9876", mock(RPCHook.class), ignored -> null);
+ factory.execute("10.0.0.1:9876", mock(RPCHook.class), ignored -> null);
+
+ assertThat(factory.created.get()).isEqualTo(2);
+ }
+
+ @Test
+ void releaseShouldEvictEveryCredentialIdentityForAnEndpoint() throws
Exception {
+ DefaultMQAdminExt admin = mock(DefaultMQAdminExt.class);
+ RecordingFactory factory = new RecordingFactory(admin);
+
+ factory.execute("10.0.0.1:9876", mock(RPCHook.class), "credential-a",
ignored -> null);
+ factory.execute("10.0.0.1:9876", mock(RPCHook.class), "credential-b",
ignored -> null);
+ factory.release("10.0.0.1:9876");
+ factory.execute("10.0.0.1:9876", mock(RPCHook.class), "credential-a",
ignored -> null);
+
+ assertThat(factory.created.get()).isEqualTo(3);
+ verify(admin, times(2)).shutdown();
+ }
+
@Test
void executeShouldRejectAddressListsWithoutUsableEntries() {
MqAdminExtFactory factory = new MqAdminExtFactory();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
index 913d03ce..3ce4df4a 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
@@ -22,9 +22,14 @@ import
org.apache.rocketmq.studio.instance.InstanceRepository;
import org.apache.rocketmq.studio.instance.InstanceVO;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.boot.context.properties.bind.Bindable;
+import org.springframework.boot.context.properties.bind.Binder;
+import
org.springframework.boot.context.properties.source.MapConfigurationPropertySource;
+import java.util.Map;
import java.util.Optional;
import static org.assertj.core.api.Assertions.assertThat;
@@ -51,14 +56,16 @@ class RuntimeAdminClientResolverTest {
instance.setId("instance-a");
when(instanceRepository.findById("instance-a")).thenReturn(Optional.of(instance));
- RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory);
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
+ new MqAdminProperties());
assertThat(resolver.resolveEndpoint("instance-a")).isEqualTo("namesrv-a:9876");
}
@Test
void rejectsUnknownOrUnconfiguredInstances() {
- RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory);
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
+ new MqAdminProperties());
when(instanceRepository.findById("missing")).thenReturn(Optional.empty());
InstanceVO noEndpoint = InstanceVO.builder().endpoint(" ").build();
when(instanceRepository.findById("no-endpoint")).thenReturn(Optional.of(noEndpoint));
@@ -75,12 +82,13 @@ class RuntimeAdminClientResolverTest {
void executesAgainstTheSelectedInstanceEndpoint() {
InstanceVO instance =
InstanceVO.builder().endpoint("namesrv-b:9876").build();
when(instanceRepository.findById("instance-b")).thenReturn(Optional.of(instance));
- when(adminFactory.execute(eq("namesrv-b:9876"), isNull(),
any())).thenReturn("done");
- RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory);
+ when(adminFactory.execute(eq("namesrv-b:9876"), isNull(), isNull(),
any())).thenReturn("done");
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
+ new MqAdminProperties());
String result = resolver.execute("instance-b", admin -> "unused");
assertThat(result).isEqualTo("done");
- verify(adminFactory).execute(eq("namesrv-b:9876"), isNull(), any());
+ verify(adminFactory).execute(eq("namesrv-b:9876"), isNull(), isNull(),
any());
}
@Test
@@ -91,7 +99,8 @@ class RuntimeAdminClientResolverTest {
.build();
instance.setId("cloud-instance");
when(instanceRepository.findById("cloud-instance")).thenReturn(Optional.of(instance));
- RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory);
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
+ new MqAdminProperties());
assertThatThrownBy(() -> resolver.resolveEndpoint("cloud-instance"))
.isInstanceOf(BusinessException.class)
@@ -101,4 +110,75 @@ class RuntimeAdminClientResolverTest {
.hasMessage("Runtime AdminClient only supports Apache
instances: cloud-instance");
verifyNoInteractions(adminFactory);
}
+
+ @Test
+ void executesWithTheSelectedInstanceCredentialReference() {
+ InstanceVO instance = InstanceVO.builder()
+ .endpoint("namesrv-b:9876")
+ .adminCredentialRef(" production-admin ")
+ .build();
+ instance.setId("instance-b");
+ MqAdminProperties properties = new MqAdminProperties();
+ MqAdminProperties.Credential credential = new
MqAdminProperties.Credential();
+ credential.setAccessKey("admin-ak");
+ credential.setSecretKey("admin-sk");
+ properties.getCredentials().put("production-admin", credential);
+
when(instanceRepository.findById("instance-b")).thenReturn(Optional.of(instance));
+ when(adminFactory.execute(eq("namesrv-b:9876"), any(),
eq("production-admin"), any()))
+ .thenReturn("done");
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
+ properties);
+
+ String result = resolver.execute("instance-b", ignored -> "unused");
+
+ assertThat(result).isEqualTo("done");
+ ArgumentCaptor<org.apache.rocketmq.remoting.RPCHook> hookCaptor =
ArgumentCaptor.forClass(
+ org.apache.rocketmq.remoting.RPCHook.class);
+ verify(adminFactory).execute(eq("namesrv-b:9876"),
hookCaptor.capture(), eq("production-admin"), any());
+ org.apache.rocketmq.acl.common.AclClientRPCHook resolvedHook =
+ (org.apache.rocketmq.acl.common.AclClientRPCHook)
hookCaptor.getValue();
+
assertThat(resolvedHook.getSessionCredentials().getAccessKey()).isEqualTo("admin-ak");
+
assertThat(resolvedHook.getSessionCredentials().getSecretKey()).isEqualTo("admin-sk");
+ }
+
+ @Test
+ void rejectsUnknownOrIncompleteCredentialReferencesBeforeNetworkCalls() {
+ MqAdminProperties properties = new MqAdminProperties();
+ MqAdminProperties.Credential credential = new
MqAdminProperties.Credential();
+ credential.setAccessKey("admin-ak");
+ credential.setSecretKey("admin-sk");
+ properties.getCredentials().put("production-admin", credential);
+ MqAdminProperties.Credential incomplete = new
MqAdminProperties.Credential();
+ incomplete.setAccessKey("admin-ak");
+ properties.getCredentials().put("incomplete", incomplete);
+ InstanceVO instance = InstanceVO.builder().endpoint("namesrv-b:9876")
+ .adminCredentialRef("missing").build();
+ instance.setId("instance-b");
+
when(instanceRepository.findById("instance-b")).thenReturn(Optional.of(instance));
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
+ properties);
+
+ assertThatThrownBy(() -> resolver.execute("instance-b", ignored ->
"unused"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Admin credential reference is not configured:
missing")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(422));
+ instance.setAdminCredentialRef("incomplete");
+ assertThatThrownBy(() -> resolver.execute("instance-b", ignored ->
"unused"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Admin credential reference is not configured:
incomplete");
+ verifyNoInteractions(adminFactory);
+ }
+
+ @Test
+ void bindsReferencedCredentialsFromExternalizedConfiguration() {
+ MqAdminProperties properties = new Binder(new
MapConfigurationPropertySource(Map.of(
+
"studio.cluster.admin.credentials.production-admin.access-key", "admin-ak",
+
"studio.cluster.admin.credentials.production-admin.secret-key", "admin-sk")))
+ .bind("studio.cluster.admin",
Bindable.of(MqAdminProperties.class)).get();
+
+ MqAdminProperties.Credential credential =
properties.getCredentials().get("production-admin");
+
+ assertThat(credential.getAccessKey()).isEqualTo("admin-ak");
+ assertThat(credential.getSecretKey()).isEqualTo("admin-sk");
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
index fe314a20..f0e6bb44 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
@@ -347,6 +347,33 @@ class InstanceServiceTest {
assertThat(instanceService.createInstance(input).getName()).isEqualTo("production");
}
+ @Test
+ void
createApacheInstanceShouldTrimAndPersistOnlyAdminCredentialReference() {
+ InstanceVO input =
InstanceVO.builder().name("production").type(InstanceType.PROXY)
+ .endpoint("namesrv:9876").adminCredentialRef("
production-admin ").build();
+
when(instanceRepository.save(any(InstanceVO.class))).thenAnswer(invocation ->
invocation.getArgument(0));
+
+ InstanceVO saved = instanceService.createInstance(input);
+
+
assertThat(saved.getAdminCredentialRef()).isEqualTo("production-admin");
+ }
+
+ @Test
+ void
updateApacheInstanceShouldReleaseCachedClientWhenCredentialReferenceChanges() {
+ InstanceVO existing =
InstanceVO.builder().name("production").type(InstanceType.PROXY)
+
.endpoint("namesrv:9876").adminCredentialRef("credential-a").build();
+ existing.setId("inst-1");
+ InstanceVO update =
InstanceVO.builder().adminCredentialRef("credential-b").build();
+ update.setId("inst-1");
+
when(instanceRepository.findById("inst-1")).thenReturn(Optional.of(existing));
+
when(instanceRepository.save(any(InstanceVO.class))).thenAnswer(invocation ->
invocation.getArgument(0));
+
+ InstanceVO saved = instanceService.updateInstance(update);
+
+ assertThat(saved.getAdminCredentialRef()).isEqualTo("credential-b");
+ verify(adminFactory).release("namesrv:9876");
+ }
+
@Test
void updateInstanceShouldMergeFieldsOntoExisting() {
LocalDateTime originalCreatedAt = LocalDateTime.of(2025, 1, 2, 3, 4,
5);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
index b1290ba5..731b26ba 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
@@ -27,6 +27,7 @@ import
org.apache.rocketmq.studio.persistence.mapper.RmqInstanceMapper;
import org.apache.rocketmq.studio.persistence.mapper.RmqTopicMapper;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
@@ -75,6 +76,7 @@ class MybatisPlusInstanceRepositoryTest {
assertThat(direct.getConsumerGroupCount()).isZero();
assertThat(proxy.getTopicCount()).isZero();
assertThat(proxy.getConsumerGroupCount()).isZero();
+
assertThat(direct.getAdminCredentialRef()).isEqualTo("admin-instance-direct-1");
verifyNoInteractions(topicMapper, groupMapper);
}
@@ -160,7 +162,9 @@ class MybatisPlusInstanceRepositoryTest {
repository.save(vo);
- verify(instanceMapper).insert(any(RmqInstance.class));
+ ArgumentCaptor<RmqInstance> entity =
ArgumentCaptor.forClass(RmqInstance.class);
+ verify(instanceMapper).insert(entity.capture());
+
assertThat(entity.getValue().getAdminCredentialRef()).isEqualTo("admin-instance-proxy-2");
verify(instanceMapper, never()).updateById(any(RmqInstance.class));
}
@@ -189,6 +193,7 @@ class MybatisPlusInstanceRepositoryTest {
entity.setType(type.name());
entity.setEndpoint("10.0.0.1:9876");
entity.setVendor(InstanceVendor.APACHE.name());
+ entity.setAdminCredentialRef("admin-" + id);
entity.setCreatedAt(LocalDateTime.of(2026, 8, 3, 0, 0));
entity.setUpdatedAt(LocalDateTime.of(2026, 8, 3, 0, 0));
return entity;
@@ -201,6 +206,7 @@ class MybatisPlusInstanceRepositoryTest {
.endpoint("10.0.0.1:9876")
.build();
vo.setId(id);
+ vo.setAdminCredentialRef("admin-" + id);
return vo;
}
}
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 d6ca17c2..773f1f2a 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
@@ -54,9 +54,41 @@ class AclControllerTest {
@MockBean
private AclService aclService;
+ @MockBean
+ private ApacheAclReadService apacheAclReadService;
+
@Autowired
private ObjectMapper objectMapper;
+ @Test
+ void capabilitiesShouldReturnApacheRemoteReadState() throws Exception {
+ when(aclService.capabilities("instance-1"))
+ .thenReturn(new AclCapabilitiesVO("instance-1",
+
org.apache.rocketmq.studio.common.domain.enums.InstanceVendor.APACHE,
+
org.apache.rocketmq.studio.common.domain.enums.InstanceType.DIRECT,
+ "APACHE_ACL2", true, false));
+
+ mockMvc.perform(get("/api/acl/capabilities").param("instanceId",
"instance-1"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.instanceId").value("instance-1"))
+ .andExpect(jsonPath("$.data.stateSource").value("APACHE_ACL2"))
+ .andExpect(jsonPath("$.data.remoteReadSupported").value(true))
+
.andExpect(jsonPath("$.data.remoteWriteSupported").value(false));
+ }
+
+ @Test
+ void listRemoteRulesShouldRequireInstanceAndDelegateToApacheProvider()
throws Exception {
+ RemoteAclReadResult result = RemoteAclReadResult.builder()
+
.source("APACHE_ACL2").policiesByBroker(Map.of()).failuresByBroker(Map.of()).build();
+ when(apacheAclReadService.listRules("instance-1", null,
null)).thenReturn(result);
+
+ mockMvc.perform(get("/api/acl/remote/rules").param("instanceId",
"instance-1"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.source").value("APACHE_ACL2"));
+
+ verify(apacheAclReadService).listRules("instance-1", null, null);
+ }
+
@Test
void listRulesShouldReturnRules() throws Exception {
AclRuleVO rule = AclRuleVO.builder()
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 dc8e571a..67100647 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
@@ -20,6 +20,10 @@ package org.apache.rocketmq.studio.instance.acl;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.audit.OperationAuditService;
import org.apache.rocketmq.studio.model.Acl2PolicyContext;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
+import org.apache.rocketmq.studio.instance.InstanceRepository;
+import org.apache.rocketmq.studio.instance.InstanceVO;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@@ -54,6 +58,9 @@ class AclServiceTest {
@Mock
private OperationAuditService operationAuditService;
+ @Mock
+ private InstanceRepository instanceRepository;
+
@InjectMocks
private AclService aclService;
@@ -86,6 +93,35 @@ class AclServiceTest {
verify(aclRepository).findRules("cluster-1", "user1");
}
+ @Test
+ void capabilitiesShouldDescribeApacheRemoteReadSupport() {
+ InstanceVO instance = InstanceVO.builder()
+ .name("instance-1")
+ .vendor(InstanceVendor.APACHE)
+ .type(InstanceType.DIRECT)
+ .build();
+ instance.setId("instance-1");
+
when(instanceRepository.findById("instance-1")).thenReturn(Optional.of(instance));
+
+ AclCapabilitiesVO capabilities = aclService.capabilities("instance-1");
+
+ assertThat(capabilities.instanceId()).isEqualTo("instance-1");
+ assertThat(capabilities.vendor()).isEqualTo(InstanceVendor.APACHE);
+ assertThat(capabilities.instanceType()).isEqualTo(InstanceType.DIRECT);
+ assertThat(capabilities.stateSource()).isEqualTo("APACHE_ACL2");
+ assertThat(capabilities.remoteReadSupported()).isTrue();
+ assertThat(capabilities.remoteWriteSupported()).isFalse();
+ }
+
+ @Test
+ void capabilitiesShouldRejectUnknownInstance() {
+
when(instanceRepository.findById("missing")).thenReturn(Optional.empty());
+
+ assertThatThrownBy(() -> aclService.capabilities("missing"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Instance not found: missing");
+ }
+
@Test
void listRulesShouldPassNullFilters() {
when(aclRepository.findRules(null, null)).thenReturn(List.of());
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/ApacheAclReadServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/ApacheAclReadServiceTest.java
new file mode 100644
index 00000000..3bf178b5
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/ApacheAclReadServiceTest.java
@@ -0,0 +1,61 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance.acl;
+
+import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
+import org.apache.rocketmq.remoting.protocol.route.BrokerData;
+import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
+import org.apache.rocketmq.tools.admin.MQAdminExt;
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+class ApacheAclReadServiceTest {
+ @Test
+ void retainsSuccessfulBrokerDataWhenAnotherBrokerFails() throws Exception {
+ RuntimeAdminClientResolver resolver =
mock(RuntimeAdminClientResolver.class);
+ MQAdminExt admin = mock(MQAdminExt.class);
+ ClusterInfo clusterInfo = new ClusterInfo();
+ BrokerData healthy = new BrokerData();
+ healthy.setBrokerAddrs(new HashMap<>(Map.of(0L, "broker-a:10911")));
+ BrokerData failing = new BrokerData();
+ failing.setBrokerAddrs(new HashMap<>(Map.of(0L, "broker-b:10911")));
+ clusterInfo.setBrokerAddrTable(Map.of("a", healthy, "b", failing));
+ when(admin.examineBrokerClusterInfo()).thenReturn(clusterInfo);
+ when(admin.listAcl("broker-a:10911", null,
null)).thenReturn(List.of());
+ when(admin.listAcl("broker-b:10911", null, null)).thenThrow(new
IllegalStateException("unavailable"));
+ when(resolver.execute(eq("instance-1"), any())).thenAnswer(invocation
-> {
+ MqAdminExtFactory.AdminAction<?> action =
invocation.getArgument(1);
+ return action.apply(admin);
+ });
+
+ RemoteAclReadResult result = new
ApacheAclReadService(resolver).listRules("instance-1", null, null);
+
+ assertThat(result.getPoliciesByBroker()).containsKey("broker-a:10911");
+
assertThat(result.getFailuresByBroker()).containsEntry("broker-b:10911",
"unavailable");
+ assertThat(result.isPartial()).isTrue();
+ }
+}
diff --git a/web/src/api/instance.ts b/web/src/api/instance.ts
index 0bc7e015..e8bb1305 100644
--- a/web/src/api/instance.ts
+++ b/web/src/api/instance.ts
@@ -29,6 +29,7 @@ export interface Instance {
vendor?: InstanceVendor;
cloudInstanceId?: string;
credentialId?: string;
+ adminCredentialRef?: string;
regionId?: string;
topicCount: number;
consumerGroupCount: number;
@@ -45,6 +46,7 @@ export interface CreateInstanceRequest {
vendor?: InstanceVendor;
cloudInstanceId?: string;
credentialId?: string;
+ adminCredentialRef?: string;
regionId?: string;
}
@@ -54,6 +56,7 @@ export interface UpdateInstanceRequest {
type?: 'PROXY' | 'DIRECT';
endpoint?: string;
remark?: string;
+ adminCredentialRef?: string;
}
export interface InstanceQuery {
diff --git a/web/src/pages/instance/index.tsx b/web/src/pages/instance/index.tsx
index f1b1b4f2..824d9297 100644
--- a/web/src/pages/instance/index.tsx
+++ b/web/src/pages/instance/index.tsx
@@ -245,7 +245,11 @@ const InstancePage = () => {
try {
const values = await editForm.validateFields();
setSubmitting(true);
- const updated = await updateInstance({ id: editingInstance.id, remark:
values.remark || '' });
+ const updated = await updateInstance({
+ id: editingInstance.id,
+ remark: values.remark || '',
+ adminCredentialRef: values.adminCredentialRef,
+ });
await loadInstances();
message.success(`实例「${updated.name}」备注已更新`);
setEditModalOpen(false);
@@ -382,7 +386,10 @@ const InstancePage = () => {
style={{ borderColor: '#1677ff', color: '#1677ff' }}
onClick={() => {
setEditingInstance(record);
- editForm.setFieldsValue({ remark: record.remark });
+ editForm.setFieldsValue({
+ remark: record.remark,
+ adminCredentialRef: record.adminCredentialRef,
+ });
setEditModalOpen(true);
}}
>
@@ -625,6 +632,13 @@ const InstancePage = () => {
message="接入地址为客户端访问入口"
description="接入地址会展示在 Topic 等页面供客户端配置使用。若客户端环境无法解析该地址(如 K8s 内部
Service 域名),可自行配置 DNS 解析或在客户端 hosts 中映射。"
/>
+ <Form.Item
+ label="管理凭据引用"
+ name="adminCredentialRef"
+ extra="可选。仅保存服务端配置中的凭据引用,不会保存或传输 AK/SK。"
+ >
+ <Input placeholder="例:production-admin" />
+ </Form.Item>
<Form.Item label="备注" name="remark">
<Input.TextArea rows={2} placeholder="可选,描述实例用途" />
</Form.Item>
@@ -663,6 +677,15 @@ const InstancePage = () => {
<Form.Item label="接入地址">
<Input value={editingInstance?.endpoint} disabled />
</Form.Item>
+ {editingInstance?.vendor === 'APACHE' && (
+ <Form.Item
+ label="管理凭据引用"
+ name="adminCredentialRef"
+ extra="仅保存服务端配置中的引用,不会保存或传输 AK/SK。"
+ >
+ <Input placeholder="例:production-admin" />
+ </Form.Item>
+ )}
<Form.Item label="备注" name="remark">
<Input.TextArea rows={3} placeholder="描述实例用途" />
</Form.Item>