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 cf0efbbca feat(server): extend the metadata, ACL and admin layer for
the v10 tools (#4304)
cf0efbbca is described below
commit cf0efbbca7a0e946ed9fa38c3f534dfa5608dd9f
Author: lizhimins <[email protected]>
AuthorDate: Tue Sep 15 11:05:52 2026 +0800
feat(server): extend the metadata, ACL and admin layer for the v10 tools
(#4304)
- MetadataService: reject topic message type changes (a creation-only
attribute), add getTopicStats, redeliver through %RETRY%<group> with
the RocketMQ system-reserved properties filtered out, and cascade
%DLQ%<group> deletion when a consumer group is deleted
- AclService: guard ACL 2.0 by broker version and add get-by-id lookups
- RocketMQAdminClientImpl: filter system-reserved properties when
sending, map tag and key explicitly, select the FIFO queue by message
group hash, and set TIMER_DELIVER_MS for delayed messages
- NameServerConfigDiffService: add read() over a safe key allowlist
- RocketMQDefaultClusterResolver: advertise the configured admin
credential reference only when it exists, and fail closed with 422
instead of silently dropping ACL credentials during default-cluster
discovery
---
.../nameserver/NameServerConfigDiffService.java | 55 ++++++
.../studio/instance/acl/AclController.java | 3 +-
.../rocketmq/studio/instance/acl/AclService.java | 135 +++++++++++++-
.../studio/instance/message/MessageProvider.java | 7 +
.../instance/message/MessageProviderStub.java | 7 +
.../studio/instance/message/MessageService.java | 12 ++
.../studio/instance/topic/MetadataService.java | 140 ++++++++++++--
.../studio/instance/topic/SendMessageDTO.java | 6 +
...{SendMessageDTO.java => TopicQueueStatsVO.java} | 19 +-
.../provider/apache/RocketMQAdminClientImpl.java | 36 +++-
.../provider/apache/RocketMQDLQProvider.java | 50 ++++-
.../apache/RocketMQDefaultClusterResolver.java | 71 +++++--
.../provider/apache/RocketMQMessageProvider.java | 90 +++++++++
.../NameServerConfigDiffServiceTest.java | 62 +++++++
.../studio/instance/acl/AclControllerTest.java | 11 +-
.../studio/instance/acl/AclServiceTest.java | 155 ++++++++++++++--
.../instance/message/MessageServiceTest.java | 31 ++++
.../studio/instance/topic/MetadataServiceTest.java | 205 ++++++++++++++++++++-
.../apache/RocketMQAdminClientImplTest.java | 98 ++++++++++
.../apache/RocketMQClusterResolverTest.java | 21 +++
.../provider/apache/RocketMQDLQProviderTest.java | 28 +++
.../apache/RocketMQMessageProviderTest.java | 92 +++++++++
22 files changed, 1257 insertions(+), 77 deletions(-)
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 cdcc6ccbe..1cc7255a9 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
@@ -27,6 +27,7 @@ import
org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.tools.admin.MQAdminExt;
import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.io.IOException;
@@ -39,6 +40,7 @@ import java.util.Map;
import java.util.Properties;
import java.util.stream.Stream;
+@Slf4j
@Service
@RequiredArgsConstructor
public class NameServerConfigDiffService {
@@ -140,6 +142,59 @@ public class NameServerConfigDiffService {
.build();
}
+ /**
+ * Reads the safe configuration keys of every reachable management
NameServer endpoint for one
+ * physical cluster. Extracted from {@link #compare} so the read path can
back the read-only
+ * {@code rmq.nameserver.config} tool (decision 15); {@link #compare} is
retained for the REST
+ * diff view. {@code clusterId} must be the physical cluster name and
{@code instanceId} the
+ * Studio instance that owns it, so the cluster-details lookup resolves
the live topology
+ * (fixes the §15.5.5 "Cluster details are unavailable" path that keyed on
the instance id).
+ * Secret-bearing keys are never exposed: only {@link #SAFE_CONFIG_KEYS}
are returned.
+ */
+ public List<NodeConfig> read(String clusterId, String instanceId) {
+ String normalizedClusterId = requireClusterId(clusterId);
+ String normalizedInstanceId = normalizeInstanceId(instanceId);
+ ClusterVO cluster = normalizedInstanceId == null
+ ? clusterService.getCluster(normalizedClusterId)
+ : clusterService.getCluster(normalizedClusterId,
normalizedInstanceId);
+ List<String> addresses = collectNameServerAddresses(cluster);
+ if (addresses.isEmpty()) {
+ throw new BusinessException(409,
+ "Cluster has no NameServer endpoints: " +
normalizedClusterId);
+ }
+ String connectionEndpoint = connectionEndpoint(cluster, addresses);
+ List<NodeConfig> read = new ArrayList<>();
+ for (String address : addresses) {
+ try {
+ Properties config = readConfig(normalizedInstanceId,
connectionEndpoint, address);
+ read.add(new NodeConfig(address, safeConfig(config)));
+ } catch (BusinessException exception) {
+ log.warn("Skipping unreachable NameServer {} while reading
config for cluster {}: {}",
+ address, normalizedClusterId, exception.getMessage());
+ }
+ }
+ if (read.isEmpty()) {
+ throw new BusinessException(502,
+ "No reachable NameServer endpoint to read config from: " +
normalizedClusterId);
+ }
+ return read;
+ }
+
+ private Map<String, String> safeConfig(Properties config) {
+ Map<String, String> safe = new LinkedHashMap<>();
+ for (String key : SAFE_CONFIG_KEYS) {
+ String value = config.getProperty(key);
+ if (value != null) {
+ safe.put(key, value);
+ }
+ }
+ return safe;
+ }
+
+ /** One NameServer endpoint's safe configuration snapshot. */
+ public record NodeConfig(String addr, Map<String, String> config) {
+ }
+
private Properties readConfig(String instanceId, String
connectionEndpoint, String address) {
if (instanceId != null) {
return runtimeAdminClientResolver.execute(instanceId, admin ->
readConfig(admin, address));
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclController.java
index e7a2a87fc..de69c376e 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclController.java
@@ -61,11 +61,10 @@ public class AclController {
@RequestParam(required = false) String resource,
@RequestParam(required = false) String scope,
@RequestParam(required = false) String decision,
- @RequestParam(required = false) String aclVersion,
@RequestParam(required = false) String instanceId,
@RequestParam(defaultValue = "1") Integer page,
@RequestParam(defaultValue = "20") Integer pageSize) {
- return Result.ok(aclService.listRules(principal, resource, scope,
decision, aclVersion,
+ return Result.ok(aclService.listRules(principal, resource, scope,
decision,
instanceId, page, pageSize));
}
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 b9a9c967c..5f7d4a0c4 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
@@ -25,6 +25,8 @@ import
org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.util.CredentialUtils;
import org.apache.rocketmq.studio.common.util.EntityIds;
import org.apache.rocketmq.studio.audit.OperationAuditService;
+import org.apache.rocketmq.studio.cluster.broker.BrokerVO;
+import org.apache.rocketmq.studio.cluster.broker.ClusterProvider;
import org.apache.rocketmq.studio.model.Acl2PolicyContext;
import org.apache.rocketmq.studio.instance.InstanceResolver;
import org.apache.rocketmq.studio.instance.InstanceVO;
@@ -35,6 +37,7 @@ import org.springframework.stereotype.Service;
import java.security.SecureRandom;
import java.time.LocalDateTime;
+import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
import java.util.Locale;
@@ -48,10 +51,18 @@ public class AclService {
private static final int DEFAULT_RULE_PAGE_SIZE = 20;
private static final int MAX_PAGE_SIZE = 100;
+ /**
+ * Minimum broker version that supports ACL 2.0 (the RocketMQ {@code auth}
module with
+ * {@code authenticationEnabled}/{@code authorizationEnabled}). Below this
threshold the ACL
+ * tools return an upgrade hint instead of failing silently.
+ */
+ private static final int[] MIN_ACL2_BROKER_VERSION = {5, 3, 0};
+
private final AclRepository aclRepository;
private final OperationAuditService operationAuditService;
private final InstanceResolver instanceResolver;
private final TencentAclService tencentAclService;
+ private final ClusterProvider clusterProvider;
public AclCapabilitiesVO capabilities(String instanceId) {
if (!StringUtils.hasText(instanceId)) {
@@ -70,26 +81,27 @@ public class AclService {
public PageResult<AclRuleVO> listRules(String principal, String resource,
String scope, String decision,
- String aclVersion, String instanceId, Integer page, Integer
pageSize) {
+ String instanceId, Integer page, Integer pageSize) {
int normalizedPage = requireValidPage(page);
int normalizedPageSize = requireValidPageSize(pageSize);
+ requireAcl2Supported(instanceId);
if (isTencentInstance(instanceId)) {
List<AclRuleVO> filtered = tencentAclService.listRules(instanceId,
principal).stream()
.filter(rule -> containsIgnoreCase(rule.getResource(),
resource))
.filter(rule -> equalsIgnoreCase(rule.getScope(), scope))
.filter(rule -> equalsIgnoreCase(rule.getDecision(),
decision))
- .filter(rule -> equalsIgnoreCase(rule.getAclVersion(),
aclVersion))
.toList();
return paginateRules(filtered, normalizedPage, normalizedPageSize);
}
- log.info("Listing ACL rules for principal={}, resource={}, scope={},
decision={}, aclVersion={}, page={}, pageSize={}",
- principal, resource, scope, decision, aclVersion,
normalizedPage, normalizedPageSize);
- return aclRepository.findRulePage(principal, resource, scope,
decision, aclVersion,
+ log.info("Listing ACL rules for principal={}, resource={}, scope={},
decision={}, page={}, pageSize={}",
+ principal, resource, scope, decision, normalizedPage,
normalizedPageSize);
+ return aclRepository.findRulePage(principal, resource, scope,
decision, null,
normalizedPage, normalizedPageSize);
}
public AclRuleVO createRule(AclRuleVO rule, String instanceId) {
+ requireAcl2Supported(instanceId);
if (isTencentInstance(instanceId)) {
return tencentAclService.createRule(instanceId, rule);
}
@@ -133,6 +145,7 @@ public class AclService {
}
public AclRuleVO updateRule(AclRuleVO rule, String instanceId) {
+ requireAcl2Supported(instanceId);
if (isTencentInstance(instanceId)) {
return tencentAclService.updateRule(instanceId, rule);
}
@@ -147,6 +160,7 @@ public class AclService {
}
public void deleteRule(String id, String instanceId) {
+ requireAcl2Supported(instanceId);
if (isTencentInstance(instanceId)) {
tencentAclService.deleteRule(instanceId, id);
return;
@@ -198,6 +212,7 @@ public class AclService {
public AclUserVO createUser(AclUserVO user, String instanceId) {
+ requireAcl2Supported(instanceId);
if (isTencentInstance(instanceId)) {
return tencentAclService.createUser(instanceId, user);
}
@@ -214,6 +229,7 @@ public class AclService {
}
public AclUserVO updateUser(UpdateAclUserDTO user, String instanceId) {
+ requireAcl2Supported(instanceId);
if (isTencentInstance(instanceId)) {
return tencentAclService.updateUser(instanceId,
user.toAclUserVO());
}
@@ -244,6 +260,7 @@ public class AclService {
}
public void deleteUser(String id, String instanceId) {
+ requireAcl2Supported(instanceId);
if (isTencentInstance(instanceId)) {
// For Tencent roles the id is the role name.
tencentAclService.deleteUser(instanceId, id);
@@ -363,6 +380,114 @@ public class AclService {
.orElse(false);
}
+ /**
+ * Guards every ACL/user tool entrypoint so ACL 2.0 operations fail fast
with a clear upgrade
+ * hint when the backing broker is too old. Tencent instances are exempt
(they use role-based
+ * ACL rather than the Apache broker auth module). When the broker version
cannot be resolved
+ * the guard is permissive and lets the operation through to avoid false
positives.
+ */
+ private void requireAcl2Supported(String instanceId) {
+ if (isTencentInstance(instanceId)) {
+ return;
+ }
+ int[] detected = detectLowestBrokerVersion(instanceId);
+ if (detected == null) {
+ return;
+ }
+ if (compareVersions(detected, MIN_ACL2_BROKER_VERSION) < 0) {
+ throw new BusinessException(426,
+ "ACL 2.0 requires broker >= " +
formatVersion(MIN_ACL2_BROKER_VERSION)
+ + "; detected " + formatVersion(detected) + "
\u2014 please upgrade the broker");
+ }
+ }
+
+ private int[] detectLowestBrokerVersion(String instanceId) {
+ List<BrokerVO> brokers;
+ try {
+ brokers = clusterProvider.discoverBrokers(instanceId, null);
+ } catch (Exception discoveryFailure) {
+ log.debug("ACL 2.0 version guard could not discover brokers for
instance {}: {}",
+ instanceId, discoveryFailure.getMessage());
+ return null;
+ }
+ if (brokers == null || brokers.isEmpty()) {
+ return null;
+ }
+ int[] lowest = null;
+ for (BrokerVO broker : brokers) {
+ if (broker == null) {
+ continue;
+ }
+ int[] parsed = parseVersion(broker.getVersion());
+ if (parsed == null) {
+ continue;
+ }
+ if (lowest == null || compareVersions(parsed, lowest) < 0) {
+ lowest = parsed;
+ }
+ }
+ return lowest;
+ }
+
+ /**
+ * Normalizes a broker version descriptor into a comparable {@code [major,
minor, patch]} tuple.
+ * Handles both {@code MQVersion.getVersionDesc} forms (e.g. {@code
V5_3_3}) and plain semantic
+ * versions (e.g. {@code 5.3.0}). Returns {@code null} when no numeric
version can be extracted.
+ */
+ static int[] parseVersion(String raw) {
+ if (!StringUtils.hasText(raw)) {
+ return null;
+ }
+ String value = raw.trim();
+ int start = 0;
+ while (start < value.length() &&
!Character.isDigit(value.charAt(start))) {
+ start++;
+ }
+ if (start == value.length()) {
+ return null;
+ }
+ value = value.substring(start);
+ List<Integer> parts = new ArrayList<>(3);
+ StringBuilder digits = new StringBuilder();
+ for (int i = 0; i < value.length() && parts.size() < 3; i++) {
+ char character = value.charAt(i);
+ if (Character.isDigit(character)) {
+ digits.append(character);
+ } else if (digits.length() > 0) {
+ parts.add(Integer.parseInt(digits.toString()));
+ digits.setLength(0);
+ if (character != '.' && character != '_') {
+ break;
+ }
+ } else {
+ break;
+ }
+ }
+ if (digits.length() > 0 && parts.size() < 3) {
+ parts.add(Integer.parseInt(digits.toString()));
+ }
+ if (parts.isEmpty()) {
+ return null;
+ }
+ int major = parts.get(0);
+ int minor = parts.size() > 1 ? parts.get(1) : 0;
+ int patch = parts.size() > 2 ? parts.get(2) : 0;
+ return new int[] {major, minor, patch};
+ }
+
+ private static int compareVersions(int[] left, int[] right) {
+ for (int i = 0; i < 3; i++) {
+ if (left[i] != right[i]) {
+ return Integer.compare(left[i], right[i]);
+ }
+ }
+ return 0;
+ }
+
+ private static String formatVersion(int[] version) {
+ return version[0] + "." + version[1] + "." + version[2];
+ }
+
private boolean isValidAcl2BoundType(String boundType) {
return switch (boundType.trim().toUpperCase(Locale.ROOT)) {
case "TOPIC", "GROUP", "*", "USER", "SERVICE_ACCOUNT" -> true;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
index c77bf7306..b7ed4ec9b 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
@@ -30,6 +30,13 @@ public interface MessageProvider {
startTime, endTime));
}
+ /**
+ * Lookup by the client-generated unique key (UNIQ_KEY index). Without a
time window the
+ * provider defaults to a recent 3-day range; no match returns an empty
list.
+ */
+ List<MessageRecordVO> queryMessageByUniqueKey(String instanceId, String
topic, String uniqueKey,
+ Long startTime, Long
endTime);
+
TraceRecordVO getMessageTrace(String instanceId, String msgId, String
topic);
List<QueueOffsetVO> getQueueOffsets(String instanceId, String topic);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
index a29a98ef7..454dde173 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
@@ -35,6 +35,13 @@ public class MessageProviderStub implements MessageProvider {
throw unsupported();
}
+ @Override
+ public List<MessageRecordVO> queryMessageByUniqueKey(String instanceId,
String topic, String uniqueKey,
+ Long startTime, Long
endTime) {
+ log.warn("MessageProviderStub.queryMessageByUniqueKey called but no
real message provider is configured");
+ throw unsupported();
+ }
+
@Override
public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic) {
log.warn("MessageProviderStub.getMessageTrace called but no real
message provider is configured");
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 0e29caf91..250f2220c 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
@@ -62,6 +62,18 @@ public class MessageService {
return result;
}
+ public List<MessageRecordVO> queryMessageByUniqueKey(
+ String instanceId, String topic, String uniqueKey, Long startTime,
Long endTime) {
+ if (!StringUtils.hasText(topic)) {
+ throw new BusinessException(400, "topic is required");
+ }
+ if (!StringUtils.hasText(uniqueKey)) {
+ throw new BusinessException(400, "uniqueKey is required");
+ }
+ log.info("Querying message by unique key: topic={}, uniqueKey={}",
topic, uniqueKey);
+ return messageProvider.queryMessageByUniqueKey(instanceId, topic,
uniqueKey, startTime, endTime);
+ }
+
public MessageQueryPageVO queryMessagesPage(String instanceId, String
topic, String msgId, String tag,
String key, Long startTime,
Long endTime, int page, int pageSize) {
if (page < 1 || pageSize < 1 || pageSize > MAX_PAGE_SIZE) {
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 50e23f040..d45b4b513 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
@@ -16,10 +16,18 @@
*/
package org.apache.rocketmq.studio.instance.topic;
+import org.apache.rocketmq.common.MixAll;
+import org.apache.rocketmq.common.message.MessageConst;
+import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.remoting.protocol.ResponseCode;
+import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
+import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
import org.apache.rocketmq.studio.audit.OperationAuditConstants.Operation;
import org.apache.rocketmq.studio.audit.OperationAuditConstants.ResourceType;
import org.apache.rocketmq.studio.audit.OperationAuditConstants.Result;
import org.apache.rocketmq.studio.audit.OperationAuditService;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
+import org.apache.rocketmq.studio.common.util.MqResponseCodes;
import org.apache.rocketmq.studio.instance.InstanceResolver;
import org.apache.rocketmq.studio.provider.apache.AdminClient;
import org.apache.rocketmq.studio.provider.apache.ConsumerLagResolver;
@@ -49,7 +57,9 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
+import java.util.Comparator;
import java.util.HashSet;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
@@ -71,6 +81,7 @@ public class MetadataService {
private final InstanceResolver instanceResolver;
private final OperationAuditService operationAuditService;
private final MessageService messageService;
+ private final RuntimeAdminClientResolver runtimeAdminClientResolver;
/**
* Canonicalizes registered instance names and legacy numeric IDs, while
preserving physical
@@ -157,11 +168,24 @@ public class MetadataService {
public TopicVO updateTopic(String instanceId, TopicVO topic) {
requireTopic(topic);
topic.setInstanceId(instanceId);
+ guardImmutableType(instanceId, topic);
InstanceProvider provider = resolve(instanceId);
return executeWithAudit(provider, Operation.UPDATE_TOPIC,
ResourceType.TOPIC, topic.getName(),
instanceId, topicDetail(topic), () ->
provider.updateTopic(instanceId, topic));
}
+ /** The registered message type of an existing topic is immutable
(creation-only attribute). */
+ private void guardImmutableType(String instanceId, TopicVO topic) {
+ if (topic.getType() == null) {
+ return;
+ }
+ findTopic(instanceId, null, topic.getName()).ifPresent(existing -> {
+ if (topic.getType() != existing.getType()) {
+ throw new BusinessException(400, "topic message type is
immutable");
+ }
+ });
+ }
+
public void deleteTopic(String name) {
deleteTopic(null, name);
}
@@ -188,6 +212,48 @@ public class MetadataService {
return metadataProvider.getTopicRoutes(instanceId, topicName);
}
+ /**
+ * Per-queue offset stats (≈ admin topicStatus) read through the pooled
admin client.
+ * A missing topic route is an empty business state, not an RPC error;
other failures
+ * surface as 502 so callers can decide whether to degrade.
+ */
+ public List<TopicQueueStatsVO> getTopicStats(String instanceId, String
name) {
+ String target = normalizeInstanceId(instanceId);
+ String topicName = requireName(name, "topic name");
+ TopicStatsTable statsTable =
runtimeAdminClientResolver.execute(target, admin -> {
+ try {
+ return admin.examineTopicStats(topicName);
+ } catch (Exception e) {
+ if (MqResponseCodes.hasResponseCode(e,
ResponseCode.TOPIC_NOT_EXIST)) {
+ log.info("Topic {} has no broker route yet; returning
empty queue stats: {}",
+ topicName, e.getMessage());
+ return null;
+ }
+ throw e;
+ }
+ });
+ if (statsTable == null || statsTable.getOffsetTable() == null) {
+ return List.of();
+ }
+ return statsTable.getOffsetTable().entrySet().stream()
+ .filter(entry -> entry.getKey() != null && entry.getValue() !=
null)
+ .map(entry -> toQueueStats(entry.getKey(), entry.getValue()))
+ .sorted(Comparator.comparing(TopicQueueStatsVO::getBrokerName,
+
Comparator.nullsLast(Comparator.naturalOrder()))
+ .thenComparingInt(TopicQueueStatsVO::getQueueId))
+ .toList();
+ }
+
+ private static TopicQueueStatsVO toQueueStats(MessageQueue queue,
TopicOffset offset) {
+ return TopicQueueStatsVO.builder()
+ .brokerName(queue.getBrokerName())
+ .queueId(queue.getQueueId())
+ .minOffset(offset.getMinOffset())
+ .maxOffset(offset.getMaxOffset())
+ .lastUpdateTimestamp(offset.getLastUpdateTimestamp())
+ .build();
+ }
+
public List<TopicConsumerVO> getTopicConsumers(String name) {
return getTopicConsumers(null, name);
@@ -221,33 +287,62 @@ public class MetadataService {
}
/**
- * Re-publishes one stored message to a target topic. The original broker
message is read
- * through the instance-aware message service and only the application
payload/properties are
- * copied; broker offsets and delivery metadata are never reused.
+ * Re-publishes one stored message towards a consumer group. By default
the copy goes to
+ * {@code %RETRY%<groupName>} so only that group re-consumes it; an
explicit targetTopic
+ * overrides the destination (visible to all its subscribers) while
groupName still scopes
+ * audit and trace. System-reserved properties are never copied — tag/key
travel as
+ * first-class DTO fields instead.
*/
- public SendMessageVO resendMessage(String instanceId, String sourceTopic,
String msgId, String targetTopic) {
- MessageRecordVO original = findMessageForResend(instanceId,
sourceTopic, msgId);
- return resendMessage(instanceId, original, targetTopic);
- }
-
- public SendMessageVO resendMessage(String instanceId, MessageRecordVO
original, String targetTopic) {
- String destination = StringUtils.hasText(targetTopic) ?
targetTopic.trim() : original.getTopic();
- if (!StringUtils.hasText(destination)) {
- throw new BusinessException(400, "target topic is required when
the source message has no topic");
- }
+ public SendMessageVO redeliverMessage(String instanceId, String groupName,
String sourceTopic,
+ String msgId, String targetTopic) {
+ String group = requireName(groupName, "group name");
+ MessageRecordVO original = findMessageForRedelivery(instanceId,
sourceTopic, msgId);
+ return redeliverMessage(instanceId, group, original, targetTopic);
+ }
+
+ public SendMessageVO redeliverMessage(String instanceId, String groupName,
MessageRecordVO original,
+ String targetTopic) {
+ String group = requireName(groupName, "group name");
+ String destination = StringUtils.hasText(targetTopic)
+ ? targetTopic.trim()
+ : MixAll.getRetryTopic(group);
SendMessageDTO request = SendMessageDTO.builder()
.instanceId(normalizeInstanceId(instanceId))
.topic(destination)
.tag(original.getTag())
.key(original.getKey())
.body(original.getBody())
- .properties(original.getProperties() == null ? null :
Map.copyOf(original.getProperties()))
+ .properties(redeliveryProperties(original.getProperties()))
.build();
return sendMessage(request);
}
- /** Loads the exact source message used by resend without exposing or
publishing its body. */
- public MessageRecordVO findMessageForResend(String instanceId, String
sourceTopic, String msgId) {
+ /** Drops the system-reserved keys {@code Message.putUserProperty} would
reject (§15.5.1 KEYS defect). */
+ private static Map<String, String> redeliveryProperties(Map<String,
String> properties) {
+ if (properties == null || properties.isEmpty()) {
+ return null;
+ }
+ Map<String, String> filtered = new LinkedHashMap<>();
+ properties.forEach((key, value) -> {
+ if (!isSystemProperty(key)) {
+ filtered.put(key, value);
+ }
+ });
+ return filtered;
+ }
+
+ private static boolean isSystemProperty(String key) {
+ if (!StringUtils.hasText(key)) {
+ return true;
+ }
+ return MessageConst.STRING_HASH_SET.contains(key)
+ || key.startsWith("TIMER_")
+ || key.startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)
+ || key.startsWith(MixAll.DLQ_GROUP_TOPIC_PREFIX);
+ }
+
+ /** Loads the exact source message used by redelivery without exposing or
publishing its body. */
+ public MessageRecordVO findMessageForRedelivery(String instanceId, String
sourceTopic, String msgId) {
String messageId = requireName(msgId, "message id");
List<MessageRecordVO> matches = messageService.queryMessages(
instanceId, normalizeFilter(sourceTopic), messageId, null,
null, null, null);
@@ -420,6 +515,19 @@ public class MetadataService {
InstanceProvider provider = resolve(instanceId);
executeWithAudit(provider, Operation.DELETE_GROUP, ResourceType.GROUP,
groupName, instanceId, null, () -> mutation.accept(provider,
groupName));
+ cascadeDeleteDlqTopic(instanceId, groupName);
+ }
+
+ /** Best-effort DLQ cascade (decision 17): a missing or undeletable %DLQ%
topic never blocks group deletion. */
+ private void cascadeDeleteDlqTopic(String instanceId, String groupName) {
+ String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + groupName;
+ try {
+ deleteTopic(instanceId, dlqTopic);
+ log.info("Cascaded DLQ topic deletion for consumer group {}: {}",
groupName, dlqTopic);
+ } catch (Exception e) {
+ log.warn("Failed to cascade delete DLQ topic {} for consumer group
{}: {}",
+ dlqTopic, groupName, e.getMessage());
+ }
}
public void resetOffset(String name, long timestamp, String topic) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/SendMessageDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/SendMessageDTO.java
index 21f343a62..8680495d2 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/SendMessageDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/SendMessageDTO.java
@@ -37,4 +37,10 @@ public class SendMessageDTO {
private String key;
private String body;
private Map<String, String> properties;
+
+ /** FIFO sharding key; messages of the same group are sent to the same
queue. */
+ private String messageGroup;
+
+ /** Absolute delivery time in epoch milliseconds for DELAY (timer)
messages. */
+ private Long deliveryTimestamp;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/SendMessageDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicQueueStatsVO.java
similarity index 76%
copy from
server/src/main/java/org/apache/rocketmq/studio/instance/topic/SendMessageDTO.java
copy to
server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicQueueStatsVO.java
index 21f343a62..4857741dc 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/SendMessageDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicQueueStatsVO.java
@@ -16,25 +16,20 @@
*/
package org.apache.rocketmq.studio.instance.topic;
-import jakarta.validation.constraints.NotBlank;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
-import java.util.Map;
-
+/** Per-queue offset stats of one topic, as reported by broker-side topic
statistics. */
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
-public class SendMessageDTO {
- private String instanceId;
-
- @NotBlank(message = "topic is required")
- private String topic;
- private String tag;
- private String key;
- private String body;
- private Map<String, String> properties;
+public class TopicQueueStatsVO {
+ private String brokerName;
+ private int queueId;
+ private long minOffset;
+ private long maxOffset;
+ private long lastUpdateTimestamp;
}
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 e105ccfe4..2e3d0dbca 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
@@ -26,6 +26,7 @@ import org.apache.rocketmq.client.producer.SendStatus;
import org.apache.rocketmq.common.TopicConfig;
import org.apache.rocketmq.common.TopicAttributes;
import org.apache.rocketmq.common.message.Message;
+import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
import org.apache.rocketmq.remoting.protocol.ResponseCode;
@@ -472,16 +473,38 @@ public class RocketMQAdminClientImpl implements
AdminClient {
}
MqClientPool.ClientAction<DefaultMQProducer, SendMessageVO> sendAction
= producer -> {
- Message msg = new Message(topic, tag, key, bodyBytes);
+ Message msg = new Message(topic, bodyBytes);
+ if (StringUtils.hasText(tag)) {
+ msg.setTags(tag);
+ }
+ if (StringUtils.hasText(key)) {
+ msg.setKeys(key);
+ }
- // Add custom properties
+ // Custom properties must skip system-reserved keys
(KEYS/TAGS/UNIQ_KEY/WAIT/
+ // TIMER_*/RETRY_TOPIC/...): putUserProperty rejects them with
+ // "The Property<X> is used by system", which is how redelivering
a keyed
+ // message used to fail when the source properties were copied
verbatim.
if (request.getProperties() != null) {
for (Map.Entry<String, String> entry :
request.getProperties().entrySet()) {
+ if (isSystemReservedProperty(entry.getKey())) {
+ log.debug("Skipping system-reserved message property:
{}", entry.getKey());
+ continue;
+ }
msg.putUserProperty(entry.getKey(), entry.getValue());
}
}
- SendResult sendResult = producer.send(msg);
+ // Timer delivery is a system property; it must be set through the
typed API,
+ // never forwarded via user properties.
+ if (request.getDeliveryTimestamp() != null) {
+ msg.setDeliverTimeMs(request.getDeliveryTimestamp());
+ }
+
+ SendResult sendResult =
StringUtils.hasText(request.getMessageGroup())
+ ? producer.send(msg, (queues, message, arg) ->
+ queues.get(Math.floorMod(arg.hashCode(),
queues.size())), request.getMessageGroup())
+ : 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);
@@ -511,6 +534,13 @@ public class RocketMQAdminClientImpl implements
AdminClient {
}
}
+ private static boolean isSystemReservedProperty(String key) {
+ return key == null
+ || MessageConst.STRING_HASH_SET.contains(key)
+ || key.startsWith("%RETRY%")
+ || key.startsWith("%DLQ%");
+ }
+
@Override
public ConsumerGroupVO createConsumerGroup(ConsumerGroupVO group) {
String instanceId = group != null ? group.getInstanceId() : null;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
index ae6d5b2f9..cdb4b09fd 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
@@ -198,6 +198,17 @@ public class RocketMQDLQProvider implements DLQProvider {
recordAudit(groupName, detail, "FAILED");
throw e;
}
+ if (scanResult.topicMissing()) {
+ // The group has no %DLQ% topic yet (never dead-lettered a
message). Resending is a
+ // mutation against a non-existent target, so surface a clean
NOT_FOUND (legacy
+ // dlq.resend semantics) instead of silently reporting a
zero-message success.
+ String detail = String.format("instanceId=%s, group=%s,
dlqTopic=%s, targetTopic=%s, "
+ + "matched=0, resent=0, failed=0,
dlqTopicMissing=true",
+ instanceId, groupName, dlqTopic,
+ StringUtils.hasText(targetTopic) ? targetTopic :
"<original>");
+ recordAudit(groupName, detail, "NOT_FOUND");
+ throw new BusinessException(404, "No dead-letter queue found for
consumer group: " + groupName);
+ }
List<MessageExt> deadLetters = scanResult.messages();
int[] counts = {0, 0};
if (!deadLetters.isEmpty()) {
@@ -441,7 +452,7 @@ public class RocketMQDLQProvider implements DLQProvider {
try {
Set<MessageQueue> queues =
consumer.fetchSubscribeMessageQueues(dlqTopic);
if (queues == null || queues.isEmpty()) {
- return new DeadLetterScanResult(result, 0, false);
+ return new DeadLetterScanResult(result, 0, false, false);
}
outer:
for (MessageQueue queue : queues) {
@@ -514,10 +525,42 @@ public class RocketMQDLQProvider implements DLQProvider {
if (e instanceof BusinessException businessException) {
throw businessException;
}
+ if (isDlqTopicMissing(e)) {
+ // A group that has never exceeded its retry budget has no
%DLQ% topic yet; treat it
+ // as an empty dead-letter set instead of failing the scan
(matches legacy dlq.list,
+ // which returned items=[]/total=0 for a group without a DLQ).
Callers that mutate
+ // (redelivery_dlq) turn this flag into a clean NOT_FOUND
rather than a 502 crash.
+ log.info("DLQ topic {} has no route/queue yet; returning empty
scan result", dlqTopic);
+ return new DeadLetterScanResult(Collections.emptyList(), 0,
false, true);
+ }
log.warn("Failed to collect dead letters from {}: {}", dlqTopic,
e.getMessage());
throw new BusinessException(502, "Failed to scan DLQ topic " +
dlqTopic + ": " + e.getMessage());
}
- return new DeadLetterScanResult(result, failedQueueCount, truncated);
+ return new DeadLetterScanResult(result, failedQueueCount, truncated,
false);
+ }
+
+ /**
+ * True when the failure only means "the {@code %DLQ%} topic does not
exist / has no route or
+ * message queue yet" — the expected state for a consumer group that has
never dead-lettered a
+ * message. Any other cause is a real scan failure and must not be
silently degraded to empty.
+ */
+ private static boolean isDlqTopicMissing(Throwable e) {
+ Throwable cause = e;
+ while (cause != null) {
+ String message = cause.getMessage();
+ if (message != null) {
+ String lower = message.toLowerCase(Locale.ROOT);
+ if (lower.contains("can not find message queue")
+ || lower.contains("no topic route info")) {
+ return true;
+ }
+ }
+ if (cause.getCause() == cause) {
+ break;
+ }
+ cause = cause.getCause();
+ }
+ return false;
}
private boolean resendOne(DefaultMQProducer producer, MessageExt
deadLetter, String targetTopic) {
@@ -645,7 +688,8 @@ public class RocketMQDLQProvider implements DLQProvider {
}
}
- private record DeadLetterScanResult(List<MessageExt> messages, int
failedQueueCount, boolean truncated) {
+ private record DeadLetterScanResult(List<MessageExt> messages, int
failedQueueCount, boolean truncated,
+ boolean topicMissing) {
boolean scanIncomplete() {
return failedQueueCount > 0 || truncated;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDefaultClusterResolver.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDefaultClusterResolver.java
index bdd7b13cb..2bc64b9d2 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDefaultClusterResolver.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDefaultClusterResolver.java
@@ -35,6 +35,16 @@ import java.util.Optional;
@Component
@RequiredArgsConstructor
public class RocketMQDefaultClusterResolver {
+
+ /**
+ * Well-known key of the externally supplied admin credential that
protects the configured
+ * default cluster ({@code studio.cluster.admin.credentials.admin.*}).
+ *
+ * <p>The default cluster has no database record, so unlike a registered
instance it cannot
+ * carry a per-instance credential reference: a physical cluster name is
never a credential key.
+ */
+ private static final String DEFAULT_ADMIN_CREDENTIAL_REF = "admin";
+
private final RocketMQProperties properties;
private final MqAdminProperties adminProperties;
private final MqAdminExtFactory adminFactory;
@@ -50,37 +60,66 @@ public class RocketMQDefaultClusterResolver {
if (!StringUtils.hasText(properties.getNamesrvAddr())) {
return List.of();
}
+ // Discovery runs before any instance is resolved and must keep
working on deployments
+ // without ACL, so it falls back to an anonymous admin connection when
no default admin
+ // credential is configured. Every other entry point fails closed
instead.
return execute(admin -> {
var info = admin.examineBrokerClusterInfo();
return info == null || info.getClusterAddrTable() == null ?
List.of()
:
info.getClusterAddrTable().keySet().stream().sorted().toList();
- });
+ }, configuredAdminCredential());
}
public InstanceVO instance(String cluster) {
- if (!StringUtils.hasText(properties.getNamesrvAddr())) {
- throw new BusinessException(503, "RocketMQ admin not connected");
- }
+ String endpoint = requireEndpoint();
return
InstanceVO.builder().name(cluster).vendor(InstanceVendor.APACHE).type(InstanceType.DIRECT)
- .endpoint(properties.getNamesrvAddr().trim())
- .adminCredentialRef(cluster) // Use the cluster name as the
default credential reference.
+ .endpoint(endpoint)
+ // Advertise the configured default admin credential so
RuntimeAdminClientResolver
+ // authenticates with it; stay anonymous when ACL is not
configured at all.
+ .adminCredentialRef(configuredAdminCredential() == null ? null
: DEFAULT_ADMIN_CREDENTIAL_REF)
.build();
}
+ /**
+ * Runs an action against the configured default cluster using its
configured admin credential.
+ *
+ * <p>A configured ACL identity is never dropped silently: when the
default admin credential is
+ * absent or incomplete this fails with 422 instead of reconnecting
anonymously.
+ */
public <T> T execute(MqAdminExtFactory.AdminAction<T> action) {
- InstanceVO instance = instance(null);
- String reference = instance.getAdminCredentialRef();
- if (!StringUtils.hasText(reference)) {
- return adminFactory.execute(instance.getEndpoint(), null, action);
+ MqAdminProperties.Credential credential = configuredAdminCredential();
+ if (credential == null) {
+ throw new BusinessException(422,
+ "Admin credential reference is not configured: " +
DEFAULT_ADMIN_CREDENTIAL_REF);
}
- reference = reference.trim();
- MqAdminProperties.Credential credential =
adminProperties.getCredentials().get(reference);
- if (credential == null ||
!StringUtils.hasText(credential.getAccessKey())
- || !StringUtils.hasText(credential.getSecretKey())) {
- throw new BusinessException(422, "Admin credential reference is
not configured: " + reference);
+ return execute(action, credential);
+ }
+
+ private <T> T execute(MqAdminExtFactory.AdminAction<T> action,
MqAdminProperties.Credential credential) {
+ String endpoint = requireEndpoint();
+ if (credential == null) {
+ return adminFactory.execute(endpoint, null, action);
}
RPCHook hook = new AclClientRPCHook(new SessionCredentials(
credential.getAccessKey().trim(),
credential.getSecretKey().trim()));
- return adminFactory.execute(instance.getEndpoint(), hook, reference,
action);
+ return adminFactory.execute(endpoint, hook,
DEFAULT_ADMIN_CREDENTIAL_REF, action);
+ }
+
+ /** Returns the usable default admin credential, or {@code null} when ACL
is not configured. */
+ private MqAdminProperties.Credential configuredAdminCredential() {
+ MqAdminProperties.Credential credential =
adminProperties.getCredentials().get(DEFAULT_ADMIN_CREDENTIAL_REF);
+ if (credential == null ||
!StringUtils.hasText(credential.getAccessKey())
+ || !StringUtils.hasText(credential.getSecretKey())) {
+ return null;
+ }
+ return credential;
+ }
+
+ private String requireEndpoint() {
+ String namesrvAddr = properties.getNamesrvAddr();
+ if (!StringUtils.hasText(namesrvAddr)) {
+ throw new BusinessException(503, "RocketMQ admin not connected");
+ }
+ return namesrvAddr.trim();
}
}
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 c6d6066c6..640f500bd 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
@@ -21,6 +21,7 @@ import
org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
import org.apache.rocketmq.client.consumer.PullResult;
import org.apache.rocketmq.client.consumer.PullStatus;
import org.apache.rocketmq.client.trace.TraceConstants;
+import org.apache.rocketmq.common.MixAll;
import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageId;
@@ -59,6 +60,7 @@ import java.util.Collections;
import java.util.Comparator;
import java.util.Base64;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.PriorityQueue;
import java.util.Set;
@@ -76,6 +78,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
private static final String TRACE_TOPIC = "RMQ_SYS_TRACE_TOPIC";
private static final int KEY_QUERY_MAX = 64;
+ private static final int UNIQUE_KEY_QUERY_MAX = 1;
private static final int TRACE_QUERY_MAX = 64;
private static final int DEFAULT_TOPIC_LIMIT = 200;
private static final int TOPIC_QUERY_HARD_CAP = 2000;
@@ -86,6 +89,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
private static final long ONE_HOUR_MILLIS = 3600_000L;
private static final long ONE_DAY_MILLIS = 24 * ONE_HOUR_MILLIS;
private static final long MAX_TOPIC_QUERY_WINDOW_MILLIS = 7 *
ONE_DAY_MILLIS;
+ private static final long UNIQUE_KEY_DEFAULT_WINDOW_MILLIS = 3 *
ONE_DAY_MILLIS;
private static final int MAX_PULLS_PER_QUEUE = 32;
private static final int MAX_CONSECUTIVE_OFFSET_ILLEGAL = 3;
private static final int MAX_TOPIC_SCAN_MESSAGES_PER_QUEUE =
MAX_PULLS_PER_QUEUE * TOPIC_PULL_BATCH_SIZE;
@@ -199,6 +203,45 @@ public class RocketMQMessageProvider implements
MessageProvider {
}
}
+ @Override
+ public List<MessageRecordVO> queryMessageByUniqueKey(String instanceId,
String topic, String uniqueKey,
+ Long startTime, Long
endTime) {
+ return runtimeAdminClientResolver.execute(instanceId,
+ adminExt -> queryMessageByUniqueKey((DefaultMQAdminExt)
adminExt, topic, uniqueKey,
+ startTime, endTime));
+ }
+
+ private List<MessageRecordVO> queryMessageByUniqueKey(DefaultMQAdminExt
adminExt, String topic,
+ String uniqueKey,
Long startTime, Long endTime) {
+ try {
+ if (startTime == null && endTime == null) {
+ // Two-arg MQAdmin lookup: UNIQ_KEY index over a default
recent 3-day window.
+ MessageExt messageExt = adminExt.getDefaultMQAdminExtImpl()
+ .getMqClientInstance()
+ .getMQAdminImpl()
+ .queryMessageByUniqKey(topic, uniqueKey);
+ return messageExt == null ? Collections.emptyList() :
List.of(toRecordVO(messageExt));
+ }
+ long end = endTime != null ? endTime : System.currentTimeMillis();
+ long begin = startTime != null ? startTime : end -
UNIQUE_KEY_DEFAULT_WINDOW_MILLIS;
+ QueryResult queryResult = adminExt.queryMessageByUniqKey(null,
topic, uniqueKey,
+ UNIQUE_KEY_QUERY_MAX, begin, end);
+ if (queryResult == null || queryResult.getMessageList() == null
+ || queryResult.getMessageList().isEmpty()) {
+ return Collections.emptyList();
+ }
+ return
List.of(toRecordVO(queryResult.getMessageList().getFirst()));
+ } catch (Exception e) {
+ if (MqResponseCodes.hasResponseCode(e, ResponseCode.NO_MESSAGE,
ResponseCode.QUERY_NOT_FOUND)) {
+ // The index query completed but matched nothing: empty
result, not a gateway error.
+ log.info("queryMessageByUniqKey(topic={}, uniqueKey={})
matched nothing", topic, uniqueKey);
+ return Collections.emptyList();
+ }
+ log.warn("queryMessageByUniqKey(topic={}, uniqueKey={}) failed:
{}", topic, uniqueKey, e.getMessage());
+ throw new BusinessException(502, "Failed to query message by
unique key: " + e.getMessage());
+ }
+ }
+
@Override
public List<QueueOffsetVO> getQueueOffsets(String instanceId, String
topic) {
return runtimeAdminClientResolver.execute(instanceId, adminExt -> {
@@ -242,6 +285,11 @@ public class RocketMQMessageProvider implements
MessageProvider {
}
return toRecordVO(pullResult.getMsgFoundList().get(0),
brokerName);
} catch (Exception e) {
+ if (isRetryTopicReadBlocked(topic, e)) {
+ log.warn("pullMessageAtOffset(topic={}) skipped: reading a
%RETRY% topic requires a "
+ + "group-matched pull consumer; returning empty.
cause={}", topic, e.getMessage());
+ return null;
+ }
log.warn("pullMessageAtOffset(topic={}, broker={}, queue={},
offset={}) failed: {}",
topic, brokerName, queueId, offset, e.getMessage());
throw new BusinessException(502, "Failed to pull message at
offset: " + e.getMessage());
@@ -325,6 +373,16 @@ public class RocketMQMessageProvider implements
MessageProvider {
}
}
} catch (Exception e) {
+ if (isRetryTopicReadBlocked(topic, e)) {
+ // Reading a %RETRY%<group> topic through the shared
pooled pull consumer is
+ // rejected by broker ACL (the pull consumer group must
equal the retry topic's
+ // group). We deliberately do NOT spin up a group-matched
pull consumer: it would
+ // register as a member of that real group and take part
in its push-consumer
+ // rebalance, stalling queues. Degrade to an empty result
instead of failing.
+ log.warn("queryByTopic(topic={}) skipped: reading a
%RETRY% topic requires a "
+ + "group-matched pull consumer; returning empty.
cause={}", topic, e.getMessage());
+ return Collections.emptyList();
+ }
log.warn("queryByTopic(topic={}) failed: {}", topic,
e.getMessage());
throw new BusinessException(502, "Failed to query messages by
topic: " + e.getMessage());
}
@@ -334,6 +392,38 @@ public class RocketMQMessageProvider implements
MessageProvider {
});
}
+ /**
+ * True only when {@code topic} is a {@code %RETRY%<group>} system topic
and {@code e} carries one
+ * of the broker signals that the shared pooled pull consumer cannot read
it: the ACL rejection
+ * {@code retry topic does not match consumer group} (CODE:16, because the
pull consumer group must
+ * equal the retry topic's embedded group) or a missing route/queue.
Reading such a topic would
+ * require a group-matched pull consumer, which we intentionally avoid (it
would join that real
+ * group's rebalance); callers degrade to an empty result instead. Normal
topics and unrelated
+ * errors return false so genuine failures still surface.
+ */
+ private static boolean isRetryTopicReadBlocked(String topic, Throwable e) {
+ if (topic == null ||
!topic.startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
+ return false;
+ }
+ Throwable cause = e;
+ while (cause != null) {
+ String message = cause.getMessage();
+ if (message != null) {
+ String lower = message.toLowerCase(Locale.ROOT);
+ if (lower.contains("retry topic does not match consumer group")
+ || lower.contains("can not find message queue")
+ || lower.contains("no topic route info")) {
+ return true;
+ }
+ }
+ if (cause.getCause() == cause) {
+ break;
+ }
+ cause = cause.getCause();
+ }
+ return false;
+ }
+
private TopicQueueScanPlan buildTopicQueueScanPlan(DefaultMQPullConsumer
consumer, MessageQueue queue,
long begin, long end)
throws Exception {
long minOffset = consumer.minOffset(queue);
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 0f7fd7ab1..fb9f48256 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
@@ -299,6 +299,68 @@ class NameServerConfigDiffServiceTest {
.isEqualTo(400));
}
+ @Test
+ void readShouldReturnSafeConfigPerEndpointTest() throws Exception {
+ when(clusterService.getCluster("cluster-a",
"instance-a")).thenReturn(cluster(
+ "ns-a:9876;ns-b:9876",
+ List.of(nameServer("ns-a:9876"), nameServer("ns-b:9876"))));
+ when(runtimeAdminClientResolver.execute(eq("instance-a"),
any())).thenAnswer(invocation -> {
+ MqAdminExtFactory.AdminAction<Object> action =
invocation.getArgument(1);
+ return action.apply(admin);
+ });
+ when(admin.getNameServerConfig(List.of("ns-a:9876")))
+ .thenReturn(Map.of("ns-a:9876", properties(
+ "listenPort", "9876",
+ "serverWorkerThreads", "8",
+ "password", "secret")));
+ when(admin.getNameServerConfig(List.of("ns-b:9876")))
+ .thenReturn(Map.of("ns-b:9876", properties("listenPort",
"9876")));
+
+ List<NameServerConfigDiffService.NodeConfig> nodes =
service.read("cluster-a", "instance-a");
+
+ assertThat(nodes)
+ .extracting(NameServerConfigDiffService.NodeConfig::addr)
+ .containsExactly("ns-a:9876", "ns-b:9876");
+ assertThat(nodes.get(0).config())
+ .containsEntry("listenPort", "9876")
+ .containsEntry("serverWorkerThreads", "8")
+ .doesNotContainKey("password");
+ assertThat(nodes.get(1).config()).containsEntry("listenPort", "9876");
+ }
+
+ @Test
+ void readShouldSkipUnreachableEndpointsTest() throws Exception {
+ stubAdminFactory();
+ when(clusterService.getCluster("cluster-a")).thenReturn(cluster(
+ "ns-a:9876;ns-b:9876",
+ List.of(nameServer("ns-a:9876"), nameServer("ns-b:9876"))));
+ when(admin.getNameServerConfig(List.of("ns-a:9876")))
+ .thenReturn(Map.of("ns-a:9876", properties("listenPort",
"9876")));
+ when(admin.getNameServerConfig(List.of("ns-b:9876")))
+ .thenThrow(new IllegalStateException("unreachable"));
+
+ List<NameServerConfigDiffService.NodeConfig> nodes =
service.read("cluster-a", null);
+
+ assertThat(nodes).singleElement()
+ .extracting(NameServerConfigDiffService.NodeConfig::addr)
+ .isEqualTo("ns-a:9876");
+ }
+
+ @Test
+ void readShouldFailWhenNoEndpointIsReachableTest() throws Exception {
+ stubAdminFactory();
+ when(clusterService.getCluster("cluster-a")).thenReturn(cluster(
+ "ns-a:9876", List.of(nameServer("ns-a:9876"))));
+ when(admin.getNameServerConfig(List.of("ns-a:9876")))
+ .thenThrow(new IllegalStateException("unreachable"));
+
+ assertThatThrownBy(() -> service.read("cluster-a", null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("No reachable NameServer endpoint to read config
from: cluster-a")
+ .satisfies(exception -> assertThat(((BusinessException)
exception).getCode())
+ .isEqualTo(502));
+ }
+
private ClusterVO cluster(String endpoint, List<NameServerVO> nameServers)
{
ClusterVO cluster = ClusterVO.builder()
.name("cluster-a")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclControllerTest.java
index 7e38e7e0c..388d6e7de 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclControllerTest.java
@@ -109,7 +109,7 @@ class AclControllerTest extends WebMvcAuthTestSupport {
rule.setId(1L);
rule.setGmtCreate(LocalDateTime.of(2026, 1, 1, 0, 0));
- when(aclService.listRules(isNull(), isNull(), isNull(), isNull(),
isNull(), isNull(), eq(1), eq(20)))
+ when(aclService.listRules(isNull(), isNull(), isNull(), isNull(),
isNull(), eq(1), eq(20)))
.thenReturn(PageResult.of(java.util.List.of(rule), 1, 1, 20));
mockMvc.perform(get("/api/acl/rules"))
@@ -127,14 +127,13 @@ class AclControllerTest extends WebMvcAuthTestSupport {
@Test
void listRulesShouldPassQueryParams() throws Exception {
when(aclService.listRules(eq("user1"), eq("topic-a"), eq("namespace"),
eq("DENY"),
- eq("1.0"), isNull(), eq(3),
eq(5))).thenReturn(PageResult.empty(3, 5));
+ isNull(), eq(3), eq(5))).thenReturn(PageResult.empty(3, 5));
mockMvc.perform(get("/api/acl/rules")
.param("principal", "user1")
.param("resource", "topic-a")
.param("scope", "namespace")
.param("decision", "DENY")
- .param("aclVersion", "1.0")
.param("page", "3")
.param("pageSize", "5"))
.andExpect(status().isOk())
@@ -143,12 +142,12 @@ class AclControllerTest extends WebMvcAuthTestSupport {
.andExpect(jsonPath("$.data.size").value(5));
verify(aclService).listRules(eq("user1"), eq("topic-a"),
eq("namespace"), eq("DENY"),
- eq("1.0"), isNull(), eq(3), eq(5));
+ isNull(), eq(3), eq(5));
}
@Test
void listRulesShouldRejectPageSizeAboveTheInventoryLimit() throws
Exception {
- when(aclService.listRules(isNull(), isNull(), isNull(), isNull(),
isNull(), isNull(),
+ when(aclService.listRules(isNull(), isNull(), isNull(), isNull(),
isNull(),
eq(1), eq(101))).thenThrow(new BusinessException(400,
"page must be >= 1 and pageSize must be between 1 and 100"));
@@ -160,7 +159,7 @@ class AclControllerTest extends WebMvcAuthTestSupport {
.andExpect(jsonPath("$.message")
.value("page must be >= 1 and pageSize must be between
1 and 100"));
- verify(aclService).listRules(isNull(), isNull(), isNull(), isNull(),
isNull(), isNull(),
+ verify(aclService).listRules(isNull(), isNull(), isNull(), isNull(),
isNull(),
eq(1), eq(101));
}
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 2d7e25833..e5bddf384 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
@@ -23,6 +23,8 @@ import org.apache.rocketmq.studio.audit.OperationAuditService;
import org.apache.rocketmq.studio.model.Acl2PolicyContext;
import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
+import org.apache.rocketmq.studio.cluster.broker.BrokerVO;
+import org.apache.rocketmq.studio.cluster.broker.ClusterProvider;
import org.apache.rocketmq.studio.instance.InstanceResolver;
import org.apache.rocketmq.studio.instance.InstanceVO;
import org.apache.rocketmq.studio.provider.tencent.TencentAclService;
@@ -71,6 +73,9 @@ class AclServiceTest {
@Mock
private TencentAclService tencentAclService;
+ @Mock
+ private ClusterProvider clusterProvider;
+
@InjectMocks
private AclService aclService;
@@ -94,16 +99,16 @@ class AclServiceTest {
AclRuleVO.builder().principal("user1").resource("topic-1").decision("ALLOW").build(),
AclRuleVO.builder().principal("user2").resource("topic-2").decision("DENY").build()
);
- when(aclRepository.findRulePage("user1", "topic", "cluster", "ALLOW",
"2.0", 2, 5))
+ when(aclRepository.findRulePage("user1", "topic", "cluster", "ALLOW",
null, 2, 5))
.thenReturn(PageResult.of(rules, 12, 2, 5));
PageResult<AclRuleVO> result = aclService.listRules("user1", "topic",
"cluster",
- "ALLOW", "2.0", null, 2, 5);
+ "ALLOW", null, 2, 5);
assertThat(result.getItems()).hasSize(2);
assertThat(result.getItems().get(0).getPrincipal()).isEqualTo("user1");
assertThat(result.getTotal()).isEqualTo(12);
- verify(aclRepository).findRulePage("user1", "topic", "cluster",
"ALLOW", "2.0", 2, 5);
+ verify(aclRepository).findRulePage("user1", "topic", "cluster",
"ALLOW", null, 2, 5);
}
@Test
@@ -141,7 +146,7 @@ class AclServiceTest {
.thenReturn(PageResult.empty(1, 20));
PageResult<AclRuleVO> result = aclService.listRules(null, null, null,
null, null,
- null, null, null);
+ null, null);
assertThat(result.getItems()).isEmpty();
verify(aclRepository).findRulePage(null, null, null, null, null, 1,
20);
@@ -150,16 +155,16 @@ class AclServiceTest {
@Test
void listRulesShouldRejectInvalidPaginationBeforeQueryingRules() {
assertThatThrownBy(() -> aclService.listRules(null, null, null, null,
null,
- null, 0, 20))
+ 0, 20))
.isInstanceOf(BusinessException.class)
.hasMessage("page must be >= 1 and pageSize must be between 1
and 100")
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(400));
assertThatThrownBy(() -> aclService.listRules(null, null, null, null,
null,
- null, 1, 0))
+ 1, 0))
.isInstanceOf(BusinessException.class)
.hasMessage("page must be >= 1 and pageSize must be between 1
and 100");
assertThatThrownBy(() -> aclService.listRules(null, null, null, null,
null,
- null, 1, 101))
+ 1, 101))
.isInstanceOf(BusinessException.class)
.hasMessage("page must be >= 1 and pageSize must be between 1
and 100");
@@ -171,14 +176,14 @@ class AclServiceTest {
when(aclRepository.findRulePage(null, null, null, null, null, 1, 100))
.thenReturn(PageResult.empty(1, 100));
- aclService.listRules(null, null, null, null, null, null, 1, 100);
+ aclService.listRules(null, null, null, null, null, 1, 100);
verify(aclRepository).findRulePage(null, null, null, null, null, 1,
100);
}
@Test
void listRulesShouldRejectInvalidPaginationBeforeTencentRuleDiscovery() {
- assertThatThrownBy(() -> aclService.listRules(null, null, null, null,
null,
+ assertThatThrownBy(() -> aclService.listRules(null, null, null, null,
"tencent-instance", 1, 101))
.isInstanceOf(BusinessException.class)
.hasMessage("page must be >= 1 and pageSize must be between 1
and 100");
@@ -225,7 +230,7 @@ class AclServiceTest {
when(tencentAclService.listRules("tencent-instance",
null)).thenReturn(List.of(
AclRuleVO.builder().principal("role-a").resource("topic-a").build()));
- PageResult<AclRuleVO> result = aclService.listRules(null, null, null,
null, null,
+ PageResult<AclRuleVO> result = aclService.listRules(null, null, null,
null,
"tencent-instance", Integer.MAX_VALUE, 100);
assertThat(result.getItems()).isEmpty();
@@ -322,7 +327,7 @@ class AclServiceTest {
.isInstanceOf(BusinessException.class)
.hasMessage("ACL rule not found: 999")
.satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(404));
- assertThat(aclService.listRules(null, null, null, null, null, null, 1,
20).getItems()).isEmpty();
+ assertThat(aclService.listRules(null, null, null, null, null, 1,
20).getItems()).isEmpty();
verify(aclRepository, never()).saveRule(any(AclRuleVO.class));
}
@@ -974,4 +979,132 @@ class AclServiceTest {
.isInstanceOf(BusinessException.class)
.hasMessageContaining("whiteSet entry is not a valid IP/CIDR
range");
}
+
+ // ── ACL 2.0 broker version guard (§15) ─────────────────────────────
+
+ private InstanceVO apacheInstance(String identifier) {
+ InstanceVO instance = InstanceVO.builder()
+ .name(identifier)
+ .vendor(InstanceVendor.APACHE)
+ .type(InstanceType.DIRECT)
+ .build();
+ instance.setId(1L);
+ return instance;
+ }
+
+ @Test
+ void listRulesShouldRejectOldBrokerWithUpgradeHint() {
+ when(instanceResolver.findByIdentifier("apache-instance"))
+ .thenReturn(Optional.of(apacheInstance("apache-instance")));
+ when(clusterProvider.discoverBrokers("apache-instance", null))
+
.thenReturn(List.of(BrokerVO.builder().name("broker-a").version("V5_1_0").build()));
+
+ assertThatThrownBy(() -> aclService.listRules(null, null, null, null,
+ "apache-instance", 1, 20))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("ACL 2.0 requires broker >= 5.3.0")
+ .hasMessageContaining("please upgrade the broker")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(426));
+
+ verifyNoInteractions(aclRepository);
+ }
+
+ @Test
+ void createRuleShouldRejectOldBrokerWithUpgradeHint() {
+ when(instanceResolver.findByIdentifier("apache-instance"))
+ .thenReturn(Optional.of(apacheInstance("apache-instance")));
+ when(clusterProvider.discoverBrokers("apache-instance", null))
+
.thenReturn(List.of(BrokerVO.builder().name("broker-a").version("5.2.0").build()));
+
+ AclRuleVO input =
AclRuleVO.builder().principal("user1").resource("topic-1").build();
+
+ assertThatThrownBy(() -> aclService.createRule(input,
"apache-instance"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("please upgrade the broker")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(426));
+
+ verifyNoInteractions(aclRepository);
+ }
+
+ @Test
+ void listRulesShouldUseTheLowestBrokerVersionWhenSeveralBrokersExist() {
+ when(instanceResolver.findByIdentifier("apache-instance"))
+ .thenReturn(Optional.of(apacheInstance("apache-instance")));
+ when(clusterProvider.discoverBrokers("apache-instance", null))
+ .thenReturn(List.of(
+
BrokerVO.builder().name("broker-a").version("V5_3_1").build(),
+
BrokerVO.builder().name("broker-b").version("V5_2_9").build()));
+
+ assertThatThrownBy(() -> aclService.listRules(null, null, null, null,
+ "apache-instance", 1, 20))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("detected 5.2.9");
+ }
+
+ @Test
+ void listRulesShouldAllowModernBrokerVersion() {
+ when(instanceResolver.findByIdentifier("apache-instance"))
+ .thenReturn(Optional.of(apacheInstance("apache-instance")));
+ when(clusterProvider.discoverBrokers("apache-instance", null))
+
.thenReturn(List.of(BrokerVO.builder().name("broker-a").version("V5_3_0").build()));
+ when(aclRepository.findRulePage(null, null, null, null, null, 1, 20))
+ .thenReturn(PageResult.empty(1, 20));
+
+ PageResult<AclRuleVO> result = aclService.listRules(null, null, null,
null,
+ "apache-instance", 1, 20);
+
+ assertThat(result.getItems()).isEmpty();
+ verify(aclRepository).findRulePage(null, null, null, null, null, 1,
20);
+ }
+
+ @Test
+ void listRulesShouldPassWhenBrokerVersionCannotBeResolved() {
+ when(instanceResolver.findByIdentifier("apache-instance"))
+ .thenReturn(Optional.of(apacheInstance("apache-instance")));
+ when(clusterProvider.discoverBrokers("apache-instance", null))
+
.thenReturn(List.of(BrokerVO.builder().name("broker-a").version(null).build()));
+ when(aclRepository.findRulePage(null, null, null, null, null, 1, 20))
+ .thenReturn(PageResult.empty(1, 20));
+
+ PageResult<AclRuleVO> result = aclService.listRules(null, null, null,
null,
+ "apache-instance", 1, 20);
+
+ assertThat(result.getItems()).isEmpty();
+ verify(aclRepository).findRulePage(null, null, null, null, null, 1,
20);
+ }
+
+ @Test
+ void listRulesShouldSkipVersionGuardForTencentInstances() {
+ InstanceVO tencent = InstanceVO.builder()
+ .name("tencent-instance")
+ .vendor(InstanceVendor.TENCENT)
+ .type(InstanceType.CLOUD)
+ .build();
+
when(instanceResolver.findByIdentifier("tencent-instance")).thenReturn(Optional.of(tencent));
+ when(tencentAclService.listRules("tencent-instance",
null)).thenReturn(List.of(
+
AclRuleVO.builder().principal("role-a").resource("topic-a").build()));
+
+ PageResult<AclRuleVO> result = aclService.listRules(null, null, null,
null,
+ "tencent-instance", 1, 20);
+
+ assertThat(result.getItems()).hasSize(1);
+ verify(clusterProvider, never()).discoverBrokers(any(), any());
+ }
+
+ @ParameterizedTest
+ @MethodSource("brokerVersionDescriptors")
+ void parseVersionShouldNormalizeBrokerVersionDescriptors(String raw, int[]
expected) {
+ assertThat(AclService.parseVersion(raw)).isEqualTo(expected);
+ }
+
+ private static Stream<Arguments> brokerVersionDescriptors() {
+ return Stream.of(
+ Arguments.of("V5_3_3", new int[] {5, 3, 3}),
+ Arguments.of("5.3.0", new int[] {5, 3, 0}),
+ Arguments.of("V4_9_8", new int[] {4, 9, 8}),
+ Arguments.of("5.3", new int[] {5, 3, 0}),
+ Arguments.of(null, null),
+ Arguments.of("", null),
+ Arguments.of("unknown", null));
+ }
}
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 c2683a7de..e92108109 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
@@ -323,4 +323,35 @@ class MessageServiceTest {
verify(fallback).getMessageTraceByKey("instance-a", "ORDER-1",
"orders", "CUSTOM_TRACE");
verifyNoInteractions(history);
}
+
+ @Test
+ void rejectsBlankUniqueKeyQueryBeforeCallingProviderTest() {
+ MessageProvider provider = mock(MessageProvider.class);
+ MessageService service = new MessageService(provider,
mock(InstanceProviderRegistry.class),
+ mock(QueryHistoryService.class),
mock(OperationAuditService.class));
+
+ assertThatThrownBy(() -> service.queryMessageByUniqueKey("instance-a",
null, "uniq-1", null, null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("topic is required");
+ assertThatThrownBy(() -> service.queryMessageByUniqueKey("instance-a",
"TopicA", " ", null, null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("uniqueKey is required");
+
+ verifyNoInteractions(provider);
+ }
+
+ @Test
+ void delegatesUniqueKeyQueryToMessageProviderTest() {
+ MessageProvider provider = mock(MessageProvider.class);
+ MessageService service = new MessageService(provider,
mock(InstanceProviderRegistry.class),
+ mock(QueryHistoryService.class),
mock(OperationAuditService.class));
+ MessageRecordVO record =
MessageRecordVO.builder().msgId("msg-1").build();
+ when(provider.queryMessageByUniqueKey("instance-a", "TopicA",
"uniq-1", 100L, 200L))
+ .thenReturn(List.of(record));
+
+ assertThat(service.queryMessageByUniqueKey("instance-a", "TopicA",
"uniq-1", 100L, 200L))
+ .containsExactly(record);
+
+ verify(provider).queryMessageByUniqueKey("instance-a", "TopicA",
"uniq-1", 100L, 200L);
+ }
}
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 8bb752fe5..a86498582 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
@@ -17,7 +17,16 @@
package org.apache.rocketmq.studio.instance.topic;
+import org.apache.rocketmq.client.exception.MQClientException;
+import org.apache.rocketmq.common.message.MessageConst;
+import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.remoting.protocol.ResponseCode;
+import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
+import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
+import org.apache.rocketmq.tools.admin.MQAdminExt;
import org.apache.rocketmq.studio.audit.OperationAuditService;
+import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
@@ -49,15 +58,19 @@ import org.mockito.junit.jupiter.MockitoExtension;
import java.time.LocalDateTime;
import java.util.ArrayList;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatCode;
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.doThrow;
import static org.mockito.Mockito.lenient;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verifyNoMoreInteractions;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
@@ -90,11 +103,14 @@ class MetadataServiceTest {
@Mock
private MessageService messageService;
+ @Mock
+ private RuntimeAdminClientResolver runtimeAdminClientResolver;
+
@InjectMocks
private MetadataService metadataService;
@Test
- void
resendMessageShouldCopyApplicationPayloadWithoutExposingOriginalMetadata() {
+ void
redeliverMessageShouldCopyApplicationPayloadWithoutExposingOriginalMetadata() {
MessageRecordVO original = MessageRecordVO.builder()
.msgId("msg-original")
.topic("orders")
@@ -109,17 +125,87 @@ class MetadataServiceTest {
when(adminClient.sendMessage(any(SendMessageDTO.class)))
.thenReturn(SendMessageVO.builder().msgId("msg-new").build());
- SendMessageVO result = metadataService.resendMessage(
- "instance-a", "orders", "msg-original", "orders-retry");
+ SendMessageVO result = metadataService.redeliverMessage(
+ "instance-a", "group-a", "orders", "msg-original",
"orders-retry");
assertThat(result.getMsgId()).isEqualTo("msg-new");
ArgumentCaptor<SendMessageDTO> request =
ArgumentCaptor.forClass(SendMessageDTO.class);
verify(adminClient).sendMessage(request.capture());
assertThat(request.getValue().getTopic()).isEqualTo("orders-retry");
+ assertThat(request.getValue().getTag()).isEqualTo("paid");
+ assertThat(request.getValue().getKey()).isEqualTo("order-1");
assertThat(request.getValue().getBody()).isEqualTo("payload");
assertThat(request.getValue().getProperties()).containsEntry("tenant",
"alpha");
}
+ @Test
+ void redeliverMessageShouldDefaultToGroupRetryTopicTest() {
+ MessageRecordVO original = MessageRecordVO.builder()
+ .msgId("msg-original")
+ .topic("orders")
+ .body("payload")
+ .build();
+ when(messageService.queryMessages(
+ "instance-a", "orders", "msg-original", null, null, null,
null))
+ .thenReturn(List.of(original));
+ when(adminClient.sendMessage(any(SendMessageDTO.class)))
+ .thenReturn(SendMessageVO.builder().msgId("msg-new").build());
+
+ metadataService.redeliverMessage("instance-a", "group-a", "orders",
"msg-original", null);
+
+ ArgumentCaptor<SendMessageDTO> request =
ArgumentCaptor.forClass(SendMessageDTO.class);
+ verify(adminClient).sendMessage(request.capture());
+ assertThat(request.getValue().getTopic()).isEqualTo("%RETRY%group-a");
+ }
+
+ @Test
+ void redeliverMessageShouldFilterSystemReservedPropertiesTest() {
+ // §15.5.1: copying KEYS/TAGS/UNIQ_KEY verbatim makes putUserProperty
reject the send.
+ MessageRecordVO original = MessageRecordVO.builder()
+ .msgId("msg-original")
+ .topic("orders")
+ .tag("paid")
+ .key("order-1")
+ .body("payload")
+ .properties(Map.ofEntries(
+ Map.entry(MessageConst.PROPERTY_KEYS, "order-1"),
+ Map.entry(MessageConst.PROPERTY_TAGS, "paid"),
+
Map.entry(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX, "uniq"),
+ Map.entry(MessageConst.PROPERTY_WAIT_STORE_MSG_OK,
"true"),
+ Map.entry(MessageConst.PROPERTY_DELAY_TIME_LEVEL, "3"),
+ Map.entry(MessageConst.PROPERTY_RETRY_TOPIC, "orders"),
+ Map.entry(MessageConst.PROPERTY_TIMER_DELIVER_MS,
"1700000000000"),
+ Map.entry("%RETRY%group-a", "x"),
+ Map.entry("%DLQ%group-a", "y"),
+ Map.entry("tenant", "alpha")))
+ .build();
+ when(messageService.queryMessages(
+ "instance-a", "orders", "msg-original", null, null, null,
null))
+ .thenReturn(List.of(original));
+ when(adminClient.sendMessage(any(SendMessageDTO.class)))
+ .thenReturn(SendMessageVO.builder().msgId("msg-new").build());
+
+ metadataService.redeliverMessage("instance-a", "group-a", "orders",
"msg-original", "orders-copy");
+
+ ArgumentCaptor<SendMessageDTO> request =
ArgumentCaptor.forClass(SendMessageDTO.class);
+ verify(adminClient).sendMessage(request.capture());
+ assertThat(request.getValue().getProperties())
+ .containsExactlyInAnyOrderEntriesOf(Map.of("tenant", "alpha"));
+ assertThat(request.getValue().getTag()).isEqualTo("paid");
+ assertThat(request.getValue().getKey()).isEqualTo("order-1");
+ }
+
+ @Test
+ void redeliverMessageShouldRejectBlankGroupNameTest() {
+ assertThatThrownBy(() -> metadataService.redeliverMessage(
+ "instance-a", " ", "orders", "msg-original", null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("group name is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verifyNoInteractions(messageService, adminClient);
+ }
+
@Test
void skipAccumulatedShouldAdvanceEveryDistinctSubscription() {
when(apacheProvider.getGroupSubscriptions("instance-a", "group-a"))
@@ -251,6 +337,8 @@ class MetadataServiceTest {
verify(apacheProvider).createConsumerGroup("instance-a", group);
verify(apacheProvider).updateConsumerGroup("instance-a", group);
verify(apacheProvider).deleteConsumerGroup("instance-a", "consumers");
+ // decision 17: group deletion cascades to the dead-letter topic
+ verify(apacheProvider).deleteTopic("instance-a", "%DLQ%consumers");
verifyNoInteractions(operationAuditService);
}
@@ -451,6 +539,115 @@ class MetadataServiceTest {
verifyNoInteractions(apacheProvider);
}
+ @Test
+ void updateTopicShouldRejectMessageTypeChangeTest() {
+ // Decision 9: the registered topic type is immutable.
+ TopicVO existing = topic("orders", null, TopicType.NORMAL);
+ when(apacheProvider.listTopics("instance-a", null,
"orders")).thenReturn(List.of(existing));
+ TopicVO update = topic("orders", null, TopicType.FIFO);
+
+ assertThatThrownBy(() -> metadataService.updateTopic("instance-a",
update))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("topic message type is immutable")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verify(apacheProvider, never()).updateTopic(any(), any());
+ }
+
+ @Test
+ void updateTopicShouldAllowSameMessageTypeTest() {
+ TopicVO existing = topic("orders", null, TopicType.NORMAL);
+ when(apacheProvider.listTopics("instance-a", null,
"orders")).thenReturn(List.of(existing));
+ TopicVO update = topic("orders", null, TopicType.NORMAL);
+ when(apacheProvider.updateTopic("instance-a",
update)).thenReturn(update);
+
+ assertThat(metadataService.updateTopic("instance-a",
update)).isSameAs(update);
+ verify(apacheProvider).updateTopic("instance-a", update);
+ }
+
+ @Test
+ void getTopicStatsShouldMapAndSortOffsetTableTest() throws Exception {
+ MQAdminExt admin = mock(MQAdminExt.class);
+ TopicStatsTable table = new TopicStatsTable();
+ Map<MessageQueue, TopicOffset> offsetTable = new HashMap<>();
+ offsetTable.put(queue("broker-b", 0), offset(30, 40, 1700000000003L));
+ offsetTable.put(queue("broker-a", 1), offset(10, 20, 1700000000002L));
+ offsetTable.put(queue("broker-a", 0), offset(1, 2, 1700000000001L));
+ table.setOffsetTable(offsetTable);
+ when(admin.examineTopicStats("orders")).thenReturn(table);
+ stubAdminAction(admin);
+
+ List<TopicQueueStatsVO> stats =
metadataService.getTopicStats("instance-a", "orders");
+
+ assertThat(stats).containsExactly(
+ TopicQueueStatsVO.builder().brokerName("broker-a").queueId(0)
+
.minOffset(1).maxOffset(2).lastUpdateTimestamp(1700000000001L).build(),
+ TopicQueueStatsVO.builder().brokerName("broker-a").queueId(1)
+
.minOffset(10).maxOffset(20).lastUpdateTimestamp(1700000000002L).build(),
+ TopicQueueStatsVO.builder().brokerName("broker-b").queueId(0)
+
.minOffset(30).maxOffset(40).lastUpdateTimestamp(1700000000003L).build());
+ }
+
+ @Test
+ void getTopicStatsShouldReturnEmptyWhenRouteMissingTest() throws Exception
{
+ MQAdminExt admin = mock(MQAdminExt.class);
+ when(admin.examineTopicStats("orders")).thenThrow(
+ new MQClientException(ResponseCode.TOPIC_NOT_EXIST,
+ "No topic route info in name server for the topic:
orders"));
+ stubAdminAction(admin);
+
+ assertThat(metadataService.getTopicStats("instance-a",
"orders")).isEmpty();
+ }
+
+ @Test
+ void getTopicStatsShouldPropagateRpcFailureAsBadGatewayTest() throws
Exception {
+ MQAdminExt admin = mock(MQAdminExt.class);
+ when(admin.examineTopicStats("orders")).thenThrow(
+ new MQClientException(ResponseCode.SYSTEM_ERROR, "broker not
available"));
+ stubAdminAction(admin);
+
+ assertThatThrownBy(() -> metadataService.getTopicStats("instance-a",
"orders"))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(502));
+ }
+
+ @Test
+ void deleteConsumerGroupShouldNotFailWhenDlqDeletionFailsTest() {
+ // Decision 17: DLQ cascade is best-effort — a missing/undeletable DLQ
never blocks group deletion.
+ doThrow(new BusinessException(502, "no route for %DLQ%group-a"))
+ .when(apacheProvider).deleteTopic("instance-a",
"%DLQ%group-a");
+
+ assertThatCode(() -> metadataService.deleteConsumerGroup("instance-a",
"group-a"))
+ .doesNotThrowAnyException();
+
+ verify(apacheProvider).deleteConsumerGroup("instance-a", "group-a");
+ verify(apacheProvider).deleteTopic("instance-a", "%DLQ%group-a");
+ }
+
+ /** Runs the resolver action against the given admin mock, wrapping
failures like MqAdminExtFactory does. */
+ private void stubAdminAction(MQAdminExt admin) {
+ when(runtimeAdminClientResolver.execute(eq("instance-a"),
any())).thenAnswer(invocation -> {
+ MqAdminExtFactory.AdminAction<Object> action =
invocation.getArgument(1);
+ try {
+ return action.apply(admin);
+ } catch (Exception e) {
+ throw new BusinessException(502, "RocketMQ admin call failed:
" + e.getMessage());
+ }
+ });
+ }
+
+ private static MessageQueue queue(String brokerName, int queueId) {
+ return new MessageQueue("orders", brokerName, queueId);
+ }
+
+ private static TopicOffset offset(long min, long max, long lastUpdate) {
+ TopicOffset topicOffset = new TopicOffset();
+ topicOffset.setMinOffset(min);
+ topicOffset.setMaxOffset(max);
+ topicOffset.setLastUpdateTimestamp(lastUpdate);
+ return topicOffset;
+ }
+
@Test
void topicRuntimeDiagnosticsShouldDelegateWithSelectedInstance() {
BrokerRouteVO route =
BrokerRouteVO.builder().brokerName("broker-a").build();
@@ -642,6 +839,8 @@ class MetadataServiceTest {
"cloud-instance", "consumeType=-, subscriptionMode=-,
retryMaxTimes=16", "SUCCESS", null);
verify(operationAuditService).record("DELETE_GROUP", "GROUP",
"cg-orders",
"cloud-instance", null, "SUCCESS", null);
+ verify(operationAuditService).record("DELETE_TOPIC", "TOPIC",
"%DLQ%cg-orders",
+ "cloud-instance", null, "SUCCESS", null);
verify(operationAuditService).record("RESET_OFFSET", "GROUP",
"cg-orders",
"cloud-instance", "topic=orders, timestamp=1784246400000",
"SUCCESS", null);
verifyNoMoreInteractions(operationAuditService);
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 b900b60f8..b07bddb62 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,9 +19,11 @@ import
org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.remoting.protocol.body.ConsumerConnection;
import org.apache.rocketmq.common.TopicConfig;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
+import org.apache.rocketmq.client.producer.MessageQueueSelector;
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.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.remoting.exception.RemotingTimeoutException;
import org.apache.rocketmq.remoting.protocol.admin.ConsumeStats;
@@ -1350,4 +1352,100 @@ class RocketMQAdminClientImplTest {
verify(auditService).record(eq("SEND_MESSAGE"), eq("MESSAGE"),
eq("TopicA"),
eq(null), eq("Message send did not succeed: null"),
eq("FAILED"));
}
+
+ @Test
+ void sendMessageFiltersSystemReservedPropertiesTest() throws Exception {
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.SEND_OK);
+ sendResult.setMsgId("msg-1");
+ when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+ // Source properties copied verbatim from a stored message used to
fail the send with
+ // "The Property<KEYS> is used by system"; reserved keys must be
dropped instead.
+ Map<String, String> properties = new HashMap<>();
+ properties.put(MessageConst.PROPERTY_KEYS, "injected-keys");
+ properties.put(MessageConst.PROPERTY_TAGS, "injected-tags");
+ properties.put(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX,
"injected-uniq");
+ properties.put(MessageConst.PROPERTY_WAIT_STORE_MSG_OK, "false");
+ properties.put(MessageConst.PROPERTY_TIMER_DELIVER_MS, "123");
+ properties.put(MessageConst.PROPERTY_RETRY_TOPIC, "orders");
+ properties.put("%RETRY%group-a", "x");
+ properties.put("%DLQ%group-a", "y");
+ properties.put("bizKey", "bizValue");
+
+ SendMessageDTO request = new SendMessageDTO();
+ request.setTopic("TopicA");
+ request.setTag("tagA");
+ request.setKey("keyA");
+ request.setBody("hello");
+ request.setProperties(properties);
+
+ adminClient.sendMessage(request);
+
+ ArgumentCaptor<Message> captor =
ArgumentCaptor.forClass(Message.class);
+ verify(sendProducer).send(captor.capture());
+ Message sent = captor.getValue();
+ assertThat(sent.getTags()).isEqualTo("tagA");
+ assertThat(sent.getKeys()).isEqualTo("keyA");
+ assertThat(sent.getUserProperty("bizKey")).isEqualTo("bizValue");
+ assertThat(sent.getProperties())
+ .containsEntry(MessageConst.PROPERTY_KEYS, "keyA")
+ .containsEntry(MessageConst.PROPERTY_TAGS, "tagA")
+ .doesNotContainKeys(
+ MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX,
+ MessageConst.PROPERTY_RETRY_TOPIC,
+ MessageConst.PROPERTY_TIMER_DELIVER_MS,
+ "%RETRY%group-a",
+ "%DLQ%group-a");
+ }
+
+ @Test
+ void sendMessageSelectsQueueByMessageGroupHashTest() throws Exception {
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.SEND_OK);
+ sendResult.setMsgId("msg-1");
+ when(sendProducer.send(any(Message.class),
any(MessageQueueSelector.class), any()))
+ .thenReturn(sendResult);
+
+ SendMessageDTO request = new SendMessageDTO();
+ request.setTopic("TopicA");
+ request.setBody("hello");
+ request.setMessageGroup("group-x");
+
+ adminClient.sendMessage(request);
+
+ ArgumentCaptor<MessageQueueSelector> selectorCaptor =
+ ArgumentCaptor.forClass(MessageQueueSelector.class);
+ verify(sendProducer).send(any(Message.class),
selectorCaptor.capture(), eq("group-x"));
+ verify(sendProducer, never()).send(any(Message.class));
+ List<MessageQueue> queues = List.of(
+ new MessageQueue("TopicA", "broker-a", 0),
+ new MessageQueue("TopicA", "broker-a", 1),
+ new MessageQueue("TopicA", "broker-a", 2));
+ MessageQueue selected = selectorCaptor.getValue().select(queues, null,
"group-x");
+
assertThat(selected).isEqualTo(queues.get(Math.floorMod("group-x".hashCode(),
queues.size())));
+ // The same group must always land on the same queue.
+ assertThat(selectorCaptor.getValue().select(queues, null,
"group-x")).isEqualTo(selected);
+ }
+
+ @Test
+ void sendMessageSetsTimerDeliverMsForDeliveryTimestampTest() throws
Exception {
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.SEND_OK);
+ sendResult.setMsgId("msg-1");
+ when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+ SendMessageDTO request = new SendMessageDTO();
+ request.setTopic("TopicA");
+ request.setBody("hello");
+ request.setDeliveryTimestamp(1_900_000_000_000L);
+
+ adminClient.sendMessage(request);
+
+ ArgumentCaptor<Message> captor =
ArgumentCaptor.forClass(Message.class);
+ verify(sendProducer).send(captor.capture());
+
assertThat(captor.getValue().getDeliverTimeMs()).isEqualTo(1_900_000_000_000L);
+
assertThat(captor.getValue().getProperty(MessageConst.PROPERTY_TIMER_DELIVER_MS))
+ .isEqualTo("1900000000000");
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterResolverTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterResolverTest.java
index 62426d273..8f8136a12 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterResolverTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterResolverTest.java
@@ -82,4 +82,25 @@ class RocketMQClusterResolverTest {
credentials.getCredentials().clear();
assertThatThrownBy(() -> service.execute(admin ->
null)).isInstanceOf(BusinessException.class);
}
+
+ @Test
+ void instanceAdvertisesTheConfiguredAdminCredentialOnlyWhenItExistsTest() {
+ properties.setNamesrvAddr("configured:9876");
+ // ACL-less deployment: no credential reference, so runtime admin
calls stay anonymous
+ // instead of failing with "Admin credential reference is not
configured".
+
assertThat(service.instance("DefaultCluster").getAdminCredentialRef()).isNull();
+
+ MqAdminProperties.Credential credential = new
MqAdminProperties.Credential();
+ credential.setAccessKey("ak");
+ credential.setSecretKey("sk");
+ credentials.getCredentials().put("admin", credential);
+
+ assertThat(service.instance("DefaultCluster"))
+ .satisfies(instance -> {
+ assertThat(instance.getName()).isEqualTo("DefaultCluster");
+
assertThat(instance.getEndpoint()).isEqualTo("configured:9876");
+
assertThat(instance.getAdminCredentialRef()).isEqualTo("admin");
+ });
+ verifyNoInteractions(factory);
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
index e436f987b..5530bfca2 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.provider.apache;
import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
import org.apache.rocketmq.client.consumer.PullResult;
import org.apache.rocketmq.client.consumer.PullStatus;
+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;
@@ -474,6 +475,33 @@ class RocketMQDLQProviderTest {
eq("NO_MESSAGES"));
}
+ @Test
+ void listMessagesDegradesToEmptyWhenDlqTopicMissingTest() throws Exception
{
+ String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+ when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic))
+ .thenThrow(new MQClientException("Can not find Message Queue
for this topic, " + dlqTopic, null));
+
+ PageResult<DLQMessageVO> page = provider.listMessages("instance-a",
"group-a", 100L, 200L, 1, 20);
+
+ assertThat(page.getTotal()).isZero();
+ assertThat(page.getItems()).isEmpty();
+ verify(pullConsumer, never()).pull(any(MessageQueue.class),
anyString(), anyLong(), anyInt());
+ }
+
+ @Test
+ void resendMessagesThrowsNotFoundWhenDlqTopicMissingTest() throws
Exception {
+ String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+ when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic))
+ .thenThrow(new MQClientException("Can not find Message Queue
for this topic, " + dlqTopic, null));
+
+ assertThatThrownBy(() -> provider.resendMessages("instance-a",
"group-a", 100L, 200L, null))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(404));
+ verify(auditService).record(eq("RESEND_DLQ"), eq("DLQ"),
eq("group-a"), isNull(),
+ contains("dlqTopicMissing=true"), eq("NOT_FOUND"));
+ verify(runtimeAdminClientResolver,
never()).executeProducer(anyString(), any());
+ }
+
@Test
void resendMessagesUsesPooledClientsForScanAndResend() throws Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
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 7812c46f2..c794ed919 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
@@ -21,9 +21,11 @@ import
org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
import org.apache.rocketmq.client.consumer.PullResult;
import org.apache.rocketmq.client.consumer.PullStatus;
import org.apache.rocketmq.client.exception.MQClientException;
+import org.apache.rocketmq.client.impl.MQAdminImpl;
import org.apache.rocketmq.client.impl.MQClientAPIImpl;
import org.apache.rocketmq.client.impl.factory.MQClientInstance;
import org.apache.rocketmq.client.trace.TraceConstants;
+import org.apache.rocketmq.common.MixAll;
import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageId;
@@ -227,6 +229,80 @@ class RocketMQMessageProviderTest {
"instance-a", "TopicA", null, null, "order-1", 100L,
200L)).isEmpty();
}
+ @Test
+ void queryByUniqueKeyWithoutWindowUsesTwoArgAdminLookupTest() throws
Exception {
+ MQAdminImpl mqAdmin = mockUniqKeyLookupAdmin();
+ MessageExt message = new MessageExt();
+ message.setMsgId("msg-1");
+ message.setTopic("TopicA");
+ when(mqAdmin.queryMessageByUniqKey("TopicA",
"uniq-1")).thenReturn(message);
+
+ List<MessageRecordVO> result = provider.queryMessageByUniqueKey(
+ "instance-a", "TopicA", "uniq-1", null, null);
+
+ assertThat(result).hasSize(1);
+ assertThat(result.getFirst().getMsgId()).isEqualTo("msg-1");
+ verify(adminExt, never()).queryMessageByUniqKey(any(), anyString(),
anyString(),
+ anyInt(), anyLong(), anyLong());
+ }
+
+ @Test
+ void queryByUniqueKeyWithoutWindowReturnsEmptyWhenNoMatchTest() throws
Exception {
+ MQAdminImpl mqAdmin = mockUniqKeyLookupAdmin();
+ when(mqAdmin.queryMessageByUniqKey("TopicA",
"uniq-404")).thenReturn(null);
+
+ assertThat(provider.queryMessageByUniqueKey(
+ "instance-a", "TopicA", "uniq-404", null, null)).isEmpty();
+ }
+
+ @Test
+ void queryByUniqueKeyWithWindowUsesSixArgAdminLookupTest() throws
Exception {
+ MessageExt message = new MessageExt();
+ message.setMsgId("msg-9");
+ message.setTopic("TopicA");
+ when(adminExt.queryMessageByUniqKey(null, "TopicA", "uniq-1", 1, 100L,
200L))
+ .thenReturn(new QueryResult(0L, List.of(message)));
+
+ List<MessageRecordVO> result = provider.queryMessageByUniqueKey(
+ "instance-a", "TopicA", "uniq-1", 100L, 200L);
+
+ assertThat(result).hasSize(1);
+ assertThat(result.getFirst().getMsgId()).isEqualTo("msg-9");
+ verify(adminExt, never()).getDefaultMQAdminExtImpl();
+ }
+
+ @Test
+ void queryByUniqueKeyDegradesToEmptyWhenIndexHasNoMatchTest() throws
Exception {
+ when(adminExt.queryMessageByUniqKey(null, "TopicA", "uniq-404", 1,
100L, 200L))
+ .thenThrow(new MQClientException(ResponseCode.QUERY_NOT_FOUND,
+ "Can not find message"));
+
+ assertThat(provider.queryMessageByUniqueKey(
+ "instance-a", "TopicA", "uniq-404", 100L, 200L)).isEmpty();
+ }
+
+ @Test
+ void queryByUniqueKeySurfacesAdminFailureTest() throws Exception {
+ when(adminExt.queryMessageByUniqKey(null, "TopicA", "uniq-1", 1, 100L,
200L))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() -> provider.queryMessageByUniqueKey(
+ "instance-a", "TopicA", "uniq-1", 100L, 200L))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to query message by unique key: broker
unavailable")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
+ private MQAdminImpl mockUniqKeyLookupAdmin() {
+ DefaultMQAdminExtImpl adminExtImpl = mock(DefaultMQAdminExtImpl.class);
+ MQClientInstance clientInstance = mock(MQClientInstance.class);
+ MQAdminImpl mqAdmin = mock(MQAdminImpl.class);
+ when(adminExt.getDefaultMQAdminExtImpl()).thenReturn(adminExtImpl);
+ when(adminExtImpl.getMqClientInstance()).thenReturn(clientInstance);
+ when(clientInstance.getMQAdminImpl()).thenReturn(mqAdmin);
+ return mqAdmin;
+ }
+
@Test
void queryByMsgIdUsesDecodedPhysicalOffsetForFallback() throws Exception {
String msgId = "AC1E0A6400002A9F0000000001A3F2B1";
@@ -343,6 +419,22 @@ class RocketMQMessageProviderTest {
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
}
+ @Test
+ void queryByTopicDegradesToEmptyForRetryTopicReadBlockTest() throws
Exception {
+ String retryTopic = MixAll.RETRY_GROUP_TOPIC_PREFIX + "group-a";
+ MessageQueue queue = new MessageQueue(retryTopic, "broker-a", 0);
+
when(pullConsumer.fetchSubscribeMessageQueues(retryTopic)).thenReturn(Set.of(queue));
+ mockQueueWindow(pullConsumer, queue, 100L, 200L, 10L, 10L, 11L, 11L);
+ when(pullConsumer.pull(eq(queue), eq("*"), eq(10L), eq(32)))
+ .thenThrow(new MQClientException(
+ "CODE: 16 DESC: retry topic does not match consumer
group. BROKER: broker-a:10911", null));
+
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", retryTopic, null, null, null, 100L, 200L);
+
+ assertThat(messages).isEmpty();
+ }
+
@Test
@Timeout(value = 1, unit = TimeUnit.SECONDS)
void queryByTopicStopsWhenPullOffsetDoesNotAdvance() throws Exception {