mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] v2 support
This commit is contained in:
@@ -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);
|
||||
|
||||
|
||||
-29
@@ -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<PullMessageRequest, PullMessageResponse> {
|
||||
public PullMessageChannel(PullMessageResponseHandler handler) {
|
||||
super(handler);
|
||||
}
|
||||
}
|
||||
-57
@@ -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<PullMessageRequest, PullMessageResponse> {
|
||||
@Override
|
||||
public void handle(RemotingCommand responseCommand,
|
||||
InvocationContext<PullMessageRequest, PullMessageResponse> context) {
|
||||
try {
|
||||
PullMessageResponseHeader responseHeader = (PullMessageResponseHeader) responseCommand.readCustomHeader();
|
||||
PullMessageResponse.Builder builder = PullMessageResponse.newBuilder();
|
||||
if (responseCommand.getCode() == ResponseCode.SUCCESS) {
|
||||
ByteBuffer byteBuffer = ByteBuffer.wrap(responseCommand.getBody());
|
||||
List<MessageExt> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user