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 f73f91b92 fix(dlq): resolve selected messages by id (#2636)
f73f91b92 is described below
commit f73f91b925894239d0aa9bd2f16fc3487d794a8d
Author: btlqql <[email protected]>
AuthorDate: Wed Sep 2 17:03:54 2026 +0800
fix(dlq): resolve selected messages by id (#2636)
---
.../provider/apache/RocketMQDLQProvider.java | 27 ++++++++-----
.../provider/apache/RocketMQDLQProviderTest.java | 47 ++++++++++++++++++++++
2 files changed, 65 insertions(+), 9 deletions(-)
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 5f2ec556c..e6dd8f42a 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
@@ -58,6 +58,7 @@ import java.util.ArrayList;
import java.util.Base64;
import java.util.Collections;
import java.util.Comparator;
+import java.util.LinkedHashSet;
import java.util.Locale;
import java.util.List;
import java.util.Map;
@@ -241,7 +242,7 @@ public class RocketMQDLQProvider implements DLQProvider {
throw new BusinessException(400, "At least one msgId is required
for selected DLQ resend");
}
groupName = groupName.trim();
- Set<String> selected = new java.util.HashSet<>();
+ Set<String> selected = new LinkedHashSet<>();
for (String msgId : msgIds) {
if (StringUtils.hasText(msgId)) {
selected.add(msgId.trim());
@@ -253,14 +254,21 @@ public class RocketMQDLQProvider implements DLQProvider {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + groupName;
- // Scan a wide window so the selected messages are found regardless of
their store time.
- long end = System.currentTimeMillis();
- long begin = end - 7 * 24 * ONE_HOUR_MILLIS;
- DeadLetterScanResult scanResult =
- collectDeadLetters(instanceId, dlqTopic, begin, end,
RESEND_HARD_CAP);
- List<MessageExt> deadLetters = scanResult.messages().stream()
- .filter(message -> selected.contains(message.getMsgId()))
- .toList();
+ List<MessageExt> deadLetters =
runtimeAdminClientResolver.execute(instanceId, admin -> {
+ List<MessageExt> resolved = new ArrayList<>(selected.size());
+ for (String msgId : selected) {
+ try {
+ MessageExt deadLetter = admin.viewMessage(dlqTopic, msgId);
+ if (deadLetter != null) {
+ resolved.add(deadLetter);
+ }
+ } catch (Exception e) {
+ log.warn("Failed to resolve selected dead letter {} from
{}: {}",
+ msgId, dlqTopic, e.getMessage());
+ }
+ }
+ return resolved;
+ });
int[] counts = {0, 0};
if (!deadLetters.isEmpty()) {
try {
@@ -293,6 +301,7 @@ public class RocketMQDLQProvider implements DLQProvider {
.resent(resent)
.failed(failed)
.outcome(outcome)
+ .scanIncomplete(!foundAll)
.build();
}
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 3aef77a7c..12dc1a026 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
@@ -461,6 +461,53 @@ class RocketMQDLQProviderTest {
.doesNotContainEntry(MessageConst.PROPERTY_REAL_TOPIC,
"system-topic-should-not-copy");
}
+ @Test
+ void resendSelectedMessagesLooksUpIdsWithoutAStoreTimeWindow() throws
Exception {
+ String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+ MessageExt oldDeadLetter = new MessageExt();
+ oldDeadLetter.setMsgId("old-msg");
+ oldDeadLetter.setTopic(dlqTopic);
+ oldDeadLetter.setBody(new byte[] {1});
+ oldDeadLetter.setStoreTimestamp(1L);
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.SEND_OK);
+
+ when(adminExt.viewMessage(dlqTopic,
"old-msg")).thenReturn(oldDeadLetter);
+ when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+
+ assertThat(provider.resendMessages(
+ "instance-a", "group-a", List.of("old-msg"), "target-topic"))
+ .extracting("matched", "resent", "failed", "outcome",
"scanIncomplete")
+ .containsExactly(1, 1, 0, "SUCCESS", false);
+
+ verify(adminExt).viewMessage(dlqTopic, "old-msg");
+ verifyNoInteractions(pullConsumer);
+ }
+
+ @Test
+ void resendSelectedMessagesReportsMissingLookupsAsPartial() throws
Exception {
+ String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+ MessageExt found = new MessageExt();
+ found.setMsgId("found-msg");
+ found.setTopic(dlqTopic);
+ found.setBody(new byte[] {1});
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.SEND_OK);
+
+ when(adminExt.viewMessage(dlqTopic, "missing-msg"))
+ .thenThrow(new IllegalStateException("message not found"));
+ when(adminExt.viewMessage(dlqTopic, "found-msg")).thenReturn(found);
+ when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+
+ assertThat(provider.resendMessages(
+ "instance-a", "group-a", List.of("missing-msg", "found-msg"),
"target-topic"))
+ .extracting("matched", "resent", "failed", "outcome",
"scanIncomplete")
+ .containsExactly(1, 1, 0, "PARTIAL", true);
+
+ verify(adminExt).viewMessage(dlqTopic, "missing-msg");
+ verify(adminExt).viewMessage(dlqTopic, "found-msg");
+ }
+
@Test
void resendMessagesMarksAResultPartialWhenScanReachesHardCap() throws
Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";