From e14fcdbf9367bd04f20f5e34e04622195ea9b3a8 Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 16:58:42 -0700 Subject: [PATCH 1/7] fix: summarize heartbeat sync logs --- .../service/sysmessage/HeartbeatSyncer.java | 72 ++++++++++++++++--- .../sysmessage/HeartbeatSyncerTest.java | 36 +++++++++- 2 files changed, 98 insertions(+), 10 deletions(-) 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..ed5e52eda78 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 @@ -18,6 +18,7 @@ package org.apache.rocketmq.proxy.service.sysmessage; import com.alibaba.fastjson2.JSON; +import com.google.common.base.MoreObjects; import io.netty.channel.Channel; import org.apache.rocketmq.broker.client.ClientChannelInfo; import org.apache.rocketmq.broker.client.ConsumerGroupEvent; @@ -47,6 +48,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; public class HeartbeatSyncer extends AbstractSystemMessageSyncer { @@ -131,16 +133,19 @@ public void onConsumerRegister(String consumerGroup, ClientChannelInfo clientCha ); data.setSubscriptionDataSet(subList); - log.debug("sync register heart beat. topic:{}, data:{}", this.getBroadcastTopicName(), data); + log.debug("sync register heart beat. topic:{}, dataSummary:{}", + this.getBroadcastTopicName(), summarizeHeartbeatData(data)); this.sendSystemMessage(data); } catch (Throwable t) { - log.error("heartbeat register broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}, messageModel:{}, consumeFromWhere:{}, subList:{}", - consumerGroup, clientChannelInfo, consumeType, messageModel, consumeFromWhere, subList, t); + log.error("heartbeat register broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}, messageModel:{}, consumeFromWhere:{}, subscriptionSummary:{}", + consumerGroup, clientChannelInfo, consumeType, messageModel, consumeFromWhere, + summarizeSubscriptionDataSet(subList), t); } }); } catch (Throwable t) { - log.error("heartbeat submit register broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}, messageModel:{}, consumeFromWhere:{}, subList:{}", - consumerGroup, clientChannelInfo, consumeType, messageModel, consumeFromWhere, subList, t); + log.error("heartbeat submit register broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}, messageModel:{}, consumeFromWhere:{}, subscriptionSummary:{}", + consumerGroup, clientChannelInfo, consumeType, messageModel, consumeFromWhere, + summarizeSubscriptionDataSet(subList), t); } } @@ -168,7 +173,8 @@ public void onConsumerUnRegister(String consumerGroup, ClientChannelInfo clientC remoteChannel.encode() ); - log.debug("sync unregister heart beat. topic:{}, data:{}", this.getBroadcastTopicName(), data); + log.debug("sync unregister heart beat. topic:{}, dataSummary:{}", + this.getBroadcastTopicName(), summarizeHeartbeatData(data)); this.sendSystemMessage(data); } catch (Throwable t) { log.error("heartbeat unregister broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}", @@ -188,8 +194,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; } @@ -203,7 +210,8 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo data.getLanguage(), data.getVersion() ); - log.debug("start process remote channel. data:{}, clientChannelInfo:{}", data, clientChannelInfo); + log.debug("start process remote channel. dataSummary:{}, clientChannelInfo:{}", + summarizeHeartbeatData(data), clientChannelInfo); if (data.getHeartbeatType().equals(HeartbeatType.REGISTER)) { this.consumerManager.registerConsumer( data.getGroup(), @@ -222,7 +230,8 @@ 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. msg:{}, dataSummary:{}", + summarizeSystemMessage(msg), summarizeHeartbeatData(data), t); } } @@ -238,4 +247,49 @@ private String buildLocalProxyId() { private static String buildKey(String group, Channel channel) { return group + "@" + channel.id().asLongText(); } + + static String summarizeHeartbeatData(HeartbeatSyncerData data) { + if (data == null) { + return "null"; + } + return MoreObjects.toStringHelper("HeartbeatSyncerData") + .add("heartbeatType", data.getHeartbeatType()) + .add("clientId", data.getClientId()) + .add("language", data.getLanguage()) + .add("version", data.getVersion()) + .add("group", data.getGroup()) + .add("consumeType", data.getConsumeType()) + .add("messageModel", data.getMessageModel()) + .add("consumeFromWhere", data.getConsumeFromWhere()) + .add("localProxyId", data.getLocalProxyId()) + .add("channelDataPresent", data.getChannelData() != null) + .add("subscriptionSummary", summarizeSubscriptionDataSet(data.getSubscriptionDataSet())) + .toString(); + } + + static String summarizeSubscriptionDataSet(Set subscriptions) { + if (subscriptions == null) { + return "null"; + } + List topics = subscriptions.stream() + .map(SubscriptionData::getTopic) + .sorted() + .collect(Collectors.toList()); + return MoreObjects.toStringHelper("SubscriptionDataSet") + .add("count", subscriptions.size()) + .add("topics", topics) + .toString(); + } + + static String summarizeSystemMessage(MessageExt msg) { + if (msg == null) { + return "null"; + } + byte[] body = msg.getBody(); + return MoreObjects.toStringHelper("MessageExt") + .add("topic", msg.getTopic()) + .add("msgId", msg.getMsgId()) + .add("bodyBytes", body == null ? 0 : body.length) + .toString(); + } } 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..6e71c5be452 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,39 @@ public void testSyncGrpcV2Channel() throws Exception { assertSame(channelInfoList.get(0).getChannel(), syncUnRegisterChannelInfoArgumentCaptor.getValue().getChannel()); } + @Test + public void testSummarizeHeartbeatDataDoesNotExposeSubscriptionExpressions() throws Exception { + String expression = "secretTagA || secretTagB"; + HeartbeatSyncerData data = new HeartbeatSyncerData( + HeartbeatType.REGISTER, + clientId, + LanguageCode.JAVA, + 5, + "consumerGroup", + ConsumeType.CONSUME_PASSIVELY, + MessageModel.CLUSTERING, + ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET, + "proxy-0", + "raw-channel-data" + ); + data.setSubscriptionDataSet(Sets.newHashSet( + FilterAPI.buildSubscriptionData("topic-a", expression), + FilterAPI.buildSubscriptionData("topic-b", "*") + )); + + String summary = HeartbeatSyncer.summarizeHeartbeatData(data); + + assertTrue(summary.contains("heartbeatType=REGISTER")); + assertTrue(summary.contains("clientId=" + clientId)); + assertTrue(summary.contains("group=consumerGroup")); + assertTrue(summary.contains("channelDataPresent=true")); + assertTrue(summary.contains("count=2")); + assertTrue(summary.contains("topic-a")); + assertTrue(summary.contains("topic-b")); + assertFalse(summary.contains(expression)); + assertFalse(summary.contains("raw-channel-data")); + } + @Test public void testSyncRemotingChannel() throws Exception { String consumerGroup = "consumerGroup"; @@ -433,4 +467,4 @@ public int compareTo(@NotNull ChannelId o) { return this.channelId.compareTo(o.asLongText()); } } -} \ No newline at end of file +} From 9ed7cdc4075834829bd546223f82932dc9bc98bd Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 17:05:46 -0700 Subject: [PATCH 2/7] fix: summarize heartbeat system message failures --- .../sysmessage/AbstractSystemMessageSyncer.java | 13 ++++++++++--- .../proxy/service/sysmessage/HeartbeatSyncer.java | 8 ++++++++ 2 files changed, 18 insertions(+), 3 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..7392b0c91d4 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 @@ -109,18 +109,25 @@ 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: {}", + summarizeSystemMessageData(data), 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:{}", + summarizeSystemMessageData(data), getBroadcastTopicName(), result); } }); } catch (Throwable t) { - log.error("send system message failed. data: {}, topic: {}", data, targetTopic, t); + log.error("send system message failed. dataSummary: {}, topic: {}", + summarizeSystemMessageData(data), targetTopic, t); } } + protected Object summarizeSystemMessageData(Object data) { + return data; + } + protected SendMessageRequestHeader buildSendMessageRequestHeader(Message message, String producerGroup, int queueId) { SendMessageRequestHeader requestHeader = new SendMessageRequestHeader(); 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 ed5e52eda78..4caec68ac15 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 @@ -292,4 +292,12 @@ static String summarizeSystemMessage(MessageExt msg) { .add("bodyBytes", body == null ? 0 : body.length) .toString(); } + + @Override + protected Object summarizeSystemMessageData(Object data) { + if (data instanceof HeartbeatSyncerData) { + return summarizeHeartbeatData((HeartbeatSyncerData) data); + } + return super.summarizeSystemMessageData(data); + } } From ee086d938bd563f19fc9359f7753a1201d4f9b3c Mon Sep 17 00:00:00 2001 From: liuhy Date: Fri, 31 Jul 2026 05:34:52 -0700 Subject: [PATCH 3/7] [ISSUE #10726] Avoid raw heartbeat sync body logs --- .../service/sysmessage/HeartbeatSyncer.java | 18 +++++++ .../sysmessage/HeartbeatSyncerTest.java | 49 +++++++++++++++++++ 2 files changed, 67 insertions(+) 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 4caec68ac15..186bcf5f0aa 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 @@ -238,6 +238,24 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo 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 6e71c5be452..c8f53148995 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; @@ -422,6 +424,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()); From a9e59fedf1ac20b757df0ec2e5f87c1c0ffb4954 Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 2 Aug 2026 22:14:32 -0700 Subject: [PATCH 4/7] [ISSUE #10760] Avoid logging full system message data --- .../AbstractSystemMessageSyncer.java | 27 +++++++++++++---- .../sysmessage/HeartbeatSyncerTest.java | 29 +++++++++++++++++++ 2 files changed, 50 insertions(+), 6 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 7392b0c91d4..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, @@ -110,22 +111,36 @@ protected void sendSystemMessage(Object data) { ).whenCompleteAsync((result, throwable) -> { if (throwable != null) { log.error("send system message failed. dataSummary: {}, topic: {}", - summarizeSystemMessageData(data), getBroadcastTopicName(), throwable); + dataSummary, getBroadcastTopicName(), throwable); return; } if (SendStatus.SEND_OK != result.getSendStatus()) { log.error("send system message failed. dataSummary: {}, topic: {}, sendResult:{}", - summarizeSystemMessageData(data), getBroadcastTopicName(), result); + dataSummary, getBroadcastTopicName(), result); } }); } catch (Throwable t) { - log.error("send system message failed. dataSummary: {}, topic: {}", - summarizeSystemMessageData(data), targetTopic, t); + log.error("send system message failed. dataSummary: {}, topic: {}", dataSummary, targetTopic, t); } } - protected Object summarizeSystemMessageData(Object data) { - return data; + 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, 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 c8f53148995..42b9d622c98 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 @@ -253,6 +253,35 @@ public void testSummarizeHeartbeatDataDoesNotExposeSubscriptionExpressions() thr assertFalse(summary.contains("raw-channel-data")); } + @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"; From 5b0181f94b9cec67a3be231ccccef8b1a156f9b6 Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 03:46:49 -0700 Subject: [PATCH 5/7] 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 42b9d622c98..464e28b710a 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 @@ -273,15 +273,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"; From 04c608b8edee1f6c836d81a927f1ab436d62d603 Mon Sep 17 00:00:00 2001 From: liuhy Date: Sat, 15 Aug 2026 00:43:27 -0700 Subject: [PATCH 6/7] fix(proxy): use heartbeat diagnostic summary --- .../rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) 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 186bcf5f0aa..87b79904b80 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 @@ -230,8 +230,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo ); } } catch (Throwable t) { - log.error("heartbeat consume message failed. msg:{}, dataSummary:{}", - summarizeSystemMessage(msg), summarizeHeartbeatData(data), t); + log.error("heartbeat consume message failed. summary:{}", summarizeHeartbeatMessage(msg, data), t); } } From a47740ef76dbee543ab0d67442a6561600817e5d Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 16 Aug 2026 00:41:47 -0700 Subject: [PATCH 7/7] fix: remove invalid heartbeat summary override --- .../rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java | 7 ------- 1 file changed, 7 deletions(-) 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 87b79904b80..bd9ad422aaf 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 @@ -310,11 +310,4 @@ static String summarizeSystemMessage(MessageExt msg) { .toString(); } - @Override - protected Object summarizeSystemMessageData(Object data) { - if (data instanceof HeartbeatSyncerData) { - return summarizeHeartbeatData((HeartbeatSyncerData) data); - } - return super.summarizeSystemMessageData(data); - } }