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 41d579d2d fix(dlq): prevent cross-topic selected message resends 
(#5100)
41d579d2d is described below

commit 41d579d2d965905d3ddb3eb9158b367bf4c8c005
Author: coder999o <[email protected]>
AuthorDate: Thu Oct 1 18:15:01 2026 +0800

    fix(dlq): prevent cross-topic selected message resends (#5100)
---
 .../provider/apache/RocketMQDLQProvider.java       |  6 +++-
 .../provider/apache/RocketMQDLQProviderTest.java   | 36 ++++++++++++++++++++++
 2 files changed, 41 insertions(+), 1 deletion(-)

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 f5c3b1085..107ebf8e3 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
@@ -273,8 +273,12 @@ public class RocketMQDLQProvider implements DLQProvider {
                         continue;
                     }
                     MessageExt deadLetter = admin.viewMessage(dlqTopic, msgId);
-                    if (deadLetter != null) {
+                    // Offset-based broker lookup ignores the requested topic; 
check the returned message.
+                    if (deadLetter != null && 
dlqTopic.equals(deadLetter.getTopic())) {
                         resolved.add(deadLetter);
+                    } else if (deadLetter != null) {
+                        log.warn("Ignoring selected message {} from topic {} 
instead of {}",
+                                msgId, deadLetter.getTopic(), dlqTopic);
                     }
                 } catch (Exception e) {
                     log.warn("Failed to resolve selected dead letter {} from 
{}: {}",
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 bc1620922..5d85c9abf 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
@@ -434,6 +434,42 @@ class RocketMQDLQProviderTest {
         verify(adminExt).viewMessage(dlqTopic, msgId);
     }
 
+    @Test
+    void resendSelectedMessagesSkipsMessagesOutsideTheRequestedDlq() throws 
Exception {
+        String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+        String normalId = MessageDecoder.createMessageId(new 
InetSocketAddress("172.30.10.100", 10911), 12345L);
+        String otherGroupId = MessageDecoder.createMessageId(new 
InetSocketAddress("172.30.10.100", 10911), 12346L);
+        String validId = MessageDecoder.createMessageId(new 
InetSocketAddress("172.30.10.100", 10911), 12347L);
+        MessageExt normalMessage = new MessageExt();
+        normalMessage.setMsgId(normalId);
+        normalMessage.setTopic("orders");
+        normalMessage.setBody(new byte[] {1});
+        MessageExt otherGroupMessage = new MessageExt();
+        otherGroupMessage.setMsgId(otherGroupId);
+        otherGroupMessage.setTopic(MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-b");
+        otherGroupMessage.setBody(new byte[] {1});
+        MessageExt deadLetter = new MessageExt();
+        deadLetter.setTopic(dlqTopic);
+        deadLetter.setMsgId(validId);
+        deadLetter.setBody(new byte[] {1});
+        
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithBrokerAddresses("172.30.10.100:10911"));
+        when(adminExt.viewMessage(dlqTopic, 
normalId)).thenReturn(normalMessage);
+        when(adminExt.viewMessage(dlqTopic, 
otherGroupId)).thenReturn(otherGroupMessage);
+        when(adminExt.viewMessage(dlqTopic, validId)).thenReturn(deadLetter);
+        stubExistingTarget("target-topic");
+        SendResult sendResult = new SendResult();
+        sendResult.setSendStatus(SendStatus.SEND_OK);
+        when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+
+        DLQResendResultVO result = provider.resendMessages(
+                "instance-a", "group-a", List.of(normalId, otherGroupId, 
validId), "target-topic");
+
+        assertThat(result)
+                .extracting("matched", "resent", "failed", "outcome", 
"scanIncomplete")
+                .containsExactly(1, 1, 0, "PARTIAL", true);
+        verify(dlqProducer, times(1)).send(any(Message.class));
+    }
+
     @Test
     void resendSelectedMessagesPassesNonOffsetIdsThroughToViewMessage() throws 
Exception {
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";

Reply via email to