diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java index e063d79707b..fb0dd0bc098 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java @@ -188,8 +188,9 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo } for (MessageExt msg : msgs) { + HeartbeatSyncerData data = null; try { - HeartbeatSyncerData data = JSON.parseObject(new String(msg.getBody(), StandardCharsets.UTF_8), HeartbeatSyncerData.class); + data = JSON.parseObject(new String(msg.getBody(), StandardCharsets.UTF_8), HeartbeatSyncerData.class); if (data.getLocalProxyId().equals(localProxyId)) { continue; } @@ -222,13 +223,31 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo ); } } catch (Throwable t) { - log.error("heartbeat consume message failed. msg:{}, data:{}", msg, new String(msg.getBody(), StandardCharsets.UTF_8), t); + log.error("heartbeat consume message failed. summary:{}", summarizeHeartbeatMessage(msg, data), t); } } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } + static String summarizeHeartbeatMessage(MessageExt msg, HeartbeatSyncerData data) { + if (msg == null) { + return "msg=null"; + } + StringBuilder summary = new StringBuilder() + .append("topic=").append(msg.getTopic()) + .append(", msgId=").append(msg.getMsgId()) + .append(", bodySize=").append(msg.getBody() == null ? 0 : msg.getBody().length); + if (data != null) { + summary.append(", heartbeatType=").append(data.getHeartbeatType()) + .append(", group=").append(data.getGroup()) + .append(", clientId=").append(data.getClientId()) + .append(", subscriptionCount=") + .append(data.getSubscriptionDataSet() == null ? 0 : data.getSubscriptionDataSet().size()); + } + return summary.toString(); + } + private String buildLocalProxyId() { ProxyConfig proxyConfig = ConfigurationManager.getProxyConfig(); // use local address, remoting port and grpc port to build unique local proxy Id diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncerTest.java index 9a2c5e3437d..1a25b8f972e 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncerTest.java @@ -23,9 +23,11 @@ import apache.rocketmq.v2.Settings; import apache.rocketmq.v2.Subscription; import apache.rocketmq.v2.SubscriptionEntry; +import com.alibaba.fastjson2.JSON; import com.google.common.collect.Sets; import io.netty.channel.Channel; import io.netty.channel.ChannelId; +import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.Collections; import java.util.HashMap; @@ -77,6 +79,7 @@ import static org.awaitility.Awaitility.await; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotSame; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; @@ -388,6 +391,53 @@ private void testProcessConsumerGroupEvent(String consumerGroup, ClientChannelIn assertTrue(heartbeatSyncer.remoteChannelMap.isEmpty()); } + @Test + public void testSummarizeHeartbeatMessageDoesNotExposeSubscriptionOrChannelData() throws Exception { + SubscriptionData subscriptionData = FilterAPI.buildSubscriptionData("topic", "secret-tag"); + HeartbeatSyncerData data = new HeartbeatSyncerData( + HeartbeatType.REGISTER, + "client-secret", + LanguageCode.JAVA, + 5, + "consumerGroup", + ConsumeType.CONSUME_PASSIVELY, + MessageModel.CLUSTERING, + ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET, + "proxyId", + "secret-channel-data" + ); + data.setSubscriptionDataSet(Sets.newHashSet(subscriptionData)); + MessageExt msg = new MessageExt(); + msg.setTopic("heartbeatTopic"); + msg.setMsgId("msgId"); + msg.setBody(JSON.toJSONString(data).getBytes(StandardCharsets.UTF_8)); + + String summary = HeartbeatSyncer.summarizeHeartbeatMessage(msg, data); + + assertTrue(summary.contains("topic=heartbeatTopic")); + assertTrue(summary.contains("msgId=msgId")); + assertTrue(summary.contains("heartbeatType=REGISTER")); + assertTrue(summary.contains("group=consumerGroup")); + assertTrue(summary.contains("subscriptionCount=1")); + assertFalse(summary.contains("secret-tag")); + assertFalse(summary.contains("secret-channel-data")); + } + + @Test + public void testSummarizeHeartbeatMessageDoesNotExposeUnparsedBody() { + MessageExt msg = new MessageExt(); + msg.setTopic("heartbeatTopic"); + msg.setMsgId("msgId"); + msg.setBody("raw-secret-body".getBytes(StandardCharsets.UTF_8)); + + String summary = HeartbeatSyncer.summarizeHeartbeatMessage(msg, null); + + assertTrue(summary.contains("topic=heartbeatTopic")); + assertTrue(summary.contains("msgId=msgId")); + assertTrue(summary.contains("bodySize=15")); + assertFalse(summary.contains("raw-secret-body")); + } + private MessageExt convertFromMessage(Message message) { MessageExt messageExt = new MessageExt(); messageExt.setTopic(message.getTopic()); @@ -433,4 +483,4 @@ public int compareTo(@NotNull ChannelId o) { return this.channelId.compareTo(o.asLongText()); } } -} \ No newline at end of file +}