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..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,8 @@ 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; import org.apache.rocketmq.client.consumer.PopResult; @@ -39,6 +41,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 +161,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(5, TimeUnit.SECONDS); + + 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));