From f83bfea2e8bc5bff8022ef3db0689d3bbf7f6048 Mon Sep 17 00:00:00 2001 From: liuhy Date: Fri, 31 Jul 2026 05:39:25 -0700 Subject: [PATCH 1/3] [ISSUE #10728] Avoid raw POP message logs --- .../proxy/processor/ConsumerProcessor.java | 20 ++++++++++++++-- .../processor/ConsumerProcessorTest.java | 24 +++++++++++++++++++ 2 files changed, 42 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java index f77f269274b..47daeaf0578 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java @@ -23,6 +23,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.TreeSet; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArraySet; import java.util.concurrent.ExecutorService; @@ -187,7 +188,8 @@ private PopResult filterPopResult(ProxyContext ctx, PopResult popResult, Command fillUniqIDIfNeed(messageExt); String handleString = createHandle(messageExt.getProperty(MessageConst.PROPERTY_POP_CK), messageExt.getCommitLogOffset()); if (handleString == null) { - log.error("[BUG] pop message from broker but handle is empty. requestHeader:{}, msg:{}", requestHeader, messageExt); + log.error("[BUG] pop message from broker but handle is empty. requestHeader:{}, msgSummary:{}", + requestHeader, summarizeMessageExt(messageExt)); messageExtList.add(messageExt); continue; } @@ -233,7 +235,8 @@ private PopResult filterPopResult(ProxyContext ctx, PopResult popResult, Command break; } } catch (Throwable t) { - log.error("process filterMessage failed. requestHeader:{}, msg:{}", requestHeader, messageExt, t); + log.error("process filterMessage failed. requestHeader:{}, msgSummary:{}", + requestHeader, summarizeMessageExt(messageExt), t); messageExtList.add(messageExt); } } @@ -242,6 +245,19 @@ private PopResult filterPopResult(ProxyContext ctx, PopResult popResult, Command return popResult; } + static String summarizeMessageExt(MessageExt messageExt) { + if (messageExt == null) { + return "msg=null"; + } + return "topic=" + messageExt.getTopic() + + ", msgId=" + messageExt.getMsgId() + + ", queueId=" + messageExt.getQueueId() + + ", queueOffset=" + messageExt.getQueueOffset() + + ", commitLogOffset=" + messageExt.getCommitLogOffset() + + ", bodySize=" + (messageExt.getBody() == null ? 0 : messageExt.getBody().length) + + ", propertyKeys=" + (messageExt.getProperties() == null + ? "[]" : new TreeSet<>(messageExt.getProperties().keySet())); + } private void fillUniqIDIfNeed(MessageExt messageExt) { if (StringUtils.isBlank(MessageClientIDSetter.getUniqID(messageExt))) { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java index b61c22b441e..ea75cebfc4a 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java @@ -92,6 +92,30 @@ public void before() throws Throwable { this.consumerProcessor = new ConsumerProcessor(messagingProcessor, serviceManager, Executors.newCachedThreadPool()); } + @Test + public void testSummarizeMessageExtDoesNotExposeBodyOrPropertyValues() { + MessageExt messageExt = new MessageExt(); + messageExt.setTopic(TOPIC); + messageExt.setMsgId("msgId"); + messageExt.setQueueId(1); + messageExt.setQueueOffset(2L); + messageExt.setCommitLogOffset(3L); + messageExt.setBody(new byte[] {115, 101, 99, 114, 101, 116}); + messageExt.putUserProperty("secretKey", "secretValue"); + + String summary = ConsumerProcessor.summarizeMessageExt(messageExt); + + assertThat(summary).contains("topic=" + TOPIC); + assertThat(summary).contains("msgId=msgId"); + assertThat(summary).contains("queueId=1"); + assertThat(summary).contains("queueOffset=2"); + assertThat(summary).contains("commitLogOffset=3"); + assertThat(summary).contains("bodySize=6"); + assertThat(summary).contains("secretKey"); + assertThat(summary).doesNotContain("secretValue"); + assertThat(summary).doesNotContain("115, 101, 99, 114, 101, 116"); + } + @Test public void testPopMessage() throws Throwable { final String tag = "tag"; From a54066d98bf6ffeaa71f7a0df56e155435f02d87 Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 00:42:10 -0700 Subject: [PATCH 2/3] [ISSUE #10784] Drop POP messages without receipt handles --- .../proxy/processor/ConsumerProcessor.java | 1 - .../processor/ConsumerProcessorTest.java | 40 +++++++++++++++++++ 2 files changed, 40 insertions(+), 1 deletion(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java index 47daeaf0578..bdb8a395420 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ConsumerProcessor.java @@ -190,7 +190,6 @@ private PopResult filterPopResult(ProxyContext ctx, PopResult popResult, Command if (handleString == null) { log.error("[BUG] pop message from broker but handle is empty. requestHeader:{}, msgSummary:{}", requestHeader, summarizeMessageExt(messageExt)); - messageExtList.add(messageExt); continue; } MessageAccessor.putProperty(messageExt, MessageConst.PROPERTY_POP_CK, handleString); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java index ea75cebfc4a..550b9de9e6b 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java @@ -28,6 +28,7 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicBoolean; import org.apache.rocketmq.client.consumer.AckResult; import org.apache.rocketmq.client.consumer.AckStatus; import org.apache.rocketmq.client.consumer.PopResult; @@ -39,6 +40,7 @@ import org.apache.rocketmq.common.constant.ConsumeInitMode; import org.apache.rocketmq.common.consumer.ReceiptHandle; import org.apache.rocketmq.common.filter.ExpressionType; +import org.apache.rocketmq.common.message.MessageAccessor; import org.apache.rocketmq.common.message.MessageClientIDSetter; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.message.MessageExt; @@ -182,6 +184,44 @@ public void testPopMessage() throws Throwable { assertEquals(messageExtList.get(2).getMsgId(), toDLQMessageIdArgumentCaptor.getValue()); } + @Test + public void testPopMessageShouldDropMessageWithoutReceiptHandle() throws Throwable { + final long invisibleTime = Duration.ofSeconds(15).toMillis(); + MessageExt messageExt = createMessageExt(TOPIC, "tag", 0, invisibleTime); + MessageAccessor.clearProperty(messageExt, MessageConst.PROPERTY_POP_CK); + PopResult innerPopResult = new PopResult(PopStatus.FOUND, Collections.singletonList(messageExt)); + when(this.messageService.popMessage(any(), any(), any(), anyLong())) + .thenReturn(CompletableFuture.completedFuture(innerPopResult)); + when(this.topicRouteService.getCurrentMessageQueueView(any(), anyString())) + .thenReturn(mock(MessageQueueView.class)); + + AtomicBoolean filterInvoked = new AtomicBoolean(false); + PopMessageResultFilter popMessageResultFilter = (ctx, consumerGroup, subscriptionData, message) -> { + filterInvoked.set(true); + return PopMessageResultFilter.FilterResult.MATCH; + }; + + PopResult popResult = this.consumerProcessor.popMessage( + createContext(), + (ctx, messageQueueView) -> mock(AddressableMessageQueue.class), + CONSUMER_GROUP, + TOPIC, + 60, + invisibleTime, + Duration.ofSeconds(3).toMillis(), + ConsumeInitMode.MAX, + FilterAPI.build(TOPIC, "*", ExpressionType.TAG), + false, + popMessageResultFilter, + null, + Duration.ofSeconds(3).toMillis() + ).get(); + + assertEquals(PopStatus.FOUND, popResult.getPopStatus()); + assertThat(popResult.getMsgFoundList()).isEmpty(); + assertFalse(filterInvoked.get()); + } + @Test public void testAckMessage() throws Throwable { ReceiptHandle handle = create(createMessageExt(MixAll.RETRY_GROUP_TOPIC_PREFIX + TOPIC, "", 0, 3000)); From f175cea6ea4c9fee86704dd9cb4b81590797fb9b Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 03:57:41 -0700 Subject: [PATCH 3/3] test(proxy): bound pop message regression wait --- .../apache/rocketmq/proxy/processor/ConsumerProcessorTest.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java index 550b9de9e6b..f94e6d23699 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java @@ -28,6 +28,7 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import org.apache.rocketmq.client.consumer.AckResult; import org.apache.rocketmq.client.consumer.AckStatus; @@ -215,7 +216,7 @@ public void testPopMessageShouldDropMessageWithoutReceiptHandle() throws Throwab popMessageResultFilter, null, Duration.ofSeconds(3).toMillis() - ).get(); + ).get(5, TimeUnit.SECONDS); assertEquals(PopStatus.FOUND, popResult.getPopStatus()); assertThat(popResult.getMsgFoundList()).isEmpty();