From 4b4cedc980ee5ccf24eb80bbb98aafce63e1b1d3 Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Fri, 8 Apr 2022 16:09:43 +0800 Subject: [PATCH] [ISSUE #3949] Support v2 Implement v2 Local mode Add V2Converter for Cluster mode --- .../proxy/grpc/v1/adapter/V2Converter.java | 68 ++++++++++++++++ .../grpc/v1/service/ClusterGrpcService.java | 79 +++---------------- .../grpc/v2/adapter/ResponseBuilder.java | 6 +- .../grpc/v2/service/ClusterGrpcService.java | 6 +- .../grpc/v2/service/GrpcClientManager.java | 14 +--- .../grpc/v2/service/LocalGrpcService.java | 23 +++--- .../v2/service/cluster/ConsumerService.java | 10 ++- .../grpc/v2/service/cluster/RouteService.java | 11 ++- .../grpc/v2/service/LocalGrpcServiceTest.java | 57 +++++++++---- 9 files changed, 162 insertions(+), 112 deletions(-) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/adapter/V2Converter.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/adapter/V2Converter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/adapter/V2Converter.java new file mode 100644 index 0000000000..64ce1275db --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/adapter/V2Converter.java @@ -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; + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/service/ClusterGrpcService.java index 94b1e56798..1aed35126d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/service/ClusterGrpcService.java @@ -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 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 forwardMessageToDeadLetterQueue(Context ctx, - ForwardMessageToDeadLetterQueueRequest request) { + ForwardMessageToDeadLetterQueueRequest request) { return null; } @@ -179,7 +140,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc @Override public CompletableFuture 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 reportMessageConsumptionResult(Context ctx, - ReportMessageConsumptionResultRequest request) { + ReportMessageConsumptionResultRequest request) { return null; } @Override public CompletableFuture 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); - } - } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java index 8c2fc13dd6..eecc19f53c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/ResponseBuilder.java @@ -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(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java index e8454971ff..515039b263 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java @@ -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); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java index 8371d2eaee..eb86c41cfc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcClientManager.java @@ -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 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); } } 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 48f3169e30..730ef5ea0f 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 @@ -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 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 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 telemetry(Context ctx, StreamObserver responseObserver) { + String clientId = InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.CLIENT_ID); return new StreamObserver() { @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() diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java index d696dd07bd..7267b78b85 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java @@ -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 ackMessageHook; private volatile ResponseHook 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 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 resultFuture = this.producer.sendMessageBack( diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java index 83322e2400..d107b093a4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/RouteService.java @@ -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 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 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() diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java index f3e276c91f..9e2f51b79f 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/LocalGrpcServiceTest.java @@ -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 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() { + @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 grpcFuture = localGrpcService.heartbeat( - Context.current().withValue(InterceptorConstants.METADATA, metadata).attach(), request); + CompletableFuture 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 grpcFuture = localGrpcService.heartbeat( - Context.current().withValue(InterceptorConstants.METADATA, metadata).attach(), request); + CompletableFuture 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 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 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 grpcFuture = localGrpcService.queryOffset( - Context.current() - .withValue(InterceptorConstants.METADATA, metadata) - .attach(), request); + CompletableFuture grpcFuture = localGrpcService.queryOffset(Context.current(), request); QueryOffsetResponse r = grpcFuture.get(); assertThat(r.getStatus().getCode()).isEqualTo(Code.OK); assertThat(r.getOffset()).isEqualTo(0);