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 ce5bf009b fix(server): harden external provider boundaries (#2197)
ce5bf009b is described below

commit ce5bf009b01ee2c6186a11e044d51cbcab1adf0b
Author: btlqql <[email protected]>
AuthorDate: Tue Aug 18 18:54:53 2026 +0800

    fix(server): harden external provider boundaries (#2197)
---
 .../studio/cluster/metrics/MetricsService.java     |  20 +++-
 .../studio/ops/ai/OpenAiCompatibleLlmClient.java   |   6 +
 .../studio/provider/tencent/TencentAclService.java |  47 ++++++--
 .../provider/tencent/TencentCatalogService.java    |   3 +
 .../studio/cluster/metrics/MetricsServiceTest.java |  30 +++++
 .../ops/ai/OpenAiCompatibleLlmClientTest.java      |  21 ++++
 .../provider/tencent/TencentAclServiceTest.java    | 129 +++++++++++++++++++++
 .../tencent/TencentCatalogServiceTest.java         |  16 +++
 8 files changed, 263 insertions(+), 9 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
 
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
index 2ceb35f08..da2de5bba 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
@@ -40,7 +40,7 @@ public class MetricsService {
     private static final long MAX_RANGE_SECONDS = 31L * 24 * 60 * 60;
     private static final long MAX_SAMPLE_POINTS = 11_000L;
     private static final Pattern NUMBER_PATTERN = 
Pattern.compile("\\d+(?:\\.\\d+)?");
-    private static final Pattern DURATION_PART_PATTERN = 
Pattern.compile("(\\d+(?:\\.\\d+)?)(ms|s|m|h|d|w|y)");
+    private static final Pattern DURATION_PART_PATTERN = 
Pattern.compile("(\\d+)(ms|s|m|h|d|w|y)");
     private static final Map<String, BigDecimal> UNIT_TO_MILLIS = Map.of(
             "ms", BigDecimal.ONE,
             "s", BigDecimal.valueOf(1_000L),
@@ -50,6 +50,15 @@ public class MetricsService {
             "w", BigDecimal.valueOf(604_800_000L),
             "y", BigDecimal.valueOf(31_536_000_000L)
     );
+    private static final Map<String, Integer> UNIT_ORDER = Map.of(
+            "y", 0,
+            "w", 1,
+            "d", 2,
+            "h", 3,
+            "m", 4,
+            "s", 5,
+            "ms", 6
+    );
 
     private final MetricsSource metricsSource;
     private final MetricProfileService metricProfileService;
@@ -188,13 +197,20 @@ public class MetricsService {
         Matcher matcher = DURATION_PART_PATTERN.matcher(value);
         BigDecimal millis = BigDecimal.ZERO;
         int position = 0;
+        int previousUnitOrder = -1;
         while (matcher.find()) {
             if (matcher.start() != position) {
                 throw badRequest("Metric query step is invalid");
             }
+            String unit = matcher.group(2);
+            int unitOrder = UNIT_ORDER.get(unit);
+            if (unitOrder <= previousUnitOrder) {
+                throw badRequest("Metric query step is invalid");
+            }
             BigDecimal amount = new BigDecimal(matcher.group(1));
-            millis = 
millis.add(amount.multiply(UNIT_TO_MILLIS.get(matcher.group(2))));
+            millis = millis.add(amount.multiply(UNIT_TO_MILLIS.get(unit)));
             position = matcher.end();
+            previousUnitOrder = unitOrder;
         }
         if (position != value.length()) {
             throw badRequest("Metric query step is invalid");
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClient.java
 
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClient.java
index 947a9c815..1b1b6d228 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClient.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClient.java
@@ -428,6 +428,12 @@ public class OpenAiCompatibleLlmClient {
     private String parseDelta(String body) {
         try {
             JsonNode root = objectMapper.readTree(body);
+            String errorMessage = root.path("error").path("message").asText();
+            if (StringUtils.hasText(errorMessage)) {
+                throw new LlmGatewayException(502, "llm.provider.stream_error",
+                        "LLM provider stream failed: " + errorMessage,
+                        "Check the provider credentials, model name, and 
account quota.");
+            }
             return 
root.path("choices").path(0).path("delta").path("content").asText("");
         } catch (JsonProcessingException exception) {
             throw new LlmGatewayException(502, 
"llm.provider.malformed_stream_event",
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 4f1b93259..2c9985473 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
@@ -150,28 +150,61 @@ public class TencentAclService {
 
     public AclUserVO updateUser(String instanceId, AclUserVO user) {
         Context context = resolve(instanceId);
-        if (!StringUtils.hasText(user.getUsername())) {
+        // 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");
         }
+        RoleItem existing = user.getPermRead() == null || user.getPermWrite() 
== null
+                ? findRole(context, roleName) : null;
+        boolean permRead = user.getPermRead() == null
+                ? Boolean.TRUE.equals(existing.getPermRead()) : 
user.getPermRead();
+        boolean permWrite = user.getPermWrite() == null
+                ? Boolean.TRUE.equals(existing.getPermWrite()) : 
user.getPermWrite();
         ModifyRoleRequest request = new ModifyRoleRequest();
         request.setInstanceId(context.cloudInstanceId());
-        request.setRole(user.getUsername());
-        request.setPermRead(user.getPermRead() == null || user.getPermRead());
-        request.setPermWrite(user.getPermWrite() == null || 
user.getPermWrite());
+        request.setRole(roleName);
+        request.setPermRead(permRead);
+        request.setPermWrite(permWrite);
         clientFactory.call(context.credentialId(), context.regionId(),
                 client -> client.ModifyRole(request));
         return AclUserVO.builder()
-                .username(user.getUsername())
+                .username(roleName)
                 .accessKey(null)
                 .secretKey(null)
                 .admin(false)
-                .permRead(user.getPermRead() == null || user.getPermRead())
-                .permWrite(user.getPermWrite() == null || user.getPermWrite())
+                .permRead(permRead)
+                .permWrite(permWrite)
                 .clusters(context.cloudInstanceId() == null ? List.of() : 
List.of(context.cloudInstanceId()))
                 .gmtCreate(user.getGmtCreate())
                 .build();
     }
 
+    private RoleItem findRole(Context context, String roleName) {
+        for (int page = 0; page < MAX_PAGES; page++) {
+            DescribeRoleListRequest request = new DescribeRoleListRequest();
+            request.setInstanceId(context.cloudInstanceId());
+            request.setOffset((long) page * PAGE_SIZE);
+            request.setLimit((long) PAGE_SIZE);
+            DescribeRoleListResponse response = 
clientFactory.call(context.credentialId(),
+                    context.regionId(), client -> 
client.DescribeRoleList(request));
+            RoleItem[] data = response == null ? null : response.getData();
+            if (data == null || data.length == 0) {
+                break;
+            }
+            for (RoleItem role : data) {
+                if (role != null && roleName.equals(role.getRoleName())) {
+                    return role;
+                }
+            }
+            if (data.length < PAGE_SIZE) {
+                break;
+            }
+        }
+        throw new BusinessException(404, "ACL user not found: " + roleName);
+    }
+
     public void deleteUser(String instanceId, String username) {
         Context context = resolve(instanceId);
         if (!StringUtils.hasText(username)) {
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 32a13b2bd..fd84f571f 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
@@ -94,6 +94,9 @@ public class TencentCatalogService implements 
CloudCatalogProvider {
                 break;
             }
             for (InstanceItem item : data) {
+                if (item == null) {
+                    continue;
+                }
                 CloudInstanceOptionVO option = toInstanceOption(item, 
regionId);
                 if (matchesSearch(search, option)) {
                     instances.add(option);
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
index 1c9be5b2e..fd3ff5b2e 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
@@ -270,6 +270,36 @@ class MetricsServiceTest {
         verifyNoInteractions(metricsSource);
     }
 
+    @Test
+    void queryShouldRejectInvalidPrometheusDurationCombinations() {
+        for (String step : List.of("1s1h", "1m1m", "1ms1s", "1.5h")) {
+            MetricQueryDTO query = MetricQueryDTO.builder()
+                    .metric("rocketmq_messages_in_total")
+                    .start(1700000000L)
+                    .end(1700003600L)
+                    .step(step)
+                    .build();
+
+            assertBadRequest(query, "Metric query step is invalid");
+        }
+        verifyNoInteractions(metricsSource);
+    }
+
+    @Test
+    void queryShouldAcceptOrderedPrometheusDurationCombinations() {
+        MetricQueryDTO query = MetricQueryDTO.builder()
+                .metric("rocketmq_messages_in_total")
+                .start(1700000000L)
+                .end(1700003600L)
+                .step("1h30m15s500ms")
+                .build();
+        when(metricsSource.query(query)).thenReturn(emptyMetricData());
+
+        metricsService.query(query);
+
+        verify(metricsSource).query(query);
+    }
+
     @Test
     void queryShouldRejectTooManySamples() {
         MetricQueryDTO query = MetricQueryDTO.builder()
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClientTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClientTest.java
index 11b0b06c0..ff4407142 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClientTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/OpenAiCompatibleLlmClientTest.java
@@ -152,6 +152,27 @@ class OpenAiCompatibleLlmClientTest {
         assertThat(requestBody.get().path("stream").asBoolean()).isTrue();
     }
 
+    @Test
+    void streamShouldExposeErrorEnvelopeFromSuccessfulResponse() {
+        server.createContext("/v1/chat/completions", exchange -> 
respond(exchange, 200, """
+                data: {"error":{"message":"quota exceeded"}}
+
+                data: [DONE]
+
+                """, "text/event-stream"));
+
+        assertThatThrownBy(() -> client.stream(
+                config("openai", "sk-test"), "hello", null, token -> { }))
+                .isInstanceOf(LlmGatewayException.class)
+                .hasMessage("LLM provider stream failed: quota exceeded")
+                .satisfies(exception -> {
+                    LlmGatewayException gatewayException = 
(LlmGatewayException) exception;
+                    
assertThat(gatewayException.getStatusCode()).isEqualTo(502);
+                    
assertThat(gatewayException.getCode()).isEqualTo("llm.provider.stream_error");
+                    assertThat(gatewayException.getHint()).contains("account 
quota");
+                });
+    }
+
     @Test
     void streamShouldEnforceTimeoutWhileReadingResponseBody() {
         OpenAiCompatibleLlmClient timeoutClient = new 
OpenAiCompatibleLlmClient(
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
new file mode 100644
index 000000000..d1043b1a6
--- /dev/null
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentAclServiceTest.java
@@ -0,0 +1,129 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.provider.tencent;
+
+import com.tencentcloudapi.trocket.v20230308.TrocketClient;
+import com.tencentcloudapi.trocket.v20230308.models.DescribeRoleListResponse;
+import com.tencentcloudapi.trocket.v20230308.models.ModifyRoleRequest;
+import com.tencentcloudapi.trocket.v20230308.models.RoleItem;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.instance.InstanceRepository;
+import org.apache.rocketmq.studio.instance.InstanceVO;
+import org.apache.rocketmq.studio.instance.acl.AclUserVO;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import 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.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.lenient;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+@ExtendWith(MockitoExtension.class)
+class TencentAclServiceTest {
+
+    private static final String INSTANCE_ID = "instance-1";
+    private static final String CLOUD_INSTANCE_ID = "rmq-abc";
+    private static final Long CREDENTIAL_ID = 1L;
+    private static final String REGION = "ap-guangzhou";
+
+    @Mock
+    private TencentClientFactory clientFactory;
+    @Mock
+    private InstanceRepository instanceRepository;
+    @Mock
+    private TrocketClient client;
+
+    private TencentAclService service;
+
+    @BeforeEach
+    void setUp() {
+        service = new TencentAclService(clientFactory, instanceRepository);
+        
when(instanceRepository.findByIdentifier(INSTANCE_ID)).thenReturn(Optional.of(InstanceVO.builder()
+                .cloudInstanceId(CLOUD_INSTANCE_ID)
+                .credentialId(CREDENTIAL_ID)
+                .regionId(REGION)
+                .build()));
+        lenient().when(clientFactory.call(anyLong(), anyString(), 
any())).thenAnswer(invocation -> {
+            TencentClientFactory.TencentCall<Object> action = 
invocation.getArgument(2);
+            return action.execute(client);
+        });
+    }
+
+    @Test
+    void updateUserUsesUsernameAsRoleNameTest() throws Exception {
+        AclUserVO updated = service.updateUser(INSTANCE_ID, AclUserVO.builder()
+                .username("reader-role")
+                .permRead(false)
+                .permWrite(true)
+                .build());
+
+        ArgumentCaptor<ModifyRoleRequest> requestCaptor = 
ArgumentCaptor.forClass(ModifyRoleRequest.class);
+        verify(client).ModifyRole(requestCaptor.capture());
+        ModifyRoleRequest request = requestCaptor.getValue();
+        assertThat(request.getRole()).isEqualTo("reader-role");
+        assertThat(request.getPermRead()).isFalse();
+        assertThat(request.getPermWrite()).isTrue();
+        assertThat(updated.getId()).isNull();
+        assertThat(updated.getUsername()).isEqualTo("reader-role");
+    }
+
+    @Test
+    void updateUserShouldPreserveOmittedPermissionsTest() throws Exception {
+        RoleItem role = new RoleItem();
+        role.setRoleName("reader-role");
+        role.setPermRead(false);
+        role.setPermWrite(false);
+        DescribeRoleListResponse response = new DescribeRoleListResponse();
+        response.setData(new RoleItem[]{role});
+        when(client.DescribeRoleList(any())).thenReturn(response);
+
+        AclUserVO updated = service.updateUser(INSTANCE_ID, AclUserVO.builder()
+                .username("reader-role")
+                .build());
+
+        ArgumentCaptor<ModifyRoleRequest> requestCaptor = 
ArgumentCaptor.forClass(ModifyRoleRequest.class);
+        verify(client).ModifyRole(requestCaptor.capture());
+        ModifyRoleRequest request = requestCaptor.getValue();
+        assertThat(request.getPermRead()).isFalse();
+        assertThat(request.getPermWrite()).isFalse();
+        assertThat(updated.getPermRead()).isFalse();
+        assertThat(updated.getPermWrite()).isFalse();
+    }
+
+    @Test
+    void updateUserThrowsNotFoundWhenRoleMissingTest() throws Exception {
+        DescribeRoleListResponse response = new DescribeRoleListResponse();
+        response.setData(new RoleItem[0]);
+        when(client.DescribeRoleList(any())).thenReturn(response);
+
+        assertThatThrownBy(() -> service.updateUser(INSTANCE_ID, 
AclUserVO.builder()
+                .username("ghost-role")
+                .build()))
+                .isInstanceOf(BusinessException.class)
+                .hasMessageContaining("ACL user not found: ghost-role");
+    }
+}
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 44ae6c1cd..5fff85ad9 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,22 @@ class TencentCatalogServiceTest {
         assertThat(instances.get(0).getTopicCount()).isEqualTo(3);
     }
 
+    @Test
+    void listCloudInstancesShouldSkipNullItemsTest() {
+        InstanceItem item = new InstanceItem();
+        item.setInstanceId("rmq-valid");
+        item.setInstanceName("chengdu-prod");
+        DescribeInstanceListResponse response = new 
DescribeInstanceListResponse();
+        response.setData(new InstanceItem[]{null, item});
+        when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION), 
any())).thenReturn(response);
+
+        List<CloudInstanceOptionVO> instances = 
service.listCloudInstances(CREDENTIAL_ID, REGION, null);
+
+        assertThat(instances).singleElement()
+                .extracting(CloudInstanceOptionVO::getInstanceId)
+                .isEqualTo("rmq-valid");
+    }
+
     @Test
     void getCloudInstanceShouldMapOnlyOpenEndpointsTest() {
         Endpoint vpc = endpoint("VPC", "OPEN", "vpc.tencent:8080");

Reply via email to