From 78862bcd6b3e33b8beee8f432a3c7d6ae75e7d43 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Tue, 12 Apr 2022 20:49:14 +0800 Subject: [PATCH] [ISSUE #3949] v2 test cases --- .../rocketmq/proxy/config/ProxyConfig.java | 2 +- .../proxy/grpc/v2/adapter/GrpcConverter.java | 20 +- .../grpc/v2/service/GrpcClientManager.java | 7 + .../grpc/v2/service/LocalGrpcService.java | 3 +- .../v2/service/cluster/ConsumerService.java | 7 +- .../grpc/v2/service/cluster/RouteService.java | 6 +- .../service/cluster/ConsumerServiceTest.java | 13 +- .../v2/service/cluster/RouteServiceTest.java | 16 +- .../apache/rocketmq/test/base/BaseConf.java | 11 + .../test/grpc/v2/ClusterGrpcTest.java | 235 ++++++++++++++++++ .../rocketmq/test/grpc/v2/GrpcBaseTest.java | 3 - .../rocketmq/test/grpc/v2/LocalGrpcTest.java | 8 +- 12 files changed, 289 insertions(+), 42 deletions(-) create mode 100644 test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcTest.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java index 1375abc3e7..5b39f5bccc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java @@ -80,7 +80,7 @@ public class ProxyConfig { private int longPollingReserveTimeInMillis = 10000; - private int retryDelayLevelDelta = 3; + private int retryDelayLevelDelta = 2; private String messageDelayLevel = "1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h"; private boolean enableACL = false; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java index 5819ce6dcc..40256e3193 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java @@ -205,7 +205,7 @@ public class GrpcConverter { return requestHeader; } - public static PopMessageRequestHeader buildPopMessageRequestHeader(ReceiveMessageRequest request, long pollTime) { + public static PopMessageRequestHeader buildPopMessageRequestHeader(ReceiveMessageRequest request, long pollTime, boolean isFifo) { Resource group = request.getGroup(); String groupName = GrpcConverter.wrapResourceWithNamespace(group); MessageQueue messageQueue = request.getMessageQueue(); @@ -219,7 +219,7 @@ public class GrpcConverter { maxMessageNumbers = ProxyUtils.MAX_MSG_NUMS_FOR_POP_REQUEST; } long invisibleTime = Durations.toMillis(request.getInvisibleDuration()); - long bornTime = Timestamps.toMillis(request.getInitializationTimestamp()); + long bornTime = System.currentTimeMillis(); FilterExpression filterExpression = request.getFilterExpression(); String expression = filterExpression.getExpression(); @@ -236,7 +236,7 @@ public class GrpcConverter { requestHeader.setInitMode(ConsumeInitMode.MAX); requestHeader.setExpType(expressionType); requestHeader.setExp(expression); - requestHeader.setOrder(request.getFifo()); + requestHeader.setOrder(isFifo); return requestHeader; } @@ -602,13 +602,6 @@ public class GrpcConverter { systemPropertiesBuilder.setStoreHost(storeHost.toString()); } - // delay_level - // TODO: delete -// String delayLevel = messageExt.getProperty(MessageConst.PROPERTY_DELAY_TIME_LEVEL); -// if (delayLevel != null) { -// systemAttributeBuilder.setDelayLevel(Integer.parseInt(delayLevel)); -// } - // delivery_timestamp String deliverMsString; long deliverMs; @@ -645,13 +638,6 @@ public class GrpcConverter { // delivery_attempt systemPropertiesBuilder.setDeliveryAttempt(messageExt.getReconsumeTimes() + 1); - // publisher_group - // TODO: delete -// String producerGroup = messageExt.getProperty(MessageConst.PROPERTY_PRODUCER_GROUP); -// if (producerGroup != null) { -// systemAttributeBuilder.setProducerGroup(buildResource(producerGroup)); -// } - // trace context String traceContext = messageExt.getProperty(MessageConst.PROPERTY_TRACE_CONTEXT); if (traceContext != null) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java index eb86c41cfc..e3705c8532 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java @@ -18,13 +18,20 @@ package org.apache.rocketmq.proxy.grpc.v2.service; import apache.rocketmq.v2.ClientSettings; +import io.grpc.Context; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; +import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; public class GrpcClientManager { private static final Map CLIENT_SETTINGS_MAP = new ConcurrentHashMap<>(); + public ClientSettings getClientSettings(Context ctx) { + String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); + return CLIENT_SETTINGS_MAP.get(clientId); + } + public ClientSettings getClientSettings(String clientId) { return CLIENT_SETTINGS_MAP.get(clientId); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java index 7a4c964e40..eaa92a396e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcService.java @@ -252,7 +252,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo long pollTime = GrpcConverter.buildPollTimeFromContext(ctx); String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId); - PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime); + PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime, + clientSettings.getSettings().getSubscription().getFifo()); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader); command.makeCustomHeaderToNet(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java index ac744b4b3c..303cc5c25b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java @@ -130,7 +130,8 @@ public class ConsumerService extends BaseService { protected PopMessageRequestHeader buildPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) { checkSubscriptionData(request.getMessageQueue().getTopic(), request.getFilterExpression()); - return GrpcConverter.buildPopMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx)); + boolean isFifo = grpcClientManager.getClientSettings(ctx).getSettings().getSubscription().getFifo(); + return GrpcConverter.buildPopMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx), isFifo); } protected ReceiveMessageResponse convertToReceiveMessageResponse(Context ctx, ReceiveMessageRequest request, PopResult result) { @@ -254,11 +255,10 @@ public class ConsumerService extends BaseService { } }); try { - String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle()); String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); - Settings settings = grpcClientManager.getClientSettings(clientId).getSettings(); + Settings settings = grpcClientManager.getClientSettings(ctx).getSettings(); int maxDeliveryAttempts = settings.getSubscription().getDeadLetterPolicy().getMaxDeliveryAttempts(); if (request.getDeliveryAttempt() >= maxDeliveryAttempts) { CompletableFuture resultFuture = this.producer.sendMessageBack( @@ -280,6 +280,7 @@ public class ConsumerService extends BaseService { }); } else { ChangeInvisibleTimeRequestHeader requestHeader = this.buildChangeInvisibleTimeRequestHeader(ctx, request); + System.out.println(requestHeader); CompletableFuture resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader); resultFuture .thenAccept(result -> { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java index 2b240eca41..8ebea5c10a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java @@ -114,8 +114,7 @@ public class RouteService extends BaseService { List messageQueueList = new ArrayList<>(); if (ProxyMode.isClusterMode(mode.name())) { - String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); - ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId); + ClientSettings clientSettings = grpcClientManager.getClientSettings(ctx); Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) { future.complete(QueryRouteResponse.newBuilder() @@ -246,8 +245,7 @@ public class RouteService extends BaseService { } } if (ProxyMode.isClusterMode(mode)) { - String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); - ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId); + ClientSettings clientSettings = grpcClientManager.getClientSettings(ctx); Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) { future.complete(QueryAssignmentResponse.newBuilder() diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java index ccfb2a0fa0..5efa590fba 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerServiceTest.java @@ -62,6 +62,15 @@ public class ConsumerServiceTest extends BaseServiceTest { new MessageQueue("namespace%topic", "brokerName", 0), "brokerAddr"); when(readQueueSelector.select(any(), any(), any())).thenReturn(selectableMessageQueue); + ClientSettings clientSettings = ClientSettings.newBuilder() + .setSettings(Settings.newBuilder() + .setSubscription(Subscription.newBuilder() + .setFifo(false) + .build()) + .build()) + .build(); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); + List messageExtList = Lists.newArrayList( createMessageExt("msg1", "msg1"), createMessageExt("msg2", "msg2") @@ -129,7 +138,7 @@ public class ConsumerServiceTest extends BaseServiceTest { when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); ClientSettings clientSettings = createClientSettings(3); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); NackMessageResponse response = consumerService.nackMessage(Context.current(), NackMessageRequest.newBuilder() .setTopic(Resource.newBuilder() @@ -160,7 +169,7 @@ public class ConsumerServiceTest extends BaseServiceTest { when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); ClientSettings clientSettings = createClientSettings(3); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); NackMessageResponse response = consumerService.nackMessage(Context.current(), NackMessageRequest.newBuilder() .setTopic(Resource.newBuilder() diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteServiceTest.java index be88441fa3..62ca554f38 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteServiceTest.java @@ -48,7 +48,7 @@ import org.junit.Test; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; -import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.when; public class RouteServiceTest extends BaseServiceTest { @@ -159,7 +159,7 @@ public class RouteServiceTest extends BaseServiceTest { .setScheme(AddressScheme.DOMAIN_NAME) .build()) .build(); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() .setTopic(Resource.newBuilder() @@ -177,7 +177,7 @@ public class RouteServiceTest extends BaseServiceTest { public void testQueryRouteWithInvalidEndpoints() throws Exception { RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(ClientSettings.getDefaultInstance()); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(ClientSettings.getDefaultInstance()); CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() .setTopic(Resource.newBuilder() .setName("topic") @@ -201,7 +201,7 @@ public class RouteServiceTest extends BaseServiceTest { .setScheme(AddressScheme.DOMAIN_NAME) .build()) .build(); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() .setTopic(Resource.newBuilder() @@ -229,7 +229,7 @@ public class RouteServiceTest extends BaseServiceTest { .setScheme(AddressScheme.DOMAIN_NAME) .build()) .build(); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() .setTopic(Resource.newBuilder() @@ -245,7 +245,7 @@ public class RouteServiceTest extends BaseServiceTest { public void testQueryAssignmentInvalidEndpoints() throws Exception { RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(ClientSettings.getDefaultInstance()); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(ClientSettings.getDefaultInstance()); CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() .setTopic( Resource.newBuilder() @@ -271,7 +271,7 @@ public class RouteServiceTest extends BaseServiceTest { .setScheme(AddressScheme.DOMAIN_NAME) .build()) .build(); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() .setTopic(Resource.newBuilder() @@ -302,7 +302,7 @@ public class RouteServiceTest extends BaseServiceTest { .setScheme(AddressScheme.DOMAIN_NAME) .build()) .build(); - when(grpcClientManager.getClientSettings(anyString())).thenReturn(clientSettings); + when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() .setTopic(Resource.newBuilder() diff --git a/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java b/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java index 4e29c84c6b..f1d24d20bf 100644 --- a/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java +++ b/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java @@ -48,6 +48,7 @@ import org.apache.rocketmq.test.factory.ConsumerFactory; import org.apache.rocketmq.test.listener.AbstractListener; import org.apache.rocketmq.test.util.MQAdminTestUtils; import org.apache.rocketmq.test.util.MQRandomUtils; +import org.apache.rocketmq.test.util.RandomUtils; import org.apache.rocketmq.tools.admin.DefaultMQAdminExt; import org.apache.rocketmq.tools.admin.MQAdminExt; import org.junit.Assert; @@ -140,11 +141,21 @@ public class BaseConf { return initTopicWithName(topic); } + public static String initTopicOnSampleTopicBroker(String sampleTopic) { + String topic = RandomUtils.getStringWithNumber(10); + return initTopicOnSampleTopicBroker(topic, sampleTopic); + } + public static String initTopicWithName(String topicName) { IntegrationTestBase.initTopic(topicName, nsAddr, clusterName, CQType.SimpleCQ); return topicName; } + public static String initTopicOnSampleTopicBroker(String topicName, String sampleTopic) { + IntegrationTestBase.initTopic(topicName, nsAddr, sampleTopic, CQType.SimpleCQ); + return topicName; + } + public static String initConsumerGroup() { String group = MQRandomUtils.getRandomConsumerGroup(); return initConsumerGroup(group); diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcTest.java new file mode 100644 index 0000000000..3f1725c3d2 --- /dev/null +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcTest.java @@ -0,0 +1,235 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.rocketmq.test.grpc.v2; + +import apache.rocketmq.v2.AckMessageResponse; +import apache.rocketmq.v2.Address; +import apache.rocketmq.v2.AddressScheme; +import apache.rocketmq.v2.ClientOverwrittenSettings; +import apache.rocketmq.v2.ClientSettings; +import apache.rocketmq.v2.ClientType; +import apache.rocketmq.v2.DeadLetterPolicy; +import apache.rocketmq.v2.Endpoints; +import apache.rocketmq.v2.Message; +import apache.rocketmq.v2.MessagingServiceGrpc; +import apache.rocketmq.v2.NackMessageResponse; +import apache.rocketmq.v2.QueryRouteResponse; +import apache.rocketmq.v2.ReceiveMessageResponse; +import apache.rocketmq.v2.SendMessageResponse; +import apache.rocketmq.v2.Settings; +import apache.rocketmq.v2.Subscription; +import io.grpc.Channel; +import java.net.URL; +import java.time.Duration; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer; +import org.apache.rocketmq.client.consumer.PullResult; +import org.apache.rocketmq.client.consumer.PullStatus; +import org.apache.rocketmq.common.MixAll; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.message.MessageQueue; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.proxy.config.ConfigurationManager; +import org.apache.rocketmq.proxy.grpc.v2.GrpcMessagingProcessor; +import org.apache.rocketmq.proxy.grpc.v2.service.ClusterGrpcService; +import org.apache.rocketmq.proxy.grpc.v2.service.GrpcForwardService; +import org.apache.rocketmq.test.util.MQAdminTestUtils; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import static org.apache.rocketmq.common.message.MessageClientIDSetter.createUniqID; +import static org.apache.rocketmq.proxy.config.ConfigurationManager.RMQ_PROXY_HOME; +import static org.awaitility.Awaitility.await; + +public class ClusterGrpcTest extends GrpcBaseTest { + + private final int PORT = 8082; + private GrpcForwardService grpcForwardService; + private MessagingServiceGrpc.MessagingServiceBlockingStub blockingStub; + private MessagingServiceGrpc.MessagingServiceStub stub; + + @Before + public void setUp() throws Exception { + super.setUp(); + String mockProxyHome = "/mock/rmq/proxy/home"; + URL mockProxyHomeURL = getClass().getClassLoader().getResource("rmq-proxy-home"); + if (mockProxyHomeURL != null) { + mockProxyHome = mockProxyHomeURL.toURI().getPath(); + } + System.setProperty(RMQ_PROXY_HOME, mockProxyHome); + ConfigurationManager.initEnv(); + ConfigurationManager.intConfig(); + ConfigurationManager.getProxyConfig().setGrpcServerPort(PORT); + ConfigurationManager.getProxyConfig().setNameSrvAddr(nsAddr); + grpcForwardService = new ClusterGrpcService(); + grpcForwardService.start(); + GrpcMessagingProcessor processor = new GrpcMessagingProcessor(grpcForwardService); + setUpServer(processor, ConfigurationManager.getProxyConfig().getGrpcServerPort(), true); + blockingStub = createBlockingStub(createChannel(ConfigurationManager.getProxyConfig().getGrpcServerPort())); + stub = createStub(createChannel(ConfigurationManager.getProxyConfig().getGrpcServerPort())); + + System.out.println(nsAddr); + await().atMost(Duration.ofSeconds(40)).until(() -> { + Map brokerDataMap = MQAdminTestUtils.getCluster(nsAddr).getBrokerAddrTable(); + return brokerDataMap.size() == brokerNum; + }); + System.out.println(MQAdminTestUtils.getCluster(nsAddr)); + } + + @After + public void tearDown() throws Exception { + grpcForwardService.shutdown(); + shutdown(); + } + + @Test + public void testQueryRoute() throws Exception { + String topic = initTopic(); + String requestId = UUID.randomUUID().toString(); + CompletableFuture future = this.sendClientSettings(stub, ClientSettings.newBuilder() + .setNonce(requestId) + .setAccessPoint(Endpoints.newBuilder() + .setScheme(AddressScheme.IPv4) + .addAddresses(Address.newBuilder() + .setHost("127.0.0.1") + .setPort(PORT) + .build()) + .build()) + .build()); +// System.out.println(future.get()); + +// TimeUnit.SECONDS.sleep(3); + QueryRouteResponse response = blockingStub.queryRoute(buildQueryRouteRequest(topic)); + assertQueryRoute(response, brokerControllerList.size()); + } + + @Test + public void testSendReceiveMessage() throws Exception { + String topic = initTopicOnSampleTopicBroker(broker1Name); + this.sendClientSettings(stub, ClientSettings.newBuilder() + .setNonce(UUID.randomUUID().toString()) + .setClientType(ClientType.PRODUCER) + .build()) + .get(); + + String group = "group"; + String messageId = createUniqID(); + SendMessageResponse sendResponse = blockingStub.sendMessage(buildSendMessageRequest(topic, messageId)); + assertSendMessage(sendResponse, messageId); + + this.sendClientSettings(stub, ClientSettings.newBuilder() + .setNonce(UUID.randomUUID().toString()) + .setClientType(ClientType.PUSH_CONSUMER) + .setSettings(Settings.newBuilder() + .setSubscription(Subscription.newBuilder() + .setFifo(false) + .build()) + .build()) + .build()) + .get(); + + ReceiveMessageResponse receiveResponse = blockingStub.withDeadlineAfter(3, TimeUnit.SECONDS) + .receiveMessage(buildReceiveMessageRequest(group, topic)); + assertReceiveMessage(receiveResponse, messageId); + String receiptHandle = receiveResponse.getMessages(0).getSystemProperties().getReceiptHandle(); + AckMessageResponse ackMessageResponse = blockingStub.ackMessage(buildAckMessageRequest(group, topic, receiptHandle)); + assertAck(ackMessageResponse); + } + + @Test + public void testSendReceiveMessageThenToDLQ() throws Exception { + String topic = initTopicOnSampleTopicBroker(broker1Name); + this.sendClientSettings(stub, ClientSettings.newBuilder() + .setNonce(UUID.randomUUID().toString()) + .setClientType(ClientType.PRODUCER) + .build()) + .get(); + + String group = "group"; + String messageId = createUniqID(); + SendMessageResponse sendResponse = blockingStub.sendMessage(buildSendMessageRequest(topic, messageId)); + assertSendMessage(sendResponse, messageId); + + this.sendClientSettings(stub, ClientSettings.newBuilder() + .setNonce(UUID.randomUUID().toString()) + .setClientType(ClientType.PUSH_CONSUMER) + .setSettings(Settings.newBuilder() + .setSubscription(Subscription.newBuilder() + .setDeadLetterPolicy(DeadLetterPolicy.newBuilder() + .setMaxDeliveryAttempts(2) + .build()) + .setFifo(false) + .build()) + .build()) + .build()) + .get(); + + ReceiveMessageResponse receiveResponse = blockingStub.withDeadlineAfter(20, TimeUnit.SECONDS) + .receiveMessage(buildReceiveMessageRequest(group, topic)); + assertReceiveMessage(receiveResponse, messageId); + + Message message = receiveResponse.getMessages(0); + NackMessageResponse nackMessageResponse = blockingStub.nackMessage(buildNackMessageRequest( + group, topic, messageId, message.getSystemProperties().getReceiptHandle(), 1 + )); + assertNackMessageResponse(nackMessageResponse); + + AtomicReference receiveRetryResponseRef = new AtomicReference<>(); + await().atMost(Duration.ofSeconds(60)).until(() -> { + ReceiveMessageResponse receiveRetryResponse = blockingStub.withDeadlineAfter(20, TimeUnit.SECONDS) + .receiveMessage(buildReceiveMessageRequest(group, topic)); + if (receiveRetryResponse.getMessagesCount() <= 0) { + return false; + } + receiveRetryResponseRef.set(receiveRetryResponse); + return receiveRetryResponse.getMessages(0).getSystemProperties() + .getMessageId().equals(messageId); + }); + + message = receiveRetryResponseRef.get().getMessages(0); + nackMessageResponse = blockingStub.nackMessage(buildNackMessageRequest( + group, topic, messageId, message.getSystemProperties().getReceiptHandle(), 2 + )); + assertNackMessageResponse(nackMessageResponse); + + DefaultMQPullConsumer defaultMQPullConsumer = new DefaultMQPullConsumer(group); + defaultMQPullConsumer.start(); + MessageQueue dlqMQ = new MessageQueue(MixAll.getDLQTopic(group), topic, 0); + await().atMost(Duration.ofSeconds(10)).until(() -> { + try { + PullResult pullResult = defaultMQPullConsumer.pull(dlqMQ, "*", 0L, 1); + if (!PullStatus.FOUND.equals(pullResult.getPullStatus())) { + return false; + } + MessageExt messageExt = pullResult.getMsgFoundList().get(0); + return messageId.equals(messageExt.getMsgId()); + } catch (Throwable ignore) { + return false; + } + }); + + System.out.println(1); + } + + +} diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java index 35f44032d1..f94c41c3f3 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java @@ -192,9 +192,6 @@ public class GrpcBaseTest extends BaseConf { .setInvisibleDuration(Duration.newBuilder() .setSeconds(3) .build()) - .setInitializationTimestamp(Timestamp.newBuilder() - .setSeconds(TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis())) - .build()) .build(); } diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java index 3b2610848e..30c7c50270 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcTest.java @@ -22,7 +22,6 @@ import apache.rocketmq.v2.MessagingServiceGrpc; import apache.rocketmq.v2.QueryRouteResponse; import apache.rocketmq.v2.ReceiveMessageResponse; import apache.rocketmq.v2.SendMessageResponse; -import io.grpc.Channel; import java.net.URL; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.proxy.config.ConfigurationManager; @@ -37,10 +36,12 @@ import static org.apache.rocketmq.proxy.config.ConfigurationManager.RMQ_PROXY_HO public class LocalGrpcTest extends GrpcBaseTest { private MessagingServiceGrpc.MessagingServiceBlockingStub blockingStub; + private MessagingServiceGrpc.MessagingServiceStub stub; private LocalGrpcService localGrpcService; @Before public void setUp() throws Exception { + super.setUp(); String mockProxyHome = "/mock/rmq/proxy/home"; URL mockProxyHomeURL = getClass().getClassLoader().getResource("rmq-proxy-home"); if (mockProxyHomeURL != null) { @@ -54,8 +55,9 @@ public class LocalGrpcTest extends GrpcBaseTest { localGrpcService = new LocalGrpcService(brokerController1); localGrpcService.start(); GrpcMessagingProcessor processor = new GrpcMessagingProcessor(localGrpcService); - Channel channel = setUpServer(processor, ConfigurationManager.getProxyConfig().getGrpcServerPort(), true); - blockingStub = MessagingServiceGrpc.newBlockingStub(channel); + setUpServer(processor, ConfigurationManager.getProxyConfig().getGrpcServerPort(), true); + blockingStub = createBlockingStub(createChannel(ConfigurationManager.getProxyConfig().getGrpcServerPort())); + stub = createStub(createChannel(ConfigurationManager.getProxyConfig().getGrpcServerPort())); } @After