From 2c53e88f6cd8e6301da7ae5ca6f58ae5ebfd2d68 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Thu, 21 Apr 2022 17:42:44 +0800 Subject: [PATCH] [ISSUE #3949] v2 support --- .../proxy/grpc/v2/adapter/RequestMapping.java | 4 -- .../adapter/channel/PullMessageChannel.java | 29 ---------- .../handler/PullMessageResponseHandler.java | 57 ------------------- .../grpc/v2/service/LocalGrpcService.java | 15 ----- 4 files changed, 105 deletions(-) delete mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/PullMessageChannel.java delete mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/PullMessageResponseHandler.java 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 add450e5a1..304d23e9ec 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 @@ -24,9 +24,7 @@ import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; import apache.rocketmq.v2.HeartbeatRequest; import apache.rocketmq.v2.NackMessageRequest; import apache.rocketmq.v2.NotifyClientTerminationRequest; -import apache.rocketmq.v2.PullMessageRequest; import apache.rocketmq.v2.QueryAssignmentRequest; -import apache.rocketmq.v2.QueryOffsetRequest; import apache.rocketmq.v2.QueryRouteRequest; import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.SendMessageRequest; @@ -47,8 +45,6 @@ public class RequestMapping { put(NackMessageRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); put(ForwardMessageToDeadLetterQueueResponse.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); put(EndTransactionRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); - put(QueryOffsetRequest.getDescriptor().getFullName(), RequestCode.SEARCH_OFFSET_BY_TIMESTAMP); - put(PullMessageRequest.getDescriptor().getFullName(), RequestCode.PULL_MESSAGE); put(NotifyClientTerminationRequest.getDescriptor().getFullName(), RequestCode.UNREGISTER_CLIENT); put(ChangeInvisibleDurationRequest.getDescriptor().getFullName(), RequestCode.CONSUMER_SEND_MSG_BACK); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/PullMessageChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/PullMessageChannel.java deleted file mode 100644 index d65c74aba9..0000000000 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/PullMessageChannel.java +++ /dev/null @@ -1,29 +0,0 @@ -/* - * 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.adapter.channel; - -import apache.rocketmq.v2.PullMessageRequest; -import apache.rocketmq.v2.PullMessageResponse; -import org.apache.rocketmq.proxy.channel.InvocationChannel; -import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.PullMessageResponseHandler; - -public class PullMessageChannel extends InvocationChannel { - public PullMessageChannel(PullMessageResponseHandler handler) { - super(handler); - } -} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/PullMessageResponseHandler.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/PullMessageResponseHandler.java deleted file mode 100644 index d4e351121a..0000000000 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/PullMessageResponseHandler.java +++ /dev/null @@ -1,57 +0,0 @@ -/* - * 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.adapter.handler; - -import apache.rocketmq.v2.PullMessageRequest; -import apache.rocketmq.v2.PullMessageResponse; -import java.nio.ByteBuffer; -import java.util.List; -import org.apache.rocketmq.common.message.MessageDecoder; -import org.apache.rocketmq.common.message.MessageExt; -import org.apache.rocketmq.common.protocol.ResponseCode; -import org.apache.rocketmq.common.protocol.header.PullMessageResponseHeader; -import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; -import org.apache.rocketmq.proxy.channel.InvocationContext; -import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; -import org.apache.rocketmq.remoting.protocol.RemotingCommand; - -public class PullMessageResponseHandler implements ResponseHandler { - @Override - public void handle(RemotingCommand responseCommand, - InvocationContext context) { - try { - PullMessageResponseHeader responseHeader = (PullMessageResponseHeader) responseCommand.readCustomHeader(); - PullMessageResponse.Builder builder = PullMessageResponse.newBuilder(); - if (responseCommand.getCode() == ResponseCode.SUCCESS) { - ByteBuffer byteBuffer = ByteBuffer.wrap(responseCommand.getBody()); - List msgFoundList = MessageDecoder.decodes(byteBuffer); - for (MessageExt messageExt : msgFoundList) { - builder.addMessages(GrpcConverter.buildMessage(messageExt)); - } - } - PullMessageResponse response = builder.setStatus(ResponseBuilder.buildStatus(responseCommand.getCode(), responseCommand.getRemark())) - .setMinOffset(responseHeader.getMinOffset()) - .setNextOffset(responseHeader.getNextBeginOffset()) - .setMaxOffset(responseHeader.getMaxOffset()) - .build(); - context.getResponse().complete(response); - } catch (Exception e) { - context.getResponse().completeExceptionally(e); - } - } -} 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 bca408ee58..fe42e56c1d 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 @@ -21,10 +21,7 @@ import apache.rocketmq.v2.AckMessageRequest; import apache.rocketmq.v2.AckMessageResponse; import apache.rocketmq.v2.ChangeInvisibleDurationRequest; import apache.rocketmq.v2.ChangeInvisibleDurationResponse; -import apache.rocketmq.v2.ClientOverwrittenSettings; -import apache.rocketmq.v2.ClientSettings; import apache.rocketmq.v2.Code; -import apache.rocketmq.v2.Direction; import apache.rocketmq.v2.EndTransactionRequest; import apache.rocketmq.v2.EndTransactionResponse; import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest; @@ -35,14 +32,8 @@ import apache.rocketmq.v2.NackMessageRequest; import apache.rocketmq.v2.NackMessageResponse; import apache.rocketmq.v2.NotifyClientTerminationRequest; import apache.rocketmq.v2.NotifyClientTerminationResponse; -import apache.rocketmq.v2.Publishing; -import apache.rocketmq.v2.PullMessageRequest; -import apache.rocketmq.v2.PullMessageResponse; import apache.rocketmq.v2.QueryAssignmentRequest; import apache.rocketmq.v2.QueryAssignmentResponse; -import apache.rocketmq.v2.QueryOffsetPolicy; -import apache.rocketmq.v2.QueryOffsetRequest; -import apache.rocketmq.v2.QueryOffsetResponse; import apache.rocketmq.v2.QueryRouteRequest; import apache.rocketmq.v2.QueryRouteResponse; import apache.rocketmq.v2.ReceiveMessageRequest; @@ -50,12 +41,9 @@ import apache.rocketmq.v2.ReceiveMessageResponse; import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.SendMessageRequest; import apache.rocketmq.v2.SendMessageResponse; -import apache.rocketmq.v2.Settings; -import apache.rocketmq.v2.Subscription; import apache.rocketmq.v2.TelemetryCommand; import apache.rocketmq.v2.ThreadStackTrace; import apache.rocketmq.v2.VerifyMessageResult; -import com.google.protobuf.util.Timestamps; import io.grpc.Context; import io.grpc.stub.StreamObserver; import io.netty.channel.Channel; @@ -86,7 +74,6 @@ import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeResponseHea import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader; import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; -import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.UnregisterClientRequestHeader; import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData; @@ -105,10 +92,8 @@ import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel; -import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.PullMessageChannel; import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.ReceiveMessageChannel; import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.SendMessageChannel; -import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.PullMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.ReceiveMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.v2.adapter.handler.SendMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.v2.service.cluster.RouteService;