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 b5a7d6c0b refactor(server): replace short-lived admin clients with
pooled long-lived instances (#2544)
b5a7d6c0b is described below
commit b5a7d6c0b2263b790730b6ff2b01c25d9f470f02
Author: lizhimins <[email protected]>
AuthorDate: Mon Aug 24 14:58:31 2026 +0800
refactor(server): replace short-lived admin clients with pooled long-lived
instances (#2544)
Introduce MqClientPool to cache pull-consumer and producer instances by
(endpoint, credential) pair, replacing the previous pattern of creating
and shutting down clients on every request. RuntimeAdminClientResolver
now exposes executePullConsumer/executeProducer entry points backed by
the pool. Remove ShortLivedClientName utility which is no longer needed.
All provider implementations (RocketMQAdminClientImpl,
RocketMQMessageProvider,
RocketMQDLQProvider) are migrated to use pooled clients. Tests updated
accordingly.
---
.../studio/cluster/broker/MqClientPool.java | 219 +++++++++++
.../cluster/broker/RuntimeAdminClientResolver.java | 28 ++
.../provider/apache/RocketMQAdminClientImpl.java | 32 +-
.../provider/apache/RocketMQDLQProvider.java | 68 ++--
.../provider/apache/RocketMQMessageProvider.java | 240 +++++++-----
.../provider/apache/ShortLivedClientName.java | 29 --
.../broker/RuntimeAdminClientResolverTest.java | 19 +-
.../apache/RocketMQAdminClientImplTest.java | 216 +++++------
.../apache/RocketMQClusterProviderTest.java | 3 +-
.../provider/apache/RocketMQDLQProviderTest.java | 421 ++++++++-------------
.../apache/RocketMQMessageProviderTest.java | 312 ++++++---------
.../provider/apache/ShortLivedClientNameTest.java | 37 --
12 files changed, 803 insertions(+), 821 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqClientPool.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqClientPool.java
new file mode 100644
index 000000000..d9fcf292f
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqClientPool.java
@@ -0,0 +1,219 @@
+/*
+ * 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.client.consumer.DefaultMQPullConsumer;
+import org.apache.rocketmq.client.producer.DefaultMQProducer;
+import org.apache.rocketmq.remoting.RPCHook;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+
+import jakarta.annotation.PreDestroy;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Component;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicInteger;
+
+/**
+ * Long-lived pool of {@link DefaultMQPullConsumer} and {@link
DefaultMQProducer} clients, keyed by
+ * normalized NameServer address plus a non-secret authentication identity.
Mirrors the lifecycle
+ * rules of {@link MqAdminExtFactory}: clients are created lazily, reused
across requests, and shut
+ * down only on application shutdown or explicit endpoint release. Per-request
client creation is
+ * forbidden because each {@code start()} registers with NameServer/brokers
and tears down
+ * connections and threads again on {@code shutdown()}.
+ */
+@Slf4j
+@Component
+public class MqClientPool {
+
+ private static final String PULL_CONSUMER_GROUP =
"studio-pool-pull-consumer";
+ private static final String PRODUCER_GROUP = "studio-pool-producer";
+ private static final long PRODUCER_SEND_TIMEOUT_MILLIS = 5000L;
+
+ private enum Kind { PULL_CONSUMER, PRODUCER }
+
+ private record ClientKey(String namesrvAddr, String
authenticationIdentity, Kind kind) {
+ }
+
+ private final Map<ClientKey, Object> cache = new ConcurrentHashMap<>();
+ private final AtomicInteger instanceCounter = new AtomicInteger();
+ private volatile boolean closed = false;
+
+ @FunctionalInterface
+ public interface ClientAction<C, T> {
+ T apply(C client) throws Exception;
+ }
+
+ public <T> T withPullConsumer(String namesrvAddr, RPCHook rpcHook, String
authenticationIdentity,
+ ClientAction<DefaultMQPullConsumer, T>
action) {
+ return execute(new ClientKey(normalize(namesrvAddr),
identity(authenticationIdentity), Kind.PULL_CONSUMER),
+ rpcHook, this::createPullConsumer, action);
+ }
+
+ public <T> T withProducer(String namesrvAddr, RPCHook rpcHook, String
authenticationIdentity,
+ ClientAction<DefaultMQProducer, T> action) {
+ return execute(new ClientKey(normalize(namesrvAddr),
identity(authenticationIdentity), Kind.PRODUCER),
+ rpcHook, this::createProducer, action);
+ }
+
+ private <C, T> T execute(ClientKey cacheKey, RPCHook rpcHook,
+ ClientCreator<C> creator, ClientAction<C, T>
action) {
+ if (cacheKey.namesrvAddr().isEmpty()) {
+ throw new BusinessException(400, "NameServer address is required");
+ }
+ if (closed) {
+ throw new BusinessException(503, "RocketMQ client pool is shutting
down");
+ }
+ @SuppressWarnings("unchecked")
+ C client = (C) cache.computeIfAbsent(cacheKey, key -> {
+ // Re-check under the cache lock so a request that passed the
initial closed check
+ // cannot create a fresh connection while the pool is shutting
down.
+ if (closed) {
+ throw new BusinessException(503, "RocketMQ client pool is
shutting down");
+ }
+ return creator.create(key.namesrvAddr(), rpcHook);
+ });
+ try {
+ return action.apply(client);
+ } catch (BusinessException ex) {
+ throw ex;
+ } catch (Exception ex) {
+ log.warn("RocketMQ client action failed against namesrv {}: {}",
cacheKey.namesrvAddr(), ex.getMessage());
+ throw new BusinessException(502, "RocketMQ client call failed: " +
rootMessage(ex));
+ }
+ }
+
+ /** Stops and removes pooled clients for an endpoint that is no longer
referenced by Studio. */
+ public void release(String namesrvAddr) {
+ String normalized = normalize(namesrvAddr);
+ if (normalized.isEmpty()) {
+ return;
+ }
+ cache.entrySet().removeIf(entry -> {
+ if (!entry.getKey().namesrvAddr().equals(normalized)) {
+ return false;
+ }
+ safeShutdown(entry.getValue());
+ return true;
+ });
+ log.info("Released pooled RocketMQ clients for namesrv {}",
normalized);
+ }
+
+ /** Stops one credential-scoped client without interrupting other
identities on the endpoint. */
+ public void release(String namesrvAddr, String authenticationIdentity) {
+ String normalized = normalize(namesrvAddr);
+ if (normalized.isEmpty()) {
+ return;
+ }
+ for (Kind kind : Kind.values()) {
+ ClientKey key = new ClientKey(normalized,
identity(authenticationIdentity), kind);
+ Object client = cache.remove(key);
+ if (client != null) {
+ safeShutdown(client);
+ }
+ }
+ log.info("Released pooled RocketMQ clients for namesrv {} and identity
{}", normalized,
+ identity(authenticationIdentity));
+ }
+
+ @FunctionalInterface
+ private interface ClientCreator<C> {
+ C create(String namesrvAddr, RPCHook rpcHook);
+ }
+
+ private DefaultMQPullConsumer createPullConsumer(String namesrvAddr,
RPCHook rpcHook) {
+ DefaultMQPullConsumer consumer = new
DefaultMQPullConsumer(PULL_CONSUMER_GROUP, rpcHook);
+ consumer.setNamesrvAddr(namesrvAddr);
+ consumer.setInstanceName(buildInstanceName(namesrvAddr));
+ try {
+ consumer.start();
+ log.info("Started pooled RocketMQ pull consumer for namesrv {}",
namesrvAddr);
+ return consumer;
+ } catch (Exception ex) {
+ safeShutdown(consumer);
+ throw new BusinessException(502,
+ "Failed to connect NameServer " + namesrvAddr + ": " +
rootMessage(ex));
+ }
+ }
+
+ private DefaultMQProducer createProducer(String namesrvAddr, RPCHook
rpcHook) {
+ DefaultMQProducer producer = new DefaultMQProducer(PRODUCER_GROUP,
rpcHook);
+ producer.setNamesrvAddr(namesrvAddr);
+ producer.setInstanceName(buildInstanceName(namesrvAddr));
+ producer.setSendMsgTimeout((int) PRODUCER_SEND_TIMEOUT_MILLIS);
+ producer.setRetryTimesWhenSendFailed(2);
+ // Clusters without a trace topic drop trace dispatch asynchronously,
so enabling
+ // it unconditionally lets Studio-sent messages produce traces
wherever trace is on.
+ producer.setEnableTrace(true);
+ try {
+ producer.start();
+ log.info("Started pooled RocketMQ producer for namesrv {}",
namesrvAddr);
+ return producer;
+ } catch (Exception ex) {
+ safeShutdown(producer);
+ throw new BusinessException(502,
+ "Failed to connect NameServer " + namesrvAddr + ": " +
rootMessage(ex));
+ }
+ }
+
+ private String buildInstanceName(String namesrvAddr) {
+ return "rmq-studio-pool-" + Integer.toHexString(namesrvAddr.hashCode())
+ + "-" + instanceCounter.incrementAndGet();
+ }
+
+ private static String normalize(String namesrvAddr) {
+ if (namesrvAddr == null || namesrvAddr.isBlank()) {
+ return "";
+ }
+ return MqAdminExtFactory.normalizeNamesrvAddr(namesrvAddr);
+ }
+
+ private static String identity(String authenticationIdentity) {
+ return authenticationIdentity == null ||
authenticationIdentity.isBlank()
+ ? "anonymous" : authenticationIdentity.trim();
+ }
+
+ private void safeShutdown(Object client) {
+ try {
+ if (client instanceof DefaultMQPullConsumer pullConsumer) {
+ pullConsumer.shutdown();
+ } else if (client instanceof DefaultMQProducer producer) {
+ producer.shutdown();
+ }
+ } catch (Exception ex) {
+ log.debug("Ignoring client shutdown error: {}", ex.getMessage());
+ }
+ }
+
+ private String rootMessage(Throwable ex) {
+ Throwable cause = ex;
+ while (cause.getCause() != null && cause.getCause() != cause) {
+ cause = cause.getCause();
+ }
+ String message = cause.getMessage();
+ return message == null ? cause.getClass().getSimpleName() : message;
+ }
+
+ @PreDestroy
+ public void shutdown() {
+ closed = true;
+ cache.values().forEach(this::safeShutdown);
+ cache.clear();
+ log.info("Shut down all pooled RocketMQ clients");
+ }
+}
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
index 37cdf1c75..f689e3acd 100644
---
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
@@ -9,6 +9,8 @@ package org.apache.rocketmq.studio.cluster.broker;
import lombok.RequiredArgsConstructor;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
+import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
+import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.exception.BusinessException;
@@ -24,6 +26,7 @@ public class RuntimeAdminClientResolver {
private final InstanceRepository instanceRepository;
private final MqAdminExtFactory adminFactory;
private final MqAdminProperties adminProperties;
+ private final MqClientPool clientPool;
public InstanceVO resolveInstance(String instanceId) {
if (!StringUtils.hasText(instanceId)) {
@@ -64,6 +67,31 @@ public class RuntimeAdminClientResolver {
credentialRef, action);
}
+ /** Runs an action on a pooled, long-lived pull consumer bound to the
instance endpoint. */
+ public <T> T executePullConsumer(String instanceId,
+
MqClientPool.ClientAction<DefaultMQPullConsumer, T> action) {
+ InstanceVO instance =
requireApacheInstance(resolveInstance(instanceId));
+ String credentialRef = credentialRef(instance);
+ return clientPool.withPullConsumer(requireEndpoint(instance),
resolveCredential(credentialRef),
+ credentialRef, action);
+ }
+
+ /** Runs an action on a pooled, long-lived producer bound to the instance
endpoint. */
+ public <T> T executeProducer(String instanceId,
+ MqClientPool.ClientAction<DefaultMQProducer,
T> action) {
+ InstanceVO instance =
requireApacheInstance(resolveInstance(instanceId));
+ String credentialRef = credentialRef(instance);
+ return clientPool.withProducer(requireEndpoint(instance),
resolveCredential(credentialRef),
+ credentialRef, action);
+ }
+
+ private String requireEndpoint(InstanceVO instance) {
+ if (instance == null || !StringUtils.hasText(instance.getEndpoint())) {
+ throw new BusinessException(400, "Instance endpoint is required");
+ }
+ return instance.getEndpoint().trim();
+ }
+
private String credentialRef(InstanceVO instance) {
return StringUtils.hasText(instance.getAdminCredentialRef())
? instance.getAdminCredentialRef().trim() : null;
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 8c223a3d8..6448be000 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,12 +26,12 @@ 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.remoting.RPCHook;
import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
import org.apache.rocketmq.remoting.protocol.ResponseCode;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
@@ -71,7 +71,6 @@ import java.util.Set;
@RequiredArgsConstructor
public class RocketMQAdminClientImpl implements AdminClient {
- private static final String MESSAGE_SENDER_GROUP_PREFIX =
"studio-msg-sender";
private static final int MAX_MESSAGE_SIZE = 4 * 1024 * 1024; // 4 MB
default broker limit
private static final String LEGACY_METADATA_SCOPE = "";
@@ -81,6 +80,7 @@ public class RocketMQAdminClientImpl implements AdminClient {
private final RmqGroupMapper groupMapper;
private final AuditService auditService;
private final RuntimeAdminClientResolver runtimeAdminClientResolver;
+ private final MqClientPool clientPool;
@org.springframework.beans.factory.annotation.Autowired(required = false)
private ProxyConsumerResolver proxyConsumerResolver;
@@ -441,16 +441,7 @@ public class RocketMQAdminClientImpl implements
AdminClient {
throw new BusinessException(400, message);
}
- String namesrvAddr = namesrvAddr(request.getInstanceId());
- RPCHook credentialHook = credentialHook(request.getInstanceId());
-
- DefaultMQProducer producer = new
DefaultMQProducer(nextMessageSenderGroup(), credentialHook);
- producer.setNamesrvAddr(namesrvAddr);
- producer.setSendMsgTimeout(5000);
-
- try {
- producer.start();
-
+ MqClientPool.ClientAction<DefaultMQProducer, SendMessageVO> sendAction
= producer -> {
Message msg = new Message(topic, tag, key, bodyBytes);
// Add custom properties
@@ -476,21 +467,20 @@ public class RocketMQAdminClientImpl implements
AdminClient {
.sendTime(System.currentTimeMillis())
.offsetMsgId(sendResult.getOffsetMsgId())
.build();
+ };
+ try {
+ return StringUtils.hasText(request.getInstanceId())
+ ?
runtimeAdminClientResolver.executeProducer(request.getInstanceId(), sendAction)
+ : clientPool.withProducer(namesrvAddr(), null, null,
sendAction);
} catch (BusinessException e) {
recordAudit("SEND_MESSAGE", request.getTopic(), e.getMessage(),
"FAILED");
throw e;
} catch (Exception e) {
recordAudit("SEND_MESSAGE", request.getTopic(), e.getMessage(),
"FAILED");
throw new BusinessException(500, "Failed to send message: " +
e.getMessage());
- } finally {
- producer.shutdown();
}
}
- static String nextMessageSenderGroup() {
- return ShortLivedClientName.next(MESSAGE_SENDER_GROUP_PREFIX);
- }
-
@Override
public ConsumerGroupVO createConsumerGroup(ConsumerGroupVO group) {
if (group != null && group.getInstanceId() != null) {
@@ -749,12 +739,6 @@ public class RocketMQAdminClientImpl implements
AdminClient {
: namesrvAddr();
}
- private RPCHook credentialHook(String instanceId) {
- return StringUtils.hasText(instanceId)
- ? runtimeAdminClientResolver.resolveCredentialHook(instanceId)
- : null;
- }
-
private <T> T executeForInstance(String instanceId,
MqAdminExtFactory.AdminAction<T> action) {
if (StringUtils.hasText(instanceId)) {
return runtimeAdminClientResolver.execute(instanceId, action);
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 7c4aca641..2058a77a0 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
@@ -27,7 +27,6 @@ import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
import org.apache.rocketmq.remoting.protocol.body.TopicList;
@@ -174,13 +173,11 @@ public class RocketMQDLQProvider implements DLQProvider {
throw new BusinessException(400, "DLQ resend start time must be
before end time");
}
- String endpoint =
runtimeAdminClientResolver.resolveEndpoint(instanceId);
- RPCHook credentialHook =
runtimeAdminClientResolver.resolveCredentialHook(instanceId);
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + groupName;
DeadLetterScanResult scanResult;
try {
- scanResult = collectDeadLetters(endpoint, credentialHook,
dlqTopic, begin, end, RESEND_HARD_CAP);
+ scanResult = collectDeadLetters(instanceId, dlqTopic, begin, end,
RESEND_HARD_CAP);
} catch (BusinessException e) {
String detail = String.format("instanceId=%s, group=%s,
dlqTopic=%s, targetTopic=%s, "
+ "matched=0, resent=0, failed=0,
scanIncomplete=true, scanFailedQueues=all",
@@ -190,26 +187,26 @@ public class RocketMQDLQProvider implements DLQProvider {
throw e;
}
List<MessageExt> deadLetters = scanResult.messages();
- int resent = 0;
- int failed = 0;
+ int[] counts = {0, 0};
if (!deadLetters.isEmpty()) {
- DefaultMQProducer producer = newProducer(endpoint, credentialHook);
try {
- producer.start();
- for (MessageExt deadLetter : deadLetters) {
- if (resendOne(producer, deadLetter, targetTopic)) {
- resent++;
- } else {
- failed++;
+ runtimeAdminClientResolver.executeProducer(instanceId,
producer -> {
+ for (MessageExt deadLetter : deadLetters) {
+ if (resendOne(producer, deadLetter, targetTopic)) {
+ counts[0]++;
+ } else {
+ counts[1]++;
+ }
}
- }
+ return null;
+ });
} catch (Exception e) {
- log.warn("Failed to start resend producer for group {}: {}",
groupName, e.getMessage());
- failed += deadLetters.size() - resent;
- } finally {
- producer.shutdown();
+ log.warn("Failed to resend dead letters for group {}: {}",
groupName, e.getMessage());
+ counts[1] += deadLetters.size() - counts[0];
}
}
+ int resent = counts[0];
+ int failed = counts[1];
String outcome = classifyOutcome(deadLetters.size(), resent, failed,
scanResult.scanIncomplete());
String detail = String.format("instanceId=%s, group=%s, dlqTopic=%s,
targetTopic=%s, matched=%d, resent=%d, "
@@ -232,13 +229,11 @@ public class RocketMQDLQProvider implements DLQProvider {
@Override
public List<DLQMessageVO> exportMessages(String instanceId, String
groupName, Long startTime, Long endTime,
Integer maxCount) {
- String endpoint =
runtimeAdminClientResolver.resolveEndpoint(instanceId);
- RPCHook credentialHook =
runtimeAdminClientResolver.resolveCredentialHook(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;
int cap = maxCount == null || maxCount <= 0 ? RESEND_HARD_CAP :
Math.min(maxCount, RESEND_HARD_CAP);
- DeadLetterScanResult scanResult = collectDeadLetters(endpoint,
credentialHook, dlqTopic, begin, end, cap);
+ DeadLetterScanResult scanResult = collectDeadLetters(instanceId,
dlqTopic, begin, end, cap);
return scanResult.messages().stream().map(this::toExportVO).toList();
}
@@ -267,14 +262,18 @@ public class RocketMQDLQProvider implements DLQProvider {
}
}
- private DeadLetterScanResult collectDeadLetters(String endpoint, RPCHook
credentialHook, String dlqTopic,
+ private DeadLetterScanResult collectDeadLetters(String instanceId, String
dlqTopic,
long begin, long end, int
cap) {
- DefaultMQPullConsumer consumer = newPullConsumer(endpoint,
credentialHook);
+ return runtimeAdminClientResolver.executePullConsumer(instanceId,
+ consumer -> scanDeadLetters(consumer, dlqTopic, begin, end,
cap));
+ }
+
+ private DeadLetterScanResult scanDeadLetters(DefaultMQPullConsumer
consumer, String dlqTopic,
+ long begin, long end, int
cap) {
List<MessageExt> result = new ArrayList<>();
int failedQueueCount = 0;
boolean truncated = false;
try {
- consumer.start();
Set<MessageQueue> queues =
consumer.fetchSubscribeMessageQueues(dlqTopic);
if (queues == null || queues.isEmpty()) {
return new DeadLetterScanResult(result, 0, false);
@@ -352,8 +351,6 @@ public class RocketMQDLQProvider implements DLQProvider {
}
log.warn("Failed to collect dead letters from {}: {}", dlqTopic,
e.getMessage());
throw new BusinessException(502, "Failed to scan DLQ topic " +
dlqTopic + ": " + e.getMessage());
- } finally {
- consumer.shutdown();
}
return new DeadLetterScanResult(result, failedQueueCount, truncated);
}
@@ -423,25 +420,6 @@ public class RocketMQDLQProvider implements DLQProvider {
return null;
}
- private DefaultMQPullConsumer newPullConsumer(String endpoint, RPCHook
credentialHook) {
- DefaultMQPullConsumer consumer = new
DefaultMQPullConsumer("studio-dlq-query-group", credentialHook);
-
consumer.setInstanceName(ShortLivedClientName.next("studio-dlq-query"));
- consumer.setNamesrvAddr(endpoint);
- return consumer;
- }
-
- private DefaultMQProducer newProducer(String endpoint, RPCHook
credentialHook) {
- DefaultMQProducer producer = new
DefaultMQProducer(nextResendProducerGroup(), credentialHook);
-
producer.setInstanceName(ShortLivedClientName.next("studio-dlq-resend"));
- producer.setRetryTimesWhenSendFailed(2);
- producer.setNamesrvAddr(endpoint);
- return producer;
- }
-
- static String nextResendProducerGroup() {
- return ShortLivedClientName.next("studio-dlq-resend");
- }
-
private String classifyOutcome(int matched, int resent, int failed,
boolean scanIncomplete) {
if (scanIncomplete) {
return "PARTIAL";
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 fc324929b..b1b3c00de 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
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.provider.apache;
import org.apache.rocketmq.client.QueryResult;
+import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
import org.apache.rocketmq.client.consumer.PullResult;
import org.apache.rocketmq.client.consumer.PullStatus;
@@ -25,15 +26,18 @@ import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageId;
import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.remoting.RPCHook;
+import org.apache.rocketmq.remoting.protocol.ResponseCode;
import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
+import org.apache.rocketmq.remoting.protocol.route.QueueData;
+import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
import org.apache.rocketmq.studio.instance.message.ConsumerStatusVO;
import org.apache.rocketmq.studio.instance.message.MessageProvider;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
+import org.apache.rocketmq.studio.instance.message.QueueOffsetVO;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
@@ -87,7 +91,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 int MAX_PULLS_PER_QUEUE = 1_000;
+ 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;
private static final int MAX_PULL_ATTEMPTS_PER_QUEUE = MAX_PULLS_PER_QUEUE
+ MAX_CONSECUTIVE_OFFSET_ILLEGAL;
@@ -100,15 +104,12 @@ public class RocketMQMessageProvider implements
MessageProvider {
@Override
public List<MessageRecordVO> queryMessages(String instanceId, String
topic, String msgId, String tag, String key,
Long startTime, Long endTime) {
- String endpoint =
runtimeAdminClientResolver.resolveEndpoint(instanceId);
- RPCHook credentialHook =
runtimeAdminClientResolver.resolveCredentialHook(instanceId);
return runtimeAdminClientResolver.execute(instanceId,
- adminExt -> queryMessages(instanceId, (DefaultMQAdminExt)
adminExt, endpoint,
- credentialHook, topic, msgId, tag, key, startTime,
endTime));
+ adminExt -> queryMessages(instanceId, (DefaultMQAdminExt)
adminExt, topic, msgId, tag, key,
+ startTime, endTime));
}
- private List<MessageRecordVO> queryMessages(String instanceId,
DefaultMQAdminExt adminExt, String endpoint,
- RPCHook credentialHook,
+ private List<MessageRecordVO> queryMessages(String instanceId,
DefaultMQAdminExt adminExt,
String topic, String msgId,
String tag, String key,
Long startTime, Long endTime)
{
@@ -127,7 +128,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
if (begin >= 0 && end >= 0 && end - begin >
MAX_TOPIC_QUERY_WINDOW_MILLIS) {
throw new BusinessException(400, "Topic message query time
range must not exceed 7 days");
}
- result = queryByTopic(endpoint, credentialHook, topic, tag, begin,
end, DEFAULT_TOPIC_LIMIT);
+ result = queryByTopic(instanceId, topic, tag, begin, end,
DEFAULT_TOPIC_LIMIT);
} else {
log.warn("queryMessages requires at least one of msgId/topic,
returning empty list");
return Collections.emptyList();
@@ -240,91 +241,139 @@ public class RocketMQMessageProvider implements
MessageProvider {
}
}
+ @Override
+ public List<QueueOffsetVO> getQueueOffsets(String instanceId, String
topic) {
+ return runtimeAdminClientResolver.execute(instanceId, adminExt -> {
+ List<QueueOffsetVO> result = new ArrayList<>();
+ try {
+ TopicRouteData route = adminExt.examineTopicRouteInfo(topic);
+ if (route == null || route.getQueueDatas() == null) {
+ return Collections.emptyList();
+ }
+ for (QueueData queueData : route.getQueueDatas()) {
+ for (int queueId = 0; queueId <
queueData.getWriteQueueNums(); queueId++) {
+ MessageQueue queue = new MessageQueue(topic,
queueData.getBrokerName(), queueId);
+ result.add(QueueOffsetVO.builder()
+ .brokerName(queue.getBrokerName())
+ .queueId(queue.getQueueId())
+ .minOffset(adminExt.minOffset(queue))
+ .maxOffset(adminExt.maxOffset(queue))
+ .build());
+ }
+ }
+ result.sort(Comparator.comparing(QueueOffsetVO::getBrokerName)
+ .thenComparingInt(QueueOffsetVO::getQueueId));
+ } catch (Exception e) {
+ log.warn("getQueueOffsets(topic={}) failed: {}", topic,
e.getMessage());
+ throw new BusinessException(502, "Failed to get queue offsets:
" + e.getMessage());
+ }
+ return result;
+ });
+ }
+
+ @Override
+ public MessageRecordVO pullMessageAtOffset(String instanceId, String
topic, String brokerName,
+ int queueId, long offset) {
+ return runtimeAdminClientResolver.executePullConsumer(instanceId,
consumer -> {
+ try {
+ MessageQueue queue = new MessageQueue(topic, brokerName,
queueId);
+ PullResult pullResult = consumer.pull(queue, "*", offset, 1);
+ if (pullResult == null || pullResult.getPullStatus() !=
PullStatus.FOUND
+ || pullResult.getMsgFoundList() == null ||
pullResult.getMsgFoundList().isEmpty()) {
+ return null;
+ }
+ return toRecordVO(pullResult.getMsgFoundList().get(0),
brokerName);
+ } catch (Exception e) {
+ 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());
+ }
+ });
+ }
+
/**
- * Scan a topic within a time range using a short-lived pull consumer,
mirroring the approach
- * used by the RocketMQ dashboard for time-range topic queries.
+ * Scan a topic within a time range on the pooled long-lived pull
consumer, mirroring the
+ * approach used by the RocketMQ dashboard for time-range topic queries.
*/
- private List<MessageRecordVO> queryByTopic(String endpoint, RPCHook
credentialHook, String topic, String tag,
+ private List<MessageRecordVO> queryByTopic(String instanceId, String
topic, String tag,
long begin, long end, int
limit) {
- DefaultMQPullConsumer consumer = newPullConsumer("studio-msg-query",
endpoint, credentialHook);
int resultLimit = Math.min(limit, TOPIC_QUERY_HARD_CAP);
- PriorityQueue<MessageRecordVO> newestMessages = new
PriorityQueue<>(TOPIC_QUERY_ORDER);
- try {
- consumer.start();
- Set<MessageQueue> queues =
consumer.fetchSubscribeMessageQueues(topic);
- if (queues == null || queues.isEmpty()) {
- return Collections.emptyList();
- }
- for (MessageQueue queue : queues) {
- TopicQueueScanPlan scanPlan =
buildTopicQueueScanPlan(consumer, queue, begin, end);
- if (scanPlan.isEmpty()) {
- continue;
- }
- if (scanPlan.truncated()) {
- log.info("Truncate topic query for {} queue {} to offsets
[{}..{}) within the guarded tail budget",
- topic, queue, scanPlan.startOffset(),
scanPlan.endOffsetExclusive());
+ return runtimeAdminClientResolver.executePullConsumer(instanceId,
consumer -> {
+ PriorityQueue<MessageRecordVO> newestMessages = new
PriorityQueue<>(TOPIC_QUERY_ORDER);
+ try {
+ Set<MessageQueue> queues =
consumer.fetchSubscribeMessageQueues(topic);
+ if (queues == null || queues.isEmpty()) {
+ return Collections.emptyList();
}
- int consecutiveIllegalOffsets = 0;
- int pullAttempts = 0;
- for (long offset = scanPlan.startOffset(); offset <
scanPlan.endOffsetExclusive(); ) {
- if (++pullAttempts > MAX_PULL_ATTEMPTS_PER_QUEUE) {
- log.warn("Stop topic query for {} because queue {}
exhausted the guarded pull budget at offset {}",
- topic, queue, offset);
- break;
- }
- PullResult pullResult = consumer.pull(queue, "*", offset,
TOPIC_PULL_BATCH_SIZE);
- if (pullResult == null) {
- log.warn("Stop topic query for {} because queue {}
returned no pull result", topic, queue);
- break;
+ for (MessageQueue queue : queues) {
+ TopicQueueScanPlan scanPlan =
buildTopicQueueScanPlan(consumer, queue, begin, end);
+ if (scanPlan.isEmpty()) {
+ continue;
}
- long nextOffset = pullResult.getNextBeginOffset();
- if (nextOffset <= offset) {
- log.warn("Stop topic query for {} because queue {} did
not advance offset {}", topic, queue, offset);
- break;
+ if (scanPlan.truncated()) {
+ log.info("Truncate topic query for {} queue {} to
offsets [{}..{}) within the guarded tail budget",
+ topic, queue, scanPlan.startOffset(),
scanPlan.endOffsetExclusive());
}
- offset = Math.min(nextOffset,
scanPlan.endOffsetExclusive());
- if (pullResult.getPullStatus() ==
PullStatus.OFFSET_ILLEGAL) {
- // The broker returned a corrected offset in
nextBeginOffset because
- // the requested offset is no longer valid (expired,
compacted, or
- // before the queue's minimum offset). Retry from the
corrected
- // position instead of abandoning the queue --
otherwise messages
- // that still exist after the corrected offset are
silently dropped.
- if (++consecutiveIllegalOffsets >
MAX_CONSECUTIVE_OFFSET_ILLEGAL) {
- log.warn("Stop topic query for {} because queue {}
returned OFFSET_ILLEGAL "
- + "{} times consecutively, giving up at
offset {}", topic, queue,
- consecutiveIllegalOffsets, offset);
+ int consecutiveIllegalOffsets = 0;
+ int pullAttempts = 0;
+ for (long offset = scanPlan.startOffset(); offset <
scanPlan.endOffsetExclusive(); ) {
+ if (++pullAttempts > MAX_PULL_ATTEMPTS_PER_QUEUE) {
+ log.warn("Stop topic query for {} because queue {}
exhausted the guarded pull budget at offset {}",
+ topic, queue, offset);
break;
}
- log.debug("Offset was illegal for queue {} in topic
{}, retrying from {}",
- queue, topic, offset);
- continue;
- }
- if (pullResult.getPullStatus() != PullStatus.FOUND
- || pullResult.getMsgFoundList() == null) {
- break;
- }
- consecutiveIllegalOffsets = 0;
- for (MessageExt messageExt : pullResult.getMsgFoundList())
{
- if (messageExt.getStoreTimestamp() < begin
- || messageExt.getStoreTimestamp() > end) {
- continue;
+ PullResult pullResult = consumer.pull(queue, "*",
offset, TOPIC_PULL_BATCH_SIZE);
+ if (pullResult == null) {
+ log.warn("Stop topic query for {} because queue {}
returned no pull result", topic, queue);
+ break;
+ }
+ long nextOffset = pullResult.getNextBeginOffset();
+ if (nextOffset <= offset) {
+ log.warn("Stop topic query for {} because queue {}
did not advance offset {}", topic, queue, offset);
+ break;
}
- if (!matchesTag(messageExt, tag)) {
+ offset = Math.min(nextOffset,
scanPlan.endOffsetExclusive());
+ if (pullResult.getPullStatus() ==
PullStatus.OFFSET_ILLEGAL) {
+ // The broker returned a corrected offset in
nextBeginOffset because
+ // the requested offset is no longer valid
(expired, compacted, or
+ // before the queue's minimum offset). Retry from
the corrected
+ // position instead of abandoning the queue --
otherwise messages
+ // that still exist after the corrected offset are
silently dropped.
+ if (++consecutiveIllegalOffsets >
MAX_CONSECUTIVE_OFFSET_ILLEGAL) {
+ log.warn("Stop topic query for {} because
queue {} returned OFFSET_ILLEGAL "
+ + "{} times consecutively, giving up
at offset {}", topic, queue,
+ consecutiveIllegalOffsets, offset);
+ break;
+ }
+ log.debug("Offset was illegal for queue {} in
topic {}, retrying from {}",
+ queue, topic, offset);
continue;
}
- addTopicQueryCandidate(newestMessages,
toRecordVO(messageExt), resultLimit);
+ if (pullResult.getPullStatus() != PullStatus.FOUND
+ || pullResult.getMsgFoundList() == null) {
+ break;
+ }
+ consecutiveIllegalOffsets = 0;
+ for (MessageExt messageExt :
pullResult.getMsgFoundList()) {
+ if (messageExt.getStoreTimestamp() < begin
+ || messageExt.getStoreTimestamp() > end) {
+ continue;
+ }
+ if (!matchesTag(messageExt, tag)) {
+ continue;
+ }
+ addTopicQueryCandidate(newestMessages,
toRecordVO(messageExt, queue.getBrokerName()), resultLimit);
+ }
}
}
+ } catch (Exception e) {
+ log.warn("queryByTopic(topic={}) failed: {}", topic,
e.getMessage());
+ throw new BusinessException(502, "Failed to query messages by
topic: " + e.getMessage());
}
- } catch (Exception e) {
- log.warn("queryByTopic(topic={}) failed: {}", topic,
e.getMessage());
- throw new BusinessException(502, "Failed to query messages by
topic: " + e.getMessage());
- } finally {
- consumer.shutdown();
- }
- return newestMessages.stream()
- .sorted(TOPIC_QUERY_ORDER.reversed())
- .toList();
+ return newestMessages.stream()
+ .sorted(TOPIC_QUERY_ORDER.reversed())
+ .toList();
+ });
}
private TopicQueueScanPlan buildTopicQueueScanPlan(DefaultMQPullConsumer
consumer, MessageQueue queue,
@@ -410,6 +459,13 @@ public class RocketMQMessageProvider implements
MessageProvider {
} catch (BusinessException e) {
throw e;
} catch (Exception e) {
+ if (isTraceTopicAbsent(e)) {
+ // The cluster has no trace topic route (trace dispatch
disabled): the RPC
+ // succeeded but there is no business data, so return an empty
trace instead
+ // of surfacing an error (exception-grading convention).
+ log.info("Trace topic not available on this cluster
(msgId={}), returning empty trace", msgId);
+ return emptyTrace();
+ }
log.warn("Trace query for msgId={} failed: {}", msgId,
e.getMessage());
throw new BusinessException(502, "Failed to query message trace: "
+ e.getMessage());
}
@@ -420,6 +476,18 @@ public class RocketMQMessageProvider implements
MessageProvider {
.build();
}
+ private static boolean isTraceTopicAbsent(Throwable error) {
+ Throwable cause = error;
+ while (cause != null) {
+ if (cause instanceof MQClientException clientException
+ && clientException.getResponseCode() ==
ResponseCode.TOPIC_NOT_EXIST) {
+ return true;
+ }
+ cause = cause.getCause() == cause ? null : cause.getCause();
+ }
+ return false;
+ }
+
/**
* Attempts to resolve the store timestamp of the original message so the
trace
* query window can be derived from the message's own timeline rather than
the
@@ -601,6 +669,10 @@ public class RocketMQMessageProvider implements
MessageProvider {
}
MessageRecordVO toRecordVO(MessageExt messageExt) {
+ return toRecordVO(messageExt, null);
+ }
+
+ MessageRecordVO toRecordVO(MessageExt messageExt, String brokerName) {
byte[] body = messageExt.getBody();
DisplayBody displayBody = displayBody(body);
Map<String, String> properties = messageExt.getProperties();
@@ -610,6 +682,9 @@ public class RocketMQMessageProvider implements
MessageProvider {
.topic(messageExt.getTopic())
.tag(messageExt.getTags())
.key(messageExt.getKeys())
+ .brokerName(brokerName)
+ .queueId(messageExt.getQueueId())
+ .queueOffset(messageExt.getQueueOffset())
.body(displayBody.value())
.bodyEncoding(displayBody.encoding())
.bodyTruncated(displayBody.truncated())
@@ -685,13 +760,6 @@ public class RocketMQMessageProvider implements
MessageProvider {
return tag.equals(messageExt.getTags());
}
- private DefaultMQPullConsumer newPullConsumer(String groupPrefix, String
endpoint, RPCHook credentialHook) {
- DefaultMQPullConsumer consumer = new DefaultMQPullConsumer(groupPrefix
+ "-group", credentialHook);
- consumer.setInstanceName(ShortLivedClientName.next(groupPrefix));
- consumer.setNamesrvAddr(endpoint);
- return consumer;
- }
-
private static TraceRecordVO emptyTrace() {
return TraceRecordVO.builder()
.nodes(Collections.emptyList())
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientName.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientName.java
deleted file mode 100644
index 3259856dc..000000000
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientName.java
+++ /dev/null
@@ -1,29 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package org.apache.rocketmq.studio.provider.apache;
-
-import java.util.UUID;
-
-final class ShortLivedClientName {
-
- private ShortLivedClientName() {
- }
-
- static String next(String prefix) {
- return prefix + "-" + UUID.randomUUID();
- }
-}
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
index f54831898..f6b18b4de 100644
---
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
@@ -50,6 +50,9 @@ class RuntimeAdminClientResolverTest {
@Mock
private MqAdminExtFactory adminFactory;
+ @Mock
+ private MqClientPool clientPool;
+
@Test
void resolvesTrimmedEndpointFromSelectedInstance() {
InstanceVO instance = InstanceVO.builder().endpoint(" namesrv-a:9876
").build();
@@ -57,7 +60,7 @@ class RuntimeAdminClientResolverTest {
when(instanceRepository.findByIdentifier("instance-a")).thenReturn(Optional.of(instance));
RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
- new MqAdminProperties());
+ new MqAdminProperties(), clientPool);
assertThat(resolver.resolveEndpoint("instance-a")).isEqualTo("namesrv-a:9876");
}
@@ -65,7 +68,7 @@ class RuntimeAdminClientResolverTest {
@Test
void rejectsUnknownOrUnconfiguredInstances() {
RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
- new MqAdminProperties());
+ new MqAdminProperties(), clientPool);
when(instanceRepository.findByIdentifier("missing")).thenReturn(Optional.empty());
InstanceVO noEndpoint = InstanceVO.builder().endpoint(" ").build();
when(instanceRepository.findByIdentifier("no-endpoint")).thenReturn(Optional.of(noEndpoint));
@@ -84,7 +87,7 @@ class RuntimeAdminClientResolverTest {
when(instanceRepository.findByIdentifier("instance-b")).thenReturn(Optional.of(instance));
when(adminFactory.execute(eq("namesrv-b:9876"), isNull(), isNull(),
any())).thenReturn("done");
RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
- new MqAdminProperties());
+ new MqAdminProperties(), clientPool);
String result = resolver.execute("instance-b", admin -> "unused");
assertThat(result).isEqualTo("done");
@@ -100,7 +103,7 @@ class RuntimeAdminClientResolverTest {
instance.setId(2L);
when(instanceRepository.findByIdentifier("cloud-instance")).thenReturn(Optional.of(instance));
RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
- new MqAdminProperties());
+ new MqAdminProperties(), clientPool);
assertThatThrownBy(() -> resolver.resolveEndpoint("cloud-instance"))
.isInstanceOf(BusinessException.class)
@@ -127,7 +130,7 @@ class RuntimeAdminClientResolverTest {
when(adminFactory.execute(eq("namesrv-b:9876"), any(),
eq("production-admin"), any()))
.thenReturn("done");
RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
- properties);
+ properties, clientPool);
String result = resolver.execute("instance-b", ignored -> "unused");
@@ -155,7 +158,7 @@ class RuntimeAdminClientResolverTest {
properties.getCredentials().put("production-admin", credential);
when(instanceRepository.findByIdentifier("instance-b")).thenReturn(Optional.of(instance));
RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
- properties);
+ properties, clientPool);
org.apache.rocketmq.acl.common.AclClientRPCHook hook =
(org.apache.rocketmq.acl.common.AclClientRPCHook)
resolver.resolveCredentialHook("instance-b");
@@ -171,7 +174,7 @@ class RuntimeAdminClientResolverTest {
instance.setId(5L);
when(instanceRepository.findByIdentifier("instance-b")).thenReturn(Optional.of(instance));
RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
- new MqAdminProperties());
+ new MqAdminProperties(), clientPool);
assertThat(resolver.resolveCredentialHook("instance-b")).isNull();
verifyNoInteractions(adminFactory);
@@ -192,7 +195,7 @@ class RuntimeAdminClientResolverTest {
instance.setId(3L);
when(instanceRepository.findByIdentifier("instance-b")).thenReturn(Optional.of(instance));
RuntimeAdminClientResolver resolver = new
RuntimeAdminClientResolver(instanceRepository, adminFactory,
- properties);
+ properties, clientPool);
assertThatThrownBy(() -> resolver.execute("instance-b", ignored ->
"unused"))
.isInstanceOf(BusinessException.class)
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 fba406cac..df80091c7 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
@@ -22,13 +22,13 @@ import
org.apache.rocketmq.client.producer.DefaultMQProducer;
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.remoting.RPCHook;
import org.apache.rocketmq.remoting.exception.RemotingTimeoutException;
import org.apache.rocketmq.remoting.protocol.ResponseCode;
import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
@@ -48,11 +48,9 @@ import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
-import org.mockito.MockedConstruction;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
-import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -65,13 +63,13 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.lenient;
-import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.doThrow;
-import static org.mockito.Mockito.mockConstruction;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@@ -93,6 +91,10 @@ class RocketMQAdminClientImplTest {
private AuditService auditService;
@Mock
private RuntimeAdminClientResolver runtimeAdminClientResolver;
+ @Mock
+ private MqClientPool clientPool;
+ @Mock
+ private DefaultMQProducer sendProducer;
private RocketMQAdminClientImpl adminClient;
@@ -101,8 +103,12 @@ class RocketMQAdminClientImplTest {
lenient().when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
lenient().when(adminFactory.execute(anyString(), any(),
any())).thenAnswer(invocation ->
invocation.<MqAdminExtFactory.AdminAction<Object>>getArgument(2).apply(adminExt));
+ lenient().when(clientPool.withProducer(any(), any(), any(),
any())).thenAnswer(invocation ->
+ invocation.<MqClientPool.ClientAction<DefaultMQProducer,
Object>>getArgument(3).apply(sendProducer));
+ lenient().when(runtimeAdminClientResolver.executeProducer(any(),
any())).thenAnswer(invocation ->
+ invocation.<MqClientPool.ClientAction<DefaultMQProducer,
Object>>getArgument(1).apply(sendProducer));
adminClient = new RocketMQAdminClientImpl(adminFactory, properties,
topicMapper, groupMapper, auditService,
- runtimeAdminClientResolver);
+ runtimeAdminClientResolver, clientPool);
}
@Test
@@ -729,26 +735,20 @@ class RocketMQAdminClientImplTest {
@Test
void sendMessageShouldNotFailWhenAuditRecordingFails() throws Exception {
- when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
doThrow(new RuntimeException("audit db down")).when(auditService)
.record(anyString(), anyString(), anyString(), any(),
anyString(), anyString());
- try (MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- doNothing().when(producer).start();
- SendResult sendResult = new SendResult();
- sendResult.setSendStatus(SendStatus.SEND_OK);
- sendResult.setMsgId("msg-1");
- sendResult.setOffsetMsgId("offset-1");
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- SendMessageDTO request = new SendMessageDTO();
- request.setTopic("TopicA");
- request.setBody("hello");
- SendMessageVO result = adminClient.sendMessage(request);
- // The message was already delivered; an audit failure must not
turn this into an error.
- assertThat(result.getMsgId()).isEqualTo("msg-1");
- }
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.SEND_OK);
+ sendResult.setMsgId("msg-1");
+ sendResult.setOffsetMsgId("offset-1");
+ when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+ SendMessageDTO request = new SendMessageDTO();
+ request.setTopic("TopicA");
+ request.setBody("hello");
+ SendMessageVO result = adminClient.sendMessage(request);
+ // The message was already delivered; an audit failure must not turn
this into an error.
+ assertThat(result.getMsgId()).isEqualTo("msg-1");
}
@Test
@@ -758,16 +758,13 @@ class RocketMQAdminClientImplTest {
request.setTopic("TopicA");
request.setBody("\u754c".repeat((4 * 1024 * 1024 / 3) + 1));
- try (MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class)) {
- assertThatThrownBy(() -> adminClient.sendMessage(request))
- .isInstanceOf(BusinessException.class)
- .hasMessageContaining("exceeds the maximum")
- .satisfies(exception -> assertThat(((BusinessException)
exception).getCode()).isEqualTo(400));
+ assertThatThrownBy(() -> adminClient.sendMessage(request))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("exceeds the maximum")
+ .satisfies(exception -> assertThat(((BusinessException)
exception).getCode()).isEqualTo(400));
- assertThat(mockedProducers.constructed()).isEmpty();
- }
verifyNoInteractions(runtimeAdminClientResolver);
+ verifyNoInteractions(clientPool);
verify(properties, never()).getNamesrvAddr();
verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", null,
"Message body size 4194306 exceeds the maximum of 4194304
bytes", "FAILED");
@@ -775,112 +772,75 @@ class RocketMQAdminClientImplTest {
@Test
void sendMessageShouldAllowBodyAtMaximumSize() throws Exception {
- when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
- List<List<?>> constructorArguments = new ArrayList<>();
- try (MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- constructorArguments.add(context.arguments());
- doNothing().when(producer).start();
- SendResult sendResult = new SendResult();
- sendResult.setSendStatus(SendStatus.SEND_OK);
- sendResult.setMsgId("msg-1");
- sendResult.setOffsetMsgId("offset-1");
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- SendMessageDTO request = new SendMessageDTO();
- request.setTopic("TopicA");
- request.setBody("x".repeat(4 * 1024 * 1024));
-
- SendMessageVO result = adminClient.sendMessage(request);
-
- assertThat(result.getMsgId()).isEqualTo("msg-1");
- DefaultMQProducer producer =
mockedProducers.constructed().getFirst();
- ArgumentCaptor<Message> messageCaptor =
ArgumentCaptor.forClass(Message.class);
- verify(producer).send(messageCaptor.capture());
- assertThat(messageCaptor.getValue().getBody()).hasSize(4 * 1024 *
1024);
- assertThat(constructorArguments).singleElement();
- assertThat(constructorArguments.get(0)).hasSize(2);
- assertThat(constructorArguments.get(0).get(1)).isNull();
- }
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.SEND_OK);
+ sendResult.setMsgId("msg-1");
+ sendResult.setOffsetMsgId("offset-1");
+ when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+ SendMessageDTO request = new SendMessageDTO();
+ request.setTopic("TopicA");
+ request.setBody("x".repeat(4 * 1024 * 1024));
+
+ SendMessageVO result = adminClient.sendMessage(request);
+
+ assertThat(result.getMsgId()).isEqualTo("msg-1");
+ ArgumentCaptor<Message> messageCaptor =
ArgumentCaptor.forClass(Message.class);
+ verify(sendProducer).send(messageCaptor.capture());
+ assertThat(messageCaptor.getValue().getBody()).hasSize(4 * 1024 *
1024);
+ verify(clientPool).withProducer(eq("10.0.0.1:9876"), isNull(),
isNull(), any());
}
@Test
- void sendMessageUsesSelectedInstanceEndpointAndCredentialHook() throws
Exception {
- RPCHook credentialHook = mock(RPCHook.class);
- List<List<?>> constructorArguments = new ArrayList<>();
-
when(runtimeAdminClientResolver.resolveEndpoint("instance-a")).thenReturn("10.0.0.2:9876");
-
when(runtimeAdminClientResolver.resolveCredentialHook("instance-a")).thenReturn(credentialHook);
- try (MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- constructorArguments.add(context.arguments());
- doNothing().when(producer).start();
- SendResult sendResult = new SendResult();
- sendResult.setSendStatus(SendStatus.SEND_OK);
- sendResult.setMsgId("msg-1");
- sendResult.setOffsetMsgId("offset-1");
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- SendMessageDTO request = new SendMessageDTO();
- request.setTopic("TopicA");
- request.setBody("hello");
- request.setInstanceId("instance-a");
-
- adminClient.sendMessage(request);
-
- DefaultMQProducer producer =
mockedProducers.constructed().getFirst();
- verify(producer).setNamesrvAddr("10.0.0.2:9876");
- assertThat(constructorArguments).singleElement();
- assertThat(constructorArguments.get(0)).hasSize(2);
-
assertThat(constructorArguments.get(0).get(1)).isSameAs(credentialHook);
- }
- verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
+ void sendMessageUsesSelectedInstancePooledProducer() throws Exception {
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.SEND_OK);
+ sendResult.setMsgId("msg-1");
+ sendResult.setOffsetMsgId("offset-1");
+ when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+ SendMessageDTO request = new SendMessageDTO();
+ request.setTopic("TopicA");
+ request.setBody("hello");
+ request.setInstanceId("instance-a");
+
+ adminClient.sendMessage(request);
+
+ verify(runtimeAdminClientResolver).executeProducer(eq("instance-a"),
any());
+ verify(clientPool, never()).withProducer(any(), any(), any(), any());
}
@Test
void sendMessageShouldRejectNonSuccessfulSendStatus() throws Exception {
- when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
- try (MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- doNothing().when(producer).start();
- SendResult sendResult = new SendResult();
-
sendResult.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT);
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- SendMessageDTO request = new SendMessageDTO();
- request.setTopic("TopicA");
- request.setBody("hello");
-
- assertThatThrownBy(() -> adminClient.sendMessage(request))
- .isInstanceOf(BusinessException.class)
- .hasMessageContaining("FLUSH_DISK_TIMEOUT");
-
- verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA",
null,
- "Message send did not succeed: FLUSH_DISK_TIMEOUT",
"FAILED");
- }
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT);
+ when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+ SendMessageDTO request = new SendMessageDTO();
+ request.setTopic("TopicA");
+ request.setBody("hello");
+
+ assertThatThrownBy(() -> adminClient.sendMessage(request))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("FLUSH_DISK_TIMEOUT");
+
+ verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", null,
+ "Message send did not succeed: FLUSH_DISK_TIMEOUT", "FAILED");
}
@Test
void sendMessageShouldRejectNullSendResult() throws Exception {
- when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
- try (MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- doNothing().when(producer).start();
-
when(producer.send(any(Message.class))).thenReturn(null);
- doNothing().when(producer).shutdown();
- })) {
- SendMessageDTO request = new SendMessageDTO();
- request.setTopic("TopicA");
- request.setBody("hello");
-
- assertThatThrownBy(() -> adminClient.sendMessage(request))
- .isInstanceOf(BusinessException.class)
- .hasMessageContaining("null");
-
- verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA",
null,
- "Message send did not succeed: null", "FAILED");
- }
+ when(sendProducer.send(any(Message.class))).thenReturn(null);
+
+ SendMessageDTO request = new SendMessageDTO();
+ request.setTopic("TopicA");
+ request.setBody("hello");
+
+ assertThatThrownBy(() -> adminClient.sendMessage(request))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("null");
+
+ verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", null,
+ "Message send did not succeed: null", "FAILED");
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
index 5177df3e9..5ba2fb094 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
@@ -300,7 +300,8 @@ class RocketMQClusterProviderTest {
credential.setSecretKey("admin-sk");
adminProperties.getCredentials().put("cluster-admin", credential);
RuntimeAdminClientResolver resolver =
- new RuntimeAdminClientResolver(instanceRepository,
adminFactory, adminProperties);
+ new RuntimeAdminClientResolver(instanceRepository,
adminFactory, adminProperties,
+
org.mockito.Mockito.mock(org.apache.rocketmq.studio.cluster.broker.MqClientPool.class));
when(adminFactory.execute(eq("10.0.0.2:9876"), any(),
eq("cluster-admin"), any()))
.thenAnswer(invocation ->
invocation.<MqAdminExtFactory.AdminAction<Object>>getArgument(3)
.apply(adminExt));
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 607d1bddd..c22b07036 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
@@ -27,10 +27,10 @@ import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
import org.apache.rocketmq.remoting.protocol.body.TopicList;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.exception.BusinessException;
@@ -44,11 +44,9 @@ import org.junit.jupiter.api.Timeout;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
-import org.mockito.MockedConstruction;
import org.mockito.junit.jupiter.MockitoExtension;
import java.nio.charset.StandardCharsets;
-import java.util.ArrayList;
import java.util.Base64;
import java.util.List;
import java.util.Set;
@@ -64,11 +62,7 @@ import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.contains;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.isNull;
-import static org.mockito.Mockito.doNothing;
-import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.lenient;
-import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.mockConstruction;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
@@ -87,6 +81,12 @@ class RocketMQDLQProviderTest {
@Mock
private MQAdminExt adminExt;
+ @Mock
+ private DefaultMQPullConsumer pullConsumer;
+
+ @Mock
+ private DefaultMQProducer dlqProducer;
+
private RocketMQDLQProvider provider;
@BeforeEach
@@ -96,6 +96,18 @@ class RocketMQDLQProviderTest {
MqAdminExtFactory.AdminAction<Object> action =
invocation.getArgument(1);
return action.apply(adminExt);
});
+
lenient().when(runtimeAdminClientResolver.executePullConsumer(anyString(),
any()))
+ .thenAnswer(invocation -> {
+ MqClientPool.ClientAction<DefaultMQPullConsumer, Object>
action =
+ invocation.getArgument(1);
+ return action.apply(pullConsumer);
+ });
+ lenient().when(runtimeAdminClientResolver.executeProducer(anyString(),
any()))
+ .thenAnswer(invocation -> {
+ MqClientPool.ClientAction<DefaultMQProducer, Object>
action =
+ invocation.getArgument(1);
+ return action.apply(dlqProducer);
+ });
provider = new RocketMQDLQProvider(runtimeAdminClientResolver,
auditService);
}
@@ -226,20 +238,12 @@ class RocketMQDLQProviderTest {
@Test
void
resendMessagesShouldNormalizeGroupNameBeforeBuildingDlqTopicAndAuditing()
throws Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(null);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class)) {
- provider.resendMessages("instance-a", " group-a ", 100L, 200L,
"target-topic");
-
- assertThat(mockedConsumers.constructed()).hasSize(1);
-
verify(mockedConsumers.constructed().get(0)).fetchSubscribeMessageQueues(dlqTopic);
- assertThat(mockedProducers.constructed()).isEmpty();
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(null);
+ provider.resendMessages("instance-a", " group-a ", 100L, 200L,
"target-topic");
+
+
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
+ verify(pullConsumer).fetchSubscribeMessageQueues(dlqTopic);
+ verify(runtimeAdminClientResolver,
never()).executeProducer(anyString(), any());
verify(auditService).record(eq("RESEND_DLQ"), eq("DLQ"),
eq("group-a"), eq(null),
contains("group=group-a, dlqTopic=%DLQ%group-a"),
eq("NO_MESSAGES"));
}
@@ -247,30 +251,13 @@ class RocketMQDLQProviderTest {
@Test
void resendMessagesDoesNotPullWhenDlqQueueSetIsNull() throws Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
- List<List<?>> consumerConstructorArguments = new ArrayList<>();
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- consumerConstructorArguments.add(context.arguments());
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(null);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class)) {
- provider.resendMessages("instance-a", "group-a", 100L, 200L,
"target-topic");
-
- assertThat(mockedConsumers.constructed()).hasSize(1);
- assertThat(consumerConstructorArguments).singleElement();
- assertThat(consumerConstructorArguments.get(0)).hasSize(2);
- assertThat(consumerConstructorArguments.get(0).get(1)).isNull();
- 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());
- verify(consumer).shutdown();
- assertThat(mockedProducers.constructed()).isEmpty();
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(null);
+ provider.resendMessages("instance-a", "group-a", 100L, 200L,
"target-topic");
+
+
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
+ verify(pullConsumer).fetchSubscribeMessageQueues(dlqTopic);
+ verify(pullConsumer, never()).pull(any(MessageQueue.class),
anyString(), anyLong(), anyInt());
+ verify(runtimeAdminClientResolver,
never()).executeProducer(anyString(), any());
verify(auditService).record(
eq("RESEND_DLQ"),
eq("DLQ"),
@@ -278,12 +265,10 @@ class RocketMQDLQProviderTest {
isNull(),
contains("matched=0, resent=0, failed=0"),
eq("NO_MESSAGES"));
- verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
- verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
}
@Test
- void resendMessagesUsesSelectedInstanceCredentialHookForScanAndResend()
throws Exception {
+ void resendMessagesUsesPooledClientsForScanAndResend() throws Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
MessageQueue queue = new MessageQueue(dlqTopic, "broker-a", 0);
MessageExt deadLetter = new MessageExt();
@@ -294,56 +279,27 @@ class RocketMQDLQProviderTest {
PullResult pullResult = new PullResult(PullStatus.FOUND, 1L, 0L, 0L,
List.of(deadLetter));
SendResult sendResult = new SendResult();
sendResult.setSendStatus(SendStatus.SEND_OK);
- RPCHook credentialHook = mock(RPCHook.class);
- List<List<?>> consumerConstructorArguments = new ArrayList<>();
- List<List<?>> producerConstructorArguments = new ArrayList<>();
-
when(runtimeAdminClientResolver.resolveCredentialHook("instance-a")).thenReturn(credentialHook);
-
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- consumerConstructorArguments.add(context.arguments());
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
- when(consumer.searchOffset(queue,
100L)).thenReturn(0L);
- when(consumer.searchOffset(queue,
200L)).thenReturn(0L);
- when(consumer.pull(queue, "*", 0L,
32)).thenReturn(pullResult);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- producerConstructorArguments.add(context.arguments());
- doNothing().when(producer).start();
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- provider.resendMessages("instance-a", "group-a", 100L, 200L,
"target-topic");
-
- assertThat(mockedConsumers.constructed()).singleElement();
- assertThat(mockedProducers.constructed()).singleElement();
- }
-
- assertThat(consumerConstructorArguments).singleElement();
-
assertThat(consumerConstructorArguments.get(0).get(1)).isSameAs(credentialHook);
- assertThat(producerConstructorArguments).singleElement();
-
assertThat(producerConstructorArguments.get(0).get(1)).isSameAs(credentialHook);
- verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+ when(pullConsumer.searchOffset(queue, 200L)).thenReturn(0L);
+ when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+ when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+ provider.resendMessages("instance-a", "group-a", 100L, 200L,
"target-topic");
+
+
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
+ verify(runtimeAdminClientResolver).executeProducer(eq("instance-a"),
any());
}
@Test
void resendMessagesRejectsAnAllFailedDlqScan() throws Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doThrow(new IllegalStateException("broker
unavailable")).when(consumer).start();
- doNothing().when(consumer).shutdown();
- })) {
- assertThatThrownBy(() -> provider.resendMessages("instance-a",
"group-a", 100L, 200L, "target-topic"))
- .isInstanceOf(BusinessException.class)
- .hasMessageContaining("Failed to scan DLQ topic " +
dlqTopic)
- .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
-
- verify(mockedConsumers.constructed().get(0)).shutdown();
- }
+ when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+ assertThatThrownBy(() -> provider.resendMessages("instance-a",
"group-a", 100L, 200L, "target-topic"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("Failed to scan DLQ topic " + dlqTopic)
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+
verify(auditService).record(
eq("RESEND_DLQ"),
eq("DLQ"),
@@ -359,25 +315,17 @@ class RocketMQDLQProviderTest {
MessageQueue unavailableQueue = new MessageQueue(dlqTopic, "broker-a",
0);
MessageQueue emptyQueue = new MessageQueue(dlqTopic, "broker-b", 0);
PullResult emptyResult = new PullResult(PullStatus.NO_NEW_MSG, 1L, 0L,
0L, List.of());
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
- when(consumer.fetchSubscribeMessageQueues(dlqTopic))
- .thenReturn(Set.of(unavailableQueue,
emptyQueue));
- when(consumer.searchOffset(eq(unavailableQueue),
anyLong()))
- .thenThrow(new IllegalStateException("broker
unavailable"));
- when(consumer.searchOffset(eq(emptyQueue),
anyLong())).thenReturn(0L);
- when(consumer.pull(eq(emptyQueue), eq("*"), eq(0L),
eq(32))).thenReturn(emptyResult);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class)) {
- assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
- .extracting("matched", "resent", "failed", "outcome",
"scanIncomplete", "failedQueueCount")
- .containsExactly(0, 0, 0, "PARTIAL", true, 1);
-
- assertThat(mockedProducers.constructed()).isEmpty();
- }
+ when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic))
+ .thenReturn(Set.of(unavailableQueue, emptyQueue));
+ when(pullConsumer.searchOffset(eq(unavailableQueue), anyLong()))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+ when(pullConsumer.searchOffset(eq(emptyQueue),
anyLong())).thenReturn(0L);
+ when(pullConsumer.pull(eq(emptyQueue), eq("*"), eq(0L),
eq(32))).thenReturn(emptyResult);
+ assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
+ .extracting("matched", "resent", "failed", "outcome",
"scanIncomplete", "failedQueueCount")
+ .containsExactly(0, 0, 0, "PARTIAL", true, 1);
+
+ verify(runtimeAdminClientResolver,
never()).executeProducer(anyString(), any());
verify(auditService).record(
eq("RESEND_DLQ"),
eq("DLQ"),
@@ -393,21 +341,13 @@ class RocketMQDLQProviderTest {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
MessageQueue queue = new MessageQueue(dlqTopic, "broker-a", 0);
PullResult stalledResult = new PullResult(PullStatus.FOUND, 10, 0, 10,
List.of());
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
- when(consumer.searchOffset(eq(queue),
anyLong())).thenReturn(10L);
- when(consumer.pull(eq(queue), eq("*"), eq(10L),
eq(32))).thenReturn(stalledResult);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class)) {
- provider.resendMessages("instance-a", "group-a", 100L, 200L,
"target-topic");
-
- verify(mockedConsumers.constructed().get(0), times(1)).pull(queue,
"*", 10L, 32);
- assertThat(mockedProducers.constructed()).isEmpty();
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(eq(queue), anyLong())).thenReturn(10L);
+ when(pullConsumer.pull(eq(queue), eq("*"), eq(10L),
eq(32))).thenReturn(stalledResult);
+ provider.resendMessages("instance-a", "group-a", 100L, 200L,
"target-topic");
+
+ verify(pullConsumer, times(1)).pull(queue, "*", 10L, 32);
+ verify(runtimeAdminClientResolver,
never()).executeProducer(anyString(), any());
}
@Test
@@ -425,33 +365,22 @@ class RocketMQDLQProviderTest {
SendResult sendResult = new SendResult();
sendResult.setSendStatus(SendStatus.SEND_OK);
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
- when(consumer.searchOffset(queue,
100L)).thenReturn(10L);
- when(consumer.searchOffset(queue,
200L)).thenReturn(50L);
- when(consumer.pull(queue, "*", 10L,
32)).thenReturn(illegalOffset);
- when(consumer.pull(queue, "*", 20L,
32)).thenReturn(foundAfterCorrection);
- when(consumer.pull(queue, "*", 40L,
32)).thenReturn(endOfQueue);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- doNothing().when(producer).start();
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "orders"))
- .extracting("matched", "resent", "failed", "outcome")
- .containsExactly(1, 1, 0, "SUCCESS");
-
- DefaultMQPullConsumer consumer =
mockedConsumers.constructed().get(0);
- verify(consumer).pull(queue, "*", 10L, 32);
- verify(consumer).pull(queue, "*", 20L, 32);
- verify(consumer).pull(queue, "*", 40L, 32);
-
verify(mockedProducers.constructed().get(0)).send(any(Message.class));
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(queue, 100L)).thenReturn(10L);
+ when(pullConsumer.searchOffset(queue, 200L)).thenReturn(50L);
+ when(pullConsumer.pull(queue, "*", 10L, 32)).thenReturn(illegalOffset);
+ when(pullConsumer.pull(queue, "*", 20L,
32)).thenReturn(foundAfterCorrection);
+ when(pullConsumer.pull(queue, "*", 40L, 32)).thenReturn(endOfQueue);
+ when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+ assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "orders"))
+ .extracting("matched", "resent", "failed", "outcome")
+ .containsExactly(1, 1, 0, "SUCCESS");
+
+ DefaultMQPullConsumer consumer = pullConsumer;
+ verify(consumer).pull(queue, "*", 10L, 32);
+ verify(consumer).pull(queue, "*", 20L, 32);
+ verify(consumer).pull(queue, "*", 40L, 32);
+ verify(dlqProducer).send(any(Message.class));
verify(auditService).record(
eq("RESEND_DLQ"),
eq("DLQ"),
@@ -474,29 +403,18 @@ class RocketMQDLQProviderTest {
SendResult sendResult = new SendResult();
sendResult.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT);
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
- when(consumer.searchOffset(queue,
100L)).thenReturn(0L);
- when(consumer.searchOffset(queue,
200L)).thenReturn(0L);
- when(consumer.pull(queue, "*", 0L,
32)).thenReturn(pullResult);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- doNothing().when(producer).start();
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
- .extracting("matched", "resent", "failed", "outcome")
- .containsExactly(1, 0, 1, "FAILED");
-
- assertThat(mockedConsumers.constructed()).hasSize(1);
- assertThat(mockedProducers.constructed()).hasSize(1);
-
verify(mockedProducers.constructed().get(0)).send(any(Message.class));
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+ when(pullConsumer.searchOffset(queue, 200L)).thenReturn(0L);
+ when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+ when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+ assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
+ .extracting("matched", "resent", "failed", "outcome")
+ .containsExactly(1, 0, 1, "FAILED");
+
+
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
+ verify(runtimeAdminClientResolver).executeProducer(eq("instance-a"),
any());
+ verify(dlqProducer).send(any(Message.class));
verify(auditService).record(
eq("RESEND_DLQ"),
eq("DLQ"),
@@ -522,35 +440,24 @@ class RocketMQDLQProviderTest {
SendResult sendResult = new SendResult();
sendResult.setSendStatus(SendStatus.SEND_OK);
- try (MockedConstruction<DefaultMQPullConsumer> ignoredConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
- when(consumer.searchOffset(queue,
100L)).thenReturn(0L);
- when(consumer.searchOffset(queue,
200L)).thenReturn(0L);
- when(consumer.pull(queue, "*", 0L,
32)).thenReturn(pullResult);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- doNothing().when(producer).start();
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
- .extracting("matched", "resent", "failed", "outcome")
- .containsExactly(1, 1, 0, "SUCCESS");
-
- ArgumentCaptor<Message> messageCaptor =
ArgumentCaptor.forClass(Message.class);
-
verify(mockedProducers.constructed().get(0)).send(messageCaptor.capture());
- Message resent = messageCaptor.getValue();
- assertThat(resent.getProperties()).containsEntry("traceId",
"trace-123")
- .containsEntry("tenantId", "tenant-a")
- .containsEntry("studio_dlq_origin_message_id",
"msg-with-properties")
- .containsEntry("studio_dlq_origin_topic", dlqTopic);
- assertThat(resent.getProperties())
- .doesNotContainEntry(MessageConst.PROPERTY_REAL_TOPIC,
"system-topic-should-not-copy");
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+ when(pullConsumer.searchOffset(queue, 200L)).thenReturn(0L);
+ when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+ when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+ assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
+ .extracting("matched", "resent", "failed", "outcome")
+ .containsExactly(1, 1, 0, "SUCCESS");
+
+ ArgumentCaptor<Message> messageCaptor =
ArgumentCaptor.forClass(Message.class);
+ verify(dlqProducer).send(messageCaptor.capture());
+ Message resent = messageCaptor.getValue();
+ assertThat(resent.getProperties()).containsEntry("traceId",
"trace-123")
+ .containsEntry("tenantId", "tenant-a")
+ .containsEntry("studio_dlq_origin_message_id",
"msg-with-properties")
+ .containsEntry("studio_dlq_origin_topic", dlqTopic);
+ assertThat(resent.getProperties())
+ .doesNotContainEntry(MessageConst.PROPERTY_REAL_TOPIC,
"system-topic-should-not-copy");
}
@Test
@@ -571,27 +478,16 @@ class RocketMQDLQProviderTest {
SendResult sendResult = new SendResult();
sendResult.setSendStatus(SendStatus.SEND_OK);
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
- when(consumer.searchOffset(queue,
100L)).thenReturn(0L);
- when(consumer.searchOffset(queue,
200L)).thenReturn(5001L);
- when(consumer.pull(queue, "*", 0L,
32)).thenReturn(pullResult);
- doNothing().when(consumer).shutdown();
- });
- MockedConstruction<DefaultMQProducer> mockedProducers =
- mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
- doNothing().when(producer).start();
-
when(producer.send(any(Message.class))).thenReturn(sendResult);
- doNothing().when(producer).shutdown();
- })) {
- assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
- .extracting("matched", "resent", "failed", "outcome",
"scanIncomplete", "failedQueueCount")
- .containsExactly(5000, 5000, 0, "PARTIAL", true, 0);
-
- verify(mockedProducers.constructed().get(0),
times(5000)).send(any(Message.class));
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+ when(pullConsumer.searchOffset(queue, 200L)).thenReturn(5001L);
+ when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+ when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+ assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
+ .extracting("matched", "resent", "failed", "outcome",
"scanIncomplete", "failedQueueCount")
+ .containsExactly(5000, 5000, 0, "PARTIAL", true, 0);
+
+ verify(dlqProducer, times(5000)).send(any(Message.class));
verify(auditService).record(
eq("RESEND_DLQ"),
eq("DLQ"),
@@ -601,15 +497,6 @@ class RocketMQDLQProviderTest {
eq("PARTIAL"));
}
- @Test
- void createsUniqueProducerGroupsForConcurrentDlqResends() {
- String first = RocketMQDLQProvider.nextResendProducerGroup();
- String second = RocketMQDLQProvider.nextResendProducerGroup();
-
- assertThat(first).startsWith("studio-dlq-resend-");
-
assertThat(second).startsWith("studio-dlq-resend-").isNotEqualTo(first);
- }
-
@Test
void exportMessagesReturnsMappedDeadLetterMessages() throws Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
@@ -623,48 +510,36 @@ class RocketMQDLQProviderTest {
deadLetter.setKeys("key-a,key-b");
deadLetter.setBody("hello dlq".getBytes(StandardCharsets.UTF_8));
PullResult pullResult = new PullResult(PullStatus.FOUND, 1L, 0L, 0L,
List.of(deadLetter));
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
- when(consumer.searchOffset(eq(queue),
anyLong())).thenReturn(0L);
- when(consumer.pull(eq(queue), eq("*"), eq(0L),
eq(32))).thenReturn(pullResult);
- doNothing().when(consumer).shutdown();
- })) {
- List<DLQMessageVO> exported =
- provider.exportMessages("instance-a", "group-a", 100L,
200L, 1000);
-
- assertThat(exported).hasSize(1);
- DLQMessageVO vo = exported.get(0);
- assertThat(vo.getMsgId()).isEqualTo("msg-1");
- assertThat(vo.getTopic()).isEqualTo(dlqTopic);
- assertThat(vo.getQueueId()).isEqualTo(0);
- assertThat(vo.getOffset()).isEqualTo(5L);
- assertThat(vo.getStoreTime()).isEqualTo(150L);
- assertThat(vo.getKeys()).isEqualTo("key-a,key-b");
- assertThat(vo.getBody()).isEqualTo("hello dlq");
- assertThat(vo.getBodyBase64())
- .isEqualTo(Base64.getEncoder().encodeToString("hello
dlq".getBytes(StandardCharsets.UTF_8)));
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(eq(queue), anyLong())).thenReturn(0L);
+ when(pullConsumer.pull(eq(queue), eq("*"), eq(0L),
eq(32))).thenReturn(pullResult);
+ List<DLQMessageVO> exported =
+ provider.exportMessages("instance-a", "group-a", 100L, 200L,
1000);
+
+ assertThat(exported).hasSize(1);
+ DLQMessageVO vo = exported.get(0);
+ assertThat(vo.getMsgId()).isEqualTo("msg-1");
+ assertThat(vo.getTopic()).isEqualTo(dlqTopic);
+ assertThat(vo.getQueueId()).isEqualTo(0);
+ assertThat(vo.getOffset()).isEqualTo(5L);
+ assertThat(vo.getStoreTime()).isEqualTo(150L);
+ assertThat(vo.getKeys()).isEqualTo("key-a,key-b");
+ assertThat(vo.getBody()).isEqualTo("hello dlq");
+ assertThat(vo.getBodyBase64())
+ .isEqualTo(Base64.getEncoder().encodeToString("hello
dlq".getBytes(StandardCharsets.UTF_8)));
}
@Test
void exportMessagesHonorsMaxCountCap() throws Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
MessageQueue queue = new MessageQueue(dlqTopic, "broker-a", 0);
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
- when(consumer.searchOffset(eq(queue),
anyLong())).thenReturn(0L);
- when(consumer.pull(eq(queue), eq("*"), eq(0L),
eq(32)))
- .thenReturn(new
PullResult(PullStatus.NO_NEW_MSG, 1L, 0L, 0L, List.of()));
- doNothing().when(consumer).shutdown();
- })) {
- // maxCount=0 falls back to the hard cap instead of failing; scan
still completes.
- List<DLQMessageVO> exported =
- provider.exportMessages("instance-a", "group-a", 100L,
200L, 0);
- assertThat(exported).isEmpty();
- }
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(eq(queue), anyLong())).thenReturn(0L);
+ when(pullConsumer.pull(eq(queue), eq("*"), eq(0L), eq(32)))
+ .thenReturn(new PullResult(PullStatus.NO_NEW_MSG, 1L, 0L, 0L,
List.of()));
+ // maxCount=0 falls back to the hard cap instead of failing; scan
still completes.
+ List<DLQMessageVO> exported =
+ provider.exportMessages("instance-a", "group-a", 100L, 200L,
0);
+ assertThat(exported).isEmpty();
}
}
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 2eb00365b..2501eda4b 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
@@ -27,10 +27,10 @@ import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.message.MessageId;
import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
@@ -45,7 +45,6 @@ import org.junit.jupiter.api.Timeout;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
-import org.mockito.MockedConstruction;
import org.mockito.junit.jupiter.MockitoExtension;
import java.net.InetSocketAddress;
@@ -67,11 +66,8 @@ import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
-import static org.mockito.Mockito.doNothing;
-import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.mockConstruction;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
@@ -86,6 +82,9 @@ class RocketMQMessageProviderTest {
@Mock
private RuntimeAdminClientResolver runtimeAdminClientResolver;
+ @Mock
+ private DefaultMQPullConsumer pullConsumer;
+
private RocketMQMessageProvider provider;
@BeforeEach
@@ -105,61 +104,36 @@ class RocketMQMessageProviderTest {
});
lenient().when(adminExt.examineBrokerClusterInfo())
.thenReturn(clusterInfoWithBrokerAddresses("172.30.10.100:10911"));
+
lenient().when(runtimeAdminClientResolver.executePullConsumer(anyString(),
any()))
+ .thenAnswer(invocation -> {
+ MqClientPool.ClientAction<DefaultMQPullConsumer, Object>
action =
+ invocation.getArgument(1);
+ return action.apply(pullConsumer);
+ });
provider = new RocketMQMessageProvider(runtimeAdminClientResolver);
}
@Test
void queryByTopicReturnsEmptyListWhenQueueSetIsNull() throws Exception {
- List<List<?>> constructorArguments = new ArrayList<>();
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- constructorArguments.add(context.arguments());
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(null);
- doNothing().when(consumer).shutdown();
- })) {
- List<MessageRecordVO> messages =
provider.queryMessages("instance-a", "TopicA", null, null, null, 100L, 200L);
-
- assertThat(messages).isEmpty();
- assertThat(mockedConsumers.constructed()).hasSize(1);
- DefaultMQPullConsumer consumer =
mockedConsumers.constructed().get(0);
- assertThat(constructorArguments).singleElement();
- assertThat(constructorArguments.get(0)).hasSize(2);
-
assertThat(constructorArguments.get(0).get(0)).isEqualTo("studio-msg-query-group");
- assertThat(constructorArguments.get(0).get(1)).isNull();
- verify(consumer).setNamesrvAddr("namesrv-a:9876");
- verify(consumer).start();
- verify(consumer).fetchSubscribeMessageQueues("TopicA");
- verify(consumer, never()).pull(any(MessageQueue.class),
anyString(), anyLong(), anyInt());
- verify(consumer).shutdown();
- }
- verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
- verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
+
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(null);
+
+ List<MessageRecordVO> messages = provider.queryMessages("instance-a",
"TopicA", null, null, null, 100L, 200L);
+
+ assertThat(messages).isEmpty();
+ verify(pullConsumer).fetchSubscribeMessageQueues("TopicA");
+ verify(pullConsumer, never()).pull(any(MessageQueue.class),
anyString(), anyLong(), anyInt());
+
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
}
@Test
- void queryByTopicUsesSelectedInstanceCredentialHook() throws Exception {
- RPCHook credentialHook = mock(RPCHook.class);
- List<List<?>> constructorArguments = new ArrayList<>();
-
when(runtimeAdminClientResolver.resolveCredentialHook("instance-a")).thenReturn(credentialHook);
-
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- constructorArguments.add(context.arguments());
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of());
- doNothing().when(consumer).shutdown();
- })) {
- provider.queryMessages("instance-a", "TopicA", null, null, null,
100L, 200L);
-
- assertThat(mockedConsumers.constructed()).singleElement();
- assertThat(constructorArguments).singleElement();
- assertThat(constructorArguments.get(0)).hasSize(2);
-
assertThat(constructorArguments.get(0).get(0)).isEqualTo("studio-msg-query-group");
-
assertThat(constructorArguments.get(0).get(1)).isSameAs(credentialHook);
- }
- verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
+ void queryByTopicReturnsEmptyListWhenNoQueuesExist() throws Exception {
+
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of());
+
+ List<MessageRecordVO> messages = provider.queryMessages("instance-a",
"TopicA", null, null, null, 100L, 200L);
+
+ assertThat(messages).isEmpty();
+
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
}
@Test
@@ -170,7 +144,6 @@ class RocketMQMessageProviderTest {
.hasMessage("Message query start time must be before end time")
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(400));
- verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
verify(adminExt, never()).queryMessage(anyString(), anyString(),
anyInt(), anyLong(), anyLong());
}
@@ -183,7 +156,6 @@ class RocketMQMessageProviderTest {
.hasMessage("Message query start time must be before end time")
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(400));
- verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
verify(adminExt, never()).queryMessage(anyString(), anyString(),
anyInt(), anyLong(), anyLong());
}
@@ -244,17 +216,14 @@ class RocketMQMessageProviderTest {
@Test
void queryByTopicSurfacesPullConsumerFailure() throws Exception {
- try (MockedConstruction<DefaultMQPullConsumer> ignored =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doThrow(new IllegalStateException("nameserver
unavailable")).when(consumer).start();
- doNothing().when(consumer).shutdown();
- })) {
- assertThatThrownBy(() -> provider.queryMessages(
- "instance-a", "TopicA", null, null, null, 100L, 200L))
- .isInstanceOf(BusinessException.class)
- .hasMessage("Failed to query messages by topic: nameserver
unavailable")
- .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
- }
+ when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+ .thenThrow(new IllegalStateException("nameserver
unavailable"));
+
+ assertThatThrownBy(() -> provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 200L))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to query messages by topic: nameserver
unavailable")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
}
@Test
@@ -262,20 +231,15 @@ class RocketMQMessageProviderTest {
void queryByTopicStopsWhenPullOffsetDoesNotAdvance() throws Exception {
MessageQueue queue = new MessageQueue("TopicA", "broker-a", 0);
PullResult stalledResult = new PullResult(PullStatus.FOUND, 10, 0, 10,
List.of());
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
- mockQueueWindow(consumer, queue, 100L, 200L, 10L,
10L, 11L, 11L);
- when(consumer.pull(eq(queue), eq("*"), eq(10L),
eq(32))).thenReturn(stalledResult);
- doNothing().when(consumer).shutdown();
- })) {
- List<MessageRecordVO> messages = provider.queryMessages(
- "instance-a", "TopicA", null, null, null, 100L, 200L);
-
- assertThat(messages).isEmpty();
- verify(mockedConsumers.constructed().get(0), times(1)).pull(queue,
"*", 10L, 32);
- }
+
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+ mockQueueWindow(pullConsumer, queue, 100L, 200L, 10L, 10L, 11L, 11L);
+ when(pullConsumer.pull(eq(queue), eq("*"), eq(10L),
eq(32))).thenReturn(stalledResult);
+
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 200L);
+
+ assertThat(messages).isEmpty();
+ verify(pullConsumer, times(1)).pull(queue, "*", 10L, 32);
}
@Test
@@ -286,23 +250,18 @@ class RocketMQMessageProviderTest {
List.of(topicMessage("older", 150L)));
PullResult newerPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
List.of(topicMessage("newer", 250L)));
- try (MockedConstruction<DefaultMQPullConsumer> ignored =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
- when(consumer.fetchSubscribeMessageQueues("TopicA"))
- .thenReturn(new
LinkedHashSet<>(List.of(olderQueue, newerQueue)));
- mockQueueWindow(consumer, olderQueue, 100L, 300L,
10L, 10L, 11L, 11L);
- mockQueueWindow(consumer, newerQueue, 100L, 300L,
10L, 10L, 11L, 11L);
- when(consumer.pull(olderQueue, "*", 10L,
32)).thenReturn(olderPullResult);
- when(consumer.pull(newerQueue, "*", 10L,
32)).thenReturn(newerPullResult);
- doNothing().when(consumer).shutdown();
- })) {
- List<MessageRecordVO> messages = provider.queryMessages(
- "instance-a", "TopicA", null, null, null, 100L, 300L);
-
- assertThat(messages).extracting(MessageRecordVO::getMsgId)
- .containsExactly("newer", "older");
- }
+ when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+ .thenReturn(new LinkedHashSet<>(List.of(olderQueue,
newerQueue)));
+ mockQueueWindow(pullConsumer, olderQueue, 100L, 300L, 10L, 10L, 11L,
11L);
+ mockQueueWindow(pullConsumer, newerQueue, 100L, 300L, 10L, 10L, 11L,
11L);
+ when(pullConsumer.pull(olderQueue, "*", 10L,
32)).thenReturn(olderPullResult);
+ when(pullConsumer.pull(newerQueue, "*", 10L,
32)).thenReturn(newerPullResult);
+
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+ assertThat(messages).extracting(MessageRecordVO::getMsgId)
+ .containsExactly("newer", "older");
}
@Test
@@ -316,25 +275,19 @@ class RocketMQMessageProviderTest {
PullResult illegalOffset = new PullResult(PullStatus.OFFSET_ILLEGAL,
20L, 0L, 30L, null);
PullResult foundAfterCorrection = new PullResult(PullStatus.FOUND,
40L, 20L, 30L, List.of(message));
PullResult endOfQueue = new PullResult(PullStatus.NO_NEW_MSG, 50L,
40L, 40L, List.of());
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
- mockQueueWindow(consumer, queue, 100L, 200L, 10L,
10L, 50L, 50L);
- when(consumer.pull(eq(queue), eq("*"), eq(10L),
eq(32))).thenReturn(illegalOffset);
- when(consumer.pull(eq(queue), eq("*"), eq(20L),
eq(32))).thenReturn(foundAfterCorrection);
- when(consumer.pull(eq(queue), eq("*"), eq(40L),
eq(32))).thenReturn(endOfQueue);
- doNothing().when(consumer).shutdown();
- })) {
- List<MessageRecordVO> messages = provider.queryMessages(
- "instance-a", "TopicA", null, null, null, 100L, 200L);
-
-
assertThat(messages).extracting(MessageRecordVO::getMsgId).containsExactly("msg-after-correction");
- DefaultMQPullConsumer consumer =
mockedConsumers.constructed().get(0);
- verify(consumer).pull(queue, "*", 10L, 32);
- verify(consumer).pull(queue, "*", 20L, 32);
- verify(consumer).pull(queue, "*", 40L, 32);
- }
+
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+ mockQueueWindow(pullConsumer, queue, 100L, 200L, 10L, 10L, 50L, 50L);
+ when(pullConsumer.pull(eq(queue), eq("*"), eq(10L),
eq(32))).thenReturn(illegalOffset);
+ when(pullConsumer.pull(eq(queue), eq("*"), eq(20L),
eq(32))).thenReturn(foundAfterCorrection);
+ when(pullConsumer.pull(eq(queue), eq("*"), eq(40L),
eq(32))).thenReturn(endOfQueue);
+
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 200L);
+
+
assertThat(messages).extracting(MessageRecordVO::getMsgId).containsExactly("msg-after-correction");
+ verify(pullConsumer).pull(queue, "*", 10L, 32);
+ verify(pullConsumer).pull(queue, "*", 20L, 32);
+ verify(pullConsumer).pull(queue, "*", 40L, 32);
}
@Test
@@ -347,25 +300,20 @@ class RocketMQMessageProviderTest {
.toList());
PullResult newerPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
List.of(topicMessage("newer", 250L)));
- try (MockedConstruction<DefaultMQPullConsumer> ignored =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
- when(consumer.fetchSubscribeMessageQueues("TopicA"))
- .thenReturn(new
LinkedHashSet<>(List.of(olderQueue, newerQueue)));
- mockQueueWindow(consumer, olderQueue, 100L, 300L,
10L, 10L, 11L, 11L);
- mockQueueWindow(consumer, newerQueue, 100L, 300L,
10L, 10L, 11L, 11L);
- when(consumer.pull(olderQueue, "*", 10L,
32)).thenReturn(olderPullResult);
- when(consumer.pull(newerQueue, "*", 10L,
32)).thenReturn(newerPullResult);
- doNothing().when(consumer).shutdown();
- })) {
- List<MessageRecordVO> messages = provider.queryMessages(
- "instance-a", "TopicA", null, null, null, 100L, 300L);
-
- assertThat(messages).hasSize(200);
- assertThat(messages.get(0).getMsgId()).isEqualTo("newer");
- assertThat(messages).extracting(MessageRecordVO::getMsgId)
- .contains("newer");
- }
+ when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+ .thenReturn(new LinkedHashSet<>(List.of(olderQueue,
newerQueue)));
+ mockQueueWindow(pullConsumer, olderQueue, 100L, 300L, 10L, 10L, 11L,
11L);
+ mockQueueWindow(pullConsumer, newerQueue, 100L, 300L, 10L, 10L, 11L,
11L);
+ when(pullConsumer.pull(olderQueue, "*", 10L,
32)).thenReturn(olderPullResult);
+ when(pullConsumer.pull(newerQueue, "*", 10L,
32)).thenReturn(newerPullResult);
+
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+ assertThat(messages).hasSize(200);
+ assertThat(messages.get(0).getMsgId()).isEqualTo("newer");
+ assertThat(messages).extracting(MessageRecordVO::getMsgId)
+ .contains("newer");
}
@Test
@@ -376,23 +324,18 @@ class RocketMQMessageProviderTest {
List.of(topicMessage("message-a", 250L)));
PullResult secondPullResult = new PullResult(PullStatus.FOUND, 11L,
10L, 10L,
List.of(topicMessage("message-b", 250L)));
- try (MockedConstruction<DefaultMQPullConsumer> ignored =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
- when(consumer.fetchSubscribeMessageQueues("TopicA"))
- .thenReturn(new
LinkedHashSet<>(List.of(firstQueue, secondQueue)));
- mockQueueWindow(consumer, firstQueue, 100L, 300L,
10L, 10L, 11L, 11L);
- mockQueueWindow(consumer, secondQueue, 100L, 300L,
10L, 10L, 11L, 11L);
- when(consumer.pull(firstQueue, "*", 10L,
32)).thenReturn(firstPullResult);
- when(consumer.pull(secondQueue, "*", 10L,
32)).thenReturn(secondPullResult);
- doNothing().when(consumer).shutdown();
- })) {
- List<MessageRecordVO> messages = provider.queryMessages(
- "instance-a", "TopicA", null, null, null, 100L, 300L);
-
- assertThat(messages).extracting(MessageRecordVO::getMsgId)
- .containsExactly("message-b", "message-a");
- }
+ when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+ .thenReturn(new LinkedHashSet<>(List.of(firstQueue,
secondQueue)));
+ mockQueueWindow(pullConsumer, firstQueue, 100L, 300L, 10L, 10L, 11L,
11L);
+ mockQueueWindow(pullConsumer, secondQueue, 100L, 300L, 10L, 10L, 11L,
11L);
+ when(pullConsumer.pull(firstQueue, "*", 10L,
32)).thenReturn(firstPullResult);
+ when(pullConsumer.pull(secondQueue, "*", 10L,
32)).thenReturn(secondPullResult);
+
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+ assertThat(messages).extracting(MessageRecordVO::getMsgId)
+ .containsExactly("message-b", "message-a");
}
@Test
@@ -637,25 +580,19 @@ class RocketMQMessageProviderTest {
long begin = 10_000L;
long end = 49_999L;
long maxOffsetExclusive = 40_000L;
- try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
- mockQueueWindow(consumer, queue, begin, end, 0L, 0L,
maxOffsetExclusive, maxOffsetExclusive);
- when(consumer.pull(eq(queue), eq("*"), anyLong(),
eq(32)))
- .thenAnswer(invocation ->
topicPullBatch(invocation.getArgument(2), maxOffsetExclusive, begin));
- doNothing().when(consumer).shutdown();
- })) {
- List<MessageRecordVO> messages = provider.queryMessages(
- "instance-a", "TopicA", null, null, null, begin, end);
-
- assertThat(messages).hasSize(200);
- assertThat(messages.get(0).getMsgId()).isEqualTo("msg-39999");
- assertThat(messages.get(199).getMsgId()).isEqualTo("msg-39800");
- DefaultMQPullConsumer consumer =
mockedConsumers.constructed().get(0);
- verify(consumer).pull(queue, "*", 8_000L, 32);
- verify(consumer, never()).pull(queue, "*", 0L, 32);
- }
+
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+ mockQueueWindow(pullConsumer, queue, begin, end, 0L, 0L,
maxOffsetExclusive, maxOffsetExclusive);
+ when(pullConsumer.pull(eq(queue), eq("*"), anyLong(), eq(32)))
+ .thenAnswer(invocation ->
topicPullBatch(invocation.getArgument(2), maxOffsetExclusive, begin));
+
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, begin, end);
+
+ assertThat(messages).hasSize(200);
+ assertThat(messages.get(0).getMsgId()).isEqualTo("msg-39999");
+ assertThat(messages.get(199).getMsgId()).isEqualTo("msg-39800");
+ verify(pullConsumer).pull(queue, "*", 8_000L, 32);
+ verify(pullConsumer, never()).pull(queue, "*", 0L, 32);
}
@Test
@@ -665,27 +602,22 @@ class RocketMQMessageProviderTest {
long end = 69_999L;
long maxOffsetExclusive = 50_000L;
List<Long> pulledOffsets = new ArrayList<>();
- try (MockedConstruction<DefaultMQPullConsumer> ignored =
- mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
- doNothing().when(consumer).start();
-
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
- mockQueueWindow(consumer, queue, begin, end, 0L, 0L,
maxOffsetExclusive, maxOffsetExclusive);
- when(consumer.pull(eq(queue), eq("*"), anyLong(),
eq(32)))
- .thenAnswer(invocation -> {
- long offset = invocation.getArgument(2);
- pulledOffsets.add(offset);
- return topicPullBatch(offset,
maxOffsetExclusive, begin);
- });
- doNothing().when(consumer).shutdown();
- })) {
- List<MessageRecordVO> messages = provider.queryMessages(
- "instance-a", "TopicA", null, null, null, begin, end);
-
- assertThat(messages).hasSize(200);
- assertThat(pulledOffsets).isNotEmpty();
- assertThat(pulledOffsets.get(0)).isEqualTo(18_000L);
- assertThat(pulledOffsets).allMatch(offset -> offset >= 18_000L);
- }
+
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+ mockQueueWindow(pullConsumer, queue, begin, end, 0L, 0L,
maxOffsetExclusive, maxOffsetExclusive);
+ when(pullConsumer.pull(eq(queue), eq("*"), anyLong(), eq(32)))
+ .thenAnswer(invocation -> {
+ long offset = invocation.getArgument(2);
+ pulledOffsets.add(offset);
+ return topicPullBatch(offset, maxOffsetExclusive, begin);
+ });
+
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, begin, end);
+
+ assertThat(messages).hasSize(200);
+ assertThat(pulledOffsets).isNotEmpty();
+ assertThat(pulledOffsets.get(0)).isEqualTo(18_000L);
+ assertThat(pulledOffsets).allMatch(offset -> offset >= 18_000L);
}
private MQClientAPIImpl mockOffsetLookupClient() {
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientNameTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientNameTest.java
deleted file mode 100644
index b57a322ce..000000000
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientNameTest.java
+++ /dev/null
@@ -1,37 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package org.apache.rocketmq.studio.provider.apache;
-
-import static org.assertj.core.api.Assertions.assertThat;
-
-import java.util.Set;
-import java.util.stream.Collectors;
-import java.util.stream.IntStream;
-import org.junit.jupiter.api.Test;
-
-class ShortLivedClientNameTest {
-
- @Test
- void generatesUniqueNamesWithoutClockDelays() {
- Set<String> names = IntStream.range(0, 100)
- .parallel()
- .mapToObj(ignored ->
RocketMQAdminClientImpl.nextMessageSenderGroup())
- .collect(Collectors.toSet());
-
- assertThat(names).hasSize(100).allMatch(name ->
name.startsWith("studio-msg-sender-"));
- }
-}