From 62c14bfdb942c1cbae2c8ab94fc01aec61486616 Mon Sep 17 00:00:00 2001 From: zhouxiang Date: Wed, 6 Apr 2022 16:58:07 +0800 Subject: [PATCH] [ISSUE #3949] Add telemetry --- .../proxy/grpc/GrpcMessagingProcessorV2.java | 65 ++--- .../proxy/grpc/adapter/GrpcConverterV2.java | 20 +- .../adapter/channel/GrpcClientChannelV2.java | 161 ++++++++++++ .../proxy/grpc/service/GrpcClientManager.java | 21 +- .../grpc/service/GrpcForwardServiceV2.java | 6 +- .../proxy/grpc/service/LocalGrpcService.java | 241 +++++++++++------- 6 files changed, 362 insertions(+), 152 deletions(-) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannelV2.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessorV2.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessorV2.java index 6a282ff55a..3f5dd2e089 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessorV2.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessorV2.java @@ -17,34 +17,34 @@ package org.apache.rocketmq.proxy.grpc; +import apache.rocketmq.v2.AckMessageRequest; +import apache.rocketmq.v2.AckMessageResponse; +import apache.rocketmq.v2.ChangeInvisibleDurationRequest; +import apache.rocketmq.v2.ChangeInvisibleDurationResponse; import apache.rocketmq.v2.Code; +import apache.rocketmq.v2.EndTransactionRequest; +import apache.rocketmq.v2.EndTransactionResponse; +import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest; +import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; +import apache.rocketmq.v2.HeartbeatRequest; +import apache.rocketmq.v2.HeartbeatResponse; +import apache.rocketmq.v2.MessagingServiceGrpc; import apache.rocketmq.v2.NackMessageRequest; import apache.rocketmq.v2.NackMessageResponse; -import apache.rocketmq.v2.QueryRouteResponse; -import apache.rocketmq.v2.QueryRouteRequest; -import apache.rocketmq.v2.HeartbeatResponse; -import apache.rocketmq.v2.HeartbeatRequest; -import apache.rocketmq.v2.SendMessageResponse; -import apache.rocketmq.v2.SendMessageRequest; -import apache.rocketmq.v2.QueryAssignmentResponse; -import apache.rocketmq.v2.QueryAssignmentRequest; -import apache.rocketmq.v2.ReceiveMessageResponse; -import apache.rocketmq.v2.ReceiveMessageRequest; -import apache.rocketmq.v2.AckMessageResponse; -import apache.rocketmq.v2.AckMessageRequest; -import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueResponse; -import apache.rocketmq.v2.ForwardMessageToDeadLetterQueueRequest; -import apache.rocketmq.v2.EndTransactionResponse; -import apache.rocketmq.v2.EndTransactionRequest; -import apache.rocketmq.v2.QueryOffsetResponse; -import apache.rocketmq.v2.QueryOffsetRequest; -import apache.rocketmq.v2.PullMessageResponse; -import apache.rocketmq.v2.PullMessageRequest; -import apache.rocketmq.v2.NotifyClientTerminationResponse; import apache.rocketmq.v2.NotifyClientTerminationRequest; -import apache.rocketmq.v2.ChangeInvisibleDurationResponse; -import apache.rocketmq.v2.ChangeInvisibleDurationRequest; -import apache.rocketmq.v2.MessagingServiceGrpc; +import apache.rocketmq.v2.NotifyClientTerminationResponse; +import apache.rocketmq.v2.PullMessageRequest; +import apache.rocketmq.v2.PullMessageResponse; +import apache.rocketmq.v2.QueryAssignmentRequest; +import apache.rocketmq.v2.QueryAssignmentResponse; +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; +import apache.rocketmq.v2.ReceiveMessageResponse; +import apache.rocketmq.v2.SendMessageRequest; +import apache.rocketmq.v2.SendMessageResponse; import apache.rocketmq.v2.Status; import apache.rocketmq.v2.TelemetryCommand; import io.grpc.Context; @@ -249,21 +249,6 @@ public class GrpcMessagingProcessorV2 extends MessagingServiceGrpc.MessagingServ @Override public StreamObserver telemetry(StreamObserver responseObserver) { - return new StreamObserver() { - @Override - public void onNext(TelemetryCommand request) { - - } - - @Override - public void onError(Throwable t) { - - } - - @Override - public void onCompleted() { - - } - }; + return grpcForwardService.telemetry(Context.current(), responseObserver); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverterV2.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverterV2.java index e48f062183..7679fb4d8b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverterV2.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverterV2.java @@ -18,6 +18,7 @@ package org.apache.rocketmq.proxy.grpc.adapter; import apache.rocketmq.v2.AckMessageRequest; +import apache.rocketmq.v2.ChangeInvisibleDurationRequest; import apache.rocketmq.v2.ClientSettings; import apache.rocketmq.v2.ClientType; import apache.rocketmq.v2.Code; @@ -59,6 +60,7 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; import org.apache.rocketmq.common.constant.ConsumeInitMode; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.consumer.ConsumeFromWhere; @@ -267,7 +269,23 @@ public class GrpcConverterV2 { return changeInvisibleTimeRequestHeader; } - public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader(NackMessageRequest request, int maxReconsumeTimes) { + public static ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(ChangeInvisibleDurationRequest request) { + String groupName = GrpcConverterV2.wrapResourceWithNamespace(request.getGroup()); + String topicName = GrpcConverterV2.wrapResourceWithNamespace(request.getTopic()); + ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); + + ChangeInvisibleTimeRequestHeader changeInvisibleTimeRequestHeader = new ChangeInvisibleTimeRequestHeader(); + changeInvisibleTimeRequestHeader.setConsumerGroup(groupName); + changeInvisibleTimeRequestHeader.setTopic(handle.getRealTopic(topicName, groupName)); + changeInvisibleTimeRequestHeader.setQueueId(handle.getQueueId()); + changeInvisibleTimeRequestHeader.setExtraInfo(handle.getReceiptHandle()); + changeInvisibleTimeRequestHeader.setOffset(handle.getOffset()); + changeInvisibleTimeRequestHeader.setInvisibleTime(Durations.toMillis(request.getInvisibleDuration())); + return changeInvisibleTimeRequestHeader; + } + + public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader( + NackMessageRequest request, int maxReconsumeTimes) { String groupName = GrpcConverterV2.wrapResourceWithNamespace(request.getGroup()); String topicName = GrpcConverterV2.wrapResourceWithNamespace(request.getTopic()); ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannelV2.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannelV2.java new file mode 100644 index 0000000000..793df7b679 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannelV2.java @@ -0,0 +1,161 @@ +/* + * 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.adapter.channel; + +import apache.rocketmq.v2.PrintThreadStackTraceCommand; +import apache.rocketmq.v2.RecoverOrphanedTransactionCommand; +import apache.rocketmq.v2.TelemetryCommand; +import io.grpc.Context; +import io.grpc.stub.StreamObserver; +import io.netty.channel.ChannelFuture; +import java.nio.ByteBuffer; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.rocketmq.common.message.MessageDecoder; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.protocol.RequestCode; +import org.apache.rocketmq.common.protocol.header.CheckTransactionStateRequestHeader; +import org.apache.rocketmq.common.protocol.header.GetConsumerRunningInfoRequestHeader; +import org.apache.rocketmq.proxy.channel.ChannelManager; +import org.apache.rocketmq.proxy.channel.SimpleChannel; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverterV2; +import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager; +import org.apache.rocketmq.remoting.protocol.RemotingCommand; + +public class GrpcClientChannelV2 extends SimpleChannel { + private final AtomicReference> telemetryCommandRef = new AtomicReference<>(); + + private final String group; + private final String clientId; + private final PollResponseManager manager; + + private GrpcClientChannelV2(Context ctx, String group, String clientId, PollResponseManager manager) { + super(ChannelManager.createSimpleChannelDirectly(ctx)); + this.group = group; + this.clientId = clientId; + this.manager = manager; + } + + public void setClientObserver(StreamObserver future) { + this.telemetryCommandRef.set(future); + } + + public static GrpcClientChannelV2 create( + ChannelManager channelManager, + String group, + String clientId, + PollResponseManager manager + ) { + return create(Context.current(), channelManager, group, clientId, manager); + } + + public static GrpcClientChannelV2 create( + Context ctx, + ChannelManager channelManager, + String group, + String clientId, + PollResponseManager manager + ) { + GrpcClientChannelV2 channel = channelManager.createChannel( + buildKey(group, clientId), + () -> new GrpcClientChannelV2(ctx, group, clientId, manager), + GrpcClientChannelV2.class + ); + + channelManager.addGroupClientId(group, clientId); + return channel; + } + + public static GrpcClientChannelV2 getChannel(ChannelManager channelManager, String group, String clientId) { + return channelManager.getChannel(buildKey(group, clientId), GrpcClientChannelV2.class); + } + + public static GrpcClientChannelV2 removeChannel(ChannelManager channelManager, String group, String clientId) { + return channelManager.removeChannel(buildKey(group, clientId), GrpcClientChannelV2.class); + } + + private static String buildKey(String group, String clientId) { + return group + "@" + clientId; + } + + @Override + public boolean isWritable() { + if (this.telemetryCommandRef.get() == null) { + return false; + } + return true; + } + + /** + * Write response to corresponding remote client + * + * @param msg Target write object, {@link RemotingCommand} or {@link TelemetryCommand} + * @return Always success {@link ChannelFuture} + *

+ * Case {@link RequestCode#CHECK_TRANSACTION_STATE} + * @see org.apache.rocketmq.broker.client.net.Broker2Client#checkProducerTransactionState + */ + @Override + public ChannelFuture writeAndFlush(Object msg) { + StreamObserver streamObserver = telemetryCommandRef.get(); + if (msg instanceof RemotingCommand) { + RemotingCommand command = (RemotingCommand) msg; + try { + switch (command.getCode()) { + case RequestCode.CHECK_TRANSACTION_STATE: { + final CheckTransactionStateRequestHeader requestHeader = command.decodeCommandCustomHeader(CheckTransactionStateRequestHeader.class); + MessageExt messageExt = MessageDecoder.decode(ByteBuffer.wrap(command.getBody()), true, false, false); + streamObserver.onNext(TelemetryCommand.newBuilder() + .setRecoverOrphanedTransactionCommand(RecoverOrphanedTransactionCommand.newBuilder() + .setTransactionId(requestHeader.getTransactionId()) + .setOrphanedTransactionalMessage(GrpcConverterV2.buildMessage(messageExt)) + .build()) + .build()); + break; + } + case RequestCode.GET_CONSUMER_RUNNING_INFO: { + final GetConsumerRunningInfoRequestHeader requestHeader = command.decodeCommandCustomHeader(GetConsumerRunningInfoRequestHeader.class); + if (!requestHeader.isJstackEnable()) { + break; + } + String nonce = manager.putResponse(command.getOpaque()); + streamObserver.onNext(TelemetryCommand.newBuilder() + .setPrintThreadStackTraceCommand(PrintThreadStackTraceCommand.newBuilder() + .setNonce(nonce) + .build()) + .build()); + break; + } + } + } catch (Exception ignore) { + + } + } + if (msg instanceof TelemetryCommand) { + TelemetryCommand response = (TelemetryCommand) msg; + streamObserver.onNext(response); + } + return super.writeAndFlush(msg); + } + + public String getGroup() { + return group; + } + + public String getClientId() { + return clientId; + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcClientManager.java index 3785f6c265..2911a72ae6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcClientManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcClientManager.java @@ -18,37 +18,24 @@ package org.apache.rocketmq.proxy.grpc.service; import apache.rocketmq.v2.ClientSettings; -import apache.rocketmq.v2.TelemetryCommand; import io.grpc.Context; -import io.grpc.stub.StreamObserver; 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_DATA = new ConcurrentHashMap<>(); + private static final Map CLIENT_SETTINGS_MAP = new ConcurrentHashMap<>(); public static ClientSettings getClientSettings(Context ctx) { - return CLIENT_DATA.get(getClientId(ctx)).clientSettings; + return CLIENT_SETTINGS_MAP.get(getClientId(ctx)); } - public static void updateClientData(Context ctx, ClientSettings clientSettings, StreamObserver responseStreamObserver) { - CLIENT_DATA.put(getClientId(ctx), new ClientData(clientSettings, responseStreamObserver)); + 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 static class ClientData { - private final ClientSettings clientSettings; - private final StreamObserver responseStreamObserver; - - public ClientData(ClientSettings clientSettings, - StreamObserver responseStreamObserver) { - this.clientSettings = clientSettings; - this.responseStreamObserver = responseStreamObserver; - } - } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcForwardServiceV2.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcForwardServiceV2.java index 0e58a68188..ac231b8bdf 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcForwardServiceV2.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/GrpcForwardServiceV2.java @@ -43,7 +43,9 @@ import apache.rocketmq.v2.ReceiveMessageRequest; import apache.rocketmq.v2.ReceiveMessageResponse; import apache.rocketmq.v2.SendMessageRequest; import apache.rocketmq.v2.SendMessageResponse; +import apache.rocketmq.v2.TelemetryCommand; import io.grpc.Context; +import io.grpc.stub.StreamObserver; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.proxy.common.StartAndShutdown; @@ -63,8 +65,6 @@ public interface GrpcForwardServiceV2 extends StartAndShutdown { CompletableFuture ackMessage(Context ctx, AckMessageRequest request); - CompletableFuture nackMessage(Context ctx, NackMessageRequest request); - CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request); CompletableFuture endTransaction(Context ctx, EndTransactionRequest request); @@ -76,4 +76,6 @@ public interface GrpcForwardServiceV2 extends StartAndShutdown { CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request); CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request); + + StreamObserver telemetry(Context ctx, StreamObserver responseObserver); } 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 b7c1841313..4a9daf6c75 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 @@ -21,8 +21,10 @@ 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; @@ -33,6 +35,7 @@ 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; @@ -44,10 +47,17 @@ import apache.rocketmq.v2.QueryRouteRequest; import apache.rocketmq.v2.QueryRouteResponse; import apache.rocketmq.v2.ReceiveMessageRequest; 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; import java.util.List; import java.util.concurrent.CompletableFuture; @@ -63,6 +73,8 @@ import org.apache.rocketmq.common.message.MessageBatch; import org.apache.rocketmq.common.message.MessageClientIDSetter; import org.apache.rocketmq.common.protocol.RequestCode; import org.apache.rocketmq.common.protocol.ResponseCode; +import org.apache.rocketmq.common.protocol.body.ConsumeMessageDirectlyResult; +import org.apache.rocketmq.common.protocol.body.ConsumerRunningInfo; import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader; import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeResponseHeader; @@ -83,10 +95,11 @@ import org.apache.rocketmq.proxy.grpc.adapter.DelayPolicy; import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverterV2; import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext; +import org.apache.rocketmq.proxy.grpc.adapter.PollResponseFuture; import org.apache.rocketmq.proxy.grpc.adapter.PollResponseManager; import org.apache.rocketmq.proxy.grpc.adapter.ProxyMode; import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilderV2; -import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; +import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannelV2; import org.apache.rocketmq.proxy.grpc.adapter.channel.PullMessageChannel; import org.apache.rocketmq.proxy.grpc.adapter.channel.ReceiveMessageChannel; import org.apache.rocketmq.proxy.grpc.adapter.channel.SendMessageChannel; @@ -95,6 +108,8 @@ import org.apache.rocketmq.proxy.grpc.adapter.handler.ReceiveMessageResponseHand import org.apache.rocketmq.proxy.grpc.adapter.handler.SendMessageResponseHandler; import org.apache.rocketmq.proxy.grpc.interceptor.InterceptorConstants; import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService; +import org.apache.rocketmq.remoting.RemotingServer; +import org.apache.rocketmq.remoting.netty.NettyRemotingAbstract; import org.apache.rocketmq.remoting.protocol.LanguageCode; import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.slf4j.Logger; @@ -137,24 +152,49 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo ClientSettings clientSettings = GrpcClientManager.getClientSettings(ctx); HeartbeatData heartbeatData = GrpcConverterV2.buildHeartbeatData(clientId, request, clientSettings); - - CompletableFuture future = new CompletableFuture<>(); - String groupName = GrpcConverterV2.wrapResourceWithNamespace(request.getGroup()); - - GrpcClientChannel channel = GrpcClientChannel.create(channelManager, groupName, clientId, pollCommandResponseManager); - SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.HEART_BEAT, null); command.setLanguage(languageCode); command.setVersion(MQVersion.Version.V5_0_0.ordinal()); command.setBody(heartbeatData.encode()); command.makeCustomHeaderToNet(); - RemotingCommand response = this.brokerController.getClientManageProcessor() - .heartBeat(simpleChannelHandlerContext, command); - HeartbeatResponse heartbeatResponse = HeartbeatResponse.newBuilder() - .setStatus(ResponseBuilderV2.buildStatus(response.getCode(), response.getRemark())) - .build(); - future.complete(heartbeatResponse); + CompletableFuture future = new CompletableFuture<>(); + switch (clientSettings.getClientType()) { + case PRODUCER: { + for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) { + String topicName = GrpcConverterV2.wrapResourceWithNamespace(topic); + GrpcClientChannelV2 channel = GrpcClientChannelV2.create(channelManager, topicName, clientId, pollCommandResponseManager); + SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel); + + this.brokerController.getClientManageProcessor() + .heartBeat(simpleChannelHandlerContext, command); + } + HeartbeatResponse heartbeatResponse = HeartbeatResponse.newBuilder() + .setStatus(ResponseBuilderV2.buildStatus(Code.OK, "Producer heartbeat")) + .build(); + future.complete(heartbeatResponse); + break; + } + case PULL_CONSUMER: + case PUSH_CONSUMER: + case SIMPLE_CONSUMER: { + String groupName = GrpcConverterV2.wrapResourceWithNamespace(request.getGroup()); + GrpcClientChannelV2 channel = GrpcClientChannelV2.create(channelManager, groupName, clientId, pollCommandResponseManager); + SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel); + + RemotingCommand response = this.brokerController.getClientManageProcessor() + .heartBeat(simpleChannelHandlerContext, command); + HeartbeatResponse heartbeatResponse = HeartbeatResponse.newBuilder() + .setStatus(ResponseBuilderV2.buildStatus(response.getCode(), response.getRemark())) + .build(); + future.complete(heartbeatResponse); + break; + } + default: { + throw new IllegalArgumentException("ClientType not exist " + clientSettings.getClientType()); + } + } + return future; } @@ -376,85 +416,42 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo return future; } -// @Override -// public CompletableFuture pollCommand(Context ctx, PollCommandRequest request) { -// String clientId = request.getClientId(); -// CompletableFuture future = new CompletableFuture<>(); -// switch (request.getGroupCase()) { -// case PRODUCER_GROUP: -// Resource producerGroup = request.getProducerGroup(); -// String producerGroupName = GrpcConverter.wrapResourceWithNamespace(producerGroup); -// GrpcClientChannel producerChannel = GrpcClientChannel.getChannel(channelManager, producerGroupName, clientId); -// if (producerChannel == null) { -// future.complete(PollCommandResponse.newBuilder() -// .setNoopCommand(NoopCommand.newBuilder().build()) -// .build()); -// break; -// } -// producerChannel.setClientObserver(future); -// break; -// case CONSUMER_GROUP: -// Resource consumerGroup = request.getConsumerGroup(); -// String consumerGroupName = GrpcConverter.wrapResourceWithNamespace(consumerGroup); -// GrpcClientChannel consumerChannel = GrpcClientChannel.getChannel(channelManager, consumerGroupName, clientId); -// if (consumerChannel == null) { -// future.complete(PollCommandResponse.newBuilder() -// .setNoopCommand(NoopCommand.newBuilder().build()) -// .build()); -// break; -// } -// consumerChannel.setClientObserver(future); -// break; -// default: -// break; -// } -// return future; -// } -// -// @Override -// public CompletableFuture reportThreadStackTrace(Context ctx, -// ReportThreadStackTraceRequest request) { -// String commandId = request.getCommandId(); -// String threadStack = request.getThreadStackTrace(); -// PollResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(commandId); -// if (pollCommandResponseFuture != null) { -// RemotingServer remotingServer = this.brokerController.getRemotingServer(); -// if (remotingServer instanceof NettyRemotingAbstract) { -// NettyRemotingAbstract nettyRemotingAbstract = (NettyRemotingAbstract) remotingServer; -// RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "From gRPC client"); -// remotingCommand.setOpaque(pollCommandResponseFuture.getOpaque()); -// ConsumerRunningInfo runningInfo = new ConsumerRunningInfo(); -// runningInfo.setJstack(threadStack); -// remotingCommand.setBody(runningInfo.encode()); -// nettyRemotingAbstract.processResponseCommand(new SimpleChannelHandlerContext(channelManager.createChannel()), remotingCommand); -// } -// } -// return CompletableFuture.completedFuture(ReportThreadStackTraceResponse.newBuilder() -// .setCommon(ResponseBuilder.buildSuccessCommon()) -// .build()); -// } -// -// @Override -// public CompletableFuture reportMessageConsumptionResult(Context ctx, -// ReportMessageConsumptionResultRequest request) { -// -// String commandId = request.getCommandId(); -// PollResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(commandId); -// if (pollCommandResponseFuture != null) { -// RemotingServer remotingServer = this.brokerController.getRemotingServer(); -// if (remotingServer instanceof NettyRemotingAbstract) { -// NettyRemotingAbstract nettyRemotingAbstract = (NettyRemotingAbstract) remotingServer; -// RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "From gRPC client"); -// remotingCommand.setOpaque(pollCommandResponseFuture.getOpaque()); -// ConsumeMessageDirectlyResult result = GrpcConverter.buildConsumeMessageDirectlyResult(request); -// remotingCommand.setBody(result.encode()); -// nettyRemotingAbstract.processResponseCommand(new SimpleChannelHandlerContext(channelManager.createChannel()), remotingCommand); -// } -// } -// return CompletableFuture.completedFuture(ReportMessageConsumptionResultResponse.newBuilder() -// .setCommon(ResponseBuilder.buildSuccessCommon()) -// .build()); -// } + public void reportThreadStackTrace(ThreadStackTrace request) { + String nonce = request.getNonce(); + String threadStack = request.getThreadStackTrace(); + PollResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(nonce); + if (pollCommandResponseFuture != null) { + RemotingServer remotingServer = this.brokerController.getRemotingServer(); + if (remotingServer instanceof NettyRemotingAbstract) { + NettyRemotingAbstract nettyRemotingAbstract = (NettyRemotingAbstract) remotingServer; + RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "From gRPC client"); + remotingCommand.setOpaque(pollCommandResponseFuture.getOpaque()); + ConsumerRunningInfo runningInfo = new ConsumerRunningInfo(); + runningInfo.setJstack(threadStack); + remotingCommand.setBody(runningInfo.encode()); + nettyRemotingAbstract.processResponseCommand(new SimpleChannelHandlerContext(channelManager.createChannel()), remotingCommand); + } + } + } + + public void reportVerifyMessageResult(VerifyMessageResult request) { + String nonce = request.getNonce(); + PollResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(nonce); + if (pollCommandResponseFuture != null) { + Integer opaque = pollCommandResponseFuture.getOpaque(); + if (opaque != null) { + RemotingServer remotingServer = this.brokerController.getRemotingServer(); + if (remotingServer instanceof NettyRemotingAbstract) { + NettyRemotingAbstract nettyRemotingAbstract = (NettyRemotingAbstract) remotingServer; + RemotingCommand remotingCommand = RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "From gRPC client"); + remotingCommand.setOpaque(pollCommandResponseFuture.getOpaque()); + ConsumeMessageDirectlyResult result = GrpcConverterV2.buildConsumeMessageDirectlyResult(request); + remotingCommand.setBody(result.encode()); + nettyRemotingAbstract.processResponseCommand(new SimpleChannelHandlerContext(channelManager.createChannel()), remotingCommand); + } + } + } + } @Override public CompletableFuture notifyClientTermination(Context ctx, @@ -514,6 +511,66 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo return future; } + @Override + public StreamObserver telemetry(Context ctx, StreamObserver responseObserver) { + 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); + Settings settings = clientSettings.getSettings(); + if (settings.hasPublishing()) { + Publishing publishing = settings.getPublishing(); + for (Resource topic : publishing.getTopicsList()) { + String topicName = GrpcConverterV2.wrapResourceWithNamespace(topic); + GrpcClientChannelV2 producerChannel = GrpcClientChannelV2.getChannel(channelManager, topicName, clientId); + producerChannel.setClientObserver(responseObserver); + } + } + if (settings.hasSubscription()) { + Subscription subscription = settings.getSubscription(); + String groupName = GrpcConverterV2.wrapResourceWithNamespace(subscription.getGroup()); + GrpcClientChannelV2 consumerChannel = GrpcClientChannelV2.getChannel(channelManager, groupName, clientId); + consumerChannel.setClientObserver(responseObserver); + } + responseObserver.onNext(TelemetryCommand.newBuilder() + .setClientOverwrittenSettings(ClientOverwrittenSettings.newBuilder() + .setNonce(clientSettings.getNonce()) + .setDirection(Direction.RESPONSE) + .setSettings(settings) + .build()) + .build()); + break; + } + case THREAD_STACK_TRACE: { + reportThreadStackTrace(request.getThreadStackTrace()); + break; + } + case VERIFY_MESSAGE_RESULT: { + reportVerifyMessageResult(request.getVerifyMessageResult()); + break; + } + default: { + throw new IllegalArgumentException("Request type is illegal"); + } + } + } + + @Override + public void onError(Throwable t) { + + } + + @Override + public void onCompleted() { + responseObserver.onCompleted(); + } + }; + } + private class LocalGrpcServiceStartAndShutdown implements StartAndShutdown { @Override public void start() throws Exception { LocalGrpcService.this.scheduledExecutorService.scheduleWithFixedDelay(LocalGrpcService.this::scanAndCleanChannels, 5, 5, TimeUnit.MINUTES);