From ad5dc3fd46cdca7cb83f479a8926635183ca8a18 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Wed, 25 May 2022 17:41:52 +0800 Subject: [PATCH] [ISSUE #3949] add QueueSelector test cases --- .../consumer/ReceiveMessageActivityTest.java | 57 +++++++++++ .../v2/producer/SendMessageActivityTest.java | 94 +++++++++++++++++++ 2 files changed, 151 insertions(+) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java index ae5bef6eec..e9a5457cdd 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java @@ -28,15 +28,27 @@ import apache.rocketmq.v2.Settings; import io.grpc.stub.ServerCallStreamObserver; import io.grpc.stub.StreamObserver; import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.PopResult; import org.apache.rocketmq.client.consumer.PopStatus; +import org.apache.rocketmq.common.MixAll; +import org.apache.rocketmq.common.constant.PermName; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.common.protocol.route.QueueData; +import org.apache.rocketmq.common.protocol.route.TopicRouteData; +import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.grpc.v2.BaseActivityTest; +import org.apache.rocketmq.proxy.service.route.MessageQueueView; +import org.apache.rocketmq.proxy.service.route.SelectableMessageQueue; +import org.assertj.core.util.Lists; import org.junit.Before; import org.junit.Test; import org.mockito.ArgumentCaptor; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotEquals; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyInt; @@ -48,6 +60,9 @@ import static org.mockito.Mockito.when; public class ReceiveMessageActivityTest extends BaseActivityTest { + protected static final String BROKER_NAME = "broker"; + protected static final String CLUSTER_NAME = "cluster"; + protected static final String BROKER_ADDR = "127.0.0.1:10911"; private static final String TOPIC = "topic"; private static final String CONSUMER_GROUP = "consumerGroup"; private ReceiveMessageActivity receiveMessageActivity; @@ -120,4 +135,46 @@ public class ReceiveMessageActivityTest extends BaseActivityTest { ); assertEquals(Code.MESSAGE_NOT_FOUND, responseArgumentCaptor.getValue().getStatus().getCode()); } + + @Test + public void testReceiveMessageQueueSelector() { + TopicRouteData topicRouteData = new TopicRouteData(); + List queueDatas = new ArrayList<>(); + for (int i = 0; i < 2; i++) { + QueueData queueData = new QueueData(); + queueData.setBrokerName(BROKER_NAME + i); + queueData.setReadQueueNums(1); + queueData.setPerm(PermName.PERM_READ); + queueDatas.add(queueData); + } + topicRouteData.setQueueDatas(queueDatas); + + List brokerDatas = new ArrayList<>(); + for (int i = 0; i < 2; i++) { + BrokerData brokerData = new BrokerData(); + brokerData.setCluster(CLUSTER_NAME); + brokerData.setBrokerName(BROKER_NAME + i); + HashMap brokerAddrs = new HashMap<>(); + brokerAddrs.put(MixAll.MASTER_ID, BROKER_ADDR); + brokerData.setBrokerAddrs(brokerAddrs); + brokerDatas.add(brokerData); + } + topicRouteData.setBrokerDatas(brokerDatas); + + MessageQueueView messageQueueView = new MessageQueueView(TOPIC, topicRouteData); + ReceiveMessageActivity.ReceiveMessageQueueSelector selector = new ReceiveMessageActivity.ReceiveMessageQueueSelector(""); + + SelectableMessageQueue firstSelect = selector.select(ProxyContext.create(), messageQueueView); + SelectableMessageQueue secondSelect = selector.select(ProxyContext.create(), messageQueueView); + SelectableMessageQueue thirdSelect = selector.select(ProxyContext.create(), messageQueueView); + + assertEquals(firstSelect, thirdSelect); + assertNotEquals(firstSelect, secondSelect); + + for (int i = 0; i < 2; i++) { + ReceiveMessageActivity.ReceiveMessageQueueSelector selectorBrokerName = + new ReceiveMessageActivity.ReceiveMessageQueueSelector(BROKER_NAME + i); + assertEquals(BROKER_NAME + i, selectorBrokerName.select(ProxyContext.create(), messageQueueView).getBrokerName()); + } + } } \ No newline at end of file diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java index a23fdd5612..782aaa5e51 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivityTest.java @@ -28,29 +28,41 @@ import apache.rocketmq.v2.SystemProperties; import com.google.protobuf.ByteString; import com.google.protobuf.util.Durations; import com.google.protobuf.util.Timestamps; +import java.util.HashMap; import java.util.concurrent.CompletableFuture; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendStatus; +import org.apache.rocketmq.common.MixAll; +import org.apache.rocketmq.common.constant.PermName; import org.apache.rocketmq.common.message.MessageClientIDSetter; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.common.protocol.route.QueueData; +import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.common.sysflag.MessageSysFlag; import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.grpc.v2.BaseActivityTest; import org.apache.rocketmq.proxy.grpc.v2.common.GrpcProxyException; +import org.apache.rocketmq.proxy.service.route.MessageQueueView; +import org.apache.rocketmq.proxy.service.route.SelectableMessageQueue; import org.apache.rocketmq.remoting.common.RemotingUtil; import org.assertj.core.util.Lists; import org.junit.Before; import org.junit.Test; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotEquals; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.when; public class SendMessageActivityTest extends BaseActivityTest { + protected static final String BROKER_NAME = "broker"; + protected static final String CLUSTER_NAME = "cluster"; + protected static final String BROKER_ADDR = "127.0.0.1:10911"; private static final String TOPIC = "topic"; private static final String CONSUMER_GROUP = "consumerGroup"; @@ -239,4 +251,86 @@ public class SendMessageActivityTest extends BaseActivityTest { assertEquals(MessageClientIDSetter.getUniqID(messageExt), msgId); assertEquals(MessageSysFlag.TRANSACTION_PREPARED_TYPE | MessageSysFlag.COMPRESSED_FLAG, messageExt.getSysFlag()); } + + @Test + public void testSendOrderMessageQueueSelector() { + TopicRouteData topicRouteData = new TopicRouteData(); + QueueData queueData = new QueueData(); + BrokerData brokerData = new BrokerData(); + queueData.setBrokerName(BROKER_NAME); + queueData.setWriteQueueNums(8); + queueData.setPerm(PermName.PERM_WRITE); + topicRouteData.setQueueDatas(Lists.newArrayList(queueData)); + brokerData.setCluster(CLUSTER_NAME); + brokerData.setBrokerName(BROKER_NAME); + HashMap brokerAddrs = new HashMap<>(); + brokerAddrs.put(MixAll.MASTER_ID, BROKER_ADDR); + brokerData.setBrokerAddrs(brokerAddrs); + topicRouteData.setBrokerDatas(Lists.newArrayList(brokerData)); + + MessageQueueView messageQueueView = new MessageQueueView(TOPIC, topicRouteData); + SendMessageActivity.SendMessageQueueSelector selector1 = new SendMessageActivity.SendMessageQueueSelector( + SendMessageRequest.newBuilder() + .addMessages(Message.newBuilder() + .setSystemProperties(SystemProperties.newBuilder() + .setMessageGroup(String.valueOf(1)) + .build()) + .build()) + .build() + ); + + SendMessageActivity.SendMessageQueueSelector selector2 = new SendMessageActivity.SendMessageQueueSelector( + SendMessageRequest.newBuilder() + .addMessages(Message.newBuilder() + .setSystemProperties(SystemProperties.newBuilder() + .setMessageGroup(String.valueOf(1)) + .build()) + .build()) + .build() + ); + + SendMessageActivity.SendMessageQueueSelector selector3 = new SendMessageActivity.SendMessageQueueSelector( + SendMessageRequest.newBuilder() + .addMessages(Message.newBuilder() + .setSystemProperties(SystemProperties.newBuilder() + .setMessageGroup(String.valueOf(2)) + .build()) + .build()) + .build() + ); + + assertEquals(selector1.select(ProxyContext.create(), messageQueueView), selector2.select(ProxyContext.create(), messageQueueView)); + assertNotEquals(selector1.select(ProxyContext.create(), messageQueueView), selector3.select(ProxyContext.create(), messageQueueView)); + } + + @Test + public void testSendNormalMessageQueueSelector() { + TopicRouteData topicRouteData = new TopicRouteData(); + QueueData queueData = new QueueData(); + BrokerData brokerData = new BrokerData(); + queueData.setBrokerName(BROKER_NAME); + queueData.setWriteQueueNums(2); + queueData.setPerm(PermName.PERM_WRITE); + topicRouteData.setQueueDatas(Lists.newArrayList(queueData)); + brokerData.setCluster(CLUSTER_NAME); + brokerData.setBrokerName(BROKER_NAME); + HashMap brokerAddrs = new HashMap<>(); + brokerAddrs.put(MixAll.MASTER_ID, BROKER_ADDR); + brokerData.setBrokerAddrs(brokerAddrs); + topicRouteData.setBrokerDatas(Lists.newArrayList(brokerData)); + + MessageQueueView messageQueueView = new MessageQueueView(TOPIC, topicRouteData); + SendMessageActivity.SendMessageQueueSelector selector = new SendMessageActivity.SendMessageQueueSelector( + SendMessageRequest.newBuilder() + .addMessages(Message.newBuilder().build()) + .build() + ); + + SelectableMessageQueue firstSelect = selector.select(ProxyContext.create(), messageQueueView); + SelectableMessageQueue secondSelect = selector.select(ProxyContext.create(), messageQueueView); + SelectableMessageQueue thirdSelect = selector.select(ProxyContext.create(), messageQueueView); + + assertEquals(firstSelect, thirdSelect); + assertNotEquals(firstSelect, secondSelect); + } } \ No newline at end of file