This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new 5dbbdb14d fix(provider): sanitize malformed external data in apache
runtime and tencent catalog/ACL (#2409)
5dbbdb14d is described below
commit 5dbbdb14d108729eb69b9010ee56b4bbd5fa89a2
Author: yyqdbngt <[email protected]>
AuthorDate: Sat Aug 22 14:46:53 2026 +0800
fix(provider): sanitize malformed external data in apache runtime and
tencent catalog/ACL (#2409)
Co-authored-by: Yue Wang <[email protected]>
---
.../provider/apache/RocketMQClusterProvider.java | 31 +++++++----
.../studio/provider/tencent/TencentAclService.java | 58 ++++++++++----------
.../provider/tencent/TencentCatalogService.java | 20 +++++--
.../apache/RocketMQClusterProviderTest.java | 61 ++++++++++++++++++++++
.../provider/tencent/TencentAclServiceTest.java | 26 ++++++++-
.../tencent/TencentCatalogServiceTest.java | 38 ++++++++++++++
6 files changed, 188 insertions(+), 46 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProvider.java
index fd6771135..8cd304741 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProvider.java
@@ -209,12 +209,16 @@ public class RocketMQClusterProvider implements
ClusterProvider {
// Use master address (brokerId = 0) preferentially
String masterAddr = brokerData.getBrokerAddrs().get(0L);
- if (masterAddr == null && !brokerData.getBrokerAddrs().isEmpty()) {
- masterAddr =
brokerData.getBrokerAddrs().values().iterator().next();
+ if (!StringUtils.hasText(masterAddr)) {
+ masterAddr = brokerData.getBrokerAddrs().values().stream()
+ .filter(StringUtils::hasText)
+ .findFirst()
+ .orElse(null);
}
- if (masterAddr == null) {
+ if (!StringUtils.hasText(masterAddr)) {
continue;
}
+ masterAddr = masterAddr.trim();
BrokerVO.BrokerVOBuilder builder = BrokerVO.builder()
.name(brokerName)
@@ -270,7 +274,10 @@ public class RocketMQClusterProvider implements
ClusterProvider {
String diskRatio = table.get("commitLogDiskRatio");
if (diskRatio != null && !diskRatio.isEmpty()) {
try {
- builder.diskUsage(Double.parseDouble(diskRatio));
+ double parsedRatio = Double.parseDouble(diskRatio);
+ if (Double.isFinite(parsedRatio) && parsedRatio >= 0) {
+ builder.diskUsage(parsedRatio);
+ }
} catch (NumberFormatException ignored) {
// keep default
}
@@ -283,7 +290,8 @@ public class RocketMQClusterProvider implements
ClusterProvider {
}
private boolean hasUnavailableRuntimeStats(List<BrokerVO> brokers) {
- return brokers != null && brokers.stream().anyMatch(broker ->
!broker.isRuntimeStatsAvailable());
+ return brokers != null && brokers.stream()
+ .anyMatch(broker -> broker == null ||
!broker.isRuntimeStatsAvailable());
}
/**
@@ -293,12 +301,12 @@ public class RocketMQClusterProvider implements
ClusterProvider {
private long parseTpsValue(String tpsStr) {
try {
String[] parts = tpsStr.trim().split("\\s+");
- if (parts.length >= 2) {
- return (long) Double.parseDouble(parts[1]);
- }
- if (parts.length == 1) {
- return (long) Double.parseDouble(parts[0]);
+ String selected = parts.length >= 2 ? parts[1] : parts[0];
+ double value = Double.parseDouble(selected);
+ if (!Double.isFinite(value) || value <= 0) {
+ return 0;
}
+ return value >= Long.MAX_VALUE ? Long.MAX_VALUE : (long) value;
} catch (NumberFormatException ignored) {
// fall through
}
@@ -329,6 +337,9 @@ public class RocketMQClusterProvider implements
ClusterProvider {
}
Set<String> proxyIps = new TreeSet<>();
for (Connection conn : connection.getConnectionSet()) {
+ if (conn == null) {
+ continue;
+ }
String clientAddr = conn.getClientAddr();
if (clientAddr == null || clientAddr.isBlank()) {
continue;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentAclService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentAclService.java
index 2c9985473..f57a1e930 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentAclService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentAclService.java
@@ -79,7 +79,7 @@ public class TencentAclService {
break;
}
for (RoleItem role : data) {
- if (role != null) {
+ if (role != null && StringUtils.hasText(role.getRoleName())) {
users.add(toUser(role, context.cloudInstanceId()));
}
}
@@ -92,6 +92,7 @@ public class TencentAclService {
public List<AclRuleVO> listRules(String instanceId, String principal) {
Context context = resolve(instanceId);
+ String requestedPrincipal = StringUtils.hasText(principal) ?
principal.trim() : null;
List<AclRuleVO> rules = new ArrayList<>();
for (int page = 0; page < MAX_PAGES; page++) {
DescribeRoleListRequest request = new DescribeRoleListRequest();
@@ -105,11 +106,11 @@ public class TencentAclService {
break;
}
for (RoleItem role : data) {
- if (role == null) {
+ if (role == null || !StringUtils.hasText(role.getRoleName())) {
continue;
}
- if (StringUtils.hasText(principal)
- && !principal.equals(role.getRoleName())) {
+ if (requestedPrincipal != null
+ && !requestedPrincipal.equals(role.getRoleName())) {
continue;
}
rules.add(toRule(role));
@@ -123,21 +124,19 @@ public class TencentAclService {
public AclUserVO createUser(String instanceId, AclUserVO user) {
Context context = resolve(instanceId);
- if (!StringUtils.hasText(user.getUsername())) {
- throw new BusinessException(400, "ACL username is required");
- }
+ String roleName = requireRoleName(user == null ? null :
user.getUsername(), "ACL username");
CreateRoleRequest request = new CreateRoleRequest();
request.setInstanceId(context.cloudInstanceId());
- request.setRole(user.getUsername());
+ request.setRole(roleName);
request.setPermRead(user.getPermRead() == null || user.getPermRead());
request.setPermWrite(user.getPermWrite() == null ||
user.getPermWrite());
- request.setRemark(user.getUsername());
+ request.setRemark(roleName);
clientFactory.call(context.credentialId(), context.regionId(),
client -> client.CreateRole(request));
// The created role is not returned by the API; reconstruct from the
known inputs.
// Tencent roles have no database row; the role name lives in
username, id stays null.
return AclUserVO.builder()
- .username(user.getUsername())
+ .username(roleName)
.accessKey(null)
.secretKey(null)
.admin(false)
@@ -152,10 +151,7 @@ public class TencentAclService {
Context context = resolve(instanceId);
// Tencent roles have no database row; username is the stable
role-name source (id is the
// numeric ACL-user primary key and stays null for Tencent roles).
- String roleName = user.getUsername();
- if (!StringUtils.hasText(roleName)) {
- throw new BusinessException(400, "ACL username is required");
- }
+ String roleName = requireRoleName(user == null ? null :
user.getUsername(), "ACL username");
RoleItem existing = user.getPermRead() == null || user.getPermWrite()
== null
? findRole(context, roleName) : null;
boolean permRead = user.getPermRead() == null
@@ -207,12 +203,10 @@ public class TencentAclService {
public void deleteUser(String instanceId, String username) {
Context context = resolve(instanceId);
- if (!StringUtils.hasText(username)) {
- throw new BusinessException(400, "ACL username is required");
- }
+ String roleName = requireRoleName(username, "ACL username");
DeleteRoleRequest request = new DeleteRoleRequest();
request.setInstanceId(context.cloudInstanceId());
- request.setRole(username);
+ request.setRole(roleName);
clientFactory.call(context.credentialId(), context.regionId(),
client -> client.DeleteRole(request));
}
@@ -224,17 +218,17 @@ public class TencentAclService {
*/
public AclRuleVO createRule(String instanceId, AclRuleVO rule) {
Context context = resolve(instanceId);
- requireRulePrincipal(rule);
+ String principal = requireRulePrincipal(rule);
boolean permRead = hasAction(rule, "SUB");
boolean permWrite = hasAction(rule, "PUB");
ModifyRoleRequest request = new ModifyRoleRequest();
request.setInstanceId(context.cloudInstanceId());
- request.setRole(rule.getPrincipal());
+ request.setRole(principal);
request.setPermRead(permRead);
request.setPermWrite(permWrite);
clientFactory.call(context.credentialId(), context.regionId(),
client -> client.ModifyRole(request));
- return toRule(rule.getPrincipal(), permRead, permWrite);
+ return toRule(principal, permRead, permWrite);
}
public AclRuleVO updateRule(String instanceId, AclRuleVO rule) {
@@ -243,20 +237,23 @@ public class TencentAclService {
public void deleteRule(String instanceId, String principal) {
Context context = resolve(instanceId);
- if (!StringUtils.hasText(principal)) {
- throw new BusinessException(400, "ACL principal is required");
- }
+ String roleName = requireRoleName(principal, "ACL principal");
DeleteRoleRequest request = new DeleteRoleRequest();
request.setInstanceId(context.cloudInstanceId());
- request.setRole(principal);
+ request.setRole(roleName);
clientFactory.call(context.credentialId(), context.regionId(),
client -> client.DeleteRole(request));
}
- private static void requireRulePrincipal(AclRuleVO rule) {
- if (rule == null || !StringUtils.hasText(rule.getPrincipal())) {
- throw new BusinessException(400, "ACL principal is required");
+ private static String requireRulePrincipal(AclRuleVO rule) {
+ return requireRoleName(rule == null ? null : rule.getPrincipal(), "ACL
principal");
+ }
+
+ private static String requireRoleName(String value, String label) {
+ if (!StringUtils.hasText(value)) {
+ throw new BusinessException(400, label + " is required");
}
+ return value.trim();
}
private static boolean hasAction(AclRuleVO rule, String action) {
@@ -289,10 +286,11 @@ public class TencentAclService {
* available from DescribeRoleList, so re-fetch and match by role name.
*/
public AclUserVO getUserCredentials(String instanceId, String username) {
+ String roleName = requireRoleName(username, "ACL username");
return listUsers(instanceId).stream()
- .filter(user -> user.getUsername().equals(username))
+ .filter(user -> roleName.equals(user.getUsername()))
.findFirst()
- .orElseThrow(() -> new BusinessException(404, "ACL user not
found: " + username));
+ .orElseThrow(() -> new BusinessException(404, "ACL user not
found: " + roleName));
}
private Context resolve(String instanceId) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentCatalogService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentCatalogService.java
index fd84f571f..886b907bc 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentCatalogService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentCatalogService.java
@@ -118,7 +118,7 @@ public class TencentCatalogService implements
CloudCatalogProvider {
request.setInstanceId(cloudInstanceId);
DescribeInstanceResponse response = clientFactory.call(credentialId,
regionId,
client -> client.DescribeInstance(request));
- if (response == null || response.getInstanceId() == null) {
+ if (response == null || response.getInstanceId() == null ||
response.getInstanceId().isBlank()) {
throw new BusinessException(404, "Tencent Cloud RocketMQ 5.x
instance not found: " + cloudInstanceId);
}
CloudInstanceDetailVO detail = new CloudInstanceDetailVO();
@@ -149,9 +149,10 @@ public class TencentCatalogService implements
CloudCatalogProvider {
return result;
}
for (Endpoint endpoint : endpoints) {
- if (endpoint != null &&
"OPEN".equalsIgnoreCase(endpoint.getStatus())) {
+ if (endpoint != null &&
"OPEN".equalsIgnoreCase(endpoint.getStatus())
+ && endpoint.getEndpointUrl() != null &&
!endpoint.getEndpointUrl().isBlank()) {
result.add(new CloudInstanceDetailVO.CloudEndpoint(
- endpointType(endpoint.getType()),
endpoint.getEndpointUrl()));
+ endpointType(endpoint.getType()),
endpoint.getEndpointUrl().trim()));
}
}
return result;
@@ -168,14 +169,23 @@ public class TencentCatalogService implements
CloudCatalogProvider {
}
private static Integer toInteger(Long value) {
- return value == null ? null : Math.toIntExact(value);
+ if (value == null) {
+ return null;
+ }
+ if (value > Integer.MAX_VALUE) {
+ return Integer.MAX_VALUE;
+ }
+ if (value < 0) {
+ return 0;
+ }
+ return value.intValue();
}
private static boolean matchesSearch(String search, CloudInstanceOptionVO
option) {
if (search == null || search.isBlank()) {
return true;
}
- String needle = search.toLowerCase(Locale.ROOT);
+ String needle = search.trim().toLowerCase(Locale.ROOT);
return containsIgnoreCase(option.getInstanceId(), needle)
|| containsIgnoreCase(option.getInstanceName(), needle);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
index 1c3e12c14..5177df3e9 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
@@ -36,6 +36,7 @@ import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
import org.junit.jupiter.api.Test;
import java.util.HashMap;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Optional;
import java.util.Set;
@@ -155,6 +156,44 @@ class RocketMQClusterProviderTest {
});
}
+ @Test
+ void discoverClustersShouldRejectNonFiniteRuntimeMetrics() throws
Exception {
+ DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+ RocketMQClusterProvider provider = newProvider(adminExt);
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
+ KVTable runtime = runtimeStats("1 NaN 3", "1 Infinity 3");
+ runtime.getTable().put("commitLogDiskRatio", "NaN");
+
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtime);
+
+ ClusterVO cluster = provider.discoverClusters().get(0);
+
+ assertThat(cluster.getBrokers()).singleElement().satisfies(broker -> {
+ assertThat(broker.getTpsIn()).isZero();
+ assertThat(broker.getTpsOut()).isZero();
+ assertThat(broker.getDiskUsage()).isZero();
+ });
+ }
+
+ @Test
+ void discoverClustersShouldUseFirstNonBlankBrokerAddress() throws
Exception {
+ DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+ RocketMQClusterProvider provider = newProvider(adminExt);
+ ClusterInfo info = clusterInfo();
+ LinkedHashMap<Long, String> addrs = new LinkedHashMap<>();
+ addrs.put(0L, " ");
+ addrs.put(1L, null);
+ addrs.put(2L, " 10.0.0.12:10911 ");
+ info.getBrokerAddrTable().put(
+ "broker-a", new BrokerData("DefaultCluster", "broker-a",
addrs));
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(info);
+
+ ClusterVO cluster = provider.discoverClusters().get(0);
+
+ assertThat(cluster.getBrokers()).singleElement()
+ .extracting(broker -> broker.getAddr())
+ .isEqualTo("10.0.0.12:10911");
+ }
+
@Test
void discoverClustersShouldDiscoverProxiesViaHeartbeatSyncerTest() throws
Exception {
DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
@@ -179,6 +218,28 @@ class RocketMQClusterProviderTest {
assertThat(clusters.get(0).getProxies().get(0).getStatus()).isEqualTo(ClusterStatus.healthy);
}
+ @Test
+ void discoverClustersShouldIgnoreNullHeartbeatConnectionsTest() throws
Exception {
+ DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+ RocketMQClusterProvider provider = newProvider(adminExt);
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
+ Connection valid = new Connection();
+ valid.setClientAddr("10.0.3.5:54321");
+ java.util.HashSet<Connection> connections = new java.util.HashSet<>();
+ connections.add(null);
+ connections.add(valid);
+ ConsumerConnection consumerConnection = new ConsumerConnection();
+ consumerConnection.setConnectionSet(connections);
+
when(adminExt.examineConsumerConnectionInfo("CID_DefaultHeartBeatSyncerTopic"))
+ .thenReturn(consumerConnection);
+
+ List<ClusterVO> clusters = provider.discoverClusters();
+
+ assertThat(clusters.get(0).getProxies()).singleElement()
+ .extracting(proxy -> proxy.getAddr())
+ .isEqualTo("10.0.3.5:8080");
+ }
+
private RocketMQClusterProvider newProvider(DefaultMQAdminExt adminExt) {
MqAdminExtFactory adminFactory = mock(MqAdminExtFactory.class);
when(adminFactory.execute(anyString(), any(),
any())).thenAnswer(invocation ->
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
index d1043b1a6..8e431af52 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
@@ -76,7 +76,7 @@ class TencentAclServiceTest {
@Test
void updateUserUsesUsernameAsRoleNameTest() throws Exception {
AclUserVO updated = service.updateUser(INSTANCE_ID, AclUserVO.builder()
- .username("reader-role")
+ .username(" reader-role ")
.permRead(false)
.permWrite(true)
.build());
@@ -91,6 +91,30 @@ class TencentAclServiceTest {
assertThat(updated.getUsername()).isEqualTo("reader-role");
}
+ @Test
+ void listUsersShouldSkipRolesWithoutNamesTest() throws Exception {
+ RoleItem blank = new RoleItem();
+ blank.setRoleName(" ");
+ RoleItem valid = new RoleItem();
+ valid.setRoleName("reader-role");
+ DescribeRoleListResponse response = new DescribeRoleListResponse();
+ response.setData(new RoleItem[]{null, blank, valid});
+ when(client.DescribeRoleList(any())).thenReturn(response);
+
+ assertThat(service.listUsers(INSTANCE_ID))
+ .extracting(AclUserVO::getUsername)
+ .containsExactly("reader-role");
+ assertThat(service.getUserCredentials(INSTANCE_ID, " reader-role
").getUsername())
+ .isEqualTo("reader-role");
+ }
+
+ @Test
+ void createUserShouldRejectNullPayloadTest() {
+ assertThatThrownBy(() -> service.createUser(INSTANCE_ID, null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("ACL username is required");
+ }
+
@Test
void updateUserShouldPreserveOmittedPermissionsTest() throws Exception {
RoleItem role = new RoleItem();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentCatalogServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentCatalogServiceTest.java
index 5fff85ad9..1d4714a1d 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentCatalogServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentCatalogServiceTest.java
@@ -85,6 +85,26 @@ class TencentCatalogServiceTest {
assertThat(instances.get(0).getTopicCount()).isEqualTo(3);
}
+ @Test
+ void listCloudInstancesShouldNormalizeOutOfRangeCountsAndTrimSearchTest() {
+ InstanceItem item = new InstanceItem();
+ item.setInstanceId("rmq-abc");
+ item.setInstanceName("chengdu-prod");
+ item.setTopicNum(Long.MAX_VALUE);
+ item.setGroupNum(-1L);
+ DescribeInstanceListResponse response = new
DescribeInstanceListResponse();
+ response.setData(new InstanceItem[]{item});
+ when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION),
any())).thenReturn(response);
+
+ List<CloudInstanceOptionVO> instances = service.listCloudInstances(
+ CREDENTIAL_ID, REGION, " PROD ");
+
+ assertThat(instances).singleElement().satisfies(instance -> {
+ assertThat(instance.getTopicCount()).isEqualTo(Integer.MAX_VALUE);
+ assertThat(instance.getGroupCount()).isZero();
+ });
+ }
+
@Test
void listCloudInstancesShouldSkipNullItemsTest() {
InstanceItem item = new InstanceItem();
@@ -120,6 +140,24 @@ class TencentCatalogServiceTest {
assertThat(detail.getEndpoints().get(1).getEndpointType()).isEqualTo("TCP_INTERNET");
}
+ @Test
+ void getCloudInstanceShouldSkipBlankOpenEndpointsTest() {
+ Endpoint blank = endpoint("PUBLIC", "OPEN", " ");
+ Endpoint missing = endpoint("VPC", "OPEN", null);
+ Endpoint valid = endpoint("VPC", "OPEN", " vpc.tencent:8080 ");
+ DescribeInstanceResponse response = new DescribeInstanceResponse();
+ response.setInstanceId("rmq-abc");
+ response.setEndpointList(new Endpoint[]{blank, missing, valid});
+ when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION),
any())).thenReturn(response);
+
+ CloudInstanceDetailVO detail = service.getCloudInstance(CREDENTIAL_ID,
REGION, "rmq-abc");
+
+ assertThat(detail.getEndpoints()).singleElement().satisfies(endpoint
-> {
+ assertThat(endpoint.getEndpointType()).isEqualTo("TCP_VPC");
+
assertThat(endpoint.getEndpointUrl()).isEqualTo("vpc.tencent:8080");
+ });
+ }
+
private static Endpoint endpoint(String type, String status, String url) {
Endpoint endpoint = new Endpoint();
endpoint.setType(type);