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