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";

Reply via email to