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 33fe90ff0 fix(server): consolidate boundary correctness (#2849)
33fe90ff0 is described below
commit 33fe90ff03095c0645f7764400e47cb0f36cb3b5
Author: shown <[email protected]>
AuthorDate: Wed Sep 2 15:34:11 2026 +0800
fix(server): consolidate boundary correctness (#2849)
* fix(server): normalize service boundary inputs
Signed-off-by: yuluo-yx <[email protected]>
* fix(aliyun): normalize catalog lookup inputs
Signed-off-by: yuluo-yx <[email protected]>
* fix(metrics): validate query windows without overflow
Signed-off-by: yuluo-yx <[email protected]>
* fix(message): ignore time bounds for id lookup
Signed-off-by: yuluo-yx <[email protected]>
* fix(broker): trim flush disk configuration
Signed-off-by: yuluo-yx <[email protected]>
* [ISSUE #2710] fix(auth): prefer explicit Bearer credentials
Signed-off-by: yuluo-yx <[email protected]>
* [ISSUE #2709] fix(api): return 406 for unacceptable response types
Signed-off-by: yuluo-yx <[email protected]>
---------
Signed-off-by: yuluo-yx <[email protected]>
---
.../apache/rocketmq/studio/auth/AuthCookie.java | 12 +++-
.../apache/rocketmq/studio/auth/AuthService.java | 3 +
.../studio/cluster/metrics/MetricsService.java | 9 ++-
.../nameserver/NameServerConfigDiffService.java | 19 ++++++-
.../common/exception/GlobalExceptionHandler.java | 13 +++++
.../rocketmq/studio/instance/dlq/DLQService.java | 65 ++++++++++++++++------
.../provider/alibaba/AliyunCatalogService.java | 14 +++--
.../apache/RocketMQBrokerConfigService.java | 2 +-
.../provider/apache/RocketMQMessageProvider.java | 23 ++++----
.../rocketmq/studio/auth/AuthCookieTest.java | 33 +++++++++--
.../studio/auth/AuthServiceDatabaseTest.java | 15 +++++
.../studio/cluster/metrics/MetricsServiceTest.java | 41 ++++++++++++++
.../NameServerConfigDiffServiceTest.java | 34 +++++++++++
.../exception/GlobalExceptionHandlerTest.java | 16 ++++++
.../studio/instance/dlq/DLQServiceTest.java | 41 ++++++++++++++
.../provider/alibaba/AliyunCatalogServiceTest.java | 24 ++++++++
.../apache/RocketMQBrokerConfigServiceTest.java | 20 +++++++
.../apache/RocketMQMessageProviderTest.java | 17 ++++++
18 files changed, 356 insertions(+), 45 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/auth/AuthCookie.java
b/server/src/main/java/org/apache/rocketmq/studio/auth/AuthCookie.java
index 8a7c2ad22..d51a67339 100644
--- a/server/src/main/java/org/apache/rocketmq/studio/auth/AuthCookie.java
+++ b/server/src/main/java/org/apache/rocketmq/studio/auth/AuthCookie.java
@@ -36,6 +36,10 @@ final class AuthCookie {
}
static String authorization(HttpServletRequest request, AuthProperties
properties) {
+ String authorization = request.getHeader(HttpHeaders.AUTHORIZATION);
+ if (hasBearerToken(authorization)) {
+ return authorization;
+ }
Cookie[] cookies = request.getCookies();
if (cookies != null) {
for (Cookie cookie : cookies) {
@@ -46,7 +50,13 @@ final class AuthCookie {
}
}
// API clients that explicitly requested a bearer token at login
authenticate with this header.
- return request.getHeader(HttpHeaders.AUTHORIZATION);
+ return authorization;
+ }
+
+ private static boolean hasBearerToken(String authorization) {
+ return authorization != null
+ && authorization.regionMatches(true, 0, TOKEN_PREFIX, 0,
TOKEN_PREFIX.length())
+ && !authorization.substring(TOKEN_PREFIX.length()).isBlank();
}
static boolean requestsBearerToken(HttpServletRequest request) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/auth/AuthService.java
b/server/src/main/java/org/apache/rocketmq/studio/auth/AuthService.java
index 148f12761..41409d03c 100644
--- a/server/src/main/java/org/apache/rocketmq/studio/auth/AuthService.java
+++ b/server/src/main/java/org/apache/rocketmq/studio/auth/AuthService.java
@@ -213,6 +213,9 @@ public class AuthService {
public RmqStudioUser setUserEnabled(Long userId, boolean enabled) {
requireDatabaseBacked();
RmqStudioUser user = getUser(userId);
+ if (Boolean.valueOf(enabled).equals(user.getEnabled())) {
+ return user;
+ }
if (!enabled && Boolean.TRUE.equals(user.getAdmin())) {
// Lock the enabled administrator rows so concurrent disables
serialize. A plain
// count would let two requests both observe a count of 2 and
disable everyone.
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 da2de5bba..2f6b5206f 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
@@ -138,13 +138,16 @@ public class MetricsService {
}
private void validateQueryWindow(MetricQueryDTO query) {
- long rangeSeconds = query.getEnd() - query.getStart();
- if (rangeSeconds <= 0) {
+ long start = query.getStart();
+ long end = query.getEnd();
+ if (end <= start) {
throw badRequest("Metric query end must be later than start");
}
- if (rangeSeconds > MAX_RANGE_SECONDS) {
+ if (start <= Long.MAX_VALUE - MAX_RANGE_SECONDS
+ && end > start + MAX_RANGE_SECONDS) {
throw badRequest("Metric query range must not exceed 31 days");
}
+ long rangeSeconds = end - start;
BigDecimal stepMillis = parseStepMillis(query.getStep());
if (stepMillis.signum() <= 0) {
throw badRequest("Metric query step must be positive");
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerConfigDiffService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerConfigDiffService.java
index cf27e7876..7179ed8ac 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerConfigDiffService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerConfigDiffService.java
@@ -34,6 +34,7 @@ import java.io.InputStream;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.Properties;
import java.util.stream.Stream;
@@ -192,14 +193,28 @@ public class NameServerConfigDiffService {
.filter(node -> node != null)
.map(NameServerVO::getAddr);
Stream<String> endpointAddresses =
splitEndpoint(cluster.getEndpoint()).stream();
- return Stream.concat(declared, endpointAddresses)
+ Map<String, String> uniqueAddresses = new LinkedHashMap<>();
+ Stream.concat(declared, endpointAddresses)
.filter(address -> address != null && !address.isBlank())
.map(String::trim)
- .distinct()
+ .forEach(address ->
uniqueAddresses.putIfAbsent(canonicalAddressKey(address), address));
+ return uniqueAddresses.values().stream()
.sorted()
.toList();
}
+ private String canonicalAddressKey(String address) {
+ int separator = address.lastIndexOf(':');
+ if (separator <= 0 || address.indexOf(':') != separator) {
+ return address;
+ }
+ String host = address.substring(0, separator);
+ if (host.indexOf('%') >= 0) {
+ return address;
+ }
+ return host.toLowerCase(Locale.ROOT) + address.substring(separator);
+ }
+
private String connectionEndpoint(ClusterVO cluster, List<String>
addresses) {
if (cluster.getEndpoint() != null && !cluster.getEndpoint().isBlank())
{
return cluster.getEndpoint().trim();
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/exception/GlobalExceptionHandler.java
b/server/src/main/java/org/apache/rocketmq/studio/common/exception/GlobalExceptionHandler.java
index 2fe73fffe..8134d7bca 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/exception/GlobalExceptionHandler.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/common/exception/GlobalExceptionHandler.java
@@ -23,6 +23,7 @@ import org.apache.rocketmq.studio.ops.ai.LlmGatewayException;
import org.springframework.http.HttpStatus;
import org.springframework.http.ResponseEntity;
import org.springframework.http.converter.HttpMessageNotReadableException;
+import org.springframework.web.HttpMediaTypeNotAcceptableException;
import org.springframework.web.HttpMediaTypeNotSupportedException;
import org.springframework.web.HttpRequestMethodNotSupportedException;
import org.springframework.web.bind.MethodArgumentNotValidException;
@@ -145,4 +146,16 @@ public class GlobalExceptionHandler {
log.warn("Unsupported media type: {}", ex.getMessage());
return Result.error(HttpStatus.UNSUPPORTED_MEDIA_TYPE.value(),
ex.getMessage());
}
+
+ /**
+ * Requests that cannot accept any available response representation must
be
+ * reported as 406, not 500.
+ */
+ @ExceptionHandler(HttpMediaTypeNotAcceptableException.class)
+ @ResponseStatus(HttpStatus.NOT_ACCEPTABLE)
+ public Result<?> handleHttpMediaTypeNotAcceptableException(
+ HttpMediaTypeNotAcceptableException ex) {
+ log.warn("No acceptable response media type: {}", ex.getMessage());
+ return Result.error(HttpStatus.NOT_ACCEPTABLE.value(),
ex.getMessage());
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
index 87434d9a1..8217622aa 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
@@ -56,51 +56,65 @@ public class DLQService {
public DLQResendResultVO resendMessages(String instanceId, String
groupName, Long startTime, Long endTime,
String targetTopic) {
requireApacheInstance(instanceId);
- validateResendRequest(groupName, startTime, endTime);
- log.info("Resending DLQ messages: group={}, targetTopic={}",
groupName, targetTopic);
- return dlqProvider.resendMessages(instanceId, groupName, startTime,
endTime, targetTopic);
+ String normalizedGroupName = requireGroupName(groupName);
+ validateTimeRange(startTime, endTime);
+ String normalizedTargetTopic = normalizeOptional(targetTopic);
+ log.info("Resending DLQ messages: group={}, targetTopic={}",
+ normalizedGroupName, normalizedTargetTopic);
+ return dlqProvider.resendMessages(
+ instanceId, normalizedGroupName, startTime, endTime,
normalizedTargetTopic);
}
public DLQExportResultVO exportMessages(String instanceId, String
groupName, Long startTime, Long endTime,
Integer maxCount) {
requireApacheInstance(instanceId);
- validateResendRequest(groupName, startTime, endTime);
- log.info("Exporting DLQ messages: group={}, maxCount={}", groupName,
maxCount);
- return dlqProvider.exportMessages(instanceId, groupName, startTime,
endTime, maxCount);
+ String normalizedGroupName = requireGroupName(groupName);
+ validateTimeRange(startTime, endTime);
+ log.info("Exporting DLQ messages: group={}, maxCount={}",
normalizedGroupName, maxCount);
+ return dlqProvider.exportMessages(instanceId, normalizedGroupName,
startTime, endTime, maxCount);
}
public PageResult<DLQMessageVO> listMessages(String instanceId, String
groupName, Long startTime, Long endTime,
int page, int pageSize) {
requireApacheInstance(instanceId);
- validateResendRequest(groupName, startTime, endTime);
- log.info("Listing DLQ messages: group={}, page={}, pageSize={}",
groupName, page, pageSize);
- return dlqProvider.listMessages(instanceId, groupName, startTime,
endTime, page, pageSize);
+ String normalizedGroupName = requireGroupName(groupName);
+ validateTimeRange(startTime, endTime);
+ log.info("Listing DLQ messages: group={}, page={}, pageSize={}",
normalizedGroupName, page, pageSize);
+ return dlqProvider.listMessages(
+ instanceId, normalizedGroupName, startTime, endTime, page,
pageSize);
}
public DLQResendResultVO resendSelectedMessages(String instanceId, String
groupName, List<String> msgIds,
String targetTopic) {
requireApacheInstance(instanceId);
+ String normalizedGroupName = requireGroupName(groupName);
if (msgIds == null || msgIds.isEmpty()) {
throw new BusinessException(400, "At least one msgId is required");
}
if (msgIds.size() > MAX_SELECTED_MESSAGES) {
throw new BusinessException(400, "At most 100 msgIds are allowed
per resend");
}
+ List<String> normalizedMsgIds = normalizeMsgIds(msgIds);
+ String normalizedTargetTopic = normalizeOptional(targetTopic);
log.info("Resending selected DLQ messages: group={}, count={},
targetTopic={}",
- groupName, msgIds.size(), targetTopic);
- return dlqProvider.resendMessages(instanceId, groupName, msgIds,
targetTopic);
+ normalizedGroupName, normalizedMsgIds.size(),
normalizedTargetTopic);
+ return dlqProvider.resendMessages(
+ instanceId, normalizedGroupName, normalizedMsgIds,
normalizedTargetTopic);
}
public DLQExcelExportResultVO exportExcel(String instanceId, String
groupName, Long startTime, Long endTime,
List<String> msgIds) {
requireApacheInstance(instanceId);
- validateResendRequest(groupName, startTime, endTime);
+ String normalizedGroupName = requireGroupName(groupName);
+ validateTimeRange(startTime, endTime);
if (msgIds != null && msgIds.size() > MAX_SELECTED_MESSAGES) {
throw new BusinessException(400, "At most 100 msgIds are allowed
per export");
}
- log.info("Exporting DLQ messages as Excel: group={}, selected={}",
groupName,
- msgIds == null ? 0 : msgIds.size());
- return dlqProvider.exportExcel(instanceId, groupName, startTime,
endTime, msgIds);
+ List<String> normalizedMsgIds = msgIds == null ? null :
normalizeMsgIds(msgIds);
+ log.info("Exporting DLQ messages as Excel: group={}, selected={}",
normalizedGroupName,
+ normalizedMsgIds == null ? 0 : normalizedMsgIds.size());
+ return dlqProvider.exportExcel(
+ instanceId, normalizedGroupName, startTime, endTime,
normalizedMsgIds);
}
private void requireApacheInstance(String instanceId) {
@@ -111,10 +125,29 @@ public class DLQService {
});
}
- private void validateResendRequest(String groupName, Long startTime, Long
endTime) {
+ private String requireGroupName(String groupName) {
if (!StringUtils.hasText(groupName)) {
throw new BusinessException(400, "groupName is required");
}
+ return groupName.trim();
+ }
+
+ private String normalizeOptional(String value) {
+ return StringUtils.hasText(value) ? value.trim() : null;
+ }
+
+ private List<String> normalizeMsgIds(List<String> msgIds) {
+ return msgIds.stream()
+ .map(msgId -> {
+ if (!StringUtils.hasText(msgId)) {
+ throw new BusinessException(400, "msgId must not be
blank");
+ }
+ return msgId.trim();
+ })
+ .toList();
+ }
+
+ private void validateTimeRange(Long startTime, Long endTime) {
if ((startTime == null) != (endTime == null)) {
throw new BusinessException(400, "startTime and endTime must be
provided together");
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
index cc011df88..69696884c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
@@ -83,14 +83,16 @@ public class AliyunCatalogService implements
CloudCatalogProvider {
public List<CloudInstanceOptionVO> listCloudInstances(Long credentialId,
String regionId, String search) {
requireId(credentialId, "credentialId");
requireNonBlank(regionId, "regionId");
- List<ListInstancesResponseBody.List> all =
fetchAllInstances(credentialId, regionId);
+ String normalizedRegionId = regionId.strip();
+ String normalizedSearch = search == null ? null : search.strip();
+ List<ListInstancesResponseBody.List> all =
fetchAllInstances(credentialId, normalizedRegionId);
List<CloudInstanceOptionVO> options = new ArrayList<>();
for (ListInstancesResponseBody.List item : all) {
if (item == null) {
continue;
}
CloudInstanceOptionVO vo =
AliyunConverters.toInstanceOptionVO(item);
- if (matchesSearch(search, vo)) {
+ if (matchesSearch(normalizedSearch, vo)) {
options.add(vo);
}
}
@@ -102,13 +104,15 @@ public class AliyunCatalogService implements
CloudCatalogProvider {
requireId(credentialId, "credentialId");
requireNonBlank(regionId, "regionId");
requireNonBlank(cloudInstanceId, "cloudInstanceId");
- GetInstanceRequest request =
GetInstanceRequest.builder().instanceId(cloudInstanceId).build();
- GetInstanceResponse response = clientFactory.call(credentialId,
regionId,
+ String normalizedRegionId = regionId.strip();
+ String normalizedCloudInstanceId = cloudInstanceId.strip();
+ GetInstanceRequest request =
GetInstanceRequest.builder().instanceId(normalizedCloudInstanceId).build();
+ GetInstanceResponse response = clientFactory.call(credentialId,
normalizedRegionId,
client -> client.getInstance(request));
GetInstanceResponseBody body = response == null ? null :
response.getBody();
GetInstanceResponseBody.Data data = body == null ? null :
body.getData();
if (data == null) {
- throw new BusinessException(404, "Aliyun instance not found: " +
cloudInstanceId);
+ throw new BusinessException(404, "Aliyun instance not found: " +
normalizedCloudInstanceId);
}
return AliyunConverters.toInstanceDetailVO(data);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigService.java
index 668266c98..de4dd6989 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigService.java
@@ -116,7 +116,7 @@ public class RocketMQBrokerConfigService {
private FlushDiskType parseFlushDiskType(String value) {
try {
- return FlushDiskType.valueOf(value);
+ return FlushDiskType.valueOf(value.trim());
} catch (IllegalArgumentException e) {
return FlushDiskType.ASYNC_FLUSH;
}
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 612c93720..883ad052d 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
@@ -115,28 +115,27 @@ public class RocketMQMessageProvider implements
MessageProvider {
String topic, String msgId,
String tag, String key,
Long startTime, Long endTime)
{
+ if (StringUtils.hasText(msgId)) {
+ return queryByMsgId(adminExt, topic, msgId);
+ }
+
long end = endTime != null ? endTime : System.currentTimeMillis();
long begin = startTime != null ? startTime : end - ONE_HOUR_MILLIS;
if (begin >= end) {
throw new BusinessException(400, "Message query start time must be
before end time");
}
-
- List<MessageRecordVO> result;
- if (StringUtils.hasText(msgId)) {
- result = queryByMsgId(adminExt, topic, msgId);
- } else if (StringUtils.hasText(topic) && StringUtils.hasText(key)) {
- result = queryByKey(adminExt, topic, key, tag, begin, end);
- } else if (StringUtils.hasText(topic)) {
+ if (StringUtils.hasText(topic) && StringUtils.hasText(key)) {
+ return queryByKey(adminExt, topic, key, tag, begin, end);
+ }
+ if (StringUtils.hasText(topic)) {
if (begin >= 0 && end >= 0 && end - begin >
MAX_TOPIC_QUERY_WINDOW_MILLIS) {
throw new BusinessException(400, "Topic message query time
range must not exceed 7 days");
}
- result = queryByTopic(instanceId, topic, tag, begin, end,
DEFAULT_TOPIC_LIMIT);
- } else {
- log.warn("queryMessages requires at least one of msgId/topic,
returning empty list");
- return Collections.emptyList();
+ return queryByTopic(instanceId, topic, tag, begin, end,
DEFAULT_TOPIC_LIMIT);
}
- return result;
+ log.warn("queryMessages requires at least one of msgId/topic,
returning empty list");
+ return Collections.emptyList();
}
private List<MessageRecordVO> queryByMsgId(DefaultMQAdminExt adminExt,
String topic, String msgId) {
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/auth/AuthCookieTest.java
b/server/src/test/java/org/apache/rocketmq/studio/auth/AuthCookieTest.java
index f85e04f90..3f53cbc41 100644
--- a/server/src/test/java/org/apache/rocketmq/studio/auth/AuthCookieTest.java
+++ b/server/src/test/java/org/apache/rocketmq/studio/auth/AuthCookieTest.java
@@ -28,19 +28,42 @@ import static org.assertj.core.api.Assertions.assertThat;
class AuthCookieTest {
@Test
- void usesHttpOnlySecureCookieAndPrefersItOverAuthorizationHeader() {
+ void usesHttpOnlySecureCookie() {
AuthProperties properties = new AuthProperties();
- MockHttpServletRequest request = new MockHttpServletRequest();
- request.setCookies(new Cookie("rmq_studio_session", "cookie-token"));
- request.addHeader("Authorization", "Bearer header-token");
MockHttpServletResponse response = new MockHttpServletResponse();
AuthCookie.write(response, properties, "cookie-token",
Duration.ofMinutes(30));
- assertThat(AuthCookie.authorization(request,
properties)).isEqualTo("Bearer cookie-token");
assertThat(response.getHeader("Set-Cookie"))
.contains("HttpOnly")
.contains("Secure")
.contains("SameSite=Strict");
}
+
+ @Test
+ void prefersExplicitBearerTokenOverStaleCookie() {
+ AuthProperties properties = new AuthProperties();
+ MockHttpServletRequest request = new MockHttpServletRequest();
+ request.setCookies(new Cookie("rmq_studio_session",
"stale-cookie-token"));
+ request.addHeader("Authorization", "bearer current-header-token");
+
+ assertThat(AuthCookie.authorization(request, properties))
+ .isEqualTo("bearer current-header-token");
+ }
+
+ @Test
+ void fallsBackToCookieForEmptyOrNonBearerAuthorization() {
+ AuthProperties properties = new AuthProperties();
+ MockHttpServletRequest emptyHeaderRequest = new
MockHttpServletRequest();
+ emptyHeaderRequest.setCookies(new Cookie("rmq_studio_session",
"cookie-token"));
+ emptyHeaderRequest.addHeader("Authorization", "Bearer ");
+ MockHttpServletRequest nonBearerRequest = new MockHttpServletRequest();
+ nonBearerRequest.setCookies(new Cookie("rmq_studio_session",
"cookie-token"));
+ nonBearerRequest.addHeader("Authorization", "Basic credentials");
+
+ assertThat(AuthCookie.authorization(emptyHeaderRequest, properties))
+ .isEqualTo("Bearer cookie-token");
+ assertThat(AuthCookie.authorization(nonBearerRequest, properties))
+ .isEqualTo("Bearer cookie-token");
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/auth/AuthServiceDatabaseTest.java
b/server/src/test/java/org/apache/rocketmq/studio/auth/AuthServiceDatabaseTest.java
index afe69409f..5ef2f9eb1 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/auth/AuthServiceDatabaseTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/auth/AuthServiceDatabaseTest.java
@@ -183,6 +183,21 @@ class AuthServiceDatabaseTest {
verify(userMapper).updateById(any(RmqStudioUser.class));
}
+ @Test
+ void settingTheCurrentEnabledStateShouldBeIdempotentTest() {
+ RmqStudioUser disabledAdmin = user(1L, "retired-admin", true, false,
"password-1");
+ RmqStudioUser enabledOperator = user(2L, "operator", false, true,
"password-1");
+ when(userMapper.selectById(1L)).thenReturn(disabledAdmin);
+ when(userMapper.selectById(2L)).thenReturn(enabledOperator);
+
+ assertThat(authService.setUserEnabled(1L,
false)).isSameAs(disabledAdmin);
+ assertThat(authService.setUserEnabled(2L,
true)).isSameAs(enabledOperator);
+
+ verify(userMapper, never()).selectList(any(Wrapper.class));
+ verify(userMapper, never()).updateById(any(RmqStudioUser.class));
+ verify(sessionMapper, never()).update(isNull(), any(Wrapper.class));
+ }
+
@Test
void databaseAuthenticationThrottlesLastSeenWrites() {
RmqStudioUser user = user(1L, "operator", false, true, "password-1");
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 fd3ff5b2e..86f2fb6f1 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
@@ -244,6 +244,47 @@ class MetricsServiceTest {
verifyNoInteractions(metricsSource);
}
+ @Test
+ void
queryShouldRejectReversedWindowWhenTimestampSubtractionWouldOverflow() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .metric("rocketmq_messages_in_total")
+ .start(Long.MAX_VALUE)
+ .end(Long.MIN_VALUE)
+ .step("1m")
+ .build();
+
+ assertBadRequest(query, "Metric query end must be later than start");
+ verifyNoInteractions(metricsSource);
+ }
+
+ @Test
+ void
queryShouldRejectOversizedWindowWhenTimestampSubtractionWouldOverflow() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .metric("rocketmq_messages_in_total")
+ .start(Long.MIN_VALUE)
+ .end(Long.MAX_VALUE)
+ .step("1h")
+ .build();
+
+ assertBadRequest(query, "Metric query range must not exceed 31 days");
+ verifyNoInteractions(metricsSource);
+ }
+
+ @Test
+ void queryShouldAcceptMaximumWindowNearTimestampUpperBound() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .metric("rocketmq_messages_in_total")
+ .start(Long.MAX_VALUE - 31L * 24 * 60 * 60)
+ .end(Long.MAX_VALUE)
+ .step("1h")
+ .build();
+ when(metricsSource.query(query)).thenReturn(emptyMetricData());
+
+ metricsService.query(query);
+
+ verify(metricsSource).query(query);
+ }
+
@Test
void queryShouldRejectOversizedWindow() {
MetricQueryDTO query = MetricQueryDTO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerConfigDiffServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerConfigDiffServiceTest.java
index 62344bd1d..9a1f632fa 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerConfigDiffServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerConfigDiffServiceTest.java
@@ -217,6 +217,40 @@ class NameServerConfigDiffServiceTest {
assertThat(result.getDifferences()).isEmpty();
}
+ @Test
+ void compareShouldDeduplicateDnsHostnamesIgnoringCase() throws Exception {
+ stubAdminFactory();
+ when(clusterService.getCluster("cluster-a")).thenReturn(cluster(
+ "ns.example.com:9876",
List.of(nameServer("NS.EXAMPLE.COM:9876"))));
+ when(admin.getNameServerConfig(List.of("NS.EXAMPLE.COM:9876")))
+ .thenReturn(Map.of("NS.EXAMPLE.COM:9876",
properties("listenPort", "9876")));
+
+ NameServerConfigDiffVO result = service.compare("cluster-a");
+
+ assertThat(result.getNodeCount()).isEqualTo(1);
+ assertThat(result.getReachableNodeCount()).isEqualTo(1);
+ verify(admin).getNameServerConfig(List.of("NS.EXAMPLE.COM:9876"));
+ }
+
+ @Test
+ void compareShouldPreserveIpv6ZoneCase() throws Exception {
+ stubAdminFactory();
+ String lowerZone = "[fe80::1%en0]:9876";
+ String upperZone = "[fe80::1%EN0]:9876";
+ when(clusterService.getCluster("cluster-a")).thenReturn(cluster(
+ upperZone, List.of(nameServer(lowerZone))));
+ when(admin.getNameServerConfig(List.of(lowerZone)))
+ .thenReturn(Map.of(lowerZone, properties("listenPort",
"9876")));
+ when(admin.getNameServerConfig(List.of(upperZone)))
+ .thenReturn(Map.of(upperZone, properties("listenPort",
"9876")));
+
+ NameServerConfigDiffVO result = service.compare("cluster-a");
+
+ assertThat(result.getNodeCount()).isEqualTo(2);
+ verify(admin).getNameServerConfig(List.of(lowerZone));
+ verify(admin).getNameServerConfig(List.of(upperZone));
+ }
+
@Test
void compareShouldMarkTheCheckIncompleteWhenEveryNodeIsUnavailable()
throws Exception {
stubAdminFactory();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/common/exception/GlobalExceptionHandlerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/common/exception/GlobalExceptionHandlerTest.java
index af2b16961..905da0d9b 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/common/exception/GlobalExceptionHandlerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/common/exception/GlobalExceptionHandlerTest.java
@@ -20,12 +20,16 @@ import org.apache.rocketmq.studio.common.domain.Result;
import org.apache.rocketmq.studio.ops.ai.LlmGatewayException;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.springframework.http.MediaType;
import org.springframework.test.web.servlet.MockMvc;
import org.springframework.test.web.servlet.setup.MockMvcBuilders;
+import org.springframework.web.HttpMediaTypeNotAcceptableException;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RestController;
+import java.util.List;
+
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;
@@ -88,6 +92,13 @@ class GlobalExceptionHandlerTest {
.andExpect(jsonPath("$.code").value(405));
}
+ @Test
+ void unacceptableResponseMediaTypeReturns406WithEnvelopeInsteadOf500()
throws Exception {
+ mockMvc.perform(get("/test/not-acceptable"))
+ .andExpect(status().isNotAcceptable())
+ .andExpect(jsonPath("$.code").value(406));
+ }
+
@RestController
static class FailingController {
@@ -101,5 +112,10 @@ class GlobalExceptionHandlerTest {
throw new LlmGatewayException(504, "llm.provider.timeout",
"LLM provider request timed out", "Retry later.");
}
+
+ @GetMapping("/test/not-acceptable")
+ Result<Void> failNotAcceptable() throws
HttpMediaTypeNotAcceptableException {
+ throw new
HttpMediaTypeNotAcceptableException(List.of(MediaType.APPLICATION_JSON));
+ }
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQServiceTest.java
index b80bdd16b..cf3b16917 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQServiceTest.java
@@ -89,6 +89,47 @@ class DLQServiceTest {
verify(dlqProvider).resendMessages("instance-1", "group-1", 1000L,
2000L, "target-topic");
}
+ @Test
+ void actionsShouldNormalizeIdentifiersBeforeDelegatingTest() {
+ List<String> msgIds = List.of(" msg-1 ", " msg-2 ");
+
+ dlqService.resendMessages("instance-1", " group-1 ", 1000L, 2000L, "
target-topic ");
+ dlqService.exportMessages("instance-1", " group-1 ", 1000L, 2000L,
100);
+ dlqService.listMessages("instance-1", " group-1 ", 1000L, 2000L, 1,
20);
+ dlqService.resendSelectedMessages("instance-1", " group-1 ", msgIds, "
target-topic ");
+ dlqService.exportExcel("instance-1", " group-1 ", 1000L, 2000L,
msgIds);
+
+ verify(dlqProvider).resendMessages(
+ "instance-1", "group-1", 1000L, 2000L, "target-topic");
+ verify(dlqProvider).exportMessages("instance-1", "group-1", 1000L,
2000L, 100);
+ verify(dlqProvider).listMessages("instance-1", "group-1", 1000L,
2000L, 1, 20);
+ verify(dlqProvider).resendMessages(
+ "instance-1", "group-1", List.of("msg-1", "msg-2"),
"target-topic");
+ verify(dlqProvider).exportExcel(
+ "instance-1", "group-1", 1000L, 2000L, List.of("msg-1",
"msg-2"));
+ }
+
+ @Test
+ void resendMessagesShouldTreatBlankTargetTopicAsAbsentTest() {
+ dlqService.resendMessages("instance-1", "group-1", 1000L, 2000L, "
");
+
+ verify(dlqProvider).resendMessages("instance-1", "group-1", 1000L,
2000L, null);
+ }
+
+ @Test
+ void selectedActionsShouldRejectBlankMsgIdsTest() {
+ assertThatThrownBy(() -> dlqService.resendSelectedMessages(
+ "instance-1", "group-1", List.of("msg-1", " "), null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("msgId must not be blank");
+ assertThatThrownBy(() -> dlqService.exportExcel(
+ "instance-1", "group-1", null, null, List.of("msg-1", " ")))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("msgId must not be blank");
+
+ verifyNoInteractions(dlqProvider);
+ }
+
@Test
void resendMessagesShouldAcceptNullTimeRange() {
dlqService.resendMessages("instance-1", "group-1", null, null,
"target-topic");
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
index 2936e2ebb..b0ce00c28 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
@@ -166,6 +166,19 @@ class AliyunCatalogServiceTest {
assertThat(byName.get(0).getInstanceName()).isEqualTo("Staging");
}
+ @Test
+ void listCloudInstancesShouldNormalizeRegionAndSearchTest() {
+ when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION), any()))
+
.thenReturn(instancesResponse(List.of(instanceRow("rmq-prod-001",
"Production"))));
+
+ List<CloudInstanceOptionVO> result = service.listCloudInstances(
+ CREDENTIAL_ID, " cn-hangzhou ", " production ");
+
+ assertThat(result).extracting(CloudInstanceOptionVO::getInstanceName)
+ .containsExactly("Production");
+ verify(clientFactory).call(eq(CREDENTIAL_ID), eq(REGION), any());
+ }
+
@Test
void listCloudInstancesShouldRequireRegionTest() {
assertThatThrownBy(() -> service.listCloudInstances(CREDENTIAL_ID, "
", null))
@@ -212,6 +225,17 @@ class AliyunCatalogServiceTest {
.isEqualTo("rmq-cn-001-vpc.rmq.aliyuncs.com:8080");
}
+ @Test
+ void getCloudInstanceShouldNormalizeLookupIdentifiersTest() {
+ when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION),
any())).thenReturn(null);
+
+ assertThatThrownBy(() -> service.getCloudInstance(
+ CREDENTIAL_ID, " cn-hangzhou ", " rmq-missing "))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Aliyun instance not found: rmq-missing");
+ verify(clientFactory).call(eq(CREDENTIAL_ID), eq(REGION), any());
+ }
+
private static List<ListInstancesResponseBody.List> instanceRows(int
count, int idOffset) {
List<ListInstancesResponseBody.List> rows = new ArrayList<>();
for (int i = 0; i < count; i++) {
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigServiceTest.java
index 71a6356d7..a46ed26ab 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigServiceTest.java
@@ -8,6 +8,7 @@ package org.apache.rocketmq.studio.provider.apache;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.common.domain.enums.FlushDiskType;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.ops.audit.AuditService;
import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
@@ -19,6 +20,7 @@ import org.mockito.junit.jupiter.MockitoExtension;
import java.util.Properties;
+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.anyString;
@@ -26,6 +28,7 @@ import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
class RocketMQBrokerConfigServiceTest {
@@ -87,4 +90,21 @@ class RocketMQBrokerConfigServiceTest {
"UPDATE_BROKER_CONFIG", "BROKER", "CLUSTER:cluster-a",
"cluster-a",
"brokerAddr=broker-a:10911, config={}", "SUCCESS");
}
+
+ @Test
+ void trimsFlushDiskTypeWithoutChangingUnknownValueFallback() throws
Exception {
+ Properties padded = new Properties();
+ padded.setProperty("flushDiskType", " SYNC_FLUSH ");
+ when(adminExt.getBrokerConfig("broker-a:10911")).thenReturn(padded);
+
+
assertThat(brokerConfigService.getBrokerConfig("broker-a:10911").getFlushDiskType())
+ .isEqualTo(FlushDiskType.SYNC_FLUSH);
+
+ Properties unknown = new Properties();
+ unknown.setProperty("flushDiskType", "future-mode");
+ when(adminExt.getBrokerConfig("broker-b:10911")).thenReturn(unknown);
+
+
assertThat(brokerConfigService.getBrokerConfig("broker-b:10911").getFlushDiskType())
+ .isEqualTo(FlushDiskType.ASYNC_FLUSH);
+ }
}
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 b8fb72876..91af89385 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
@@ -229,6 +229,23 @@ class RocketMQMessageProviderTest {
verify(clientApi).viewMessage("172.30.10.100:10911", "TopicA",
27521713L, 3000L);
}
+ @Test
+ void queryByMsgIdIgnoresUnrelatedTimeBounds() throws Exception {
+ MessageExt message = new MessageExt();
+ message.setMsgId("msg-1");
+ message.setTopic("TopicA");
+ when(adminExt.viewMessage("TopicA", "msg-1")).thenReturn(message);
+
+ assertThat(provider.queryMessages(
+ "instance-a", "TopicA", "msg-1", null, null, 200L, 100L))
+
.singleElement().extracting(MessageRecordVO::getMsgId).isEqualTo("msg-1");
+ assertThat(provider.queryMessages(
+ "instance-a", "TopicA", "msg-1", null, null, 100L, 100L))
+
.singleElement().extracting(MessageRecordVO::getMsgId).isEqualTo("msg-1");
+
+ verify(adminExt, times(2)).viewMessage("TopicA", "msg-1");
+ }
+
@Test
void queryByMsgIdRejectsDecodedBrokerOutsideKnownTopology() throws
Exception {
String msgId = MessageDecoder.createMessageId(new
InetSocketAddress("10.2.3.4", 10911), 12345L);