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 9cf836cc1 fix(message): enforce broker topology guard on the primary
msgId lookup path (#2833)
9cf836cc1 is described below
commit 9cf836cc1c89331d239ae6090c85915cc8bb2763
Author: 烤化の初雪 <[email protected]>
AuthorDate: Fri Sep 4 12:01:26 2026 +0800
fix(message): enforce broker topology guard on the primary msgId lookup
path (#2833)
Signed-off-by: unbridled-41
<[email protected]>
Co-authored-by: unbridled-41
<[email protected]>
---
.../provider/apache/RocketMQMessageProvider.java | 32 +++++++++++-
.../apache/RocketMQMessageProviderTest.java | 61 ++++++++++++++++++++--
2 files changed, 87 insertions(+), 6 deletions(-)
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 883ad052d..789aa2bc0 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
@@ -142,7 +142,9 @@ public class RocketMQMessageProvider implements
MessageProvider {
MessageExt messageExt = null;
if (StringUtils.hasText(topic)) {
try {
- messageExt = adminExt.viewMessage(topic, msgId);
+ if (isWithinKnownBrokerTopology(adminExt, msgId)) {
+ messageExt = adminExt.viewMessage(topic, msgId);
+ }
} catch (Exception e) {
log.warn("viewMessage(topic={}, msgId={}) failed: {}", topic,
msgId, e.getMessage());
}
@@ -190,6 +192,29 @@ public class RocketMQMessageProvider implements
MessageProvider {
return null;
}
+ /**
+ * Offset-style message ids embed a broker address that {@code
MQAdminImpl#viewMessage} connects
+ * to directly, so ids whose embedded address is outside the selected
instance topology must be
+ * rejected before that call. Ids that do not decode as offset ids take
MQAdminImpl's unique-key
+ * lookup, which resolves brokers from the topic route and needs no guard.
When the topology
+ * itself cannot be verified, reject too — mirroring the fallback path's
behavior — instead of
+ * handing an unverified address to remoting.
+ */
+ private boolean isWithinKnownBrokerTopology(DefaultMQAdminExt adminExt,
String msgId) {
+ MessageId messageId;
+ try {
+ messageId = MessageDecoder.decodeMessageId(msgId);
+ } catch (Exception e) {
+ return true;
+ }
+ try {
+ return validatedBrokerAddr(adminExt, msgId, messageId) != null;
+ } catch (Exception e) {
+ log.warn("Could not verify broker topology for msgId={}: {}",
msgId, e.getMessage());
+ return false;
+ }
+ }
+
private String decodedBrokerAddr(MessageId messageId) {
SocketAddress address = messageId.getAddress();
if (!(address instanceof InetSocketAddress)) {
@@ -566,7 +591,10 @@ public class RocketMQMessageProvider implements
MessageProvider {
private long resolveMessageStoreTimestamp(DefaultMQAdminExt adminExt,
String msgId, String topic) {
if (StringUtils.hasText(topic)) {
try {
- MessageExt messageExt = adminExt.viewMessage(topic, msgId);
+ MessageExt messageExt = null;
+ if (isWithinKnownBrokerTopology(adminExt, msgId)) {
+ messageExt = adminExt.viewMessage(topic, msgId);
+ }
if (messageExt != null) {
return messageExt.getStoreTimestamp();
}
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 91af89385..4a4949d8e 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
@@ -250,16 +250,70 @@ class RocketMQMessageProviderTest {
void queryByMsgIdRejectsDecodedBrokerOutsideKnownTopology() throws
Exception {
String msgId = MessageDecoder.createMessageId(new
InetSocketAddress("10.2.3.4", 10911), 12345L);
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithBrokerAddresses("172.30.10.100:10911"));
- when(adminExt.viewMessage("TopicA", msgId))
- .thenThrow(new IllegalStateException("primary lookup failed"));
List<MessageRecordVO> result = provider.queryMessages(
"instance-a", "TopicA", msgId, null, null, 100L, 200L);
assertThat(result).isEmpty();
+ verify(adminExt, never()).viewMessage(anyString(), anyString());
+ verify(adminExt, never()).getDefaultMQAdminExtImpl();
+ }
+
+ @Test
+ void queryByMsgIdRejectsOffsetIdWhenTopologyCannotBeVerified() throws
Exception {
+ String msgId = MessageDecoder.createMessageId(new
InetSocketAddress("172.30.10.100", 10911), 12345L);
+ when(adminExt.examineBrokerClusterInfo()).thenThrow(new
IllegalStateException("nameserver unreachable"));
+
+ List<MessageRecordVO> result = provider.queryMessages(
+ "instance-a", "TopicA", msgId, null, null, 100L, 200L);
+
+ assertThat(result).isEmpty();
+ verify(adminExt, never()).viewMessage(anyString(), anyString());
verify(adminExt, never()).getDefaultMQAdminExtImpl();
}
+ @Test
+ void queryByMsgIdRejectsOffsetIdWhenTopologyIsEmpty() throws Exception {
+ String msgId = MessageDecoder.createMessageId(new
InetSocketAddress("172.30.10.100", 10911), 12345L);
+
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithBrokerAddresses());
+
+ List<MessageRecordVO> result = provider.queryMessages(
+ "instance-a", "TopicA", msgId, null, null, 100L, 200L);
+
+ assertThat(result).isEmpty();
+ verify(adminExt, never()).viewMessage(anyString(), anyString());
+ }
+
+ @Test
+ void queryByMsgIdPassesNonOffsetIdsThroughToViewMessage() throws Exception
{
+ when(adminExt.viewMessage("TopicA", "uniq-key-1"))
+ .thenThrow(new IllegalStateException("unique key lookup
handled by MQAdminImpl"));
+
+ List<MessageRecordVO> result = provider.queryMessages(
+ "instance-a", "TopicA", "uniq-key-1", null, null, 100L, 200L);
+
+ assertThat(result).isEmpty();
+ verify(adminExt).viewMessage("TopicA", "uniq-key-1");
+ verify(adminExt, never()).examineBrokerClusterInfo();
+ }
+
+ @Test
+ void queryByMsgIdStillViewsMessagesInsideKnownTopology() throws Exception {
+ String msgId = MessageDecoder.createMessageId(
+ new InetSocketAddress("172.30.10.100", 10911), 12345L);
+ MessageExt message = new MessageExt();
+ message.setMsgId(msgId);
+ message.setTopic("TopicA");
+
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithBrokerAddresses("172.30.10.100:10911"));
+ when(adminExt.viewMessage("TopicA", msgId)).thenReturn(message);
+
+ List<MessageRecordVO> result = provider.queryMessages(
+ "instance-a", "TopicA", msgId, null, null, 100L, 200L);
+
+
assertThat(result).singleElement().extracting(MessageRecordVO::getMsgId).isEqualTo(msgId);
+ verify(adminExt).viewMessage("TopicA", msgId);
+ }
+
@Test
void queryByTopicSurfacesPullConsumerFailure() throws Exception {
when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
@@ -661,13 +715,12 @@ class RocketMQMessageProviderTest {
void getMessageTraceDoesNotUseDecodedBrokerOutsideKnownTopology() throws
Exception {
String msgId = MessageDecoder.createMessageId(new
InetSocketAddress("10.2.3.4", 10911), 12345L);
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithBrokerAddresses("172.30.10.100:10911"));
- when(adminExt.viewMessage("TopicA", msgId))
- .thenThrow(new IllegalStateException("topic lookup failed"));
when(adminExt.queryMessage(anyString(), anyString(), anyInt(),
anyLong(), anyLong()))
.thenReturn(new QueryResult(0L, List.of()));
provider.getMessageTrace("instance-a", msgId, "TopicA");
+ verify(adminExt, never()).viewMessage(anyString(), anyString());
verify(adminExt, never()).getDefaultMQAdminExtImpl();
ArgumentCaptor<Long> beginCaptor = ArgumentCaptor.forClass(Long.class);
ArgumentCaptor<Long> endCaptor = ArgumentCaptor.forClass(Long.class);