From 908f977e7e0f0e442d2f9be21114b651b5463de5 Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Tue, 24 May 2022 16:23:48 +0800 Subject: [PATCH] [ISSUE #3949] Add unit test for ProducerProcessor --- .../proxy/processor/ProducerProcessor.java | 10 +++-- .../proxy/processor/BaseProcessorTest.java | 4 ++ .../processor/ProducerProcessorTest.java | 42 +++++++++++++++++++ 3 files changed, 53 insertions(+), 3 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java index 2007fa53cd..d9666559f2 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java @@ -29,6 +29,7 @@ import org.apache.rocketmq.common.message.MessageAccessor; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.message.MessageDecoder; import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.protocol.NamespaceUtil; import org.apache.rocketmq.common.protocol.ResponseCode; import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; @@ -65,9 +66,12 @@ public class ProducerProcessor extends AbstractProcessor { String topic = messageExt0.getTopic(); if (ConfigurationManager.getProxyConfig().isEnableTopicMessageTypeCheck()) { if (topicMessageTypeValidator != null) { - TopicMessageType topicMessageType = serviceManager.getMetadataService().getTopicMessageType(topic); - TopicMessageType messageType = parseFromMessageExt(messageExt0); - topicMessageTypeValidator.validate(topicMessageType, messageType); + // Do not check retry or dlq topic + if (!NamespaceUtil.isRetryTopic(topic) && !NamespaceUtil.isDLQTopic(topic)) { + TopicMessageType topicMessageType = serviceManager.getMetadataService().getTopicMessageType(topic); + TopicMessageType messageType = parseFromMessageExt(messageExt0); + topicMessageTypeValidator.validate(topicMessageType, messageType); + } } } SelectableMessageQueue messageQueue = queueSelector.select(ctx, diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/BaseProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/BaseProcessorTest.java index 31954705b7..f67d861158 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/BaseProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/BaseProcessorTest.java @@ -31,6 +31,7 @@ import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest; import org.apache.rocketmq.proxy.service.ServiceManager; import org.apache.rocketmq.proxy.service.message.MessageService; +import org.apache.rocketmq.proxy.service.metadata.MetadataService; import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.TransactionService; @@ -63,6 +64,8 @@ public class BaseProcessorTest extends InitConfigAndLoggerTest { @Mock protected ProxyRelayService proxyRelayService; @Mock + protected MetadataService metadataService; + @Mock protected ProducerProcessor producerProcessor; @Mock protected ConsumerProcessor consumerProcessor; @@ -79,6 +82,7 @@ public class BaseProcessorTest extends InitConfigAndLoggerTest { when(serviceManager.getConsumerManager()).thenReturn(consumerManager); when(serviceManager.getTransactionService()).thenReturn(transactionService); when(serviceManager.getProxyRelayService()).thenReturn(proxyRelayService); + when(serviceManager.getMetadataService()).thenReturn(metadataService); } protected static ProxyContext createContext() { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java index eac0abe553..7208c6da83 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java @@ -25,6 +25,7 @@ import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.common.KeyBuilder; import org.apache.rocketmq.common.MixAll; +import org.apache.rocketmq.common.attribute.TopicMessageType; import org.apache.rocketmq.common.consumer.ReceiptHandle; import org.apache.rocketmq.common.message.MessageAccessor; import org.apache.rocketmq.common.message.MessageClientIDSetter; @@ -46,6 +47,7 @@ import static org.junit.Assert.assertNotNull; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -65,6 +67,7 @@ public class ProducerProcessorTest extends BaseProcessorTest { @Test public void testSendMessage() throws Throwable { + when(metadataService.getTopicMessageType(eq(TOPIC))).thenReturn(TopicMessageType.NORMAL); String txId = MessageClientIDSetter.createUniqID(); String msgId = MessageClientIDSetter.createUniqID(); @@ -76,6 +79,45 @@ public class ProducerProcessorTest extends BaseProcessorTest { when(this.messageService.sendMessage(any(), any(), any(), requestHeaderArgumentCaptor.capture(), anyLong())) .thenReturn(CompletableFuture.completedFuture(Lists.newArrayList(sendResult))); + List messageExtList = new ArrayList<>(); + MessageExt messageExt = createMessageExt(TOPIC, "tag", 0, 0); + messageExt.setSysFlag(MessageSysFlag.TRANSACTION_PREPARED_TYPE); + messageExtList.add(messageExt); + SelectableMessageQueue messageQueue = mock(SelectableMessageQueue.class); + when(messageQueue.getBrokerName()).thenReturn("mockBroker"); + + List sendResultList = this.producerProcessor.sendMessage( + createContext(), + (ctx, messageQueueView) -> messageQueue, + PRODUCER_GROUP, + messageExtList, + 3000 + ).get(); + + assertNotNull(sendResultList); + TransactionId transactionId = TransactionId.decode(sendResultList.get(0).getTransactionId()); + assertNotNull(transactionId); + assertEquals(txId, transactionId.getBrokerTransactionId()); + assertEquals("mockBroker", transactionId.getBrokerName()); + + SendMessageRequestHeader requestHeader = requestHeaderArgumentCaptor.getValue(); + assertEquals(PRODUCER_GROUP, requestHeader.getProducerGroup()); + assertEquals(TOPIC, requestHeader.getTopic()); + } + + @Test + public void testSendRetryMessage() throws Throwable { + String txId = MessageClientIDSetter.createUniqID(); + String msgId = MessageClientIDSetter.createUniqID(); + + SendResult sendResult = new SendResult(); + sendResult.setSendStatus(SendStatus.SEND_OK); + sendResult.setTransactionId(txId); + sendResult.setMsgId(msgId); + ArgumentCaptor requestHeaderArgumentCaptor = ArgumentCaptor.forClass(SendMessageRequestHeader.class); + when(this.messageService.sendMessage(any(), any(), any(), requestHeaderArgumentCaptor.capture(), anyLong())) + .thenReturn(CompletableFuture.completedFuture(Lists.newArrayList(sendResult))); + List messageExtList = new ArrayList<>(); MessageExt messageExt = createMessageExt(MixAll.getRetryTopic(CONSUMER_GROUP), "tag", 0, 0); messageExt.setSysFlag(MessageSysFlag.TRANSACTION_PREPARED_TYPE);