[ISSUE #3949] v2 test cases

This commit is contained in:
kaiyi.lk
2022-07-13 11:29:18 +08:00
committed by zhouxiang
parent 58241f374a
commit 78862bcd6b
12 changed files with 289 additions and 42 deletions
@@ -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;
@@ -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) {
@@ -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<String, ClientSettings> 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);
}
@@ -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();
@@ -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<RemotingCommand> 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<AckResult> resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader);
resultFuture
.thenAccept(result -> {
@@ -114,8 +114,7 @@ public class RouteService extends BaseService {
List<MessageQueue> 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()
@@ -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<MessageExt> 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()
@@ -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<QueryRouteResponse> 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<QueryRouteResponse> 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<QueryRouteResponse> 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<QueryRouteResponse> 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<QueryAssignmentResponse> 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<QueryAssignmentResponse> 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<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
.setTopic(Resource.newBuilder()
@@ -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);
@@ -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<String, BrokerData> 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<ClientOverwrittenSettings> 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<ReceiveMessageResponse> 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);
}
}
@@ -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();
}
@@ -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