From 67d749b569d054e7aabb952ca015e00c5f82c662 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Thu, 5 May 2022 17:44:04 +0800 Subject: [PATCH] [ISSUE #3949] v2 support --- .../acl/plain/PlainAccessValidator.java | 9 -- ...ractTransactionalMessageCheckListener.java | 1 + .../CheckTransactionStateRequestHeader.java | 12 +- .../ProxyClientRemotingProcessor.java | 1 + .../TransactionStateCheckRequest.java | 11 ++ .../proxy/grpc/v2/GrpcMessagingProcessor.java | 15 --- .../proxy/grpc/v2/adapter/GrpcConverter.java | 65 ++++------- .../proxy/grpc/v2/adapter/RequestMapping.java | 2 - .../v2/adapter/channel/GrpcClientChannel.java | 1 + ...seReceiveMessageResponseStreamWriter.java} | 6 +- .../BaseReceiveMessageResultFilter.java | 67 ++++++++++++ .../grpc/v2/service/ClusterGrpcService.java | 7 -- .../grpc/v2/service/GrpcForwardService.java | 4 - .../grpc/v2/service/LocalGrpcService.java | 51 +-------- .../v2/service/cluster/ConsumerService.java | 92 +--------------- ...ultReceiveMessageResponseStreamWriter.java | 4 +- .../DefaultReceiveMessageResultFilter.java | 103 ++++++------------ .../service/cluster/TransactionService.java | 5 +- ...calReceiveMessageResponseStreamWriter.java | 4 +- .../LocalReceiveMessageResultFilter.java | 45 ++------ .../grpc/v2/service/LocalGrpcServiceTest.java | 62 ----------- .../service/cluster/ConsumerServiceTest.java | 63 +---------- .../cluster/TransactionServiceTest.java | 7 +- .../test/grpc/v2/ClusterGrpcTest.java | 10 -- .../rocketmq/test/grpc/v2/GrpcBaseTest.java | 100 ----------------- .../rocketmq/test/grpc/v2/LocalGrpcTest.java | 10 -- 26 files changed, 181 insertions(+), 576 deletions(-) rename proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/{ReceiveMessageResponseStreamWriter.java => BaseReceiveMessageResponseStreamWriter.java} (97%) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/BaseReceiveMessageResultFilter.java diff --git a/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainAccessValidator.java b/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainAccessValidator.java index da49e56106..6243e3fdf4 100644 --- a/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainAccessValidator.java +++ b/acl/src/main/java/org/apache/rocketmq/acl/plain/PlainAccessValidator.java @@ -21,7 +21,6 @@ import apache.rocketmq.v2.EndTransactionRequest; import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest; import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.Message; -import apache.rocketmq.v2.NackMessageRequest; import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.SendMessageRequest; @@ -210,14 +209,6 @@ public class PlainAccessValidator implements AccessValidator { Resource topic = request.getTopic(); String topicName = NamespaceUtil.wrapNamespace(topic.getResourceNamespace(), topic.getName()); accessResource.addResourceAndPerm(topicName, Permission.SUB); - } else if (NackMessageRequest.getDescriptor().getFullName().equals(rpcFullName)) { - NackMessageRequest request = (NackMessageRequest) messageV3; - Resource group = request.getGroup(); - String groupName = NamespaceUtil.wrapNamespace(group.getResourceNamespace(), group.getName()); - accessResource.addResourceAndPerm(groupName, Permission.SUB); - Resource topic = request.getTopic(); - String topicName = NamespaceUtil.wrapNamespace(topic.getResourceNamespace(), topic.getName()); - accessResource.addResourceAndPerm(topicName, Permission.SUB); } else if (ForwardMessageToDeadLetterQueueRequest.getDescriptor().getFullName().equals(rpcFullName)) { ForwardMessageToDeadLetterQueueRequest request = (ForwardMessageToDeadLetterQueueRequest) messageV3; Resource group = request.getGroup(); diff --git a/broker/src/main/java/org/apache/rocketmq/broker/transaction/AbstractTransactionalMessageCheckListener.java b/broker/src/main/java/org/apache/rocketmq/broker/transaction/AbstractTransactionalMessageCheckListener.java index 2ed0d9d1cd..613fe0f589 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/transaction/AbstractTransactionalMessageCheckListener.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/transaction/AbstractTransactionalMessageCheckListener.java @@ -56,6 +56,7 @@ public abstract class AbstractTransactionalMessageCheckListener { checkTransactionStateRequestHeader.setMsgId(msgExt.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX)); checkTransactionStateRequestHeader.setTransactionId(checkTransactionStateRequestHeader.getMsgId()); checkTransactionStateRequestHeader.setTranStateTableOffset(msgExt.getQueueOffset()); + checkTransactionStateRequestHeader.setBrokerName(brokerController.getBrokerConfig().getBrokerName()); msgExt.setTopic(msgExt.getUserProperty(MessageConst.PROPERTY_REAL_TOPIC)); msgExt.setQueueId(Integer.parseInt(msgExt.getUserProperty(MessageConst.PROPERTY_REAL_QUEUE_ID))); msgExt.setStoreSize(0); diff --git a/common/src/main/java/org/apache/rocketmq/common/protocol/header/CheckTransactionStateRequestHeader.java b/common/src/main/java/org/apache/rocketmq/common/protocol/header/CheckTransactionStateRequestHeader.java index 149de9b579..6671a9d773 100644 --- a/common/src/main/java/org/apache/rocketmq/common/protocol/header/CheckTransactionStateRequestHeader.java +++ b/common/src/main/java/org/apache/rocketmq/common/protocol/header/CheckTransactionStateRequestHeader.java @@ -25,6 +25,7 @@ import org.apache.rocketmq.remoting.annotation.CFNotNull; import org.apache.rocketmq.remoting.exception.RemotingCommandException; public class CheckTransactionStateRequestHeader implements CommandCustomHeader { + private String brokerName; @CFNotNull private Long tranStateTableOffset; @CFNotNull @@ -37,6 +38,14 @@ public class CheckTransactionStateRequestHeader implements CommandCustomHeader { public void checkFields() throws RemotingCommandException { } + public String getBrokerName() { + return brokerName; + } + + public void setBrokerName(String brokerName) { + this.brokerName = brokerName; + } + public Long getTranStateTableOffset() { return tranStateTableOffset; } @@ -80,7 +89,8 @@ public class CheckTransactionStateRequestHeader implements CommandCustomHeader { @Override public String toString() { return "CheckTransactionStateRequestHeader{" + - "tranStateTableOffset=" + tranStateTableOffset + + "brokerName='" + brokerName + '\'' + + ", tranStateTableOffset=" + tranStateTableOffset + ", commitLogOffset=" + commitLogOffset + ", msgId='" + msgId + '\'' + ", transactionId='" + transactionId + '\'' + diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java index 41001e8f6e..3c03488916 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java @@ -68,6 +68,7 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor { requestHeader.getTransactionId(), requestHeader.getCommitLogOffset(), requestHeader.getTranStateTableOffset()), + requestHeader.getBrokerName(), messageExt ) ); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionStateCheckRequest.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionStateCheckRequest.java index ad6b5ac0d6..c5269e64f9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionStateCheckRequest.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionStateCheckRequest.java @@ -24,6 +24,7 @@ public class TransactionStateCheckRequest { private Long commitLogOffset; private String msgId; private TransactionId transactionId; + private String brokerName; private MessageExt messageExt; public TransactionStateCheckRequest( @@ -32,6 +33,7 @@ public class TransactionStateCheckRequest { Long commitLogOffset, String msgId, TransactionId transactionId, + String brokerName, MessageExt messageExt ) { this.groupId = groupId; @@ -39,6 +41,7 @@ public class TransactionStateCheckRequest { this.commitLogOffset = commitLogOffset; this.msgId = msgId; this.transactionId = transactionId; + this.brokerName = brokerName; this.messageExt = messageExt; } @@ -82,6 +85,14 @@ public class TransactionStateCheckRequest { this.transactionId = transactionId; } + public String getBrokerName() { + return brokerName; + } + + public void setBrokerName(String brokerName) { + this.brokerName = brokerName; + } + public MessageExt getMessageExt() { return messageExt; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java index ad06e23074..89666c1a80 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java @@ -28,8 +28,6 @@ import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.HeartbeatResponse; import apache.rocketmq.v2.MessagingServiceGrpc; -import apache.rocketmq.v2.NackMessageRequest; -import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.NotifyClientTerminationRequest; import apache.rocketmq.v2.NotifyClientTerminationResponse; import apache.rocketmq.v2.QueryAssignmentRequest; @@ -119,19 +117,6 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic grpcForwardService.receiveMessage(Context.current(), request, responseObserver); } - @Override - public void nackMessage(NackMessageRequest request, StreamObserver responseObserver) { - CompletableFuture future = grpcForwardService.nackMessage(Context.current(), request); - future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) - .exceptionally(e -> { - ResponseWriter.write( - responseObserver, - NackMessageResponse.newBuilder().setStatus(convertExceptionToStatus(e)).build() - ); - return null; - }); - } - @Override public void ackMessage(AckMessageRequest request, StreamObserver responseObserver) { CompletableFuture future = grpcForwardService.ackMessage(Context.current(), request); 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 883abf6291..52fb47d950 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 @@ -26,7 +26,6 @@ import apache.rocketmq.v2.Digest; import apache.rocketmq.v2.DigestType; import apache.rocketmq.v2.Encoding; import apache.rocketmq.v2.EndTransactionRequest; -import apache.rocketmq.v2.ExponentialBackoff; import apache.rocketmq.v2.FilterExpression; import apache.rocketmq.v2.FilterType; import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest; @@ -34,12 +33,10 @@ import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.Message; import apache.rocketmq.v2.MessageQueue; import apache.rocketmq.v2.MessageType; -import apache.rocketmq.v2.NackMessageRequest; import apache.rocketmq.v2.NotifyClientTerminationRequest; import apache.rocketmq.v2.Permission; import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.Resource; -import apache.rocketmq.v2.RetryPolicy; import apache.rocketmq.v2.SendMessageRequest; import apache.rocketmq.v2.Settings; import apache.rocketmq.v2.SubscriptionEntry; @@ -63,6 +60,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.TimeUnit; +import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.constant.ConsumeInitMode; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.constant.PermName; @@ -245,11 +243,6 @@ public class GrpcConverter { return buildAckMessageRequestHeader(request.getTopic(), request.getGroup(), handle); } - public static AckMessageRequestHeader buildAckMessageRequestHeader(NackMessageRequest request) { - ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); - return buildAckMessageRequestHeader(request.getTopic(), request.getGroup(), handle); - } - public static AckMessageRequestHeader buildAckMessageRequestHeader(Resource topic, Resource group, ReceiptHandle handle) { String groupName = GrpcConverter.wrapResourceWithNamespace(group); String topicName = GrpcConverter.wrapResourceWithNamespace(topic); @@ -263,37 +256,6 @@ public class GrpcConverter { return ackMessageRequestHeader; } - public static ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(NackMessageRequest request, - RetryPolicy retryPolicy) { - String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); - String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); - ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); - - ChangeInvisibleTimeRequestHeader changeInvisibleTimeRequestHeader = new ChangeInvisibleTimeRequestHeader(); - changeInvisibleTimeRequestHeader.setConsumerGroup(groupName); - changeInvisibleTimeRequestHeader.setTopic(handle.getRealTopic(topicName, groupName)); - changeInvisibleTimeRequestHeader.setQueueId(handle.getQueueId()); - changeInvisibleTimeRequestHeader.setExtraInfo(handle.getReceiptHandle()); - changeInvisibleTimeRequestHeader.setOffset(handle.getOffset()); - changeInvisibleTimeRequestHeader.setInvisibleTime( - Durations.toMillis(calculateNextDeliveryDurations(retryPolicy, request.getDeliveryAttempt()))); - return changeInvisibleTimeRequestHeader; - } - - public static Duration calculateNextDeliveryDurations(RetryPolicy retryPolicy, int deliveryAttempt) { - if (retryPolicy.hasCustomizedBackoff()) { - int nextCount = retryPolicy.getCustomizedBackoff().getNextCount(); - return retryPolicy.getCustomizedBackoff().getNext(Math.min(nextCount, deliveryAttempt)); - } - ExponentialBackoff exponentialBackoff = retryPolicy.getExponentialBackoff(); - long nextDurationMillis = (long) (Math.pow(exponentialBackoff.getMultiplier(), deliveryAttempt) * - Durations.toMillis(exponentialBackoff.getInitial())); - nextDurationMillis = Math.min( - Durations.toMillis(exponentialBackoff.getMax()), - nextDurationMillis); - return Durations.fromMillis(nextDurationMillis); - } - public static ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(ChangeInvisibleDurationRequest request) { String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); @@ -328,13 +290,6 @@ public class GrpcConverter { return buildConsumerSendMsgBackRequestHeader(request.getMessageQueue().getTopic(), request.getGroup(), handle, messageId, maxReconsumeTimes); } - public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader( - NackMessageRequest request, int maxReconsumeTimes) { - ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); - return buildConsumerSendMsgBackRequestHeader(request.getTopic(), request.getGroup(), handle, - request.getMessageId(), maxReconsumeTimes); - } - public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackRequestHeader( ForwardMessageToDeadLetterQueueRequest request) { ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); @@ -494,6 +449,24 @@ public class GrpcConverter { return message; } + public static MessageQueue buildMessageQueue(MessageExt messageExt, String brokerName) { + Broker broker = Broker.getDefaultInstance(); + if (!StringUtils.isEmpty(brokerName)) { + broker = Broker.newBuilder() + .setName(brokerName) + .setId(0) + .build(); + } + return MessageQueue.newBuilder() + .setId(messageExt.getQueueId()) + .setTopic(Resource.newBuilder() + .setName(NamespaceUtil.withoutNamespace(messageExt.getTopic())) + .setResourceNamespace(NamespaceUtil.getNamespaceFromResource(messageExt.getTopic())) + .build()) + .setBroker(broker) + .build(); + } + public static String buildExpressionType(FilterType filterType) { switch (filterType) { case SQL: diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/RequestMapping.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/RequestMapping.java index eafca24f49..7f5ee2a13e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/RequestMapping.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/RequestMapping.java @@ -22,7 +22,6 @@ import apache.rocketmq.v2.ChangeInvisibleDurationRequest; import apache.rocketmq.v2.EndTransactionRequest; import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; import apache.rocketmq.v2.HeartbeatRequest; -import apache.rocketmq.v2.NackMessageRequest; import apache.rocketmq.v2.NotifyClientTerminationRequest; import apache.rocketmq.v2.QueryAssignmentRequest; import apache.rocketmq.v2.QueryRouteRequest; @@ -42,7 +41,6 @@ public class RequestMapping { put(QueryAssignmentRequest.getDescriptor().getFullName(), RequestCode.GET_ROUTEINFO_BY_TOPIC); put(ReceiveMessageRequest.getDescriptor().getFullName(), RequestCode.PULL_MESSAGE); put(AckMessageRequest.getDescriptor().getFullName(), RequestCode.UPDATE_CONSUMER_OFFSET); - put(NackMessageRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); put(ForwardMessageToDeadLetterQueueResponse.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); put(EndTransactionRequest.getDescriptor().getFullName(), RequestCode.END_TRANSACTION); put(NotifyClientTerminationRequest.getDescriptor().getFullName(), RequestCode.UNREGISTER_CLIENT); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java index cd82c34d1b..a14dd4e3f8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java @@ -130,6 +130,7 @@ public class GrpcClientChannel extends SimpleChannel { .setRecoverOrphanedTransactionCommand(RecoverOrphanedTransactionCommand.newBuilder() .setTransactionId(transactionId.getProxyTransactionId()) .setOrphanedTransactionalMessage(GrpcConverter.buildMessage(messageExt)) + .setMessageQueue(GrpcConverter.buildMessageQueue(messageExt, header.getBrokerName())) .build()) .build()); break; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ReceiveMessageResponseStreamWriter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/BaseReceiveMessageResponseStreamWriter.java similarity index 97% rename from proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ReceiveMessageResponseStreamWriter.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/BaseReceiveMessageResponseStreamWriter.java index 32998280f2..1f1cadb746 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ReceiveMessageResponseStreamWriter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/BaseReceiveMessageResponseStreamWriter.java @@ -30,19 +30,19 @@ 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.adapter.ResponseWriter; -public abstract class ReceiveMessageResponseStreamWriter { +public abstract class BaseReceiveMessageResponseStreamWriter { protected final StreamObserver streamObserver; protected final ResponseHook receiveMessageHook; protected final ReceiveMessageResultFilter receiveMessageResultFilter; public interface Builder { - ReceiveMessageResponseStreamWriter build( + BaseReceiveMessageResponseStreamWriter build( StreamObserver observer, ResponseHook hook); } - public ReceiveMessageResponseStreamWriter( + public BaseReceiveMessageResponseStreamWriter( StreamObserver observer, ResponseHook hook, ReceiveMessageResultFilter messageResultFilter) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/BaseReceiveMessageResultFilter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/BaseReceiveMessageResultFilter.java new file mode 100644 index 0000000000..350335f1fd --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/BaseReceiveMessageResultFilter.java @@ -0,0 +1,67 @@ +/* + * 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; + +import apache.rocketmq.v2.Message; +import apache.rocketmq.v2.ReceiveMessageRequest; +import apache.rocketmq.v2.Settings; +import io.grpc.Context; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; +import org.apache.rocketmq.proxy.common.utils.FilterUtils; +import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; + +public abstract class BaseReceiveMessageResultFilter implements ReceiveMessageResultFilter { + + protected final GrpcClientManager grpcClientManager; + + public BaseReceiveMessageResultFilter(GrpcClientManager manager) { + grpcClientManager = manager; + } + + @Override + public List filterMessage(Context ctx, ReceiveMessageRequest request, List messageExtList) { + if (messageExtList == null || messageExtList.isEmpty()) { + return Collections.emptyList(); + } + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getMessageQueue().getTopic()); + SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(topicName, request.getFilterExpression()); + Settings settings = grpcClientManager.getClientSettings(ctx); + int maxAttempts = settings.getBackoffPolicy().getMaxAttempts(); + List resMessageList = new ArrayList<>(); + for (MessageExt messageExt : messageExtList) { + if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) { + processNoMatchMessage(ctx, request, messageExt); + continue; + } + if (messageExt.getReconsumeTimes() >= maxAttempts) { + processExceedMaxAttemptsMessage(ctx, request, messageExt, maxAttempts); + continue; + } + resMessageList.add(GrpcConverter.buildMessage(messageExt)); + } + return resMessageList; + } + + protected abstract void processNoMatchMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt); + + protected abstract void processExceedMaxAttemptsMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt, int maxAttempts); +} 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 6ebc350747..f47dcc9071 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 @@ -27,8 +27,6 @@ import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest; import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.HeartbeatResponse; -import apache.rocketmq.v2.NackMessageRequest; -import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.NotifyClientTerminationRequest; import apache.rocketmq.v2.NotifyClientTerminationResponse; import apache.rocketmq.v2.QueryAssignmentRequest; @@ -130,11 +128,6 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc consumerService.receiveMessage(ctx, request, responseObserver); } - @Override - public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { - return consumerService.nackMessage(ctx, request); - } - @Override public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { return consumerService.ackMessage(ctx, request); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java index 48d2f2504f..3f0b459612 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java @@ -27,8 +27,6 @@ import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest; import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.HeartbeatResponse; -import apache.rocketmq.v2.NackMessageRequest; -import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.NotifyClientTerminationRequest; import apache.rocketmq.v2.NotifyClientTerminationResponse; import apache.rocketmq.v2.QueryAssignmentRequest; @@ -57,8 +55,6 @@ public interface GrpcForwardService extends StartAndShutdown { void receiveMessage(Context ctx, ReceiveMessageRequest request, StreamObserver responseObserver); - CompletableFuture nackMessage(Context ctx, NackMessageRequest request); - CompletableFuture ackMessage(Context ctx, AckMessageRequest request); CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request); 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 6cb84a5c4b..414c480343 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 @@ -30,8 +30,6 @@ import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest; import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.HeartbeatResponse; -import apache.rocketmq.v2.NackMessageRequest; -import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.NotifyClientTerminationRequest; import apache.rocketmq.v2.NotifyClientTerminationResponse; import apache.rocketmq.v2.QueryAssignmentRequest; @@ -41,7 +39,6 @@ import apache.rocketmq.v2.QueryRouteResponse; import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.ReceiveMessageResponse; import apache.rocketmq.v2.Resource; -import apache.rocketmq.v2.RetryPolicy; import apache.rocketmq.v2.SendMessageRequest; import apache.rocketmq.v2.SendMessageResponse; import apache.rocketmq.v2.Settings; @@ -120,7 +117,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo private final RouteService routeService; private final ClientSettingsService clientSettingsService; private final LocalWriteQueueSelector localWriteQueueSelector; - private final ReceiveMessageResponseStreamWriter.Builder streamWriterBuilder; + private final BaseReceiveMessageResponseStreamWriter.Builder streamWriterBuilder; private volatile ResponseHook receiveMessageHook; @@ -271,7 +268,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo long pollTime = ctx.getDeadline().timeRemaining(TimeUnit.MILLISECONDS); // TODO: get fifo config from subscriptionGroupManager boolean fifo = false; - ReceiveMessageResponseStreamWriter writer = streamWriterBuilder.build(responseObserver, receiveMessageHook); + BaseReceiveMessageResponseStreamWriter writer = streamWriterBuilder.build(responseObserver, receiveMessageHook); ReceiveMessageResponseHandler handler = new ReceiveMessageResponseHandler(brokerController.getBrokerConfig().getBrokerName(), fifo); ReceiveMessageChannel channel = channelManager.createChannel(ctx, context -> new ReceiveMessageChannel(context, handler), ReceiveMessageChannel.class); CompletableFuture> future = new CompletableFuture<>(); @@ -345,50 +342,6 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo return future; } - @Override - public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { - Channel channel = channelManager.createChannel(ctx); - SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); - CompletableFuture future = new CompletableFuture<>(); - - RetryPolicy retryPolicy = grpcClientManager.getClientSettings(ctx).getBackoffPolicy(); - int maxReconsumeTimes = retryPolicy.getMaxAttempts(); - if (request.getDeliveryAttempt() >= maxReconsumeTimes) { - ConsumerSendMsgBackRequestHeader requestHeader = GrpcConverter.buildConsumerSendMsgBackToDLQRequestHeader(request, maxReconsumeTimes); - RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader); - command.makeCustomHeaderToNet(); - - try { - RemotingCommand responseCommand = brokerController.getSendMessageProcessor() - .processRequest(channelHandlerContext, command); - NackMessageResponse response = NackMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(responseCommand.getCode(), responseCommand.getRemark())) - .build(); - future.complete(response); - } catch (Exception e) { - log.error("Exception raised while nackMessage", e); - future.completeExceptionally(e); - } - } else { - ChangeInvisibleTimeRequestHeader requestHeader = GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, retryPolicy); - RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, requestHeader); - command.makeCustomHeaderToNet(); - - try { - RemotingCommand responseCommand = brokerController.getChangeInvisibleTimeProcessor() - .processRequest(channelHandlerContext, command); - NackMessageResponse response = NackMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(responseCommand.getCode(), responseCommand.getRemark())) - .build(); - future.complete(response); - } catch (Exception e) { - log.error("Exception raised while nackMessage", e); - future.completeExceptionally(e); - } - } - return future; - } - @Override public CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request) { 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 8df9fd14f7..474ed29f15 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 @@ -23,12 +23,8 @@ import apache.rocketmq.v2.AckMessageResultEntry; import apache.rocketmq.v2.ChangeInvisibleDurationRequest; import apache.rocketmq.v2.ChangeInvisibleDurationResponse; import apache.rocketmq.v2.Code; -import apache.rocketmq.v2.NackMessageRequest; -import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.ReceiveMessageResponse; -import apache.rocketmq.v2.RetryPolicy; -import apache.rocketmq.v2.Settings; import io.grpc.Context; import io.grpc.stub.StreamObserver; import java.util.ArrayList; @@ -39,7 +35,6 @@ import org.apache.rocketmq.client.consumer.AckStatus; import org.apache.rocketmq.common.consumer.ReceiptHandle; 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.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.ForwardProducer; @@ -52,8 +47,7 @@ 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.BaseService; import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; -import org.apache.rocketmq.proxy.grpc.v2.service.ReceiveMessageResponseStreamWriter; -import org.apache.rocketmq.remoting.protocol.RemotingCommand; +import org.apache.rocketmq.proxy.grpc.v2.service.BaseReceiveMessageResponseStreamWriter; public class ConsumerService extends BaseService { protected final ForwardReadConsumer readConsumer; @@ -65,11 +59,10 @@ public class ConsumerService extends BaseService { protected final GrpcClientManager grpcClientManager; private volatile ReadQueueSelector readQueueSelector; - private volatile ReceiveMessageResponseStreamWriter.Builder receiveMessageWriterBuilder; + private volatile BaseReceiveMessageResponseStreamWriter.Builder receiveMessageWriterBuilder; private volatile ResponseHook receiveMessageHook; private volatile ResponseHook ackMessageHook; - private volatile ResponseHook nackMessageHook; private volatile ResponseHook changeInvisibleDurationHook; public ConsumerService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) { @@ -92,7 +85,7 @@ public class ConsumerService extends BaseService { public void receiveMessage(Context ctx, ReceiveMessageRequest request, StreamObserver responseObserver) { - ReceiveMessageResponseStreamWriter writer = receiveMessageWriterBuilder.build(responseObserver, receiveMessageHook); + BaseReceiveMessageResponseStreamWriter writer = receiveMessageWriterBuilder.build(responseObserver, receiveMessageHook); try { PopMessageRequestHeader requestHeader = this.buildPopMessageRequestHeader(ctx, request); SelectableMessageQueue messageQueue = this.readQueueSelector.select(ctx, request, requestHeader); @@ -203,72 +196,6 @@ public class ConsumerService extends BaseService { .build(); } - public CompletableFuture nackMessage(Context ctx, NackMessageRequest request) { - CompletableFuture future = new CompletableFuture<>(); - try { - ReceiptHandle receiptHandle = resolveReceiptHandle(ctx, request.getReceiptHandle()); - String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); - - Settings settings = grpcClientManager.getClientSettings(ctx); - int maxDeliveryAttempts = settings.getBackoffPolicy().getMaxAttempts(); - if (request.getDeliveryAttempt() >= maxDeliveryAttempts) { - future = this.producer.sendMessageBackThenAckOrg( - ctx, - brokerAddr, - 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(), request.getMessageId(), requestHeader) - .thenApply(result -> convertToNackMessageResponse(ctx, request, result)); - } - } catch (Throwable t) { - future.completeExceptionally(t); - } - future.whenComplete((response, throwable) -> { - if (nackMessageHook != null) { - nackMessageHook.beforeResponse(ctx, request, response, throwable); - } - }); - return future; - } - - protected ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(Context ctx, - NackMessageRequest request) { - RetryPolicy retryPolicy = grpcClientManager.getClientSettings(ctx).getBackoffPolicy(); - return GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, retryPolicy); - } - - protected AckMessageRequestHeader buildAckMessageRequestHeader(Context ctx, NackMessageRequest request) { - return GrpcConverter.buildAckMessageRequestHeader(request); - } - - protected ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader(Context ctx, - NackMessageRequest request, - int maxReconsumeTimes) { - return GrpcConverter.buildConsumerSendMsgBackToDLQRequestHeader(request, maxReconsumeTimes); - } - - protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request, - AckResult ackResult) { - if (AckStatus.OK.equals(ackResult.getStatus())) { - return NackMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.OK, Code.OK.name())) - .build(); - } - return NackMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(Code.INTERNAL_SERVER_ERROR, "nack failed: status is abnormal")) - .build(); - } - - protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request, - RemotingCommand sendMsgBackToDLQResult) { - return NackMessageResponse.newBuilder() - .setStatus(ResponseBuilder.buildStatus(sendMsgBackToDLQResult.getCode(), sendMsgBackToDLQResult.getRemark())) - .build(); - } - public CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request) { CompletableFuture future = new CompletableFuture<>(); @@ -318,12 +245,12 @@ public class ConsumerService extends BaseService { this.readQueueSelector = readQueueSelector; } - public ReceiveMessageResponseStreamWriter.Builder getReceiveMessageWriterBuilder() { + public BaseReceiveMessageResponseStreamWriter.Builder getReceiveMessageWriterBuilder() { return receiveMessageWriterBuilder; } public void setReceiveMessageWriterBuilder( - ReceiveMessageResponseStreamWriter.Builder receiveMessageWriterBuilder) { + BaseReceiveMessageResponseStreamWriter.Builder receiveMessageWriterBuilder) { this.receiveMessageWriterBuilder = receiveMessageWriterBuilder; } @@ -345,15 +272,6 @@ public class ConsumerService extends BaseService { this.ackMessageHook = ackMessageHook; } - public ResponseHook getNackMessageHook() { - return nackMessageHook; - } - - public void setNackMessageHook( - ResponseHook nackMessageHook) { - this.nackMessageHook = nackMessageHook; - } - public ResponseHook getChangeInvisibleDurationHook() { return changeInvisibleDurationHook; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/DefaultReceiveMessageResponseStreamWriter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/DefaultReceiveMessageResponseStreamWriter.java index bf071719f3..10dfda6fcc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/DefaultReceiveMessageResponseStreamWriter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/DefaultReceiveMessageResponseStreamWriter.java @@ -33,10 +33,10 @@ 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.BaseService; -import org.apache.rocketmq.proxy.grpc.v2.service.ReceiveMessageResponseStreamWriter; +import org.apache.rocketmq.proxy.grpc.v2.service.BaseReceiveMessageResponseStreamWriter; import org.apache.rocketmq.proxy.grpc.v2.service.ReceiveMessageResultFilter; -public class DefaultReceiveMessageResponseStreamWriter extends ReceiveMessageResponseStreamWriter { +public class DefaultReceiveMessageResponseStreamWriter extends BaseReceiveMessageResponseStreamWriter { protected static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); protected static final long NACK_INVISIBLE_TIME = Duration.ofSeconds(1).toMillis(); 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 index b15951eb74..adf662e87b 100644 --- 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 @@ -17,38 +17,29 @@ package org.apache.rocketmq.proxy.grpc.v2.service.cluster; -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.BaseReceiveMessageResultFilter; import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; -import org.apache.rocketmq.proxy.grpc.v2.service.ReceiveMessageResultFilter; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import static org.apache.rocketmq.proxy.grpc.v2.service.BaseService.getBrokerAddr; -public class DefaultReceiveMessageResultFilter implements ReceiveMessageResultFilter { +public class DefaultReceiveMessageResultFilter extends BaseReceiveMessageResultFilter { protected final ForwardProducer producer; protected final ForwardWriteConsumer writeConsumer; - protected final GrpcClientManager grpcClientManager; protected final TopicRouteCache topicRouteCache; private volatile ResponseHook ackNoMatchedMessageHook; @@ -56,69 +47,14 @@ public class DefaultReceiveMessageResultFilter implements ReceiveMessageResultFi public DefaultReceiveMessageResultFilter(ForwardProducer producer, ForwardWriteConsumer writeConsumer, GrpcClientManager grpcClientManager, TopicRouteCache topicRouteCache) { + super(grpcClientManager); 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); - int maxAttempts = settings.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 (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()); - ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader( - request, - handle, - messageExt.getMsgId(), - maxReconsumeTimes); - AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader(request, 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) { + protected void processNoMatchMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) { CompletableFuture future = new CompletableFuture<>(); ReceiptHandle handle = ReceiptHandle.create(messageExt); @@ -140,6 +76,37 @@ public class DefaultReceiveMessageResultFilter implements ReceiveMessageResultFi }); } + @Override + protected void processExceedMaxAttemptsMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt, + int maxAttempts) { + 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()); + ConsumerSendMsgBackRequestHeader sendMsgBackRequestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader( + request, + handle, + messageExt.getMsgId(), + maxAttempts); + AckMessageRequestHeader ackMessageRequestHeader = GrpcConverter.buildAckMessageRequestHeader(request, 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); + } + }); + } + public ResponseHook getAckNoMatchedMessageHook() { return ackNoMatchedMessageHook; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java index d142e2de98..9600d7acc2 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionService.java @@ -27,6 +27,7 @@ import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ThreadLocalRandom; import org.apache.commons.collections.CollectionUtils; +import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader; import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; @@ -68,12 +69,14 @@ public class TransactionService extends BaseService implements TransactionStateC GrpcClientChannel channel = GrpcClientChannel.getChannel(this.channelManager, checkData.getGroupId(), clientId); String transactionId = checkData.getTransactionId().getProxyTransactionId(); - Message message = GrpcConverter.buildMessage(checkData.getMessageExt()); + MessageExt messageExt = checkData.getMessageExt(); + Message message = GrpcConverter.buildMessage(messageExt); TelemetryCommand response = TelemetryCommand.newBuilder() .setRecoverOrphanedTransactionCommand( RecoverOrphanedTransactionCommand.newBuilder() .setOrphanedTransactionalMessage(message) .setTransactionId(transactionId) + .setMessageQueue(GrpcConverter.buildMessageQueue(messageExt, checkData.getBrokerName())) .build() ).build(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/LocalReceiveMessageResponseStreamWriter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/LocalReceiveMessageResponseStreamWriter.java index 50fe0e1a59..0ec9913e22 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/LocalReceiveMessageResponseStreamWriter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/LocalReceiveMessageResponseStreamWriter.java @@ -32,14 +32,14 @@ import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.channel.SimpleChannelHandlerContext; 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.ReceiveMessageResponseStreamWriter; +import org.apache.rocketmq.proxy.grpc.v2.service.BaseReceiveMessageResponseStreamWriter; import org.apache.rocketmq.proxy.grpc.v2.service.ReceiveMessageResultFilter; import org.apache.rocketmq.remoting.exception.RemotingCommandException; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -public class LocalReceiveMessageResponseStreamWriter extends ReceiveMessageResponseStreamWriter { +public class LocalReceiveMessageResponseStreamWriter extends BaseReceiveMessageResponseStreamWriter { private final static Logger log = LoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); private final ChannelManager channelManager; private final BrokerController brokerController; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/LocalReceiveMessageResultFilter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/LocalReceiveMessageResultFilter.java index 62ea8fb78c..c69aab2567 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/LocalReceiveMessageResultFilter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/local/LocalReceiveMessageResultFilter.java @@ -17,14 +17,9 @@ package org.apache.rocketmq.proxy.grpc.v2.service.local; -import apache.rocketmq.v2.Message; import apache.rocketmq.v2.ReceiveMessageRequest; -import apache.rocketmq.v2.Settings; import io.grpc.Context; import io.netty.channel.Channel; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.consumer.ReceiptHandle; @@ -33,56 +28,30 @@ import org.apache.rocketmq.common.protocol.RequestCode; import org.apache.rocketmq.common.protocol.ResponseCode; 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.channel.ChannelManager; import org.apache.rocketmq.proxy.channel.SimpleChannelHandlerContext; -import org.apache.rocketmq.proxy.common.utils.FilterUtils; import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; +import org.apache.rocketmq.proxy.grpc.v2.service.BaseReceiveMessageResultFilter; import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; -import org.apache.rocketmq.proxy.grpc.v2.service.ReceiveMessageResultFilter; import org.apache.rocketmq.remoting.exception.RemotingCommandException; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -public class LocalReceiveMessageResultFilter implements ReceiveMessageResultFilter { +public class LocalReceiveMessageResultFilter extends BaseReceiveMessageResultFilter { private final static Logger log = LoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); private final ChannelManager channelManager; private final BrokerController brokerController; - private final GrpcClientManager grpcClientManager; public LocalReceiveMessageResultFilter(ChannelManager channelManager, BrokerController brokerController, GrpcClientManager grpcClientManager) { + super(grpcClientManager); this.channelManager = channelManager; this.brokerController = brokerController; - this.grpcClientManager = grpcClientManager; } @Override - public List filterMessage(Context ctx, ReceiveMessageRequest request, List messageExtList) { - if (messageExtList == null || messageExtList.isEmpty()) { - return Collections.emptyList(); - } - String topicName = GrpcConverter.wrapResourceWithNamespace(request.getMessageQueue().getTopic()); - SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(topicName, request.getFilterExpression()); - Settings settings = grpcClientManager.getClientSettings(ctx); - int maxAttempts = settings.getBackoffPolicy().getMaxAttempts(); - List resMessageList = new ArrayList<>(); - for (MessageExt messageExt : messageExtList) { - if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) { - ackMessage(ctx, request, messageExt); - continue; - } - if (messageExt.getReconsumeTimes() >= maxAttempts) { - forwardMessageToDLQ(ctx, request, messageExt, maxAttempts); - continue; - } - resMessageList.add(GrpcConverter.buildMessage(messageExt)); - } - return resMessageList; - } - - private void ackMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) { + protected void processNoMatchMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) { ReceiptHandle handle = ReceiptHandle.create(messageExt); if (handle == null) { return; @@ -98,7 +67,9 @@ public class LocalReceiveMessageResultFilter implements ReceiveMessageResultFilt } } - private void forwardMessageToDLQ(Context ctx, ReceiveMessageRequest request, MessageExt messageExt, int maxAttempt) { + @Override + protected void processExceedMaxAttemptsMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt, + int maxAttempts) { try { ReceiptHandle handle = ReceiptHandle.create(messageExt); if (handle == null) { @@ -106,7 +77,7 @@ public class LocalReceiveMessageResultFilter implements ReceiveMessageResultFilt } Channel channel = channelManager.createChannel(ctx); SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel); - ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader(request, handle, messageExt.getMsgId(), maxAttempt); + ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = GrpcConverter.buildConsumerSendMsgBackRequestHeader(request, handle, messageExt.getMsgId(), maxAttempts); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, consumerSendMsgBackRequestHeader); command.makeCustomHeaderToNet(); RemotingCommand response = brokerController.getSendMessageProcessor().processRequest(simpleChannelHandlerContext, command); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java index dc07124b75..13a87b25d1 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java @@ -32,8 +32,6 @@ import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.HeartbeatResponse; import apache.rocketmq.v2.Message; import apache.rocketmq.v2.MessageQueue; -import apache.rocketmq.v2.NackMessageRequest; -import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.NotifyClientTerminationRequest; import apache.rocketmq.v2.Publishing; import apache.rocketmq.v2.ReceiveMessageRequest; @@ -379,66 +377,6 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); } - @Test - public void testNackMessage() throws Exception { - ChangeInvisibleTimeResponseHeader responseHeader = new ChangeInvisibleTimeResponseHeader(); - responseHeader.setInvisibleTime(1000L); - responseHeader.setPopTime(0L); - responseHeader.setReviveQid(0); - RemotingCommand response = RemotingCommand.createResponseCommandWithHeader(ResponseCode.SUCCESS, responseHeader); - - ChangeInvisibleTimeProcessor changeInvisibleTimeProcessor = Mockito.mock(ChangeInvisibleTimeProcessor.class); - Mockito.when(brokerControllerMock.getChangeInvisibleTimeProcessor()).thenReturn(changeInvisibleTimeProcessor); - Mockito.when(changeInvisibleTimeProcessor.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) - .thenReturn(response); - NackMessageRequest request = NackMessageRequest.newBuilder().setReceiptHandle( - ReceiptHandle.builder() - .startOffset(0L) - .retrieveTime(0L) - .invisibleTime(1000L) - .nextVisibleTime(1000L) - .reviveQueueId(0) - .topicType("topic") - .brokerName("brokerName") - .queueId(0) - .offset(0L) - .build().encode() - ).build(); - CompletableFuture grpcFuture = localGrpcService.nackMessage(Context.current(), request); - NackMessageResponse r = grpcFuture.get(); - assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); - } - - @Test - public void testNackMessageWhenDLQ() throws Exception { - ConsumerSendMsgBackRequestHeader responseHeader = new ConsumerSendMsgBackRequestHeader(); - RemotingCommand response = RemotingCommand.createResponseCommandWithHeader(ResponseCode.SUCCESS, responseHeader); - - SendMessageProcessor sendMessageProcessor = Mockito.mock(SendMessageProcessor.class); - Mockito.when(brokerControllerMock.getSendMessageProcessor()).thenReturn(sendMessageProcessor); - Mockito.when(sendMessageProcessor.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class))) - .thenReturn(response); - NackMessageRequest request = NackMessageRequest.newBuilder() - .setDeliveryAttempt(3) - .setReceiptHandle( - ReceiptHandle.builder() - .startOffset(0L) - .retrieveTime(0L) - .invisibleTime(1000L) - .nextVisibleTime(1000L) - .reviveQueueId(0) - .topicType("topic") - .brokerName("brokerName") - .queueId(0) - .offset(0L) - .build().encode() - ).build(); - CompletableFuture grpcFuture = localGrpcService.nackMessage( - Context.current(), request); - NackMessageResponse r = grpcFuture.get(); - assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); - } - @Test public void testForwardMessageToDeadLetterQueue() throws Exception { RemotingCommand response = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, null); 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 7857f207d2..b8d045dc8c 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 @@ -25,8 +25,6 @@ import apache.rocketmq.v2.ClientType; import apache.rocketmq.v2.Code; import apache.rocketmq.v2.FilterExpression; import apache.rocketmq.v2.FilterType; -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; @@ -175,6 +173,8 @@ public class ConsumerServiceTest extends BaseServiceTest { ArgumentCaptor.forClass(ConsumerSendMsgBackRequestHeader.class); when(producerClient.sendMessageBackThenAckOrg(any(), anyString(), sendMsgBackRequestHeaderArgumentCaptor.capture(), any())) .thenReturn(CompletableFuture.completedFuture(RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, ""))); + when(writeConsumerClient.ackMessage(any(), anyString(), anyString(), any())) + .thenReturn(CompletableFuture.completedFuture(new AckResult())); Context ctx = Context.current().withDeadlineAfter(3, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor()); consumerService.receiveMessage(ctx, @@ -187,7 +187,7 @@ public class ConsumerServiceTest extends BaseServiceTest { .build()) .setFilterExpression(FilterExpression.newBuilder() .setType(FilterType.TAG) - .setExpression("msg1") + .setExpression("*") .build()) .build(), receiveMessageResponseStreamObserver @@ -228,63 +228,6 @@ public class ConsumerServiceTest extends BaseServiceTest { assertEquals(Code.OK, response.getStatus().getCode()); } - @Test - public void testNackMessageToDLQ() throws Exception { - ReceiptHandle receiptHandle = createReceiptHandle(); - ArgumentCaptor headerArgumentCaptor = ArgumentCaptor.forClass(ConsumerSendMsgBackRequestHeader.class); - when(producerClient.sendMessageBackThenAckOrg(any(), anyString(), headerArgumentCaptor.capture(), any())) - .thenReturn(CompletableFuture.completedFuture(RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, ""))); - when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); - - Settings clientSettings = createClientSettings(3); - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); - - NackMessageResponse response = consumerService.nackMessage(Context.current(), NackMessageRequest.newBuilder() - .setTopic(Resource.newBuilder() - .setName("topic") - .build()) - .setGroup(Resource.newBuilder() - .setName("group") - .build()) - .setReceiptHandle(receiptHandle.encode()) - .setDeliveryAttempt(3) - .build()) - .get(); - - assertEquals(Code.OK, response.getStatus().getCode()); - assertEquals(receiptHandle.getCommitLogOffset(), headerArgumentCaptor.getValue().getOffset().longValue()); - } - - @Test - public void testNackMessage() throws Exception { - ReceiptHandle receiptHandle = createReceiptHandle(); - ArgumentCaptor headerArgumentCaptor = ArgumentCaptor.forClass(ChangeInvisibleTimeRequestHeader.class); - AckResult ackResult = new AckResult(); - ackResult.setStatus(AckStatus.OK); - when(writeConsumerClient.changeInvisibleTimeAsync(any(), anyString(), anyString(), anyString(), headerArgumentCaptor.capture())) - .thenReturn(CompletableFuture.completedFuture(ackResult)); - when(topicRouteCache.getBrokerAddr(anyString())).thenReturn("brokerAddr"); - - Settings clientSettings = createClientSettings(3); - when(grpcClientManager.getClientSettings(any(Context.class))).thenReturn(clientSettings); - - NackMessageResponse response = consumerService.nackMessage(Context.current(), NackMessageRequest.newBuilder() - .setTopic(Resource.newBuilder() - .setName("topic") - .build()) - .setGroup(Resource.newBuilder() - .setName("group") - .build()) - .setReceiptHandle(receiptHandle.encode()) - .setDeliveryAttempt(1) - .build()) - .get(); - - assertEquals(Code.OK, response.getStatus().getCode()); - assertEquals(receiptHandle.getOffset(), headerArgumentCaptor.getValue().getOffset().longValue()); - assertEquals(receiptHandle.encode(), headerArgumentCaptor.getValue().getExtraInfo()); - } - @Test public void testChangeInvisibleDuration() throws Exception { Duration newDuration = Duration.newBuilder() diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionServiceTest.java index 5c762e3f12..eafc3a4325 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/TransactionServiceTest.java @@ -19,8 +19,10 @@ package org.apache.rocketmq.proxy.grpc.v2.service.cluster; import apache.rocketmq.v2.Code; import apache.rocketmq.v2.EndTransactionRequest; import apache.rocketmq.v2.EndTransactionResponse; +import apache.rocketmq.v2.RecoverOrphanedTransactionCommand; import apache.rocketmq.v2.TelemetryCommand; import io.grpc.Context; +import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader; import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; @@ -69,13 +71,16 @@ public class TransactionServiceTest extends BaseServiceTest { 2L, "msgId", transactionId, + "brokerName", createMessageExt("msgId", "msgId") )); Object flushData = flushDataCaptor.getValue(); assertTrue(flushData instanceof TelemetryCommand); TelemetryCommand response = (TelemetryCommand) flushData; - assertEquals(transactionId.getProxyTransactionId(), response.getRecoverOrphanedTransactionCommand().getTransactionId()); + RecoverOrphanedTransactionCommand command = response.getRecoverOrphanedTransactionCommand(); + assertEquals(transactionId.getProxyTransactionId(), command.getTransactionId()); + assertEquals("brokerName", command.getMessageQueue().getBroker().getName()); } @Test 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 index 431b236f83..82e6dc659c 100644 --- 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 @@ -76,21 +76,11 @@ public class ClusterGrpcTest extends GrpcBaseTest { assertQueryAssignment(response, brokerNum); } - @Test - public void testSendReceiveMessage() throws Exception { - super.testSendReceiveMessage(); - } - @Test public void testTransactionCheckThenCommit() { super.testTransactionCheckThenCommit(); } - @Test - public void testSendReceiveMessageThenToDLQ() throws Exception { - super.testSendReceiveMessageThenToDLQ(); - } - @Test public void testSimpleConsumerSendAndRecv() throws Exception { super.testSimpleConsumerSendAndRecv(); 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 b008ead8d0..a1a93620ba 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 @@ -35,8 +35,6 @@ import apache.rocketmq.v2.Message; import apache.rocketmq.v2.MessageQueue; import apache.rocketmq.v2.MessageType; import apache.rocketmq.v2.MessagingServiceGrpc; -import apache.rocketmq.v2.NackMessageRequest; -import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.Publishing; import apache.rocketmq.v2.QueryAssignmentRequest; import apache.rocketmq.v2.QueryAssignmentResponse; @@ -208,84 +206,6 @@ public class GrpcBaseTest extends BaseConf { .build()); } - public void testSendReceiveMessage() throws Exception { - String topic = initTopicOnSampleTopicBroker(broker1Name); - String group = MQRandomUtils.getRandomConsumerGroup(); - - // init consumer offset - this.sendClientSettings(stub, buildPushConsumerClientSettings()).get(); - receiveMessage(blockingStub, topic, group, 1); - - String messageId = createUniqID(); - this.sendClientSettings(stub, buildProducerClientSettings(topic)).get(); - SendMessageResponse sendResponse = blockingStub.sendMessage(buildSendMessageRequest(topic, messageId)); - assertSendMessage(sendResponse, messageId); - - this.sendClientSettings(stub, buildPushConsumerClientSettings()).get(); - - Message responseMessage = assertAndGetReceiveMessage(receiveMessage(blockingStub, topic, group), messageId); - String receiptHandle = responseMessage.getSystemProperties().getReceiptHandle(); - AckMessageResponse ackMessageResponse = blockingStub.ackMessage(buildAckMessageRequest(topic, group, messageId, receiptHandle)); - assertAllAckOk(ackMessageResponse); - } - - public void testSendReceiveMessageThenToDLQ() throws Exception { - String topic = initTopicOnSampleTopicBroker(broker1Name); - String group = MQRandomUtils.getRandomConsumerGroup(); - - // init consumer offset - this.sendClientSettings(stub, buildPushConsumerClientSettings()).get(); - receiveMessage(blockingStub, topic, group, 1); - - this.sendClientSettings(stub, buildProducerClientSettings(topic)).get(); - String messageId = createUniqID(); - SendMessageResponse sendResponse = blockingStub.sendMessage(buildSendMessageRequest(topic, messageId)); - assertSendMessage(sendResponse, messageId); - - this.sendClientSettings(stub, buildPushConsumerClientSettings()).get(); - - Message message = assertAndGetReceiveMessage(receiveMessage(blockingStub, topic, group), messageId); - - NackMessageResponse nackMessageResponse = blockingStub.nackMessage(buildNackMessageRequest( - topic, group, messageId, message.getSystemProperties().getReceiptHandle(), 1 - )); - assertNackMessageResponse(nackMessageResponse); - - AtomicReference receiveRetryMessageRef = new AtomicReference<>(); - await().atMost(java.time.Duration.ofSeconds(30)).until(() -> { - List messageList = getMessageFromReceiveMessageResponse(receiveMessage(blockingStub, topic, group, 1)); - if (messageList.isEmpty()) { - return false; - } - - receiveRetryMessageRef.set(messageList.get(0)); - return messageList.get(0).getSystemProperties() - .getMessageId().equals(messageId); - }); - - message = receiveRetryMessageRef.get(); - nackMessageResponse = blockingStub.nackMessage(buildNackMessageRequest( - topic, group, messageId, message.getSystemProperties().getReceiptHandle(), 2 - )); - assertNackMessageResponse(nackMessageResponse); - - DefaultMQPullConsumer defaultMQPullConsumer = new DefaultMQPullConsumer(group); - defaultMQPullConsumer.start(); - org.apache.rocketmq.common.message.MessageQueue dlqMQ = new org.apache.rocketmq.common.message.MessageQueue(MixAll.getDLQTopic(group), broker1Name, 0); - await().atMost(java.time.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; - } - }); - } - public void testTransactionCheckThenCommit() { String topic = initTopicOnSampleTopicBroker(broker1Name); String group = MQRandomUtils.getRandomConsumerGroup(); @@ -590,22 +510,6 @@ public class GrpcBaseTest extends BaseConf { .build(); } - public NackMessageRequest buildNackMessageRequest(String topic, String group, String messageId, - String receiptHandle, - int deliveryAttempt) { - return NackMessageRequest.newBuilder() - .setDeliveryAttempt(deliveryAttempt) - .setMessageId(messageId) - .setReceiptHandle(receiptHandle) - .setTopic(Resource.newBuilder() - .setName(topic) - .build()) - .setGroup(Resource.newBuilder() - .setName(group) - .build()) - .build(); - } - public EndTransactionRequest buildEndTransactionRequest(String topic, String messageId, String transactionId, TransactionResolution resolution) { return EndTransactionRequest.newBuilder() @@ -666,10 +570,6 @@ public class GrpcBaseTest extends BaseConf { } } - public void assertNackMessageResponse(NackMessageResponse response) { - assertThat(response.getStatus().getCode()).isEqualTo(Code.OK); - } - public void assertRecoverOrphanedTransactionCommand(RecoverOrphanedTransactionCommand command, String messageId) { assertThat(command.getOrphanedTransactionalMessage().getSystemProperties().getMessageId()) .isEqualTo(messageId); 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 d28f4311fd..d99b28ea99 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 @@ -62,21 +62,11 @@ public class LocalGrpcTest extends GrpcBaseTest { assertQueryAssignment(response, brokerNum); } - @Test - public void testSendReceiveMessage() throws Exception { - super.testSendReceiveMessage(); - } - @Test public void testTransactionCheckThenCommit() { super.testTransactionCheckThenCommit(); } - @Test - public void testSendReceiveMessageThenToDLQ() throws Exception { - super.testSendReceiveMessageThenToDLQ(); - } - @Test public void testSimpleConsumerSendAndRecv() throws Exception { super.testSimpleConsumerSendAndRecv();