From df161afd2ab096b638b449b70f1d3a506b91ce3b Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 2 Aug 2026 22:14:32 -0700 Subject: [PATCH 1/2] [ISSUE #10760] Avoid logging full system message data --- .../AbstractSystemMessageSyncer.java | 28 ++++++++++++++-- .../sysmessage/HeartbeatSyncerTest.java | 32 ++++++++++++++++++- 2 files changed, 56 insertions(+), 4 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java index 05eb6726188..9520f65c35f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java @@ -93,6 +93,7 @@ public RPCHook getRpcHook() { protected void sendSystemMessage(Object data) { String targetTopic = this.getBroadcastTopicName(); + String dataSummary = summarizeSystemMessageData(data); try { Message message = new Message( targetTopic, @@ -109,18 +110,39 @@ protected void sendSystemMessage(Object data) { Duration.ofSeconds(3).toMillis() ).whenCompleteAsync((result, throwable) -> { if (throwable != null) { - log.error("send system message failed. data: {}, topic: {}", data, getBroadcastTopicName(), throwable); + log.error("send system message failed. dataSummary: {}, topic: {}", + dataSummary, getBroadcastTopicName(), throwable); return; } if (SendStatus.SEND_OK != result.getSendStatus()) { - log.error("send system message failed. data: {}, topic: {}, sendResult:{}", data, getBroadcastTopicName(), result); + log.error("send system message failed. dataSummary: {}, topic: {}, sendResult:{}", + dataSummary, getBroadcastTopicName(), result); } }); } catch (Throwable t) { - log.error("send system message failed. data: {}, topic: {}", data, targetTopic, t); + log.error("send system message failed. dataSummary: {}, topic: {}", dataSummary, targetTopic, t); } } + static String summarizeSystemMessageData(Object data) { + if (data == null) { + return "null"; + } + if (data instanceof HeartbeatSyncerData) { + HeartbeatSyncerData heartbeatData = (HeartbeatSyncerData) data; + int subscriptionCount = heartbeatData.getSubscriptionDataSet() == null + ? 0 : heartbeatData.getSubscriptionDataSet().size(); + return "HeartbeatSyncerData{" + + "heartbeatType=" + heartbeatData.getHeartbeatType() + + ", clientId=" + heartbeatData.getClientId() + + ", group=" + heartbeatData.getGroup() + + ", subscriptionCount=" + subscriptionCount + + ", channelDataPresent=" + (heartbeatData.getChannelData() != null) + + '}'; + } + return data.getClass().getSimpleName(); + } + protected SendMessageRequestHeader buildSendMessageRequestHeader(Message message, String producerGroup, int queueId) { SendMessageRequestHeader requestHeader = new SendMessageRequestHeader(); 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..85f4ba835c9 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 @@ -77,6 +77,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; @@ -217,6 +218,35 @@ public void testSyncGrpcV2Channel() throws Exception { assertSame(channelInfoList.get(0).getChannel(), syncUnRegisterChannelInfoArgumentCaptor.getValue().getChannel()); } + @Test + public void testSummarizeSystemMessageDataAvoidsDetailedHeartbeatPayload() throws Exception { + HeartbeatSyncerData data = new HeartbeatSyncerData( + HeartbeatType.REGISTER, + clientId, + LanguageCode.JAVA, + 5, + "sensitiveGroup", + ConsumeType.CONSUME_PASSIVELY, + MessageModel.CLUSTERING, + ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET, + "localProxyId", + "sensitiveChannelData" + ); + data.setSubscriptionDataSet(Sets.newHashSet(FilterAPI.buildSubscriptionData("sensitiveTopic", "sensitiveTag"))); + + String summary = AbstractSystemMessageSyncer.summarizeSystemMessageData(data); + + assertTrue(summary.contains("HeartbeatSyncerData")); + assertTrue(summary.contains("heartbeatType=REGISTER")); + assertTrue(summary.contains("clientId=" + clientId)); + assertTrue(summary.contains("group=sensitiveGroup")); + assertTrue(summary.contains("subscriptionCount=1")); + assertTrue(summary.contains("channelDataPresent=true")); + assertFalse(summary.contains("sensitiveTopic")); + assertFalse(summary.contains("sensitiveTag")); + assertFalse(summary.contains("sensitiveChannelData")); + } + @Test public void testSyncRemotingChannel() throws Exception { String consumerGroup = "consumerGroup"; @@ -433,4 +463,4 @@ public int compareTo(@NotNull ChannelId o) { return this.channelId.compareTo(o.asLongText()); } } -} \ No newline at end of file +} From bb03dc2cce953f0580db8157b8dc762177bb43a8 Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 03:46:49 -0700 Subject: [PATCH 2/2] fix(proxy): harden system message log summaries --- .../AbstractSystemMessageSyncer.java | 19 ++++++++++----- .../sysmessage/HeartbeatSyncerTest.java | 24 +++++++++++++++++-- 2 files changed, 35 insertions(+), 8 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java index 9520f65c35f..5e0d82f9822 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java @@ -93,7 +93,7 @@ public RPCHook getRpcHook() { protected void sendSystemMessage(Object data) { String targetTopic = this.getBroadcastTopicName(); - String dataSummary = summarizeSystemMessageData(data); + String dataSummary = safelySummarizeSystemMessageData(data); try { Message message = new Message( targetTopic, @@ -111,12 +111,12 @@ protected void sendSystemMessage(Object data) { ).whenCompleteAsync((result, throwable) -> { if (throwable != null) { log.error("send system message failed. dataSummary: {}, topic: {}", - dataSummary, getBroadcastTopicName(), throwable); + dataSummary, targetTopic, throwable); return; } if (SendStatus.SEND_OK != result.getSendStatus()) { log.error("send system message failed. dataSummary: {}, topic: {}, sendResult:{}", - dataSummary, getBroadcastTopicName(), result); + dataSummary, targetTopic, result); } }); } catch (Throwable t) { @@ -134,13 +134,20 @@ static String summarizeSystemMessageData(Object data) { ? 0 : heartbeatData.getSubscriptionDataSet().size(); return "HeartbeatSyncerData{" + "heartbeatType=" + heartbeatData.getHeartbeatType() - + ", clientId=" + heartbeatData.getClientId() - + ", group=" + heartbeatData.getGroup() + ", subscriptionCount=" + subscriptionCount + ", channelDataPresent=" + (heartbeatData.getChannelData() != null) + '}'; } - return data.getClass().getSimpleName(); + String simpleName = data.getClass().getSimpleName(); + return StringUtils.isEmpty(simpleName) ? data.getClass().getName() : simpleName; + } + + static String safelySummarizeSystemMessageData(Object data) { + try { + return summarizeSystemMessageData(data); + } catch (Throwable ignored) { + return "unavailable"; + } } protected SendMessageRequestHeader buildSendMessageRequestHeader(Message message, 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 85f4ba835c9..cff69bdf583 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 @@ -238,15 +238,35 @@ public void testSummarizeSystemMessageDataAvoidsDetailedHeartbeatPayload() throw assertTrue(summary.contains("HeartbeatSyncerData")); assertTrue(summary.contains("heartbeatType=REGISTER")); - assertTrue(summary.contains("clientId=" + clientId)); - assertTrue(summary.contains("group=sensitiveGroup")); assertTrue(summary.contains("subscriptionCount=1")); assertTrue(summary.contains("channelDataPresent=true")); + assertFalse(summary.contains(clientId)); + assertFalse(summary.contains("sensitiveGroup")); assertFalse(summary.contains("sensitiveTopic")); assertFalse(summary.contains("sensitiveTag")); assertFalse(summary.contains("sensitiveChannelData")); } + @Test + public void testSummarizeSystemMessageDataUsesClassNameForAnonymousPayload() { + Object data = new Object() { + }; + + assertEquals(data.getClass().getName(), AbstractSystemMessageSyncer.summarizeSystemMessageData(data)); + } + + @Test + public void testSafelySummarizeSystemMessageDataDoesNotPropagateFailures() { + HeartbeatSyncerData data = new HeartbeatSyncerData() { + @Override + public Set getSubscriptionDataSet() { + throw new IllegalStateException("summary failure"); + } + }; + + assertEquals("unavailable", AbstractSystemMessageSyncer.safelySummarizeSystemMessageData(data)); + } + @Test public void testSyncRemotingChannel() throws Exception { String consumerGroup = "consumerGroup";