[ISSUE #3949] Add telemetry

This commit is contained in:
zhouxiang
2022-07-13 11:29:17 +08:00
parent f0edf9964f
commit 62c14bfdb9
6 changed files with 362 additions and 152 deletions
@@ -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<TelemetryCommand> telemetry(StreamObserver<TelemetryCommand> responseObserver) {
return new StreamObserver<TelemetryCommand>() {
@Override
public void onNext(TelemetryCommand request) {
}
@Override
public void onError(Throwable t) {
}
@Override
public void onCompleted() {
}
};
return grpcForwardService.telemetry(Context.current(), responseObserver);
}
}
@@ -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());
@@ -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<StreamObserver<TelemetryCommand>> 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<TelemetryCommand> 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}
* <p>
* Case {@link RequestCode#CHECK_TRANSACTION_STATE}
* @see org.apache.rocketmq.broker.client.net.Broker2Client#checkProducerTransactionState
*/
@Override
public ChannelFuture writeAndFlush(Object msg) {
StreamObserver<TelemetryCommand> 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;
}
}
@@ -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<String, ClientData> CLIENT_DATA = new ConcurrentHashMap<>();
private static final Map<String, ClientSettings> 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<TelemetryCommand> 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<TelemetryCommand> responseStreamObserver;
public ClientData(ClientSettings clientSettings,
StreamObserver<TelemetryCommand> responseStreamObserver) {
this.clientSettings = clientSettings;
this.responseStreamObserver = responseStreamObserver;
}
}
}
@@ -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<AckMessageResponse> ackMessage(Context ctx, AckMessageRequest request);
CompletableFuture<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request);
CompletableFuture<ForwardMessageToDeadLetterQueueResponse> forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request);
CompletableFuture<EndTransactionResponse> endTransaction(Context ctx, EndTransactionRequest request);
@@ -76,4 +76,6 @@ public interface GrpcForwardServiceV2 extends StartAndShutdown {
CompletableFuture<NotifyClientTerminationResponse> notifyClientTermination(Context ctx, NotifyClientTerminationRequest request);
CompletableFuture<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request);
StreamObserver<TelemetryCommand> telemetry(Context ctx, StreamObserver<TelemetryCommand> responseObserver);
}
@@ -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<HeartbeatResponse> 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<HeartbeatResponse> 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<PollCommandResponse> pollCommand(Context ctx, PollCommandRequest request) {
// String clientId = request.getClientId();
// CompletableFuture<PollCommandResponse> 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<ReportThreadStackTraceResponse> 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<ReportMessageConsumptionResultResponse> 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<NotifyClientTerminationResponse> notifyClientTermination(Context ctx,
@@ -514,6 +511,66 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
return future;
}
@Override
public StreamObserver<TelemetryCommand> telemetry(Context ctx, StreamObserver<TelemetryCommand> responseObserver) {
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);
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);