From 42ee31ec4df6145456bb6420fa7c5c858dc462e7 Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Wed, 16 Mar 2022 14:55:27 +0800 Subject: [PATCH] [ISSUE #3949] Implement ForwardMessageToDeadLetterQueue --- .../rocketmq/proxy/grpc/common/Converter.java | 19 ++++++++++++ .../proxy/grpc/service/LocalGrpcService.java | 29 ++++++++++++++++++- 2 files changed, 47 insertions(+), 1 deletion(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java index 229e0a4daa..bfc005d259 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java @@ -27,6 +27,7 @@ import apache.rocketmq.v1.DigestType; import apache.rocketmq.v1.Encoding; import apache.rocketmq.v1.FilterExpression; import apache.rocketmq.v1.FilterType; +import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueRequest; import apache.rocketmq.v1.HeartbeatRequest; import apache.rocketmq.v1.Message; import apache.rocketmq.v1.MessageType; @@ -65,6 +66,7 @@ import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.common.protocol.NamespaceUtil; import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader; +import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType; @@ -179,6 +181,23 @@ public class Converter { return changeInvisibleTimeRequestHeader; } + public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackRequestHeader( + ForwardMessageToDeadLetterQueueRequest request) { + String groupName = Converter.getResourceNameWithNamespace(request.getGroup()); + String topicName = Converter.getResourceNameWithNamespace(request.getTopic()); + String receiptHandleStr = request.getReceiptHandle(); + ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); + + ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader(); + consumerSendMsgBackRequestHeader.setOffset(handle.getOffset()); + consumerSendMsgBackRequestHeader.setGroup(groupName); + consumerSendMsgBackRequestHeader.setDelayLevel(-1); + consumerSendMsgBackRequestHeader.setOriginMsgId(request.getMessageId()); + consumerSendMsgBackRequestHeader.setOriginMsgId(topicName); + consumerSendMsgBackRequestHeader.setMaxReconsumeTimes(request.getMaxDeliveryAttempts()); + return consumerSendMsgBackRequestHeader; + } + public static Map buildMessageProperty(Message message) { org.apache.rocketmq.common.message.Message messageWithHeader = new org.apache.rocketmq.common.message.Message(); // set user properties diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index 06dafea8e7..9169f191ee 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -66,10 +66,12 @@ import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.protocol.RequestCode; import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader; +import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData; import org.apache.rocketmq.proxy.channel.ChannelManager; +import org.apache.rocketmq.proxy.channel.SimpleChannel; import org.apache.rocketmq.proxy.channel.SimpleChannelHandlerContext; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext; @@ -252,7 +254,32 @@ public class LocalGrpcService implements GrpcForwardService { @Override public CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request) { - return null; + SimpleChannel channel = channelManager.createChannel(); + SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel); + + ConsumerSendMsgBackRequestHeader requestHeader = Converter.buildConsumerSendMsgBackRequestHeader(request); + RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader); + command.makeCustomHeaderToNet(); + + CompletableFuture future = new CompletableFuture<>(); + try { + CompletableFuture processorFuture = brokerController.getSendMessageProcessor() + .asyncProcessRequest(channelHandlerContext, command); + processorFuture.thenAccept(r -> { + ForwardMessageToDeadLetterQueueResponse.Builder builder = ForwardMessageToDeadLetterQueueResponse.newBuilder(); + if (null != r) { + builder.setCommon(ResponseBuilder.buildCommon(r.getCode(), r.getRemark())); + } else { + builder.setCommon(ResponseBuilder.buildCommon(Code.INTERNAL, "Response command is null")); + } + ForwardMessageToDeadLetterQueueResponse response = builder.build(); + future.complete(response); + }); + } catch (Exception e) { + LOGGER.error("Exception raised when forwardMessageToDeadLetterQueue", e); + future.completeExceptionally(e); + } + return future; } @Override public CompletableFuture endTransaction(Context ctx, EndTransactionRequest request) {