[ISSUE #3949] v2 support

This commit is contained in:
kaiyi.lk
2022-07-13 11:29:29 +08:00
committed by zhouxiang
parent d08dfce4fc
commit 67d749b569
26 changed files with 181 additions and 576 deletions
@@ -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();
@@ -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);
@@ -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 + '\'' +
@@ -68,6 +68,7 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor {
requestHeader.getTransactionId(),
requestHeader.getCommitLogOffset(),
requestHeader.getTranStateTableOffset()),
requestHeader.getBrokerName(),
messageExt
)
);
@@ -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;
}
@@ -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<NackMessageResponse> responseObserver) {
CompletableFuture<NackMessageResponse> 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<AckMessageResponse> responseObserver) {
CompletableFuture<AckMessageResponse> future = grpcForwardService.ackMessage(Context.current(), request);
@@ -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:
@@ -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);
@@ -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;
@@ -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<ReceiveMessageResponse> streamObserver;
protected final ResponseHook<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook;
protected final ReceiveMessageResultFilter receiveMessageResultFilter;
public interface Builder {
ReceiveMessageResponseStreamWriter build(
BaseReceiveMessageResponseStreamWriter build(
StreamObserver<ReceiveMessageResponse> observer,
ResponseHook<ReceiveMessageRequest, ReceiveMessageResponse> hook);
}
public ReceiveMessageResponseStreamWriter(
public BaseReceiveMessageResponseStreamWriter(
StreamObserver<ReceiveMessageResponse> observer,
ResponseHook<ReceiveMessageRequest, ReceiveMessageResponse> hook,
ReceiveMessageResultFilter messageResultFilter) {
@@ -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<Message> filterMessage(Context ctx, ReceiveMessageRequest request, List<MessageExt> 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<Message> 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);
}
@@ -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<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request) {
return consumerService.nackMessage(ctx, request);
}
@Override
public CompletableFuture<AckMessageResponse> ackMessage(Context ctx, AckMessageRequest request) {
return consumerService.ackMessage(ctx, request);
@@ -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<ReceiveMessageResponse> responseObserver);
CompletableFuture<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request);
CompletableFuture<AckMessageResponse> ackMessage(Context ctx, AckMessageRequest request);
CompletableFuture<ForwardMessageToDeadLetterQueueResponse> forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request);
@@ -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<ReceiveMessageRequest, ReceiveMessageResponse> 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<List<MessageExt>> future = new CompletableFuture<>();
@@ -345,50 +342,6 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
return future;
}
@Override
public CompletableFuture<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request) {
Channel channel = channelManager.createChannel(ctx);
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
CompletableFuture<NackMessageResponse> 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<ForwardMessageToDeadLetterQueueResponse> forwardMessageToDeadLetterQueue(Context ctx,
ForwardMessageToDeadLetterQueueRequest request) {
@@ -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<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook;
private volatile ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook;
private volatile ResponseHook<NackMessageRequest, NackMessageResponse> nackMessageHook;
private volatile ResponseHook<ChangeInvisibleDurationRequest, ChangeInvisibleDurationResponse> changeInvisibleDurationHook;
public ConsumerService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) {
@@ -92,7 +85,7 @@ public class ConsumerService extends BaseService {
public void receiveMessage(Context ctx, ReceiveMessageRequest request,
StreamObserver<ReceiveMessageResponse> 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<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request) {
CompletableFuture<NackMessageResponse> 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<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx,
ChangeInvisibleDurationRequest request) {
CompletableFuture<ChangeInvisibleDurationResponse> 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<NackMessageRequest, NackMessageResponse> getNackMessageHook() {
return nackMessageHook;
}
public void setNackMessageHook(
ResponseHook<NackMessageRequest, NackMessageResponse> nackMessageHook) {
this.nackMessageHook = nackMessageHook;
}
public ResponseHook<ChangeInvisibleDurationRequest, ChangeInvisibleDurationResponse> getChangeInvisibleDurationHook() {
return changeInvisibleDurationHook;
}
@@ -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();
@@ -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<AckMessageRequestHeader, AckResult> 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<Message> filterMessage(Context ctx, ReceiveMessageRequest request, List<MessageExt> 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<Message> 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<RemotingCommand> 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<AckResult> 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<RemotingCommand> 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<AckMessageRequestHeader, AckResult> getAckNoMatchedMessageHook() {
return ackNoMatchedMessageHook;
}
@@ -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();
@@ -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;
@@ -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<Message> filterMessage(Context ctx, ReceiveMessageRequest request, List<MessageExt> 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<Message> 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);
@@ -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<NackMessageResponse> 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<NackMessageResponse> 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);
@@ -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<ConsumerSendMsgBackRequestHeader> 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<ChangeInvisibleTimeRequestHeader> 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()
@@ -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
@@ -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();
@@ -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<Message> receiveRetryMessageRef = new AtomicReference<>();
await().atMost(java.time.Duration.ofSeconds(30)).until(() -> {
List<Message> 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);
@@ -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();