mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Refector telemetry command
* Add unit test * Rename variable
This commit is contained in:
+8
-8
@@ -21,17 +21,17 @@ import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
public class PollResponseManager {
|
||||
private final ConcurrentMap<String, PollResponseFuture> futureTable = new ConcurrentHashMap<>();
|
||||
public class TelemetryCommandManager {
|
||||
private final ConcurrentMap<String, TelemetryCommandRecord> commandTable = new ConcurrentHashMap<>();
|
||||
private final AtomicLong commandIdGenerator = new AtomicLong(0);
|
||||
|
||||
public String putResponse(int opaque) {
|
||||
String commandId = String.valueOf(commandIdGenerator.incrementAndGet());
|
||||
futureTable.put(commandId, new PollResponseFuture(commandId, opaque));
|
||||
return commandId;
|
||||
public String putCommand(int opaque) {
|
||||
String nonce = String.valueOf(commandIdGenerator.incrementAndGet());
|
||||
commandTable.put(nonce, new TelemetryCommandRecord(nonce, opaque));
|
||||
return nonce;
|
||||
}
|
||||
|
||||
public PollResponseFuture getResponse(String commandId) {
|
||||
return futureTable.get(commandId);
|
||||
public TelemetryCommandRecord getCommand(String commandId) {
|
||||
return commandTable.get(commandId);
|
||||
}
|
||||
}
|
||||
+8
-8
@@ -17,22 +17,22 @@
|
||||
|
||||
package org.apache.rocketmq.proxy.common;
|
||||
|
||||
public class PollResponseFuture {
|
||||
private final String commandId;
|
||||
public class TelemetryCommandRecord {
|
||||
private final String nonce;
|
||||
private final Integer opaque;
|
||||
|
||||
public PollResponseFuture(String commandId, int opaque) {
|
||||
this.commandId = commandId;
|
||||
public TelemetryCommandRecord(String nonce, int opaque) {
|
||||
this.nonce = nonce;
|
||||
this.opaque = opaque;
|
||||
}
|
||||
|
||||
public PollResponseFuture(String commandId) {
|
||||
this.commandId = commandId;
|
||||
public TelemetryCommandRecord(String nonce) {
|
||||
this.nonce = nonce;
|
||||
this.opaque = null;
|
||||
}
|
||||
|
||||
public String getCommandId() {
|
||||
return commandId;
|
||||
public String getNonce() {
|
||||
return nonce;
|
||||
}
|
||||
|
||||
public Integer getOpaque() {
|
||||
+7
-7
@@ -32,7 +32,7 @@ import org.apache.rocketmq.common.protocol.header.GetConsumerRunningInfoRequestH
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.v1.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.common.PollResponseManager;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public class GrpcClientChannel extends SimpleChannel {
|
||||
@@ -40,13 +40,13 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
|
||||
private final String group;
|
||||
private final String clientId;
|
||||
private final PollResponseManager manager;
|
||||
private final TelemetryCommandManager manager;
|
||||
|
||||
private GrpcClientChannel(String group, String clientId, PollResponseManager manager) {
|
||||
private GrpcClientChannel(String group, String clientId, TelemetryCommandManager manager) {
|
||||
this(Context.current(), group, clientId, manager);
|
||||
}
|
||||
|
||||
private GrpcClientChannel(Context ctx, String group, String clientId, PollResponseManager manager) {
|
||||
private GrpcClientChannel(Context ctx, String group, String clientId, TelemetryCommandManager manager) {
|
||||
super(ChannelManager.createSimpleChannelDirectly(ctx));
|
||||
this.group = group;
|
||||
this.clientId = clientId;
|
||||
@@ -61,7 +61,7 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
ChannelManager channelManager,
|
||||
String group,
|
||||
String clientId,
|
||||
PollResponseManager manager
|
||||
TelemetryCommandManager manager
|
||||
) {
|
||||
return create(Context.current(), channelManager, group, clientId, manager);
|
||||
}
|
||||
@@ -71,7 +71,7 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
ChannelManager channelManager,
|
||||
String group,
|
||||
String clientId,
|
||||
PollResponseManager manager
|
||||
TelemetryCommandManager manager
|
||||
) {
|
||||
GrpcClientChannel channel = channelManager.createChannel(
|
||||
buildKey(group, clientId),
|
||||
@@ -135,7 +135,7 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
if (!requestHeader.isJstackEnable()) {
|
||||
break;
|
||||
}
|
||||
String commandId = manager.putResponse(command.getOpaque());
|
||||
String commandId = manager.putCommand(command.getOpaque());
|
||||
future.complete(PollCommandResponse.newBuilder()
|
||||
.setPrintThreadStackTraceCommand(PrintThreadStackTraceCommand.newBuilder()
|
||||
.setCommandId(commandId)
|
||||
|
||||
+6
-6
@@ -32,7 +32,7 @@ import org.apache.rocketmq.common.protocol.header.GetConsumerRunningInfoRequestH
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.common.PollResponseManager;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
|
||||
public class GrpcClientChannel extends SimpleChannel {
|
||||
@@ -40,9 +40,9 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
|
||||
private final String group;
|
||||
private final String clientId;
|
||||
private final PollResponseManager manager;
|
||||
private final TelemetryCommandManager manager;
|
||||
|
||||
private GrpcClientChannel(Context ctx, String group, String clientId, PollResponseManager manager) {
|
||||
private GrpcClientChannel(Context ctx, String group, String clientId, TelemetryCommandManager manager) {
|
||||
super(ChannelManager.createSimpleChannelDirectly(ctx));
|
||||
this.group = group;
|
||||
this.clientId = clientId;
|
||||
@@ -57,7 +57,7 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
ChannelManager channelManager,
|
||||
String group,
|
||||
String clientId,
|
||||
PollResponseManager manager
|
||||
TelemetryCommandManager manager
|
||||
) {
|
||||
return create(Context.current(), channelManager, group, clientId, manager);
|
||||
}
|
||||
@@ -67,7 +67,7 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
ChannelManager channelManager,
|
||||
String group,
|
||||
String clientId,
|
||||
PollResponseManager manager
|
||||
TelemetryCommandManager manager
|
||||
) {
|
||||
GrpcClientChannel channel = channelManager.createChannel(
|
||||
buildKey(group, clientId),
|
||||
@@ -131,7 +131,7 @@ public class GrpcClientChannel extends SimpleChannel {
|
||||
if (!requestHeader.isJstackEnable()) {
|
||||
break;
|
||||
}
|
||||
String nonce = manager.putResponse(command.getOpaque());
|
||||
String nonce = manager.putCommand(command.getOpaque());
|
||||
streamObserver.onNext(TelemetryCommand.newBuilder()
|
||||
.setPrintThreadStackTraceCommand(PrintThreadStackTraceCommand.newBuilder()
|
||||
.setNonce(nonce)
|
||||
|
||||
+3
-3
@@ -57,7 +57,7 @@ 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.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ConsumerService;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.service.cluster.ForwardClientService;
|
||||
@@ -82,13 +82,13 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
private final ForwardClientService clientService;
|
||||
private final PullMessageService pullMessageService;
|
||||
private final TransactionService transactionService;
|
||||
private final PollResponseManager pollCommandResponseManager;
|
||||
private final TelemetryCommandManager pollCommandResponseManager;
|
||||
private final GrpcClientManager grpcClientManager;
|
||||
|
||||
public ClusterGrpcService() {
|
||||
this.channelManager = new ChannelManager();
|
||||
this.grpcClientManager = new GrpcClientManager();
|
||||
this.pollCommandResponseManager = new PollResponseManager();
|
||||
this.pollCommandResponseManager = new TelemetryCommandManager();
|
||||
this.connectorManager = new ConnectorManager(new GrpcTransactionStateChecker());
|
||||
this.consumerService = new ConsumerService(connectorManager, grpcClientManager);
|
||||
this.producerService = new ProducerService(connectorManager);
|
||||
|
||||
+31
-19
@@ -94,8 +94,8 @@ import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.common.DelayPolicy;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter;
|
||||
import org.apache.rocketmq.proxy.channel.InvocationContext;
|
||||
import org.apache.rocketmq.proxy.common.PollResponseFuture;
|
||||
import org.apache.rocketmq.proxy.common.PollResponseManager;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandRecord;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyMode;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel;
|
||||
@@ -121,17 +121,26 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
private final ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(
|
||||
new ThreadFactoryImpl("LocalGrpcServiceScheduledThread"));
|
||||
private final ChannelManager channelManager;
|
||||
private final PollResponseManager pollCommandResponseManager;
|
||||
private final TelemetryCommandManager telemetryCommandManager;
|
||||
private final GrpcClientManager grpcClientManager;
|
||||
private final RouteService routeService;
|
||||
private final DelayPolicy delayPolicy;
|
||||
|
||||
public LocalGrpcService(BrokerController brokerController) {
|
||||
this(brokerController, new TelemetryCommandManager());
|
||||
}
|
||||
|
||||
/**
|
||||
* For unit test
|
||||
* @param brokerController BrokerController works in local mode
|
||||
* @param telemetryCommandManager Used to manage telemetry command
|
||||
*/
|
||||
LocalGrpcService(BrokerController brokerController, TelemetryCommandManager telemetryCommandManager) {
|
||||
this.brokerController = brokerController;
|
||||
this.channelManager = new ChannelManager();
|
||||
// TransactionStateChecker is not used in Local mode.
|
||||
ConnectorManager connectorManager = new ConnectorManager(null);
|
||||
this.pollCommandResponseManager = new PollResponseManager();
|
||||
this.telemetryCommandManager = telemetryCommandManager;
|
||||
this.grpcClientManager = new GrpcClientManager();
|
||||
this.routeService = new RouteService(ProxyMode.LOCAL, connectorManager, grpcClientManager);
|
||||
this.delayPolicy = DelayPolicy.build(brokerController.getMessageStoreConfig().getMessageDelayLevel());
|
||||
@@ -164,7 +173,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
case PRODUCER: {
|
||||
for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) {
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, topicName, clientId, pollCommandResponseManager);
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
|
||||
SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel);
|
||||
|
||||
this.brokerController.getClientManageProcessor()
|
||||
@@ -180,7 +189,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
case PUSH_CONSUMER:
|
||||
case SIMPLE_CONSUMER: {
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, groupName, clientId, pollCommandResponseManager);
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, groupName, clientId, telemetryCommandManager);
|
||||
SimpleChannelHandlerContext simpleChannelHandlerContext = new SimpleChannelHandlerContext(channel);
|
||||
|
||||
RemotingCommand response = this.brokerController.getClientManageProcessor()
|
||||
@@ -421,24 +430,27 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
public void reportThreadStackTrace(ThreadStackTrace request) {
|
||||
String nonce = request.getNonce();
|
||||
String threadStack = request.getThreadStackTrace();
|
||||
PollResponseFuture pollCommandResponseFuture = pollCommandResponseManager.getResponse(nonce);
|
||||
TelemetryCommandRecord pollCommandResponseFuture = telemetryCommandManager.getCommand(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);
|
||||
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());
|
||||
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);
|
||||
TelemetryCommandRecord pollCommandResponseFuture = telemetryCommandManager.getCommand(nonce);
|
||||
if (pollCommandResponseFuture != null) {
|
||||
Integer opaque = pollCommandResponseFuture.getOpaque();
|
||||
if (opaque != null) {
|
||||
@@ -528,14 +540,14 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo
|
||||
Publishing publishing = settings.getPublishing();
|
||||
for (Resource topic : publishing.getTopicsList()) {
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, pollCommandResponseManager);
|
||||
GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
|
||||
producerChannel.setClientObserver(responseObserver);
|
||||
}
|
||||
}
|
||||
if (settings.hasSubscription()) {
|
||||
Subscription subscription = settings.getSubscription();
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(subscription.getGroup());
|
||||
GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, pollCommandResponseManager);
|
||||
GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, telemetryCommandManager);
|
||||
consumerChannel.setClientObserver(responseObserver);
|
||||
}
|
||||
responseObserver.onNext(TelemetryCommand.newBuilder()
|
||||
|
||||
+8
-8
@@ -44,8 +44,8 @@ import org.apache.rocketmq.common.MQVersion;
|
||||
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.common.PollResponseManager;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
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;
|
||||
@@ -62,15 +62,15 @@ public class ForwardClientService extends BaseService {
|
||||
private final ChannelManager channelManager;
|
||||
private final ConsumerManager consumerManager;
|
||||
private final ProducerManager producerManager;
|
||||
private final PollResponseManager pollCommandResponseManager;
|
||||
private final GrpcClientManager grpcClientManager;
|
||||
private final TelemetryCommandManager telemetryCommandManager;
|
||||
|
||||
public ForwardClientService(
|
||||
ConnectorManager connectorManager,
|
||||
ScheduledExecutorService scheduledExecutorService,
|
||||
ChannelManager channelManager,
|
||||
GrpcClientManager grpcClientManager,
|
||||
PollResponseManager pollCommandResponseManager
|
||||
TelemetryCommandManager telemetryCommandManager
|
||||
) {
|
||||
super(connectorManager);
|
||||
scheduledExecutorService.scheduleWithFixedDelay(
|
||||
@@ -80,7 +80,7 @@ public class ForwardClientService extends BaseService {
|
||||
TimeUnit.MILLISECONDS);
|
||||
this.channelManager = channelManager;
|
||||
this.grpcClientManager = grpcClientManager;
|
||||
this.pollCommandResponseManager = pollCommandResponseManager;
|
||||
this.telemetryCommandManager = telemetryCommandManager;
|
||||
|
||||
this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListener() {
|
||||
@Override
|
||||
@@ -108,7 +108,7 @@ public class ForwardClientService extends BaseService {
|
||||
case PRODUCER: {
|
||||
for (Resource topic : clientSettings.getSettings().getPublishing().getTopicsList()) {
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, topicName, clientId, pollCommandResponseManager);
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
|
||||
ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal());
|
||||
// use topic name as producer group
|
||||
producerManager.registerProducer(topicName, clientChannelInfo);
|
||||
@@ -122,7 +122,7 @@ public class ForwardClientService extends BaseService {
|
||||
throw new ProxyException(Code.ILLEGAL_CONSUMER_GROUP, "group cannot be empty for consumer");
|
||||
}
|
||||
String consumerGroup = GrpcConverter.wrapResourceWithNamespace(request.getGroup());
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(ctx, channelManager, consumerGroup, clientId, pollCommandResponseManager);
|
||||
GrpcClientChannel channel = GrpcClientChannel.create(ctx, channelManager, consumerGroup, clientId, telemetryCommandManager);
|
||||
ClientChannelInfo clientChannelInfo = new ClientChannelInfo(channel, clientId, languageCode, MQVersion.Version.V5_0_0.ordinal());
|
||||
|
||||
consumerManager.registerConsumer(
|
||||
@@ -207,14 +207,14 @@ public class ForwardClientService extends BaseService {
|
||||
Publishing publishing = settings.getPublishing();
|
||||
for (Resource topic : publishing.getTopicsList()) {
|
||||
String topicName = GrpcConverter.wrapResourceWithNamespace(topic);
|
||||
GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, pollCommandResponseManager);
|
||||
GrpcClientChannel producerChannel = GrpcClientChannel.create(channelManager, topicName, clientId, telemetryCommandManager);
|
||||
producerChannel.setClientObserver(responseObserver);
|
||||
}
|
||||
}
|
||||
if (settings.hasSubscription()) {
|
||||
Subscription subscription = settings.getSubscription();
|
||||
String groupName = GrpcConverter.wrapResourceWithNamespace(subscription.getGroup());
|
||||
GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, pollCommandResponseManager);
|
||||
GrpcClientChannel consumerChannel = GrpcClientChannel.create(channelManager, groupName, clientId, telemetryCommandManager);
|
||||
consumerChannel.setClientObserver(responseObserver);
|
||||
}
|
||||
responseObserver.onNext(TelemetryCommand.newBuilder()
|
||||
|
||||
+42
-1
@@ -49,6 +49,8 @@ import apache.rocketmq.v2.SendMessageResponse;
|
||||
import apache.rocketmq.v2.Settings;
|
||||
import apache.rocketmq.v2.SystemProperties;
|
||||
import apache.rocketmq.v2.TelemetryCommand;
|
||||
import apache.rocketmq.v2.ThreadStackTrace;
|
||||
import apache.rocketmq.v2.VerifyMessageResult;
|
||||
import com.google.protobuf.Timestamp;
|
||||
import com.google.protobuf.util.Durations;
|
||||
import io.grpc.Context;
|
||||
@@ -79,11 +81,15 @@ import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.PopMessageResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.PullMessageResponseHeader;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandRecord;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.proxy.config.InitConfigAndLoggerTest;
|
||||
import org.apache.rocketmq.proxy.connector.transaction.TransactionId;
|
||||
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.ResponseBuilder;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRemotingServer;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.store.MessageStore;
|
||||
import org.apache.rocketmq.store.config.MessageStoreConfig;
|
||||
@@ -109,6 +115,9 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
@Mock
|
||||
private BrokerController brokerControllerMock;
|
||||
|
||||
@Mock
|
||||
private TelemetryCommandManager telemetryCommandManager;
|
||||
|
||||
private Metadata metadata;
|
||||
|
||||
private StreamObserver<TelemetryCommand> streamObserver;
|
||||
@@ -121,7 +130,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
Mockito.when(brokerControllerMock.getPullMessageProcessor()).thenReturn(pullMessageProcessorMock);
|
||||
Mockito.when(brokerControllerMock.getBrokerConfig()).thenReturn(new BrokerConfig());
|
||||
Mockito.when(brokerControllerMock.getMessageStoreConfig()).thenReturn(new MessageStoreConfig());
|
||||
localGrpcService = new LocalGrpcService(brokerControllerMock);
|
||||
localGrpcService = new LocalGrpcService(brokerControllerMock, telemetryCommandManager);
|
||||
metadata = new Metadata();
|
||||
metadata.put(InterceptorConstants.REMOTE_ADDRESS, "1.1.1.1");
|
||||
metadata.put(InterceptorConstants.LOCAL_ADDRESS, "0.0.0.0");
|
||||
@@ -485,7 +494,39 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
|
||||
|
||||
@Test
|
||||
public void testReportThreadStackTrace() throws Exception {
|
||||
int opaque = 1;
|
||||
String nonce = "123";
|
||||
NettyRemotingServer remotingServerMock = Mockito.mock(NettyRemotingServer.class);
|
||||
Mockito.when(brokerControllerMock.getRemotingServer()).thenReturn(remotingServerMock);
|
||||
Mockito.doNothing().when(remotingServerMock).processResponseCommand(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class));
|
||||
Mockito.when(telemetryCommandManager.getCommand(Mockito.eq(nonce))).thenReturn(new TelemetryCommandRecord(nonce, opaque));
|
||||
String jstack = "jstack";
|
||||
|
||||
streamObserver.onNext(TelemetryCommand.newBuilder()
|
||||
.setThreadStackTrace(ThreadStackTrace.newBuilder()
|
||||
.setNonce(nonce)
|
||||
.setThreadStackTrace(jstack).build())
|
||||
.build());
|
||||
Mockito.verify(remotingServerMock, Mockito.times(1))
|
||||
.processResponseCommand(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testReportVerifyMessageResult() {
|
||||
int opaque = 1;
|
||||
String nonce = "123";
|
||||
NettyRemotingServer remotingServerMock = Mockito.mock(NettyRemotingServer.class);
|
||||
Mockito.when(brokerControllerMock.getRemotingServer()).thenReturn(remotingServerMock);
|
||||
Mockito.doNothing().when(remotingServerMock).processResponseCommand(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class));
|
||||
Mockito.when(telemetryCommandManager.getCommand(Mockito.eq(nonce))).thenReturn(new TelemetryCommandRecord(nonce, opaque));
|
||||
|
||||
streamObserver.onNext(TelemetryCommand.newBuilder()
|
||||
.setVerifyMessageResult(VerifyMessageResult.newBuilder()
|
||||
.setNonce(nonce)
|
||||
.setStatus(ResponseBuilder.buildStatus(Code.OK, "ok")).build())
|
||||
.build());
|
||||
Mockito.verify(remotingServerMock, Mockito.times(1))
|
||||
.processResponseCommand(Mockito.any(ChannelHandlerContext.class), Mockito.any(RemotingCommand.class));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
+3
-3
@@ -22,7 +22,7 @@ import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.common.PollResponseManager;
|
||||
import org.apache.rocketmq.proxy.common.TelemetryCommandManager;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.remoting.protocol.LanguageCode;
|
||||
import org.junit.Test;
|
||||
@@ -37,7 +37,7 @@ import static org.mockito.Mockito.when;
|
||||
public class ForwardClientServiceTest extends BaseServiceTest {
|
||||
|
||||
private ChannelManager channelManager = new ChannelManager();
|
||||
private PollResponseManager pollResponseManager = new PollResponseManager();
|
||||
private TelemetryCommandManager telemetryCommandManager = new TelemetryCommandManager();
|
||||
private ForwardClientService clientService;
|
||||
|
||||
@Override
|
||||
@@ -47,7 +47,7 @@ public class ForwardClientServiceTest extends BaseServiceTest {
|
||||
Executors.newSingleThreadScheduledExecutor(),
|
||||
this.channelManager,
|
||||
this.grpcClientManager,
|
||||
this.pollResponseManager);
|
||||
this.telemetryCommandManager);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -107,7 +107,6 @@ public class GrpcBaseTest extends BaseConf {
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName(topic)
|
||||
.build())
|
||||
.setEndpoints(endpoints)
|
||||
.build();
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user