From 5588873cb22f9e7f118619b76b20ac400df309c8 Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 00:42:10 -0700 Subject: [PATCH 1/2] [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 f77f269274b..9c12e787f8d 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 @@ -188,7 +188,6 @@ private PopResult filterPopResult(ProxyContext ctx, PopResult popResult, Command 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); - 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 b61c22b441e..27e099aed43 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; @@ -158,6 +160,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 de992d861ded345fafcaab0a6b30fbd4957b0ac1 Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 03:57:41 -0700 Subject: [PATCH 2/2] 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 27e099aed43..9829ac22f69 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; @@ -191,7 +192,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();