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
commit 958e3962110c954ff3b463a69a2adabab77f327f Author: Aias00 <[email protected]> AuthorDate: Mon Aug 10 12:57:56 2026 +0800 fix: instance-scoped hardening, mutation auditing and validation (#1274, #1280, #1282, #1288, #1293, #1296, #1303, #1305, #1307, #1309, #1311, #1319) * fix: reject non-SEND_OK topic test messages Signed-off-by: liuhy <[email protected]> * fix: surface message trace query failures Signed-off-by: liuhy <[email protected]> * fix: require a topic for message ID queries * fix: reject blank message trace identifiers * fix: clear stale DLQ state when switching instances Signed-off-by: liuhy <[email protected]> * fix: handle client instance bootstrap failures Signed-off-by: liuhy <[email protected]> * fix: scope producer topics by selected instance Signed-off-by: liuhy <[email protected]> * fix: require instance IDs for client queries * fix: ignore stale dashboard responses Signed-off-by: liuhy <[email protected]> * test: cover dashboard instance response ordering Signed-off-by: liuhy <[email protected]> * fix: initialize the H2 schema for development * fix: reject cloud instances from remoting admin resolution * feat: audit cloud credential mutations * feat: audit managed instance mutations * fix: normalize managed instance names * feat: audit ACL mutations * fix: clarify that ACL records are Studio-local * feat: audit K8s certificate mutations * fix: clarify that K8s certificate records are local * feat: audit alert mutations * fix: exclude disabled alert rules from export * fix: validate runtime Topic and Group names * fix: adapt batch to studio conventions and stabilize producer instance switch - isolate audit write failures (try/catch + warn) in Instance/K8sCert/Acl/Alert/CloudCredential services, matching the control-plane audit isolation precedent - replace private isBlank/JDK blank checks with Spring StringUtils in MetadataService and ClientService - Producer page: clear stale topic/group state in the instance change handler instead of synchronously in an effect (also fixes antd Select breakage across remounts in tests) - Producer test: scope stale-state assertions to the render container (antd keeps closed dropdown portals in document.body) --------- Signed-off-by: liuhy <[email protected]> Co-authored-by: liuhy <[email protected]> Co-authored-by: lizhimins <[email protected]> --- .../cluster/broker/RuntimeAdminClientResolver.java | 12 ++- .../studio/cluster/client/ClientService.java | 11 +- .../studio/cluster/k8s/K8sCertService.java | 31 +++++- .../rocketmq/studio/instance/InstanceService.java | 43 ++++++-- .../rocketmq/studio/instance/acl/AclService.java | 41 ++++++- .../studio/instance/message/MessageService.java | 6 ++ .../studio/instance/topic/MetadataService.java | 33 +++--- .../rocketmq/studio/ops/alert/AlertService.java | 43 +++++++- .../provider/apache/RocketMQAdminClientImpl.java | 8 ++ .../provider/apache/RocketMQMessageProvider.java | 3 + .../credential/CloudCredentialService.java | 27 ++++- server/src/main/resources/application-dev.yml | 4 + server/src/main/resources/db/schema.sql | 14 ++- .../rocketmq/studio/StudioApplicationTest.java | 20 +++- .../broker/RuntimeAdminClientResolverTest.java | 21 ++++ .../studio/cluster/client/ClientServiceTest.java | 20 +++- .../studio/cluster/k8s/K8sCertServiceTest.java | 18 +++- .../studio/instance/InstanceServiceTest.java | 35 ++++++ .../studio/instance/acl/AclServiceTest.java | 20 ++++ .../instance/message/MessageServiceTest.java | 27 +++-- .../studio/instance/topic/MetadataServiceTest.java | 21 ++++ .../ops/alert/AlertServiceDefaultRulesTest.java | 8 +- .../studio/ops/alert/AlertServiceTest.java | 70 +++++++++++- .../apache/RocketMQAdminClientImplTest.java | 49 +++++++++ .../apache/RocketMQMessageProviderTest.java | 13 +++ .../credential/CloudCredentialServiceTest.java | 14 +++ web/src/api/producer.test.ts | 11 +- web/src/api/producer.ts | 6 +- web/src/i18n/translations.ts | 8 ++ .../pages/cluster/__tests__/ClientsPage.test.tsx | 28 +++++ .../pages/cluster/__tests__/K8sCertsPage.test.tsx | 8 ++ web/src/pages/cluster/certs.tsx | 9 ++ web/src/pages/cluster/clients.tsx | 40 +++++-- .../pages/home/__tests__/DashboardPage.test.tsx | 118 +++++++++++++++++++++ web/src/pages/home/dashboard.tsx | 17 ++- web/src/pages/instance/__tests__/AclPage.test.tsx | 8 ++ web/src/pages/instance/__tests__/DLQPage.test.tsx | 46 ++++++++ .../pages/instance/__tests__/MessagePage.test.tsx | 22 +++- web/src/pages/instance/acl.tsx | 10 ++ web/src/pages/instance/dlq.tsx | 17 ++- web/src/pages/instance/message.tsx | 51 ++++++--- web/src/pages/studio/Producer.tsx | 22 +++- web/src/pages/studio/__tests__/Producer.test.tsx | 64 ++++++++++- 43 files changed, 994 insertions(+), 103 deletions(-) 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 5acead9a..2bd87db7 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,7 @@ package org.apache.rocketmq.studio.cluster.broker; import lombok.RequiredArgsConstructor; +import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor; import org.apache.rocketmq.studio.common.exception.BusinessException; import org.apache.rocketmq.studio.instance.InstanceRepository; import org.apache.rocketmq.studio.instance.InstanceVO; @@ -29,7 +30,7 @@ public class RuntimeAdminClientResolver { } public String resolveEndpoint(String instanceId) { - InstanceVO instance = resolveInstance(instanceId); + InstanceVO instance = requireApacheInstance(resolveInstance(instanceId)); if (!StringUtils.hasText(instance.getEndpoint())) { throw new BusinessException(400, "Instance has no endpoint: " + instanceId); } @@ -41,9 +42,18 @@ public class RuntimeAdminClientResolver { } public <T> T execute(InstanceVO instance, MqAdminExtFactory.AdminAction<T> action) { + requireApacheInstance(instance); if (instance == null || !StringUtils.hasText(instance.getEndpoint())) { throw new BusinessException(400, "Instance endpoint is required"); } return adminFactory.execute(instance.getEndpoint().trim(), null, action); } + + private InstanceVO requireApacheInstance(InstanceVO instance) { + if (instance != null && instance.getVendor() != null && instance.getVendor() != InstanceVendor.APACHE) { + throw new BusinessException(400, + "Runtime AdminClient only supports Apache instances: " + instance.getId()); + } + return instance; + } } diff --git a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientService.java b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientService.java index ed9723e4..9d04ebcf 100644 --- a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientService.java +++ b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientService.java @@ -18,7 +18,9 @@ package org.apache.rocketmq.studio.cluster.client; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.apache.rocketmq.studio.common.exception.BusinessException; import org.springframework.stereotype.Service; +import org.springframework.util.StringUtils; import java.util.List; @@ -31,7 +33,14 @@ public class ClientService { public List<ClientConnectionVO> listConnections(String instanceId, String clusterId, String type) { log.info("Listing client connections, instanceId={}, clusterId={}, type={}", instanceId, clusterId, type); - return clientProvider.findConnections(normalizeFilter(instanceId), normalizeFilter(clusterId), normalizeFilter(type)); + return clientProvider.findConnections(requireInstanceId(instanceId), normalizeFilter(clusterId), normalizeFilter(type)); + } + + private String requireInstanceId(String instanceId) { + if (!StringUtils.hasText(instanceId)) { + throw new BusinessException(400, "instanceId is required"); + } + return instanceId.trim(); } private String normalizeFilter(String value) { diff --git a/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertService.java b/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertService.java index 6ca6debf..b13f9804 100644 --- a/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertService.java +++ b/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertService.java @@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.cluster.k8s; import org.apache.rocketmq.studio.common.domain.enums.CertStatus; import org.apache.rocketmq.studio.common.domain.enums.CertType; import org.apache.rocketmq.studio.common.exception.BusinessException; +import org.apache.rocketmq.studio.audit.OperationAuditService; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -38,15 +39,17 @@ public class K8sCertService { private static final int EXPIRING_THRESHOLD_DAYS = 30; private final K8sCertRepository k8sCertRepository; + private final OperationAuditService operationAuditService; private final Clock clock; @Autowired - public K8sCertService(K8sCertRepository k8sCertRepository) { - this(k8sCertRepository, Clock.systemDefaultZone()); + public K8sCertService(K8sCertRepository k8sCertRepository, OperationAuditService operationAuditService) { + this(k8sCertRepository, operationAuditService, Clock.systemDefaultZone()); } - K8sCertService(K8sCertRepository k8sCertRepository, Clock clock) { + K8sCertService(K8sCertRepository k8sCertRepository, OperationAuditService operationAuditService, Clock clock) { this.k8sCertRepository = k8sCertRepository; + this.operationAuditService = operationAuditService; this.clock = clock; } @@ -82,6 +85,7 @@ public class K8sCertService { cert.setUpdatedAt(now); K8sCertVO saved = k8sCertRepository.save(cert); + auditCertificate("CREATE_K8S_CERTIFICATE", saved); log.info("K8s certificate created: {} (id={})", saved.getName(), saved.getId()); return saved; } @@ -115,6 +119,7 @@ public class K8sCertService { updated.setUpdatedAt(now); K8sCertVO saved = k8sCertRepository.save(refreshExpirationState(updated, now)); + auditCertificate("UPDATE_K8S_CERTIFICATE", saved); log.info("K8s certificate updated: {} (id={})", saved.getName(), saved.getId()); return saved; } @@ -136,6 +141,7 @@ public class K8sCertService { renewed.setUpdatedAt(now); K8sCertVO saved = k8sCertRepository.save(renewed); + auditCertificate("RENEW_K8S_CERTIFICATE", saved); log.info("K8s certificate renewed: {} (id={}), new expiry: {}", saved.getName(), saved.getId(), notAfter); return saved; } @@ -146,6 +152,8 @@ public class K8sCertService { k8sCertRepository.findById(command.getId()) .orElseThrow(() -> new BusinessException(404, "Certificate not found: " + command.getId())); k8sCertRepository.deleteById(command.getId()); + recordAudit("DELETE_K8S_CERTIFICATE", "K8S_CERTIFICATE", command.getId(), null, + null); log.info("K8s certificate deleted: {}", command.getId()); } @@ -192,4 +200,21 @@ public class K8sCertService { copy.setUpdatedAt(cert.getUpdatedAt()); return copy; } + + private void auditCertificate(String operation, K8sCertVO certificate) { + recordAudit(operation, "K8S_CERTIFICATE", certificate.getId(), null, + "name=" + certificate.getName() + ", namespace=" + certificate.getNamespace() + + ", cluster=" + certificate.getCluster()); + } + + private void recordAudit(String operation, String resourceType, String resourceName, + String clusterId, String detail) { + try { + operationAuditService.record(operation, resourceType, resourceName, clusterId, detail, "SUCCESS", null); + } catch (Exception auditFailure) { + log.warn("Failed to record audit operation={} resource={}: {}", operation, resourceName, + auditFailure.getMessage()); + } + } + } 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 753cc464..4971f38f 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 @@ -22,6 +22,7 @@ import org.springframework.util.StringUtils; import org.apache.rocketmq.studio.provider.credential.CloudCredentialRepository; import org.apache.rocketmq.studio.provider.credential.CloudCredentialVO; import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory; +import org.apache.rocketmq.studio.audit.OperationAuditService; import org.apache.rocketmq.studio.common.domain.enums.InstanceType; import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor; import org.apache.rocketmq.studio.common.exception.BusinessException; @@ -45,6 +46,7 @@ public class InstanceService { private final CloudCredentialRepository cloudCredentialRepository; private final InstanceProviderRegistry providerRegistry; private final MqAdminExtFactory adminFactory; + private final OperationAuditService operationAuditService; public List<InstanceVO> listInstances(InstanceType type, String search) { log.debug("Listing instances, type={}, search={}", type, search); @@ -96,14 +98,15 @@ public class InstanceService { instance.setId(UUID.randomUUID().toString()); instance.setCreatedAt(LocalDateTime.now()); instance.setUpdatedAt(LocalDateTime.now()); - return instanceRepository.save(instance); + InstanceVO saved = instanceRepository.save(instance); + recordAudit("CREATE_INSTANCE", "INSTANCE", saved.getId(), null, + instanceAuditDetail(saved)); + return saved; } private void createApacheInstance(InstanceVO instance) { instance.setVendor(InstanceVendor.APACHE); - if (instance.getName() == null || instance.getName().isBlank()) { - throw new BusinessException(400, "InstanceVO name is required"); - } + instance.setName(requireInstanceName(instance.getName())); instance.setEndpoint(requireValidEndpoint(instance.getEndpoint())); if (instance.getType() == null) { throw new BusinessException(400, "InstanceVO type is required"); @@ -131,10 +134,11 @@ public class InstanceService { } CloudInstanceDetailVO detail = providerRegistry.catalogFor(InstanceVendor.ALIYUN) .getCloudInstance(instance.getCredentialId(), instance.getRegionId(), instance.getCloudInstanceId()); - if (instance.getName() == null || instance.getName().isBlank()) { + if (!StringUtils.hasText(instance.getName())) { instance.setName(detail.getInstanceName() != null && !detail.getInstanceName().isBlank() ? detail.getInstanceName() : detail.getInstanceId()); } + instance.setName(requireInstanceName(instance.getName())); instance.setType(InstanceType.PROXY); instance.setEndpoint(resolveEndpoint(detail)); } @@ -176,6 +180,13 @@ public class InstanceService { return normalized; } + private String requireInstanceName(String name) { + if (!StringUtils.hasText(name)) { + throw new BusinessException(400, "InstanceVO name is required"); + } + return name.trim(); + } + public InstanceVO updateInstance(InstanceVO instance) { requireInstance(instance); log.info("Updating instance: {}", instance.getId()); @@ -194,7 +205,7 @@ public class InstanceService { InstanceVO updated = copyOf(existing); boolean cloudInstance = existing.getVendor() != null && existing.getVendor() != InstanceVendor.APACHE; if (instance.getName() != null) { - updated.setName(instance.getName()); + updated.setName(requireInstanceName(instance.getName())); } if (!cloudInstance) { if (instance.getType() != null) { @@ -211,6 +222,8 @@ public class InstanceService { InstanceVO saved = instanceRepository.save(updated); releaseApacheEndpointIfUnused(existing, saved.getEndpoint()); + recordAudit("UPDATE_INSTANCE", "INSTANCE", saved.getId(), null, + instanceAuditDetail(saved)); return saved; } @@ -235,6 +248,8 @@ public class InstanceService { } instanceRepository.deleteById(id); releaseApacheEndpointIfUnused(existing, null); + recordAudit("DELETE_INSTANCE", "INSTANCE", id, null, + instanceAuditDetail(existing)); } private void requireInstance(InstanceVO instance) { @@ -243,6 +258,11 @@ public class InstanceService { } } + private String instanceAuditDetail(InstanceVO instance) { + InstanceVendor vendor = instance.getVendor() == null ? InstanceVendor.APACHE : instance.getVendor(); + return "name=" + instance.getName() + ", vendor=" + vendor + ", type=" + instance.getType(); + } + private InstanceVO copyOf(InstanceVO instance) { InstanceVO copy = InstanceVO.builder() .name(instance.getName()) @@ -290,4 +310,15 @@ public class InstanceService { String normalizedEndpoint = MqAdminExtFactory.normalizeNamesrvAddr(endpoint); return normalizedEndpoint.isEmpty() ? null : normalizedEndpoint; } + + private void recordAudit(String operation, String resourceType, String resourceName, + String clusterId, String detail) { + try { + operationAuditService.record(operation, resourceType, resourceName, clusterId, detail, "SUCCESS", null); + } catch (Exception auditFailure) { + log.warn("Failed to record audit operation={} resource={}: {}", operation, resourceName, + auditFailure.getMessage()); + } + } + } 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 c3ada9d7..1916f62b 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 @@ -20,6 +20,7 @@ import org.springframework.util.StringUtils; import org.apache.rocketmq.studio.common.exception.BusinessException; import org.apache.rocketmq.studio.common.util.CredentialUtils; +import org.apache.rocketmq.studio.audit.OperationAuditService; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -34,6 +35,7 @@ import java.util.UUID; public class AclService { private final AclRepository aclRepository; + private final OperationAuditService operationAuditService; public List<AclRuleVO> listRules(String clusterId, String principal) { @@ -52,7 +54,9 @@ public class AclService { } rule.setId(UUID.randomUUID().toString()); rule.setCreatedAt(LocalDateTime.now()); - return aclRepository.saveRule(rule); + AclRuleVO saved = aclRepository.saveRule(rule); + auditRule("CREATE_ACL_RULE", saved); + return saved; } public AclRuleVO updateRule(AclRuleVO rule) { @@ -60,8 +64,10 @@ public class AclService { throw new BusinessException(400, "ACL rule id is required"); } log.info("Updating ACL rule id={}, principal={}", rule.getId(), rule.getPrincipal()); - return aclRepository.replaceRule(rule) + AclRuleVO saved = aclRepository.replaceRule(rule) .orElseThrow(() -> new BusinessException(404, "ACL rule not found: " + rule.getId())); + auditRule("UPDATE_ACL_RULE", saved); + return saved; } public void deleteRule(String id) { @@ -69,6 +75,7 @@ public class AclService { if (!aclRepository.deleteRule(id)) { throw new BusinessException(404, "ACL rule not found: " + id); } + recordAudit("DELETE_ACL_RULE", "ACL_RULE", id, null, null); } @@ -89,7 +96,9 @@ public class AclService { user.setAccessKey(UUID.randomUUID().toString().replace("-", "")); user.setSecretKey(UUID.randomUUID().toString().replace("-", "")); user.setCreatedAt(LocalDateTime.now()); - return aclRepository.saveUser(user); + AclUserVO saved = aclRepository.saveUser(user); + auditUser("CREATE_ACL_USER", saved); + return saved; } public AclUserVO updateUser(UpdateAclUserDTO user) { @@ -108,7 +117,9 @@ public class AclService { .clusters(user.getClusters() == null ? existing.getClusters() : user.getClusters()) .createdAt(existing.getCreatedAt()) .build(); - return maskCredentials(aclRepository.saveUser(merged)); + AclUserVO saved = aclRepository.saveUser(merged); + auditUser("UPDATE_ACL_USER", saved); + return maskCredentials(saved); } public void deleteUser(String id) { @@ -116,6 +127,7 @@ public class AclService { if (!aclRepository.deleteUser(id)) { throw new BusinessException(404, "ACL user not found: " + id); } + recordAudit("DELETE_ACL_USER", "ACL_USER", id, null, null); } /** @@ -143,4 +155,25 @@ public class AclService { .build(); } + private void auditRule(String operation, AclRuleVO rule) { + recordAudit(operation, "ACL_RULE", rule.getId(), null, + "principal=" + rule.getPrincipal()); + } + + private void auditUser(String operation, AclUserVO user) { + recordAudit(operation, "ACL_USER", user.getId(), null, + "username=" + user.getUsername() + ", admin=" + user.isAdmin()); + } + + + private void recordAudit(String operation, String resourceType, String resourceName, + String clusterId, String detail) { + try { + operationAuditService.record(operation, resourceType, resourceName, clusterId, detail, "SUCCESS", null); + } catch (Exception auditFailure) { + log.warn("Failed to record audit operation={} resource={}: {}", operation, resourceName, + auditFailure.getMessage()); + } + } + } diff --git a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java index b67b9557..5e5a7da5 100644 --- a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java +++ b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java @@ -45,6 +45,9 @@ public class MessageService { } public TraceRecordVO getMessageTrace(String instanceId, String msgId) { + if (!StringUtils.hasText(msgId)) { + throw new BusinessException(400, "msgId is required"); + } log.info("Getting message trace: msgId={}", msgId); return providerRegistry.byInstanceId(instanceId) .map(provider -> provider.getMessageTrace(instanceId, msgId)) @@ -58,6 +61,9 @@ public class MessageService { if (hasKey && !hasTopic) { throw new BusinessException(400, "topic is required when key is specified"); } + if (hasMessageId && !hasTopic) { + throw new BusinessException(400, "topic is required when msgId is specified"); + } if (!hasTopic && !hasMessageId) { throw new BusinessException(400, "topic or msgId is required"); } diff --git a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java index d67ace2a..31b70cb7 100644 --- a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java +++ b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java @@ -20,6 +20,7 @@ import org.apache.rocketmq.studio.provider.apache.AdminClient; import org.apache.rocketmq.studio.provider.apache.MetadataProvider; import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor; import org.apache.rocketmq.studio.common.exception.BusinessException; +import org.springframework.util.StringUtils; import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO; import org.apache.rocketmq.studio.instance.group.QueueProgressVO; import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO; @@ -49,7 +50,7 @@ public class MetadataService { } public List<TopicVO> listTopics(String instanceId, String clusterId, String type, String search) { - if (isBlank(instanceId) && !isBlank(clusterId)) { + if (!StringUtils.hasText(instanceId) && StringUtils.hasText(clusterId)) { // legacy cluster-scoped read kept for AI tool handlers return metadataProvider.listTopics( normalizeFilter(clusterId), normalizeFilter(type), normalizeFilter(search)); @@ -84,11 +85,12 @@ public class MetadataService { } public List<BrokerRouteVO> getTopicRoutes(String instanceId, String name) { + String topicName = requireName(name, "topic name"); if (resolve(instanceId).vendor() != InstanceVendor.APACHE) { // broker routing does not apply to serverless cloud instances return List.of(); } - return metadataProvider.getTopicRoutes(instanceId, name); + return metadataProvider.getTopicRoutes(instanceId, topicName); } @@ -97,7 +99,8 @@ public class MetadataService { } public List<TopicConsumerVO> getTopicConsumers(String instanceId, String name) { - return resolve(instanceId).getTopicConsumers(instanceId, name); + String topicName = requireName(name, "topic name"); + return resolve(instanceId).getTopicConsumers(instanceId, topicName); } @@ -117,7 +120,7 @@ public class MetadataService { } public List<ConsumerGroupVO> listConsumerGroups(String instanceId, String clusterId, String search) { - if (isBlank(instanceId) && !isBlank(clusterId)) { + if (!StringUtils.hasText(instanceId) && StringUtils.hasText(clusterId)) { return metadataProvider.listConsumerGroups(normalizeFilter(clusterId), normalizeFilter(search)); } return resolve(instanceId).listConsumerGroups(instanceId, normalizeFilter(search)); @@ -129,10 +132,11 @@ public class MetadataService { } public ConsumerGroupVO getConsumerGroup(String instanceId, String name) { + String groupName = requireName(name, "consumer group name"); if (resolve(instanceId).vendor() != InstanceVendor.APACHE) { throw new BusinessException(501, "Consumer group detail is not supported for cloud instances"); } - return adminClient.getConsumerGroup(instanceId, name); + return adminClient.getConsumerGroup(instanceId, groupName); } @@ -141,7 +145,8 @@ public class MetadataService { } public List<QueueProgressVO> getGroupProgress(String instanceId, String name) { - return resolve(instanceId).getGroupProgress(instanceId, name); + String groupName = requireName(name, "consumer group name"); + return resolve(instanceId).getGroupProgress(instanceId, groupName); } @@ -150,7 +155,8 @@ public class MetadataService { } public List<SubscriptionEntryVO> getGroupSubscriptions(String instanceId, String name) { - return resolve(instanceId).getGroupSubscriptions(instanceId, name); + String groupName = requireName(name, "consumer group name"); + return resolve(instanceId).getGroupSubscriptions(instanceId, groupName); } @@ -187,12 +193,8 @@ public class MetadataService { .orElseGet(() -> providerRegistry.forVendor(InstanceVendor.APACHE)); } - private static boolean isBlank(String value) { - return value == null || value.isBlank(); - } - private String normalizeFilter(String value) { - return value == null || value.isBlank() ? null : value.trim(); + return !StringUtils.hasText(value) ? null : value.trim(); } private void requireTopic(TopicVO topic) { @@ -206,4 +208,11 @@ public class MetadataService { throw new BusinessException(400, "Topic send message request is required"); } } + + private String requireName(String value, String fieldName) { + if (!StringUtils.hasText(value)) { + throw new BusinessException(400, fieldName + " is required"); + } + return value.trim(); + } } diff --git a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java index 053bd567..512371a9 100644 --- a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java +++ b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java @@ -17,6 +17,7 @@ package org.apache.rocketmq.studio.ops.alert; import org.apache.rocketmq.studio.common.exception.BusinessException; +import org.apache.rocketmq.studio.audit.OperationAuditService; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -34,6 +35,7 @@ public class AlertService { private final AlertRepository alertRepository; private final AlertRuleAssetService alertRuleAssetService; + private final OperationAuditService operationAuditService; public List<AlertRuleVO> listRules() { @@ -42,7 +44,9 @@ public class AlertService { } public String exportPrometheusRulesYaml() { - List<AlertRuleVO> rules = alertRepository.findAllRules(); + List<AlertRuleVO> rules = alertRepository.findAllRules().stream() + .filter(AlertRuleVO::isEnabled) + .toList(); List<PrometheusAlertRule> prometheusRules = rules.isEmpty() ? defaultPrometheusRules() : rules.stream().map(this::toPrometheusRule).toList(); @@ -83,7 +87,9 @@ public class AlertService { } log.info("Creating alert rule: {}", rule.getName()); rule.setId(UUID.randomUUID().toString()); - return alertRepository.saveRule(rule); + AlertRuleVO saved = alertRepository.saveRule(rule); + auditRule("CREATE_ALERT_RULE", saved, null); + return saved; } @@ -97,6 +103,7 @@ public class AlertService { if (!alertRepository.replaceRule(rule)) { throw ruleNotFound(id); } + auditRule("UPDATE_ALERT_RULE", rule, null); return rule; } @@ -109,7 +116,9 @@ public class AlertService { .findFirst() .orElseThrow(() -> new org.apache.rocketmq.studio.common.exception.BusinessException(404, "Alert rule not found: " + id)); rule.setEnabled(enabled); - return alertRepository.saveRule(rule); + AlertRuleVO saved = alertRepository.saveRule(rule); + auditRule("TOGGLE_ALERT_RULE", saved, "enabled=" + enabled); + return saved; } @@ -118,6 +127,7 @@ public class AlertService { if (!alertRepository.deleteRule(id)) { throw ruleNotFound(id); } + recordAudit("DELETE_ALERT_RULE", "ALERT_RULE", id, null, null); } @@ -135,13 +145,19 @@ public class AlertService { .findFirst() .orElseThrow(() -> new org.apache.rocketmq.studio.common.exception.BusinessException(404, "System alert not found: " + id)); alert.setAcknowledged(true); - return alertRepository.saveAlert(alert); + SystemAlertVO saved = alertRepository.saveAlert(alert); + recordAudit("ACKNOWLEDGE_SYSTEM_ALERT", "SYSTEM_ALERT", saved.getId(), null, + "acknowledged=true"); + return saved; } public int clearAcknowledged() { log.info("Clearing acknowledged system alerts"); - return alertRepository.deleteAcknowledgedAlerts(); + int deleted = alertRepository.deleteAcknowledgedAlerts(); + recordAudit("CLEAR_ACKNOWLEDGED_SYSTEM_ALERTS", "SYSTEM_ALERT", null, null, + "deleted=" + deleted); + return deleted; } private List<PrometheusAlertRule> defaultPrometheusRules() { @@ -202,6 +218,12 @@ public class AlertService { selector.append(label).append("=\"").append(escapeDoubleQuotedValue(value.trim())).append('"'); } + private void auditRule(String operation, AlertRuleVO rule, String detail) { + String auditDetail = detail == null ? "name=" + rule.getName() : detail; + recordAudit(operation, "ALERT_RULE", rule.getId(), null, + auditDetail); + } + private String severity(AlertRuleVO rule) { String severity = rule.getSeverity(); if (hasText(severity)) { @@ -296,4 +318,15 @@ public class AlertService { private BusinessException ruleNotFound(String id) { return new BusinessException(404, "Alert rule not found: " + id); } + + private void recordAudit(String operation, String resourceType, String resourceName, + String clusterId, String detail) { + try { + operationAuditService.record(operation, resourceType, resourceName, clusterId, detail, "SUCCESS", null); + } catch (Exception auditFailure) { + log.warn("Failed to record audit operation={} resource={}: {}", operation, resourceName, + auditFailure.getMessage()); + } + } + } diff --git a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java index 2f59fb11..e7a89ea6 100644 --- a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java +++ b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java @@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.provider.apache; import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.client.producer.SendResult; +import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.common.TopicConfig; import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.remoting.protocol.body.ClusterInfo; @@ -335,6 +336,10 @@ public class RocketMQAdminClientImpl implements AdminClient { } SendResult sendResult = producer.send(msg); + if (sendResult == null || sendResult.getSendStatus() != SendStatus.SEND_OK) { + String status = sendResult == null ? "null" : String.valueOf(sendResult.getSendStatus()); + throw new BusinessException(502, "Message send did not succeed: " + status); + } // The message is already delivered by now; an audit write failure must not turn a // successful send into an error, or callers would retry and duplicate the message. @@ -346,6 +351,9 @@ public class RocketMQAdminClientImpl implements AdminClient { .sendTime(System.currentTimeMillis()) .offsetMsgId(sendResult.getOffsetMsgId()) .build(); + } catch (BusinessException e) { + recordAudit("SEND_MESSAGE", request.getTopic(), e.getMessage(), "FAILED"); + throw e; } catch (Exception e) { recordAudit("SEND_MESSAGE", request.getTopic(), e.getMessage(), "FAILED"); throw new BusinessException(500, "Failed to send message: " + e.getMessage()); diff --git a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java index 1ef0d70d..420a796a 100644 --- a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java +++ b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java @@ -263,8 +263,11 @@ public class RocketMQMessageProvider implements MessageProvider { parseTraceBody(traceMessage.getBody(), msgId, nodes, consumerStatus); } } + } catch (BusinessException e) { + throw e; } catch (Exception e) { log.warn("Trace query for msgId={} failed: {}", msgId, e.getMessage()); + throw new BusinessException(502, "Failed to query message trace: " + e.getMessage()); } recordTraceQuery(instanceId, msgId, null, nodes.size(), consumerStatus.size()); diff --git a/server/src/main/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialService.java b/server/src/main/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialService.java index 10d7ca2a..0decc7ec 100644 --- a/server/src/main/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialService.java +++ b/server/src/main/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialService.java @@ -20,6 +20,7 @@ import org.springframework.util.StringUtils; import org.apache.rocketmq.studio.common.exception.BusinessException; import org.apache.rocketmq.studio.common.util.CredentialUtils; +import org.apache.rocketmq.studio.audit.OperationAuditService; import org.apache.rocketmq.studio.instance.InstanceRepository; import org.apache.rocketmq.studio.provider.alibaba.AliyunClientFactory; import org.springframework.stereotype.Service; @@ -38,6 +39,7 @@ public class CloudCredentialService { private final CloudCredentialRepository credentialRepository; private final InstanceRepository instanceRepository; private final AliyunClientFactory aliyunClientFactory; + private final OperationAuditService operationAuditService; public List<CloudCredentialVO> listMasked() { log.info("Listing cloud credentials (masked)"); @@ -69,7 +71,10 @@ public class CloudCredentialService { credential.setId(UUID.randomUUID().toString()); credential.setCreatedAt(LocalDateTime.now()); credential.setUpdatedAt(LocalDateTime.now()); - return maskAccessKey(credentialRepository.save(credential)); + CloudCredentialVO saved = credentialRepository.save(credential); + recordAudit("CREATE_CLOUD_CREDENTIAL", "CLOUD_CREDENTIAL", saved.getId(), null, + credentialAuditDetail(saved)); + return maskAccessKey(saved); } public CloudCredentialVO update(UpdateCloudCredentialDTO request) { @@ -94,6 +99,8 @@ public class CloudCredentialService { existing.setUpdatedAt(LocalDateTime.now()); CloudCredentialVO saved = credentialRepository.save(existing); invalidateAliyunClients(saved); + recordAudit("UPDATE_CLOUD_CREDENTIAL", "CLOUD_CREDENTIAL", saved.getId(), null, + credentialAuditDetail(saved)); return maskAccessKey(saved); } @@ -109,6 +116,8 @@ public class CloudCredentialService { } credentialRepository.deleteById(id); invalidateAliyunClients(existing); + recordAudit("DELETE_CLOUD_CREDENTIAL", "CLOUD_CREDENTIAL", id, null, + credentialAuditDetail(existing)); } public CloudCredentialVO reveal(String id) { @@ -131,9 +140,25 @@ public class CloudCredentialService { masked.setUpdatedAt(credential.getUpdatedAt()); return masked; } + + private String credentialAuditDetail(CloudCredentialVO credential) { + return "name=" + credential.getName() + ", vendor=" + credential.getVendor(); + } + private void invalidateAliyunClients(CloudCredentialVO credential) { if (credential.getVendor() == org.apache.rocketmq.studio.common.domain.enums.InstanceVendor.ALIYUN) { aliyunClientFactory.invalidateCredential(credential.getId()); } } + + private void recordAudit(String operation, String resourceType, String resourceName, + String clusterId, String detail) { + try { + operationAuditService.record(operation, resourceType, resourceName, clusterId, detail, "SUCCESS", null); + } catch (Exception auditFailure) { + log.warn("Failed to record audit operation={} resource={}: {}", operation, resourceName, + auditFailure.getMessage()); + } + } + } diff --git a/server/src/main/resources/application-dev.yml b/server/src/main/resources/application-dev.yml index f7efb59e..c6af6056 100644 --- a/server/src/main/resources/application-dev.yml +++ b/server/src/main/resources/application-dev.yml @@ -5,6 +5,10 @@ spring: username: sa password: driver-class-name: org.h2.Driver + sql: + init: + mode: always + schema-locations: classpath:db/schema.sql h2: console: enabled: true diff --git a/server/src/main/resources/db/schema.sql b/server/src/main/resources/db/schema.sql index 6fc1f01d..a9a4c15a 100644 --- a/server/src/main/resources/db/schema.sql +++ b/server/src/main/resources/db/schema.sql @@ -24,6 +24,10 @@ CREATE TABLE IF NOT EXISTS rmq_instance ( remark VARCHAR(255), type VARCHAR(32) NOT NULL COMMENT 'PROXY/DIRECT', endpoint VARCHAR(512) NOT NULL, + vendor VARCHAR(32), + cloud_instance_id VARCHAR(128), + credential_id VARCHAR(64), + region_id VARCHAR(128), created_at DATETIME DEFAULT CURRENT_TIMESTAMP, updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; @@ -44,7 +48,7 @@ CREATE TABLE IF NOT EXISTS rmq_topic ( created_at DATETIME DEFAULT CURRENT_TIMESTAMP, updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY uk_cluster_topic (cluster_id, name), - INDEX idx_instance (instance_id) + INDEX idx_topic_instance (instance_id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; -- 4. Consumer Group 管理记录 @@ -61,7 +65,7 @@ CREATE TABLE IF NOT EXISTS rmq_group ( created_at DATETIME DEFAULT CURRENT_TIMESTAMP, updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY uk_cluster_group (cluster_id, name), - INDEX idx_instance (instance_id) + INDEX idx_group_instance (instance_id) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; -- 5. K8s 证书管理 @@ -95,7 +99,7 @@ CREATE TABLE IF NOT EXISTS rmq_message_query ( cluster_id VARCHAR(255), queried_by VARCHAR(64), queried_at DATETIME DEFAULT CURRENT_TIMESTAMP, - INDEX idx_queried_at (queried_at), + INDEX idx_message_query_queried_at (queried_at), INDEX idx_topic (topic) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; @@ -110,7 +114,7 @@ CREATE TABLE IF NOT EXISTS rmq_trace_query ( queried_by VARCHAR(64), queried_at DATETIME DEFAULT CURRENT_TIMESTAMP, INDEX idx_msg_id (msg_id), - INDEX idx_queried_at (queried_at) + INDEX idx_trace_query_queried_at (queried_at) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; -- 8. 操作审计日志(所有写操作) @@ -132,7 +136,7 @@ CREATE TABLE IF NOT EXISTS rmq_operation_audit ( -- 9. 通用设置(单行) CREATE TABLE IF NOT EXISTS rmq_settings ( - id VARCHAR(16) PRIMARY KEY DEFAULT 'singleton', + id VARCHAR(16) DEFAULT 'singleton' PRIMARY KEY, json TEXT NOT NULL COMMENT 'GeneralSettingsVO JSON', updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; 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 2881946c..125d65be 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/StudioApplicationTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/StudioApplicationTest.java @@ -18,13 +18,20 @@ 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.persistence.mapper.RmqInstanceMapper; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.autoconfigure.web.servlet.AutoConfigureMockMvc; +import org.springframework.test.web.servlet.MockMvc; import static org.assertj.core.api.Assertions.assertThat; +import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @SpringBootTest +@AutoConfigureMockMvc class StudioApplicationTest { @Autowired @@ -33,9 +40,20 @@ class StudioApplicationTest { @Autowired private ToolGatewayService toolGatewayService; + @Autowired + private RmqInstanceMapper instanceMapper; + + @Autowired + private MockMvc mockMvc; + @Test - void applicationContextLoadsWithToolGateway() { + void applicationContextLoadsWithInitializedDevSchema() throws Exception { assertThat(toolCatalog.getVersion()).isEqualTo("1.0.0"); assertThat(toolGatewayService.discover(null)).isNotEmpty(); + assertThat(instanceMapper.selectList(null)).isEmpty(); + + mockMvc.perform(get("/api/instances")) + .andExpect(status().isOk()) + .andExpect(jsonPath("$.data").isEmpty()); } } 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 efc118d3..913d03ce 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 @@ -17,6 +17,7 @@ package org.apache.rocketmq.studio.cluster.broker; import org.apache.rocketmq.studio.common.exception.BusinessException; +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.Test; @@ -32,6 +33,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.isNull; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; @ExtendWith(MockitoExtension.class) @@ -80,4 +82,23 @@ class RuntimeAdminClientResolverTest { assertThat(result).isEqualTo("done"); verify(adminFactory).execute(eq("namesrv-b:9876"), isNull(), any()); } + + @Test + void rejectsCloudInstancesBeforeResolvingOrExecutingAdminClient() { + InstanceVO instance = InstanceVO.builder() + .vendor(InstanceVendor.ALIYUN) + .endpoint("cloud-endpoint:9876") + .build(); + instance.setId("cloud-instance"); + when(instanceRepository.findById("cloud-instance")).thenReturn(Optional.of(instance)); + RuntimeAdminClientResolver resolver = new RuntimeAdminClientResolver(instanceRepository, adminFactory); + + assertThatThrownBy(() -> resolver.resolveEndpoint("cloud-instance")) + .isInstanceOf(BusinessException.class) + .hasMessage("Runtime AdminClient only supports Apache instances: cloud-instance"); + assertThatThrownBy(() -> resolver.execute(instance, admin -> "unused")) + .isInstanceOf(BusinessException.class) + .hasMessage("Runtime AdminClient only supports Apache instances: cloud-instance"); + verifyNoInteractions(adminFactory); + } } diff --git a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientServiceTest.java b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientServiceTest.java index ce472901..e6152373 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientServiceTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientServiceTest.java @@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.cluster.client; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; +import org.apache.rocketmq.studio.common.exception.BusinessException; import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; @@ -25,7 +26,9 @@ import org.mockito.junit.jupiter.MockitoExtension; import java.util.List; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; @ExtendWith(MockitoExtension.class) @@ -52,12 +55,21 @@ class ClientServiceTest { } @Test - void listConnectionsShouldTreatBlankFiltersAsUnspecified() { - when(clientProvider.findConnections(null, null, null)).thenReturn(List.of()); + void listConnectionsShouldTreatBlankOptionalFiltersAsUnspecified() { + when(clientProvider.findConnections("instance-1", null, null)).thenReturn(List.of()); - List<ClientConnectionVO> result = clientService.listConnections(" ", " ", "\t"); + List<ClientConnectionVO> result = clientService.listConnections("instance-1", " ", "\t"); assertThat(result).isEmpty(); - verify(clientProvider).findConnections(null, null, null); + verify(clientProvider).findConnections("instance-1", null, null); + } + + @Test + void listConnectionsShouldRejectBlankInstanceIdBeforeQueryingProvider() { + assertThatThrownBy(() -> clientService.listConnections(" ", null, null)) + .isInstanceOf(BusinessException.class) + .hasMessage("instanceId is required"); + + verifyNoInteractions(clientProvider); } } diff --git a/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertServiceTest.java b/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertServiceTest.java index d93bc32e..72382060 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertServiceTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertServiceTest.java @@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.cluster.k8s; import org.apache.rocketmq.studio.common.domain.enums.CertStatus; import org.apache.rocketmq.studio.common.domain.enums.CertType; import org.apache.rocketmq.studio.common.exception.BusinessException; +import org.apache.rocketmq.studio.audit.OperationAuditService; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -38,6 +39,7 @@ import java.util.Optional; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; @@ -51,13 +53,16 @@ class K8sCertServiceTest { @Mock private K8sCertRepository k8sCertRepository; + @Mock + private OperationAuditService operationAuditService; + private K8sCertService k8sCertService; private K8sCertVO sampleCert; @BeforeEach void setUp() { - k8sCertService = new K8sCertService(k8sCertRepository, CLOCK); + k8sCertService = new K8sCertService(k8sCertRepository, operationAuditService, CLOCK); sampleCert = K8sCertVO.builder() .name("rocketmq-tls") .namespace("mq-system") @@ -169,6 +174,9 @@ class K8sCertServiceTest { assertThat(result.getCreatedAt()).isEqualTo(now); assertThat(result.getUpdatedAt()).isEqualTo(now); verify(k8sCertRepository).save(any(K8sCertVO.class)); + verify(operationAuditService).record(eq("CREATE_K8S_CERTIFICATE"), eq("K8S_CERTIFICATE"), + eq(result.getId()), eq(null), eq("name=new-tls-cert, namespace=default, cluster=test-cluster"), + eq("SUCCESS"), eq(null)); } @Test @@ -246,6 +254,9 @@ class K8sCertServiceTest { assertThat(sampleCert.getType()).isEqualTo(CertType.TLS); assertThat(sampleCert.getUpdatedAt()).isEqualTo(LocalDateTime.of(2025, 1, 2, 0, 0)); verify(k8sCertRepository).save(any(K8sCertVO.class)); + verify(operationAuditService).record(eq("UPDATE_K8S_CERTIFICATE"), eq("K8S_CERTIFICATE"), + eq("cert-1"), eq(null), eq("name=updated-name, namespace=new-namespace, cluster=new-cluster"), + eq("SUCCESS"), eq(null)); } @Test @@ -351,6 +362,9 @@ class K8sCertServiceTest { assertThat(sampleCert.getNotBefore()).isEqualTo(originalNotBefore); assertThat(sampleCert.getNotAfter()).isEqualTo(originalNotAfter); verify(k8sCertRepository).save(any(K8sCertVO.class)); + verify(operationAuditService).record(eq("RENEW_K8S_CERTIFICATE"), eq("K8S_CERTIFICATE"), + eq("cert-1"), eq(null), eq("name=rocketmq-tls, namespace=mq-system, cluster=prod-cluster"), + eq("SUCCESS"), eq(null)); } @Test @@ -396,6 +410,8 @@ class K8sCertServiceTest { k8sCertService.deleteCert(command); verify(k8sCertRepository).deleteById("cert-1"); + verify(operationAuditService).record(eq("DELETE_K8S_CERTIFICATE"), eq("K8S_CERTIFICATE"), + eq("cert-1"), eq(null), eq(null), eq("SUCCESS"), eq(null)); } @Test 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 01065024..ec2e0ec1 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 @@ -20,6 +20,7 @@ package org.apache.rocketmq.studio.instance; import org.apache.rocketmq.studio.provider.credential.CloudCredentialRepository; import org.apache.rocketmq.studio.provider.credential.CloudCredentialVO; import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory; +import org.apache.rocketmq.studio.audit.OperationAuditService; import org.apache.rocketmq.studio.common.domain.enums.InstanceType; import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor; import org.apache.rocketmq.studio.common.exception.BusinessException; @@ -40,6 +41,8 @@ import java.util.Optional; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.argThat; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; @@ -63,6 +66,9 @@ class InstanceServiceTest { @Mock private MqAdminExtFactory adminFactory; + @Mock + private OperationAuditService operationAuditService; + @InjectMocks private InstanceService instanceService; @@ -259,6 +265,10 @@ class InstanceServiceTest { assertThat(result.getUpdatedAt()).isNotNull(); assertThat(result.getName()).isEqualTo("new-instance"); verify(instanceRepository).save(any(InstanceVO.class)); + verify(operationAuditService).record(eq("CREATE_INSTANCE"), eq("INSTANCE"), eq(result.getId()), eq(null), + argThat(detail -> detail.equals("name=new-instance, vendor=APACHE, type=PROXY") + && !detail.contains("10.0.1.1:8080")), + eq("SUCCESS"), eq(null)); } @Test @@ -327,6 +337,15 @@ class InstanceServiceTest { assertThat(instanceService.createInstance(input).getEndpoint()).isEqualTo("namesrv:9876"); } + @Test + void createInstanceShouldTrimNameBeforeSaving() { + InstanceVO input = InstanceVO.builder().name(" production ").type(InstanceType.PROXY) + .endpoint("namesrv:9876").build(); + when(instanceRepository.save(any(InstanceVO.class))).thenAnswer(invocation -> invocation.getArgument(0)); + + assertThat(instanceService.createInstance(input).getName()).isEqualTo("production"); + } + @Test void updateInstanceShouldMergeFieldsOntoExisting() { LocalDateTime originalCreatedAt = LocalDateTime.of(2025, 1, 2, 3, 4, 5); @@ -367,6 +386,20 @@ class InstanceServiceTest { assertThat(existing.getName()).isEqualTo("old-name"); assertThat(existing.getRemark()).isEqualTo("old remark"); assertThat(existing.getUpdatedAt()).isEqualTo(originalUpdatedAt); + verify(operationAuditService).record(eq("UPDATE_INSTANCE"), eq("INSTANCE"), eq("inst-1"), eq(null), + eq("name=new-name, vendor=APACHE, type=PROXY"), eq("SUCCESS"), eq(null)); + } + + @Test + void updateInstanceShouldTrimNameBeforeSaving() { + InstanceVO existing = InstanceVO.builder().name("old-name").endpoint("namesrv:9876").build(); + existing.setId("inst-1"); + InstanceVO update = InstanceVO.builder().name(" production ").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)); + + assertThat(instanceService.updateInstance(update).getName()).isEqualTo("production"); } @Test @@ -592,6 +625,8 @@ class InstanceServiceTest { instanceService.deleteInstance("inst-1"); verify(instanceRepository).deleteById("inst-1"); + verify(operationAuditService).record(eq("DELETE_INSTANCE"), eq("INSTANCE"), eq("inst-1"), eq(null), + eq("name=to-delete, vendor=APACHE, type=null"), eq("SUCCESS"), eq(null)); } @Test 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 0a78fbcc..b3cd90c4 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 @@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.instance.acl; import org.apache.rocketmq.studio.common.exception.BusinessException; +import org.apache.rocketmq.studio.audit.OperationAuditService; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -37,6 +38,8 @@ import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.argThat; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -47,6 +50,9 @@ class AclServiceTest { @Mock private AclRepository aclRepository; + @Mock + private OperationAuditService operationAuditService; + @InjectMocks private AclService aclService; @@ -107,6 +113,8 @@ class AclServiceTest { assertThat(result.getPrincipal()).isEqualTo("user1"); assertThat(result.getResource()).isEqualTo("topic-1"); verify(aclRepository).saveRule(any(AclRuleVO.class)); + verify(operationAuditService).record(eq("CREATE_ACL_RULE"), eq("ACL_RULE"), eq(result.getId()), eq(null), + eq("principal=user1"), eq("SUCCESS"), eq(null)); } @Test @@ -115,6 +123,8 @@ class AclServiceTest { aclService.deleteRule("rule-1"); verify(aclRepository).deleteRule("rule-1"); + verify(operationAuditService).record(eq("DELETE_ACL_RULE"), eq("ACL_RULE"), eq("rule-1"), eq(null), + eq(null), eq("SUCCESS"), eq(null)); } @Test @@ -148,6 +158,8 @@ class AclServiceTest { assertThat(result.getDecision()).isEqualTo("DENY"); verify(aclRepository).replaceRule(input); verify(aclRepository, never()).saveRule(any(AclRuleVO.class)); + verify(operationAuditService).record(eq("UPDATE_ACL_RULE"), eq("ACL_RULE"), eq("rule-1"), eq(null), + eq("principal=user1"), eq("SUCCESS"), eq(null)); } @Test @@ -307,6 +319,10 @@ class AclServiceTest { assertThat(result.getCreatedAt()).isNotNull(); assertThat(result.getUsername()).isEqualTo("newuser"); verify(aclRepository).saveUser(any(AclUserVO.class)); + verify(operationAuditService).record(eq("CREATE_ACL_USER"), eq("ACL_USER"), eq(result.getId()), eq(null), + argThat(detail -> detail.equals("username=newuser, admin=false") + && !detail.contains(result.getAccessKey()) && !detail.contains(result.getSecretKey())), + eq("SUCCESS"), eq(null)); } @Test @@ -315,6 +331,8 @@ class AclServiceTest { aclService.deleteUser("user-1"); verify(aclRepository).deleteUser("user-1"); + verify(operationAuditService).record(eq("DELETE_ACL_USER"), eq("ACL_USER"), eq("user-1"), eq(null), + eq(null), eq("SUCCESS"), eq(null)); } @Test @@ -365,6 +383,8 @@ class AclServiceTest { verify(aclRepository).saveUser(captor.capture()); assertThat(captor.getValue().getAccessKey()).isEqualTo("access-key-123456"); assertThat(captor.getValue().getSecretKey()).isEqualTo("secret-key-987654"); + verify(operationAuditService).record(eq("UPDATE_ACL_USER"), eq("ACL_USER"), eq("user-1"), eq(null), + eq("username=newuser, admin=true"), eq("SUCCESS"), eq(null)); } @Test diff --git a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java index 8d237eb6..8d96cc8b 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java @@ -14,13 +14,9 @@ import org.apache.rocketmq.studio.common.exception.BusinessException; import org.apache.rocketmq.studio.provider.InstanceProviderRegistry; import org.junit.jupiter.api.Test; -import java.util.Collections; - import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; -import static org.mockito.Mockito.when; class MessageServiceTest { @@ -38,16 +34,29 @@ class MessageServiceTest { } @Test - void keepsMessageIdQueryWithoutTopicValid() { + void rejectsMessageIdQueryWithoutTopicBeforeCallingProvider() { MessageProvider provider = mock(MessageProvider.class); InstanceProviderRegistry registry = mock(InstanceProviderRegistry.class); - when(provider.queryMessages(null, null, "msg-001", null, null, null, null)) - .thenReturn(Collections.emptyList()); MessageService service = new MessageService(provider, registry); - service.queryMessages(null, null, "msg-001", null, null, null, null); + assertThatThrownBy(() -> service.queryMessages(null, null, "msg-001", null, null, null, null)) + .isInstanceOf(BusinessException.class) + .hasMessage("topic is required when msgId is specified"); + + verifyNoInteractions(provider, registry); + } + + @Test + void rejectsBlankMessageTraceIdBeforeCallingProvider() { + MessageProvider provider = mock(MessageProvider.class); + InstanceProviderRegistry registry = mock(InstanceProviderRegistry.class); + MessageService service = new MessageService(provider, registry); + + assertThatThrownBy(() -> service.getMessageTrace("instance-a", " ")) + .isInstanceOf(BusinessException.class) + .hasMessage("msgId is required"); - verify(provider).queryMessages(null, null, "msg-001", null, null, null, null); + verifyNoInteractions(provider, registry); } @Test diff --git a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java index 4d6d02e1..918a7fad 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java @@ -171,6 +171,27 @@ class MetadataServiceTest { verify(apacheProvider).getTopicConsumers("instance-a", "orders"); } + @Test + void runtimeDiagnosticsShouldRejectBlankTopicAndGroupNamesBeforeProviderResolution() { + assertThatThrownBy(() -> metadataService.getTopicRoutes("instance-a", " ")) + .isInstanceOf(BusinessException.class) + .hasMessage("topic name is required"); + assertThatThrownBy(() -> metadataService.getTopicConsumers("instance-a", " ")) + .isInstanceOf(BusinessException.class) + .hasMessage("topic name is required"); + assertThatThrownBy(() -> metadataService.getConsumerGroup("instance-a", " ")) + .isInstanceOf(BusinessException.class) + .hasMessage("consumer group name is required"); + assertThatThrownBy(() -> metadataService.getGroupProgress("instance-a", " ")) + .isInstanceOf(BusinessException.class) + .hasMessage("consumer group name is required"); + assertThatThrownBy(() -> metadataService.getGroupSubscriptions("instance-a", " ")) + .isInstanceOf(BusinessException.class) + .hasMessage("consumer group name is required"); + + verifyNoInteractions(metadataProvider, adminClient, providerRegistry, apacheProvider); + } + @Test void sendMessageShouldReturnResult() { SendMessageDTO request = SendMessageDTO.builder() diff --git a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceDefaultRulesTest.java b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceDefaultRulesTest.java index 786ebf65..6d0acfda 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceDefaultRulesTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceDefaultRulesTest.java @@ -16,7 +16,9 @@ */ package org.apache.rocketmq.studio.ops.alert; +import org.apache.rocketmq.studio.audit.OperationAuditService; import org.junit.jupiter.api.Test; +import org.mockito.Mockito; import java.util.Collections; import java.util.List; @@ -36,7 +38,8 @@ class AlertServiceDefaultRulesTest { AlertRepository repository = mock(AlertRepository.class); when(repository.findAllRules()).thenReturn(Collections.emptyList()); - AlertService service = new AlertService(repository, new AlertRuleAssetService()); + AlertService service = new AlertService(repository, new AlertRuleAssetService(), + Mockito.mock(OperationAuditService.class)); String yaml = service.exportPrometheusRulesYaml(); int ruleCount = countRules(yaml); @@ -48,7 +51,8 @@ class AlertServiceDefaultRulesTest { AlertRepository repository = mock(AlertRepository.class); when(repository.findAllRules()).thenReturn(List.of()); - AlertService service = new AlertService(repository, new AlertRuleAssetService()); + AlertService service = new AlertService(repository, new AlertRuleAssetService(), + Mockito.mock(OperationAuditService.class)); String yaml = service.exportPrometheusRulesYaml(); assertTrue(yaml.contains("rocketmq-broker.rules")); diff --git a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java index 35a46555..73b440ce 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java @@ -21,6 +21,7 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.dataformat.yaml.YAMLFactory; import org.apache.rocketmq.studio.common.domain.enums.AlertLevel; import org.apache.rocketmq.studio.common.exception.BusinessException; +import org.apache.rocketmq.studio.audit.OperationAuditService; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.Mock; @@ -33,6 +34,7 @@ import java.util.List; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -43,11 +45,14 @@ class AlertServiceTest { @Mock private AlertRepository alertRepository; + @Mock + private OperationAuditService operationAuditService; + private AlertService alertService; @org.junit.jupiter.api.BeforeEach void setUp() { - alertService = new AlertService(alertRepository, new AlertRuleAssetService()); + alertService = new AlertService(alertRepository, new AlertRuleAssetService(), operationAuditService); } @Test @@ -100,6 +105,7 @@ class AlertServiceTest { .threshold(5000) .duration("3m") .description("Lag too high") + .enabled(true) .build(); when(alertRepository.findAllRules()).thenReturn(List.of(rule)); @@ -122,6 +128,7 @@ class AlertServiceTest { .threshold(1000) .duration("3m") .description("First lag alert") + .enabled(true) .build(); AlertRuleVO second = AlertRuleVO.builder() .name("Lag Alert B") @@ -130,6 +137,7 @@ class AlertServiceTest { .threshold(2000) .duration("5m") .description("Second lag alert") + .enabled(true) .build(); when(alertRepository.findAllRules()).thenReturn(List.of(first, second)); @@ -156,6 +164,7 @@ class AlertServiceTest { .clusterName("DefaultCluster") .severity("critical") .description("Slave falls behind master") + .enabled(true) .build(); when(alertRepository.findAllRules()).thenReturn(List.of(rule)); @@ -178,6 +187,7 @@ class AlertServiceTest { .brokerName("*") .clusterName(" ") .severity("fatal") + .enabled(true) .build(); when(alertRepository.findAllRules()).thenReturn(List.of(rule)); @@ -198,6 +208,7 @@ class AlertServiceTest { .clusterName("prod\"east\\dc\nline") .brokerName("broker\tone") .description("Summary \"quoted\" \\ path\nnext - Details\twith\rline") + .enabled(true) .build(); when(alertRepository.findAllRules()).thenReturn(List.of(rule)); @@ -221,6 +232,7 @@ class AlertServiceTest { .metric("rocketmq_consumer_lag_messages") .operator(">") .threshold(1) + .enabled(true) .build(); when(alertRepository.findAllRules()).thenReturn(List.of(rule)); @@ -231,6 +243,50 @@ class AlertServiceTest { assertThat(exportedRule.path("alert").asText()).isEqualTo("RocketMQAlert"); } + @Test + void exportPrometheusRulesYamlShouldExcludeDisabledRules() { + AlertRuleVO enabled = AlertRuleVO.builder() + .name("Enabled Lag Alert") + .metric("rocketmq_consumer_lag_messages") + .operator(">") + .threshold(1000) + .enabled(true) + .build(); + AlertRuleVO disabled = AlertRuleVO.builder() + .name("Disabled Lag Alert") + .metric("rocketmq_consumer_lag_messages") + .operator(">") + .threshold(2000) + .enabled(false) + .build(); + when(alertRepository.findAllRules()).thenReturn(List.of(enabled, disabled)); + + String result = alertService.exportPrometheusRulesYaml(); + + assertThat(result) + .contains("EnabledLagAlert") + .doesNotContain("DisabledLagAlert") + .doesNotContain(" > 2000"); + } + + @Test + void exportPrometheusRulesYamlShouldUseDefaultRulesWhenAllConfiguredRulesAreDisabled() { + AlertRuleVO disabled = AlertRuleVO.builder() + .name("Disabled Lag Alert") + .metric("rocketmq_consumer_lag_messages") + .operator(">") + .threshold(2000) + .enabled(false) + .build(); + when(alertRepository.findAllRules()).thenReturn(List.of(disabled)); + + String result = alertService.exportPrometheusRulesYaml(); + + assertThat(result) + .contains("RocketMQBrokerDown") + .doesNotContain("DisabledLagAlert"); + } + @Test void createRuleShouldPreserveReplicationScopeFields() { AlertRuleVO input = AlertRuleVO.builder() @@ -252,6 +308,8 @@ class AlertServiceTest { assertThat(result.getClusterName()).isEqualTo("DefaultCluster"); assertThat(result.getSeverity()).isEqualTo("critical"); verify(alertRepository).saveRule(result); + verify(operationAuditService).record(eq("CREATE_ALERT_RULE"), eq("ALERT_RULE"), eq(result.getId()), + eq(null), eq("name=Replication Lag High"), eq("SUCCESS"), eq(null)); } @Test @@ -300,6 +358,8 @@ class AlertServiceTest { assertThat(result.getId()).isEqualTo("rule-1"); assertThat(result.getThreshold()).isEqualTo(90.0); verify(alertRepository).replaceRule(update); + verify(operationAuditService).record(eq("UPDATE_ALERT_RULE"), eq("ALERT_RULE"), eq("rule-1"), + eq(null), eq("name=CPU Alert"), eq("SUCCESS"), eq(null)); } @Test @@ -357,6 +417,8 @@ class AlertServiceTest { assertThat(result.isEnabled()).isTrue(); verify(alertRepository).saveRule(result); + verify(operationAuditService).record(eq("TOGGLE_ALERT_RULE"), eq("ALERT_RULE"), eq("rule-1"), + eq(null), eq("enabled=true"), eq("SUCCESS"), eq(null)); } @Test @@ -386,6 +448,8 @@ class AlertServiceTest { alertService.deleteRule("rule-1"); verify(alertRepository).deleteRule("rule-1"); + verify(operationAuditService).record(eq("DELETE_ALERT_RULE"), eq("ALERT_RULE"), eq("rule-1"), + eq(null), eq(null), eq("SUCCESS"), eq(null)); } @Test @@ -437,6 +501,8 @@ class AlertServiceTest { assertThat(result.isAcknowledged()).isTrue(); verify(alertRepository).saveAlert(result); + verify(operationAuditService).record(eq("ACKNOWLEDGE_SYSTEM_ALERT"), eq("SYSTEM_ALERT"), eq("a1"), + eq(null), eq("acknowledged=true"), eq("SUCCESS"), eq(null)); } @Test @@ -456,6 +522,8 @@ class AlertServiceTest { assertThat(result).isEqualTo(3); verify(alertRepository).deleteAcknowledgedAlerts(); + verify(operationAuditService).record(eq("CLEAR_ACKNOWLEDGED_SYSTEM_ALERTS"), eq("SYSTEM_ALERT"), + eq(null), eq(null), eq("deleted=3"), eq("SUCCESS"), eq(null)); } @Test diff --git a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java index 2bd23a1d..804ec59b 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java @@ -19,6 +19,7 @@ import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.common.TopicConfig; import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.client.producer.SendResult; +import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.remoting.exception.RemotingTimeoutException; import org.apache.rocketmq.remoting.protocol.ResponseCode; @@ -338,6 +339,7 @@ class RocketMQAdminClientImplTest { mockConstruction(DefaultMQProducer.class, (producer, context) -> { doNothing().when(producer).start(); SendResult sendResult = new SendResult(); + sendResult.setSendStatus(SendStatus.SEND_OK); sendResult.setMsgId("msg-1"); sendResult.setOffsetMsgId("offset-1"); when(producer.send(any(Message.class))).thenReturn(sendResult); @@ -359,6 +361,7 @@ class RocketMQAdminClientImplTest { mockConstruction(DefaultMQProducer.class, (producer, context) -> { doNothing().when(producer).start(); SendResult sendResult = new SendResult(); + sendResult.setSendStatus(SendStatus.SEND_OK); sendResult.setMsgId("msg-1"); sendResult.setOffsetMsgId("offset-1"); when(producer.send(any(Message.class))).thenReturn(sendResult); @@ -375,4 +378,50 @@ class RocketMQAdminClientImplTest { verify(producer).setNamesrvAddr("10.0.0.2:9876"); } } + + @Test + void sendMessageShouldRejectNonSuccessfulSendStatus() throws Exception { + when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876"); + try (MockedConstruction<DefaultMQProducer> mockedProducers = + mockConstruction(DefaultMQProducer.class, (producer, context) -> { + doNothing().when(producer).start(); + SendResult sendResult = new SendResult(); + sendResult.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT); + when(producer.send(any(Message.class))).thenReturn(sendResult); + doNothing().when(producer).shutdown(); + })) { + SendMessageDTO request = new SendMessageDTO(); + request.setTopic("TopicA"); + request.setBody("hello"); + + assertThatThrownBy(() -> adminClient.sendMessage(request)) + .isInstanceOf(BusinessException.class) + .hasMessageContaining("FLUSH_DISK_TIMEOUT"); + + verify(auditService).record("SEND_MESSAGE", "TopicA", + "Message send did not succeed: FLUSH_DISK_TIMEOUT", "FAILED"); + } + } + + @Test + void sendMessageShouldRejectNullSendResult() throws Exception { + when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876"); + try (MockedConstruction<DefaultMQProducer> mockedProducers = + mockConstruction(DefaultMQProducer.class, (producer, context) -> { + doNothing().when(producer).start(); + when(producer.send(any(Message.class))).thenReturn(null); + doNothing().when(producer).shutdown(); + })) { + SendMessageDTO request = new SendMessageDTO(); + request.setTopic("TopicA"); + request.setBody("hello"); + + assertThatThrownBy(() -> adminClient.sendMessage(request)) + .isInstanceOf(BusinessException.class) + .hasMessageContaining("null"); + + verify(auditService).record("SEND_MESSAGE", "TopicA", + "Message send did not succeed: null", "FAILED"); + } + } } diff --git a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java index d5c8aa51..6028ceb0 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java @@ -256,4 +256,17 @@ class RocketMQMessageProviderTest { assertThat(transaction.getDescription()).contains("tx-group").contains("COMMIT_MESSAGE"); assertThat(record.getConsumerStatus()).isEmpty(); } + + @Test + void getMessageTraceSurfacesAdminFailure() throws Exception { + when(adminExt.queryMessage(anyString(), anyString(), anyInt(), anyLong(), anyLong())) + .thenThrow(new IllegalStateException("broker unavailable")); + + assertThatThrownBy(() -> provider.getMessageTrace("instance-a", "msg-123")) + .isInstanceOf(BusinessException.class) + .hasMessage("Failed to query message trace: broker unavailable") + .satisfies(error -> assertThat(((BusinessException) error).getCode()).isEqualTo(502)); + + verify(queryHistoryService, never()).recordTraceQuery(anyString(), anyString(), any(), anyInt(), anyInt()); + } } diff --git a/server/src/test/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialServiceTest.java b/server/src/test/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialServiceTest.java index 9c796c78..d62c1869 100644 --- a/server/src/test/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialServiceTest.java +++ b/server/src/test/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialServiceTest.java @@ -20,6 +20,7 @@ import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor; import org.apache.rocketmq.studio.common.util.CredentialUtils; import org.apache.rocketmq.studio.common.exception.BusinessException; import org.apache.rocketmq.studio.instance.InstanceRepository; +import org.apache.rocketmq.studio.audit.OperationAuditService; import org.apache.rocketmq.studio.provider.alibaba.AliyunClientFactory; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -33,6 +34,8 @@ import java.util.Optional; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.argThat; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -49,6 +52,9 @@ class CloudCredentialServiceTest { @Mock private AliyunClientFactory aliyunClientFactory; + @Mock + private OperationAuditService operationAuditService; + @InjectMocks private CloudCredentialService service; @@ -124,6 +130,10 @@ class CloudCredentialServiceTest { assertThat(created.getId()).isNotBlank(); assertThat(created.getAccessKey()).isEqualTo("LTAI****0001"); assertThat(created.getSecretKey()).isNull(); + verify(operationAuditService).record(eq("CREATE_CLOUD_CREDENTIAL"), eq("CLOUD_CREDENTIAL"), + eq(created.getId()), eq(null), argThat(detail -> detail.equals("name=ok, vendor=ALIYUN") + && !detail.contains("LTAI5tGoodKey00000000001") && !detail.contains("sk-value")), + eq("SUCCESS"), eq(null)); } @Test @@ -157,6 +167,8 @@ class CloudCredentialServiceTest { service.update(request); verify(aliyunClientFactory).invalidateCredential("cred-1"); + verify(operationAuditService).record(eq("UPDATE_CLOUD_CREDENTIAL"), eq("CLOUD_CREDENTIAL"), + eq("cred-1"), eq(null), eq("name=null, vendor=ALIYUN"), eq("SUCCESS"), eq(null)); } @Test @@ -171,6 +183,8 @@ class CloudCredentialServiceTest { verify(credentialRepository).deleteById("cred-1"); verify(aliyunClientFactory).invalidateCredential("cred-1"); + verify(operationAuditService).record(eq("DELETE_CLOUD_CREDENTIAL"), eq("CLOUD_CREDENTIAL"), + eq("cred-1"), eq(null), eq("name=null, vendor=ALIYUN"), eq("SUCCESS"), eq(null)); } @Test diff --git a/web/src/api/producer.test.ts b/web/src/api/producer.test.ts index e01257e2..6303b075 100644 --- a/web/src/api/producer.test.ts +++ b/web/src/api/producer.test.ts @@ -34,12 +34,15 @@ describe('Producer API', () => { }); it('fetches Studio topic records sorted alphabetically', async () => { - mock.onGet('/topics').reply(200, { + mock.onGet('/topics').reply((config) => { + expect(config.params.instanceId).toBe('instance-1'); + return [200, { code: 200, data: [{ name: 'order-events' }, { name: 'user-signup' }, { name: 'batch-process' }], + }]; }); - const result = await fetchTopicList(); + const result = await fetchTopicList('instance-1'); expect(result).toEqual(['batch-process', 'order-events', 'user-signup']); }); @@ -48,14 +51,14 @@ describe('Producer API', () => { topicList: ['order-events', 'user-signup', 'batch-process'], }); - const result = await fetchTopicList(); + const result = await fetchTopicList('instance-1'); expect(result).toEqual(['batch-process', 'order-events', 'user-signup']); }); it('handles empty topic list', async () => { mock.onGet('/topics').reply(200, { topicList: [] }); - const result = await fetchTopicList(); + const result = await fetchTopicList('instance-1'); expect(result).toEqual([]); }); diff --git a/web/src/api/producer.ts b/web/src/api/producer.ts index e01c0a0c..61138103 100644 --- a/web/src/api/producer.ts +++ b/web/src/api/producer.ts @@ -36,9 +36,9 @@ interface TopicListResponse { // ─── API ──────────────────────────────────────────────────────── -/** Fetch all topic names */ -export async function fetchTopicList(): Promise<string[]> { - const res = await client.get<TopicListResponse>('/topics'); +/** Fetch topic names for a managed instance. */ +export async function fetchTopicList(instanceId: string): Promise<string[]> { + const res = await client.get<TopicListResponse>('/topics', { params: { instanceId } }); const topics = res.data.data?.map((topic) => topic.name) ?? res.data.topicList ?? []; return topics.sort(); } diff --git a/web/src/i18n/translations.ts b/web/src/i18n/translations.ts index 46dcdb3a..66050329 100644 --- a/web/src/i18n/translations.ts +++ b/web/src/i18n/translations.ts @@ -305,6 +305,14 @@ const translations: Record<string, Record<Lang, string>> = { zh: '访问控制规则与用户权限管理,共 {rules} 条规则、{users} 个用户', en: 'Access control rules and user permissions, {rules} rules, {users} users', }, + 'acl.localMetadataNotice': { + zh: '当前 ACL 规则和用户仅保存为 Studio 本地元数据', + en: 'Current ACL rules and users are stored as Studio-local metadata', + }, + 'acl.localMetadataDescription': { + 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.addRule': { zh: '添加规则', en: 'Add Rule' }, 'acl.addUser': { zh: '添加用户', en: 'Add User' }, 'acl.ruleTab': { zh: 'ACL 规则', en: 'ACL Rules' }, diff --git a/web/src/pages/cluster/__tests__/ClientsPage.test.tsx b/web/src/pages/cluster/__tests__/ClientsPage.test.tsx index 3ffae0b4..88c486de 100644 --- a/web/src/pages/cluster/__tests__/ClientsPage.test.tsx +++ b/web/src/pages/cluster/__tests__/ClientsPage.test.tsx @@ -23,6 +23,7 @@ import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from 'vite import type { ClientConnection } from '../../../api/connections'; import { LangProvider } from '../../../i18n/LangContext'; import * as connectionsService from '../../../services/connectionsService'; +import * as instanceService from '../../../services/instanceService'; import ClientsPage from '../clients'; vi.mock('../../../services/connectionsService', () => ({ @@ -221,6 +222,33 @@ describe('Clients page', () => { expect(within(screen.getByTestId('connection-total')).getByText('0')).toBeInTheDocument(); }); + it('surfaces instance discovery failures and allows retrying', async () => { + vi.mocked(instanceService.listInstances) + .mockRejectedValueOnce(new Error('Unable to load managed instances')) + .mockResolvedValueOnce([ + { + id: 'instance-1', + name: 'Instance 1', + endpoint: 'namesrv-1:9876', + type: 'DIRECT', + remark: '', + topicCount: 0, + consumerGroupCount: 0, + createdAt: '', + updatedAt: '', + }, + ]); + const user = userEvent.setup(); + renderWithProviders(<ClientsPage />); + + expect(await screen.findByText('Unable to load managed instances')).toBeInTheDocument(); + expect(connectionsService.listConnections).not.toHaveBeenCalled(); + + await user.click(screen.getByRole('button', { name: /重\s*试/ })); + await screen.findByText('[email protected]:49152'); + expect(connectionsService.listConnections).toHaveBeenCalledWith({ instanceId: 'instance-1' }); + }); + it('opens a client detail dialog from the connection table', async () => { const user = userEvent.setup(); renderWithProviders(<ClientsPage />); diff --git a/web/src/pages/cluster/__tests__/K8sCertsPage.test.tsx b/web/src/pages/cluster/__tests__/K8sCertsPage.test.tsx index 6b4903b5..aa9aa43f 100644 --- a/web/src/pages/cluster/__tests__/K8sCertsPage.test.tsx +++ b/web/src/pages/cluster/__tests__/K8sCertsPage.test.tsx @@ -101,6 +101,14 @@ describe('K8sCertsPage', () => { expect(screen.getByText('-')).toBeInTheDocument(); }); + it('explains that certificate records are Studio-local metadata', async () => { + renderPage(); + + expect(await screen.findByTestId('k8s-cert-local-metadata-notice')).toHaveTextContent( + '当前证书记录仅保存为 Studio 本地元数据', + ); + }); + it.each([ ['platform', 'rocketmq-staging-tls', 'rocketmq-prod-tls'], ['broker.prod.example.com', 'rocketmq-prod-tls', 'rocketmq-staging-tls'], diff --git a/web/src/pages/cluster/certs.tsx b/web/src/pages/cluster/certs.tsx index 400db782..9a971891 100644 --- a/web/src/pages/cluster/certs.tsx +++ b/web/src/pages/cluster/certs.tsx @@ -28,6 +28,7 @@ import { Space, Typography, Card, + Alert, message, } from 'antd'; import type { ColumnsType } from 'antd/es/table'; @@ -356,6 +357,14 @@ const K8sCertsPage = () => { </Button> } /> + <Alert + data-testid="k8s-cert-local-metadata-notice" + type="warning" + showIcon + message="当前证书记录仅保存为 Studio 本地元数据" + description="创建、续期和删除操作尚不会应用到 Kubernetes 集群或 cert-manager。请在集群侧管理实际证书,直到 Kubernetes Provider 接入完成。" + style={{ marginBottom: 16 }} + /> <Flex justify="space-between" style={{ marginBottom: 16 }}> <Space> <Input.Search diff --git a/web/src/pages/cluster/clients.tsx b/web/src/pages/cluster/clients.tsx index f4febdfb..92efe529 100644 --- a/web/src/pages/cluster/clients.tsx +++ b/web/src/pages/cluster/clients.tsx @@ -114,18 +114,34 @@ const ClientsPage = () => { const [clusterFilter, setClusterFilter] = useState<string>('ALL'); const [selectedConnection, setSelectedConnection] = useState<ClientConnection | null>(null); const [loadError, setLoadError] = useState<string | null>(null); + const [instanceLoadKey, setInstanceLoadKey] = useState(0); useEffect(() => { let cancelled = false; - void listInstances().then((nextInstances) => { - if (cancelled) return; - setInstances(nextInstances); - setSelectedInstanceId((current) => current || nextInstances[0]?.id || ''); - }); + + setLoading(true); + void listInstances() + .then((nextInstances) => { + if (cancelled) return; + setInstances(nextInstances); + setSelectedInstanceId((current) => current || nextInstances[0]?.id || ''); + setLoadError(null); + }) + .catch((error) => { + if (cancelled) return; + setInstances([]); + setSelectedInstanceId(''); + setConnections([]); + setLoadError(getLoadErrorMessage(error)); + }) + .finally(() => { + if (!cancelled) setLoading(false); + }); + return () => { cancelled = true; }; - }, []); + }, [instanceLoadKey]); useEffect(() => { let cancelled = false; @@ -355,7 +371,17 @@ const ClientsPage = () => { /> {loadError && ( - <Alert showIcon type="warning" message={loadError} style={{ marginBottom: 16 }} /> + <Alert + showIcon + type="warning" + message={loadError} + style={{ marginBottom: 16 }} + action={ + <Button size="small" onClick={() => setInstanceLoadKey((key) => key + 1)}> + 重试 + </Button> + } + /> )} {connections.some((connection) => connection.partial) && ( <Alert diff --git a/web/src/pages/home/__tests__/DashboardPage.test.tsx b/web/src/pages/home/__tests__/DashboardPage.test.tsx new file mode 100644 index 00000000..4d28874c --- /dev/null +++ b/web/src/pages/home/__tests__/DashboardPage.test.tsx @@ -0,0 +1,118 @@ +/* + * 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. + */ + +import { App } from 'antd'; +import { render, screen, waitFor } from '@testing-library/react'; +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 { DashboardData } from '../../../api/metrics'; +import { LangProvider } from '../../../i18n/LangContext'; +import * as dashboardService from '../../../services/dashboardService'; +import * as instanceService from '../../../services/instanceService'; +import DashboardPage from '../dashboard'; + +vi.mock('../../../services/dashboardService', () => ({ getDashboard: vi.fn() })); +vi.mock('../../../services/instanceService', () => ({ listInstances: vi.fn() })); + +const dashboard = (name: string): DashboardData => ({ + stats: { + totalClusters: 1, + healthyClusters: 1, + totalBrokers: 1, + totalProxies: 0, + totalNameServers: 1, + totalTopics: 1, + totalConsumerGroups: 1, + totalMessagesToday: 1, + messagesPerSecond: 1, + tpsIn: 1, + tpsOut: 1, + }, + clusters: [ + { + id: name, + name, + type: 'V4_NAMESRV', + status: 'healthy', + brokers: 1, + proxies: 0, + topics: 1, + groups: 1, + tpsIn: 1, + tpsOut: 1, + version: '5.0.0', + throughput: [1], + }, + ], +}); + +const deferred = <T,>() => { + let resolve!: (value: T) => void; + const promise = new Promise<T>((promiseResolve) => { + resolve = promiseResolve; + }); + return { promise, resolve }; +}; + +const renderWithProviders = (ui: React.ReactElement) => + render( + <App> + <LangProvider> + <MemoryRouter>{ui}</MemoryRouter> + </LangProvider> + </App>, + ); + +beforeAll(() => { + Object.defineProperty(window, 'matchMedia', { + writable: true, + value: vi.fn().mockImplementation(() => ({ + matches: false, + addEventListener: vi.fn(), + removeEventListener: vi.fn(), + })), + }); +}); + +beforeEach(() => { + vi.clearAllMocks(); + vi.mocked(instanceService.listInstances).mockResolvedValue([ + { id: 'instance-a', name: 'Instance A', endpoint: 'a:9876', type: 'DIRECT', remark: '', topicCount: 0, consumerGroupCount: 0, createdAt: '', updatedAt: '' }, + { id: 'instance-b', name: 'Instance B', endpoint: 'b:9876', type: 'DIRECT', remark: '', topicCount: 0, consumerGroupCount: 0, createdAt: '', updatedAt: '' }, + ]); +}); + +describe('DashboardPage', () => { + it('does not let a stale instance response overwrite the latest selection', async () => { + const instanceA = deferred<DashboardData>(); + const instanceB = deferred<DashboardData>(); + vi.mocked(dashboardService.getDashboard) + .mockResolvedValueOnce(dashboard('initial-cluster')) + .mockReturnValueOnce(instanceA.promise) + .mockReturnValueOnce(instanceB.promise); + const user = userEvent.setup(); + renderWithProviders(<DashboardPage />); + + await screen.findByText('initial-cluster'); + const selector = screen.getByRole('combobox', { name: 'Dashboard instance' }); + await user.click(selector); + await user.click(await screen.findByText('Instance A', { selector: '.ant-select-item-option-content' })); + await user.click(selector); + await user.click(await screen.findByText('Instance B', { selector: '.ant-select-item-option-content' })); + + instanceB.resolve(dashboard('instance-b-cluster')); + expect(await screen.findByText('instance-b-cluster')).toBeInTheDocument(); + + instanceA.resolve(dashboard('instance-a-cluster')); + await waitFor(() => { + expect(screen.queryByText('instance-a-cluster')).not.toBeInTheDocument(); + }); + expect(screen.getByText('instance-b-cluster')).toBeInTheDocument(); + }); +}); diff --git a/web/src/pages/home/dashboard.tsx b/web/src/pages/home/dashboard.tsx index 8eced776..f6ec900c 100644 --- a/web/src/pages/home/dashboard.tsx +++ b/web/src/pages/home/dashboard.tsx @@ -1,4 +1,4 @@ -import { useCallback, useEffect, useState } from 'react'; +import { useCallback, useEffect, useRef, useState } from 'react'; import { useNavigate } from 'react-router-dom'; import { Alert, @@ -37,16 +37,25 @@ const DashboardPage = () => { const [selectedInstanceId, setSelectedInstanceId] = useState<string>(); const [loading, setLoading] = useState(true); const [loadError, setLoadError] = useState(false); + const dashboardRequestIdRef = useRef(0); const loadDashboard = useCallback(async () => { + const requestId = ++dashboardRequestIdRef.current; setLoading(true); setLoadError(false); try { - setDashboard(await getDashboard(selectedInstanceId)); + const nextDashboard = await getDashboard(selectedInstanceId); + if (requestId === dashboardRequestIdRef.current) { + setDashboard(nextDashboard); + } } catch { - setLoadError(true); + if (requestId === dashboardRequestIdRef.current) { + setLoadError(true); + } } finally { - setLoading(false); + if (requestId === dashboardRequestIdRef.current) { + setLoading(false); + } } }, [selectedInstanceId]); diff --git a/web/src/pages/instance/__tests__/AclPage.test.tsx b/web/src/pages/instance/__tests__/AclPage.test.tsx index a67a9725..379535de 100644 --- a/web/src/pages/instance/__tests__/AclPage.test.tsx +++ b/web/src/pages/instance/__tests__/AclPage.test.tsx @@ -121,6 +121,14 @@ describe('ACL page', () => { expect(screen.getByText('cluster-a')).toBeInTheDocument(); }); + it('explains that ACL records are Studio-local metadata', async () => { + renderWithProviders(<AclPage />); + + expect(await screen.findByTestId('acl-local-metadata-notice')).toHaveTextContent( + '当前 ACL 规则和用户仅保存为 Studio 本地元数据', + ); + }); + it('renders backend users on the user tab', async () => { const user = userEvent.setup(); renderWithProviders(<AclPage />); diff --git a/web/src/pages/instance/__tests__/DLQPage.test.tsx b/web/src/pages/instance/__tests__/DLQPage.test.tsx index fbd4f7e6..b55778af 100644 --- a/web/src/pages/instance/__tests__/DLQPage.test.tsx +++ b/web/src/pages/instance/__tests__/DLQPage.test.tsx @@ -43,6 +43,17 @@ vi.mock('../../../services/instanceService', () => ({ createdAt: '', updatedAt: '', }, + { + id: 'instance-2', + name: 'Instance 2', + endpoint: 'namesrv-2:9876', + type: 'DIRECT', + remark: '', + topicCount: 0, + consumerGroupCount: 0, + createdAt: '', + updatedAt: '', + }, ]), })); @@ -253,4 +264,39 @@ describe('DLQ page', () => { expect(await screen.findByText('DLQ provider is not configured')).toBeInTheDocument(); }); + + it('clears retry state before loading groups for a newly selected instance', async () => { + let resolveSecondInstance!: (groups: DLQGroup[]) => void; + vi.mocked(messageService.listDLQGroups) + .mockResolvedValueOnce([dlqGroup]) + .mockImplementationOnce( + () => + new Promise<DLQGroup[]>((resolve) => { + resolveSecondInstance = resolve; + }), + ); + const user = userEvent.setup(); + renderWithProviders(<DLQPage />); + + const orderRow = (await screen.findByText('cg-order')).closest('tr'); + if (!orderRow) throw new Error('DLQ group row not found'); + + await user.click(within(orderRow).getByRole('button', { name: '重投消息' })); + expect(await screen.findByText('重投死信消息')).toBeInTheDocument(); + + await user.click(screen.getAllByRole('combobox')[0]); + await user.click( + await screen.findByText('Instance 2', { selector: '.ant-select-item-option-content' }), + ); + + await waitFor(() => { + expect(messageService.listDLQGroups).toHaveBeenLastCalledWith('instance-2'); + }); + await waitFor(() => { + expect(screen.queryByRole('row', { name: /cg-order/ })).not.toBeInTheDocument(); + }); + + resolveSecondInstance([secondDlqGroup]); + expect(await screen.findByText('-cg-"payment"')).toBeInTheDocument(); + }); }); diff --git a/web/src/pages/instance/__tests__/MessagePage.test.tsx b/web/src/pages/instance/__tests__/MessagePage.test.tsx index 1ae39b55..36dd97a3 100644 --- a/web/src/pages/instance/__tests__/MessagePage.test.tsx +++ b/web/src/pages/instance/__tests__/MessagePage.test.tsx @@ -79,6 +79,8 @@ const renderWithProviders = (ui: React.ReactElement) => </App>, ); +const lastElement = <T,>(elements: T[]): T => elements[elements.length - 1]!; + describe('Message page query history', () => { beforeEach(() => { localStorage.clear(); @@ -96,11 +98,14 @@ describe('Message page query history', () => { expect(screen.getByRole('button', { name: /最近查询/ })).toBeDisabled(); await user.click(screen.getByText('按 Message ID')); + await user.click(lastElement(screen.getAllByRole('combobox'))); + await user.click(lastElement(await screen.findAllByText('order-create'))); await user.type(screen.getByPlaceholderText('输入 Message ID'), 'MID-001'); await user.click(screen.getByRole('button', { name: /^search查询$/ })); await waitFor(() => { expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({ + topic: 'order-create', msgId: 'MID-001', instanceId: '', }); @@ -112,11 +117,12 @@ describe('Message page query history', () => { renderWithProviders(<MessagePage />); await user.click(screen.getByRole('button', { name: /最近查询/ })); - await user.click(await screen.findByText('Message ID: MID-001')); + await user.click(await screen.findByText('Message ID: MID-001 · Topic: order-create')); expect(screen.getByPlaceholderText('输入 Message ID')).toHaveValue('MID-001'); await waitFor(() => { expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({ + topic: 'order-create', msgId: 'MID-001', instanceId: '', }); @@ -134,11 +140,14 @@ describe('Message page query history', () => { renderWithProviders(<MessagePage />); await user.click(screen.getByText('按 Message ID')); + await user.click(lastElement(screen.getAllByRole('combobox'))); + await user.click(lastElement(await screen.findAllByText('order-create'))); await user.type(screen.getByPlaceholderText('输入 Message ID'), 'MID-FAILED'); await user.click(screen.getByRole('button', { name: /^search查询$/ })); await waitFor(() => { expect(messageServiceMocks.queryMessages).toHaveBeenCalledWith({ + topic: 'order-create', msgId: 'MID-FAILED', instanceId: '', }); @@ -151,6 +160,8 @@ describe('Message page query history', () => { const user = userEvent.setup(); renderWithProviders(<MessagePage />); await user.click(screen.getByText('按 Message ID')); + await user.click(lastElement(screen.getAllByRole('combobox'))); + await user.click(lastElement(await screen.findAllByText('order-create'))); const messageIdInput = screen.getByPlaceholderText('输入 Message ID'); const queryButton = screen.getByRole('button', { name: /^search查询$/ }); @@ -240,15 +251,16 @@ describe('Message page query history', () => { JSON.stringify([ { mode: 'topic', params: { topic: ['invalid'] } }, { mode: 'unknown', params: { topic: 'order-create' } }, - { mode: 'msgid', params: { msgId: 'MID-VALID' } }, + { mode: 'msgid', params: { msgId: 'MID-VALID', topic: 'order-create' } }, ]), ); renderWithProviders(<MessagePage />); await user.click(screen.getByRole('button', { name: /最近查询/ })); - expect(await screen.findByText('Message ID: MID-VALID')).toBeInTheDocument(); + expect( + await screen.findByText('Message ID: MID-VALID · Topic: order-create'), + ).toBeInTheDocument(); expect(screen.queryByText(/invalid/)).not.toBeInTheDocument(); - expect(screen.queryByText(/order-create/)).not.toBeInTheDocument(); }); it('does not report consume verification success without a backend API', async () => { @@ -257,6 +269,8 @@ describe('Message page query history', () => { renderWithProviders(<MessagePage />); await user.click(screen.getByText('按 Message ID')); + await user.click(lastElement(screen.getAllByRole('combobox'))); + await user.click(lastElement(await screen.findAllByText('order-create'))); await user.type(screen.getByPlaceholderText('输入 Message ID'), 'MID-CONSUME-VERIFY-001'); await user.click(screen.getByRole('button', { name: /^search查询$/ })); diff --git a/web/src/pages/instance/acl.tsx b/web/src/pages/instance/acl.tsx index e6e3e469..374320c0 100644 --- a/web/src/pages/instance/acl.tsx +++ b/web/src/pages/instance/acl.tsx @@ -33,6 +33,7 @@ import { Badge, Typography, Flex, + Alert, message, } from 'antd'; import { Plus, MagnifyingGlass, ShieldCheck, User, Eye, EyeSlash } from '@phosphor-icons/react'; @@ -655,6 +656,15 @@ const AclPage = () => { } /> + <Alert + data-testid="acl-local-metadata-notice" + type="warning" + showIcon + message={t('acl.localMetadataNotice')} + description={t('acl.localMetadataDescription')} + style={{ marginBottom: 16 }} + /> + <Card bordered={false} bodyStyle={{ padding: 0 }}> <Tabs activeKey={activeTab} diff --git a/web/src/pages/instance/dlq.tsx b/web/src/pages/instance/dlq.tsx index 7f3a5873..dfd574fc 100644 --- a/web/src/pages/instance/dlq.tsx +++ b/web/src/pages/instance/dlq.tsx @@ -125,12 +125,22 @@ const DLQPage = () => { useEffect(() => { let cancelled = false; + // The retry dialog owns a group name that is meaningful only for the + // currently selected instance. Clear all instance-scoped state before + // starting the next request so an old group cannot be retried on a new + // instance while that request is in flight. + setGroups([]); + setSelectedGroupNames([]); + setDetailGroup(null); + setRetryModalOpen(false); + setRetryGroup(null); + setRetryTargetTopic(''); + setRetryError(null); + setLoadError(null); + if (!selectedInstanceId) { void Promise.resolve().then(() => { if (cancelled) return; - setGroups([]); - setSelectedGroupNames([]); - setLoadError(null); setLoading(false); }); return () => { @@ -138,6 +148,7 @@ const DLQPage = () => { }; } + setLoading(true); void listDLQGroups(selectedInstanceId) .then((nextGroups) => { if (!cancelled) { diff --git a/web/src/pages/instance/message.tsx b/web/src/pages/instance/message.tsx index 96bf564c..c4a7d43a 100644 --- a/web/src/pages/instance/message.tsx +++ b/web/src/pages/instance/message.tsx @@ -161,21 +161,23 @@ const isMessageQuery = (value: unknown): value is MessageQuery => { ); }; +const isRecentQuery = (value: unknown): value is RecentQuery => { + if (typeof value !== 'object' || value === null) return false; + const query = value as RecentQuery; + return ( + isQueryMode(query.mode) && + isMessageQuery(query.params) && + (query.mode !== 'msgid' || Boolean(query.params.topic?.trim())) + ); +}; + const loadRecentQueries = (): RecentQuery[] => { try { const stored = localStorage.getItem(QUERY_HISTORY_STORAGE_KEY); if (!stored) return []; const parsed: unknown = JSON.parse(stored); if (!Array.isArray(parsed)) return []; - return parsed - .filter( - (item): item is RecentQuery => - typeof item === 'object' && - item !== null && - isQueryMode((item as RecentQuery).mode) && - isMessageQuery((item as RecentQuery).params), - ) - .slice(0, MAX_QUERY_HISTORY); + return parsed.filter(isRecentQuery).slice(0, MAX_QUERY_HISTORY); } catch { return []; } @@ -184,7 +186,8 @@ const loadRecentQueries = (): RecentQuery[] => { const querySignature = (query: RecentQuery): string => JSON.stringify(query); const queryLabel = ({ mode, params }: RecentQuery): string => { - if (mode === 'msgid') return `Message ID: ${params.msgId || '全部'}`; + if (mode === 'msgid') + return `Message ID: ${params.msgId || '全部'} · Topic: ${params.topic || '全部'}`; if (mode === 'key') { return `Key: ${params.key || '全部'}${params.topic ? ` · Topic: ${params.topic}` : ''}`; } @@ -318,7 +321,7 @@ const MessagePage = () => { } : queryMode === 'key' ? { topic: selectedTopic, key: keyInput || undefined } - : { msgId: msgIdInput || undefined }; + : { topic: selectedTopic, msgId: msgIdInput || undefined }; await executeQuery(queryMode, params); }; @@ -749,12 +752,26 @@ const MessagePage = () => { )} {queryMode === 'msgid' && ( - <Input - placeholder="输入 Message ID" - style={{ width: 400 }} - value={msgIdInput} - onChange={(e) => setMsgIdInput(e.target.value)} - /> + <> + <Select + placeholder="选择 Topic" + style={{ width: 360 }} + value={selectedTopic} + onChange={setSelectedTopic} + allowClear + showSearch + options={topicOptions.map((t) => ({ + value: t, + label: t, + }))} + /> + <Input + placeholder="输入 Message ID" + style={{ width: 400 }} + value={msgIdInput} + onChange={(e) => setMsgIdInput(e.target.value)} + /> + </> )} <Button diff --git a/web/src/pages/studio/Producer.tsx b/web/src/pages/studio/Producer.tsx index e0cfc486..132ccd37 100644 --- a/web/src/pages/studio/Producer.tsx +++ b/web/src/pages/studio/Producer.tsx @@ -60,10 +60,26 @@ const ProducerPage = () => { }; }, []); + const handleInstanceChange = (instanceId: string) => { + setSelectedInstanceId(instanceId); + setTopicList([]); + setProducerGroups([]); + setConnectionList([]); + form.setFieldsValue({ selectedTopic: undefined, producerGroup: undefined }); + }; + useEffect(() => { let cancelled = false; - void fetchTopicList() + if (!selectedInstanceId) { + return () => { + cancelled = true; + }; + } + + form.setFieldsValue({ selectedTopic: undefined, producerGroup: undefined }); + + void fetchTopicList(selectedInstanceId) .then((topics) => { if (!cancelled) { setTopicList(topics); @@ -78,7 +94,7 @@ const ProducerPage = () => { return () => { cancelled = true; }; - }, [fetchTopicFailedMessage, message]); + }, [fetchTopicFailedMessage, form, message, selectedInstanceId]); useEffect(() => { let cancelled = false; @@ -170,7 +186,7 @@ const ProducerPage = () => { <Select aria-label="Instance" value={selectedInstanceId || undefined} - onChange={setSelectedInstanceId} + onChange={handleInstanceChange} placeholder="Select instance" style={{ width: 220 }} options={instances.map((instance) => ({ value: instance.id, label: instance.name }))} diff --git a/web/src/pages/studio/__tests__/Producer.test.tsx b/web/src/pages/studio/__tests__/Producer.test.tsx index 8d43bee5..5bc65599 100644 --- a/web/src/pages/studio/__tests__/Producer.test.tsx +++ b/web/src/pages/studio/__tests__/Producer.test.tsx @@ -16,7 +16,7 @@ */ import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'; -import { fireEvent, render, screen, waitFor } from '@testing-library/react'; +import { fireEvent, render, screen, waitFor, within } from '@testing-library/react'; import userEvent from '@testing-library/user-event'; import { App } from 'antd'; import { LangProvider } from '../../../i18n/LangContext'; @@ -87,7 +87,7 @@ describe('ProducerPage', () => { renderWithProviders(<ProducerPage />); await waitFor(() => { - expect(fetchTopicList).toHaveBeenCalledTimes(1); + expect(fetchTopicList).toHaveBeenCalledWith('instance-1'); }); }); @@ -186,4 +186,64 @@ describe('ProducerPage', () => { ); }); }); + + it('clears stale producer query state before loading a new instance', async () => { + vi.mocked(listInstances).mockResolvedValue([ + { + id: 'instance-1', + name: 'Primary instance', + remark: '', + type: 'DIRECT', + endpoint: '127.0.0.1:9876', + topicCount: 0, + consumerGroupCount: 0, + createdAt: '2026-08-01T00:00:00', + updatedAt: '2026-08-01T00:00:00', + }, + { + id: 'instance-2', + name: 'Secondary instance', + remark: '', + type: 'DIRECT', + endpoint: '127.0.0.2:9876', + topicCount: 0, + consumerGroupCount: 0, + createdAt: '2026-08-01T00:00:00', + updatedAt: '2026-08-01T00:00:00', + }, + ]); + vi.mocked(fetchTopicList).mockResolvedValueOnce(['order-events']).mockResolvedValueOnce([]); + vi.mocked(queryProducerConnection).mockResolvedValue([ + { + clientId: 'producer-1', + clientAddr: '192.168.1.10', + language: 'JAVA', + versionDesc: '5.1.0', + }, + ]); + const user = userEvent.setup(); + const { container } = renderWithProviders(<ProducerPage />); + + await waitFor(() => expect(fetchTopicList).toHaveBeenCalledWith('instance-1')); + const [instanceSelect, topicSelect, groupInput] = screen.getAllByRole('combobox'); + fireEvent.mouseDown(topicSelect.parentElement!); + await user.click( + await screen.findByText('order-events', { selector: '.ant-select-item-option-content' }), + ); + await user.type(groupInput, 'order-producer'); + await user.click(screen.getByRole('button', { name: /搜索/ })); + expect(await screen.findByText('producer-1')).toBeInTheDocument(); + + fireEvent.mouseDown(instanceSelect.parentElement!); + await user.click( + await screen.findByText('Secondary instance', { + selector: '.ant-select-item-option-content', + }), + ); + + await waitFor(() => expect(fetchTopicList).toHaveBeenLastCalledWith('instance-2')); + // scope to the page container: antd keeps closed dropdown portals in document.body + expect(within(container).queryByText('producer-1')).not.toBeInTheDocument(); + expect(within(container).queryByText('order-events')).not.toBeInTheDocument(); + }); });
