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 e06c8878 feat: resolve runtime admin client for DLQ and client
connections (#1058)
e06c8878 is described below
commit e06c88784d553236986b21c6e57e9721a50f056a
Author: aias00 <[email protected]>
AuthorDate: Thu Aug 6 01:17:20 2026 -0700
feat: resolve runtime admin client for DLQ and client connections (#1058)
* fix: retain query history with request context
* feat: resolve runtime admin client for DLQ and client connections
---
.../cluster/broker/RuntimeAdminClientResolver.java | 38 ++++++++++
.../studio/cluster/client/ClientController.java | 3 +-
.../studio/cluster/client/ClientProvider.java | 2 +-
.../studio/cluster/client/ClientProviderStub.java | 6 +-
.../studio/cluster/client/ClientService.java | 6 +-
.../cluster/client/ProducerConnectionService.java | 2 +-
.../studio/instance/dlq/DLQController.java | 8 +--
.../rocketmq/studio/instance/dlq/DLQProvider.java | 5 +-
.../studio/instance/dlq/DLQProviderStub.java | 9 +--
.../studio/instance/dlq/DLQResendRequestDTO.java | 3 +
.../rocketmq/studio/instance/dlq/DLQService.java | 11 +--
.../studio/rocketmq/RocketMQClientProvider.java | 21 ++++--
.../studio/rocketmq/RocketMQDLQProvider.java | 65 +++++++----------
.../broker/RuntimeAdminClientResolverTest.java | 83 ++++++++++++++++++++++
.../cluster/client/ClientControllerTest.java | 11 +--
.../cluster/client/ClientProviderStubTest.java | 2 +-
.../studio/cluster/client/ClientServiceTest.java | 12 ++--
.../client/ProducerConnectionServiceTest.java | 4 +-
.../studio/instance/dlq/DLQControllerTest.java | 20 +++---
.../studio/instance/dlq/DLQProviderStubTest.java | 4 +-
.../studio/instance/dlq/DLQServiceTest.java | 28 ++++----
.../rocketmq/RocketMQClientProviderTest.java | 25 ++++---
.../studio/rocketmq/RocketMQDLQProviderTest.java | 16 ++---
web/src/api/connections.test.ts | 12 +++-
web/src/api/connections.ts | 1 +
web/src/api/dlq.test.ts | 7 +-
web/src/api/message.ts | 5 +-
.../pages/cluster/__tests__/ClientsPage.test.tsx | 28 +++++++-
web/src/pages/cluster/clients.tsx | 33 ++++++++-
web/src/pages/instance/__tests__/DLQPage.test.tsx | 18 ++++-
web/src/pages/instance/dlq.tsx | 20 +++++-
web/src/services/connectionsService.test.ts | 12 +++-
web/src/services/messageService.test.ts | 8 +--
web/src/services/messageService.ts | 5 +-
34 files changed, 382 insertions(+), 151 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
new file mode 100644
index 00000000..58fd473b
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
@@ -0,0 +1,38 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0.
+ */
+package org.apache.rocketmq.studio.cluster.broker;
+
+import lombok.RequiredArgsConstructor;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.instance.InstanceRepository;
+import org.apache.rocketmq.studio.instance.InstanceVO;
+import org.springframework.stereotype.Component;
+import org.springframework.util.StringUtils;
+
+@Component
+@RequiredArgsConstructor
+public class RuntimeAdminClientResolver {
+
+ private final InstanceRepository instanceRepository;
+ private final MqAdminExtFactory adminFactory;
+
+ public String resolveEndpoint(String instanceId) {
+ if (!StringUtils.hasText(instanceId)) {
+ throw new BusinessException(400, "instanceId is required");
+ }
+ InstanceVO instance = instanceRepository.findById(instanceId)
+ .orElseThrow(() -> new BusinessException(404, "Instance not
found: " + instanceId));
+ if (!StringUtils.hasText(instance.getEndpoint())) {
+ throw new BusinessException(400, "Instance has no endpoint: " +
instanceId);
+ }
+ return instance.getEndpoint().trim();
+ }
+
+ public <T> T execute(String instanceId, MqAdminExtFactory.AdminAction<T>
action) {
+ return adminFactory.execute(resolveEndpoint(instanceId), null, action);
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientController.java
index dd83c9f1..3f81d541 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientController.java
@@ -34,8 +34,9 @@ public class ClientController {
@GetMapping
public Result<List<ClientConnectionVO>> listConnections(
+ @RequestParam String instanceId,
@RequestParam(required = false) String clusterId,
@RequestParam(required = false) String type) {
- return Result.ok(clientService.listConnections(clusterId, type));
+ return Result.ok(clientService.listConnections(instanceId, clusterId,
type));
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
index 86a58c50..8011ee8b 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProvider.java
@@ -20,7 +20,7 @@ package org.apache.rocketmq.studio.cluster.client;
import java.util.List;
public interface ClientProvider {
- List<ClientConnectionVO> findConnections(String clusterId, String type);
+ List<ClientConnectionVO> findConnections(String instanceId, String
clusterId, String type);
List<ClientConnectionVO> findProducerConnections(String topic, String
producerGroup);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStub.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStub.java
index c5bb198a..bec9b9a7 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStub.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStub.java
@@ -25,9 +25,9 @@ import java.util.List;
public class ClientProviderStub implements ClientProvider {
@Override
- public List<ClientConnectionVO> findConnections(String clusterId, String
type) {
- log.warn("ClientProviderStub.findConnections called without a real
client provider. clusterId={}, type={}",
- clusterId, type);
+ public List<ClientConnectionVO> findConnections(String instanceId, String
clusterId, String type) {
+ log.warn("ClientProviderStub.findConnections called without a real
client provider. instanceId={}, clusterId={}, type={}",
+ instanceId, clusterId, type);
throw new BusinessException(501, "Client connection provider is not
configured");
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientService.java
index 7f39325a..ed9723e4 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ClientService.java
@@ -29,9 +29,9 @@ public class ClientService {
private final ClientProvider clientProvider;
- public List<ClientConnectionVO> listConnections(String clusterId, String
type) {
- log.info("Listing client connections, clusterId={}, type={}",
clusterId, type);
- return clientProvider.findConnections(normalizeFilter(clusterId),
normalizeFilter(type));
+ public List<ClientConnectionVO> listConnections(String instanceId, String
clusterId, String type) {
+ log.info("Listing client connections, instanceId={}, clusterId={},
type={}", instanceId, clusterId, type);
+ return clientProvider.findConnections(normalizeFilter(instanceId),
normalizeFilter(clusterId), normalizeFilter(type));
}
private String normalizeFilter(String value) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
index d13227ad..41d6e6e0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionService.java
@@ -41,7 +41,7 @@ public class ProducerConnectionService {
}
public List<String> listProducerGroups() {
- return clientProvider.findConnections(null,
ClientType.Producer.name()).stream()
+ return clientProvider.findConnections(null, null,
ClientType.Producer.name()).stream()
.map(ClientConnectionVO::getProducerGroup)
.filter(this::hasText)
.map(String::trim)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQController.java
index 872cdd6a..096b3b87 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQController.java
@@ -37,15 +37,15 @@ public class DLQController {
private final DLQService dlqService;
@GetMapping
- public Result<List<DLQGroupVO>> listDLQGroups(@RequestParam(required =
false) String clusterId) {
- return Result.ok(dlqService.listDLQGroups(clusterId));
+ public Result<List<DLQGroupVO>> listDLQGroups(@RequestParam String
instanceId) {
+ return Result.ok(dlqService.listDLQGroups(instanceId));
}
@PostMapping("/resend")
public Result<DLQResendResultVO> resendMessages(@Valid
@RequestBody(required = false) DLQResendRequestDTO request) {
requireRequest(request);
- return Result.ok(dlqService.resendMessages(
- request.getGroupName(), request.getStartTime(),
request.getEndTime(), request.getTargetTopic()));
+ return Result.ok(dlqService.resendMessages(request.getInstanceId(),
request.getGroupName(),
+ request.getStartTime(), request.getEndTime(),
request.getTargetTopic()));
}
private void requireRequest(DLQResendRequestDTO request) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
index f143a68c..e91cf484 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
@@ -20,6 +20,7 @@ package org.apache.rocketmq.studio.instance.dlq;
import java.util.List;
public interface DLQProvider {
- List<DLQGroupVO> listDLQGroups(String clusterId);
- DLQResendResultVO resendMessages(String groupName, Long startTime, Long
endTime, String targetTopic);
+ List<DLQGroupVO> listDLQGroups(String instanceId);
+ DLQResendResultVO resendMessages(String instanceId, String groupName, Long
startTime, Long endTime,
+ String targetTopic);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
index ad6e3bb3..2011a966 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
@@ -29,14 +29,15 @@ import java.util.List;
public class DLQProviderStub implements DLQProvider {
@Override
- public List<DLQGroupVO> listDLQGroups(String clusterId) {
- log.warn("DLQProviderStub.listDLQGroups called but no real DLQ
provider is configured. clusterId={}",
- clusterId);
+ public List<DLQGroupVO> listDLQGroups(String instanceId) {
+ log.warn("DLQProviderStub.listDLQGroups called but no real DLQ
provider is configured. instanceId={}",
+ instanceId);
throw unsupported();
}
@Override
- public DLQResendResultVO resendMessages(String groupName, Long startTime,
Long endTime, String targetTopic) {
+ public DLQResendResultVO resendMessages(String instanceId, String
groupName, Long startTime, Long endTime,
+ String targetTopic) {
log.warn("DLQProviderStub.resendMessages called but no real DLQ
provider is configured. group={}, targetTopic={}",
groupName, targetTopic);
throw unsupported();
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendRequestDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendRequestDTO.java
index b91e1f48..f99100d0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendRequestDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendRequestDTO.java
@@ -27,6 +27,9 @@ import lombok.NoArgsConstructor;
@NoArgsConstructor
@AllArgsConstructor
public class DLQResendRequestDTO {
+ @NotBlank(message = "instanceId is required")
+ private String instanceId;
+
@NotBlank(message = "groupName is required")
private String groupName;
private Long startTime;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
index a18e4e5a..a5a6baaf 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
@@ -32,15 +32,16 @@ public class DLQService {
private final DLQProvider dlqProvider;
- public List<DLQGroupVO> listDLQGroups(String clusterId) {
- log.info("Listing DLQ groups for cluster: {}", clusterId);
- return dlqProvider.listDLQGroups(clusterId);
+ public List<DLQGroupVO> listDLQGroups(String instanceId) {
+ log.info("Listing DLQ groups for instance: {}", instanceId);
+ return dlqProvider.listDLQGroups(instanceId);
}
- public DLQResendResultVO resendMessages(String groupName, Long startTime,
Long endTime, String targetTopic) {
+ public DLQResendResultVO resendMessages(String instanceId, String
groupName, Long startTime, Long endTime,
+ String targetTopic) {
validateResendRequest(groupName, startTime, endTime);
log.info("Resending DLQ messages: group={}, targetTopic={}",
groupName, targetTopic);
- return dlqProvider.resendMessages(groupName, startTime, endTime,
targetTopic);
+ return dlqProvider.resendMessages(instanceId, groupName, startTime,
endTime, targetTopic);
}
private void validateResendRequest(String groupName, Long startTime, Long
endTime) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQClientProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQClientProvider.java
index 7d5cf55a..48760cf6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQClientProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQClientProvider.java
@@ -29,11 +29,13 @@ import
org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import org.apache.rocketmq.studio.cluster.client.ClientConnectionVO;
import org.apache.rocketmq.studio.cluster.client.ClientProvider;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.ClientLanguage;
import org.apache.rocketmq.studio.common.domain.enums.ClientType;
import org.apache.rocketmq.studio.common.domain.enums.Protocol;
import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
+import org.apache.rocketmq.tools.admin.MQAdminExt;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.context.annotation.Primary;
@@ -61,15 +63,20 @@ public class RocketMQClientProvider implements
ClientProvider {
private static final long SUBSCRIPTION_GROUP_TIMEOUT_MILLIS = 5000L;
private final ObjectProvider<DefaultMQAdminExt> adminExtProvider;
+ private final RuntimeAdminClientResolver runtimeAdminClientResolver;
- public RocketMQClientProvider(ObjectProvider<DefaultMQAdminExt>
adminExtProvider) {
+ public RocketMQClientProvider(ObjectProvider<DefaultMQAdminExt>
adminExtProvider,
+ RuntimeAdminClientResolver
runtimeAdminClientResolver) {
this.adminExtProvider = adminExtProvider;
+ this.runtimeAdminClientResolver = runtimeAdminClientResolver;
}
@Override
- public List<ClientConnectionVO> findConnections(String clusterId, String
type) {
- DefaultMQAdminExt adminExt = requireAdminExt();
+ public List<ClientConnectionVO> findConnections(String instanceId, String
clusterId, String type) {
+ return runtimeAdminClientResolver.execute(instanceId, adminExt ->
findConnections(adminExt, clusterId, type));
+ }
+ private List<ClientConnectionVO> findConnections(MQAdminExt adminExt,
String clusterId, String type) {
ClientType clientType = parseType(type);
List<ClientConnectionVO> connections = new ArrayList<>();
if (clientType == null || clientType == ClientType.Producer) {
@@ -119,7 +126,7 @@ public class RocketMQClientProvider implements
ClientProvider {
}
private List<ClientConnectionVO> findAllProducerConnections(
- DefaultMQAdminExt adminExt, String clusterId) {
+ MQAdminExt adminExt, String clusterId) {
Set<String> brokerAddresses = collectProducerBrokerAddresses(adminExt);
Map<String, ClientConnectionVO> connections = new LinkedHashMap<>();
int successfulBrokers = 0;
@@ -138,7 +145,7 @@ public class RocketMQClientProvider implements
ClientProvider {
return new ArrayList<>(connections.values());
}
- private Set<String> collectProducerBrokerAddresses(DefaultMQAdminExt
adminExt) {
+ private Set<String> collectProducerBrokerAddresses(MQAdminExt adminExt) {
ClusterInfo clusterInfo;
try {
clusterInfo = adminExt.examineBrokerClusterInfo();
@@ -200,7 +207,7 @@ public class RocketMQClientProvider implements
ClientProvider {
.build();
}
- private List<ClientConnectionVO> findConsumerConnections(DefaultMQAdminExt
adminExt, String clusterId) {
+ private List<ClientConnectionVO> findConsumerConnections(MQAdminExt
adminExt, String clusterId) {
List<ClientConnectionVO> result = new ArrayList<>();
Set<String> groups = collectSubscriptionGroups(adminExt);
for (String group : groups) {
@@ -225,7 +232,7 @@ public class RocketMQClientProvider implements
ClientProvider {
return result;
}
- private Set<String> collectSubscriptionGroups(DefaultMQAdminExt adminExt) {
+ private Set<String> collectSubscriptionGroups(MQAdminExt adminExt) {
Set<String> groups = new LinkedHashSet<>();
ClusterInfo clusterInfo;
try {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
index 3af832f2..9018d9dd 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
@@ -29,14 +29,14 @@ import org.apache.rocketmq.common.message.MessageQueue;
import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
import org.apache.rocketmq.remoting.protocol.body.TopicList;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.instance.dlq.DLQGroupVO;
import org.apache.rocketmq.studio.instance.dlq.DLQProvider;
import org.apache.rocketmq.studio.instance.dlq.DLQResendResultVO;
import org.apache.rocketmq.studio.ops.audit.AuditService;
-import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
+import org.apache.rocketmq.tools.admin.MQAdminExt;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import org.springframework.beans.factory.ObjectProvider;
import org.springframework.context.annotation.Primary;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
@@ -63,34 +63,24 @@ public class RocketMQDLQProvider implements DLQProvider {
private static final long ONE_HOUR_MILLIS = 3600_000L;
private static final int RESEND_HARD_CAP = 5000;
- private final ObjectProvider<DefaultMQAdminExt> adminExtProvider;
+ private final RuntimeAdminClientResolver runtimeAdminClientResolver;
private final AuditService auditService;
- private final RocketMQProperties properties;
- public RocketMQDLQProvider(ObjectProvider<DefaultMQAdminExt>
adminExtProvider,
- AuditService auditService,
- RocketMQProperties properties) {
- this.adminExtProvider = adminExtProvider;
+ public RocketMQDLQProvider(RuntimeAdminClientResolver
runtimeAdminClientResolver,
+ AuditService auditService) {
+ this.runtimeAdminClientResolver = runtimeAdminClientResolver;
this.auditService = auditService;
- this.properties = properties;
}
@Override
- public List<DLQGroupVO> listDLQGroups(String clusterId) {
- DefaultMQAdminExt adminExt = adminExtProvider.getIfAvailable();
- if (adminExt == null) {
- log.warn("DefaultMQAdminExt is not configured, returning empty DLQ
group list");
- return Collections.emptyList();
- }
+ public List<DLQGroupVO> listDLQGroups(String instanceId) {
+ return runtimeAdminClientResolver.execute(instanceId,
this::listDLQGroups);
+ }
+ private List<DLQGroupVO> listDLQGroups(MQAdminExt adminExt) throws
Exception {
Set<String> topics;
- try {
- TopicList topicList = adminExt.fetchAllTopicList();
- topics = topicList == null ? Collections.emptySet() :
topicList.getTopicList();
- } catch (Exception e) {
- log.warn("Failed to fetch topic list for DLQ scan: {}",
e.getMessage());
- return Collections.emptyList();
- }
+ TopicList topicList = adminExt.fetchAllTopicList();
+ topics = topicList == null ? Collections.emptySet() :
topicList.getTopicList();
List<DLQGroupVO> groups = new ArrayList<>();
for (String topic : topics) {
@@ -106,7 +96,7 @@ public class RocketMQDLQProvider implements DLQProvider {
return groups;
}
- private DLQGroupVO buildDLQGroup(DefaultMQAdminExt adminExt, String
groupName, String dlqTopic) {
+ private DLQGroupVO buildDLQGroup(MQAdminExt adminExt, String groupName,
String dlqTopic) {
long messageCount = 0L;
LocalDateTime lastEnqueueTime = null;
try {
@@ -143,18 +133,19 @@ public class RocketMQDLQProvider implements DLQProvider {
}
@Override
- public DLQResendResultVO resendMessages(String groupName, Long startTime,
Long endTime, String targetTopic) {
- DefaultMQAdminExt adminExt = adminExtProvider.getIfAvailable();
+ public DLQResendResultVO resendMessages(String instanceId, String
groupName, Long startTime, Long endTime,
+ String targetTopic) {
+ String endpoint =
runtimeAdminClientResolver.resolveEndpoint(instanceId);
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + groupName;
long end = endTime != null ? endTime : System.currentTimeMillis();
long begin = startTime != null ? startTime : end - ONE_HOUR_MILLIS;
- List<MessageExt> deadLetters = collectDeadLetters(dlqTopic, begin,
end);
+ List<MessageExt> deadLetters = collectDeadLetters(endpoint, dlqTopic,
begin, end);
int resent = 0;
int failed = 0;
if (!deadLetters.isEmpty()) {
- DefaultMQProducer producer = newProducer();
+ DefaultMQProducer producer = newProducer(endpoint);
try {
producer.start();
for (MessageExt deadLetter : deadLetters) {
@@ -172,8 +163,8 @@ public class RocketMQDLQProvider implements DLQProvider {
}
}
- String detail = String.format("group=%s, dlqTopic=%s, targetTopic=%s,
matched=%d, resent=%d, failed=%d",
- groupName, dlqTopic, StringUtils.hasText(targetTopic) ?
targetTopic : "<original>",
+ String detail = String.format("instanceId=%s, group=%s, dlqTopic=%s,
targetTopic=%s, matched=%d, resent=%d, failed=%d",
+ instanceId, groupName, dlqTopic,
StringUtils.hasText(targetTopic) ? targetTopic : "<original>",
deadLetters.size(), resent, failed);
recordAudit(groupName, detail, failed == 0 ? "SUCCESS" : "PARTIAL");
log.info("DLQ resend completed: {}", detail);
@@ -185,8 +176,8 @@ public class RocketMQDLQProvider implements DLQProvider {
.build();
}
- private List<MessageExt> collectDeadLetters(String dlqTopic, long begin,
long end) {
- DefaultMQPullConsumer consumer = newPullConsumer();
+ private List<MessageExt> collectDeadLetters(String endpoint, String
dlqTopic, long begin, long end) {
+ DefaultMQPullConsumer consumer = newPullConsumer(endpoint);
List<MessageExt> result = new ArrayList<>();
try {
consumer.start();
@@ -272,22 +263,18 @@ public class RocketMQDLQProvider implements DLQProvider {
return null;
}
- private DefaultMQPullConsumer newPullConsumer() {
+ private DefaultMQPullConsumer newPullConsumer(String endpoint) {
DefaultMQPullConsumer consumer = new
DefaultMQPullConsumer("studio-dlq-query-group");
consumer.setInstanceName(ShortLivedClientName.next("studio-dlq-query"));
- if (StringUtils.hasText(properties.getNamesrvAddr())) {
- consumer.setNamesrvAddr(properties.getNamesrvAddr());
- }
+ consumer.setNamesrvAddr(endpoint);
return consumer;
}
- private DefaultMQProducer newProducer() {
+ private DefaultMQProducer newProducer(String endpoint) {
DefaultMQProducer producer = new
DefaultMQProducer(nextResendProducerGroup());
producer.setInstanceName(ShortLivedClientName.next("studio-dlq-resend"));
producer.setRetryTimesWhenSendFailed(2);
- if (StringUtils.hasText(properties.getNamesrvAddr())) {
- producer.setNamesrvAddr(properties.getNamesrvAddr());
- }
+ producer.setNamesrvAddr(endpoint);
return producer;
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
new file mode 100644
index 00000000..efc118d3
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
@@ -0,0 +1,83 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.cluster.broker;
+
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.instance.InstanceRepository;
+import org.apache.rocketmq.studio.instance.InstanceVO;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+
+import java.util.Optional;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.isNull;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+@ExtendWith(MockitoExtension.class)
+class RuntimeAdminClientResolverTest {
+
+ @Mock
+ private InstanceRepository instanceRepository;
+
+ @Mock
+ private MqAdminExtFactory adminFactory;
+
+ @Test
+ void resolvesTrimmedEndpointFromSelectedInstance() {
+ InstanceVO instance = InstanceVO.builder().endpoint(" namesrv-a:9876
").build();
+ instance.setId("instance-a");
+
when(instanceRepository.findById("instance-a")).thenReturn(Optional.of(instance));
+
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory);
+
+
assertThat(resolver.resolveEndpoint("instance-a")).isEqualTo("namesrv-a:9876");
+ }
+
+ @Test
+ void rejectsUnknownOrUnconfiguredInstances() {
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory);
+
when(instanceRepository.findById("missing")).thenReturn(Optional.empty());
+ InstanceVO noEndpoint = InstanceVO.builder().endpoint(" ").build();
+
when(instanceRepository.findById("no-endpoint")).thenReturn(Optional.of(noEndpoint));
+
+ assertThatThrownBy(() -> resolver.resolveEndpoint("missing"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Instance not found: missing");
+ assertThatThrownBy(() -> resolver.resolveEndpoint("no-endpoint"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Instance has no endpoint: no-endpoint");
+ }
+
+ @Test
+ void executesAgainstTheSelectedInstanceEndpoint() {
+ InstanceVO instance =
InstanceVO.builder().endpoint("namesrv-b:9876").build();
+
when(instanceRepository.findById("instance-b")).thenReturn(Optional.of(instance));
+ when(adminFactory.execute(eq("namesrv-b:9876"), isNull(),
any())).thenReturn("done");
+ RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory);
+
+ String result = resolver.execute("instance-b", admin -> "unused");
+ assertThat(result).isEqualTo("done");
+ verify(adminFactory).execute(eq("namesrv-b:9876"), isNull(), any());
+ }
+}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientControllerTest.java
index 5961efcf..5bd54a3a 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientControllerTest.java
@@ -69,9 +69,9 @@ class ClientControllerTest {
.connectedAt(LocalDateTime.of(2026, 1, 1, 12, 5))
.clusterName("production-cluster")
.build();
- when(clientService.listConnections(null,
null)).thenReturn(List.of(grpcClient, remotingClient));
+ when(clientService.listConnections("instance-1", null,
null)).thenReturn(List.of(grpcClient, remotingClient));
- mockMvc.perform(get("/api/clients"))
+ mockMvc.perform(get("/api/clients").param("instanceId", "instance-1"))
.andExpect(status().isOk())
.andExpect(jsonPath("$.code").value(200))
.andExpect(jsonPath("$.data").isArray())
@@ -84,14 +84,15 @@ class ClientControllerTest {
.andExpect(jsonPath("$.data[1].type").value("Producer"))
.andExpect(jsonPath("$.data[1].protocol").value("Remoting"));
- verify(clientService).listConnections(null, null);
+ verify(clientService).listConnections("instance-1", null, null);
}
@Test
void listConnectionsShouldPassClusterAndTypeFilters() throws Exception {
- when(clientService.listConnections("production-cluster",
"Consumer")).thenReturn(List.of());
+ when(clientService.listConnections("instance-1", "production-cluster",
"Consumer")).thenReturn(List.of());
mockMvc.perform(get("/api/clients")
+ .param("instanceId", "instance-1")
.param("clusterId", "production-cluster")
.param("type", "Consumer"))
.andExpect(status().isOk())
@@ -99,6 +100,6 @@ class ClientControllerTest {
.andExpect(jsonPath("$.data").isArray())
.andExpect(jsonPath("$.data").isEmpty());
- verify(clientService).listConnections("production-cluster",
"Consumer");
+ verify(clientService).listConnections("instance-1",
"production-cluster", "Consumer");
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStubTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStubTest.java
index bfefbe90..ed9c39e6 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStubTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientProviderStubTest.java
@@ -27,7 +27,7 @@ class ClientProviderStubTest {
@Test
void findConnectionsShouldFailWhenRealProviderIsMissing() {
- assertThatThrownBy(() ->
provider.findConnections("production-cluster", "Producer"))
+ assertThatThrownBy(() -> provider.findConnections("instance-1",
"production-cluster", "Producer"))
.isInstanceOf(BusinessException.class)
.hasMessage("Client connection provider is not configured")
.satisfies(ex -> assertThatBusinessExceptionCode(ex, 501));
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientServiceTest.java
index 6825cbf1..ce472901 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ClientServiceTest.java
@@ -43,21 +43,21 @@ class ClientServiceTest {
.clientId("client-1")
.clusterName("production-cluster")
.build();
- when(clientProvider.findConnections("production-cluster",
"Producer")).thenReturn(List.of(connection));
+ when(clientProvider.findConnections("instance-1",
"production-cluster", "Producer")).thenReturn(List.of(connection));
- List<ClientConnectionVO> result = clientService.listConnections("
production-cluster ", " Producer ");
+ List<ClientConnectionVO> result = clientService.listConnections("
instance-1 ", " production-cluster ", " Producer ");
assertThat(result).containsExactly(connection);
- verify(clientProvider).findConnections("production-cluster",
"Producer");
+ verify(clientProvider).findConnections("instance-1",
"production-cluster", "Producer");
}
@Test
void listConnectionsShouldTreatBlankFiltersAsUnspecified() {
- when(clientProvider.findConnections(null, null)).thenReturn(List.of());
+ when(clientProvider.findConnections(null, null,
null)).thenReturn(List.of());
- List<ClientConnectionVO> result = clientService.listConnections(" ",
"\t");
+ List<ClientConnectionVO> result = clientService.listConnections(" ", "
", "\t");
assertThat(result).isEmpty();
- verify(clientProvider).findConnections(null, null);
+ verify(clientProvider).findConnections(null, null, null);
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
index f83a06f9..c49a940e 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionServiceTest.java
@@ -97,7 +97,7 @@ class ProducerConnectionServiceTest {
@Test
void listProducerGroupsShouldReturnSortedUniqueActiveGroups() {
- when(clientProvider.findConnections(null, ClientType.Producer.name()))
+ when(clientProvider.findConnections(null, null,
ClientType.Producer.name()))
.thenReturn(List.of(
ClientConnectionVO.builder().producerGroup("
pg-payment ").build(),
ClientConnectionVO.builder().producerGroup("pg-order").build(),
@@ -107,6 +107,6 @@ class ProducerConnectionServiceTest {
assertThat(producerConnectionService.listProducerGroups())
.containsExactly("pg-order", "pg-payment");
- verify(clientProvider).findConnections(null,
ClientType.Producer.name());
+ verify(clientProvider).findConnections(null, null,
ClientType.Producer.name());
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
index 2071c348..00aca6d3 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
@@ -64,9 +64,9 @@ class DLQControllerTest {
.status("ACTIVE")
.build();
- when(dlqService.listDLQGroups(isNull())).thenReturn(List.of(group));
+
when(dlqService.listDLQGroups("instance-1")).thenReturn(List.of(group));
- mockMvc.perform(get("/api/dlq"))
+ mockMvc.perform(get("/api/dlq").param("instanceId", "instance-1"))
.andExpect(status().isOk())
.andExpect(jsonPath("$.code").value(200))
.andExpect(jsonPath("$.data").isArray())
@@ -76,20 +76,21 @@ class DLQControllerTest {
}
@Test
- void listDLQGroupsShouldPassClusterId() throws Exception {
- when(dlqService.listDLQGroups(eq("cluster-1"))).thenReturn(List.of());
+ void listDLQGroupsShouldPassInstanceId() throws Exception {
+ when(dlqService.listDLQGroups(eq("instance-1"))).thenReturn(List.of());
mockMvc.perform(get("/api/dlq")
- .param("clusterId", "cluster-1"))
+ .param("instanceId", "instance-1"))
.andExpect(status().isOk())
.andExpect(jsonPath("$.data").isArray());
- verify(dlqService).listDLQGroups(eq("cluster-1"));
+ verify(dlqService).listDLQGroups(eq("instance-1"));
}
@Test
void resendMessagesShouldReturnSuccess() throws Exception {
Map<String, Object> body = Map.of(
+ "instanceId", "instance-1",
"groupName", "test-group",
"startTime", 1000,
"endTime", 2000,
@@ -104,12 +105,13 @@ class DLQControllerTest {
.andExpect(jsonPath("$.message").value("success"));
verify(dlqService).resendMessages(
- eq("test-group"), eq(1000L), eq(2000L), eq("target-topic"));
+ eq("instance-1"), eq("test-group"), eq(1000L), eq(2000L),
eq("target-topic"));
}
@Test
void resendMessagesShouldHandleNullTimeRange() throws Exception {
Map<String, Object> body = Map.of(
+ "instanceId", "instance-1",
"groupName", "test-group",
"targetTopic", "target-topic"
);
@@ -121,7 +123,7 @@ class DLQControllerTest {
.andExpect(jsonPath("$.code").value(200));
verify(dlqService).resendMessages(
- eq("test-group"), isNull(), isNull(), eq("target-topic"));
+ eq("instance-1"), eq("test-group"), isNull(), isNull(),
eq("target-topic"));
}
@Test
@@ -139,6 +141,7 @@ class DLQControllerTest {
@Test
void resendMessagesShouldRejectMissingGroupName() throws Exception {
Map<String, Object> body = Map.of(
+ "instanceId", "instance-1",
"startTime", 1000,
"endTime", 2000,
"targetTopic", "target-topic"
@@ -157,6 +160,7 @@ class DLQControllerTest {
@Test
void resendMessagesShouldRejectInvalidTimeType() throws Exception {
Map<String, Object> body = Map.of(
+ "instanceId", "instance-1",
"groupName", "test-group",
"startTime", "invalid",
"endTime", 2000,
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStubTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStubTest.java
index 8143c2f9..69700556 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStubTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStubTest.java
@@ -28,7 +28,7 @@ class DLQProviderStubTest {
@Test
void listDLQGroupsShouldFailExplicitlyWhenRealProviderIsMissing() {
- assertThatThrownBy(() -> provider.listDLQGroups("cluster-1"))
+ assertThatThrownBy(() -> provider.listDLQGroups("instance-1"))
.isInstanceOf(BusinessException.class)
.hasMessage("DLQ provider is not configured")
.extracting("code")
@@ -37,7 +37,7 @@ class DLQProviderStubTest {
@Test
void resendMessagesShouldFailExplicitlyWhenRealProviderIsMissing() {
- assertThatThrownBy(() -> provider.resendMessages("group-1", 1000L,
2000L, "target-topic"))
+ assertThatThrownBy(() -> provider.resendMessages("instance-1",
"group-1", 1000L, 2000L, "target-topic"))
.isInstanceOf(BusinessException.class)
.hasMessage("DLQ provider is not configured")
.extracting("code")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQServiceTest.java
index 9dfeff4c..2cb65468 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQServiceTest.java
@@ -56,43 +56,43 @@ class DLQServiceTest {
.status("ACTIVE")
.build()
);
- when(dlqProvider.listDLQGroups("cluster-1")).thenReturn(groups);
+ when(dlqProvider.listDLQGroups("instance-1")).thenReturn(groups);
- List<DLQGroupVO> result = dlqService.listDLQGroups("cluster-1");
+ List<DLQGroupVO> result = dlqService.listDLQGroups("instance-1");
assertThat(result).hasSize(2);
assertThat(result.get(0).getGroupName()).isEqualTo("group-1");
assertThat(result.get(0).getDlqTopic()).isEqualTo("%DLQ%group-1");
- verify(dlqProvider).listDLQGroups("cluster-1");
+ verify(dlqProvider).listDLQGroups("instance-1");
}
@Test
void listDLQGroupsShouldReturnEmptyWhenNone() {
- when(dlqProvider.listDLQGroups("cluster-2")).thenReturn(List.of());
+ when(dlqProvider.listDLQGroups("instance-2")).thenReturn(List.of());
- List<DLQGroupVO> result = dlqService.listDLQGroups("cluster-2");
+ List<DLQGroupVO> result = dlqService.listDLQGroups("instance-2");
assertThat(result).isEmpty();
- verify(dlqProvider).listDLQGroups("cluster-2");
+ verify(dlqProvider).listDLQGroups("instance-2");
}
@Test
void resendMessagesShouldDelegateToProvider() {
- dlqService.resendMessages("group-1", 1000L, 2000L, "target-topic");
+ dlqService.resendMessages("instance-1", "group-1", 1000L, 2000L,
"target-topic");
- verify(dlqProvider).resendMessages("group-1", 1000L, 2000L,
"target-topic");
+ verify(dlqProvider).resendMessages("instance-1", "group-1", 1000L,
2000L, "target-topic");
}
@Test
void resendMessagesShouldAcceptNullTimeRange() {
- dlqService.resendMessages("group-1", null, null, "target-topic");
+ dlqService.resendMessages("instance-1", "group-1", null, null,
"target-topic");
- verify(dlqProvider).resendMessages("group-1", null, null,
"target-topic");
+ verify(dlqProvider).resendMessages("instance-1", "group-1", null,
null, "target-topic");
}
@Test
void resendMessagesShouldRejectBlankGroupName() {
- assertThatThrownBy(() -> dlqService.resendMessages(" ", 1000L, 2000L,
"target-topic"))
+ assertThatThrownBy(() -> dlqService.resendMessages("instance-1", " ",
1000L, 2000L, "target-topic"))
.hasMessage("groupName is required");
verifyNoInteractions(dlqProvider);
@@ -100,7 +100,7 @@ class DLQServiceTest {
@Test
void resendMessagesShouldRejectPartialTimeRange() {
- assertThatThrownBy(() -> dlqService.resendMessages("group-1", 1000L,
null, "target-topic"))
+ assertThatThrownBy(() -> dlqService.resendMessages("instance-1",
"group-1", 1000L, null, "target-topic"))
.hasMessage("startTime and endTime must be provided together");
verifyNoInteractions(dlqProvider);
@@ -108,7 +108,7 @@ class DLQServiceTest {
@Test
void resendMessagesShouldRejectNonPositiveTimeRange() {
- assertThatThrownBy(() -> dlqService.resendMessages("group-1", 0L,
2000L, "target-topic"))
+ assertThatThrownBy(() -> dlqService.resendMessages("instance-1",
"group-1", 0L, 2000L, "target-topic"))
.hasMessage("startTime and endTime must be positive");
verifyNoInteractions(dlqProvider);
@@ -116,7 +116,7 @@ class DLQServiceTest {
@Test
void resendMessagesShouldRejectReversedTimeRange() {
- assertThatThrownBy(() -> dlqService.resendMessages("group-1", 2000L,
1000L, "target-topic"))
+ assertThatThrownBy(() -> dlqService.resendMessages("instance-1",
"group-1", 2000L, 1000L, "target-topic"))
.hasMessage("endTime must not be earlier than startTime");
verifyNoInteractions(dlqProvider);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQClientProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQClientProviderTest.java
index ef27a801..97dc243e 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQClientProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQClientProviderTest.java
@@ -29,6 +29,7 @@ import
org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
import org.apache.rocketmq.studio.cluster.client.ClientConnectionVO;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
import org.junit.jupiter.api.BeforeEach;
@@ -47,7 +48,9 @@ import java.util.concurrent.ConcurrentHashMap;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -60,12 +63,18 @@ class RocketMQClientProviderTest {
@Mock
private DefaultMQAdminExt adminExt;
+ @Mock
+ private RuntimeAdminClientResolver runtimeAdminClientResolver;
+
private RocketMQClientProvider provider;
@BeforeEach
void setUp() {
- when(adminExtProvider.getIfAvailable()).thenReturn(adminExt);
- provider = new RocketMQClientProvider(adminExtProvider);
+ lenient().when(adminExtProvider.getIfAvailable()).thenReturn(adminExt);
+ provider = new RocketMQClientProvider(adminExtProvider,
runtimeAdminClientResolver);
+ lenient().when(runtimeAdminClientResolver.execute(anyString(),
any())).thenAnswer(invocation ->
+
invocation.<org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory.AdminAction<Object>>
+ getArgument(1).apply(adminExt));
}
@Test
@@ -74,7 +83,7 @@ class RocketMQClientProviderTest {
clusterInfo.setBrokerAddrTable(null);
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo);
- List<ClientConnectionVO> connections =
provider.findConnections("cluster-a", "Producer");
+ List<ClientConnectionVO> connections =
provider.findConnections("instance-a", "cluster-a", "Producer");
assertThat(connections).isEmpty();
verify(adminExt).examineBrokerClusterInfo();
@@ -95,7 +104,7 @@ class RocketMQClientProviderTest {
"pg-order", List.of(shared),
"pg-payment", List.of(another))));
- List<ClientConnectionVO> connections =
provider.findConnections("cluster-a", "Producer");
+ List<ClientConnectionVO> connections =
provider.findConnections("instance-a", "cluster-a", "Producer");
assertThat(connections).hasSize(2);
assertThat(connections)
@@ -117,7 +126,7 @@ class RocketMQClientProviderTest {
.thenReturn(new ProducerTableInfo(Map.of(
"pg-order", List.of(producerInfo("producer-client",
"10.0.0.1:1000")))));
- List<ClientConnectionVO> connections =
provider.findConnections("cluster-a", "Producer");
+ List<ClientConnectionVO> connections =
provider.findConnections("instance-a", "cluster-a", "Producer");
assertThat(connections).singleElement().satisfies(connection -> {
assertThat(connection.getClientId()).isEqualTo("producer-client");
@@ -132,7 +141,7 @@ class RocketMQClientProviderTest {
when(adminExt.getAllProducerInfo(anyString()))
.thenThrow(new IllegalStateException("broker unavailable"));
- assertThatThrownBy(() -> provider.findConnections("cluster-a",
"Producer"))
+ assertThatThrownBy(() -> provider.findConnections("instance-a",
"cluster-a", "Producer"))
.isInstanceOf(BusinessException.class)
.hasMessage("Failed to query producer connections from all
brokers")
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
@@ -203,7 +212,7 @@ class RocketMQClientProviderTest {
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo);
when(adminExt.getAllSubscriptionGroup("127.0.0.1:10911",
5000L)).thenReturn(wrapper);
- List<ClientConnectionVO> connections =
provider.findConnections("cluster-a", "Consumer");
+ List<ClientConnectionVO> connections =
provider.findConnections("instance-a", "cluster-a", "Consumer");
assertThat(connections).isEmpty();
verify(adminExt).examineBrokerClusterInfo();
@@ -228,7 +237,7 @@ class RocketMQClientProviderTest {
when(adminExt.getAllSubscriptionGroup("127.0.0.1:10911",
5000L)).thenReturn(wrapper);
when(adminExt.examineConsumerConnectionInfo("group-a")).thenReturn(consumerConnection);
- List<ClientConnectionVO> connections =
provider.findConnections("cluster-a", "Consumer");
+ List<ClientConnectionVO> connections =
provider.findConnections("instance-a", "cluster-a", "Consumer");
assertThat(connections).hasSize(1);
assertThat(connections.get(0).getClientId()).isEqualTo("consumer-client");
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProviderTest.java
index ce7d767e..34b5352c 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProviderTest.java
@@ -20,15 +20,14 @@ import
org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.common.MixAll;
import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.ops.audit.AuditService;
-import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.MockedConstruction;
import org.mockito.junit.jupiter.MockitoExtension;
-import org.springframework.beans.factory.ObjectProvider;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
@@ -48,10 +47,7 @@ import static org.mockito.Mockito.when;
class RocketMQDLQProviderTest {
@Mock
- private ObjectProvider<DefaultMQAdminExt> adminExtProvider;
-
- @Mock
- private DefaultMQAdminExt adminExt;
+ private RuntimeAdminClientResolver runtimeAdminClientResolver;
@Mock
private AuditService auditService;
@@ -60,8 +56,8 @@ class RocketMQDLQProviderTest {
@BeforeEach
void setUp() {
- lenient().when(adminExtProvider.getIfAvailable()).thenReturn(adminExt);
- provider = new RocketMQDLQProvider(adminExtProvider, auditService, new
RocketMQProperties());
+
lenient().when(runtimeAdminClientResolver.resolveEndpoint("instance-a")).thenReturn("namesrv-a:9876");
+ provider = new RocketMQDLQProvider(runtimeAdminClientResolver,
auditService);
}
@Test
@@ -75,10 +71,11 @@ class RocketMQDLQProviderTest {
});
MockedConstruction<DefaultMQProducer> mockedProducers =
mockConstruction(DefaultMQProducer.class)) {
- provider.resendMessages("group-a", 100L, 200L, "target-topic");
+ provider.resendMessages("instance-a", "group-a", 100L, 200L,
"target-topic");
assertThat(mockedConsumers.constructed()).hasSize(1);
DefaultMQPullConsumer consumer =
mockedConsumers.constructed().get(0);
+ verify(consumer).setNamesrvAddr("namesrv-a:9876");
verify(consumer).start();
verify(consumer).fetchSubscribeMessageQueues(dlqTopic);
verify(consumer, never()).pull(any(MessageQueue.class),
anyString(), anyLong(), anyInt());
@@ -90,6 +87,7 @@ class RocketMQDLQProviderTest {
eq("group-a"),
contains("matched=0, resent=0, failed=0"),
eq("SUCCESS"));
+ verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
}
@Test
diff --git a/web/src/api/connections.test.ts b/web/src/api/connections.test.ts
index 5f5fc5e2..590c28ad 100644
--- a/web/src/api/connections.test.ts
+++ b/web/src/api/connections.test.ts
@@ -47,12 +47,20 @@ describe('client connections API', () => {
it('uses the query parameters supported by the backend', async () => {
mock.onGet('/clients').reply((config) => {
- expect(config.params).toEqual({ clusterId: 'production-cluster', type:
'Producer' });
+ expect(config.params).toEqual({
+ instanceId: 'instance-1',
+ clusterId: 'production-cluster',
+ type: 'Producer',
+ });
return [200, { code: 200, data: [connection] }];
});
await expect(
- listConnections({ clusterId: 'production-cluster', type: 'Producer' }),
+ listConnections({
+ instanceId: 'instance-1',
+ clusterId: 'production-cluster',
+ type: 'Producer',
+ }),
).resolves.toEqual([connection]);
});
});
diff --git a/web/src/api/connections.ts b/web/src/api/connections.ts
index ef9d365c..ae9de444 100644
--- a/web/src/api/connections.ts
+++ b/web/src/api/connections.ts
@@ -15,6 +15,7 @@ export interface ClientConnection {
}
export interface ClientConnectionQuery {
+ instanceId: string;
clusterId?: string;
type?: string;
}
diff --git a/web/src/api/dlq.test.ts b/web/src/api/dlq.test.ts
index 6adca0af..c77e97b6 100644
--- a/web/src/api/dlq.test.ts
+++ b/web/src/api/dlq.test.ts
@@ -43,13 +43,16 @@ describe('DLQ API', () => {
});
it('loads and unwraps DLQ groups', async () => {
- mock.onGet('/dlq').reply(200, { code: 200, data: [group] });
+ mock
+ .onGet('/dlq', { params: { instanceId: 'instance-1' } })
+ .reply(200, { code: 200, data: [group] });
- await expect(listDLQGroups()).resolves.toEqual([group]);
+ await expect(listDLQGroups('instance-1')).resolves.toEqual([group]);
});
it('sends epoch milliseconds for the resend time range', async () => {
const payload = {
+ instanceId: 'instance-1',
groupName: group.groupName,
startTime: 1784246400000,
endTime: 1784332800000,
diff --git a/web/src/api/message.ts b/web/src/api/message.ts
index b103745c..c21509b4 100644
--- a/web/src/api/message.ts
+++ b/web/src/api/message.ts
@@ -84,12 +84,13 @@ export async function getMessageTrace(msgId: string) {
}
// ─── DLQ ────────────────────────────────────────────────────────
-export async function listDLQGroups() {
- const res = await client.get<{ data: DLQGroup[] }>('/dlq');
+export async function listDLQGroups(instanceId: string) {
+ const res = await client.get<{ data: DLQGroup[] }>('/dlq', { params: {
instanceId } });
return res.data.data;
}
export async function resendDLQ(data: {
+ instanceId: string;
groupName: string;
startTime: number;
endTime: number;
diff --git a/web/src/pages/cluster/__tests__/ClientsPage.test.tsx
b/web/src/pages/cluster/__tests__/ClientsPage.test.tsx
index d72367a5..3ffae0b4 100644
--- a/web/src/pages/cluster/__tests__/ClientsPage.test.tsx
+++ b/web/src/pages/cluster/__tests__/ClientsPage.test.tsx
@@ -28,6 +28,21 @@ import ClientsPage from '../clients';
vi.mock('../../../services/connectionsService', () => ({
listConnections: vi.fn(),
}));
+vi.mock('../../../services/instanceService', () => ({
+ listInstances: vi.fn().mockResolvedValue([
+ {
+ id: 'instance-1',
+ name: 'Instance 1',
+ endpoint: 'namesrv-1:9876',
+ type: 'DIRECT',
+ remark: '',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ createdAt: '',
+ updatedAt: '',
+ },
+ ]),
+}));
const connection: ClientConnection = {
clientId: '[email protected]:49152',
@@ -99,13 +114,20 @@ const renderWithProviders = (ui: React.ReactElement) =>
);
describe('Clients page', () => {
+ it('loads connections for the selected instance', async () => {
+ renderWithProviders(<ClientsPage />);
+
+ await screen.findByText('[email protected]:49152');
+ expect(connectionsService.listConnections).toHaveBeenCalledWith({
instanceId: 'instance-1' });
+ });
+
it('summarizes connection types, protocols, and language versions', async ()
=> {
vi.mocked(connectionsService.listConnections).mockResolvedValue(connections);
renderWithProviders(<ClientsPage />);
- expect(
- within(await screen.findByTestId('connection-total')).getByText('3'),
- ).toBeInTheDocument();
+ await waitFor(() => {
+
expect(within(screen.getByTestId('connection-total')).getByText('3')).toBeInTheDocument();
+ });
expect(within(screen.getByTestId('producer-total')).getByText('1')).toBeInTheDocument();
expect(within(screen.getByTestId('consumer-total')).getByText('2')).toBeInTheDocument();
diff --git a/web/src/pages/cluster/clients.tsx
b/web/src/pages/cluster/clients.tsx
index 5d9cfb1f..8dab3d21 100644
--- a/web/src/pages/cluster/clients.tsx
+++ b/web/src/pages/cluster/clients.tsx
@@ -39,6 +39,8 @@ import PageHeader from '../../components/PageHeader';
import { useLang } from '../../i18n/LangContext';
import type { ClientConnection } from '../../api/connections';
import { listConnections } from '../../services/connectionsService';
+import { listInstances } from '../../services/instanceService';
+import type { Instance } from '../../api/instance';
import { formatDateTime } from '../../utils/format';
const { Text } = Typography;
@@ -105,6 +107,8 @@ const ClientsPage = () => {
const { t } = useLang();
const { token } = theme.useToken();
const [connections, setConnections] = useState<ClientConnection[]>([]);
+ const [instances, setInstances] = useState<Instance[]>([]);
+ const [selectedInstanceId, setSelectedInstanceId] = useState('');
const [loading, setLoading] = useState(true);
const [search, setSearch] = useState('');
const [clusterFilter, setClusterFilter] = useState<string>('ALL');
@@ -113,8 +117,25 @@ const ClientsPage = () => {
useEffect(() => {
let cancelled = false;
+ void listInstances().then((nextInstances) => {
+ if (cancelled) return;
+ setInstances(nextInstances);
+ setSelectedInstanceId((current) => current || nextInstances[0]?.id ||
'');
+ });
+ return () => {
+ cancelled = true;
+ };
+ }, []);
- void listConnections()
+ useEffect(() => {
+ let cancelled = false;
+ if (!selectedInstanceId) {
+ return () => {
+ cancelled = true;
+ };
+ }
+
+ void listConnections({ instanceId: selectedInstanceId })
.then((nextConnections) => {
if (!cancelled) {
setConnections(nextConnections);
@@ -133,7 +154,7 @@ const ClientsPage = () => {
return () => {
cancelled = true;
};
- }, []);
+ }, [selectedInstanceId]);
/* ─── Cluster options using nsClusterName ─── */
const clusterOptions = useMemo(() => {
@@ -348,6 +369,14 @@ const ClientsPage = () => {
{/* ─── Filter Bar ─── */}
<Flex justify="space-between" align="center" style={{ marginBottom: 16
}}>
<Space size={12} wrap>
+ <Select
+ aria-label="Instance"
+ value={selectedInstanceId || undefined}
+ onChange={setSelectedInstanceId}
+ placeholder="Select instance"
+ style={{ width: 180 }}
+ options={instances.map((instance) => ({ value: instance.id, label:
instance.name }))}
+ />
<Select
aria-label={t('clients.cluster')}
value={clusterFilter}
diff --git a/web/src/pages/instance/__tests__/DLQPage.test.tsx
b/web/src/pages/instance/__tests__/DLQPage.test.tsx
index 4c1b8adb..efc9df66 100644
--- a/web/src/pages/instance/__tests__/DLQPage.test.tsx
+++ b/web/src/pages/instance/__tests__/DLQPage.test.tsx
@@ -31,7 +31,19 @@ vi.mock('../../../services/messageService', () => ({
resendDLQ: vi.fn(),
}));
vi.mock('../../../services/instanceService', () => ({
- listInstances: vi.fn().mockResolvedValue([]),
+ listInstances: vi.fn().mockResolvedValue([
+ {
+ id: 'instance-1',
+ name: 'Instance 1',
+ endpoint: 'namesrv-1:9876',
+ type: 'DIRECT',
+ remark: '',
+ topicCount: 0,
+ consumerGroupCount: 0,
+ createdAt: '',
+ updatedAt: '',
+ },
+ ]),
}));
const dlqGroup: DLQGroup = {
@@ -56,7 +68,7 @@ const renderWithProviders = (ui: React.ReactElement) =>
render(
<App>
<LangProvider>
- <MemoryRouter>{ui}</MemoryRouter>
+ <MemoryRouter
initialEntries={['/instance/instance-1/dlq']}>{ui}</MemoryRouter>
</LangProvider>
</App>,
);
@@ -107,7 +119,7 @@ describe('DLQ page', () => {
expect(await screen.findByText('cg-order')).toBeInTheDocument();
expect(screen.getByText('%DLQ%cg-order')).toBeInTheDocument();
- expect(messageService.listDLQGroups).toHaveBeenCalledTimes(1);
+ expect(messageService.listDLQGroups).toHaveBeenCalledWith('instance-1');
});
it('surfaces unavailable DLQ provider errors when loading groups', async ()
=> {
diff --git a/web/src/pages/instance/dlq.tsx b/web/src/pages/instance/dlq.tsx
index 0f02bddf..66900131 100644
--- a/web/src/pages/instance/dlq.tsx
+++ b/web/src/pages/instance/dlq.tsx
@@ -123,7 +123,20 @@ const DLQPage = () => {
useEffect(() => {
let cancelled = false;
- void listDLQGroups()
+ if (!selectedInstanceId) {
+ void Promise.resolve().then(() => {
+ if (cancelled) return;
+ setGroups([]);
+ setSelectedGroupNames([]);
+ setLoadError(null);
+ setLoading(false);
+ });
+ return () => {
+ cancelled = true;
+ };
+ }
+
+ void listDLQGroups(selectedInstanceId)
.then((nextGroups) => {
if (!cancelled) {
setGroups(nextGroups);
@@ -146,7 +159,7 @@ const DLQPage = () => {
return () => {
cancelled = true;
};
- }, [refreshKey]);
+ }, [refreshKey, selectedInstanceId]);
/* ─── Filtering ─── */
const filtered = useMemo(() => {
@@ -176,12 +189,13 @@ const DLQPage = () => {
message.warning('请输入目标 Topic');
return;
}
- if (!retryGroup) return;
+ if (!retryGroup || !selectedInstanceId) return;
setRetrySubmitting(true);
setRetryError(null);
try {
const result = await resendDLQ({
+ instanceId: selectedInstanceId,
groupName: retryGroup.groupName,
startTime: retryRange[0].valueOf(),
endTime: retryRange[1].valueOf(),
diff --git a/web/src/services/connectionsService.test.ts
b/web/src/services/connectionsService.test.ts
index 6c461108..c9056abb 100644
--- a/web/src/services/connectionsService.test.ts
+++ b/web/src/services/connectionsService.test.ts
@@ -26,14 +26,22 @@ import { listConnections } from './connectionsService';
describe('connectionsService mock connections', () => {
it('returns defensive copies after applying filters', async () => {
- const connections = await listConnections({ clusterId: 'ns-prod', type:
'Consumer' });
+ const connections = await listConnections({
+ instanceId: 'instance-1',
+ clusterId: 'ns-prod',
+ type: 'Consumer',
+ });
const originalClientId = connections[0].clientId;
const originalAddress = connections[0].address;
connections[0].clientId = 'mutated-client';
connections[0].address = '127.0.0.1:8081';
- const fresh = await listConnections({ clusterId: 'ns-prod', type:
'Consumer' });
+ const fresh = await listConnections({
+ instanceId: 'instance-1',
+ clusterId: 'ns-prod',
+ type: 'Consumer',
+ });
expect(fresh[0].clientId).toBe(originalClientId);
expect(fresh[0].address).toBe(originalAddress);
diff --git a/web/src/services/messageService.test.ts
b/web/src/services/messageService.test.ts
index 377eb5d5..b3f12704 100644
--- a/web/src/services/messageService.test.ts
+++ b/web/src/services/messageService.test.ts
@@ -46,9 +46,7 @@ describe('message service mock data', () => {
endTime: Date.parse('2026-07-01T10:26:00.000Z'),
});
- expect(messages.map((message) => message.msgId)).toEqual([
- 'AC1E0A6400002A9F0000000001A3F7C2',
- ]);
+ expect(messages.map((message) =>
message.msgId)).toEqual(['AC1E0A6400002A9F0000000001A3F7C2']);
});
it('returns copied message trace rows', async () => {
@@ -67,12 +65,12 @@ describe('message service mock data', () => {
});
it('returns copied DLQ group rows', async () => {
- const first = await listDLQGroups();
+ const first = await listDLQGroups('instance-1');
expect(first[0].groupName).toBe('cg-order-processor');
first[0].groupName = 'mutated-group';
- const second = await listDLQGroups();
+ const second = await listDLQGroups('instance-1');
expect(second[0].groupName).toBe('cg-order-processor');
expect(second[0]).not.toBe(first[0]);
});
diff --git a/web/src/services/messageService.ts
b/web/src/services/messageService.ts
index d4abcf84..31ccf040 100644
--- a/web/src/services/messageService.ts
+++ b/web/src/services/messageService.ts
@@ -56,12 +56,13 @@ export async function getMessageTrace(msgId: string):
Promise<TraceRecord | null
return messageApi.getMessageTrace(msgId);
}
-export async function listDLQGroups(): Promise<DLQGroup[]> {
+export async function listDLQGroups(instanceId: string): Promise<DLQGroup[]> {
if (isMockMode()) return (mockDLQGroups as unknown as
DLQGroup[]).map(cloneDLQGroup);
- return messageApi.listDLQGroups();
+ return messageApi.listDLQGroups(instanceId);
}
export async function resendDLQ(data: {
+ instanceId: string;
groupName: string;
startTime: number;
endTime: number;