From 7ce66f1793fbc4250bc2c6bdce6b232a3bb15154 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Mon, 25 Apr 2022 17:10:54 +0800 Subject: [PATCH] [ISSUE #3949] v2 support --- .../grpc/v2/service/ClusterGrpcService.java | 7 +- .../grpc/v2/service/cluster/BaseService.java | 24 ++- .../v2/service/cluster/ConsumerService.java | 163 +++-------------- .../DefaultReceiveMessageResultFilter.java | 173 ++++++++++++++++++ .../service/cluster/ForwardClientService.java | 4 +- .../v2/service/cluster/ProducerService.java | 8 +- .../cluster/ReceiveMessageResultFilter.java | 29 +++ .../grpc/v2/service/cluster/RouteService.java | 168 +++++------------ .../service/cluster/ConsumerServiceTest.java | 9 +- .../cluster/ForwardClientServiceTest.java | 1 + .../service/cluster/ProducerServiceTest.java | 9 +- .../v2/service/cluster/RouteServiceTest.java | 57 +----- 12 files changed, 324 insertions(+), 328 deletions(-) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/DefaultReceiveMessageResultFilter.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ReceiveMessageResultFilter.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java index 9c4d663d67..bea689d077 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java @@ -91,13 +91,18 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc this.connectorManager = new ConnectorManager(new GrpcTransactionStateChecker()); this.consumerService = new ConsumerService(connectorManager, grpcClientManager); this.producerService = new ProducerService(connectorManager); - this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager, grpcClientManager); + this.routeService = new RouteService(connectorManager, grpcClientManager); this.clientService = new ForwardClientService(connectorManager, scheduledExecutorService, channelManager, grpcClientManager, pollCommandResponseManager); this.transactionService = new TransactionService(connectorManager, channelManager); this.appendStartAndShutdown(new ClusterGrpcServiceStartAndShutdown()); this.appendStartAndShutdown(this.connectorManager); + this.appendStartAndShutdown(this.consumerService); + this.appendStartAndShutdown(this.producerService); + this.appendStartAndShutdown(this.routeService); + this.appendStartAndShutdown(this.clientService); + this.appendStartAndShutdown(this.transactionService); } @Override diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/BaseService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/BaseService.java index 835e31f6b8..9843cf7d3a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/BaseService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/BaseService.java @@ -22,11 +22,13 @@ import apache.rocketmq.v2.Resource; import io.grpc.Context; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.consumer.ReceiptHandle; +import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.connector.ConnectorManager; +import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException; -public class BaseService { +public abstract class BaseService implements StartAndShutdown { protected final ConnectorManager connectorManager; @@ -34,7 +36,7 @@ public class BaseService { this.connectorManager = connectorManager; } - protected ReceiptHandle resolveReceiptHandle(Context ctx, String receiptHandleStr) { + public static ReceiptHandle resolveReceiptHandle(Context ctx, String receiptHandleStr) { ReceiptHandle receiptHandle = ReceiptHandle.decode(receiptHandleStr); if (receiptHandle.isExpired()) { throw new ProxyException(Code.RECEIPT_HANDLE_EXPIRED, "handle has expired"); @@ -42,20 +44,34 @@ public class BaseService { return receiptHandle; } - protected String getBrokerAddr(Context ctx, String brokerName) throws Exception { + public static String getBrokerAddr(Context ctx, TopicRouteCache topicRouteCache, String brokerName) throws Exception { if (StringUtils.isBlank(brokerName)) { throw new ProxyException(Code.UNRECOGNIZED, "broker name is empty"); } - String addr = this.connectorManager.getTopicRouteCache().getBrokerAddr(brokerName); + String addr = topicRouteCache.getBrokerAddr(brokerName); if (StringUtils.isBlank(addr)) { throw new ProxyException(Code.UNRECOGNIZED, brokerName + " not exist"); } return addr; } + protected String getBrokerAddr(Context ctx, String brokerName) throws Exception { + return getBrokerAddr(ctx, this.connectorManager.getTopicRouteCache(), brokerName); + } + protected void checkSubscriptionData(Resource topic, FilterExpression filterExpression) { // for checking filterExpression. String topicName = GrpcConverter.wrapResourceWithNamespace(topic); GrpcConverter.buildSubscriptionData(topicName, filterExpression); } + + @Override + public void start() throws Exception { + + } + + @Override + public void shutdown() throws Exception { + + } } 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 fb8a277353..56a917ad68 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 @@ -22,20 +22,17 @@ import apache.rocketmq.v2.AckMessageResponse; import apache.rocketmq.v2.AckMessageResultEntry; import apache.rocketmq.v2.ChangeInvisibleDurationRequest; import apache.rocketmq.v2.ChangeInvisibleDurationResponse; -import apache.rocketmq.v2.ClientType; import apache.rocketmq.v2.Code; import apache.rocketmq.v2.Message; import apache.rocketmq.v2.NackMessageRequest; import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.ReceiveMessageResponse; -import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.RetryPolicy; import apache.rocketmq.v2.Settings; import io.grpc.Context; import io.grpc.stub.StreamObserver; import java.util.ArrayList; -import java.util.Collections; import java.util.List; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.AckResult; @@ -43,14 +40,10 @@ import org.apache.rocketmq.client.consumer.AckStatus; import org.apache.rocketmq.client.consumer.PopResult; import org.apache.rocketmq.client.consumer.PopStatus; import org.apache.rocketmq.common.consumer.ReceiptHandle; -import org.apache.rocketmq.common.message.MessageExt; -import org.apache.rocketmq.common.protocol.ResponseCode; import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader; import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; -import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; -import org.apache.rocketmq.proxy.common.utils.FilterUtils; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.ForwardProducer; import org.apache.rocketmq.proxy.connector.ForwardReadConsumer; @@ -74,9 +67,8 @@ public class ConsumerService extends BaseService { protected final GrpcClientManager grpcClientManager; private volatile ReadQueueSelector readQueueSelector; + private volatile ReceiveMessageResultFilter receiveMessageResultFilter; private volatile ResponseHook> receiveMessageHook; - private volatile ResponseHook ackNoMatchedMessageHook; - private volatile ResponseHook forwardToDLQInRecvMessageHook; private volatile ResponseHook ackMessageHook; private volatile ResponseHook nackMessageHook; private volatile ResponseHook changeInvisibleDurationHook; @@ -86,12 +78,16 @@ public class ConsumerService extends BaseService { this.readConsumer = connectorManager.getForwardReadConsumer(); this.writeConsumer = connectorManager.getForwardWriteConsumer(); this.producer = connectorManager.getForwardProducer(); - - this.readQueueSelector = new DefaultReadQueueSelector(connectorManager.getTopicRouteCache()); - this.grpcClientManager = grpcClientManager; } + @Override + public void start() throws Exception { + this.readQueueSelector = new DefaultReadQueueSelector(connectorManager.getTopicRouteCache()); + this.receiveMessageResultFilter = new DefaultReceiveMessageResultFilter( + producer, writeConsumer, grpcClientManager, connectorManager.getTopicRouteCache()); + } + public void receiveMessage(Context ctx, ReceiveMessageRequest request, StreamObserver responseObserver) { this.receiveMessage(ctx, request) @@ -146,10 +142,10 @@ public class ConsumerService extends BaseService { PopStatus status = result.getPopStatus(); switch (status) { case FOUND: - List messageList = filterMessage(ctx, request, result.getMsgFoundList()); + List messageList = this.receiveMessageResultFilter.filterMessage(ctx, request, result.getMsgFoundList()); if (messageList.isEmpty()) { responseList.add(ReceiveMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) + .setStatus(ResponseBuilder.buildStatus(Code.OK, "no new message")) .build()); } else { for (Message message : messageList) { @@ -176,97 +172,6 @@ public class ConsumerService extends BaseService { return responseList; } - protected List filterMessage(Context ctx, ReceiveMessageRequest request, List messageExtList) { - if (messageExtList == null || messageExtList.isEmpty()) { - return Collections.emptyList(); - } - Settings settings = grpcClientManager.getClientSettings(ctx); - ClientType clientType = settings.getClientType(); - int maxAttempts = settings.getSubscription().getBackoffPolicy().getMaxAttempts(); - Resource topic = request.getMessageQueue().getTopic(); - String topicName = GrpcConverter.wrapResourceWithNamespace(topic); - SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(topicName, request.getFilterExpression()); - - List resMessageList = new ArrayList<>(); - for (MessageExt messageExt : messageExtList) { - if (ClientType.SIMPLE_CONSUMER.equals(clientType) && messageExt.getReconsumeTimes() >= maxAttempts) { - forwardMessageToDLQ(ctx, request, messageExt, maxAttempts); - continue; - } - if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) { - this.ackNoMatchedMessage(ctx, request, messageExt); - continue; - } - resMessageList.add(GrpcConverter.buildMessage(messageExt)); - } - return resMessageList; - } - - protected void forwardMessageToDLQ(Context ctx, ReceiveMessageRequest request, MessageExt messageExt, - int maxReconsumeTimes) { - CompletableFuture future = new CompletableFuture<>(); - ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader(); - - try { - ReceiptHandle handle = ReceiptHandle.create(messageExt); - if (handle == null) { - return; - } - String brokerAddr = this.getBrokerAddr(ctx, handle.getBrokerName()); - Resource topic = request.getMessageQueue().getTopic(); - Resource group = request.getGroup(); - ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader( - topic, - group, - handle, - messageExt.getMsgId(), - maxReconsumeTimes); - AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader( - topic, - group, - handle); - - future = this.producer.sendMessageBackThenAckOrg(ctx, brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader); - } catch (Throwable t) { - future.completeExceptionally(t); - } - - future.whenComplete((result, throwable) -> { - if (forwardToDLQInRecvMessageHook != null) { - forwardToDLQInRecvMessageHook.beforeResponse(ctx, consumerSendMsgBackRequestHeader, result, throwable); - } - }); - } - - protected void ackNoMatchedMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) { - CompletableFuture future = new CompletableFuture<>(); - - AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader(); - try { - ReceiptHandle handle = ReceiptHandle.create(messageExt); - if (handle == null) { - return; - } - String brokerAddr = this.getBrokerAddr(ctx, handle.getBrokerName()); - ackMessageRequestHeader.setConsumerGroup(GrpcConverter.wrapResourceWithNamespace(request.getGroup())); - ackMessageRequestHeader.setTopic(messageExt.getTopic()); - ackMessageRequestHeader.setQueueId(handle.getQueueId()); - ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle()); - ackMessageRequestHeader.setOffset(handle.getOffset()); - - future = this.writeConsumer.ackMessage(ctx, brokerAddr, ackMessageRequestHeader); - } catch (Throwable t) { - future.completeExceptionally(t); - } - - future.whenComplete((ackResult, throwable) -> { - if (ackNoMatchedMessageHook != null) { - ackNoMatchedMessageHook.beforeResponse(ctx, ackMessageRequestHeader, ackResult, throwable); - } - }); - - } - public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { CompletableFuture future = new CompletableFuture<>(); future.whenComplete((response, throwable) -> { @@ -309,7 +214,7 @@ public class ConsumerService extends BaseService { .setReceiptHandle(ackMessageEntry.getReceiptHandle()); try { - ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, ackMessageEntry.getReceiptHandle()); + ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, ackMessageEntry.getReceiptHandle()); String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); AckMessageRequestHeader requestHeader = this.buildAckMessageRequestHeader(ctx, request, receiptHandle); @@ -350,25 +255,18 @@ public class ConsumerService extends BaseService { public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { CompletableFuture future = new CompletableFuture<>(); try { - ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle()); + ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, request.getReceiptHandle()); String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); Settings settings = grpcClientManager.getClientSettings(ctx); int maxDeliveryAttempts = settings.getSubscription().getBackoffPolicy().getMaxAttempts(); if (request.getDeliveryAttempt() >= maxDeliveryAttempts) { - future = this.producer.sendMessageBack( + future = this.producer.sendMessageBackThenAckOrg( ctx, brokerAddr, - this.buildConsumerSendMsgBackToDLQRequestHeader(ctx, request, maxDeliveryAttempts) - ).thenApply(result -> { - if (result.getCode() == ResponseCode.SUCCESS) { - writeConsumer.ackMessage( - ctx, - brokerAddr, - this.buildAckMessageRequestHeader(ctx, request)); - } - return convertToNackMessageResponse(ctx, request, result); - }); + this.buildConsumerSendMsgBackToDLQRequestHeader(ctx, request, maxDeliveryAttempts), + this.buildAckMessageRequestHeader(ctx, request) + ).thenApply(result -> convertToNackMessageResponse(ctx, request, result)); } else { ChangeInvisibleTimeRequestHeader requestHeader = this.buildChangeInvisibleTimeRequestHeader(ctx, request); future = this.writeConsumer.changeInvisibleTimeAsync(ctx, brokerAddr, receiptHandle.getBrokerName(), requestHeader) @@ -425,7 +323,7 @@ public class ConsumerService extends BaseService { CompletableFuture future = new CompletableFuture<>(); try { - ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle()); + ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, request.getReceiptHandle()); String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); ChangeInvisibleTimeRequestHeader requestHeader = convertToChangeInvisibleTimeRequestHeader(ctx, request); @@ -468,6 +366,15 @@ public class ConsumerService extends BaseService { this.readQueueSelector = readQueueSelector; } + public ReceiveMessageResultFilter getReceiveMessageResultFilter() { + return receiveMessageResultFilter; + } + + public void setReceiveMessageResultFilter( + ReceiveMessageResultFilter receiveMessageResultFilter) { + this.receiveMessageResultFilter = receiveMessageResultFilter; + } + public ResponseHook> getReceiveMessageHook() { return receiveMessageHook; } @@ -477,24 +384,6 @@ public class ConsumerService extends BaseService { this.receiveMessageHook = receiveMessageHook; } - public ResponseHook getAckNoMatchedMessageHook() { - return ackNoMatchedMessageHook; - } - - public void setAckNoMatchedMessageHook( - ResponseHook ackNoMatchedMessageHook) { - this.ackNoMatchedMessageHook = ackNoMatchedMessageHook; - } - - public ResponseHook getForwardToDLQInRecvMessageHook() { - return forwardToDLQInRecvMessageHook; - } - - public void setForwardToDLQInRecvMessageHook( - ResponseHook forwardToDLQInRecvMessageHook) { - this.forwardToDLQInRecvMessageHook = forwardToDLQInRecvMessageHook; - } - public ResponseHook getAckMessageHook() { return ackMessageHook; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/DefaultReceiveMessageResultFilter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/DefaultReceiveMessageResultFilter.java new file mode 100644 index 0000000000..15b5b35c69 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/DefaultReceiveMessageResultFilter.java @@ -0,0 +1,173 @@ +/* + * 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.proxy.grpc.v2.service.cluster; + +import apache.rocketmq.v2.ClientType; +import apache.rocketmq.v2.Message; +import apache.rocketmq.v2.ReceiveMessageRequest; +import apache.rocketmq.v2.Resource; +import apache.rocketmq.v2.Settings; +import io.grpc.Context; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import org.apache.rocketmq.client.consumer.AckResult; +import org.apache.rocketmq.common.consumer.ReceiptHandle; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader; +import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader; +import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; +import org.apache.rocketmq.proxy.common.utils.FilterUtils; +import org.apache.rocketmq.proxy.connector.ForwardProducer; +import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer; +import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; +import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook; +import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; +import org.apache.rocketmq.remoting.protocol.RemotingCommand; + +import static org.apache.rocketmq.proxy.grpc.v2.service.cluster.BaseService.getBrokerAddr; + +public class DefaultReceiveMessageResultFilter implements ReceiveMessageResultFilter { + + protected final ForwardProducer producer; + protected final ForwardWriteConsumer writeConsumer; + protected final GrpcClientManager grpcClientManager; + protected final TopicRouteCache topicRouteCache; + + private volatile ResponseHook ackNoMatchedMessageHook; + private volatile ResponseHook forwardToDLQInRecvMessageHook; + + public DefaultReceiveMessageResultFilter(ForwardProducer producer, ForwardWriteConsumer writeConsumer, + GrpcClientManager grpcClientManager, TopicRouteCache topicRouteCache) { + this.producer = producer; + this.writeConsumer = writeConsumer; + this.grpcClientManager = grpcClientManager; + this.topicRouteCache = topicRouteCache; + } + + @Override + public List filterMessage(Context ctx, ReceiveMessageRequest request, List messageExtList) { + if (messageExtList == null || messageExtList.isEmpty()) { + return Collections.emptyList(); + } + Settings settings = grpcClientManager.getClientSettings(ctx); + ClientType clientType = settings.getClientType(); + int maxAttempts = settings.getSubscription().getBackoffPolicy().getMaxAttempts(); + Resource topic = request.getMessageQueue().getTopic(); + String topicName = GrpcConverter.wrapResourceWithNamespace(topic); + SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(topicName, request.getFilterExpression()); + + List resMessageList = new ArrayList<>(); + for (MessageExt messageExt : messageExtList) { + if (ClientType.SIMPLE_CONSUMER.equals(clientType) && messageExt.getReconsumeTimes() >= maxAttempts) { + forwardMessageToDLQ(ctx, request, messageExt, maxAttempts); + continue; + } + if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) { + this.ackNoMatchedMessage(ctx, request, messageExt); + continue; + } + resMessageList.add(GrpcConverter.buildMessage(messageExt)); + } + return resMessageList; + } + + protected void forwardMessageToDLQ(Context ctx, ReceiveMessageRequest request, MessageExt messageExt, + int maxReconsumeTimes) { + CompletableFuture future = new CompletableFuture<>(); + ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader(); + + try { + ReceiptHandle handle = ReceiptHandle.create(messageExt); + if (handle == null) { + return; + } + String brokerAddr = getBrokerAddr(ctx, topicRouteCache, handle.getBrokerName()); + Resource topic = request.getMessageQueue().getTopic(); + Resource group = request.getGroup(); + ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader( + topic, + group, + handle, + messageExt.getMsgId(), + maxReconsumeTimes); + AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader( + topic, + group, + handle); + + future = this.producer.sendMessageBackThenAckOrg(ctx, brokerAddr, sendMsgBackRequestHeader, ackMessageRequestHeader); + } catch (Throwable t) { + future.completeExceptionally(t); + } + + future.whenComplete((result, throwable) -> { + if (forwardToDLQInRecvMessageHook != null) { + forwardToDLQInRecvMessageHook.beforeResponse(ctx, consumerSendMsgBackRequestHeader, result, throwable); + } + }); + } + + protected void ackNoMatchedMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) { + CompletableFuture future = new CompletableFuture<>(); + + AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader(); + try { + ReceiptHandle handle = ReceiptHandle.create(messageExt); + if (handle == null) { + return; + } + String brokerAddr = getBrokerAddr(ctx, topicRouteCache, handle.getBrokerName()); + ackMessageRequestHeader.setConsumerGroup(GrpcConverter.wrapResourceWithNamespace(request.getGroup())); + ackMessageRequestHeader.setTopic(messageExt.getTopic()); + ackMessageRequestHeader.setQueueId(handle.getQueueId()); + ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle()); + ackMessageRequestHeader.setOffset(handle.getOffset()); + + future = this.writeConsumer.ackMessage(ctx, brokerAddr, ackMessageRequestHeader); + } catch (Throwable t) { + future.completeExceptionally(t); + } + + future.whenComplete((ackResult, throwable) -> { + if (ackNoMatchedMessageHook != null) { + ackNoMatchedMessageHook.beforeResponse(ctx, ackMessageRequestHeader, ackResult, throwable); + } + }); + } + + public ResponseHook getAckNoMatchedMessageHook() { + return ackNoMatchedMessageHook; + } + + public void setAckNoMatchedMessageHook( + ResponseHook ackNoMatchedMessageHook) { + this.ackNoMatchedMessageHook = ackNoMatchedMessageHook; + } + + public ResponseHook getForwardToDLQInRecvMessageHook() { + return forwardToDLQInRecvMessageHook; + } + + public void setForwardToDLQInRecvMessageHook( + ResponseHook forwardToDLQInRecvMessageHook) { + this.forwardToDLQInRecvMessageHook = forwardToDLQInRecvMessageHook; + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java index 45b1d40705..0c7fde9b43 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java @@ -82,9 +82,11 @@ public class ForwardClientService extends BaseService { this.channelManager = channelManager; this.grpcClientManager = grpcClientManager; this.telemetryCommandManager = telemetryCommandManager; + } + @Override + public void start() throws Exception { this.clientSettingsService = new ClientSettingsService(this.channelManager, this.grpcClientManager, this.telemetryCommandManager); - this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListenerImpl()); this.producerManager = new ProducerManager(); this.producerManager.appendProducerChangeListener(new ProducerChangeListenerImpl()); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java index 38fb8adec8..8bb0b3909e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java @@ -53,7 +53,11 @@ public class ProducerService extends BaseService { public ProducerService(ConnectorManager connectorManager) { super(connectorManager); this.producer = connectorManager.getForwardProducer(); - writeQueueSelector = new DefaultWriteQueueSelector(this.connectorManager.getTopicRouteCache()); + } + + @Override + public void start() throws Exception { + this.writeQueueSelector = new DefaultWriteQueueSelector(this.connectorManager.getTopicRouteCache()); } public CompletableFuture sendMessage(Context ctx, SendMessageRequest request) { @@ -123,7 +127,7 @@ public class ProducerService extends BaseService { CompletableFuture future = new CompletableFuture<>(); try { - ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle()); + ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, request.getReceiptHandle()); String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = this.buildConsumerSendMsgBackRequestHeader(ctx, request); AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader( diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ReceiveMessageResultFilter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ReceiveMessageResultFilter.java new file mode 100644 index 0000000000..cd291b493b --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ReceiveMessageResultFilter.java @@ -0,0 +1,29 @@ +/* + * 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.proxy.grpc.v2.service.cluster; + +import apache.rocketmq.v2.Message; +import apache.rocketmq.v2.ReceiveMessageRequest; +import io.grpc.Context; +import java.util.List; +import org.apache.rocketmq.common.message.MessageExt; + +public interface ReceiveMessageResultFilter { + + List filterMessage(Context ctx, ReceiveMessageRequest request, List messageExtList); +} 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 aba9cef649..d91bcb4ccc 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 @@ -16,8 +16,6 @@ */ package org.apache.rocketmq.proxy.grpc.v2.service.cluster; -import apache.rocketmq.v2.Address; -import apache.rocketmq.v2.AddressScheme; import apache.rocketmq.v2.Assignment; import apache.rocketmq.v2.Broker; import apache.rocketmq.v2.Code; @@ -29,32 +27,23 @@ import apache.rocketmq.v2.QueryAssignmentResponse; import apache.rocketmq.v2.QueryRouteRequest; import apache.rocketmq.v2.QueryRouteResponse; import apache.rocketmq.v2.Settings; -import com.google.common.base.Preconditions; -import com.google.common.net.HostAndPort; import io.grpc.Context; import java.util.ArrayList; -import java.util.HashMap; import java.util.List; -import java.util.Map; import java.util.concurrent.CompletableFuture; -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.ParameterConverter; -import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; -import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook; import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; public class RouteService extends BaseService { - private final ProxyMode mode; - private volatile ParameterConverter queryRouteEndpointConverter; private volatile ResponseHook queryRouteHook; @@ -64,14 +53,16 @@ public class RouteService extends BaseService { protected final GrpcClientManager grpcClientManager; - public RouteService(ProxyMode mode, ConnectorManager connectorManager, GrpcClientManager grpcClientManager) { + public RouteService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) { super(connectorManager); - Preconditions.checkArgument(ProxyMode.isClusterMode(mode) || ProxyMode.isLocalMode(mode)); - this.mode = mode; + this.grpcClientManager = grpcClientManager; + } + + @Override + public void start() throws Exception { this.queryRouteEndpointConverter = (ctx, parameter) -> parameter; this.queryAssignmentEndpointConverter = (ctx, parameter) -> parameter; this.assignmentQueueSelector = new DefaultAssignmentQueueSelector(this.connectorManager.getTopicRouteCache()); - this.grpcClientManager = grpcClientManager; } public CompletableFuture queryRoute(Context ctx, QueryRouteRequest request) { @@ -87,44 +78,26 @@ public class RouteService extends BaseService { MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName); TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData(); List queueDataList = topicRouteData.getQueueDatas(); - List brokerDataList = topicRouteData.getBrokerDatas(); List messageQueueList = new ArrayList<>(); - if (ProxyMode.isClusterMode(mode.name())) { - Settings clientSettings = grpcClientManager.getClientSettings(ctx); - Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); - if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) { - future.complete(QueryRouteResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " + - clientSettings.getAccessPoint() + " is invalidate")) - .build()); - return future; - } - for (QueueData queueData : queueDataList) { - Broker broker = Broker.newBuilder() - .setName(queueData.getBrokerName()) - .setId(0) - .setEndpoints(resEndpoints) - .build(); - - messageQueueList.addAll(GrpcConverter.genMessageQueueFromQueueData(queueData, request.getTopic(), broker)); - } + Settings clientSettings = grpcClientManager.getClientSettings(ctx); + Endpoints resEndpoints = this.queryRouteEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); + if (resEndpoints == null || resEndpoints.getDefaultInstanceForType().equals(resEndpoints)) { + future.complete(QueryRouteResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " + + clientSettings.getAccessPoint() + " is invalidate")) + .build()); + return future; } - if (ProxyMode.isLocalMode(mode.name())) { - Map> brokerMap = buildBrokerMap(brokerDataList); + for (QueueData queueData : queueDataList) { + Broker broker = Broker.newBuilder() + .setName(queueData.getBrokerName()) + .setId(0) + .setEndpoints(resEndpoints) + .build(); - for (QueueData queueData : queueDataList) { - String brokerName = queueData.getBrokerName(); - Map brokerIdMap = brokerMap.get(brokerName); - if (brokerIdMap == null) { - break; - } - for (Broker broker : brokerIdMap.values()) { - messageQueueList.addAll(GrpcConverter.genMessageQueueFromQueueData(queueData, request.getTopic(), broker)); - } - } + messageQueueList.addAll(GrpcConverter.genMessageQueueFromQueueData(queueData, request.getTopic(), broker)); } - QueryRouteResponse response = QueryRouteResponse.newBuilder() .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) .addAllMessageQueues(messageQueueList) @@ -153,57 +126,32 @@ public class RouteService extends BaseService { try { List assignments = new ArrayList<>(); List messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request); - if (ProxyMode.isLocalMode(mode)) { - String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); - MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName); - TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData(); - Map> brokerMap = buildBrokerMap(topicRouteData.getBrokerDatas()); - for (SelectableMessageQueue messageQueue : messageQueueList) { - Map brokerIdMap = brokerMap.get(messageQueue.getBrokerName()); - if (brokerIdMap != null) { - Broker broker = brokerIdMap.get(0L); - - MessageQueue defaultMessageQueue = MessageQueue.newBuilder() - .setTopic(request.getTopic()) - .setId(-1) - .setPermission(Permission.READ_WRITE) - .setBroker(broker) - .build(); - - assignments.add(Assignment.newBuilder() - .setMessageQueue(defaultMessageQueue) - .build()); - } - } + Settings clientSettings = grpcClientManager.getClientSettings(ctx); + Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); + if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) { + future.complete(QueryAssignmentResponse.newBuilder() + .setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " + + clientSettings.getAccessPoint() + " is invalidate")) + .build()); + return future; } - if (ProxyMode.isClusterMode(mode)) { - Settings clientSettings = grpcClientManager.getClientSettings(ctx); - Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint()); - if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) { - future.complete(QueryAssignmentResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.ILLEGAL_ACCESS_POINT, "endpoint " + - clientSettings.getAccessPoint() + " is invalidate")) - .build()); - return future; - } - for (SelectableMessageQueue messageQueue : messageQueueList) { - Broker broker = Broker.newBuilder() - .setName(messageQueue.getBrokerName()) - .setId(0) - .setEndpoints(resEndpoints) - .build(); + for (SelectableMessageQueue messageQueue : messageQueueList) { + Broker broker = Broker.newBuilder() + .setName(messageQueue.getBrokerName()) + .setId(0) + .setEndpoints(resEndpoints) + .build(); - MessageQueue defaultMessageQueue = MessageQueue.newBuilder() - .setTopic(request.getTopic()) - .setId(-1) - .setPermission(Permission.READ_WRITE) - .setBroker(broker) - .build(); + MessageQueue defaultMessageQueue = MessageQueue.newBuilder() + .setTopic(request.getTopic()) + .setId(-1) + .setPermission(Permission.READ_WRITE) + .setBroker(broker) + .build(); - assignments.add(Assignment.newBuilder() - .setMessageQueue(defaultMessageQueue) - .build()); - } + assignments.add(Assignment.newBuilder() + .setMessageQueue(defaultMessageQueue) + .build()); } QueryAssignmentResponse response = QueryAssignmentResponse.newBuilder() @@ -217,34 +165,6 @@ public class RouteService extends BaseService { return future; } - private Map> buildBrokerMap(List brokerDataList) { - Map> brokerMap = new HashMap<>(); - for (BrokerData brokerData : brokerDataList) { - Map brokerIdMap = new HashMap<>(); - String brokerName = brokerData.getBrokerName(); - for (Map.Entry entry : brokerData.getBrokerAddrs().entrySet()) { - Long brokerId = entry.getKey(); - HostAndPort hostAndPort = HostAndPort.fromString(entry.getValue()); - Broker broker = Broker.newBuilder() - .setName(brokerName) - .setId(Math.toIntExact(brokerId)) - .setEndpoints(Endpoints.newBuilder() - .setScheme(AddressScheme.IPv4) - .addAddresses( - Address.newBuilder() - .setPort(ConfigurationManager.getProxyConfig().getGrpcServerPort()) - .setHost(hostAndPort.getHost()) - ) - .build()) - .build(); - - brokerIdMap.put(brokerId, broker); - } - brokerMap.put(brokerName, brokerIdMap); - } - return brokerMap; - } - public ParameterConverter getQueryRouteEndpointConverter() { return queryRouteEndpointConverter; } 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 ee57c45e7e..353c743bc1 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 @@ -51,10 +51,15 @@ public class ConsumerServiceTest extends BaseServiceTest { private ReadQueueSelector readQueueSelector; private ConsumerService consumerService; + private DefaultReceiveMessageResultFilter receiveMessageResultFilter; @Override public void beforeEach() throws Throwable { consumerService = new ConsumerService(this.connectorManager, this.grpcClientManager); + consumerService.start(); + + receiveMessageResultFilter = new DefaultReceiveMessageResultFilter(producerClient, writeConsumerClient, grpcClientManager, topicRouteCache); + consumerService.setReceiveMessageResultFilter(receiveMessageResultFilter); consumerService.setReadQueueSelector(readQueueSelector); } @@ -84,7 +89,7 @@ public class ConsumerServiceTest extends BaseServiceTest { Context ctx = Context.current().withDeadlineAfter(3, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor()); AtomicReference ackHandler = new AtomicReference<>(); - consumerService.setAckNoMatchedMessageHook((ctx1, request, response, t) -> ackHandler.set(request.getExtraInfo())); + receiveMessageResultFilter.setAckNoMatchedMessageHook((ctx1, request, response, t) -> ackHandler.set(request.getExtraInfo())); List responseList = consumerService.receiveMessage(ctx, ReceiveMessageRequest.newBuilder() .setMessageQueue(apache.rocketmq.v2.MessageQueue.newBuilder() @@ -189,7 +194,7 @@ public class ConsumerServiceTest extends BaseServiceTest { doAnswer(mock -> { headerRef.set(mock.getArgument(2)); return CompletableFuture.completedFuture(RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "")); - }).when(producerClient).sendMessageBack(any(), anyString(), any()); + }).when(producerClient).sendMessageBackThenAckOrg(any(), anyString(), any(), any()); when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); Settings clientSettings = createClientSettings(3); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java index 8ea04fe051..ad1bb98c0e 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java @@ -47,6 +47,7 @@ public class ForwardClientServiceTest extends BaseServiceTest { this.channelManager, this.grpcClientManager, this.telemetryCommandManager); + clientService.start(); } @Test diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerServiceTest.java index e350682276..fee74d2176 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerServiceTest.java @@ -46,6 +46,8 @@ import static org.mockito.Mockito.when; public class ProducerServiceTest extends BaseServiceTest { + private ProducerService producerService; + private static final SendMessageRequest REQUEST = SendMessageRequest.newBuilder() .addMessages(Message.newBuilder() .setTopic(Resource.newBuilder() @@ -61,6 +63,8 @@ public class ProducerServiceTest extends BaseServiceTest { @Override public void beforeEach() throws Throwable { + producerService = new ProducerService(this.connectorManager); + producerService.start(); } @Test @@ -71,7 +75,6 @@ public class ProducerServiceTest extends BaseServiceTest { sendResultFuture.complete(new SendResult(SendStatus.SEND_OK, "msgId", new MessageQueue(), 1L, "txId", "offsetMsgId", "regionId")); - ProducerService producerService = new ProducerService(this.connectorManager); producerService.setWriteQueueSelector((ctx, request) -> new SelectableMessageQueue(new MessageQueue("namespace%topic", "brokerName", 0), "brokerAddr")); @@ -88,8 +91,6 @@ public class ProducerServiceTest extends BaseServiceTest { @Test public void testSendMessageNoQueueSelect() { - ProducerService producerService = new ProducerService(this.connectorManager); - producerService.setWriteQueueSelector((ctx, request) -> null); CompletableFuture future = producerService.sendMessage(Context.current(), SendMessageRequest.newBuilder() @@ -125,7 +126,6 @@ public class ProducerServiceTest extends BaseServiceTest { .thenReturn(sendResultFuture); sendResultFuture.completeExceptionally(ex); - ProducerService producerService = new ProducerService(this.connectorManager); producerService.setWriteQueueSelector((ctx, request) -> new SelectableMessageQueue(new MessageQueue("namespace%topic", "brokerName", 0), "brokerAddr")); @@ -145,7 +145,6 @@ public class ProducerServiceTest extends BaseServiceTest { public void testSendMessageWithErrorThrow() { RuntimeException ex = new RuntimeException(); - ProducerService producerService = new ProducerService(this.connectorManager); producerService.setWriteQueueSelector((ctx, request) -> { throw ex; }); 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 e4e5e6fe4c..b013eb5965 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 @@ -30,7 +30,6 @@ import apache.rocketmq.v2.QueryRouteRequest; import apache.rocketmq.v2.QueryRouteResponse; import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.Settings; -import com.google.common.net.HostAndPort; import io.grpc.Context; import java.util.ArrayList; import java.util.HashMap; @@ -44,7 +43,6 @@ import org.apache.rocketmq.common.protocol.route.QueueData; import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; -import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode; import org.junit.Test; import static org.assertj.core.api.Assertions.assertThat; @@ -77,6 +75,8 @@ public class RouteServiceTest extends BaseServiceTest { .setAccessPoint(Endpoints.getDefaultInstance()) .build(); + private RouteService routeService; + @Override public void beforeEach() throws Exception { TopicRouteData routeData = new TopicRouteData(); @@ -106,6 +106,9 @@ public class RouteServiceTest extends BaseServiceTest { when(this.topicRouteCache.getMessageQueue("topic")).thenReturn(messageQueueWrapper); when(this.topicRouteCache.getMessageQueue("notExistTopic")).thenThrow(new MQClientException(ResponseCode.TOPIC_NOT_EXIST, "")); + + routeService = new RouteService(this.connectorManager, this.grpcClientManager); + routeService.start(); } @Test @@ -161,28 +164,8 @@ public class RouteServiceTest extends BaseServiceTest { return queueData; } - @Test - public void testLocalModeQueryRoute() throws Exception { - RouteService routeService = new RouteService(ProxyMode.LOCAL, this.connectorManager, this.grpcClientManager); - - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS); - - CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() - .setTopic(Resource.newBuilder() - .setName("topic") - .build()) - .build()); - QueryRouteResponse response = future.get(); - assertEquals(Code.OK.getNumber(), response.getStatus().getCode().getNumber()); - assertEquals(8, response.getMessageQueuesCount()); - assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getMessageQueues(0).getBroker() - .getEndpoints().getAddresses(0).getHost()); - } - @Test public void testQueryRouteWithInvalidEndpoints() throws Exception { - RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager); - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(INVALID_HOST_SETTINGS); CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() .setTopic(Resource.newBuilder() @@ -196,8 +179,6 @@ public class RouteServiceTest extends BaseServiceTest { @Test public void testQueryRoute() throws Exception { - RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager); - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS); CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() @@ -215,8 +196,6 @@ public class RouteServiceTest extends BaseServiceTest { @Test public void testQueryRouteWhenTopicNotExist() throws Exception { - RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager); - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS); CompletableFuture future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder() @@ -231,8 +210,6 @@ public class RouteServiceTest extends BaseServiceTest { @Test public void testQueryAssignmentInvalidEndpoints() throws Exception { - RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager); - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(INVALID_HOST_SETTINGS); CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() .setTopic( @@ -246,32 +223,8 @@ public class RouteServiceTest extends BaseServiceTest { assertEquals(Code.ILLEGAL_ACCESS_POINT.getNumber(), response.getStatus().getCode().getNumber()); } - @Test - public void testLocalModeQueryAssignment() throws Exception { - RouteService routeService = new RouteService(ProxyMode.LOCAL, this.connectorManager, this.grpcClientManager); - - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS); - - CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder() - .setTopic(Resource.newBuilder() - .setName("topic") - .build()) - .setGroup(Resource.newBuilder() - .setName("group") - .build()) - .build()); - - QueryAssignmentResponse response = future.get(); - assertEquals(Code.OK.getNumber(), response.getStatus().getCode().getNumber()); - assertEquals(1, response.getAssignmentsCount()); - assertEquals("brokerName", response.getAssignments(0).getMessageQueue().getBroker().getName()); - assertEquals(HostAndPort.fromString(brokerAddress).getHost(), response.getAssignments(0).getMessageQueue().getBroker().getEndpoints().getAddresses(0).getHost()); - } - @Test public void testQueryAssignment() throws Exception { - RouteService routeService = new RouteService(ProxyMode.CLUSTER, this.connectorManager, this.grpcClientManager); - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(WITH_HOST_SETTINGS); CompletableFuture future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()