diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/PollResponseManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/TelemetryCommandManager.java similarity index 69% rename from proxy/src/main/java/org/apache/rocketmq/proxy/common/PollResponseManager.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/common/TelemetryCommandManager.java index a6dfd335e3..62c9be63d1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/PollResponseManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/TelemetryCommandManager.java @@ -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 futureTable = new ConcurrentHashMap<>(); +public class TelemetryCommandManager { + private final ConcurrentMap 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); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/PollResponseFuture.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/TelemetryCommandRecord.java similarity index 76% rename from proxy/src/main/java/org/apache/rocketmq/proxy/common/PollResponseFuture.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/common/TelemetryCommandRecord.java index fdf68b2117..adc074113f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/PollResponseFuture.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/TelemetryCommandRecord.java @@ -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() { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/adapter/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/adapter/channel/GrpcClientChannel.java index 5f993ff2b5..99f9ec6601 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/adapter/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v1/adapter/channel/GrpcClientChannel.java @@ -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) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java index 8a945ef7da..2bbf6b84db 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/channel/GrpcClientChannel.java @@ -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) 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 e8caea4203..e1081ddd02 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 @@ -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); 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 730ef5ea0f..7a4c964e40 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 @@ -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() diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java index 996e58a89c..f38a95741e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientService.java @@ -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() 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 9e2f51b79f..0bb5eff536 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 @@ -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 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 diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java index 2c0aa30309..32aacfb3d4 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ForwardClientServiceTest.java @@ -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 diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java index 7760be6458..91dd0245b2 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseTest.java @@ -107,7 +107,6 @@ public class GrpcBaseTest extends BaseConf { .setTopic(Resource.newBuilder() .setName(topic) .build()) - .setEndpoints(endpoints) .build(); }