mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Support v2
Implement v2 Local mode Add V2Converter for Cluster mode
This commit is contained in:
@@ -0,0 +1,68 @@
|
||||
/*
|
||||
* 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.v1.adapter;
|
||||
|
||||
import apache.rocketmq.v1.ResponseCommon;
|
||||
import apache.rocketmq.v2.Code;
|
||||
import apache.rocketmq.v2.HeartbeatRequest;
|
||||
import apache.rocketmq.v2.HeartbeatResponse;
|
||||
import apache.rocketmq.v2.Resource;
|
||||
import apache.rocketmq.v2.Status;
|
||||
|
||||
public class V2Converter {
|
||||
public static Resource buildResource(apache.rocketmq.v1.Resource resource) {
|
||||
return Resource.newBuilder()
|
||||
.setName(resource.getName())
|
||||
.setResourceNamespace(resource.getResourceNamespace())
|
||||
.build();
|
||||
}
|
||||
|
||||
public static HeartbeatRequest buildHeartbeatRequest(apache.rocketmq.v1.HeartbeatRequest request) {
|
||||
Resource group;
|
||||
if (request.hasProducerData()) {
|
||||
group = buildResource(request.getProducerData().getGroup());
|
||||
} else if (request.hasConsumerData()) {
|
||||
group = buildResource(request.getConsumerData().getGroup());
|
||||
} else {
|
||||
throw new IllegalArgumentException("HeartbeatRequest is not valid");
|
||||
}
|
||||
return HeartbeatRequest.newBuilder()
|
||||
.setGroup(group)
|
||||
.build();
|
||||
}
|
||||
|
||||
public static apache.rocketmq.v1.HeartbeatResponse buildHeartbeatResponse(HeartbeatResponse response) {
|
||||
return apache.rocketmq.v1.HeartbeatResponse.newBuilder()
|
||||
.setCommon(ResponseCommon.newBuilder()
|
||||
.setStatus(buildStatus(response.getStatus()))
|
||||
.build())
|
||||
.build();
|
||||
}
|
||||
|
||||
public static com.google.rpc.Status buildStatus(Status status) {
|
||||
return com.google.rpc.Status.newBuilder()
|
||||
.setCode(buildCodeValue(status.getCode()))
|
||||
.setMessage(status.getMessage())
|
||||
.build();
|
||||
}
|
||||
|
||||
public static int buildCodeValue(Code code) {
|
||||
// TODO: complete code mapping
|
||||
return com.google.rpc.Code.OK_VALUE;
|
||||
}
|
||||
}
|
||||
+9
-70
@@ -54,57 +54,22 @@ import apache.rocketmq.v1.SendMessageResponse;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.common.StartAndShutdown;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionStateCheckRequest;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker;
|
||||
import org.apache.rocketmq.proxy.common.PollResponseManager;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ForwardClientService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ConsumerService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ProducerService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.PullMessageService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.RouteService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.TransactionService;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.adapter.V2Converter;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class ClusterGrpcService extends AbstractStartAndShutdown implements GrpcForwardService {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
|
||||
private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(
|
||||
new ThreadFactoryImpl("ClusterGrpcServiceScheduledThread"));
|
||||
|
||||
private final ChannelManager channelManager;
|
||||
private final ConnectorManager connectorManager;
|
||||
private final ProducerService producerService;
|
||||
private final ConsumerService receiveMessageService;
|
||||
private final RouteService routeService;
|
||||
private final ForwardClientService clientService;
|
||||
private final PullMessageService pullMessageService;
|
||||
private final TransactionService transactionService;
|
||||
private final PollResponseManager pollCommandResponseManager;
|
||||
private final org.apache.rocketmq.proxy.grpc.v2.service.ClusterGrpcService clusterGrpcService;
|
||||
|
||||
public ClusterGrpcService() {
|
||||
this.channelManager = new ChannelManager();
|
||||
this.pollCommandResponseManager = new PollResponseManager();
|
||||
this.connectorManager = new ConnectorManager(new GrpcTransactionStateChecker());
|
||||
this.receiveMessageService = new ConsumerService(connectorManager);
|
||||
this.producerService = new ProducerService(connectorManager);
|
||||
this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager);
|
||||
this.clientService = new ForwardClientService(connectorManager, scheduledExecutorService, channelManager, pollCommandResponseManager);
|
||||
this.pullMessageService = new PullMessageService(connectorManager);
|
||||
this.transactionService = new TransactionService(connectorManager, channelManager);
|
||||
this.clusterGrpcService = new org.apache.rocketmq.proxy.grpc.v2.service.ClusterGrpcService();;
|
||||
|
||||
this.appendStartAndShutdown(new ClusterGrpcServiceStartAndShutdown());
|
||||
this.appendStartAndShutdown(this.connectorManager);
|
||||
this.appendStartAndShutdown(clusterGrpcService);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -114,12 +79,8 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
|
||||
@Override
|
||||
public CompletableFuture<HeartbeatResponse> heartbeat(Context ctx, HeartbeatRequest request) {
|
||||
this.clientService.heartbeat(ctx, request);
|
||||
return CompletableFuture.completedFuture(
|
||||
HeartbeatResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
.build()
|
||||
);
|
||||
return clusterGrpcService.heartbeat(ctx, V2Converter.buildHeartbeatRequest(request))
|
||||
.thenApply(V2Converter::buildHeartbeatResponse);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -158,7 +119,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
|
||||
@Override
|
||||
public CompletableFuture<ForwardMessageToDeadLetterQueueResponse> forwardMessageToDeadLetterQueue(Context ctx,
|
||||
ForwardMessageToDeadLetterQueueRequest request) {
|
||||
ForwardMessageToDeadLetterQueueRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -179,7 +140,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
|
||||
@Override
|
||||
public CompletableFuture<PollCommandResponse> pollCommand(Context ctx, PollCommandRequest request) {
|
||||
return this.clientService.pollCommand(ctx, request);
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -190,14 +151,13 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
|
||||
@Override
|
||||
public CompletableFuture<ReportMessageConsumptionResultResponse> reportMessageConsumptionResult(Context ctx,
|
||||
ReportMessageConsumptionResultRequest request) {
|
||||
ReportMessageConsumptionResultRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<NotifyClientTerminationResponse> notifyClientTermination(Context ctx,
|
||||
NotifyClientTerminationRequest request) {
|
||||
this.clientService.unregister(ctx, request);
|
||||
return CompletableFuture.completedFuture(
|
||||
NotifyClientTerminationResponse.newBuilder()
|
||||
.setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name()))
|
||||
@@ -210,25 +170,4 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
private class ClusterGrpcServiceStartAndShutdown implements StartAndShutdown {
|
||||
|
||||
@Override
|
||||
public void start() throws Exception {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shutdown() throws Exception {
|
||||
scheduledExecutorService.shutdown();
|
||||
}
|
||||
}
|
||||
|
||||
private class GrpcTransactionStateChecker implements TransactionStateChecker {
|
||||
|
||||
@Override
|
||||
public void checkTransactionState(TransactionStateCheckRequest checkData) {
|
||||
transactionService.checkTransactionState(checkData);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -30,9 +30,13 @@ public class ResponseBuilder {
|
||||
}
|
||||
|
||||
public static Status buildStatus(int remotingResponseCode, String remark) {
|
||||
String message = remark;
|
||||
if (message == null) {
|
||||
message = String.valueOf(remotingResponseCode);
|
||||
}
|
||||
return Status.newBuilder()
|
||||
.setCode(buildCode(remotingResponseCode))
|
||||
.setMessage(remark)
|
||||
.setMessage(message)
|
||||
.build();
|
||||
}
|
||||
|
||||
|
||||
+4
-2
@@ -83,14 +83,16 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
private final PullMessageService pullMessageService;
|
||||
private final TransactionService transactionService;
|
||||
private final PollResponseManager pollCommandResponseManager;
|
||||
private final GrpcClientManager grpcClientManager;
|
||||
|
||||
public ClusterGrpcService() {
|
||||
this.channelManager = new ChannelManager();
|
||||
this.grpcClientManager = new GrpcClientManager();
|
||||
this.pollCommandResponseManager = new PollResponseManager();
|
||||
this.connectorManager = new ConnectorManager(new GrpcTransactionStateChecker());
|
||||
this.consumerService = new ConsumerService(connectorManager);
|
||||
this.consumerService = new ConsumerService(connectorManager, grpcClientManager);
|
||||
this.producerService = new ProducerService(connectorManager);
|
||||
this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager);
|
||||
this.routeService = new RouteService(ProxyMode.CLUSTER, connectorManager, grpcClientManager);
|
||||
this.clientService = new ForwardClientService(connectorManager, scheduledExecutorService, channelManager, pollCommandResponseManager);
|
||||
this.pullMessageService = new PullMessageService(connectorManager);
|
||||
this.transactionService = new TransactionService(connectorManager, channelManager);
|
||||
|
||||
+4
-10
@@ -18,24 +18,18 @@
|
||||
package org.apache.rocketmq.proxy.grpc.v2.service;
|
||||
|
||||
import apache.rocketmq.v2.ClientSettings;
|
||||
import io.grpc.Context;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
|
||||
public class GrpcClientManager {
|
||||
|
||||
private static final Map<String, ClientSettings> CLIENT_SETTINGS_MAP = new ConcurrentHashMap<>();
|
||||
|
||||
public static ClientSettings getClientSettings(Context ctx) {
|
||||
return CLIENT_SETTINGS_MAP.get(getClientId(ctx));
|
||||
public ClientSettings getClientSettings(String clientId) {
|
||||
return CLIENT_SETTINGS_MAP.get(clientId);
|
||||
}
|
||||
|
||||
public static void updateClientSettings(Context ctx, ClientSettings clientSettings) {
|
||||
CLIENT_SETTINGS_MAP.put(getClientId(ctx), clientSettings);
|
||||
}
|
||||
|
||||
public static String getClientId(Context ctx) {
|
||||
return InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
|
||||
public void updateClientSettings(String clientId, ClientSettings clientSettings) {
|
||||
CLIENT_SETTINGS_MAP.put(clientId, clientSettings);
|
||||
}
|
||||
}
|
||||
|
||||
+13
-10
@@ -122,6 +122,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
new ThreadFactoryImpl("LocalGrpcServiceScheduledThread"));
|
||||
private final ChannelManager channelManager;
|
||||
private final PollResponseManager pollCommandResponseManager;
|
||||
private final GrpcClientManager grpcClientManager;
|
||||
private final RouteService routeService;
|
||||
private final DelayPolicy delayPolicy;
|
||||
|
||||
@@ -131,7 +132,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
// TransactionStateChecker is not used in Local mode.
|
||||
ConnectorManager connectorManager = new ConnectorManager(null);
|
||||
this.pollCommandResponseManager = new PollResponseManager();
|
||||
this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager);
|
||||
this.grpcClientManager = new GrpcClientManager();
|
||||
this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager, grpcClientManager);
|
||||
this.delayPolicy = DelayPolicy.build(brokerController.getMessageStoreConfig().getMessageDelayLevel());
|
||||
this.appendStartAndShutdown(connectorManager);
|
||||
this.appendStartAndShutdown(new LocalGrpcServiceStartAndShutdown());
|
||||
@@ -146,10 +148,10 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
public CompletableFuture<HeartbeatResponse> heartbeat(Context ctx, HeartbeatRequest request) {
|
||||
LanguageCode languageCode;
|
||||
String language = InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.LANGUAGE);
|
||||
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
|
||||
languageCode = LanguageCode.valueOf(language);
|
||||
String clientId = GrpcClientManager.getClientId(ctx);
|
||||
|
||||
ClientSettings clientSettings = GrpcClientManager.getClientSettings(ctx);
|
||||
ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
|
||||
HeartbeatData heartbeatData = GrpcConverter.buildHeartbeatData(clientId, request, clientSettings);
|
||||
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.HEART_BEAT, null);
|
||||
command.setLanguage(languageCode);
|
||||
@@ -239,7 +241,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
@Override
|
||||
public CompletableFuture<ReceiveMessageResponse> receiveMessage(Context ctx, ReceiveMessageRequest request) {
|
||||
long pollTime = GrpcConverter.buildPollTimeFromContext(ctx);
|
||||
ClientSettings clientSettings = GrpcClientManager.getClientSettings(ctx);
|
||||
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
|
||||
ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
|
||||
PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime);
|
||||
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader);
|
||||
command.makeCustomHeaderToNet();
|
||||
@@ -457,8 +460,8 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
NotifyClientTerminationRequest request) {
|
||||
Channel channel = channelManager.createChannel();
|
||||
SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel);
|
||||
String clientId = GrpcClientManager.getClientId(ctx);
|
||||
ClientSettings clientSettings = GrpcClientManager.getClientSettings(ctx);
|
||||
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
|
||||
ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
|
||||
UnregisterClientRequestHeader header = GrpcConverter.buildUnregisterClientRequestHeader(clientId, clientSettings.getClientType(), request);
|
||||
|
||||
RemotingCommand remotingCommand = RemotingCommand.createRequestCommand(RequestCode.UNREGISTER_CLIENT, header);
|
||||
@@ -512,27 +515,27 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
|
||||
@Override
|
||||
public StreamObserver<TelemetryCommand> telemetry(Context ctx, StreamObserver<TelemetryCommand> responseObserver) {
|
||||
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
|
||||
return new StreamObserver<TelemetryCommand>() {
|
||||
@Override
|
||||
public void onNext(TelemetryCommand request) {
|
||||
switch (request.getCommandCase()) {
|
||||
case CLIENT_SETTINGS: {
|
||||
ClientSettings clientSettings = request.getClientSettings();
|
||||
GrpcClientManager.updateClientSettings(ctx, clientSettings);
|
||||
String clientId = GrpcClientManager.getClientId(ctx);
|
||||
grpcClientManager.updateClientSettings(clientId, clientSettings);
|
||||
Settings settings = clientSettings.getSettings();
|
||||
if (settings.hasPublishing()) {
|
||||
Publishing publishing = settings.getPublishing();
|
||||
for (Resource topic : publishing.getTopicsList()) {
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
GrpcClientChannel producerChannel = GrpcClientChannel.getChannel(channelManager, topicName, clientId);
|
||||
GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, pollCommandResponseManager);
|
||||
producerChannel.setClientObserver(responseObserver);
|
||||
}
|
||||
}
|
||||
if (settings.hasSubscription()) {
|
||||
Subscription subscription = settings.getSubscription();
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(subscription.getGroup());
|
||||
GrpcClientChannel consumerChannel = GrpcClientChannel.getChannel(channelManager, groupName, clientId);
|
||||
GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, pollCommandResponseManager);
|
||||
consumerChannel.setClientObserver(responseObserver);
|
||||
}
|
||||
responseObserver.onNext(TelemetryCommand.newBuilder()
|
||||
|
||||
+8
-2
@@ -48,6 +48,7 @@ import org.apache.rocketmq.proxy.connector.ForwardReadConsumer;
|
||||
import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.common.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
@@ -70,7 +71,9 @@ public class ConsumerService extends BaseService {
|
||||
private volatile ResponseHook<AckMessageRequest, AckMessageResponse> ackMessageHook;
|
||||
private volatile ResponseHook<NackMessageRequest, NackMessageResponse> nackMessageResponseResponseHook;
|
||||
|
||||
public ConsumerService(ConnectorManager connectorManager) {
|
||||
private final GrpcClientManager grpcClientManager;
|
||||
|
||||
public ConsumerService(ConnectorManager connectorManager, GrpcClientManager grpcClientManager) {
|
||||
super(connectorManager);
|
||||
this.readConsumer = connectorManager.getForwardReadConsumer();
|
||||
this.writeConsumer = connectorManager.getForwardWriteConsumer();
|
||||
@@ -78,6 +81,8 @@ public class ConsumerService extends BaseService {
|
||||
|
||||
this.readQueueSelector = new DefaultReadQueueSelector(connectorManager.getTopicRouteCache());
|
||||
this.delayPolicy = DelayPolicy.build(ConfigurationManager.getProxyConfig().getMessageDelayLevel());
|
||||
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
}
|
||||
|
||||
public CompletableFuture<ReceiveMessageResponse> receiveMessage(Context ctx, ReceiveMessageRequest request) {
|
||||
@@ -246,10 +251,11 @@ public class ConsumerService extends BaseService {
|
||||
}
|
||||
});
|
||||
try {
|
||||
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
|
||||
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
|
||||
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
|
||||
|
||||
Settings settings = GrpcClientManager.getClientSettings(ctx).getSettings();
|
||||
Settings settings = grpcClientManager.getClientSettings(clientId).getSettings();
|
||||
int maxDeliveryAttempts = settings.getSubscription().getDeadLetterPolicy().getMaxDeliveryAttempts();
|
||||
if (request.getDeliveryAttempt() >= maxDeliveryAttempts) {
|
||||
CompletableFuture<RemotingCommand> resultFuture = this.producer.sendMessageBack(
|
||||
|
||||
+8
-3
@@ -42,13 +42,14 @@ import org.apache.rocketmq.common.constant.PermName;
|
||||
import org.apache.rocketmq.common.protocol.route.BrokerData;
|
||||
import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.common.ParameterConverter;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper;
|
||||
import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.common.ParameterConverter;
|
||||
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.ResponseHook;
|
||||
@@ -64,13 +65,16 @@ public class RouteService extends BaseService {
|
||||
private volatile AssignmentQueueSelector assignmentQueueSelector;
|
||||
private volatile ResponseHook<QueryAssignmentRequest, QueryAssignmentResponse> queryAssignmentHook;
|
||||
|
||||
public RouteService(ProxyMode mode, ConnectorManager connectorManager) {
|
||||
private GrpcClientManager grpcClientManager;
|
||||
|
||||
public RouteService(ProxyMode mode, ConnectorManager connectorManager, GrpcClientManager grpcClientManager) {
|
||||
super(connectorManager);
|
||||
Preconditions.checkArgument(ProxyMode.isClusterMode(mode) || ProxyMode.isLocalMode(mode));
|
||||
this.mode = mode;
|
||||
queryRouteEndpointConverter = (ctx, parameter) -> parameter;
|
||||
queryAssignmentEndpointConverter = (ctx, parameter) -> parameter;
|
||||
assignmentQueueSelector = new DefaultAssignmentQueueSelector(this.connectorManager.getTopicRouteCache());
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
}
|
||||
|
||||
public void setQueryRouteEndpointConverter(ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter) {
|
||||
@@ -240,7 +244,8 @@ public class RouteService extends BaseService {
|
||||
}
|
||||
}
|
||||
if (ProxyMode.isClusterMode(mode)) {
|
||||
ClientSettings clientSettings = GrpcClientManager.getClientSettings(ctx);
|
||||
String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID);
|
||||
ClientSettings clientSettings = grpcClientManager.getClientSettings(clientId);
|
||||
Endpoints resEndpoints = this.queryAssignmentEndpointConverter.convert(ctx, clientSettings.getAccessPoint());
|
||||
if (resEndpoints == null || Endpoints.getDefaultInstance().equals(resEndpoints)) {
|
||||
future.complete(QueryAssignmentResponse.newBuilder()
|
||||
|
||||
+43
-14
@@ -21,6 +21,8 @@ import apache.rocketmq.v2.AckMessageRequest;
|
||||
import apache.rocketmq.v2.AckMessageResponse;
|
||||
import apache.rocketmq.v2.ChangeInvisibleDurationRequest;
|
||||
import apache.rocketmq.v2.ChangeInvisibleDurationResponse;
|
||||
import apache.rocketmq.v2.ClientSettings;
|
||||
import apache.rocketmq.v2.ClientType;
|
||||
import apache.rocketmq.v2.Code;
|
||||
import apache.rocketmq.v2.EndTransactionRequest;
|
||||
import apache.rocketmq.v2.EndTransactionResponse;
|
||||
@@ -33,6 +35,7 @@ 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.PullMessageRequest;
|
||||
import apache.rocketmq.v2.PullMessageResponse;
|
||||
import apache.rocketmq.v2.QueryOffsetPolicy;
|
||||
@@ -43,11 +46,14 @@ 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.SystemProperties;
|
||||
import apache.rocketmq.v2.TelemetryCommand;
|
||||
import com.google.protobuf.Timestamp;
|
||||
import com.google.protobuf.util.Durations;
|
||||
import io.grpc.Context;
|
||||
import io.grpc.Metadata;
|
||||
import io.grpc.stub.StreamObserver;
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
@@ -105,6 +111,8 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
|
||||
private Metadata metadata;
|
||||
|
||||
private StreamObserver<TelemetryCommand> streamObserver;
|
||||
|
||||
@Before
|
||||
public void setUp() throws Throwable {
|
||||
super.before();
|
||||
@@ -119,10 +127,35 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
metadata.put(InterceptorConstants.LOCAL_ADDRESS, "0.0.0.0");
|
||||
metadata.put(InterceptorConstants.LANGUAGE, "JAVA");
|
||||
metadata.put(InterceptorConstants.CLIENT_ID, "client-id");
|
||||
Context.current().withValue(InterceptorConstants.METADATA, metadata).attach();
|
||||
streamObserver = localGrpcService.telemetry(Context.current(), new StreamObserver<TelemetryCommand>() {
|
||||
@Override public void onNext(TelemetryCommand value) {
|
||||
}
|
||||
|
||||
@Override public void onError(Throwable t) {
|
||||
}
|
||||
|
||||
@Override public void onCompleted() {
|
||||
}
|
||||
});
|
||||
streamObserver.onNext(TelemetryCommand.newBuilder()
|
||||
.setClientSettings(ClientSettings.newBuilder().setSettings(Settings.getDefaultInstance()))
|
||||
.build());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testHeartbeatProducerData() throws Exception {
|
||||
streamObserver.onNext(TelemetryCommand.newBuilder()
|
||||
.setClientSettings(ClientSettings.newBuilder()
|
||||
.setSettings(Settings.newBuilder()
|
||||
.setPublishing(Publishing.newBuilder()
|
||||
.addTopics(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.build())
|
||||
.build())
|
||||
.setClientType(ClientType.PRODUCER).build())
|
||||
.build());
|
||||
RemotingCommand response = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, null);
|
||||
ClientManageProcessor clientManageProcessorMock = Mockito.mock(ClientManageProcessor.class);
|
||||
Mockito.when(clientManageProcessorMock.heartBeat(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class)))
|
||||
@@ -133,8 +166,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
.setName("group")
|
||||
.build())
|
||||
.build();
|
||||
CompletableFuture<HeartbeatResponse> grpcFuture = localGrpcService.heartbeat(
|
||||
Context.current().withValue(InterceptorConstants.METADATA, metadata).attach(), request);
|
||||
CompletableFuture<HeartbeatResponse> grpcFuture = localGrpcService.heartbeat(Context.current(), request);
|
||||
HeartbeatResponse r = grpcFuture.get();
|
||||
assertThat(r.getStatus().getCode())
|
||||
.isEqualTo(Code.OK);
|
||||
@@ -142,6 +174,10 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
|
||||
@Test
|
||||
public void testHeartbeatConsumerData() throws Exception {
|
||||
streamObserver.onNext(TelemetryCommand.newBuilder()
|
||||
.setClientSettings(ClientSettings.newBuilder()
|
||||
.setClientType(ClientType.PUSH_CONSUMER).build())
|
||||
.build());
|
||||
RemotingCommand response = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, null);
|
||||
ClientManageProcessor clientManageProcessorMock = Mockito.mock(ClientManageProcessor.class);
|
||||
Mockito.when(clientManageProcessorMock.heartBeat(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class)))
|
||||
@@ -152,8 +188,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
.setName("group")
|
||||
.build())
|
||||
.build();
|
||||
CompletableFuture<HeartbeatResponse> grpcFuture = localGrpcService.heartbeat(
|
||||
Context.current().withValue(InterceptorConstants.METADATA, metadata).attach(), request);
|
||||
CompletableFuture<HeartbeatResponse> grpcFuture = localGrpcService.heartbeat(Context.current(), request);
|
||||
HeartbeatResponse r = grpcFuture.get();
|
||||
assertThat(r.getStatus().getCode())
|
||||
.isEqualTo(Code.OK);
|
||||
@@ -166,7 +201,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class)))
|
||||
.thenReturn(response);
|
||||
SendMessageRequest request = SendMessageRequest.newBuilder()
|
||||
.setMessages(0, Message.newBuilder()
|
||||
.addMessages(0, Message.newBuilder()
|
||||
.setSystemProperties(SystemProperties.newBuilder()
|
||||
.setMessageId("123")
|
||||
.build())
|
||||
@@ -185,7 +220,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class)))
|
||||
.thenReturn(null);
|
||||
SendMessageRequest request = SendMessageRequest.newBuilder()
|
||||
.setMessages(0, Message.newBuilder()
|
||||
.addMessages(0, Message.newBuilder()
|
||||
.setSystemProperties(SystemProperties.newBuilder()
|
||||
.setMessageId("123")
|
||||
.build())
|
||||
@@ -202,7 +237,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
Mockito.when(sendMessageProcessorMock.processRequest(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class)))
|
||||
.thenThrow(new RemotingCommandException("test"));
|
||||
SendMessageRequest request = SendMessageRequest.newBuilder()
|
||||
.setMessages(0, Message.newBuilder()
|
||||
.addMessages(0, Message.newBuilder()
|
||||
.setSystemProperties(SystemProperties.newBuilder()
|
||||
.setMessageId("123")
|
||||
.build())
|
||||
@@ -249,7 +284,6 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
CompletableFuture<ReceiveMessageResponse> grpcFuture = localGrpcService.receiveMessage(
|
||||
Context.current()
|
||||
.withValue(InterceptorConstants.METADATA, metadata)
|
||||
.attach()
|
||||
.withDeadlineAfter(20, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor(
|
||||
new ThreadFactoryImpl("test"))), request);
|
||||
ReceiveMessageResponse r = grpcFuture.get();
|
||||
@@ -267,8 +301,6 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
ReceiveMessageRequest request = ReceiveMessageRequest.newBuilder().getDefaultInstanceForType();
|
||||
CompletableFuture<ReceiveMessageResponse> grpcFuture = localGrpcService.receiveMessage(
|
||||
Context.current()
|
||||
.withValue(InterceptorConstants.METADATA, metadata)
|
||||
.attach()
|
||||
.withDeadlineAfter(20, TimeUnit.SECONDS, Executors.newSingleThreadScheduledExecutor(
|
||||
new ThreadFactoryImpl("test"))), request);
|
||||
assertThat(grpcFuture.isDone()).isFalse();
|
||||
@@ -408,10 +440,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
.build())
|
||||
.setPolicy(QueryOffsetPolicy.BEGINNING)
|
||||
.build();
|
||||
CompletableFuture<QueryOffsetResponse> grpcFuture = localGrpcService.queryOffset(
|
||||
Context.current()
|
||||
.withValue(InterceptorConstants.METADATA, metadata)
|
||||
.attach(), request);
|
||||
CompletableFuture<QueryOffsetResponse> grpcFuture = localGrpcService.queryOffset(Context.current(), request);
|
||||
QueryOffsetResponse r = grpcFuture.get();
|
||||
assertThat(r.getStatus().getCode()).isEqualTo(Code.OK);
|
||||
assertThat(r.getOffset()).isEqualTo(0);
|
||||
|
||||
Reference in New Issue
Block a user